Skip to content
File

Blob: src/worker/dispatch/queue/execute-dispatched-run.ts

typescript98 lines
1import { type ProjectId, type RunId } from "@/contracts";
2import { type ExecuteRunWork } from "@/worker/contracts";
3import {
4 createRunExecutionContext,
5 ensureRunInitialized,
6 type RunExecutionOutcome,
7} from "@/worker/dispatch/shared/run-execution-context";
8import {
9 appendFailureLogBestEffort,
10 executeRunSteps,
11 finalizeExecution,
12 mapExecutionErrorToOutcome,
13 prepareExecutionEnvironment,
14 recoverTerminalActiveRun,
15 RunLease,
16} from "@/worker/dispatch/shared";
17import type { ProjectExecutionMaterial } from "@/worker/durable/project-do/types";
18 
19export interface DispatchedRunResult {
20 readonly kind: "executed" | "stale" | "recovered" | "project_missing";
21 readonly reason?: string;
22}
23 
24const executeClaimedRun = async (
25 env: Env,
26 executionMaterial: ProjectExecutionMaterial,
27 claim: ExecuteRunWork,
28): Promise<void> => {
29 await ensureRunInitialized(env, claim.snapshot);
30 
31 const context = createRunExecutionContext(env, executionMaterial, claim);
32 const lease = new RunLease(context);
33 lease.start();
34 
35 let outcome: RunExecutionOutcome;
36 
37 try {
38 await context.runStore.updateState({
39 status: "starting",
40 startedAt: context.scope.startedAt,
41 currentStep: null,
42 finishedAt: null,
43 exitCode: null,
44 errorMessage: null,
45 });
46 lease.throwIfOwnershipLost();
47 
48 const prepared = await prepareExecutionEnvironment(context, lease);
49 outcome =
50 prepared === null
51 ? {
52 kind: "canceled",
53 }
54 : await executeRunSteps(context, lease, prepared);
55 } catch (error) {
56 outcome = await mapExecutionErrorToOutcome(context, error);
57 if (outcome.kind === "failed") {
58 await appendFailureLogBestEffort(context, outcome.errorMessage);
59 }
60 }
61 
62 try {
63 await finalizeExecution(context, lease, outcome);
64 } finally {
65 context.runtime.dispose();
66 }
67};
68 
69export const executeDispatchedRun = async (
70 env: Env,
71 input: { projectId: ProjectId; runId: RunId },
72): Promise<DispatchedRunResult> => {
73 const projectStub = env.PROJECT_DO.getByName(input.projectId);
74 const claim = await projectStub.claimRunWork({
75 projectId: input.projectId,
76 runId: input.runId,
77 });
78 
79 if (claim.kind === "stale") {
80 if (claim.reason === "run_active") {
81 const recovered = await recoverTerminalActiveRun(env, input.projectId, input.runId);
82 if (recovered) {
83 return { kind: "recovered" };
84 }
85 }
86 
87 return { kind: "stale", reason: claim.reason };
88 }
89 
90 const executionMaterial = await projectStub.getProjectExecutionMaterial(input.projectId);
91 if (!executionMaterial) {
92 return { kind: "project_missing" };
93 }
94 
95 await executeClaimedRun(env, executionMaterial, claim);
96 return { kind: "executed" };
97};