File
Blob: src/worker/dispatch/shared/run-finalize.ts
| 1 | import type { UnixTimestampMs } from "@/contracts"; |
| 2 | import { isTerminalStatus, type RunMetaState } from "@/worker/contracts"; |
| 3 | |
| 4 | import { |
| 5 | CANCEL_GRACE_MS, |
| 6 | logger, |
| 7 | now, |
| 8 | type RunExecutionContext, |
| 9 | type RunExecutionOutcome, |
| 10 | } from "@/worker/dispatch/shared/run-execution-context"; |
| 11 | import { type RunLeaseControl } from "@/worker/dispatch/shared/run-lease"; |
| 12 | import { RunStateTransitionError } from "@/worker/durable/run-do/state"; |
| 13 | |
| 14 | type RunFinalizeContext = Pick< |
| 15 | RunExecutionContext, |
| 16 | "control" | "projectControl" | "runStore" | "runtime" | "scope" | "state" |
| 17 | >; |
| 18 | |
| 19 | interface TerminalRunResult { |
| 20 | terminalStatus: Extract<RunMetaState["status"], "passed" | "failed" | "canceled">; |
| 21 | exitCode: number | null; |
| 22 | errorMessage: string | null; |
| 23 | } |
| 24 | |
| 25 | const isCancelTransitionState = ( |
| 26 | status: RunMetaState["status"], |
| 27 | ): status is Extract<RunMetaState["status"], "cancel_requested" | "canceling"> => |
| 28 | status === "cancel_requested" || status === "canceling"; |
| 29 | |
| 30 | const shouldRetryTerminalization = (error: unknown): boolean => error instanceof RunStateTransitionError; |
| 31 | |
| 32 | const 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 | |
| 55 | const 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 | |
| 87 | const 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 | |
| 117 | const 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 | |
| 157 | const 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 | |
| 238 | export 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 | }; |