From 31aefa1bf2df73b58e3ca8ec53dbde89ad9ccecf Mon Sep 17 00:00:00 2001 From: Dan Sutton Date: Thu, 14 May 2026 08:43:49 +0100 Subject: [PATCH] chore(mollifier): address CodeRabbit review for phase-1 PR MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - changeset: drop "deferred" wording — phase-1 actively dual-writes + runs the drainer ack loop. - worker.server.ts: wrap mollifier drainer init in try/catch + register SIGTERM/SIGINT handlers so the polling loop stops cleanly on shutdown. - bufferedTriggerPayload: only serialise idempotencyKeyExpiresAt when an idempotencyKey is present (avoid impossible orphan-expiry payloads). - mollifierTelemetry: narrow recordDecision reason to DecisionReason union to keep OTEL attribute cardinality bounded. - mollifierGate: rename resolveOrgFlag → resolveFlag. The underlying FeatureFlag table is global by key, so the "org" prefix was misleading; per-org gating is out of scope for phase-1. - tests: drop vi.fn mocks. mollifierGate now uses plain closure spies; mollifierTripEvaluator runs against a real MollifierBuffer backed by a redisTest container (closed client exercises the fail-open path). Co-Authored-By: Claude Opus 4.7 (1M context) --- .../mollifier-redis-worker-primitives.md | 2 +- apps/webapp/app/services/worker.server.ts | 19 ++- .../bufferedTriggerPayload.server.ts | 7 +- .../app/v3/mollifier/mollifierGate.server.ts | 18 +-- .../v3/mollifier/mollifierTelemetry.server.ts | 3 +- .../test/bufferedTriggerPayload.test.ts | 10 ++ apps/webapp/test/mollifierGate.test.ts | 111 +++++++++-------- .../test/mollifierTripEvaluator.test.ts | 116 +++++++++++------- 8 files changed, 179 insertions(+), 107 deletions(-) diff --git a/.changeset/mollifier-redis-worker-primitives.md b/.changeset/mollifier-redis-worker-primitives.md index 6cd16de56..3378750a7 100644 --- a/.changeset/mollifier-redis-worker-primitives.md +++ b/.changeset/mollifier-redis-worker-primitives.md @@ -2,4 +2,4 @@ "@trigger.dev/redis-worker": patch --- -Add MollifierBuffer (with `accept`, `pop`, `ack`, `requeue`, `fail`, and `evaluateTrip`) and MollifierDrainer primitives for trigger burst smoothing. `evaluateTrip` is an atomic Lua sliding-window trip evaluator used by the webapp gate to detect per-env trigger bursts. Webapp shadow-mode logging is wired; buffer writes and drainer activation are deferred to a follow-up. +Add MollifierBuffer (with `accept`, `pop`, `ack`, `requeue`, `fail`, and `evaluateTrip`) and MollifierDrainer primitives for trigger burst smoothing. `evaluateTrip` is an atomic Lua sliding-window trip evaluator used by the webapp gate to detect per-env trigger bursts. Phase 1 wires MollifierBuffer dual-write monitoring alongside the real trigger path and runs MollifierDrainer's pop/ack loop end-to-end with a no-op handler; full buffering and replayed drainer-side triggers land in later phases. diff --git a/apps/webapp/app/services/worker.server.ts b/apps/webapp/app/services/worker.server.ts index 73524d768..038b14052 100644 --- a/apps/webapp/app/services/worker.server.ts +++ b/apps/webapp/app/services/worker.server.ts @@ -130,7 +130,24 @@ export async function init() { await workerQueue.initialize(); } - getMollifierDrainer(); + try { + const drainer = getMollifierDrainer(); + if (drainer) { + // The drainer owns a polling loop and a Redis client; let it drain + // in-flight pops on shutdown rather than tearing the process down + // mid-handler. Idempotent — `drainer.stop()` short-circuits if already + // stopped, so registering on both signals is safe. + const stopDrainer = () => { + drainer.stop().catch((error) => { + logger.error("Failed to stop mollifier drainer", { error }); + }); + }; + process.once("SIGTERM", stopDrainer); + process.once("SIGINT", stopDrainer); + } + } catch (error) { + logger.error("Failed to initialise mollifier drainer", { error }); + } } function getWorkerQueue() { diff --git a/apps/webapp/app/v3/mollifier/bufferedTriggerPayload.server.ts b/apps/webapp/app/v3/mollifier/bufferedTriggerPayload.server.ts index 340d1e9be..d251e9f98 100644 --- a/apps/webapp/app/v3/mollifier/bufferedTriggerPayload.server.ts +++ b/apps/webapp/app/v3/mollifier/bufferedTriggerPayload.server.ts @@ -92,9 +92,10 @@ export function buildBufferedTriggerPayload(input: { taskId: input.taskId, body: input.body, idempotencyKey: input.idempotencyKey, - idempotencyKeyExpiresAt: input.idempotencyKeyExpiresAt - ? input.idempotencyKeyExpiresAt.toISOString() - : null, + idempotencyKeyExpiresAt: + input.idempotencyKey && input.idempotencyKeyExpiresAt + ? input.idempotencyKeyExpiresAt.toISOString() + : null, tags: input.tags, parentRunFriendlyId: input.parentRunFriendlyId, traceContext: input.traceContext, diff --git a/apps/webapp/app/v3/mollifier/mollifierGate.server.ts b/apps/webapp/app/v3/mollifier/mollifierGate.server.ts index aa532f5a5..8bdd9757e 100644 --- a/apps/webapp/app/v3/mollifier/mollifierGate.server.ts +++ b/apps/webapp/app/v3/mollifier/mollifierGate.server.ts @@ -4,7 +4,11 @@ import { flag } from "~/v3/featureFlags.server"; import { FEATURE_FLAG } from "~/v3/featureFlags"; import { getMollifierBuffer } from "./mollifierBuffer.server"; import { createRealTripEvaluator } from "./mollifierTripEvaluator.server"; -import { recordDecision, type DecisionOutcome } from "./mollifierTelemetry.server"; +import { + recordDecision, + type DecisionOutcome, + type DecisionReason, +} from "./mollifierTelemetry.server"; // `count` is the *single-instance* sliding-window counter, not a fleet-wide // aggregate. Each webapp instance maintains its own Redis key, so the fleet @@ -37,7 +41,7 @@ export type TripEvaluator = (inputs: GateInputs) => Promise; export type GateDependencies = { isMollifierEnabled: () => boolean; isShadowModeOn: () => boolean; - resolveOrgFlag: () => Promise; + resolveFlag: () => Promise; evaluator: TripEvaluator; logShadow: ( inputs: GateInputs, @@ -47,7 +51,7 @@ export type GateDependencies = { inputs: GateInputs, decision: Extract, ) => void; - recordDecision: (outcome: DecisionOutcome, reason?: string) => void; + recordDecision: (outcome: DecisionOutcome, reason?: DecisionReason) => void; }; // `options` is a thunk so env reads happen per-evaluation, not at module load. @@ -82,7 +86,7 @@ function logDivertDecision( export const defaultGateDependencies: GateDependencies = { isMollifierEnabled: () => env.MOLLIFIER_ENABLED === "1", isShadowModeOn: () => env.MOLLIFIER_SHADOW_MODE === "1", - resolveOrgFlag: () => + resolveFlag: () => flag({ key: FEATURE_FLAG.mollifierEnabled, defaultValue: false }), evaluator: defaultEvaluator, logShadow: (inputs, decision) => @@ -103,10 +107,10 @@ export async function evaluateGate( return { action: "pass_through" }; } - const orgFlagEnabled = await d.resolveOrgFlag(); + const flagEnabled = await d.resolveFlag(); const shadowOn = d.isShadowModeOn(); - if (!orgFlagEnabled && !shadowOn) { + if (!flagEnabled && !shadowOn) { d.recordDecision("pass_through"); return { action: "pass_through" }; } @@ -117,7 +121,7 @@ export async function evaluateGate( return { action: "pass_through" }; } - if (orgFlagEnabled) { + if (flagEnabled) { d.logMollified(inputs, decision); d.recordDecision("mollify", decision.reason); return { action: "mollify", decision }; diff --git a/apps/webapp/app/v3/mollifier/mollifierTelemetry.server.ts b/apps/webapp/app/v3/mollifier/mollifierTelemetry.server.ts index fb04710bd..0fe302584 100644 --- a/apps/webapp/app/v3/mollifier/mollifierTelemetry.server.ts +++ b/apps/webapp/app/v3/mollifier/mollifierTelemetry.server.ts @@ -7,8 +7,9 @@ export const mollifierDecisionsCounter = meter.createCounter("mollifier.decision }); export type DecisionOutcome = "pass_through" | "shadow_log" | "mollify"; +export type DecisionReason = "per_env_rate"; -export function recordDecision(outcome: DecisionOutcome, reason?: string): void { +export function recordDecision(outcome: DecisionOutcome, reason?: DecisionReason): void { mollifierDecisionsCounter.add(1, { outcome, ...(reason ? { reason } : {}), diff --git a/apps/webapp/test/bufferedTriggerPayload.test.ts b/apps/webapp/test/bufferedTriggerPayload.test.ts index 4226e15d9..6280acd4c 100644 --- a/apps/webapp/test/bufferedTriggerPayload.test.ts +++ b/apps/webapp/test/bufferedTriggerPayload.test.ts @@ -50,6 +50,16 @@ describe("buildBufferedTriggerPayload", () => { const noKey = buildBufferedTriggerPayload(baseInput); expect(noKey.idempotencyKey).toBeNull(); expect(noKey.idempotencyKeyExpiresAt).toBeNull(); + + // Defensive: an expiresAt without an accompanying key is an impossible + // idempotency state — drop the expiresAt rather than serialise it. + const orphanExpiry = buildBufferedTriggerPayload({ + ...baseInput, + idempotencyKey: null, + idempotencyKeyExpiresAt: new Date("2026-05-13T10:00:00.000Z"), + }); + expect(orphanExpiry.idempotencyKey).toBeNull(); + expect(orphanExpiry.idempotencyKeyExpiresAt).toBeNull(); }); it("preserves customer body byte-equivalent (drainer replay must match Postgres)", () => { diff --git a/apps/webapp/test/mollifierGate.test.ts b/apps/webapp/test/mollifierGate.test.ts index 75374517d..bdd66b788 100644 --- a/apps/webapp/test/mollifierGate.test.ts +++ b/apps/webapp/test/mollifierGate.test.ts @@ -1,39 +1,56 @@ -import { describe, expect, it, vi } from "vitest"; +import { describe, expect, it } from "vitest"; import { evaluateGate, type GateDependencies, type GateInputs, type TripDecision, } from "~/v3/mollifier/mollifierGate.server"; +import type { DecisionOutcome, DecisionReason } from "~/v3/mollifier/mollifierTelemetry.server"; +// We deliberately don't use vi.fn here. Per repo policy tests shouldn't lean on +// mock frameworks for behaviours that are pure functions of the inputs — the +// gate is pure decision logic, so a hand-rolled "deps + spy log" wired with +// plain closures gives exactly the assertions we need without the indirection. type Spies = { - [K in keyof GateDependencies]: ReturnType; + evaluatorCalls: number; + logShadowCalls: Array<{ inputs: GateInputs; decision: Extract }>; + logMollifiedCalls: Array<{ inputs: GateInputs; decision: Extract }>; + recordDecisionCalls: Array<{ outcome: DecisionOutcome; reason?: DecisionReason }>; }; -function makeDeps(overrides: Partial = {}): { - deps: GateDependencies; - spies: Spies; -} { - const defaults: GateDependencies = { - isMollifierEnabled: () => false, - isShadowModeOn: () => false, - resolveOrgFlag: async () => false, - evaluator: async () => ({ divert: false }) as TripDecision, - logShadow: () => {}, - logMollified: () => {}, - recordDecision: () => {}, +type Toggles = { + enabled: boolean; + shadow: boolean; + flag: boolean; + decision: TripDecision; +}; + +function makeDeps(toggles: Toggles): { deps: GateDependencies; spies: Spies } { + const spies: Spies = { + evaluatorCalls: 0, + logShadowCalls: [], + logMollifiedCalls: [], + recordDecisionCalls: [], }; - const merged = { ...defaults, ...overrides }; - const spies = { - isMollifierEnabled: vi.fn(merged.isMollifierEnabled), - isShadowModeOn: vi.fn(merged.isShadowModeOn), - resolveOrgFlag: vi.fn(merged.resolveOrgFlag), - evaluator: vi.fn(merged.evaluator), - logShadow: vi.fn(merged.logShadow), - logMollified: vi.fn(merged.logMollified), - recordDecision: vi.fn(merged.recordDecision), - } satisfies Spies; - return { deps: spies, spies }; + const deps: GateDependencies = { + isMollifierEnabled: () => toggles.enabled, + isShadowModeOn: () => toggles.shadow, + resolveFlag: async () => toggles.flag, + evaluator: async () => { + spies.evaluatorCalls += 1; + return toggles.decision; + }, + logShadow: (inputs, decision) => { + spies.logShadowCalls.push({ inputs, decision }); + }, + logMollified: (inputs, decision) => { + spies.logMollifiedCalls.push({ inputs, decision }); + }, + recordDecision: (outcome, reason) => { + spies.recordDecisionCalls.push({ outcome, reason }); + }, + }; + return { deps, spies }; } const trippedDecision = { @@ -101,53 +118,49 @@ describe("evaluateGate cascade — exhaustive truth table", () => { "row $id: enabled=$enabled shadow=$shadow flag=$flag divert=$divert → action=$expected.action", async (row) => { const { deps, spies } = makeDeps({ - isMollifierEnabled: () => row.enabled, - isShadowModeOn: () => row.shadow, - resolveOrgFlag: async () => row.flag, - evaluator: async () => (row.divert ? trippedDecision : passDecision), + enabled: row.enabled, + shadow: row.shadow, + flag: row.flag, + decision: row.divert ? trippedDecision : passDecision, }); const outcome = await evaluateGate(inputs, deps); expect(outcome.action).toBe(row.expected.action); - expect(spies.evaluator).toHaveBeenCalledTimes(row.expected.evaluatorCalls); - expect(spies.logShadow).toHaveBeenCalledTimes(row.expected.logShadowCalls); - expect(spies.logMollified).toHaveBeenCalledTimes(row.expected.logMollifiedCalls); + expect(spies.evaluatorCalls).toBe(row.expected.evaluatorCalls); + expect(spies.logShadowCalls).toHaveLength(row.expected.logShadowCalls); + expect(spies.logMollifiedCalls).toHaveLength(row.expected.logMollifiedCalls); // Every evaluation records exactly one decision. - expect(spies.recordDecision).toHaveBeenCalledTimes(1); - if (row.expected.expectedReason === undefined) { - expect(spies.recordDecision).toHaveBeenCalledWith(row.expected.recordedOutcome); - } else { - expect(spies.recordDecision).toHaveBeenCalledWith( - row.expected.recordedOutcome, - row.expected.expectedReason, - ); - } + expect(spies.recordDecisionCalls).toHaveLength(1); + expect(spies.recordDecisionCalls[0].outcome).toBe(row.expected.recordedOutcome); + expect(spies.recordDecisionCalls[0].reason).toBe(row.expected.expectedReason); }, ); it("divert log carries the full decision (envId, orgId, taskId, reason, count, threshold, windowMs, holdMs)", async () => { const { deps, spies } = makeDeps({ - isMollifierEnabled: () => true, - isShadowModeOn: () => true, - evaluator: async () => trippedDecision, + enabled: true, + shadow: true, + flag: false, + decision: trippedDecision, }); await evaluateGate(inputs, deps); - expect(spies.logShadow).toHaveBeenCalledWith(inputs, trippedDecision); + expect(spies.logShadowCalls).toEqual([{ inputs, decision: trippedDecision }]); }); it("mollify log carries the full decision (mirrors shadow log)", async () => { const { deps, spies } = makeDeps({ - isMollifierEnabled: () => true, - resolveOrgFlag: async () => true, - evaluator: async () => trippedDecision, + enabled: true, + shadow: false, + flag: true, + decision: trippedDecision, }); await evaluateGate(inputs, deps); - expect(spies.logMollified).toHaveBeenCalledWith(inputs, trippedDecision); + expect(spies.logMollifiedCalls).toEqual([{ inputs, decision: trippedDecision }]); }); }); diff --git a/apps/webapp/test/mollifierTripEvaluator.test.ts b/apps/webapp/test/mollifierTripEvaluator.test.ts index e97418726..b9a9bf8c9 100644 --- a/apps/webapp/test/mollifierTripEvaluator.test.ts +++ b/apps/webapp/test/mollifierTripEvaluator.test.ts @@ -1,64 +1,90 @@ -import { describe, expect, it, vi } from "vitest"; +import { redisTest } from "@internal/testcontainers"; +import { MollifierBuffer } from "@trigger.dev/redis-worker"; +import { describe, expect, vi } from "vitest"; import { createRealTripEvaluator } from "~/v3/mollifier/mollifierTripEvaluator.server"; -import type { MollifierBuffer } from "@trigger.dev/redis-worker"; -function fakeBuffer(result: { tripped: boolean; count: number }): MollifierBuffer { - return { - evaluateTrip: vi.fn(async () => result), - } as unknown as MollifierBuffer; -} +vi.setConfig({ testTimeout: 30_000 }); + +// Use a real MollifierBuffer backed by a Redis testcontainer — repo policy +// is no mocks for Redis. Per-test envIds keep keys disjoint without explicit +// cleanup. We close() the buffer in a finally to release the client. +const inputs = { envId: "env_a", orgId: "org_1", taskId: "t1" } as const; describe("createRealTripEvaluator", () => { - it("returns divert=false when buffer reports not tripped", async () => { - const evaluator = createRealTripEvaluator({ - getBuffer: () => fakeBuffer({ tripped: false, count: 42 }), - options: () => ({ windowMs: 200, threshold: 100, holdMs: 500 }), - }); + redisTest( + "returns divert=false when the sliding window stays under threshold", + async ({ redisOptions }) => { + const buffer = new MollifierBuffer({ redisOptions, entryTtlSeconds: 600 }); + try { + const evaluator = createRealTripEvaluator({ + getBuffer: () => buffer, + options: () => ({ windowMs: 1000, threshold: 100, holdMs: 500 }), + }); - const decision = await evaluator({ envId: "env_a", orgId: "org_1", taskId: "t1" }); - expect(decision).toEqual({ divert: false }); - }); + const decision = await evaluator({ ...inputs, envId: "env_under" }); + expect(decision).toEqual({ divert: false }); + } finally { + await buffer.close(); + } + }, + ); - it("returns divert=true with reason per_env_rate when buffer reports tripped", async () => { - const evaluator = createRealTripEvaluator({ - getBuffer: () => fakeBuffer({ tripped: true, count: 150 }), - options: () => ({ windowMs: 200, threshold: 100, holdMs: 500 }), - }); + redisTest( + "returns divert=true with reason per_env_rate once the window trips", + async ({ redisOptions }) => { + const buffer = new MollifierBuffer({ redisOptions, entryTtlSeconds: 600 }); + try { + // threshold=2 → the 3rd call within windowMs is the first that trips. + const options = { windowMs: 5000, threshold: 2, holdMs: 5000 } as const; + const evaluator = createRealTripEvaluator({ + getBuffer: () => buffer, + options: () => options, + }); - const decision = await evaluator({ envId: "env_a", orgId: "org_1", taskId: "t1" }); - expect(decision).toEqual({ - divert: true, - reason: "per_env_rate", - count: 150, - threshold: 100, - windowMs: 200, - holdMs: 500, - }); - }); + const envId = "env_trip"; + await evaluator({ ...inputs, envId }); + await evaluator({ ...inputs, envId }); + const decision = await evaluator({ ...inputs, envId }); - it("returns divert=false when getBuffer returns null (fail-open)", async () => { + expect(decision.divert).toBe(true); + if (decision.divert) { + expect(decision.reason).toBe("per_env_rate"); + expect(decision.threshold).toBe(options.threshold); + expect(decision.windowMs).toBe(options.windowMs); + expect(decision.holdMs).toBe(options.holdMs); + expect(decision.count).toBeGreaterThan(options.threshold); + } + } finally { + await buffer.close(); + } + }, + ); + + redisTest("returns divert=false when getBuffer returns null (fail-open)", async () => { const evaluator = createRealTripEvaluator({ getBuffer: () => null, options: () => ({ windowMs: 200, threshold: 100, holdMs: 500 }), }); - const decision = await evaluator({ envId: "env_a", orgId: "org_1", taskId: "t1" }); + const decision = await evaluator(inputs); expect(decision).toEqual({ divert: false }); }); - it("returns divert=false when buffer throws (fail-open)", async () => { - const errorBuffer = { - evaluateTrip: vi.fn(async () => { - throw new Error("redis unavailable"); - }), - } as unknown as MollifierBuffer; + redisTest( + "returns divert=false when buffer throws (fail-open)", + async ({ redisOptions }) => { + const buffer = new MollifierBuffer({ redisOptions, entryTtlSeconds: 600 }); + // Closing the client up front means evaluateTrip will throw on the first + // Redis command — a real failure mode, not a stub. + await buffer.close(); - const evaluator = createRealTripEvaluator({ - getBuffer: () => errorBuffer, - options: () => ({ windowMs: 200, threshold: 100, holdMs: 500 }), - }); + const evaluator = createRealTripEvaluator({ + getBuffer: () => buffer, + options: () => ({ windowMs: 200, threshold: 100, holdMs: 500 }), + }); - const decision = await evaluator({ envId: "env_a", orgId: "org_1", taskId: "t1" }); - expect(decision).toEqual({ divert: false }); - }); + const decision = await evaluator(inputs); + expect(decision).toEqual({ divert: false }); + }, + ); });