chore(mollifier): address CodeRabbit review for phase-1 PR

- 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) <noreply@anthropic.com>
This commit is contained in:
Dan Sutton
2026-05-14 08:43:49 +01:00
parent ae05184673
commit 31aefa1bf2
8 changed files with 179 additions and 107 deletions
@@ -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.
+18 -1
View File
@@ -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() {
@@ -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,
@@ -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<TripDecision>;
export type GateDependencies = {
isMollifierEnabled: () => boolean;
isShadowModeOn: () => boolean;
resolveOrgFlag: () => Promise<boolean>;
resolveFlag: () => Promise<boolean>;
evaluator: TripEvaluator;
logShadow: (
inputs: GateInputs,
@@ -47,7 +51,7 @@ export type GateDependencies = {
inputs: GateInputs,
decision: Extract<TripDecision, { divert: true }>,
) => 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 };
@@ -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 } : {}),
@@ -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)", () => {
+62 -49
View File
@@ -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<typeof vi.fn>;
evaluatorCalls: number;
logShadowCalls: Array<{ inputs: GateInputs; decision: Extract<TripDecision, { divert: true }> }>;
logMollifiedCalls: Array<{ inputs: GateInputs; decision: Extract<TripDecision, { divert: true }> }>;
recordDecisionCalls: Array<{ outcome: DecisionOutcome; reason?: DecisionReason }>;
};
function makeDeps(overrides: Partial<GateDependencies> = {}): {
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 }]);
});
});
+71 -45
View File
@@ -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 });
},
);
});