diff --git a/internal-packages/run-engine/src/engine/waitpointCoordinator/scripts.ts b/internal-packages/run-engine/src/engine/waitpointCoordinator/scripts.ts index 5608b7226..820b4145f 100644 --- a/internal-packages/run-engine/src/engine/waitpointCoordinator/scripts.ts +++ b/internal-packages/run-engine/src/engine/waitpointCoordinator/scripts.ts @@ -154,7 +154,11 @@ export function registerWaitpointCommands(redis: Redis): void { }); // KEYS: pend, done, edge. - // ARGV: n, then n groups of 4 — waitpointId, edgeField, edgeJson, reportedJson (''). + // ARGV: n, then n groups of 5 — waitpointId, edgeField, edgeJson, reportedFlag + // ('1'|'0'), reportedJson (''). reportedFlag, not the emptiness of reportedJson, is what + // decides the branch: a waitpoint can be reported COMPLETED with no completion envelope + // (see the FINISHED-healing path), and that case must still take the reported branch — + // flag '1', reportedJson '' — or the run would block forever on something already done. redis.defineCommand("runAbsorbBlockers", { numberOfKeys: 3, lua: ` @@ -163,7 +167,7 @@ export function registerWaitpointCommands(redis: Redis): void { -- Guard before any write: a wrong n must not half-apply the script. HDEL/HSETNX below -- are irreversible mid-script, and Redis does not roll back a script that errors. - if #ARGV ~= 1 + n * 4 then + if #ARGV ~= 1 + n * 5 then return redis.error_reply('runAbsorbBlockers: arity mismatch') end @@ -174,18 +178,20 @@ export function registerWaitpointCommands(redis: Redis): void { local out = { '0', '0' } for i = 0, n - 1 do - local id = ARGV[2 + i * 4] - local field = ARGV[3 + i * 4] - local edgeJson = ARGV[4 + i * 4] - local reported = ARGV[5 + i * 4] + local id = ARGV[2 + i * 5] + local field = ARGV[3 + i * 5] + local edgeJson = ARGV[4 + i * 5] + local reportedFlag = ARGV[5 + i * 5] + local reported = ARGV[6 + i * 5] -- HSETNX is the ON CONFLICT DO NOTHING of the edge write: a retry must not -- overwrite the first attempt's metadata. redis.call('HSETNX', edge, field, edgeJson) requestedIds[id] = true - if reported ~= '' then - -- Already COMPLETED when the watcher registered. It never becomes pending. + if reportedFlag == '1' then + -- Already COMPLETED when the watcher registered. It never becomes pending, even + -- when reported ('' here) carries no envelope. redis.call('HSET', done, id, reported) redis.call('SREM', pend, id) if not seenDelivered[id] then diff --git a/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.test.ts b/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.test.ts index 039ad78a0..5a1258b06 100644 --- a/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.test.ts +++ b/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.test.ts @@ -386,10 +386,12 @@ describe("reply framing (direct Lua — pins the wire shape the coordinator deco "w_solo", fieldA, "{}", + "0", "", "w_solo", fieldB, "{}", + "1", envelope ); @@ -419,10 +421,12 @@ describe("reply framing (direct Lua — pins the wire shape the coordinator deco "w_solo", fieldB, "{}", + "1", envelope, "w_solo", fieldA, "{}", + "0", "" ); @@ -448,10 +452,12 @@ describe("reply framing (direct Lua — pins the wire shape the coordinator deco "w_a", edgeField("w_a", 0), "{}", + "0", "", "w_b", edgeField("w_b", 0), "{}", + "0", "" ); @@ -478,10 +484,12 @@ describe("reply framing (direct Lua — pins the wire shape the coordinator deco "w_a", edgeField("w_a", 0), "{}", + "0", "", "w_b", edgeField("w_b", 0), "{}", + "1", envelope ); @@ -509,10 +517,12 @@ describe("reply framing (direct Lua — pins the wire shape the coordinator deco "w_a", edgeField("w_a", 0), "{}", + "0", "", "w_a", edgeField("w_a", 1), "{}", + "0", "" ); @@ -543,6 +553,7 @@ describe("reply framing (direct Lua — pins the wire shape the coordinator deco "w_a", edgeField("w_a", 0), "{}", + "0", "" ); @@ -561,9 +572,20 @@ describe("reply framing (direct Lua — pins the wire shape the coordinator deco const keys = runBlockKeys("run_1"); const field = edgeField("w_solo", 0); - // n says 2 groups but only one group (4 ARGV entries) is supplied. + // n says 2 groups (1 + 2 * 5 = 11 ARGV entries expected) but only one group (5 + // ARGV entries) is supplied. await expect( - client.runAbsorbBlockers(keys.pend, keys.done, keys.edge, "2", "w_solo", field, "{}", "") + client.runAbsorbBlockers( + keys.pend, + keys.done, + keys.edge, + "2", + "w_solo", + field, + "{}", + "0", + "" + ) ).rejects.toThrow(); expect(await client.exists(keys.pend)).toBe(0); @@ -1162,3 +1184,230 @@ describe("clearBlockState", () => { } ); }); + +describe("registerBlocks: a COMPLETED waitpoint with no envelope never blocks (regression)", () => { + redisTest( + "created COMPLETED with no envelope: registerBlocks does not block the run", + async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + // No `completion` at all — the FINISHED-healing shape from Task 4's "can create an + // already-COMPLETED record with no completion envelope" test. + await store.createIfAbsent({ record: record("w_a"), status: "COMPLETED" }); + + const result = await store.registerBlocks({ runId: RUN_ID, edges: [edge("w_a")] }); + + expect(result.pendingOfRequested).toBe(0); + expect(result.storePendingTotal).toBe(0); + expect(result.alreadyDelivered.map((d) => d.waitpointId)).toEqual(["w_a"]); + } finally { + await store.quit(); + } + } + ); + + redisTest( + "created COMPLETED with an envelope: behaves identically with respect to blocking", + async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ + record: record("w_a"), + status: "COMPLETED", + completion: completion(), + }); + + const result = await store.registerBlocks({ runId: RUN_ID, edges: [edge("w_a")] }); + + expect(result.pendingOfRequested).toBe(0); + expect(result.storePendingTotal).toBe(0); + expect(result.alreadyDelivered.map((d) => d.waitpointId)).toEqual(["w_a"]); + } finally { + await store.quit(); + } + } + ); +}); + +describe("registerBlocks: the two orderings", () => { + redisTest("block first, then complete: the run blocks, then wakes", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + + const blocked = await store.registerBlocks({ runId: RUN_ID, edges: [edge("w_a")] }); + expect(blocked.pendingOfRequested).toBe(1); + expect(blocked.storePendingTotal).toBe(1); + + const completed = await store.complete({ waitpointId: "w_a", completion: completion() }); + expect(completed.watchers.map((w) => w.runId)).toEqual([RUN_ID]); + + const delivered = await store.deliverCompletion({ + runId: RUN_ID, + waitpointId: "w_a", + completion: completed.completion!, + }); + expect(delivered.storePendingTotal).toBe(0); + } finally { + await store.quit(); + } + }); + + redisTest("complete first, then block: the run never goes pending", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + await store.complete({ waitpointId: "w_a", completion: completion() }); + + const result = await store.registerBlocks({ runId: RUN_ID, edges: [edge("w_a")] }); + + expect(result.pendingOfRequested).toBe(0); + expect(result.storePendingTotal).toBe(0); + expect(result.alreadyDelivered.map((d) => d.waitpointId)).toEqual(["w_a"]); + + const state = await store.readBlockState(RUN_ID); + expect(state.pendingIds).toEqual([]); + expect(state.deliveredIds).toEqual(["w_a"]); + } finally { + await store.quit(); + } + }); + + redisTest("throws when a blocking waitpoint does not exist", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await expect( + store.registerBlocks({ runId: RUN_ID, edges: [edge("w_missing")] }) + ).rejects.toThrow(WaitpointNotFoundError); + } finally { + await store.quit(); + } + }); + + redisTest("is idempotent when run twice", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + + const first = await store.registerBlocks({ runId: RUN_ID, edges: [edge("w_a")] }); + const second = await store.registerBlocks({ runId: RUN_ID, edges: [edge("w_a")] }); + + expect(first.storePendingTotal).toBe(1); + expect(second.storePendingTotal).toBe(1); + expect((await store.readBlockState(RUN_ID)).edges).toHaveLength(1); + } finally { + await store.quit(); + } + }); + + redisTest( + "mixed set: one pending and one already complete blocks the run once", + async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_pending"), status: "PENDING" }); + await store.createIfAbsent({ record: record("w_done"), status: "PENDING" }); + await store.complete({ waitpointId: "w_done", completion: completion() }); + + const result = await store.registerBlocks({ + runId: RUN_ID, + edges: [edge("w_pending"), edge("w_done")], + }); + + expect(result.pendingOfRequested).toBe(1); + expect(result.storePendingTotal).toBe(1); + expect(result.alreadyDelivered.map((d) => d.waitpointId)).toEqual(["w_done"]); + } finally { + await store.quit(); + } + } + ); +}); + +describe("multi-index merge, end to end into the executor shape", () => { + redisTest( + "a run blocked on one waitpoint at two indexes resolves to two entries", + async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ + record: record("w_child", { type: "RUN", completedByTaskRunId: "run_child" }), + status: "PENDING", + }); + + await store.registerBlocks({ + runId: RUN_ID, + edges: [ + edge("w_child", { batchIndex: 0, batchId: "batch_1", type: "RUN" }), + edge("w_child", { batchIndex: 2, batchId: "batch_1", type: "RUN" }), + ], + }); + + const completed = await store.complete({ + waitpointId: "w_child", + completion: completion({ output: null }), + }); + await store.deliverCompletion({ + runId: RUN_ID, + waitpointId: "w_child", + completion: completed.completion!, + }); + + const state = await store.readBlockState(RUN_ID); + expect(state.pendingIds).toEqual([]); + expect(state.deliveredIds).toEqual(["w_child"]); + + // Derive the cycle's ordered id list the way the read path does: keep only edges + // that carry a batch index, sort ascending, map to id. Derived inline on purpose — + // another lane owns the order rule and its resolver, and this test's job is to + // prove the COORDINATOR preserved the edge multiplicity across two shards, not to + // own that rule. + const order = state.edges + .filter((e) => e.batchIndex !== undefined && e.batchIndex !== null) + .sort((a, b) => a.batchIndex! - b.batchIndex!) + .map((e) => e.waitpointId); + + // One waitpoint, two edges, so the id repeats — that repeat is what expands into + // two entries for the executor, and losing it would silently drop a batch item. + expect(order).toEqual(["w_child", "w_child"]); + expect(state.edges.map((e) => e.edgeId).sort()).toEqual(["w_child#0", "w_child#2"]); + expect(state.edges.every((e) => e.batchId === "batch_1")).toBe(true); + } finally { + await store.quit(); + } + } + ); +}); + +describe("the resume cycle drains and can start again", () => { + redisTest("a second wait on the same waitpoint blocks nothing", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + await store.registerBlocks({ runId: RUN_ID, edges: [edge("w_a")] }); + + const completed = await store.complete({ waitpointId: "w_a", completion: completion() }); + await store.deliverCompletion({ + runId: RUN_ID, + waitpointId: "w_a", + completion: completed.completion!, + }); + + const first = await store.readBlockState(RUN_ID); + await store.clearBlockState({ runId: RUN_ID, edgeIds: first.edges.map((e) => e.edgeId) }); + + // Cycle two. The waitpoint is COMPLETED for good, so the register reports it and the + // run is never blocked. + const second = await store.registerBlocks({ + runId: RUN_ID, + edges: [edge("w_a", { batchIndex: 5 })], + }); + + expect(second.storePendingTotal).toBe(0); + expect(second.alreadyDelivered.map((d) => d.waitpointId)).toEqual(["w_a"]); + expect((await store.readBlockState(RUN_ID)).edges.map((e) => e.edgeId)).toEqual(["w_a#5"]); + } finally { + await store.quit(); + } + }); +}); diff --git a/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.ts b/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.ts index 8c4eca86b..58724f9ff 100644 --- a/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.ts +++ b/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.ts @@ -77,6 +77,13 @@ export type WatcherEntry = { createdAt: string; }; +/** + * Marks an edge as reported COMPLETED with no completion envelope — the shape a + * FINISHED-healing create can produce. Distinct from `undefined` (never reported), so the + * pending decision can key on outcome alone rather than on envelope presence. + */ +export const EMPTY_REPORTED_MARKER = Symbol("waitpoint-reported-no-envelope"); + export type CreateIfAbsentResult = | { outcome: "created" } | { @@ -109,7 +116,7 @@ export type BlockEdge = { type: WaitpointRecordInput["type"]; completedAfter?: string; /** Set when the register step already reported this waitpoint COMPLETED. */ - reported?: WaitpointCompletion; + reported?: WaitpointCompletion | typeof EMPTY_REPORTED_MARKER; }; export type AbsorbResult = { @@ -367,11 +374,17 @@ export class WaitpointStoreCoordinator { const argv: string[] = [String(args.edges.length)]; for (const item of args.edges) { const { reported, ...stored } = item; + const reportedFlag = reported !== undefined ? "1" : "0"; + const reportedJson = + reported !== undefined && reported !== EMPTY_REPORTED_MARKER + ? JSON.stringify(reported) + : ""; argv.push( item.waitpointId, edgeField(item.waitpointId, item.batchIndex), JSON.stringify(stored), - reported ? JSON.stringify(reported) : "" + reportedFlag, + reportedJson ); } @@ -392,6 +405,40 @@ export class WaitpointStoreCoordinator { }; } + /** + * Block a run on a set of waitpoints. + * + * Register on every waitpoint's own shard FIRST, then absorb on the run's shard. The + * order is the protocol: a completion that lands in between finds the watcher already + * registered, so it delivers onto the run's shard, and the absorb sees that delivery and + * never marks the waitpoint pending. Reversing the two would open the window where a + * completion is missed by both steps. + * + * The register keys the decision to skip the pending set on OUTCOME, never on whether a + * completion envelope came back — a waitpoint can be reported COMPLETED with none. + */ + async registerBlocks(args: { runId: string; edges: BlockEdge[] }): Promise { + const registered: BlockEdge[] = []; + + for (const item of args.edges) { + const result = await this.registerOrReport({ + waitpointId: item.waitpointId, + runId: args.runId, + batchIndex: item.batchIndex, + spanIdToComplete: item.spanIdToComplete, + createdAt: item.createdAt, + }); + + registered.push( + result.outcome === "completed" + ? { ...item, reported: result.completion ?? EMPTY_REPORTED_MARKER } + : item + ); + } + + return this.absorbBlockers({ runId: args.runId, edges: registered }); + } + async deliverCompletion(args: { runId: string; waitpointId: string;