test(run-store,run-engine): enumerate the parity input space, refuse an unminted cycle
Replaces hand-picked parity sampling with an exhaustive grid over every Waitpoint column the enhance oracle reads: 6144 combinations in ~35ms, compared with isDeepStrictEqual so key presence is checked too. The combination count is pinned, so a new column the oracle reads fails the assertion instead of silently shrinking coverage. Reverting the orphaned-RUN guard makes the grid report 192 divergences. A carryForward now attaches a pointer only if this incarnation actually minted the cycle. seq can be evicted while a wp:<n> key survives, and a bare key-exists check adopted a dead incarnation's order and records under a count that agreed with them, reporting no mismatch. Also drops three tests that could not fail: two asserted a literal equalled its own construction, and one compared two structurally empty arrays.
This commit is contained in:
@@ -3,6 +3,7 @@
|
||||
// records, and asserts the two agree field for field. The waitpoint lane owns the
|
||||
// production resolver; this reference exists so the frozen shapes are checked rather
|
||||
// than asserted.
|
||||
import { isDeepStrictEqual } from "node:util";
|
||||
import { describe, expect, it } from "vitest";
|
||||
import type { Waitpoint } from "@trigger.dev/database";
|
||||
import { BatchId, RunId } from "@trigger.dev/core/v3/isomorphic";
|
||||
@@ -10,7 +11,6 @@ import type { CompletedWaitpoint } from "@trigger.dev/core/v3";
|
||||
import type {
|
||||
CompletedWaitpointRecord,
|
||||
CompletedWaitpointResolver,
|
||||
CompletedWaitpointsPointer,
|
||||
ResolveCompletedWaitpointsArgs,
|
||||
} from "@internal/run-store";
|
||||
import { enhanceExecutionSnapshotWithWaitpoints } from "./executionSnapshotSystem.js";
|
||||
@@ -166,9 +166,8 @@ async function assertParity(
|
||||
order,
|
||||
records: waitpoints.map(toRecord),
|
||||
};
|
||||
// The frozen rule: count is order.length, NOT the record count. Binding it here means every
|
||||
// parity case enforces it, not only the dedicated "the frozen pointer shape" cases.
|
||||
expect(args.pointer.count).toBe(order.length);
|
||||
// count-carried-forward behaviour (order.length, not the record count) is covered by
|
||||
// the run-store Redis suite, not here -- this line only constructs `args`, not asserts.
|
||||
const resolved = await referenceResolver(args, async (id) => runOutputs[id]);
|
||||
expect(resolved).toEqual(enhanced.completedWaitpoints);
|
||||
return { enhanced, resolved };
|
||||
@@ -311,19 +310,8 @@ describe("the frozen record shape", () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe("the frozen pointer shape", () => {
|
||||
it("pins count to order.length, not the record count", () => {
|
||||
const order = ["wp_a", "wp_b", "wp_a"];
|
||||
const pointer: CompletedWaitpointsPointer = { cycleSeq: 7, count: order.length };
|
||||
expect(pointer).toEqual({ cycleSeq: 7, count: 3 });
|
||||
});
|
||||
|
||||
it("pins count at 0 when order is empty, even if records exist", () => {
|
||||
const order: string[] = [];
|
||||
const pointer: CompletedWaitpointsPointer = { cycleSeq: 7, count: order.length };
|
||||
expect(pointer).toEqual({ cycleSeq: 7, count: 0 });
|
||||
});
|
||||
});
|
||||
// The pointer's shape is pinned by CompletedWaitpointsPointer and tsconfig.freeze-test.json,
|
||||
// not by a runtime assertion here -- a value that only echoes its own construction can't fail.
|
||||
|
||||
describe("the completed-waitpoints freeze", () => {
|
||||
it("expands a repeated id at each of its positions", async () => {
|
||||
@@ -379,11 +367,6 @@ describe("the completed-waitpoints freeze", () => {
|
||||
expect(resolved[0]!.output).toBe('{"type":"STRING_ERROR"}');
|
||||
});
|
||||
|
||||
it("returns an empty list for no waitpoints", async () => {
|
||||
const { resolved } = await assertParity([], [], "batch_1");
|
||||
expect(resolved).toEqual([]);
|
||||
});
|
||||
|
||||
it("resolves through the frozen hook signature", async () => {
|
||||
// Exercises resolverUnderTest, so the declared CompletedWaitpointResolver type is
|
||||
// proved implementable at runtime, on top of the compile-time proof at its
|
||||
@@ -496,6 +479,144 @@ describe("the completed-waitpoints freeze", () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe("the exhaustive parity grid", () => {
|
||||
// Dimensions mirror every Waitpoint column the oracle reads (type, output, outputType,
|
||||
// outputIsError, completedByTaskRunId, completedByBatchId, completedAfter,
|
||||
// userProvidedIdempotencyKey, inactiveIdempotencyKey), plus order-membership and the
|
||||
// reading entry's batchId. A new column the oracle reads must widen a dimension here,
|
||||
// so the pinned combination count below fails instead of coverage silently shrinking.
|
||||
const TYPES: Waitpoint["type"][] = ["RUN", "BATCH", "DATETIME", "MANUAL"];
|
||||
const OUTPUTS: (string | null)[] = [null, '{"value":42}'];
|
||||
const OUTPUT_TYPES = ["application/json", "application/store"];
|
||||
const OUTPUT_IS_ERRORS = [false, true];
|
||||
const TASK_RUN_IDS: (string | null)[] = [null, "run_child"];
|
||||
const BATCH_IDS: (string | null)[] = [null, "batch_child"];
|
||||
const COMPLETED_AFTERS: (Date | null)[] = [null, new Date("2026-02-02T00:00:00.000Z")];
|
||||
const IDEMPOTENCY_COMBOS: Array<[boolean, string | null]> = [
|
||||
[false, null],
|
||||
[false, "cleared"],
|
||||
[true, null],
|
||||
[true, "cleared"],
|
||||
];
|
||||
const ORDER_MEMBERSHIPS = ["absent", "once", "twice"] as const;
|
||||
const READING_BATCH_IDS: (string | null)[] = [null, "batch_reading_entry"];
|
||||
|
||||
// Only reached when type is RUN, output is set, outputIsError is false, and
|
||||
// completedByTaskRunId is "run_child": the deriveFromRun branch. The value matches
|
||||
// OUTPUTS' non-null entry so a correct resolver is byte-identical to the oracle.
|
||||
const RUN_OUTPUT_LOOKUP: Record<string, string> = { run_child: '{"value":42}' };
|
||||
|
||||
it("agrees with the oracle across every combination", async () => {
|
||||
type Combo = {
|
||||
type: Waitpoint["type"];
|
||||
output: string | null;
|
||||
outputType: string;
|
||||
outputIsError: boolean;
|
||||
completedByTaskRunId: string | null;
|
||||
completedByBatchId: string | null;
|
||||
completedAfter: Date | null;
|
||||
userProvidedIdempotencyKey: boolean;
|
||||
inactiveIdempotencyKey: string | null;
|
||||
orderMembership: (typeof ORDER_MEMBERSHIPS)[number];
|
||||
readingBatchId: string | null;
|
||||
};
|
||||
const failures: Array<{ combo: Combo; oracle: unknown; resolver: unknown }> = [];
|
||||
let cases = 0;
|
||||
|
||||
for (const type of TYPES) {
|
||||
for (const output of OUTPUTS) {
|
||||
for (const outputType of OUTPUT_TYPES) {
|
||||
for (const outputIsError of OUTPUT_IS_ERRORS) {
|
||||
for (const completedByTaskRunId of TASK_RUN_IDS) {
|
||||
for (const completedByBatchId of BATCH_IDS) {
|
||||
for (const completedAfter of COMPLETED_AFTERS) {
|
||||
for (const [
|
||||
userProvidedIdempotencyKey,
|
||||
inactiveIdempotencyKey,
|
||||
] of IDEMPOTENCY_COMBOS) {
|
||||
for (const orderMembership of ORDER_MEMBERSHIPS) {
|
||||
for (const readingBatchId of READING_BATCH_IDS) {
|
||||
cases++;
|
||||
const combo: Combo = {
|
||||
type,
|
||||
output,
|
||||
outputType,
|
||||
outputIsError,
|
||||
completedByTaskRunId,
|
||||
completedByBatchId,
|
||||
completedAfter,
|
||||
userProvidedIdempotencyKey,
|
||||
inactiveIdempotencyKey,
|
||||
orderMembership,
|
||||
readingBatchId,
|
||||
};
|
||||
|
||||
const id = "wp_grid";
|
||||
const w = makeWaitpoint({
|
||||
id,
|
||||
type,
|
||||
output,
|
||||
outputType,
|
||||
outputIsError,
|
||||
completedByTaskRunId,
|
||||
completedByBatchId,
|
||||
completedAfter,
|
||||
idempotencyKey: "idem_user",
|
||||
userProvidedIdempotencyKey,
|
||||
inactiveIdempotencyKey,
|
||||
});
|
||||
const order =
|
||||
orderMembership === "absent"
|
||||
? ["wp_other"]
|
||||
: orderMembership === "once"
|
||||
? [id]
|
||||
: [id, id];
|
||||
|
||||
const enhanced = enhanceExecutionSnapshotWithWaitpoints(
|
||||
makeSnapshot(readingBatchId),
|
||||
[w],
|
||||
order
|
||||
);
|
||||
const args: ResolveCompletedWaitpointsArgs = {
|
||||
runId: "run_1",
|
||||
batchId: readingBatchId ?? undefined,
|
||||
pointer: { cycleSeq: 1, count: order.length },
|
||||
order,
|
||||
records: [toRecord(w)],
|
||||
};
|
||||
const resolved = await referenceResolver(
|
||||
args,
|
||||
async (runId) => RUN_OUTPUT_LOOKUP[runId]
|
||||
);
|
||||
|
||||
if (!isDeepStrictEqual(resolved, enhanced.completedWaitpoints)) {
|
||||
failures.push({
|
||||
combo,
|
||||
oracle: enhanced.completedWaitpoints,
|
||||
resolver: resolved,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
expect(cases).toBe(6144);
|
||||
expect(
|
||||
failures.length,
|
||||
failures.length > 0
|
||||
? `${failures.length}/${cases} combinations diverged. First: ${JSON.stringify(failures[0], null, 2)}`
|
||||
: undefined
|
||||
).toBe(0);
|
||||
});
|
||||
});
|
||||
|
||||
describe("the freeze's two deliberate divergences", () => {
|
||||
it("pins completedAt at write time, where the oracle samples the clock", () => {
|
||||
// The oracle applies `w.completedAt ?? new Date()`, so a null value changes on every
|
||||
|
||||
@@ -326,6 +326,57 @@ describe("append", () => {
|
||||
}
|
||||
);
|
||||
|
||||
redisTest(
|
||||
"a carry-forward refuses a cycle this incarnation never minted",
|
||||
async ({ redisOptions }) => {
|
||||
const store = new RedisSnapshotStore({ redisOptions, completedTtlMs: 60_000 });
|
||||
const raw = createRedisClient(redisOptions);
|
||||
try {
|
||||
await store.append({
|
||||
entry: entry({ id: "snap_1" }),
|
||||
kind: "birth",
|
||||
isTerminal: false,
|
||||
cycle: {
|
||||
kind: "new",
|
||||
completedWaitpoints: [{ id: "w_old", index: 0 }],
|
||||
records: [
|
||||
{
|
||||
id: "w_old",
|
||||
friendlyId: "waitpoint_old",
|
||||
type: "MANUAL",
|
||||
completedAt: "2026-01-01T00:00:00.000Z",
|
||||
outputType: "application/json",
|
||||
outputIsError: false,
|
||||
output: { inline: "stale" },
|
||||
},
|
||||
],
|
||||
},
|
||||
});
|
||||
|
||||
// Lose the whole keyspace except the cycle key, as under maxmemory eviction.
|
||||
await raw.del("snap:{run_1}:e", "snap:{run_1}:idx", "snap:{run_1}:cur", "snap:{run_1}:seq");
|
||||
|
||||
const carried = await store.append({
|
||||
entry: entry({ id: "snap_2" }),
|
||||
kind: "birth",
|
||||
isTerminal: false,
|
||||
cycle: { kind: "carryForward", cycleSeq: 1 },
|
||||
});
|
||||
|
||||
// Written, flagged, and carrying NO pointer: the dead incarnation's waitpoints must not
|
||||
// be served to a fresh run under a count that agrees with them.
|
||||
expect(carried).toMatchObject({ outcome: "written", cycleMismatch: true });
|
||||
expect(carried).not.toHaveProperty("cycleSeq");
|
||||
const read = await store.getLatest("run_1");
|
||||
expect(read?.cycle).toBeUndefined();
|
||||
expect(read?.completedWaitpointIds).toBeUndefined();
|
||||
} finally {
|
||||
raw.disconnect();
|
||||
await store.quit();
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
redisTest(
|
||||
"reports a duplicate id without overwriting the original entry",
|
||||
async ({ redisOptions }) => {
|
||||
|
||||
@@ -578,11 +578,15 @@ export class RedisSnapshotStore {
|
||||
redis.call('HDEL', wpKey(cycleSeq), 'records')
|
||||
end
|
||||
elseif cycleMode == 'carry' then
|
||||
cycleSeq = cycleSeqIn
|
||||
local c = redis.call('HGET', wpKey(cycleSeq), 'count')
|
||||
if not c then
|
||||
-- Attach a pointer only if this incarnation actually minted the cycle. seq can be
|
||||
-- evicted while a wp:<n> key survives, so a bare key-exists check would adopt a dead
|
||||
-- incarnation's order and records under a consistent count, invisibly.
|
||||
local minted = tonumber(redis.call('HGET', seqKey, 'c') or '0')
|
||||
local c = redis.call('HGET', wpKey(cycleSeqIn), 'count')
|
||||
if not c or minted < cycleSeqIn then
|
||||
mismatch = 1
|
||||
else
|
||||
cycleSeq = cycleSeqIn
|
||||
orderCount = c
|
||||
end
|
||||
end
|
||||
|
||||
Reference in New Issue
Block a user