Skip to content
File

Blob: src/worker/durable/project-do/commands.ts

typescript455 lines
1import { and, eq } from "drizzle-orm";
2 
3import {
4 BranchName,
5 DispatchMode,
6 ExecutionRuntime,
7 type ProjectId,
8 RunId,
9 TriggerType,
10 UnixTimestampMs,
11 type WebhookProvider,
12} from "@/contracts";
13import {
14 type AcceptManualRunResult,
15 type AcceptManualRunInput,
16 type ClaimRunWorkInput,
17 type ClaimRunWorkResult,
18 D1SyncStatus,
19 nullableTrusted,
20 type PendingRunState,
21 type ProjectDetailState,
22 ProjectRunStatus,
23 type FinalizeRunExecutionInput,
24 type FinalizeRunExecutionResult,
25 type RecoverWorkflowDispatchFailureInput,
26 type RecoverWorkflowDispatchFailureResult,
27 type RecordRunResolvedCommitResult,
28 type RequestRunCancelInput,
29 type RecordRunResolvedCommitInput,
30 type RequestRunCancelResult,
31 type RunHeartbeatInput,
32 type RunHeartbeatResult,
33 expectTrusted,
34} from "@/worker/contracts";
35import * as projectSchema from "@/worker/db/durable/schema/project-do";
36 
37import { DISPATCH_RETRY_DELAYS_MS, SANDBOX_CLEANUP_RETRY_DELAYS_MS } from "./constants";
38import { getProjectConfigRow, getRunRow, listPendingProjectDetailRows } from "./repo";
39import {
40 getProjectConfigState,
41 getProjectExecutionMaterial as getProjectExecutionMaterialState,
42 getProjectWebhookIngressState as getProjectWebhookIngressStateState,
43 initializeProject as initializeProjectConfig,
44 updateProjectConfig as updateProjectConfigState,
45} from "./project-config";
46import { getRetryDelay, syncRunMetadataToD1 } from "./reconciliation/index";
47import {
48 ensureRunInitialized,
49 ensureRunInitializedWithPayload,
50 setRunDoTerminal,
51 updateRunDoCancelRequested,
52} from "./run-do-sync";
53import {
54 armReconciliation,
55 rescheduleAlarmInTransaction,
56 sandboxCleanupRetryKey,
57 setDispatchRetryAt,
58 setHeartbeatAt,
59} from "./sidecar-state";
60import {
61 ensureProjectState,
62 nextTerminalD1SyncStatus,
63 promoteNextPendingRun,
64 transitionAcceptManualRun,
65 transitionClaimRunWork,
66 transitionFinalizeRunExecution,
67 transitionRecordRunResolvedCommit,
68 transitionRequestRunCancel,
69} from "./transitions";
70import type {
71 InitializeProjectInput,
72 ProjectConfigState,
73 ProjectDoContext,
74 ProjectExecutionMaterial,
75 ProjectWebhookIngressState,
76 UpdateProjectConfigInput,
77 UpdateProjectConfigResult,
78} from "./types";
79 
80const now = (): UnixTimestampMs => expectTrusted(UnixTimestampMs, Date.now(), "UnixTimestampMs");
81const runBestEffortSidecar = async (
82 context: ProjectDoContext,
83 operation: string,
84 projectId: ProjectId,
85 runId: RunId | null,
86 effect: () => Promise<void>,
87): Promise<void> => {
88 try {
89 await effect();
90 } catch (error) {
91 context.logger.warn("project_do_sidecar_write_failed", {
92 operation,
93 projectId,
94 runId,
95 error: error instanceof Error ? error.message : String(error),
96 });
97 }
98};
99const runCriticalReconciliationMutation = async <T>(
100 context: ProjectDoContext,
101 projectId: ProjectId,
102 mutation: (txn: DurableObjectTransaction) => Promise<T> | T,
103): Promise<T> =>
104 context.ctx.storage.transaction(async (txn) => {
105 const result = await mutation(txn);
106 await rescheduleAlarmInTransaction(context, txn, projectId);
107 return result;
108 });
109 
110// Keep these handlers aligned with the public ProjectDO RPC contract. HTTP handlers and the queue
111// consumer call the shell methods by name, so this module only moves implementation behind that surface.
112export const getProjectDetailState = async (
113 context: ProjectDoContext,
114 projectId: ProjectId,
115): Promise<ProjectDetailState> => {
116 context.db.transaction((tx) => {
117 ensureProjectState(context, tx, projectId);
118 });
119 
120 const stateRows = await context.db
121 .select()
122 .from(projectSchema.projectState)
123 .where(eq(projectSchema.projectState.projectId, projectId))
124 .limit(1);
125 const pendingRows = await listPendingProjectDetailRows(context, projectId);
126 
127 const pendingRuns: PendingRunState[] = pendingRows.map((pendingRun) => ({
128 runId: expectTrusted(RunId, pendingRun.runId, "RunId"),
129 branch: expectTrusted(BranchName, pendingRun.branch, "BranchName"),
130 queuedAt: expectTrusted(UnixTimestampMs, pendingRun.queuedAt, "UnixTimestampMs"),
131 }));
132 
133 return {
134 activeRunId: nullableTrusted(RunId, stateRows[0]?.activeRunId ?? null, "RunId"),
135 pendingRuns,
136 };
137};
138 
139export const acceptManualRun = async (
140 context: ProjectDoContext,
141 input: AcceptManualRunInput,
142): Promise<AcceptManualRunResult> => {
143 const currentTime = Date.now();
144 const transition = await context.ctx.storage.transaction(async (txn) => {
145 const projectConfigRow = getProjectConfigRow(context.db, input.projectId);
146 if (!projectConfigRow) {
147 throw new Error(`Project config ${input.projectId} is missing during manual run acceptance.`);
148 }
149 
150 const result = transitionAcceptManualRun(
151 context,
152 context.db,
153 {
154 projectId: input.projectId,
155 triggeredByUserId: input.triggeredByUserId,
156 branch: input.branch ?? expectTrusted(BranchName, projectConfigRow.defaultBranch, "BranchName"),
157 repoUrl: projectConfigRow.repoUrl,
158 configPath: projectConfigRow.configPath,
159 dispatchMode: expectTrusted(DispatchMode, projectConfigRow.dispatchMode, "DispatchMode"),
160 executionRuntime: expectTrusted(ExecutionRuntime, projectConfigRow.executionRuntime, "ExecutionRuntime"),
161 },
162 currentTime,
163 );
164 if (result.kind === "accepted") {
165 await rescheduleAlarmInTransaction(context, txn, input.projectId);
166 }
167 
168 return result;
169 });
170 
171 if (transition.kind === "rejected") {
172 return transition;
173 }
174 
175 try {
176 await ensureRunInitializedWithPayload(context, transition.runInitialization);
177 } catch (error) {
178 context.logger.error("run_do_initialize_failed", {
179 projectId: input.projectId,
180 runId: transition.runId,
181 error: error instanceof Error ? error.message : String(error),
182 });
183 }
184 
185 return {
186 kind: "accepted",
187 runId: transition.runId,
188 queuedAt: transition.queuedAt,
189 executable: transition.executable,
190 };
191};
192 
193export const claimRunWork = async (
194 context: ProjectDoContext,
195 input: ClaimRunWorkInput,
196): Promise<ClaimRunWorkResult> => {
197 const result = await runCriticalReconciliationMutation(context, input.projectId, () =>
198 transitionClaimRunWork(context, context.db, input),
199 );
200 
201 if (result.kind === "execute") {
202 await runBestEffortSidecar(context, "set_heartbeat", input.projectId, result.snapshot.runId, async () => {
203 await setHeartbeatAt(context, result.snapshot.runId, Date.now());
204 });
205 }
206 
207 return result;
208};
209 
210export const finalizeRunExecution = async (
211 context: ProjectDoContext,
212 input: FinalizeRunExecutionInput,
213): Promise<FinalizeRunExecutionResult> => {
214 const result = await runCriticalReconciliationMutation(context, input.projectId, async (txn) => {
215 const transition = transitionFinalizeRunExecution(context, context.db, input);
216 
217 if (input.sandboxDestroyed) {
218 await txn.delete(sandboxCleanupRetryKey(input.runId));
219 } else {
220 await txn.put(sandboxCleanupRetryKey(input.runId), {
221 attempt: 1,
222 nextAt: Date.now() + SANDBOX_CLEANUP_RETRY_DELAYS_MS[0],
223 });
224 }
225 
226 return transition;
227 });
228 
229 await runBestEffortSidecar(context, "clear_heartbeat", input.projectId, input.runId, async () => {
230 await setHeartbeatAt(context, input.runId, null);
231 });
232 return result;
233};
234 
235export const recoverWorkflowDispatchFailure = async (
236 context: ProjectDoContext,
237 input: RecoverWorkflowDispatchFailureInput,
238): Promise<RecoverWorkflowDispatchFailureResult> => {
239 const result = await runCriticalReconciliationMutation(context, input.projectId, async () => {
240 const row = getRunRow(context.db, input.projectId, input.runId);
241 if (!row) {
242 return {
243 kind: "stale" as const,
244 };
245 }
246 
247 if (row.status === "active" || row.status === "cancel_requested") {
248 return {
249 kind: "already_active" as const,
250 };
251 }
252 
253 if (row.status !== "executable" || row.dispatchStatus !== "queued") {
254 return {
255 kind: "stale" as const,
256 };
257 }
258 
259 const nextAttempt = row.dispatchAttempts + 1;
260 if (nextAttempt > DISPATCH_RETRY_DELAYS_MS.length) {
261 await context.db
262 .update(projectSchema.projectRuns)
263 .set({
264 status: "failed",
265 position: null,
266 dispatchStatus: "terminal",
267 dispatchAttempts: nextAttempt,
268 d1SyncStatus: nextTerminalD1SyncStatus(expectTrusted(D1SyncStatus, row.d1SyncStatus, "D1SyncStatus")),
269 lastError: "dispatch_failed",
270 })
271 .where(eq(projectSchema.projectRuns.runId, row.runId));
272 promoteNextPendingRun(context.db, input.projectId);
273 
274 return {
275 kind: "terminal" as const,
276 row,
277 };
278 }
279 
280 await context.db
281 .update(projectSchema.projectRuns)
282 .set({
283 dispatchStatus: "pending",
284 dispatchAttempts: nextAttempt,
285 lastError: input.errorMessage,
286 })
287 .where(eq(projectSchema.projectRuns.runId, row.runId));
288 
289 return {
290 kind: "rearmed" as const,
291 nextAt: Date.now() + getRetryDelay(nextAttempt, DISPATCH_RETRY_DELAYS_MS),
292 };
293 });
294 
295 if (result.kind === "stale" || result.kind === "already_active") {
296 return result;
297 }
298 
299 if (result.kind === "terminal") {
300 await runBestEffortSidecar(context, "clear_dispatch_retry", input.projectId, input.runId, async () => {
301 await setDispatchRetryAt(context, input.runId, null);
302 });
303 try {
304 await setRunDoTerminal(context, result.row, "failed", "dispatch_failed");
305 } catch (error) {
306 context.logger.error("workflow_dispatch_failure_terminalize_failed", {
307 projectId: input.projectId,
308 runId: input.runId,
309 error: error instanceof Error ? error.message : String(error),
310 });
311 }
312 return {
313 kind: "terminal",
314 };
315 }
316 
317 await runBestEffortSidecar(context, "set_dispatch_retry", input.projectId, input.runId, async () => {
318 await setDispatchRetryAt(context, input.runId, result.nextAt);
319 });
320 return {
321 kind: "rearmed",
322 };
323};
324 
325export const recordRunResolvedCommit = async (
326 context: ProjectDoContext,
327 input: RecordRunResolvedCommitInput,
328): Promise<RecordRunResolvedCommitResult> => {
329 const transition = await runCriticalReconciliationMutation(context, input.projectId, () =>
330 transitionRecordRunResolvedCommit(context, context.db, input),
331 );
332 
333 if (transition.kind === "stale") {
334 return transition;
335 }
336 
337 await syncRunMetadataToD1(context, input.projectId, input.runId, transition.row);
338 return {
339 kind: "applied",
340 };
341};
342 
343export const requestRunCancel = async (
344 context: ProjectDoContext,
345 input: RequestRunCancelInput,
346): Promise<RequestRunCancelResult> => {
347 const transition = await runCriticalReconciliationMutation(context, input.projectId, () =>
348 transitionRequestRunCancel(context, context.db, input, now()),
349 );
350 
351 if (transition.runDoAction === "canceled") {
352 await runBestEffortSidecar(context, "clear_dispatch_retry", input.projectId, input.runId, async () => {
353 await setDispatchRetryAt(context, input.runId, null);
354 });
355 await runBestEffortSidecar(context, "clear_heartbeat", input.projectId, input.runId, async () => {
356 await setHeartbeatAt(context, input.runId, null);
357 });
358 try {
359 await setRunDoTerminal(context, transition.row, "canceled", null);
360 } catch (error) {
361 context.logger.error("pending_cancel_run_do_terminalize_failed", {
362 projectId: input.projectId,
363 runId: input.runId,
364 error: error instanceof Error ? error.message : String(error),
365 });
366 }
367 } else if (transition.runDoAction === "cancel_requested") {
368 try {
369 await updateRunDoCancelRequested(context, transition.row);
370 } catch (error) {
371 context.logger.error("active_cancel_run_do_update_failed", {
372 projectId: input.projectId,
373 runId: input.runId,
374 error: error instanceof Error ? error.message : String(error),
375 });
376 }
377 }
378 
379 return {
380 runId: input.runId,
381 status: transition.runStatus,
382 cancelRequestedAt: transition.cancelRequestedAt,
383 };
384};
385 
386export const recordRunHeartbeat = async (
387 context: ProjectDoContext,
388 input: RunHeartbeatInput,
389): Promise<RunHeartbeatResult> => {
390 const row = context.db.transaction((tx) => {
391 ensureProjectState(context, tx, input.projectId);
392 return tx
393 .select()
394 .from(projectSchema.projectRuns)
395 .where(
396 and(eq(projectSchema.projectRuns.projectId, input.projectId), eq(projectSchema.projectRuns.runId, input.runId)),
397 )
398 .limit(1)
399 .get();
400 });
401 if (!row) {
402 return null;
403 }
404 
405 const control = {
406 runId: input.runId,
407 status: expectTrusted(ProjectRunStatus, row.status, "ProjectRunStatus"),
408 cancelRequestedAt: nullableTrusted(UnixTimestampMs, row.cancelRequestedAt, "UnixTimestampMs"),
409 };
410 
411 if (row.status === "active" || row.status === "cancel_requested") {
412 await runBestEffortSidecar(context, "set_heartbeat", input.projectId, input.runId, async () => {
413 await setHeartbeatAt(context, input.runId, Date.now());
414 });
415 await runBestEffortSidecar(context, "arm_reconciliation", input.projectId, input.runId, async () => {
416 await armReconciliation(context, input.projectId);
417 });
418 }
419 
420 return control;
421};
422 
423export const initializeProject = async (context: ProjectDoContext, input: InitializeProjectInput): Promise<void> => {
424 await initializeProjectConfig(context, input);
425};
426 
427export const getProjectConfig = async (
428 context: ProjectDoContext,
429 projectId: ProjectId,
430): Promise<ProjectConfigState | null> => {
431 return await getProjectConfigState(context, projectId);
432};
433 
434export const updateProjectConfig = async (
435 context: ProjectDoContext,
436 input: UpdateProjectConfigInput,
437): Promise<UpdateProjectConfigResult> => {
438 return await updateProjectConfigState(context, input);
439};
440 
441export const getProjectExecutionMaterial = async (
442 context: ProjectDoContext,
443 projectId: ProjectId,
444): Promise<ProjectExecutionMaterial | null> => {
445 return await getProjectExecutionMaterialState(context, projectId);
446};
447 
448export const getProjectWebhookIngressState = async (
449 context: ProjectDoContext,
450 projectId: ProjectId,
451 provider: WebhookProvider,
452): Promise<ProjectWebhookIngressState | null> => {
453 return await getProjectWebhookIngressStateState(context, projectId, provider);
454};