File
Blob: src/worker/durable/project-do/reconciliation/dispatch-helper.ts
| 1 | import { type AcceptedRunSnapshot } from "@/worker/contracts"; |
| 2 | import { type DispatchMode, type ProjectId, type RunId } from "@/contracts"; |
| 3 | import type { ProjectDoContext } from "../types"; |
| 4 | |
| 5 | type WorkflowInstanceStatus = InstanceStatus["status"]; |
| 6 | type WorkflowDispatchState = "already_dispatched" | "restartable" | "unsupported"; |
| 7 | |
| 8 | export 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 | |
| 29 | export 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 | }; |