Skip to content
File

Blob: src/worker/dispatch/workflows/steps/claim.ts

typescript54 lines
1import { type WorkflowStep } from "cloudflare:workers";
2 
3import { type AcceptedRunSnapshot, isTerminalStatus } from "@/worker/contracts";
4import { recoverTerminalActiveRun } from "@/worker/dispatch/shared/active-run-recovery";
5import { toWorkflowStartedAt } from "../execution";
6import { boundedRetryStepConfig } from "../step-config";
7import type { WorkflowClaimResult } from "../types";
8 
9export 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 });