File
Blob: src/worker/dispatch/workflows/steps/claim.ts
| 1 | import { type WorkflowStep } from "cloudflare:workers"; |
| 2 | |
| 3 | import { type AcceptedRunSnapshot, isTerminalStatus } from "@/worker/contracts"; |
| 4 | import { recoverTerminalActiveRun } from "@/worker/dispatch/shared/active-run-recovery"; |
| 5 | import { toWorkflowStartedAt } from "../execution"; |
| 6 | import { boundedRetryStepConfig } from "../step-config"; |
| 7 | import type { WorkflowClaimResult } from "../types"; |
| 8 | |
| 9 | export const claimWorkflowRun = async ( |
| 10 | step: WorkflowStep, |
| 11 | env: Env, |
| 12 | snapshot: AcceptedRunSnapshot, |
| 13 | ): Promise<WorkflowClaimResult> => |
| 14 | await step.do("claim run", boundedRetryStepConfig(), async (): Promise<WorkflowClaimResult> => { |
| 15 | const projectStub = env.PROJECT_DO.getByName(snapshot.projectId); |
| 16 | const claim = await projectStub.claimRunWork({ |
| 17 | projectId: snapshot.projectId, |
| 18 | runId: snapshot.runId, |
| 19 | }); |
| 20 | |
| 21 | if (claim.kind === "execute") { |
| 22 | return { |
| 23 | kind: "claimed", |
| 24 | startedAt: toWorkflowStartedAt(Date.now()), |
| 25 | }; |
| 26 | } |
| 27 | |
| 28 | if (claim.reason === "run_active") { |
| 29 | const current = await env.RUN_DO.getByName(snapshot.runId).getRunSummary(snapshot.runId); |
| 30 | if (current && isTerminalStatus(current.status) && current.finishedAt !== null) { |
| 31 | if (await recoverTerminalActiveRun(env, snapshot.projectId, snapshot.runId)) { |
| 32 | return { |
| 33 | kind: "recovered", |
| 34 | }; |
| 35 | } |
| 36 | |
| 37 | return { |
| 38 | kind: "stale", |
| 39 | reason: "already_terminal", |
| 40 | }; |
| 41 | } |
| 42 | |
| 43 | return { |
| 44 | kind: "claimed", |
| 45 | startedAt: current?.startedAt ?? toWorkflowStartedAt(Date.now()), |
| 46 | }; |
| 47 | } |
| 48 | |
| 49 | return { |
| 50 | kind: "stale", |
| 51 | reason: claim.reason, |
| 52 | }; |
| 53 | }); |