Skip to content

Commit 685aa94

Browse files
d-csTrigger.dev RepoOps
authored andcommitted
feat(run-engine): durable write-ahead guard for manual waitpoint completion
Internal groundwork for recovering a waitpoint completion that commits before its run-resume fanout is enqueued. The guard is opt-in and has no user-facing effect on its own. Mono-RevId: a70cf1849ee8b86a05db41f09c71fa021576de7f
1 parent df6972c commit 685aa94

8 files changed

Lines changed: 532 additions & 39 deletions

File tree

‎apps/webapp/app/v3/runOpsMigration/unblockRouteCatalog.ts‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,12 @@ export const UNBLOCK_ROUTES: readonly UnblockRoute[] = [
4747
site: WAITPOINT_SYSTEM,
4848
symbol: "completeWaitpoint (sink declaration)",
4949
},
50+
{
51+
id: "wp.ensureCompleted",
52+
kind: "RUN",
53+
site: WAITPOINT_SYSTEM,
54+
symbol: "ensureWaitpointCompleted (completion-guard redelivery)",
55+
},
5056
{
5157
id: "wp.blockAndComplete",
5258
kind: "RUN",

‎internal-packages/run-engine/src/engine/errors.ts‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,36 @@ export class RunOneTimeUseTokenError extends Error {
105105
}
106106
}
107107

108+
/**
109+
* Wraps an error thrown by `completeWaitpoint` AFTER the write-ahead completion guard was successfully
110+
* armed. Its presence is proof the guard is durably persisted and now owns eventual completion, so the
111+
* API boundary may return success for a RETRYABLE cause. A failure BEFORE the guard is armed (including
112+
* the arm's own enqueue failing) is NOT wrapped, so it propagates and the caller does not report
113+
* success. Detect with {@link isWaitpointCompletionGuardArmedError} (a branded property check, robust
114+
* across module boundaries) rather than the requested `armGuard` boolean.
115+
*/
116+
export class WaitpointCompletionGuardArmedError extends Error {
117+
readonly isWaitpointCompletionGuardArmed = true as const;
118+
readonly waitpointId: string;
119+
readonly cause?: unknown;
120+
constructor(waitpointId: string, options?: { cause?: unknown }) {
121+
super(`Waitpoint ${waitpointId} completion failed after its durable guard was armed`);
122+
this.name = "WaitpointCompletionGuardArmedError";
123+
this.waitpointId = waitpointId;
124+
this.cause = options?.cause;
125+
}
126+
}
127+
128+
export function isWaitpointCompletionGuardArmedError(
129+
error: unknown
130+
): error is WaitpointCompletionGuardArmedError {
131+
return (
132+
error instanceof Error &&
133+
(error as { isWaitpointCompletionGuardArmed?: unknown }).isWaitpointCompletionGuardArmed ===
134+
true
135+
);
136+
}
137+
108138
export class ExecutionSnapshotNotFoundError extends Error {
109139
constructor(public readonly snapshotId: string) {
110140
super(`No execution snapshot found for id ${snapshotId}`);

‎internal-packages/run-engine/src/engine/index.ts‎

Lines changed: 16 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -315,6 +315,12 @@ export class RunEngine {
315315
deferCount: payload.deferCount,
316316
});
317317
},
318+
ensureWaitpointCompleted: async ({ payload }) => {
319+
await this.waitpointSystem.ensureWaitpointCompleted({
320+
waitpointId: payload.waitpointId,
321+
output: payload.output,
322+
});
323+
},
318324
enqueueDelayedRun: async ({ payload }) => {
319325
await this.delayedRunSystem.enqueueDelayedRun({ runId: payload.runId });
320326
},
@@ -423,6 +429,7 @@ export class RunEngine {
423429
resources,
424430
executionSnapshotSystem: this.executionSnapshotSystem,
425431
enqueueSystem: this.enqueueSystem,
432+
completionGuardDelayMs: options.completionGuardDelayMs,
426433
});
427434

428435
this.ttlSystem = new TtlSystem({
@@ -2120,13 +2127,21 @@ export class RunEngine {
21202127
async completeWaitpoint({
21212128
id,
21222129
output,
2130+
armGuard,
21232131
}: {
21242132
id: string;
21252133
output?: {
21262134
value: string;
21272135
type?: string;
21282136
isError: boolean;
21292137
};
2138+
/**
2139+
* Arm the durable write-ahead completion guard for this call. The engine does NOT arm implicitly:
2140+
* left undefined it defaults to false, so a caller must opt in explicitly. The runtime-flag gate
2141+
* lives at the guarded boundary (completeWaitpointWithGuard in the webapp), which passes armGuard
2142+
* only when runStoreInfraRetryEnabled is on. Internal/system callers leave it unset (unarmed).
2143+
*/
2144+
armGuard?: boolean;
21302145
}): Promise<Waitpoint> {
21312146
// Consult the cross-seam guard FIRST so an unclassifiable id fails loudly
21322147
// here (never a silent local apply). Do NOT branch on decision.store: store routing is
@@ -2136,7 +2151,7 @@ export class RunEngine {
21362151
if (guard) {
21372152
await guard({ waitpointId: id, routeKind: "RESUME_TOKEN" });
21382153
}
2139-
return this.waitpointSystem.completeWaitpoint({ id, output });
2154+
return this.waitpointSystem.completeWaitpoint({ id, output, armGuard: armGuard ?? false });
21402155
}
21412156

21422157
/**

0 commit comments

Comments
 (0)