Skip to content

Commit aba49bf

Browse files
1stvampTrigger.dev RepoOps
authored andcommitted
feat(run-engine): report a failed restore through a fenced worker action
Add a worker action that lets a supervisor report a failed checkpoint restore, so the run is requeued or failed promptly instead of waiting out its stall timeouts. Mono-RevId: 2a28ec932fcd5256c60369d14d7f793693d19142
1 parent 64c3ba1 commit aba49bf

11 files changed

Lines changed: 1140 additions & 26 deletions

File tree

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,80 @@
1+
import type { TypedResponse } from "@remix-run/server-runtime";
2+
import { json } from "@remix-run/server-runtime";
3+
import type { WorkerApiRunRestoreOutcomeResponseBody } from "@trigger.dev/core/v3/workers";
4+
import { WorkerApiRunRestoreOutcomeRequestBody } from "@trigger.dev/core/v3/workers";
5+
import { z } from "zod";
6+
import { logger } from "~/services/logger.server";
7+
import { createActionWorkerApiRoute } from "~/services/routeBuilders/apiBuilder.server";
8+
9+
export const action = createActionWorkerApiRoute(
10+
{
11+
params: z.object({
12+
runFriendlyId: z.string(),
13+
snapshotFriendlyId: z.string(),
14+
}),
15+
body: z.compile(WorkerApiRunRestoreOutcomeRequestBody),
16+
},
17+
async ({
18+
authenticatedWorker,
19+
params,
20+
body,
21+
runnerId,
22+
environmentId,
23+
}): Promise<TypedResponse<WorkerApiRunRestoreOutcomeResponseBody>> => {
24+
const { runFriendlyId, snapshotFriendlyId } = params;
25+
26+
logger.debug("Reporting restore outcome", { runFriendlyId, snapshotFriendlyId, body });
27+
28+
let result: Awaited<ReturnType<typeof authenticatedWorker.reportRestoreOutcome>>;
29+
try {
30+
result = await authenticatedWorker.reportRestoreOutcome({
31+
runFriendlyId,
32+
snapshotFriendlyId,
33+
outcome: body.outcome,
34+
reason: body.reason,
35+
message: body.message,
36+
snapshotRoute: body.snapshotRoute,
37+
runnerId,
38+
environmentId,
39+
});
40+
} catch (error) {
41+
if (error instanceof Response) {
42+
throw error;
43+
}
44+
45+
logger.warn("Failed to report restore outcome", {
46+
runFriendlyId,
47+
snapshotFriendlyId,
48+
environmentId,
49+
outcome: body.outcome,
50+
reason: body.reason,
51+
error,
52+
});
53+
54+
// Left to the generic 500 handler, which the worker's client retries.
55+
throw error;
56+
}
57+
58+
if (!result.ok) {
59+
logger.info("Restore outcome not applied, the run has moved on", {
60+
runFriendlyId,
61+
snapshotFriendlyId,
62+
environmentId,
63+
outcome: body.outcome,
64+
reason: body.reason,
65+
latestExecutionStatus: result.latestExecutionStatus,
66+
});
67+
68+
// The client retries a 409 by default; a stale report will never apply, so say not to.
69+
throw json(
70+
{
71+
error: "Snapshot is no longer the run's restore snapshot",
72+
latestExecutionStatus: result.latestExecutionStatus,
73+
},
74+
{ status: 409, headers: { "x-should-retry": "false" } }
75+
);
76+
}
77+
78+
return json({ ok: true });
79+
}
80+
);

‎apps/webapp/app/v3/services/worker/workerGroupTokenService.server.ts‎

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -687,6 +687,39 @@ export class AuthenticatedWorkerInstance extends WithRunEngine {
687687
});
688688
}
689689

