Skip to content
File

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

typescript72 lines
1import { type AcceptedRunSnapshot } from "@/worker/contracts";
2import { type DispatchMode, type ProjectId, type RunId } from "@/contracts";
3import type { ProjectDoContext } from "../types";
4 
5type WorkflowInstanceStatus = InstanceStatus["status"];
6type WorkflowDispatchState = "already_dispatched" | "restartable" | "unsupported";
7 
8export const classifyWorkflowDispatchState = (status: WorkflowInstanceStatus): WorkflowDispatchState => {
9 switch (status) {
10 case "queued":
11 case "running":
12 case "waiting":
13 case "paused":
14 case "waitingForPause":
15 return "already_dispatched";
16 case "complete":
17 case "errored":
18 case "terminated":
19 return "restartable";
20 case "unknown":
21 return "unsupported";
22 default: {
23 const _exhaustive: never = status;
24 throw new Error(`Unhandled workflow status: ${String(_exhaustive)}`);
25 }
26 }
27};
28 
29export const dispatchRun = async (
30 context: ProjectDoContext,
31 projectId: ProjectId,
32 runId: RunId,
33 dispatchMode: DispatchMode,
34 snapshot: AcceptedRunSnapshot,
35): Promise<void> => {
36 switch (dispatchMode) {
37 case "queue":
38 await context.env.RUN_QUEUE.send({ projectId, runId });
39 return;
40 case "workflows":
41 if (
42 (
43 await context.env.RUN_WORKFLOWS.createBatch([
44 {
45 id: runId,
46 params: snapshot,
47 },
48 ])
49 ).length === 0
50 ) {
51 const instance = await context.env.RUN_WORKFLOWS.get(runId);
52 const current = await instance.status();
53 switch (classifyWorkflowDispatchState(current.status)) {
54 case "already_dispatched":
55 return;
56 case "restartable":
57 await instance.restart();
58 return;
59 case "unsupported":
60 break;
61 }
62 
63 throw new Error(`Workflow instance ${runId} has unsupported status ${current.status}.`);
64 }
65 return;
66 default: {
67 const _exhaustive: never = dispatchMode;
68 throw new Error(`Unknown dispatch mode: ${String(_exhaustive)}`);
69 }
70 }
71};