Skip to content
File

Blob: src/worker/durable/project-do/reconciliation/dispatch.ts

typescript304 lines
1import { and, eq } from "drizzle-orm";
2 
3import { DispatchMode, type ProjectId, RunId } from "@/contracts";
4import { D1SyncStatus, expectTrusted } from "@/worker/contracts";
5import * as projectSchema from "@/worker/db/durable/schema/project-do";
6 
7import { DISPATCH_RETRY_DELAYS_MS } from "../constants";
8import { getNextDispatchableRun, getOldestCancelRequestedRun, getRunRow } from "../repo";
9import { setRunDoTerminal, updateRunDoCancelRequested } from "../run-do-sync";
10import { getDispatchRetryAt, setDispatchRetryAt } from "../sidecar-state";
11import { classifyWorkflowDispatchState, dispatchRun } from "./dispatch-helper";
12import { getRetryDelay } from "./shared";
13import { getSnapshot, nextTerminalD1SyncStatus, promoteNextPendingRun } from "../transitions/shared";
14import type { ProjectDoContext } from "../types";
15 
16type WorkflowDispatchCompletionState = "still_queued" | "rearmed" | "left_executable" | "missing";
17type DispatchFailureResult =
18 | { kind: "stale" }
19 | { kind: "retry"; nextAt: number }
20 | { kind: "terminal"; row: NonNullable<ReturnType<typeof getRunRow>> };
21 
22const markWorkflowRunQueued = (context: ProjectDoContext, projectId: ProjectId, runId: RunId): boolean =>
23 context.db.transaction((tx) => {
24 const currentRow = getRunRow(tx, projectId, runId);
25 if (!currentRow || currentRow.status !== "executable" || currentRow.dispatchStatus !== "pending") {
26 return false;
27 }
28 
29 tx.update(projectSchema.projectRuns)
30 .set({
31 dispatchStatus: "queued",
32 lastError: null,
33 })
34 .where(eq(projectSchema.projectRuns.runId, currentRow.runId))
35 .run();
36 return true;
37 });
38 
39const getWorkflowDispatchCompletionState = (
40 context: ProjectDoContext,
41 projectId: ProjectId,
42 runId: RunId,
43): WorkflowDispatchCompletionState =>
44 context.db.transaction((tx) => {
45 const currentRow = getRunRow(tx, projectId, runId);
46 if (!currentRow) {
47 return "missing";
48 }
49 
50 if (currentRow.status === "executable" && currentRow.dispatchStatus === "queued") {
51 return "still_queued";
52 }
53 
54 if (currentRow.status === "executable" && currentRow.dispatchStatus === "pending") {
55 return "rearmed";
56 }
57 
58 return "left_executable";
59 });
60 
61const getQueuedWorkflowRun = async (context: ProjectDoContext, projectId: ProjectId) => {
62 const rows = await context.db
63 .select()
64 .from(projectSchema.projectRuns)
65 .where(
66 and(
67 eq(projectSchema.projectRuns.projectId, projectId),
68 eq(projectSchema.projectRuns.status, "executable"),
69 eq(projectSchema.projectRuns.dispatchStatus, "queued"),
70 eq(projectSchema.projectRuns.dispatchMode, "workflows"),
71 ),
72 )
73 .limit(1);
74 
75 return rows[0];
76};
77 
78const resolveDispatchFailure = async (
79 context: ProjectDoContext,
80 projectId: ProjectId,
81 runId: RunId,
82 dispatchMode: DispatchMode,
83 errorMessage: string,
84): Promise<DispatchFailureResult> =>
85 await context.db.transaction((tx) => {
86 const currentRow = getRunRow(tx, projectId, runId);
87 if (!currentRow || currentRow.status !== "executable") {
88 return {
89 kind: "stale",
90 };
91 }
92 
93 const expectedDispatchStatus = dispatchMode === "workflows" ? "queued" : "pending";
94 if (currentRow.dispatchStatus !== expectedDispatchStatus) {
95 return {
96 kind: "stale",
97 };
98 }
99 
100 const nextAttempt = currentRow.dispatchAttempts + 1;
101 if (nextAttempt > DISPATCH_RETRY_DELAYS_MS.length) {
102 tx.update(projectSchema.projectRuns)
103 .set({
104 status: "failed",
105 position: null,
106 dispatchStatus: "terminal",
107 dispatchAttempts: nextAttempt,
108 d1SyncStatus: nextTerminalD1SyncStatus(expectTrusted(D1SyncStatus, currentRow.d1SyncStatus, "D1SyncStatus")),
109 lastError: "dispatch_failed",
110 })
111 .where(eq(projectSchema.projectRuns.runId, currentRow.runId))
112 .run();
113 promoteNextPendingRun(tx, projectId);
114 
115 return {
116 kind: "terminal",
117 row: currentRow,
118 };
119 }
120 
121 tx.update(projectSchema.projectRuns)
122 .set({
123 dispatchStatus: dispatchMode === "workflows" ? "pending" : currentRow.dispatchStatus,
124 dispatchAttempts: nextAttempt,
125 lastError: errorMessage,
126 })
127 .where(eq(projectSchema.projectRuns.runId, currentRow.runId))
128 .run();
129 
130 return {
131 kind: "retry",
132 nextAt: Date.now() + getRetryDelay(nextAttempt, DISPATCH_RETRY_DELAYS_MS),
133 };
134 });
135 
136const applyDispatchFailure = async (
137 context: ProjectDoContext,
138 projectId: ProjectId,
139 runId: RunId,
140 dispatchMode: DispatchMode,
141 dispatchAttempts: number,
142 errorMessage: string,
143): Promise<DispatchFailureResult["kind"]> => {
144 const result = await resolveDispatchFailure(context, projectId, runId, dispatchMode, errorMessage);
145 
146 if (result.kind === "stale") {
147 return result.kind;
148 }
149 
150 if (result.kind === "terminal") {
151 await setDispatchRetryAt(context, runId, null);
152 try {
153 await setRunDoTerminal(context, result.row, "failed", "dispatch_failed");
154 } catch (runDoError) {
155 context.logger.error("dispatch_failed_run_do_terminalize_failed", {
156 projectId,
157 runId,
158 error: runDoError instanceof Error ? runDoError.message : String(runDoError),
159 });
160 }
161 return result.kind;
162 }
163 
164 await setDispatchRetryAt(context, runId, result.nextAt);
165 context.logger.error("run_dispatch_failed", {
166 projectId,
167 runId,
168 attempt: dispatchAttempts + 1,
169 error: errorMessage,
170 });
171 return result.kind;
172};
173 
174const reconcileQueuedWorkflowRun = async (context: ProjectDoContext, projectId: ProjectId): Promise<void> => {
175 const row = await getQueuedWorkflowRun(context, projectId);
176 if (!row) {
177 return;
178 }
179 
180 const runId = expectTrusted(RunId, row.runId, "RunId");
181 
182 try {
183 const instance = await context.env.RUN_WORKFLOWS.get(runId);
184 const current = await instance.status();
185 
186 switch (classifyWorkflowDispatchState(current.status)) {
187 case "already_dispatched":
188 return;
189 case "restartable":
190 await applyDispatchFailure(
191 context,
192 projectId,
193 runId,
194 "workflows",
195 row.dispatchAttempts,
196 `Workflow instance ${runId} reached terminal status ${current.status} before ProjectDO accepted the dispatch.`,
197 );
198 return;
199 case "unsupported":
200 await applyDispatchFailure(
201 context,
202 projectId,
203 runId,
204 "workflows",
205 row.dispatchAttempts,
206 `Workflow instance ${runId} has unsupported status ${current.status}.`,
207 );
208 return;
209 }
210 } catch (error) {
211 await applyDispatchFailure(
212 context,
213 projectId,
214 runId,
215 "workflows",
216 row.dispatchAttempts,
217 error instanceof Error ? error.message : String(error),
218 );
219 }
220};
221 
222export const reconcileCancelRequestedRunDo = async (
223 context: ProjectDoContext,
224 projectId: ProjectId,
225): Promise<RunId | null> => {
226 const row = await getOldestCancelRequestedRun(context, projectId);
227 if (!row) {
228 return null;
229 }
230 
231 try {
232 const outcome = await updateRunDoCancelRequested(context, row);
233 return outcome === "applied" ? expectTrusted(RunId, row.runId, "RunId") : null;
234 } catch (error) {
235 context.logger.error("cancel_requested_run_do_sync_failed", {
236 projectId,
237 runId: row.runId,
238 error: error instanceof Error ? error.message : String(error),
239 });
240 return null;
241 }
242};
243 
244export const dispatchExecutableRun = async (context: ProjectDoContext, projectId: ProjectId): Promise<RunId | null> => {
245 const row = await getNextDispatchableRun(context, projectId);
246 if (!row) {
247 await reconcileQueuedWorkflowRun(context, projectId);
248 return null;
249 }
250 
251 const runId = expectTrusted(RunId, row.runId, "RunId");
252 const retryAt = await getDispatchRetryAt(context, runId);
253 if (retryAt !== null && retryAt > Date.now()) {
254 return null;
255 }
256 
257 const dispatchMode = expectTrusted(DispatchMode, row.dispatchMode, "DispatchMode");
258 
259 try {
260 if (dispatchMode === "workflows" && !markWorkflowRunQueued(context, projectId, runId)) {
261 return null;
262 }
263 
264 await dispatchRun(context, projectId, runId, dispatchMode, getSnapshot(row));
265 
266 if (dispatchMode === "queue") {
267 context.db.transaction((tx) => {
268 const currentRow = getRunRow(tx, projectId, runId);
269 if (!currentRow || currentRow.status !== "executable" || currentRow.dispatchStatus !== "pending") {
270 return;
271 }
272 
273 tx.update(projectSchema.projectRuns)
274 .set({
275 dispatchStatus: "queued",
276 lastError: null,
277 })
278 .where(eq(projectSchema.projectRuns.runId, currentRow.runId))
279 .run();
280 });
281 await setDispatchRetryAt(context, runId, null);
282 return runId;
283 }
284 
285 const completionState = getWorkflowDispatchCompletionState(context, projectId, runId);
286 if (completionState !== "rearmed") {
287 await setDispatchRetryAt(context, runId, null);
288 }
289 
290 return completionState === "rearmed" ? null : runId;
291 } catch (error) {
292 const errorMessage = error instanceof Error ? error.message : String(error);
293 const failureKind = await applyDispatchFailure(
294 context,
295 projectId,
296 runId,
297 dispatchMode,
298 row.dispatchAttempts,
299 errorMessage,
300 );
301 return failureKind === "terminal" && dispatchMode === "queue" ? runId : null;
302 }
303};