File
Blob: src/worker/dispatch/workflows/steps/execute.ts
| 1 | import { type UnixTimestampMs as UnixTimestampMsType } from "@/contracts"; |
| 2 | import { isTerminalStatus, type AcceptedRunSnapshot } from "@/worker/contracts"; |
| 3 | import { |
| 4 | createRunExecutionContext, |
| 5 | ensureRunInitialized, |
| 6 | type RunExecutionOutcome, |
| 7 | } from "@/worker/dispatch/shared/run-execution-context"; |
| 8 | import { |
| 9 | appendFailureLogBestEffort, |
| 10 | executeRunSteps, |
| 11 | finalizeExecution, |
| 12 | mapExecutionErrorToOutcome, |
| 13 | prepareExecutionEnvironment, |
| 14 | recoverTerminalActiveRun, |
| 15 | RunLease, |
| 16 | } from "@/worker/dispatch/shared"; |
| 17 | import { type WorkflowStep } from "cloudflare:workers"; |
| 18 | import { toWorkflowClaim, toWorkflowTerminalStatus, getWorkflowExecutionSessionId } from "../execution"; |
| 19 | import { noRetryStepConfig } from "../step-config"; |
| 20 | import type { WorkflowRunTerminalStatus } from "../types"; |
| 21 | |
| 22 | // This step wraps sandbox reset, checkout/config loading, repo execution, and finalization, |
| 23 | // so it needs more headroom than the repo-defined run timeout alone. |
| 24 | const WORKFLOW_EXECUTE_TIMEOUT_MS = 30 * 60 * 1_000; |
| 25 | |
| 26 | type WorkflowExecutionBootstrapResult = |
| 27 | | { |
| 28 | kind: "continue"; |
| 29 | } |
| 30 | | { |
| 31 | kind: "terminal"; |
| 32 | terminalStatus: WorkflowRunTerminalStatus; |
| 33 | } |
| 34 | | { |
| 35 | kind: "canceled"; |
| 36 | }; |
| 37 | |
| 38 | const resolveWorkflowExecutionBootstrap = async ( |
| 39 | env: Env, |
| 40 | snapshot: AcceptedRunSnapshot, |
| 41 | startedAt: UnixTimestampMsType, |
| 42 | context: ReturnType<typeof createRunExecutionContext>, |
| 43 | lease: RunLease, |
| 44 | ): Promise<WorkflowExecutionBootstrapResult> => { |
| 45 | await ensureRunInitialized(env, snapshot); |
| 46 | await lease.refreshControl(); |
| 47 | lease.throwIfOwnershipLost(); |
| 48 | |
| 49 | const runMeta = await context.runStore.getMeta(); |
| 50 | if (isTerminalStatus(runMeta.status)) { |
| 51 | if (runMeta.finishedAt === null) { |
| 52 | throw new Error(`Run ${snapshot.runId} is terminal in status ${runMeta.status} without finishedAt.`); |
| 53 | } |
| 54 | |
| 55 | await recoverTerminalActiveRun(env, snapshot.projectId, snapshot.runId); |
| 56 | return { |
| 57 | kind: "terminal", |
| 58 | terminalStatus: runMeta.status, |
| 59 | }; |
| 60 | } |
| 61 | |
| 62 | if (lease.isCancellationRequested() || runMeta.status === "cancel_requested" || runMeta.status === "canceling") { |
| 63 | context.state.cancelRequestedAt = context.state.cancelRequestedAt ?? runMeta.startedAt ?? startedAt; |
| 64 | context.state.currentStepPosition = runMeta.currentStep; |
| 65 | return { |
| 66 | kind: "canceled", |
| 67 | }; |
| 68 | } |
| 69 | |
| 70 | if (runMeta.status === "queued") { |
| 71 | await context.runStore.updateState({ |
| 72 | status: "starting", |
| 73 | startedAt, |
| 74 | currentStep: null, |
| 75 | finishedAt: null, |
| 76 | exitCode: null, |
| 77 | errorMessage: null, |
| 78 | }); |
| 79 | lease.throwIfOwnershipLost(); |
| 80 | return { |
| 81 | kind: "continue", |
| 82 | }; |
| 83 | } |
| 84 | |
| 85 | if (runMeta.status === "starting" || runMeta.status === "running") { |
| 86 | // Preserve the currently reported repo step for cancellation/finalization repair. |
| 87 | // A workflow replay still rebuilds a fresh sandbox and reruns repo-defined steps |
| 88 | // from the start of the internal execution loop. |
| 89 | context.state.currentStepPosition = runMeta.currentStep; |
| 90 | return { |
| 91 | kind: "continue", |
| 92 | }; |
| 93 | } |
| 94 | |
| 95 | throw new Error(`Run ${snapshot.runId} cannot execute from status ${runMeta.status}.`); |
| 96 | }; |
| 97 | |
| 98 | export const executeWorkflowRun = async ( |
| 99 | step: WorkflowStep, |
| 100 | env: Env, |
| 101 | snapshot: AcceptedRunSnapshot, |
| 102 | startedAt: UnixTimestampMsType, |
| 103 | ): Promise<WorkflowRunTerminalStatus> => |
| 104 | await step.do( |
| 105 | "execute run", |
| 106 | noRetryStepConfig(WORKFLOW_EXECUTE_TIMEOUT_MS), |
| 107 | async (): Promise<WorkflowRunTerminalStatus> => { |
| 108 | const executionMaterial = await env.PROJECT_DO.getByName(snapshot.projectId).getProjectExecutionMaterial( |
| 109 | snapshot.projectId, |
| 110 | ); |
| 111 | if (!executionMaterial) { |
| 112 | throw new Error(`Project ${snapshot.projectId} execution material is unavailable.`); |
| 113 | } |
| 114 | |
| 115 | const executionSessionId = getWorkflowExecutionSessionId(snapshot.runId); |
| 116 | const context = createRunExecutionContext(env, executionMaterial, toWorkflowClaim(snapshot), { |
| 117 | startedAt, |
| 118 | }); |
| 119 | const lease = new RunLease(context); |
| 120 | lease.start(); |
| 121 | |
| 122 | let outcome: RunExecutionOutcome; |
| 123 | |
| 124 | try { |
| 125 | try { |
| 126 | const bootstrap = await resolveWorkflowExecutionBootstrap(env, snapshot, startedAt, context, lease); |
| 127 | if (bootstrap.kind === "terminal") { |
| 128 | return bootstrap.terminalStatus; |
| 129 | } |
| 130 | |
| 131 | if (bootstrap.kind === "canceled") { |
| 132 | outcome = { |
| 133 | kind: "canceled", |
| 134 | }; |
| 135 | } else { |
| 136 | if (!(await context.runtime.destroySandbox())) { |
| 137 | throw new Error(`Failed to reset sandbox for run ${snapshot.runId}.`); |
| 138 | } |
| 139 | |
| 140 | const prepared = await prepareExecutionEnvironment(context, lease, { |
| 141 | executionSessionId, |
| 142 | }); |
| 143 | outcome = |
| 144 | prepared === null |
| 145 | ? { |
| 146 | kind: "canceled", |
| 147 | } |
| 148 | : await executeRunSteps(context, lease, prepared); |
| 149 | } |
| 150 | } catch (error) { |
| 151 | outcome = await mapExecutionErrorToOutcome(context, error); |
| 152 | if (outcome.kind === "failed") { |
| 153 | await appendFailureLogBestEffort(context, outcome.errorMessage); |
| 154 | } |
| 155 | } |
| 156 | |
| 157 | await finalizeExecution(context, lease, outcome); |
| 158 | |
| 159 | const runMeta = await env.RUN_DO.getByName(snapshot.runId).getRunSummary(snapshot.runId); |
| 160 | if (runMeta && isTerminalStatus(runMeta.status) && runMeta.finishedAt !== null) { |
| 161 | return runMeta.status; |
| 162 | } |
| 163 | |
| 164 | return toWorkflowTerminalStatus(outcome); |
| 165 | } finally { |
| 166 | context.runtime.dispose(); |
| 167 | } |
| 168 | }, |
| 169 | ); |