Skip to content
File

Blob: src/worker/durable/project-do/run-do-sync.ts

typescript197 lines
1import { BranchName, CommitSha, ProjectId, RunId, TriggerType, UnixTimestampMs } from "@/contracts";
2import {
3 type ProjectRunTerminalStatus,
4 type EnsureRunInput,
5 type UpdateRunStateInput,
6 expectTrusted,
7 nullableTrusted,
8} from "@/worker/contracts";
9 
10import type { ProjectDoContext, ProjectRunRow, RunDoCancelUpdateOutcome } from "./types";
11 
12const getRunStub = (context: ProjectDoContext, runId: RunId) => context.env.RUN_DO.getByName(runId);
13const now = (): UnixTimestampMs => expectTrusted(UnixTimestampMs, Date.now(), "UnixTimestampMs");
14 
15const repairActiveRunDoStep = async (
16 context: ProjectDoContext,
17 runId: RunId,
18 finishedAt: UnixTimestampMs,
19 terminalStatus: ProjectRunTerminalStatus,
20): Promise<void> => {
21 if (terminalStatus === "passed") {
22 return;
23 }
24 
25 const runStub = getRunStub(context, runId);
26 
27 try {
28 const detail = await runStub.getRunDetail(runId);
29 const position = detail.meta?.currentStep ?? null;
30 if (position === null) {
31 return;
32 }
33 
34 const step = detail.steps.find((entry) => entry.position === position);
35 if (!step || (step.status !== "queued" && step.status !== "running")) {
36 return;
37 }
38 
39 await runStub.updateStepState({
40 runId,
41 position,
42 status: "failed",
43 finishedAt,
44 exitCode: terminalStatus === "failed" ? 1 : null,
45 });
46 } catch (error) {
47 context.logger.warn("run_do_active_step_repair_failed", {
48 runId,
49 terminalStatus,
50 error: error instanceof Error ? error.message : String(error),
51 });
52 }
53};
54 
55export const ensureRunInitializedWithPayload = async (
56 context: ProjectDoContext,
57 payload: EnsureRunInput,
58): Promise<void> => {
59 await getRunStub(context, payload.runId).ensureInitialized(payload);
60};
61 
62export const ensureRunInitialized = async (context: ProjectDoContext, row: ProjectRunRow): Promise<void> => {
63 const runId = expectTrusted(RunId, row.runId, "RunId");
64 const payload: EnsureRunInput = {
65 runId,
66 projectId: expectTrusted(ProjectId, row.projectId, "ProjectId"),
67 triggerType: expectTrusted(TriggerType, row.triggerType, "TriggerType"),
68 branch: expectTrusted(BranchName, row.branch, "BranchName"),
69 commitSha: nullableTrusted(CommitSha, row.commitSha, "CommitSha"),
70 };
71 
72 await ensureRunInitializedWithPayload(context, payload);
73};
74 
75export const updateRunDoCancelRequested = async (
76 context: ProjectDoContext,
77 row: ProjectRunRow,
78): Promise<RunDoCancelUpdateOutcome> => {
79 const runId = expectTrusted(RunId, row.runId, "RunId");
80 const runStub = getRunStub(context, runId);
81 await ensureRunInitialized(context, row);
82 
83 const current = await runStub.getRunSummary(runId);
84 if (!current) {
85 return "deferred";
86 }
87 
88 if (current.status === "passed" || current.status === "failed" || current.status === "canceled") {
89 return "noop";
90 }
91 
92 if (current.status === "queued" || current.status === "cancel_requested" || current.status === "canceling") {
93 return current.status === "queued" ? "deferred" : "noop";
94 }
95 
96 const update: UpdateRunStateInput = {
97 runId,
98 status: "cancel_requested",
99 currentStep: current.currentStep,
100 startedAt: current.startedAt,
101 finishedAt: current.finishedAt,
102 exitCode: current.exitCode,
103 errorMessage: current.errorMessage,
104 };
105 await runStub.updateRunState(update);
106 return "applied";
107};
108 
109export const setRunDoTerminal = async (
110 context: ProjectDoContext,
111 row: ProjectRunRow,
112 terminalStatus: ProjectRunTerminalStatus,
113 errorMessage: string | null,
114): Promise<void> => {
115 const runId = expectTrusted(RunId, row.runId, "RunId");
116 const runStub = getRunStub(context, runId);
117 const finishedAt = now();
118 await ensureRunInitialized(context, row);
119 
120 const current = await runStub.getRunSummary(runId);
121 if (current && (current.status === "passed" || current.status === "failed" || current.status === "canceled")) {
122 return;
123 }
124 
125 await repairActiveRunDoStep(context, runId, finishedAt, terminalStatus);
126 
127 if (terminalStatus === "canceled") {
128 if (!current) {
129 throw new Error(`Run ${runId} is missing RunDO metadata.`);
130 }
131 
132 let cancelState = current;
133 if (cancelState.status === "starting" || cancelState.status === "running") {
134 await runStub.updateRunState({
135 runId,
136 status: "cancel_requested",
137 currentStep: cancelState.currentStep,
138 startedAt: cancelState.startedAt ?? now(),
139 finishedAt: cancelState.finishedAt,
140 exitCode: cancelState.exitCode,
141 errorMessage: cancelState.errorMessage,
142 });
143 const updatedCancelState = await runStub.getRunSummary(runId);
144 if (!updatedCancelState) {
145 throw new Error(`Run ${runId} is missing RunDO metadata after cancel request.`);
146 }
147 cancelState = updatedCancelState;
148 }
149 
150 if (
151 cancelState.status !== "queued" &&
152 cancelState.status !== "cancel_requested" &&
153 cancelState.status !== "canceling"
154 ) {
155 throw new Error(`Run ${runId} cannot be canceled from RunDO status ${cancelState.status}.`);
156 }
157 
158 await runStub.updateRunState({
159 runId,
160 status: "canceled",
161 currentStep: null,
162 startedAt: cancelState.status === "queued" ? cancelState.startedAt : (cancelState.startedAt ?? finishedAt),
163 finishedAt,
164 exitCode: null,
165 errorMessage,
166 });
167 return;
168 }
169 
170 if (current?.status === "cancel_requested" && terminalStatus === "failed") {
171 await runStub.updateRunState({
172 runId,
173 status: "canceling",
174 currentStep: current.currentStep,
175 startedAt: current.startedAt,
176 finishedAt: current.finishedAt,
177 exitCode: current.exitCode,
178 errorMessage: current.errorMessage,
179 });
180 }
181 
182 const update: UpdateRunStateInput = {
183 runId,
184 status: terminalStatus,
185 currentStep: null,
186 startedAt: current?.startedAt ?? null,
187 finishedAt,
188 exitCode: terminalStatus === "passed" ? (current?.exitCode ?? 0) : (current?.exitCode ?? 1),
189 errorMessage,
190 };
191 
192 const updateResult = await runStub.tryUpdateRunState(update);
193 if (updateResult.kind === "conflict") {
194 await runStub.repairTerminalState(update);
195 }
196};