File
Blob: src/worker/dispatch/shared/run-execution-context/context.ts
| 1 | import { getSandbox } from "@cloudflare/sandbox"; |
| 2 | |
| 3 | import { type UnixTimestampMs } from "@/contracts"; |
| 4 | import { type ExecuteRunWork, type ProjectRunStatus, type RunMetaState } from "@/worker/contracts"; |
| 5 | import type { ProjectExecutionMaterial } from "@/worker/durable/project-do/types"; |
| 6 | import { redactSecrets as redactSecretValues } from "@/worker/sandbox/git"; |
| 7 | |
| 8 | import { disposeRpcStub } from "@/worker/dispatch/shared/rpc"; |
| 9 | import { deleteSandboxSessionIfExists } from "@/worker/dispatch/shared/sandbox-errors"; |
| 10 | import { |
| 11 | ensureRunCanceling, |
| 12 | ensureRunCancelRequested, |
| 13 | markOwnershipLost, |
| 14 | preserveTerminalOutcome, |
| 15 | updateRunFromCurrent, |
| 16 | } from "./control"; |
| 17 | import { |
| 18 | destroySandbox, |
| 19 | getLiveCurrentProcess, |
| 20 | hardCancelProcessTree, |
| 21 | isProcessTreeAlive, |
| 22 | softCancelProcessTree, |
| 23 | waitForProcessTreeToStop, |
| 24 | waitForProcessTreeToStopSafely, |
| 25 | } from "./process-tree"; |
| 26 | import { getProjectStub, getRunStub, kickProjectReconciliation, now } from "./shared"; |
| 27 | import type { |
| 28 | ProjectControl, |
| 29 | RunControl, |
| 30 | RunExecutionContext, |
| 31 | RunExecutionContextState, |
| 32 | RunExecutionScope, |
| 33 | RunExecutionOutcome, |
| 34 | RunLogs, |
| 35 | RunRuntime, |
| 36 | RunStore, |
| 37 | } from "./types"; |
| 38 | |
| 39 | const createScope = ( |
| 40 | env: Env, |
| 41 | executionMaterial: ProjectExecutionMaterial, |
| 42 | claim: ExecuteRunWork, |
| 43 | options?: { |
| 44 | startedAt?: UnixTimestampMs; |
| 45 | }, |
| 46 | ): RunExecutionScope => ({ |
| 47 | env, |
| 48 | executionMaterial, |
| 49 | claim, |
| 50 | snapshot: claim.snapshot, |
| 51 | projectId: claim.snapshot.projectId, |
| 52 | runId: claim.snapshot.runId, |
| 53 | repoRoot: "/workspace/repo", |
| 54 | startedAt: options?.startedAt ?? now(), |
| 55 | logContext: { |
| 56 | projectId: claim.snapshot.projectId, |
| 57 | runId: claim.snapshot.runId, |
| 58 | }, |
| 59 | }); |
| 60 | |
| 61 | const createState = (): RunExecutionContextState => ({ |
| 62 | phase: "booting", |
| 63 | session: null, |
| 64 | currentProcess: null, |
| 65 | currentStepPosition: null, |
| 66 | cancelRequestedAt: null, |
| 67 | ownershipLost: false, |
| 68 | ownershipLossStatus: null, |
| 69 | softCancelIssued: false, |
| 70 | hardCancelIssued: false, |
| 71 | preservedTerminalStatus: null, |
| 72 | redactionSecrets: [], |
| 73 | }); |
| 74 | |
| 75 | const createRunStore = (scope: RunExecutionScope): RunStore => ({ |
| 76 | getFreshStub: () => getRunStub(scope.env, scope.runId), |
| 77 | async getMeta(): Promise<RunMetaState> { |
| 78 | const current = await getRunStub(scope.env, scope.runId).getRunSummary(scope.runId); |
| 79 | if (!current) { |
| 80 | throw new Error(`Run ${scope.runId} is not initialized.`); |
| 81 | } |
| 82 | |
| 83 | return current; |
| 84 | }, |
| 85 | async updateState(input) { |
| 86 | await getRunStub(scope.env, scope.runId).updateRunState({ |
| 87 | runId: scope.runId, |
| 88 | ...input, |
| 89 | }); |
| 90 | }, |
| 91 | async tryUpdateState(input) { |
| 92 | return await getRunStub(scope.env, scope.runId).tryUpdateRunState({ |
| 93 | runId: scope.runId, |
| 94 | ...input, |
| 95 | }); |
| 96 | }, |
| 97 | async repairTerminalState(input) { |
| 98 | await getRunStub(scope.env, scope.runId).repairTerminalState({ |
| 99 | runId: scope.runId, |
| 100 | ...input, |
| 101 | }); |
| 102 | }, |
| 103 | async replaceSteps(input) { |
| 104 | await getRunStub(scope.env, scope.runId).replaceSteps({ |
| 105 | runId: scope.runId, |
| 106 | ...input, |
| 107 | }); |
| 108 | }, |
| 109 | async updateStepState(input) { |
| 110 | await getRunStub(scope.env, scope.runId).updateStepState({ |
| 111 | runId: scope.runId, |
| 112 | ...input, |
| 113 | }); |
| 114 | }, |
| 115 | async appendLogs(events) { |
| 116 | await getRunStub(scope.env, scope.runId).appendLogs({ |
| 117 | runId: scope.runId, |
| 118 | events, |
| 119 | }); |
| 120 | }, |
| 121 | }); |
| 122 | |
| 123 | const createProjectControl = (scope: RunExecutionScope): ProjectControl => ({ |
| 124 | getFreshStub: () => getProjectStub(scope.env, scope.projectId), |
| 125 | async recordHeartbeat() { |
| 126 | return await getProjectStub(scope.env, scope.projectId).recordRunHeartbeat({ |
| 127 | projectId: scope.projectId, |
| 128 | runId: scope.runId, |
| 129 | }); |
| 130 | }, |
| 131 | async recordResolvedCommit(commitSha) { |
| 132 | return await getProjectStub(scope.env, scope.projectId).recordRunResolvedCommit({ |
| 133 | projectId: scope.projectId, |
| 134 | runId: scope.runId, |
| 135 | commitSha, |
| 136 | }); |
| 137 | }, |
| 138 | async finalizeRunExecution(terminalStatus, lastError, sandboxDestroyed) { |
| 139 | return await getProjectStub(scope.env, scope.projectId).finalizeRunExecution({ |
| 140 | projectId: scope.projectId, |
| 141 | runId: scope.runId, |
| 142 | terminalStatus, |
| 143 | lastError, |
| 144 | sandboxDestroyed, |
| 145 | }); |
| 146 | }, |
| 147 | async kickReconciliation(trigger) { |
| 148 | await kickProjectReconciliation(scope.env, scope.projectId, scope.runId, trigger); |
| 149 | }, |
| 150 | }); |
| 151 | |
| 152 | const createLogs = (_scope: RunExecutionScope, state: RunExecutionContextState, runStore: RunStore): RunLogs => ({ |
| 153 | redactMessage(message: string): string { |
| 154 | return redactSecretValues(message, state.redactionSecrets); |
| 155 | }, |
| 156 | async appendSystemLog(message: string): Promise<void> { |
| 157 | await runStore.appendLogs([ |
| 158 | { |
| 159 | stream: "system", |
| 160 | chunk: `${message}\n`, |
| 161 | createdAt: now(), |
| 162 | }, |
| 163 | ]); |
| 164 | }, |
| 165 | }); |
| 166 | |
| 167 | const createControl = (scope: RunExecutionScope, state: RunExecutionContextState, runStore: RunStore): RunControl => ({ |
| 168 | async getRunMeta() { |
| 169 | return await runStore.getMeta(); |
| 170 | }, |
| 171 | async updateRunFromCurrent(current, status, overrides = {}) { |
| 172 | return await updateRunFromCurrent(scope, state, runStore, current, status, overrides); |
| 173 | }, |
| 174 | preserveTerminalOutcome(outcome: Exclude<RunExecutionOutcome, { kind: "ownership_lost" | "canceled" }>): void { |
| 175 | preserveTerminalOutcome(state, outcome); |
| 176 | }, |
| 177 | async ensureRunCancelRequested() { |
| 178 | return await ensureRunCancelRequested(scope, state, runStore); |
| 179 | }, |
| 180 | async ensureRunCanceling() { |
| 181 | return await ensureRunCanceling(scope, state, runStore); |
| 182 | }, |
| 183 | markOwnershipLost(status: ProjectRunStatus | null): void { |
| 184 | markOwnershipLost(scope, state, status); |
| 185 | }, |
| 186 | }); |
| 187 | |
| 188 | const createRuntime = (scope: RunExecutionScope, state: RunExecutionContextState): RunRuntime => { |
| 189 | const sandbox = getSandbox(scope.env.Sandbox, scope.runId, { |
| 190 | enableDefaultSession: false, |
| 191 | keepAlive: true, |
| 192 | transport: "rpc", |
| 193 | }); |
| 194 | |
| 195 | return { |
| 196 | sandbox, |
| 197 | async getSession(sessionId) { |
| 198 | return await sandbox.getSession(sessionId); |
| 199 | }, |
| 200 | async deleteSession(sessionId) { |
| 201 | await deleteSandboxSessionIfExists(sandbox, sessionId); |
| 202 | }, |
| 203 | disposeSession(session) { |
| 204 | disposeRpcStub(session); |
| 205 | }, |
| 206 | getLiveCurrentProcess() { |
| 207 | return getLiveCurrentProcess(state); |
| 208 | }, |
| 209 | async isProcessTreeAlive(session) { |
| 210 | return await isProcessTreeAlive(session); |
| 211 | }, |
| 212 | async softCancelProcessTree(session, process) { |
| 213 | await softCancelProcessTree(scope, state, session, process); |
| 214 | }, |
| 215 | async hardCancelProcessTree(session, process) { |
| 216 | await hardCancelProcessTree(scope, state, session, process); |
| 217 | }, |
| 218 | async waitForProcessTreeToStop(session, timeoutMs, pollIntervalMs = 250) { |
| 219 | return await waitForProcessTreeToStop(session, timeoutMs, pollIntervalMs); |
| 220 | }, |
| 221 | async waitForProcessTreeToStopSafely(session, timeoutMs, cleanupPhase) { |
| 222 | return await waitForProcessTreeToStopSafely(scope, session, timeoutMs, cleanupPhase); |
| 223 | }, |
| 224 | async destroySandbox() { |
| 225 | return await destroySandbox(scope, state, sandbox); |
| 226 | }, |
| 227 | dispose() { |
| 228 | disposeRpcStub(sandbox); |
| 229 | }, |
| 230 | }; |
| 231 | }; |
| 232 | |
| 233 | export const createRunExecutionContext = ( |
| 234 | env: Env, |
| 235 | executionMaterial: ProjectExecutionMaterial, |
| 236 | claim: ExecuteRunWork, |
| 237 | options?: { |
| 238 | startedAt?: UnixTimestampMs; |
| 239 | }, |
| 240 | ): RunExecutionContext => { |
| 241 | const scope = createScope(env, executionMaterial, claim, options); |
| 242 | const state = createState(); |
| 243 | const runStore = createRunStore(scope); |
| 244 | const projectControl = createProjectControl(scope); |
| 245 | const runtime = createRuntime(scope, state); |
| 246 | const logs = createLogs(scope, state, runStore); |
| 247 | const control = createControl(scope, state, runStore); |
| 248 | |
| 249 | return { |
| 250 | scope, |
| 251 | state, |
| 252 | runStore, |
| 253 | projectControl, |
| 254 | runtime, |
| 255 | logs, |
| 256 | control, |
| 257 | }; |
| 258 | }; |