Skip to content
File

Blob: src/worker/dispatch/shared/run-execution-context/context.ts

typescript259 lines
1import { getSandbox } from "@cloudflare/sandbox";
2 
3import { type UnixTimestampMs } from "@/contracts";
4import { type ExecuteRunWork, type ProjectRunStatus, type RunMetaState } from "@/worker/contracts";
5import type { ProjectExecutionMaterial } from "@/worker/durable/project-do/types";
6import { redactSecrets as redactSecretValues } from "@/worker/sandbox/git";
7 
8import { disposeRpcStub } from "@/worker/dispatch/shared/rpc";
9import { deleteSandboxSessionIfExists } from "@/worker/dispatch/shared/sandbox-errors";
10import {
11 ensureRunCanceling,
12 ensureRunCancelRequested,
13 markOwnershipLost,
14 preserveTerminalOutcome,
15 updateRunFromCurrent,
16} from "./control";
17import {
18 destroySandbox,
19 getLiveCurrentProcess,
20 hardCancelProcessTree,
21 isProcessTreeAlive,
22 softCancelProcessTree,
23 waitForProcessTreeToStop,
24 waitForProcessTreeToStopSafely,
25} from "./process-tree";
26import { getProjectStub, getRunStub, kickProjectReconciliation, now } from "./shared";
27import type {
28 ProjectControl,
29 RunControl,
30 RunExecutionContext,
31 RunExecutionContextState,
32 RunExecutionScope,
33 RunExecutionOutcome,
34 RunLogs,
35 RunRuntime,
36 RunStore,
37} from "./types";
38 
39const 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 
61const 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 
75const 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 
123const 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 
152const 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 
167const 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 
188const 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 
233export 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};