690+
async reportRestoreOutcome({
691+
runFriendlyId,
692+
snapshotFriendlyId,
693+
outcome,
694+
reason,
695+
message,
696+
runnerId,
697+
environmentId,
698+
snapshotRoute,
699+
}: {
700+
runFriendlyId: string;
701+
snapshotFriendlyId: string;
702+
outcome: "requeue" | "fail";
703+
reason: string;
704+
message?: string;
705+
runnerId?: string;
706+
environmentId?: string;
707+
// Carried back from the restore DequeuedMessage on the supervisor's report; honors durable residency.
708+
snapshotRoute?: SnapshotRouteWire;
709+
}) {
710+
return await this._engine.reportRestoreOutcome({
711+
runId: fromFriendlyId(runFriendlyId),
712+
snapshotId: fromFriendlyId(snapshotFriendlyId),
713+
outcome,
714+
reason,
715+
message,
716+
workerId: this.workerInstanceId,
717+
runnerId,
718+
environmentId,
719+
snapshotRoute,
720+
});
721+
}
722+
690723
async getSnapshotsSince({
691724
runFriendlyId,
692725
snapshotId,

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

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,8 @@ import type {
105105
EngineWorker,
106106
HeartbeatTimeouts,
107107
ReportableQueue,
108+
RestoreOutcome,
109+
RestoreOutcomeResult,
108110
RunEngineOptions,
109111
TriggerParams,
110112
WorkerQueueDequeueOptions,
@@ -2362,6 +2364,46 @@ export class RunEngine {
23622364
});
23632365
}
23642366

2367+
/**
2368+
* Called by the worker that dequeued a run for a restore when the restore did not start.
2369+
* `requeue` gives the run back with the same checkpoint; `fail` fails the attempt as a lost
2370+
* checkpoint. A conflict is returned if `snapshotId` is no longer the restore's snapshot.
2371+
*/
2372+
async reportRestoreOutcome({
2373+
runId,
2374+
snapshotId,
2375+
outcome,
2376+
reason,
2377+
message,
2378+
workerId,
2379+
runnerId,
2380+
environmentId,
2381+
snapshotRoute,
2382+
}: {
2383+
runId: string;
2384+
snapshotId: string;
2385+
outcome: RestoreOutcome;
2386+
reason: string;
2387+
message?: string;
2388+
workerId?: string;
2389+
runnerId?: string;
2390+
environmentId?: string;
2391+
// Carried from the restore DequeuedMessage; resolved durably when absent.
2392+
snapshotRoute?: SnapshotRouteWire;
2393+
}): Promise<RestoreOutcomeResult> {
2394+
return this.runAttemptSystem.reportRestoreOutcome({
2395+
runId,
2396+
snapshotId,
2397+
outcome,
2398+
reason,
2399+
message,
2400+
workerId,
2401+
runnerId,
2402+
environmentId,
2403+
snapshotRoute,
2404+
});
2405+
}
2406+
23652407
/**
23662408
Send a heartbeat to signal the the run is still executing.
23672409
If a heartbeat isn't received, after a while the run is considered "stalled"

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

Lines changed: 26 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -137,33 +137,12 @@ export async function retryOutcomeFromCompletion(
137137
return { outcome: "fail_run", sanitizedError };
138138
}
139139

140-
const retryConfig = run.lockedRetryConfig;
141-
142-
if (!retryConfig) {
143-
return { outcome: "fail_run", sanitizedError };
144-
}
145-
146-
const parsedRetryConfig = NullishRetryOptions.safeParse(retryConfig);
147-
148-
if (!parsedRetryConfig.success) {
149-
return { outcome: "fail_run", sanitizedError };
150-
}
151-
152-
if (!parsedRetryConfig.data) {
153-
return { outcome: "fail_run", sanitizedError };
154-
}
155-
156-
const nextDelay = calculateNextRetryDelay(parsedRetryConfig.data, attemptNumber ?? 1);
140+
const retrySettings = retrySettingsFromLockedConfig(run.lockedRetryConfig, attemptNumber);
157141

158-
if (!nextDelay) {
142+
if (!retrySettings) {
159143
return { outcome: "fail_run", sanitizedError };
160144
}
161145

162-
const retrySettings = {
163-
timestamp: Date.now() + nextDelay,
164-
delay: nextDelay,
165-
};
166-
167146
return {
168147
outcome: "retry",
169148
method: "queue", // we'll always retry on the queue because usually having no settings means something bad happened
@@ -184,6 +163,30 @@ export async function retryOutcomeFromCompletion(
184163
};
185164
}
186165

166+
/** The next retry the run's locked retry config allows after `attemptNumber`, if any. */
167+
export function retrySettingsFromLockedConfig(
168+
lockedRetryConfig: unknown,
169+
attemptNumber: number | null
170+
): TaskRunExecutionRetry | undefined {
171+
if (!lockedRetryConfig) {
172+
return;
173+
}
174+
175+
const parsedRetryConfig = NullishRetryOptions.safeParse(lockedRetryConfig);
176+
177+
if (!parsedRetryConfig.success || !parsedRetryConfig.data) {
178+
return;
179+
}
180+
181+
const nextDelay = calculateNextRetryDelay(parsedRetryConfig.data, attemptNumber ?? 1);
182+
183+
if (!nextDelay) {
184+
return;
185+
}
186+
187+
return { timestamp: Date.now() + nextDelay, delay: nextDelay };
188+
}
189+
187190
async function retryOOMOnMachine(
188191
prisma: PrismaClientOrTransaction,
189192
runStore: RunStore,

0 commit comments

Comments
 (0)