Skip to content

Commit 7f3e607

Browse files
carderneTrigger.dev RepoOps
authored andcommitted
fix(webapp,run-engine): recover interrupted batch completion
Mono-RevId: 6f35650c0ce4a1243998ab7fbb10545b321bb4c1
1 parent 3fef4ac commit 7f3e607

6 files changed

Lines changed: 132 additions & 77 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
area: webapp
3+
type: fix
4+
---
5+
6+
Prevent batch waits from remaining suspended when batch completion is briefly interrupted

‎apps/webapp/app/v3/runEngineHandlers.server.ts‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,6 @@ import {
99
runOpsNewReplica,
1010
runOpsLegacyReplica,
1111
runOpsNewPrismaClient,
12-
runOpsNewReplicaClient,
1312
runOpsLegacyPrismaClient,
1413
runOpsShardHandles,
1514
} from "~/db.server";
@@ -1060,7 +1059,6 @@ export function setupBatchQueueCallbacks() {
10601059
engine.setBatchCompletionCallback(async (result: CompleteBatchResult) => {
10611060
await handleBatchCompletion(result, {
10621061
splitEnabled: await splitEnabledPromise,
1063-
newReplica: runOpsNewReplicaClient,
10641062
newWriter: runOpsNewPrismaClient,
10651063
legacyWriter: runOpsLegacyPrismaClient,
10661064
shards: runOpsShardHandles,

‎apps/webapp/app/v3/runEngineHandlersShared.server.ts‎

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -76,12 +76,12 @@ export async function readRunForEventOrThrow<S extends Prisma.TaskRunSelect>(
7676
* on a single run-ops DB. Length classification is INVALID here: a batch id may
7777
* be a run-ops id (cut-over orgs) or a cuid (and cuid-shaped ids can be backfilled
7878
* onto NEW), so id-shape does not reliably indicate the row's actual residency.
79-
* The existence probe is the correct signal.
79+
* The existence probe uses the writer because batch finalization is a read-after-write
80+
* path and a lagging replica can incorrectly route a NEW batch to LEGACY.
8081
*/
8182
export async function resolveBatchRunOpsWriter(
8283
batchId: string,
8384
deps: {
84-
newReplica: RunOpsPrismaClient;
8585
newWriter: RunOpsPrismaClient;
8686
legacyWriter: RunOpsPrismaClient;
8787
shards?: ReadonlyArray<{ key: string; writer: RunOpsPrismaClient }>;
@@ -101,7 +101,7 @@ export async function resolveBatchRunOpsWriter(
101101
return shard.writer;
102102
}
103103

104-
const onNew = await deps.newReplica.batchTaskRun.findFirst({
104+
const onNew = await deps.newWriter.batchTaskRun.findFirst({
105105
where: { id: batchId },
106106
select: { id: true },
107107
});
@@ -119,7 +119,6 @@ export const QUEUE_SIZE_LIMIT_EXCEEDED_ERROR_CODE = "QUEUE_SIZE_LIMIT_EXCEEDED";
119119

120120
export type BatchCompletionDeps = {
121121
splitEnabled: boolean;
122-
newReplica: RunOpsPrismaClient;
123122
newWriter: RunOpsPrismaClient;
124123
legacyWriter: RunOpsPrismaClient;
125124
shards?: ReadonlyArray<{ key: string; writer: RunOpsPrismaClient }>;
@@ -150,7 +149,6 @@ export async function handleBatchCompletion(
150149

151150
// Always probe residency — never special-case on splitEnabled (see commit msg).
152151
const runOpsWriter = await resolveBatchRunOpsWriter(batchId, {
153-
newReplica: deps.newReplica,
154152
newWriter: deps.newWriter,
155153
legacyWriter: deps.legacyWriter,
156154
shards: deps.shards,

‎apps/webapp/test/runEngineHandlers.test.ts‎

Lines changed: 1 addition & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -117,7 +117,6 @@ async function seedBatch(
117117
function makeBatchDeps(
118118
overrides: {
119119
splitEnabled?: boolean;
120-
newReplica?: PrismaClient;
121120
newWriter?: PrismaClient;
122121
legacyWriter?: PrismaClient;
123122
legacyReplica?: PrismaClient;
@@ -127,7 +126,6 @@ function makeBatchDeps(
127126
const tryCompleteBatchCalls: string[] = [];
128127
return {
129128
splitEnabled: overrides.splitEnabled ?? false,
130-
newReplica: (overrides.newReplica ?? single)!,
131129
newWriter: (overrides.newWriter ?? single)!,
132130
legacyWriter: (overrides.legacyWriter ?? single)!,
133131
tryCompleteBatch: async (batchId: string) => {
@@ -508,7 +506,6 @@ describe("runEngineHandlers batch residency routing", () => {
508506
const shards = [{ key: "a", writer: prisma14 }] as const;
509507

510508
const writer = await resolveBatchRunOpsWriter(gen2BatchId, {
511-
newReplica: prisma17,
512509
newWriter: prisma17,
513510
legacyWriter: prisma17,
514511
shards: shards as never,
@@ -526,7 +523,6 @@ describe("runEngineHandlers batch residency routing", () => {
526523
},
527524
{
528525
splitEnabled: true,
529-
newReplica: prisma17,
530526
newWriter: prisma17,
531527
legacyWriter: prisma17,
532528
shards: shards as never,
@@ -555,13 +551,6 @@ describe("runEngineHandlers batch residency routing", () => {
555551
const shardWriter = {} as never;
556552

557553
const writer = await resolveBatchRunOpsWriter(`${"a".repeat(24)}a2`, {
558-
newReplica: {
559-
batchTaskRun: {
560-
findFirst: async () => {
561-
throw new Error("a gen-2 batch id must never probe the NEW store");
562-
},
563-
},
564-
} as never,
565554
newWriter: {} as never,
566555
legacyWriter: {} as never,
567556
shards: [{ key: "a", writer: shardWriter as never }],
@@ -573,15 +562,14 @@ describe("runEngineHandlers batch residency routing", () => {
573562
it("an unconfigured shard key fails loud rather than writing elsewhere", async () => {
574563
await expect(
575564
resolveBatchRunOpsWriter(`${"a".repeat(24)}z2`, {
576-
newReplica: {} as never,
577565
newWriter: {} as never,
578566
legacyWriter: {} as never,
579567
shards: [{ key: "a", writer: {} as never }],
580568
})
581569
).rejects.toThrow(/shard/i);
582570
});
583571

584-
// True single-DB invariant: the topology's cpFallback makes newReplica and
572+
// True single-DB invariant: the topology's cpFallback makes newWriter and
585573
// legacyWriter the SAME control-plane client, so the probe always resolves to
586574
// that one client regardless of where length-classification would guess.
587575
containerTest("true single-DB resolves to the single client", async ({ prisma }) => {
@@ -594,7 +582,6 @@ describe("runEngineHandlers batch residency routing", () => {
594582
});
595583

596584
const writer = await resolveBatchRunOpsWriter(batchId, {
597-
newReplica: prisma,
598585
newWriter: prisma,
599586
legacyWriter: prisma,
600587
});
@@ -616,15 +603,13 @@ describe("runEngineHandlers batch residency routing", () => {
616603

617604
// The probe misses on new (the new DB has no such batch) and resolves the legacy writer.
618605
const writer = await resolveBatchRunOpsWriter(batchId, {
619-
newReplica: prisma17,
620606
newWriter: prisma17,
621607
legacyWriter: prisma14,
622608
});
623609
expect(writer).toBe(prisma14);
624610

625611
const deps: BatchCompletionDeps = {
626612
splitEnabled: true,
627-
newReplica: prisma17,
628613
newWriter: prisma17,
629614
legacyWriter: prisma14,
630615
tryCompleteBatch: async () => {},
@@ -674,7 +659,6 @@ describe("runEngineHandlers batch residency routing", () => {
674659

675660
const deps: BatchCompletionDeps = {
676661
splitEnabled: false,
677-
newReplica: prisma17,
678662
newWriter: prisma17,
679663
legacyWriter: prisma14,
680664
tryCompleteBatch: async () => {},
@@ -721,15 +705,13 @@ describe("runEngineHandlers batch residency routing", () => {
721705
});
722706

723707
const writer = await resolveBatchRunOpsWriter(batchId, {
724-
newReplica: prisma17,
725708
newWriter: prisma17,
726709
legacyWriter: prisma14,
727710
});
728711
expect(writer).toBe(prisma17);
729712

730713
const deps: BatchCompletionDeps = {
731714
splitEnabled: true,
732-
newReplica: prisma17,
733715
newWriter: prisma17,
734716
legacyWriter: prisma14,
735717
tryCompleteBatch: async () => {},

‎internal-packages/run-engine/src/engine/systems/batchSystem.test.ts‎

Lines changed: 80 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -113,7 +113,8 @@ class CountingPostgresRunStore extends PostgresRunStore {
113113
async function driveBatchToAllChildrenComplete(
114114
engine: RunEngine,
115115
prisma: PrismaClient,
116-
friendlyPrefix: string
116+
friendlyPrefix: string,
117+
completeChildren = true
117118
) {
118119
const environment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
119120
const parentTask = "parent-task";
@@ -183,7 +184,7 @@ async function driveBatchToAllChildrenComplete(
183184
queue: `task/${childTask}`,
184185
isTest: false,
185186
tags: [],
186-
resumeParentOnCompletion: true,
187+
resumeParentOnCompletion: completeChildren,
187188
parentTaskRunId: parentRun.id,
188189
batch: { id: batch.id, index: 0 },
189190
},
@@ -206,39 +207,41 @@ async function driveBatchToAllChildrenComplete(
206207
queue: `task/${childTask}`,
207208
isTest: false,
208209
tags: [],
209-
resumeParentOnCompletion: true,
210+
resumeParentOnCompletion: completeChildren,
210211
parentTaskRunId: parentRun.id,
211212
batch: { id: batch.id, index: 1 },
212213
},
213214
prisma
214215
);
215216

216-
for (const child of [child1, child2]) {
217+
if (completeChildren) {
218+
for (const child of [child1, child2]) {
219+
await setTimeout(500);
220+
const dequeued = await engine.dequeueFromWorkerQueue({
221+
consumerId: "test_consumer",
222+
workerQueue: "main",
223+
});
224+
const match = dequeued.find((d) => d.run.id === child.id) ?? dequeued[0];
225+
assertNonNullable(match);
226+
const attempt = await engine.startRunAttempt({
227+
runId: match.run.id,
228+
snapshotId: match.snapshot.id,
229+
});
230+
await engine.completeRunAttempt({
231+
runId: attempt.run.id,
232+
snapshotId: attempt.snapshot.id,
233+
completion: {
234+
id: attempt.run.id,
235+
ok: true,
236+
output: '{"foo":"bar"}',
237+
outputType: "application/json",
238+
},
239+
});
240+
}
241+
217242
await setTimeout(500);
218-
const dequeued = await engine.dequeueFromWorkerQueue({
219-
consumerId: "test_consumer",
220-
workerQueue: "main",
221-
});
222-
const match = dequeued.find((d) => d.run.id === child.id) ?? dequeued[0];
223-
assertNonNullable(match);
224-
const attempt = await engine.startRunAttempt({
225-
runId: match.run.id,
226-
snapshotId: match.snapshot.id,
227-
});
228-
await engine.completeRunAttempt({
229-
runId: attempt.run.id,
230-
snapshotId: attempt.snapshot.id,
231-
completion: {
232-
id: attempt.run.id,
233-
ok: true,
234-
output: '{"foo":"bar"}',
235-
outputType: "application/json",
236-
},
237-
});
238243
}
239244

240-
await setTimeout(500);
241-
242245
return { environment, batch, parentRun, child1, child2 };
243246
}
244247

@@ -308,6 +311,57 @@ describe("RunEngine #tryCompleteBatch store routing", () => {
308311
}
309312
);
310313

314+
containerTest(
315+
"retries waitpoint completion for a completed batch whose parent was not resumed",
316+
async ({ prisma, redisOptions }) => {
317+
const engine = new RunEngine(createEngineOptions(redisOptions, prisma));
318+
319+
try {
320+
const { batch, parentRun } = await driveBatchToAllChildrenComplete(
321+
engine,
322+
prisma,
323+
"run_batch_completion_recovery",
324+
false
325+
);
326+
327+
const waitpoint = await prisma.waitpoint.findFirstOrThrow({
328+
where: { completedByBatchId: batch.id },
329+
});
330+
expect(waitpoint.status).toBe("PENDING");
331+
expect(
332+
await prisma.taskRunWaitpoint.count({
333+
where: { taskRunId: parentRun.id, waitpointId: waitpoint.id },
334+
})
335+
).toBe(1);
336+
337+
await prisma.batchTaskRun.update({
338+
where: { id: batch.id },
339+
data: { status: "COMPLETED", resumedAt: null },
340+
});
341+
342+
await engine.batchSystem.performCompleteBatch({ batchId: batch.id });
343+
344+
const recoveredBatch = await prisma.batchTaskRun.findFirstOrThrow({
345+
where: { id: batch.id },
346+
});
347+
expect(recoveredBatch.resumedAt).not.toBeNull();
348+
349+
const recoveredWaitpoint = await prisma.waitpoint.findFirstOrThrow({
350+
where: { id: waitpoint.id },
351+
});
352+
expect(recoveredWaitpoint.status).toBe("COMPLETED");
353+
354+
await setTimeout(1_000);
355+
expect(await prisma.taskRunWaitpoint.count({ where: { taskRunId: parentRun.id } })).toBe(0);
356+
const parentExecution = await engine.getRunExecutionData({ runId: parentRun.id });
357+
assertNonNullable(parentExecution);
358+
expect(parentExecution.snapshot.executionStatus).not.toBe("EXECUTING_WITH_WAITPOINTS");
359+
} finally {
360+
await engine.quit();
361+
}
362+
}
363+
);
364+
311365
// The member-run read is driven by batchId only and does not rely on the
312366
// BatchTaskRun.runtimeEnvironmentId FK. A second batch (distinct batchId) must not leak members
313367
// into the first batch's batchId-scoped read.

0 commit comments

Comments
 (0)