Skip to content
File

Blob: src/worker/dispatch/shared/run-finalize.ts

typescript356 lines
1import type { UnixTimestampMs } from "@/contracts";
2import { isTerminalStatus, type RunMetaState } from "@/worker/contracts";
3 
4import {
5 CANCEL_GRACE_MS,
6 logger,
7 now,
8 type RunExecutionContext,
9 type RunExecutionOutcome,
10} from "@/worker/dispatch/shared/run-execution-context";
11import { type RunLeaseControl } from "@/worker/dispatch/shared/run-lease";
12import { RunStateTransitionError } from "@/worker/durable/run-do/state";
13 
14type RunFinalizeContext = Pick<
15 RunExecutionContext,
16 "control" | "projectControl" | "runStore" | "runtime" | "scope" | "state"
17>;
18 
19interface TerminalRunResult {
20 terminalStatus: Extract<RunMetaState["status"], "passed" | "failed" | "canceled">;
21 exitCode: number | null;
22 errorMessage: string | null;
23}
24 
25const isCancelTransitionState = (
26 status: RunMetaState["status"],
27): status is Extract<RunMetaState["status"], "cancel_requested" | "canceling"> =>
28 status === "cancel_requested" || status === "canceling";
29 
30const shouldRetryTerminalization = (error: unknown): boolean => error instanceof RunStateTransitionError;
31 
32const toBaseTerminalResult = (outcome: Exclude<RunExecutionOutcome, { kind: "ownership_lost" }>): TerminalRunResult => {
33 switch (outcome.kind) {
34 case "passed":
35 return {
36 terminalStatus: "passed",
37 exitCode: outcome.exitCode,
38 errorMessage: null,
39 };
40 case "failed":
41 return {
42 terminalStatus: "failed",
43 exitCode: outcome.exitCode,
44 errorMessage: outcome.errorMessage,
45 };
46 case "canceled":
47 return {
48 terminalStatus: "canceled",
49 exitCode: null,
50 errorMessage: null,
51 };
52 }
53};
54 
55const repairActiveStepState = async (
56 context: RunFinalizeContext,
57 outcome: RunExecutionOutcome,
58 finishedAtValue: UnixTimestampMs,
59): Promise<void> => {
60 const position = context.state.currentStepPosition;
61 if (position === null || outcome.kind === "passed" || outcome.kind === "ownership_lost") {
62 return;
63 }
64 
65 try {
66 const detail = await context.runStore.getFreshStub().getRunDetail(context.scope.runId);
67 const currentStep = detail.steps.find((step) => step.position === position);
68 if (!currentStep || (currentStep.status !== "queued" && currentStep.status !== "running")) {
69 return;
70 }
71 
72 await context.runStore.updateStepState({
73 position,
74 status: "failed",
75 finishedAt: finishedAtValue,
76 exitCode: outcome.kind === "failed" ? outcome.exitCode : null,
77 });
78 } catch (error) {
79 logger.warn("run_final_step_repair_failed", {
80 ...context.scope.logContext,
81 position,
82 error: error instanceof Error ? error.message : String(error),
83 });
84 }
85};
86 
87const finalizeRunCanceled = async (context: RunFinalizeContext, finishedAtValue: UnixTimestampMs): Promise<void> => {
88 let current = await context.control.getRunMeta();
89 if (current.status === "canceled") {
90 return;
91 }
92 
93 if (current.status === "passed" || current.status === "failed") {
94 throw new RunStateTransitionError("already_terminal", current.status, "canceled");
95 }
96 
97 if (current.status === "starting" || current.status === "running") {
98 current = await context.control.updateRunFromCurrent(current, "cancel_requested", {
99 startedAt: current.startedAt ?? context.scope.startedAt,
100 });
101 }
102 
103 if (current.status === "cancel_requested" || current.status === "canceling" || current.status === "queued") {
104 await context.control.updateRunFromCurrent(current, "canceled", {
105 currentStep: null,
106 startedAt: current.status === "queued" ? current.startedAt : (current.startedAt ?? context.scope.startedAt),
107 finishedAt: finishedAtValue,
108 exitCode: null,
109 errorMessage: null,
110 });
111 return;
112 }
113 
114 throw new Error(`Run ${context.scope.runId} cannot be canceled from status ${current.status}.`);
115};
116 
117const finalizeRunFailedAfterCancellation = async (
118 context: RunFinalizeContext,
119 finishedAtValue: UnixTimestampMs,
120 exitCodeValue: number | null,
121 errorMessageValue: string,
122): Promise<void> => {
123 let current = await context.control.getRunMeta();
124 if (current.status === "failed") {
125 return;
126 }
127 
128 if (current.status === "passed" || current.status === "canceled") {
129 throw new RunStateTransitionError("already_terminal", current.status, "failed");
130 }
131 
132 if (current.status === "starting" || current.status === "running") {
133 current = await context.control.updateRunFromCurrent(current, "cancel_requested", {
134 startedAt: current.startedAt ?? context.scope.startedAt,
135 });
136 }
137 
138 if (current.status === "cancel_requested") {
139 current = await context.control.updateRunFromCurrent(current, "canceling", {
140 startedAt: current.startedAt ?? context.scope.startedAt,
141 });
142 }
143 
144 if (current.status !== "canceling") {
145 throw new Error(`Run ${context.scope.runId} cannot fail after cancellation from status ${current.status}.`);
146 }
147 
148 await context.control.updateRunFromCurrent(current, "failed", {
149 currentStep: null,
150 startedAt: current.startedAt ?? context.scope.startedAt,
151 finishedAt: finishedAtValue,
152 exitCode: exitCodeValue ?? 1,
153 errorMessage: errorMessageValue,
154 });
155};
156 
157const finalizeRunTerminal = async (
158 context: RunFinalizeContext,
159 outcome: Exclude<RunExecutionOutcome, { kind: "ownership_lost" }>,
160 finishedAtValue: UnixTimestampMs,
161 sandboxDestroyedValue: boolean,
162): Promise<TerminalRunResult> => {
163 const desired = toBaseTerminalResult(outcome);
164 
165 for (let attempt = 0; attempt < 2; attempt += 1) {
166 const current = await context.control.getRunMeta();
167 if (isTerminalStatus(current.status)) {
168 return {
169 terminalStatus: current.status,
170 exitCode: current.exitCode,
171 errorMessage: current.errorMessage,
172 };
173 }
174 
175 const cancellationObserved = context.state.cancelRequestedAt !== null || isCancelTransitionState(current.status);
176 const effective =
177 cancellationObserved && !sandboxDestroyedValue
178 ? {
179 terminalStatus: "failed" as const,
180 exitCode: 1,
181 errorMessage: "cancel_cleanup_failed",
182 }
183 : cancellationObserved
184 ? {
185 terminalStatus: "canceled" as const,
186 exitCode: null,
187 errorMessage: null,
188 }
189 : desired;
190 
191 try {
192 if (effective.terminalStatus === "canceled") {
193 await finalizeRunCanceled(context, finishedAtValue);
194 } else if (cancellationObserved && effective.terminalStatus === "failed") {
195 await finalizeRunFailedAfterCancellation(
196 context,
197 finishedAtValue,
198 effective.exitCode,
199 effective.errorMessage ?? "cancel_cleanup_failed",
200 );
201 } else if (isCancelTransitionState(current.status)) {
202 await context.runStore.repairTerminalState({
203 status: effective.terminalStatus,
204 currentStep: null,
205 startedAt: current.startedAt ?? context.scope.startedAt,
206 finishedAt: finishedAtValue,
207 exitCode: effective.exitCode,
208 errorMessage: effective.errorMessage,
209 });
210 } else {
211 const updateResult = await context.runStore.tryUpdateState({
212 status: effective.terminalStatus,
213 currentStep: null,
214 startedAt: current.startedAt ?? context.scope.startedAt,
215 finishedAt: finishedAtValue,
216 exitCode: effective.exitCode,
217 errorMessage: effective.errorMessage,
218 });
219 
220 if (updateResult.kind === "conflict") {
221 throw new RunStateTransitionError(updateResult.reason, updateResult.current.status, effective.terminalStatus);
222 }
223 }
224 
225 return effective;
226 } catch (error) {
227 if (attempt === 0 && shouldRetryTerminalization(error)) {
228 continue;
229 }
230 
231 throw error;
232 }
233 }
234 
235 throw new Error(`Run ${context.scope.runId} terminalization exceeded retry budget.`);
236};
237 
238export const finalizeExecution = async (
239 context: RunFinalizeContext,
240 lease: RunLeaseControl,
241 outcome: RunExecutionOutcome,
242): Promise<void> => {
243 context.state.phase = "cleaning_up";
244 
245 try {
246 if (context.state.session) {
247 const cleanupSession = context.state.session;
248 let hasLiveProcesses = false;
249 try {
250 hasLiveProcesses = await context.runtime.isProcessTreeAlive(cleanupSession);
251 } catch (error) {
252 logger.warn("run_final_process_tree_check_failed", {
253 ...context.scope.logContext,
254 error: error instanceof Error ? error.message : String(error),
255 });
256 }
257 
258 if (!hasLiveProcesses) {
259 try {
260 await context.runtime.deleteSession(cleanupSession.id);
261 } catch (error) {
262 logger.warn("run_final_session_delete_failed", {
263 ...context.scope.logContext,
264 error: error instanceof Error ? error.message : String(error),
265 });
266 } finally {
267 context.runtime.disposeSession(cleanupSession);
268 context.state.session = null;
269 }
270 } else if (context.state.ownershipLost || outcome.kind === "ownership_lost") {
271 await context.runtime.softCancelProcessTree(context.state.session, context.runtime.getLiveCurrentProcess());
272 if (
273 !(await context.runtime.waitForProcessTreeToStopSafely(
274 context.state.session,
275 CANCEL_GRACE_MS,
276 "lost_ownership_grace",
277 ))
278 ) {
279 await context.runtime.hardCancelProcessTree(context.state.session, context.runtime.getLiveCurrentProcess());
280 }
281 } else if (lease.isCancellationRequested()) {
282 await lease.applyCancellationIfNeeded();
283 let processTreeStopped = await context.runtime.waitForProcessTreeToStopSafely(
284 context.state.session,
285 CANCEL_GRACE_MS,
286 "cancel_grace",
287 );
288 if (!processTreeStopped) {
289 await context.runtime.hardCancelProcessTree(context.state.session, context.runtime.getLiveCurrentProcess());
290 processTreeStopped = await context.runtime.waitForProcessTreeToStopSafely(
291 context.state.session,
292 5_000,
293 "cancel_force",
294 );
295 }
296 void processTreeStopped;
297 } else if (await context.runtime.isProcessTreeAlive(context.state.session)) {
298 await context.runtime.softCancelProcessTree(context.state.session, context.state.currentProcess);
299 }
300 
301 const activeCleanupSession = context.state.session;
302 if (
303 activeCleanupSession &&
304 !lease.isCancellationRequested() &&
305 !(await context.runtime.waitForProcessTreeToStopSafely(activeCleanupSession, 5_000, "final_cleanup"))
306 ) {
307 await context.runtime.hardCancelProcessTree(activeCleanupSession, context.state.currentProcess);
308 }
309 }
310 
311 const sandboxDestroyed = await context.runtime.destroySandbox();
312 
313 if (!context.state.ownershipLost && outcome.kind !== "ownership_lost") {
314 try {
315 await lease.refreshControl();
316 } catch (error) {
317 logger.warn("run_final_control_check_failed", {
318 ...context.scope.logContext,
319 error: error instanceof Error ? error.message : String(error),
320 });
321 }
322 }
323 
324 if (context.state.ownershipLost || outcome.kind === "ownership_lost") {
325 const observedStatus =
326 outcome.kind === "ownership_lost" ? outcome.observedStatus : context.state.ownershipLossStatus;
327 logger.warn("run_finalization_skipped_lost_ownership", {
328 ...context.scope.logContext,
329 status: observedStatus,
330 });
331 return;
332 }
333 
334 context.state.phase = "finalizing";
335 const finishedAt = now();
336 await repairActiveStepState(context, outcome, finishedAt);
337 context.state.currentStepPosition = null;
338 const terminal = await finalizeRunTerminal(context, outcome, finishedAt, sandboxDestroyed);
339 
340 await context.projectControl.finalizeRunExecution(terminal.terminalStatus, terminal.errorMessage, sandboxDestroyed);
341 await context.projectControl.kickReconciliation("finalize_run_execution");
342 } finally {
343 context.runtime.disposeSession(context.state.session);
344 context.state.session = null;
345 // Keep heartbeats alive until ProjectDO has either accepted terminal state or taken ownership away.
346 try {
347 await lease.stop();
348 } catch (error) {
349 logger.warn("run_heartbeat_join_failed", {
350 ...context.scope.logContext,
351 error: error instanceof Error ? error.message : String(error),
352 });
353 }
354 }
355};