From c84b03db4455ffd026c8ab863001b92313bb794e Mon Sep 17 00:00:00 2001 From: Daniel Sutton Date: Mon, 24 Aug 2026 09:42:26 +0100 Subject: [PATCH] fix(webapp): hold the active mint-shard list in the database, not the environment A rolling deploy takes hours, so two pods run different values of RUN_OPS_MINT_SHARDS at the same time. The grace window is sized in seconds, so new pods left it long before old pods were gone: for the rest of the rollout the two placed the same environment on different shards. That is the divergence the grace exists to close. The environment variable is now a ceiling that changes only by deploy. It says which shard keys this deployment can mint into. The live list moves to the control-plane database as runOpsMintShardSet, so every pod reads one shared value whatever config generation it is running. Resolution intersects the two, so a stored key this deployment cannot route is never minted into. RUN_OPS_MINT_SHARDS_PREV and RUN_OPS_MINT_SHARDS_FLIPPED_AT are gone. An environment variable cannot record its own flip time, and an operator cannot know a rollout's end in advance. The stamp is now written server-side against the control-plane clock, under an advisory lock, on a genuine change. Stamping generalizes to N graced flag groups in one transaction under one lock, covering the existing mint-kind trio and the new list. That closes a hole on the global admin flags page, which wrote any catalog key with a bare upsert: a graced key could be set with no stamp, or swept away by a save that omitted it. applyGlobalMintKindFlip stays as a thin wrapper so its route and its test keep working unchanged. Operational rule this creates: every change to RUN_OPS_MINT_SHARDS must land across the whole fleet before the flag selects a key it adds. Routing before minting, which is how the shard topology is already gated. --- apps/webapp/app/env.server.ts | 12 +- .../app/routes/admin.api.v1.feature-flags.ts | 11 +- .../webapp/app/routes/admin.feature-flags.tsx | 45 +---- apps/webapp/app/v3/featureFlags.server.ts | 124 ++++++++++-- apps/webapp/app/v3/featureFlags.ts | 24 +++ .../v3/runOpsMigration/mintShardGrace.test.ts | 132 +++++++++--- .../app/v3/runOpsMigration/mintShardGrace.ts | 82 ++++++-- .../runOpsMintShard.server.test.ts | 59 +++++- .../runOpsMigration/runOpsMintShard.server.ts | 90 ++++++--- .../test/runOpsMintShardSetFlip.test.ts | 190 ++++++++++++++++++ 10 files changed, 627 insertions(+), 142 deletions(-) create mode 100644 apps/webapp/test/runOpsMintShardSetFlip.test.ts diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index a7c572cd3..851454acb 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -2016,14 +2016,12 @@ const EnvironmentSchema = z // (stale or fresh) resolves to the same kind for the whole window. See mintFlipGrace.ts. RUN_OPS_MINT_FLIP_GRACE_MS: z.coerce.number().int().default(90_000), - // Gen-2 mint shards — CSV of single-char [a-z0-9] keys eligible for ROOT minting. Unset or - // empty means no gen-2 minting, which is today's behaviour. Validated at boot: an invalid - // key would mint an id that cannot be routed. _PREV + _FLIPPED_AT stamp a set change so - // every process crosses the cutover together; set both, or the grace never applies. - // Removing a key stops new roots on it and never stops routing it. See mintShardGrace.ts. + // Gen-2 mint shards — CSV of single-char [a-z0-9] keys this deployment can mint roots into. + // Unset or empty means no gen-2 minting, which is today's behaviour. Validated at boot: an + // invalid key would mint an id that cannot be routed. This is a CEILING, not the live list: + // it changes only by deploy, and the runOpsMintShardSet flag selects from it at runtime. + // A rolling deploy runs two values of this var at once, so it must never be the ramp lever. RUN_OPS_MINT_SHARDS: shardCsvString(), - RUN_OPS_MINT_SHARDS_PREV: shardCsvString(), - RUN_OPS_MINT_SHARDS_FLIPPED_AT: z.string().datetime().optional(), // Session replication (Postgres → ClickHouse sessions_v1). Shares Redis // with the runs replicator for leader locking but has its own slot and diff --git a/apps/webapp/app/routes/admin.api.v1.feature-flags.ts b/apps/webapp/app/routes/admin.api.v1.feature-flags.ts index 8cd4f7787..ad5d51592 100644 --- a/apps/webapp/app/routes/admin.api.v1.feature-flags.ts +++ b/apps/webapp/app/routes/admin.api.v1.feature-flags.ts @@ -3,7 +3,7 @@ import { json } from "@remix-run/server-runtime"; import { prisma } from "~/db.server"; import { env } from "~/env.server"; import { requireAdminApiRequest } from "~/services/personalAccessToken.server"; -import { applyGlobalMintKindFlip, makeSetMultipleFlags } from "~/v3/featureFlags.server"; +import { applyGlobalGracedFlips, makeSetMultipleFlags } from "~/v3/featureFlags.server"; import { validatePartialFeatureFlags } from "~/v3/featureFlags"; export async function action({ request }: ActionFunctionArgs) { @@ -29,14 +29,15 @@ export async function action({ request }: ActionFunctionArgs) { const { runOpsMintKindPrev: _ignoredPrev, runOpsMintKindFlippedAt: _ignoredFlippedAt, + runOpsMintShardSetPrev: _ignoredSetPrev, + runOpsMintShardSetFlippedAt: _ignoredSetFlippedAt, ...requestedFlags } = validationResult.data; - // A global mint-kind flip stamps its grace window under a lock (applyGlobalMintKindFlip); - // any other flag save writes directly. + // A change to a graced group stamps its window under a lock; any other save writes directly. const updatedFlags = - requestedFlags.runOpsMintKind !== undefined - ? await applyGlobalMintKindFlip(prisma, requestedFlags, env.RUN_OPS_MINT_FLIP_GRACE_MS) + requestedFlags.runOpsMintKind !== undefined || requestedFlags.runOpsMintShardSet !== undefined + ? await applyGlobalGracedFlips(prisma, requestedFlags, env.RUN_OPS_MINT_FLIP_GRACE_MS) : await makeSetMultipleFlags(prisma)(requestedFlags); return json({ diff --git a/apps/webapp/app/routes/admin.feature-flags.tsx b/apps/webapp/app/routes/admin.feature-flags.tsx index be0f61746..0b52d324d 100644 --- a/apps/webapp/app/routes/admin.feature-flags.tsx +++ b/apps/webapp/app/routes/admin.feature-flags.tsx @@ -5,17 +5,18 @@ import { json } from "@remix-run/server-runtime"; import { typedjson, useTypedLoaderData } from "remix-typedjson"; import { z } from "zod"; import { LockClosedIcon } from "@heroicons/react/20/solid"; -import { boundedIn, prisma } from "~/db.server"; +import { prisma } from "~/db.server"; import { env } from "~/env.server"; import { dashboardAction, dashboardLoader } from "~/services/routeBuilders/dashboardBuilder"; import { FEATURE_FLAG, GLOBAL_LOCKED_FLAGS, + type FeatureFlagKey, type FlagControlType, getAllFlagControlTypes, validatePartialFeatureFlags, } from "~/v3/featureFlags"; -import { flags as getGlobalFlags } from "~/v3/featureFlags.server"; +import { flags as getGlobalFlags, replaceGlobalFeatureFlags } from "~/v3/featureFlags.server"; import { featuresForRequest } from "~/features.server"; import { Button } from "~/components/primitives/Buttons"; import { Callout } from "~/components/primitives/Callout"; @@ -116,39 +117,15 @@ export const action = dashboardAction( ); } - const validatedFlags = validationResult.data as Record; - const controlTypes = getAllFlagControlTypes(); - const catalogKeys = Object.keys(controlTypes); + const catalogKeys = Object.keys(getAllFlagControlTypes()) as FeatureFlagKey[]; - const keysToDelete: string[] = []; - const upsertOps: ReturnType[] = []; - - for (const key of catalogKeys) { - if (key in validatedFlags) { - upsertOps.push( - prisma.featureFlag.upsert({ - where: { key }, - create: { key, value: validatedFlags[key] as any }, - update: { value: validatedFlags[key] as any }, - }) - ); - } else { - // On cloud, never delete locked flags (they're not in the payload - // because the UI doesn't include them). Locally, delete everything - // the user didn't include - full control. - const isProtected = isManagedCloud && GLOBAL_LOCKED_FLAGS.includes(key); - if (!isProtected) { - keysToDelete.push(key); - } - } - } - - await prisma.$transaction([ - ...upsertOps, - ...(keysToDelete.length > 0 - ? [prisma.featureFlag.deleteMany({ where: { key: { in: boundedIn(keysToDelete) } } })] - : []), - ]); + await replaceGlobalFeatureFlags(prisma, { + requestedFlags: validationResult.data, + catalogKeys, + // On cloud, never delete locked flags (the UI omits them). Locally, full control. + isProtected: (key) => isManagedCloud && GLOBAL_LOCKED_FLAGS.includes(key), + graceMs: env.RUN_OPS_MINT_FLIP_GRACE_MS, + }); return json({ success: true }); } diff --git a/apps/webapp/app/v3/featureFlags.server.ts b/apps/webapp/app/v3/featureFlags.server.ts index b32a45786..e30b5df4c 100644 --- a/apps/webapp/app/v3/featureFlags.server.ts +++ b/apps/webapp/app/v3/featureFlags.server.ts @@ -8,6 +8,8 @@ import { FeatureFlagCatalog, } from "~/v3/featureFlags"; import { stampMintKindFlip } from "~/v3/runOpsMigration/mintFlipGrace"; +import { stampMintShardSetFlip } from "~/v3/runOpsMigration/mintShardGrace"; +import { boundedIn } from "~/db.server"; export type FlagsOptions = { key: T; @@ -182,24 +184,41 @@ export function makeSetMultipleFlags(_prisma: PrismaClientOrTransaction = prisma // Read -> stamp -> write the global mint-kind grace metadata in one transaction. The three // FeatureFlag rows may not exist yet, so a row FOR UPDATE can't lock them; an advisory xact lock // serializes concurrent global flips so one can't clobber another's grace stamp (mirrors per-org). -export async function applyGlobalMintKindFlip( +// Every group of global flags whose value carries its own grace stamp. One transaction and one +// lock cover all of them, so a save that flips two groups can never stamp one and lose the other. +const GRACED_GLOBAL_GROUPS = [ + { + keys: [ + FEATURE_FLAG.runOpsMintKind, + FEATURE_FLAG.runOpsMintKindPrev, + FEATURE_FLAG.runOpsMintKindFlippedAt, + ] as FeatureFlagKey[], + stamp: stampMintKindFlip, + }, + { + keys: [ + FEATURE_FLAG.runOpsMintShardSet, + FEATURE_FLAG.runOpsMintShardSetPrev, + FEATURE_FLAG.runOpsMintShardSetFlippedAt, + ] as FeatureFlagKey[], + stamp: stampMintShardSetFlip, + }, +] as const; + +// Keys the graced path owns. They never take a bare upsert and never enter the replace sweep, +// because a server-computed stamp must not be written from a request body nor swept away. +const GRACED_GLOBAL_KEYS: FeatureFlagKey[] = GRACED_GLOBAL_GROUPS.flatMap((g) => g.keys); + +export async function applyGlobalGracedFlips( client: PrismaClient, requestedFlags: Partial>, graceMs: number ): Promise<{ key: string; value: any }[]> { return client.$transaction(async (tx) => { - await tx.$executeRaw`SELECT pg_advisory_xact_lock(hashtext('runops-global-mint-kind-flip'))`; + await tx.$executeRaw`SELECT pg_advisory_xact_lock(hashtext('runops-global-graced-flag-flip'))`; const existingRows = await tx.featureFlag.findMany({ - where: { - key: { - in: [ - FEATURE_FLAG.runOpsMintKind, - FEATURE_FLAG.runOpsMintKindPrev, - FEATURE_FLAG.runOpsMintKindFlippedAt, - ], - }, - }, + where: { key: { in: GRACED_GLOBAL_KEYS } }, select: { key: true, value: true }, }); const existingGlobal: Record = {}; @@ -207,16 +226,83 @@ export async function applyGlobalMintKindFlip( existingGlobal[row.key] = row.value; } - // Anchor the cutover to the control-plane DB clock, not this process's wall clock. + // Anchor the cutover to the control-plane DB clock, not this process's wall clock. A rolling + // deploy spans hours, so every pod must date the window against one shared clock. const [{ now }] = await tx.$queryRaw<{ now: Date }[]>`SELECT now() AS now`; - const stamped = stampMintKindFlip( - existingGlobal, - { ...requestedFlags }, - now.getTime(), - graceMs - ) as Partial>; + let stamped: Record = { ...requestedFlags }; + for (const group of GRACED_GLOBAL_GROUPS) { + stamped = group.stamp(existingGlobal, stamped, now.getTime(), graceMs); + } - return makeSetMultipleFlags(tx)(stamped); + return makeSetMultipleFlags(tx)(stamped as Partial>); }); } + +/** @deprecated Prefer applyGlobalGracedFlips, which stamps every graced group in one lock. */ +export async function applyGlobalMintKindFlip( + client: PrismaClient, + requestedFlags: Partial>, + graceMs: number +): Promise<{ key: string; value: any }[]> { + return applyGlobalGracedFlips(client, requestedFlags, graceMs); +} + +// Replace-semantics write for the global admin flags page: upsert submitted catalog flags, delete +// omitted ones unless protected, and route any graced group through the stamped path. +export async function replaceGlobalFeatureFlags( + client: PrismaClient, + params: { + requestedFlags: Partial>; + catalogKeys: FeatureFlagKey[]; + isProtected: (key: FeatureFlagKey) => boolean; + graceMs: number; + } +): Promise { + // Derived stamp fields are computed server-side; never trust them from the body. + const requestedFlags: Record = { ...params.requestedFlags }; + for (const group of GRACED_GLOBAL_GROUPS) { + for (const derived of group.keys.slice(1)) { + delete requestedFlags[derived]; + } + } + + const touchesGracedGroup = GRACED_GLOBAL_GROUPS.some( + (group) => requestedFlags[group.keys[0]] !== undefined + ); + if (touchesGracedGroup) { + await applyGlobalGracedFlips( + client, + requestedFlags as Partial>, + params.graceMs + ); + } + + const upsertOps: ReturnType[] = []; + const keysToDelete: string[] = []; + + for (const key of params.catalogKeys) { + if (GRACED_GLOBAL_KEYS.includes(key)) { + continue; + } + if (key in requestedFlags) { + const value = requestedFlags[key]; + upsertOps.push( + client.featureFlag.upsert({ + where: { key }, + create: { key, value: value as any }, + update: { value: value as any }, + }) + ); + } else if (!params.isProtected(key)) { + keysToDelete.push(key); + } + } + + await client.$transaction([ + ...upsertOps, + ...(keysToDelete.length > 0 + ? [client.featureFlag.deleteMany({ where: { key: { in: boundedIn(keysToDelete) } } })] + : []), + ]); +} diff --git a/apps/webapp/app/v3/featureFlags.ts b/apps/webapp/app/v3/featureFlags.ts index 740fccebb..4dcf7ab77 100644 --- a/apps/webapp/app/v3/featureFlags.ts +++ b/apps/webapp/app/v3/featureFlags.ts @@ -29,6 +29,11 @@ export const FEATURE_FLAG = { // Gen-2 mint shard pins, read from the org override blob only. See runOpsMintShard.server.ts. runOpsMintShard: "runOpsMintShard", runOpsMintShardEnvPins: "runOpsMintShardEnvPins", + // The active mint-shard list, global only. Lives here rather than in the environment because a + // rolling deploy runs two environment values at once for hours. See mintShardGrace.ts. + runOpsMintShardSet: "runOpsMintShardSet", + runOpsMintShardSetPrev: "runOpsMintShardSetPrev", + runOpsMintShardSetFlippedAt: "runOpsMintShardSetFlippedAt", queueMetricsUiEnabled: "queueMetricsUiEnabled", // Per-organization rollout for creating additional environment API keys. additionalApiKeysEnabled: "additionalApiKeysEnabled", @@ -118,6 +123,21 @@ export const FeatureFlagCatalog = { } } }), + // CSV of the shard keys eligible for root minting right now, bounded by RUN_OPS_MINT_SHARDS. + // Empty means no gen-2 minting. Reserved keys are rejected: "new" already means gen-1. + [FEATURE_FLAG.runOpsMintShardSet]: z.string().refine( + (v) => + v + .split(",") + .map((s) => s.trim()) + .filter(Boolean) + .every((k) => /^[a-z0-9]$/.test(k)), + "must be a CSV of single [a-z0-9] chars" + ), + // Grace stamp: the previously-effective list and the flip time, written by + // stampMintShardSetFlip on a genuine change. Display-only (see ORG_LOCKED_FLAGS). + [FEATURE_FLAG.runOpsMintShardSetPrev]: z.string(), + [FEATURE_FLAG.runOpsMintShardSetFlippedAt]: z.string().datetime(), // Per-org access to the Queue Metrics dashboard UI (view only; emission is global and // separate). Off unless enabled for the org. [FEATURE_FLAG.queueMetricsUiEnabled]: z.coerce.boolean(), @@ -151,6 +171,10 @@ export const ORG_LOCKED_FLAGS: FeatureFlagKey[] = [ // System-wide only — orgs must not be able to override these kill switches. FEATURE_FLAG.additionalApiKeyIssuanceEnabled, FEATURE_FLAG.additionalApiKeyLookupEnabled, + // The active mint-shard list is deployment-wide; only the pins are per-org. + FEATURE_FLAG.runOpsMintShardSet, + FEATURE_FLAG.runOpsMintShardSetPrev, + FEATURE_FLAG.runOpsMintShardSetFlippedAt, ]; // Create a Zod schema from the existing catalog diff --git a/apps/webapp/app/v3/runOpsMigration/mintShardGrace.test.ts b/apps/webapp/app/v3/runOpsMigration/mintShardGrace.test.ts index 9358685e1..3d26cb9fd 100644 --- a/apps/webapp/app/v3/runOpsMigration/mintShardGrace.test.ts +++ b/apps/webapp/app/v3/runOpsMigration/mintShardGrace.test.ts @@ -1,11 +1,12 @@ import { describe, expect, it } from "vitest"; import { generateRunOpsIdV2 } from "@trigger.dev/core/v3/isomorphic"; import { - buildMintShardResolution, effectiveMintShardSet, isValidPinValue, parseShardCsv, + readMintShardSetResolution, SHARD_KEY_PATTERN, + stampMintShardSetFlip, type MintShardSetResolution, } from "./mintShardGrace"; @@ -116,33 +117,118 @@ describe("effectiveMintShardSet", () => { }); }); -describe("buildMintShardResolution", () => { - it("omits prevSet entirely when no flip timestamp is configured", () => { - // A prevSet with no timestamp can never apply, so it MUST NOT linger. - const r = buildMintShardResolution({ shards: "a,b", prev: "a", flippedAt: undefined }); - expect(r).toEqual({ set: ["a", "b"], prevSet: undefined, flippedAtMs: undefined }); +describe("readMintShardSetResolution", () => { + it("returns an empty set for an absent record", () => { + expect(readMintShardSetResolution(undefined)).toEqual({ set: [] }); + expect(readMintShardSetResolution({})).toEqual({ set: [] }); }); - it("keeps an empty prevSet when a flip timestamp IS configured", () => { - const r = buildMintShardResolution({ - shards: "a", - prev: "", - flippedAt: new Date(T).toISOString(), - }); - expect(r).toEqual({ set: ["a"], prevSet: [], flippedAtMs: T }); - }); - - it("parses the flip timestamp and sorts both lists", () => { - const r = buildMintShardResolution({ - shards: "b,a", - prev: "c,a", - flippedAt: new Date(T).toISOString(), + it("reads and sorts the trio", () => { + const r = readMintShardSetResolution({ + runOpsMintShardSet: "b,a", + runOpsMintShardSetPrev: "c,a", + runOpsMintShardSetFlippedAt: new Date(T).toISOString(), }); expect(r).toEqual({ set: ["a", "b"], prevSet: ["a", "c"], flippedAtMs: T }); }); - it("treats an unparseable timestamp as no stamp at all", () => { - const r = buildMintShardResolution({ shards: "a", prev: "b", flippedAt: "not-a-date" }); - expect(r).toEqual({ set: ["a"], prevSet: undefined, flippedAtMs: undefined }); + it("omits prevSet when no flip timestamp is stored", () => { + // A prevSet with no timestamp can never apply, so it MUST NOT linger. + const r = readMintShardSetResolution({ + runOpsMintShardSet: "a,b", + runOpsMintShardSetPrev: "a", + }); + expect(r).toEqual({ set: ["a", "b"], prevSet: undefined, flippedAtMs: undefined }); + }); + + it("keeps an empty prevSet when a timestamp IS stored, which graces a first activation", () => { + const r = readMintShardSetResolution({ + runOpsMintShardSet: "a", + runOpsMintShardSetPrev: "", + runOpsMintShardSetFlippedAt: new Date(T).toISOString(), + }); + expect(r).toEqual({ set: ["a"], prevSet: [], flippedAtMs: T }); + }); + + it("degrades a stored value it cannot parse to an empty list instead of throwing", () => { + // Boot may throw on a bad env var. The mint path must never throw on a bad stored value. + expect(() => readMintShardSetResolution({ runOpsMintShardSet: "NOPE" })).not.toThrow(); + expect(readMintShardSetResolution({ runOpsMintShardSet: "NOPE" }).set).toEqual([]); + expect(readMintShardSetResolution({ runOpsMintShardSet: 42 }).set).toEqual([]); + expect( + readMintShardSetResolution({ + runOpsMintShardSet: "a", + runOpsMintShardSetFlippedAt: "not-a-date", + }) + ).toEqual({ set: ["a"], prevSet: undefined, flippedAtMs: undefined }); + }); +}); + +describe("stampMintShardSetFlip", () => { + it("does nothing when the save omits the set", () => { + // Omitting the set is an unrelated flag change; it must not inject a default or reset the clock. + const outgoing = { someOtherFlag: true } as Record; + expect(stampMintShardSetFlip({ runOpsMintShardSet: "a" }, outgoing, T, GRACE_MS)).toEqual({ + someOtherFlag: true, + }); + }); + + it("stamps prev and flippedAt on a genuine change", () => { + const stamped = stampMintShardSetFlip( + { runOpsMintShardSet: "a" }, + { runOpsMintShardSet: "a,b" }, + T, + GRACE_MS + ); + expect(stamped.runOpsMintShardSetPrev).toBe("a"); + expect(stamped.runOpsMintShardSetFlippedAt).toBe(new Date(T).toISOString()); + }); + + it("stamps an empty prev on a first activation", () => { + const stamped = stampMintShardSetFlip({}, { runOpsMintShardSet: "a" }, T, GRACE_MS); + expect(stamped.runOpsMintShardSetPrev).toBe(""); + expect(stamped.runOpsMintShardSetFlippedAt).toBe(new Date(T).toISOString()); + }); + + it("treats a reordered list as no change", () => { + const stamped = stampMintShardSetFlip( + { runOpsMintShardSet: "a,b" }, + { runOpsMintShardSet: "b,a" }, + T, + GRACE_MS + ); + expect(stamped.runOpsMintShardSetFlippedAt).toBeUndefined(); + }); + + it("carries an in-flight stamp forward rather than resetting the cutover clock", () => { + const existing = { + runOpsMintShardSet: "a,b", + runOpsMintShardSetPrev: "a", + runOpsMintShardSetFlippedAt: new Date(T).toISOString(), + }; + const stamped = stampMintShardSetFlip( + existing, + { runOpsMintShardSet: "a,b" }, + T + 1000, + GRACE_MS + ); + expect(stamped.runOpsMintShardSetPrev).toBe("a"); + expect(stamped.runOpsMintShardSetFlippedAt).toBe(new Date(T).toISOString()); + }); + + it("stamps prev as the CURRENTLY-EFFECTIVE set when a second flip lands mid-window", () => { + // Two flips inside one window must not strand the original prev; prev is what readers serve now. + const existing = { + runOpsMintShardSet: "a,b", + runOpsMintShardSetPrev: "a", + runOpsMintShardSetFlippedAt: new Date(T).toISOString(), + }; + const stamped = stampMintShardSetFlip( + existing, + { runOpsMintShardSet: "a,b,c" }, + T + 1000, + GRACE_MS + ); + expect(stamped.runOpsMintShardSetPrev).toBe("a"); }); }); diff --git a/apps/webapp/app/v3/runOpsMigration/mintShardGrace.ts b/apps/webapp/app/v3/runOpsMigration/mintShardGrace.ts index bfa3139cc..65d8af6f4 100644 --- a/apps/webapp/app/v3/runOpsMigration/mintShardGrace.ts +++ b/apps/webapp/app/v3/runOpsMigration/mintShardGrace.ts @@ -19,6 +19,12 @@ export type MintShardSetResolution = { flippedAtMs?: number; }; +// Flag keys holding the active set and its grace stamp. Named here so the pure module can read +// a flag record without importing the catalog. +const SET_KEY = "runOpsMintShardSet"; +const SET_PREV_KEY = "runOpsMintShardSetPrev"; +const SET_FLIPPED_AT_KEY = "runOpsMintShardSetFlippedAt"; + export function isValidPinValue(value: unknown): value is ShardKey { if (typeof value !== "string") return false; return value === GEN_1_PIN_VALUE || SHARD_KEY_PATTERN.test(value); @@ -61,35 +67,67 @@ export function effectiveMintShardSet( return nowMs < r.flippedAtMs + graceMs ? r.prevSet : r.set; } -// A prevSet with no timestamp can never apply, so it is dropped. A timestamp with an EMPTY -// prevSet is meaningful: it graces a first activation, serving no shards for the window. -export function buildMintShardResolution(source: { - shards: string | undefined; - prev: string | undefined; - flippedAt: string | undefined; -}): MintShardSetResolution { - const parsed = source.flippedAt !== undefined ? Date.parse(source.flippedAt) : NaN; +// The active set lives in the control-plane database, not in the environment. A deploy rolls +// for hours, so two pods can hold different environment values at the same time; only a shared +// row lets every pod agree on one set. Boot may reject a bad environment value, but the mint +// path must never throw on a bad stored value, so an unreadable list degrades to empty. +function readStoredCsv(value: unknown): string[] { + if (typeof value !== "string") return []; + try { + return parseShardCsv(value); + } catch { + return []; + } +} + +// Reads the { set, prevSet, flippedAtMs } trio out of one flag record. Pure. A prevSet with no +// timestamp can never apply, so it is dropped. A timestamp with an EMPTY prevSet is meaningful: +// it graces a first activation, serving no shards for the window. +export function readMintShardSetResolution( + flags: Record | null | undefined +): MintShardSetResolution { + const source = flags ?? {}; + const flippedAtRaw = source[SET_FLIPPED_AT_KEY]; + const parsed = typeof flippedAtRaw === "string" ? Date.parse(flippedAtRaw) : NaN; const flippedAtMs = Number.isNaN(parsed) ? undefined : parsed; return { - set: parseShardCsv(source.shards), - prevSet: flippedAtMs === undefined ? undefined : parseShardCsv(source.prev), + set: readStoredCsv(source[SET_KEY]), + prevSet: flippedAtMs === undefined ? undefined : readStoredCsv(source[SET_PREV_KEY]), flippedAtMs, }; } -// Returns a message to log when the stamp is half-configured, otherwise undefined. Stays quiet -// while the active set is empty, so an unconfigured deployment logs nothing at boot. -export function mintShardStampWarning(source: { - shards: string | undefined; - prev: string | undefined; - flippedAt: string | undefined; -}): string | undefined { - if (parseShardCsv(source.shards).length === 0) { - return undefined; +// Stamps a grace window only when the outgoing set differs from the stored one. prev becomes the +// set readers serve right now, so a second flip inside one window cannot strand the first. A save +// that leaves the set unchanged carries any in-flight stamp forward, so it cannot reset the clock. +export function stampMintShardSetFlip( + existingFlags: Record | null | undefined, + outgoingFlags: Record, + nowMs: number, + graceMs: number +): Record { + // Only act when the save actually SETS the list. Omitting it must not inject a default. + if (typeof outgoingFlags[SET_KEY] !== "string") { + return outgoingFlags; } - if (parseShardCsv(source.prev).length > 0 && source.flippedAt === undefined) { - return "RUN_OPS_MINT_SHARDS_PREV is set but RUN_OPS_MINT_SHARDS_FLIPPED_AT is not; the shard-set grace window will never apply"; + + const existing = existingFlags ?? {}; + const outgoingSet = readStoredCsv(outgoingFlags[SET_KEY]); + const storedSet = readStoredCsv(existing[SET_KEY]); + + if (outgoingSet.join(",") !== storedSet.join(",")) { + const effective = effectiveMintShardSet(readMintShardSetResolution(existing), nowMs, graceMs); + outgoingFlags[SET_PREV_KEY] = effective.join(","); + outgoingFlags[SET_FLIPPED_AT_KEY] = new Date(nowMs).toISOString(); + return outgoingFlags; } - return undefined; + + if (existing[SET_PREV_KEY] !== undefined) { + outgoingFlags[SET_PREV_KEY] = existing[SET_PREV_KEY]; + } + if (existing[SET_FLIPPED_AT_KEY] !== undefined) { + outgoingFlags[SET_FLIPPED_AT_KEY] = existing[SET_FLIPPED_AT_KEY]; + } + return outgoingFlags; } diff --git a/apps/webapp/app/v3/runOpsMigration/runOpsMintShard.server.test.ts b/apps/webapp/app/v3/runOpsMigration/runOpsMintShard.server.test.ts index 806a0c1ae..ff8c178ee 100644 --- a/apps/webapp/app/v3/runOpsMigration/runOpsMintShard.server.test.ts +++ b/apps/webapp/app/v3/runOpsMigration/runOpsMintShard.server.test.ts @@ -15,12 +15,15 @@ function envIds(count: number): string[] { return ids; } +const ALL_KEYS = "abcdefghijklmnopqrstuvwxyz0123456789".split(""); + function deps( resolution: MintShardSetResolution, overrides: Partial = {} ): MintShardDeps { return { resolution, + ceiling: ALL_KEYS, nowMs: T + GRACE_MS + 1, graceMs: GRACE_MS, orgFeatureFlags: undefined, @@ -41,14 +44,27 @@ function place(ids: string[], resolution: MintShardSetResolution): Map { - it("returns new when the active set is unset", () => { + it("returns new when the live list is empty", () => { expect(computeMintShard({ id: "env_1" }, deps({ set: [] }))).toBe("new"); }); - it("returns new when the active set is empty even with a stale stamp present", () => { + it("returns new when the deployment configures no ceiling", () => { + // An unconfigured deployment is an unconditional kill switch, whatever the stored list says. + const resolution: MintShardSetResolution = { set: ["a", "b"] }; + expect(computeMintShard({ id: "env_1" }, deps(resolution, { ceiling: [] }))).toBe("new"); + }); + + it("returns new when the stored list names nothing this deployment can route", () => { + const resolution: MintShardSetResolution = { set: ["z"] }; + expect(computeMintShard({ id: "env_1" }, deps(resolution, { ceiling: ["a"] }))).toBe("new"); + }); + + it("returns new when the ceiling is empty even with a stale stamp present", () => { const resolution: MintShardSetResolution = { set: [], prevSet: ["a"], flippedAtMs: T }; - // The empty check MUST run before the grace, so an unset set is an unconditional kill switch. - expect(computeMintShard({ id: "env_1" }, deps(resolution, { nowMs: T + 1 }))).toBe("new"); + // The ceiling gate MUST run before the grace, so no stored value can reopen a closed switch. + expect(computeMintShard({ id: "env_1" }, deps(resolution, { nowMs: T + 1, ceiling: [] }))).toBe( + "new" + ); }); it("returns new when the grace serves an empty prevSet", () => { @@ -246,3 +262,38 @@ describe("computeMintShard — rendezvous properties", () => { expect(computeMintShard({ id: "env_1" }, deps({ set: ["a", "b"] }, pinnedToC))).not.toBe("c"); }); }); + +describe("computeMintShard — the ceiling bounds the stored list", () => { + it("mints only into keys the deployment can route", () => { + const resolution: MintShardSetResolution = { set: ["a", "b", "c"] }; + const ids = envIds(300); + for (const id of ids) { + const shard = computeMintShard({ id }, deps(resolution, { ceiling: ["a", "b"] })); + expect(["a", "b"]).toContain(shard); + } + }); + + it("ignores a pin to a key outside the ceiling", () => { + const resolution: MintShardSetResolution = { set: ["a", "c"] }; + const rejected: string[] = []; + const shard = computeMintShard( + { id: "env_1" }, + deps(resolution, { + ceiling: ["a"], + orgFeatureFlags: { runOpsMintShard: "c" }, + onPinRejected: (info) => rejected.push(info.pin), + }) + ); + expect(shard).toBe("a"); + expect(rejected).toEqual(["c"]); + }); + + it("still honours a gen-1 pin when the ceiling is narrower than the stored list", () => { + const resolution: MintShardSetResolution = { set: ["a", "b"] }; + const shard = computeMintShard( + { id: "env_1" }, + deps(resolution, { ceiling: ["a"], orgFeatureFlags: { runOpsMintShard: "new" } }) + ); + expect(shard).toBe("new"); + }); +}); diff --git a/apps/webapp/app/v3/runOpsMigration/runOpsMintShard.server.ts b/apps/webapp/app/v3/runOpsMigration/runOpsMintShard.server.ts index 20f20fcf9..8080acce0 100644 --- a/apps/webapp/app/v3/runOpsMigration/runOpsMintShard.server.ts +++ b/apps/webapp/app/v3/runOpsMigration/runOpsMintShard.server.ts @@ -1,19 +1,23 @@ import { createHash } from "node:crypto"; import type { ShardKey } from "@trigger.dev/core/v3/isomorphic"; +import { $replica } from "~/db.server"; import { env } from "~/env.server"; import { logger } from "~/services/logger.server"; import { FEATURE_FLAG } from "~/v3/featureFlags"; import { - buildMintShardResolution, effectiveMintShardSet, GEN_1_PIN_VALUE, isValidPinValue, - mintShardStampWarning, + parseShardCsv, + readMintShardSetResolution, type MintShardSetResolution, } from "./mintShardGrace"; export type MintShardDeps = { + // The live list, from the control-plane database. resolution: MintShardSetResolution; + // The keys this deployment can route, from the environment. Bounds the live list. + ceiling: string[]; nowMs: number; graceMs: number; orgFeatureFlags: unknown; @@ -83,16 +87,20 @@ function hrwSelect(environmentId: string, activeSet: string[]): string { // PURE CORE — no env, no clock, no I/O; tests drive this directly. Deterministic for fixed // deps, which is what lets run minting and token minting agree on one answer. // -// The empty-set check runs BEFORE the grace, so an unset active list is an unconditional kill -// switch that a stale stamp cannot reopen. A pin outside the active set falls through to the -// hash rather than throwing: honouring it would leak the drain the active list performs, and -// throwing would fail customer triggers whenever a pinned shard drains. +// The ceiling gate runs BEFORE the grace, so an unconfigured deployment is an unconditional +// kill switch that no stored value can reopen. The live list is then intersected with the +// ceiling, so a stored key this deployment cannot route is never minted into. +// +// A pin outside the active set falls through to the hash rather than throwing: honouring it +// would leak the drain the active list performs, and throwing would fail customer triggers +// whenever a pinned shard drains. export function computeMintShard(environment: { id: string }, deps: MintShardDeps): ShardKey { - if (deps.resolution.set.length === 0) { + if (deps.ceiling.length === 0) { return "new"; } - const activeSet = effectiveMintShardSet(deps.resolution, deps.nowMs, deps.graceMs); + const live = effectiveMintShardSet(deps.resolution, deps.nowMs, deps.graceMs); + const activeSet = live.filter((key) => deps.ceiling.includes(key)); if (activeSet.length === 0) { return "new"; } @@ -111,24 +119,38 @@ export function computeMintShard(environment: { id: string }, deps: MintShardDep return hrwSelect(environment.id, activeSet); } -// ENV-BOUND wrapper — the only place env is read. The resolution is built once: these are -// deploy-time values, so re-parsing per mint would burn CPU on the hottest path in the system. -const shardResolution: MintShardSetResolution = buildMintShardResolution({ - shards: env.RUN_OPS_MINT_SHARDS, - prev: env.RUN_OPS_MINT_SHARDS_PREV, - flippedAt: env.RUN_OPS_MINT_SHARDS_FLIPPED_AT, -}); +// ENV-BOUND wrapper — the only place env is read. The ceiling is parsed once at boot; it is a +// deploy-time value, so re-parsing per mint would burn CPU on the hottest path in the system. +const ceiling: string[] = parseShardCsv(env.RUN_OPS_MINT_SHARDS); -const stampWarning = mintShardStampWarning({ - shards: env.RUN_OPS_MINT_SHARDS, - prev: env.RUN_OPS_MINT_SHARDS_PREV, - flippedAt: env.RUN_OPS_MINT_SHARDS_FLIPPED_AT, -}); -if (stampWarning) { - logger.warn(`[runOpsMintShard] ${stampWarning}`, { - RUN_OPS_MINT_SHARDS: env.RUN_OPS_MINT_SHARDS, - RUN_OPS_MINT_SHARDS_PREV: env.RUN_OPS_MINT_SHARDS_PREV, +// The live list is org-independent, so one process-wide entry serves every mint. One query per +// process per TTL, folded into a single round-trip over the three keys. The TTL bounds how long +// two processes can disagree, which is what the grace window is sized against. +const SET_KEYS = [ + FEATURE_FLAG.runOpsMintShardSet, + FEATURE_FLAG.runOpsMintShardSetPrev, + FEATURE_FLAG.runOpsMintShardSetFlippedAt, +]; + +let cachedResolution: { value: MintShardSetResolution; expiresAt: number } | undefined; + +async function readLiveResolution(): Promise { + if (cachedResolution && cachedResolution.expiresAt > Date.now()) { + return cachedResolution.value; + } + + const rows = await $replica.featureFlag.findMany({ + where: { key: { in: SET_KEYS } }, + select: { key: true, value: true }, }); + const flags: Record = {}; + for (const row of rows) { + flags[row.key] = row.value; + } + + const value = readMintShardSetResolution(flags); + cachedResolution = { value, expiresAt: Date.now() + env.RUN_OPS_MINT_FLAG_CACHE_TTL_MS }; + return value; } // Once per environment per process: a stale pin sits on the root-trigger path and would @@ -149,9 +171,6 @@ function reportPinRejected(info: { * Which shard an environment mints new roots into. Call only after resolveRunIdMintKind has * returned "runOpsId". Returns "new" to mean a gen-1 run-ops id, which is today's behaviour. * - * Async despite doing no I/O, so the deploy-free active-set layer can add a read later without - * changing every call site. - * * @knipignore the gen-2 write-path change is the first production caller; drop this tag there. */ export async function resolveMintShard(environment: { @@ -159,8 +178,23 @@ export async function resolveMintShard(environment: { // Pass environment.organization.featureFlags from the trigger call site. orgFeatureFlags?: unknown; }): Promise { + // No ceiling means no gen-2 minting, so skip the read entirely. + if (ceiling.length === 0) { + return "new"; + } + + let resolution: MintShardSetResolution; + try { + resolution = await readLiveResolution(); + } catch (error) { + // Fail safe to gen-1, mirroring the mint-kind gate's fail-safe to cuid. + logger.error("[runOpsMintShard] shard-set read failed; minting gen-1 (fail-safe)", { error }); + return "new"; + } + return computeMintShard(environment, { - resolution: shardResolution, + resolution, + ceiling, nowMs: Date.now(), graceMs: env.RUN_OPS_MINT_FLIP_GRACE_MS, orgFeatureFlags: environment.orgFeatureFlags, diff --git a/apps/webapp/test/runOpsMintShardSetFlip.test.ts b/apps/webapp/test/runOpsMintShardSetFlip.test.ts new file mode 100644 index 000000000..55ebe43a2 --- /dev/null +++ b/apps/webapp/test/runOpsMintShardSetFlip.test.ts @@ -0,0 +1,190 @@ +// The active mint-shard list lives in the control-plane database, not in the environment: a +// rolling deploy runs two environment values at once for hours, so only a shared row lets every +// pod agree on one list. A change must therefore read -> stamp -> write under an advisory lock, +// and must never be writable as a bare upsert from a request body. Real testcontainers Postgres. +import type { PrismaClient } from "@trigger.dev/database"; +import { postgresTest } from "@internal/testcontainers"; +import { describe, expect, vi } from "vitest"; +import { FEATURE_FLAG, type FeatureFlagKey } from "~/v3/featureFlags"; +import { + applyGlobalGracedFlips, + makeSetMultipleFlags, + replaceGlobalFeatureFlags, +} from "~/v3/featureFlags.server"; + +vi.setConfig({ testTimeout: 60_000 }); + +const SET_KEYS: FeatureFlagKey[] = [ + FEATURE_FLAG.runOpsMintShardSet, + FEATURE_FLAG.runOpsMintShardSetPrev, + FEATURE_FLAG.runOpsMintShardSetFlippedAt, +]; + +const MINT_KIND_KEYS: FeatureFlagKey[] = [ + FEATURE_FLAG.runOpsMintKind, + FEATURE_FLAG.runOpsMintKindPrev, + FEATURE_FLAG.runOpsMintKindFlippedAt, +]; + +const CATALOG_KEYS: FeatureFlagKey[] = [ + ...SET_KEYS, + ...MINT_KIND_KEYS, + FEATURE_FLAG.mollifierEnabled, +]; + +const NEVER_PROTECTED = () => false; + +async function readFlags( + prisma: PrismaClient, + keys: FeatureFlagKey[] +): Promise> { + const rows = await prisma.featureFlag.findMany({ + where: { key: { in: keys } }, + select: { key: true, value: true }, + }); + const m: Record = {}; + for (const row of rows) m[row.key] = row.value; + return m; +} + +describe("applyGlobalGracedFlips — the shard-set list is stamped, not bare-written", () => { + postgresTest("a genuine list change stamps prev + flippedAt", async ({ prisma }) => { + await makeSetMultipleFlags(prisma)({ [FEATURE_FLAG.runOpsMintShardSet]: "a" }); + + await applyGlobalGracedFlips(prisma, { [FEATURE_FLAG.runOpsMintShardSet]: "a,b" }, 60_000); + + const m = await readFlags(prisma, SET_KEYS); + expect(m[FEATURE_FLAG.runOpsMintShardSet]).toBe("a,b"); + expect(m[FEATURE_FLAG.runOpsMintShardSetPrev]).toBe("a"); + expect(typeof m[FEATURE_FLAG.runOpsMintShardSetFlippedAt]).toBe("string"); + }); + + postgresTest("a first activation stamps an empty prev, which graces it", async ({ prisma }) => { + await applyGlobalGracedFlips(prisma, { [FEATURE_FLAG.runOpsMintShardSet]: "a" }, 60_000); + + const m = await readFlags(prisma, SET_KEYS); + expect(m[FEATURE_FLAG.runOpsMintShardSetPrev]).toBe(""); + expect(typeof m[FEATURE_FLAG.runOpsMintShardSetFlippedAt]).toBe("string"); + }); + + postgresTest( + "a reordered list is not a change, so the clock is not reset", + async ({ prisma }) => { + await applyGlobalGracedFlips(prisma, { [FEATURE_FLAG.runOpsMintShardSet]: "a,b" }, 60_000); + const first = await readFlags(prisma, SET_KEYS); + + await applyGlobalGracedFlips(prisma, { [FEATURE_FLAG.runOpsMintShardSet]: "b,a" }, 60_000); + const second = await readFlags(prisma, SET_KEYS); + + expect(second[FEATURE_FLAG.runOpsMintShardSetFlippedAt]).toBe( + first[FEATURE_FLAG.runOpsMintShardSetFlippedAt] + ); + } + ); + + postgresTest("both graced groups stamp in ONE save", async ({ prisma }) => { + // A save that flips the kind and the list must not stamp one and lose the other. + await makeSetMultipleFlags(prisma)({ + [FEATURE_FLAG.runOpsMintKind]: "cuid", + [FEATURE_FLAG.runOpsMintShardSet]: "a", + }); + + await applyGlobalGracedFlips( + prisma, + { + [FEATURE_FLAG.runOpsMintKind]: "runOpsId", + [FEATURE_FLAG.runOpsMintShardSet]: "a,b", + }, + 60_000 + ); + + const m = await readFlags(prisma, [...SET_KEYS, ...MINT_KIND_KEYS]); + expect(m[FEATURE_FLAG.runOpsMintKindPrev]).toBe("cuid"); + expect(m[FEATURE_FLAG.runOpsMintShardSetPrev]).toBe("a"); + expect(typeof m[FEATURE_FLAG.runOpsMintKindFlippedAt]).toBe("string"); + expect(typeof m[FEATURE_FLAG.runOpsMintShardSetFlippedAt]).toBe("string"); + }); + + postgresTest("concurrent list changes serialize on the lock", async ({ prisma }) => { + await makeSetMultipleFlags(prisma)({ [FEATURE_FLAG.runOpsMintShardSet]: "a" }); + + await Promise.all([ + applyGlobalGracedFlips(prisma, { [FEATURE_FLAG.runOpsMintShardSet]: "a,b" }, 60_000), + applyGlobalGracedFlips(prisma, { [FEATURE_FLAG.runOpsMintShardSet]: "a,c" }, 60_000), + ]); + + const m = await readFlags(prisma, SET_KEYS); + // Whichever won, the stamp must describe a real predecessor, never be absent. + expect(["a", "a,b", "a,c"]).toContain(m[FEATURE_FLAG.runOpsMintShardSetPrev]); + expect(typeof m[FEATURE_FLAG.runOpsMintShardSetFlippedAt]).toBe("string"); + }); +}); + +describe("replaceGlobalFeatureFlags — the admin page cannot bypass the stamp", () => { + postgresTest("a list change through the page is stamped", async ({ prisma }) => { + await makeSetMultipleFlags(prisma)({ [FEATURE_FLAG.runOpsMintShardSet]: "a" }); + + await replaceGlobalFeatureFlags(prisma, { + requestedFlags: { [FEATURE_FLAG.runOpsMintShardSet]: "a,b" }, + catalogKeys: CATALOG_KEYS, + isProtected: NEVER_PROTECTED, + graceMs: 60_000, + }); + + const m = await readFlags(prisma, SET_KEYS); + expect(m[FEATURE_FLAG.runOpsMintShardSet]).toBe("a,b"); + expect(m[FEATURE_FLAG.runOpsMintShardSetPrev]).toBe("a"); + }); + + postgresTest("a body-supplied stamp is ignored and recomputed", async ({ prisma }) => { + await makeSetMultipleFlags(prisma)({ [FEATURE_FLAG.runOpsMintShardSet]: "a" }); + + await replaceGlobalFeatureFlags(prisma, { + requestedFlags: { + [FEATURE_FLAG.runOpsMintShardSet]: "a,b", + [FEATURE_FLAG.runOpsMintShardSetPrev]: "zzz", + [FEATURE_FLAG.runOpsMintShardSetFlippedAt]: "1999-01-01T00:00:00.000Z", + }, + catalogKeys: CATALOG_KEYS, + isProtected: NEVER_PROTECTED, + graceMs: 60_000, + }); + + const m = await readFlags(prisma, SET_KEYS); + expect(m[FEATURE_FLAG.runOpsMintShardSetPrev]).toBe("a"); + expect(m[FEATURE_FLAG.runOpsMintShardSetFlippedAt]).not.toBe("1999-01-01T00:00:00.000Z"); + }); + + postgresTest("the set trio survives a save that omits the set keys", async ({ prisma }) => { + await replaceGlobalFeatureFlags(prisma, { + requestedFlags: { [FEATURE_FLAG.runOpsMintShardSet]: "a,b" }, + catalogKeys: CATALOG_KEYS, + isProtected: NEVER_PROTECTED, + graceMs: 60_000, + }); + + await replaceGlobalFeatureFlags(prisma, { + requestedFlags: { [FEATURE_FLAG.mollifierEnabled]: true }, + catalogKeys: CATALOG_KEYS, + isProtected: NEVER_PROTECTED, + graceMs: 60_000, + }); + + const m = await readFlags(prisma, SET_KEYS); + expect(m[FEATURE_FLAG.runOpsMintShardSet]).toBe("a,b"); + }); + + postgresTest("an ordinary flag keeps replace semantics", async ({ prisma }) => { + await makeSetMultipleFlags(prisma)({ [FEATURE_FLAG.mollifierEnabled]: true }); + + await replaceGlobalFeatureFlags(prisma, { + requestedFlags: {}, + catalogKeys: CATALOG_KEYS, + isProtected: NEVER_PROTECTED, + graceMs: 60_000, + }); + + const m = await readFlags(prisma, [FEATURE_FLAG.mollifierEnabled]); + expect(m[FEATURE_FLAG.mollifierEnabled]).toBeUndefined(); + }); +});