Skip to content
File

Blob: src/worker/dispatch/workflows/run-workflows.ts

typescript127 lines
1import { WorkflowEntrypoint, type WorkflowEvent, type WorkflowStep } from "cloudflare:workers";
2 
3import { type UnixTimestampMs as UnixTimestampMsType } from "@/contracts";
4import {
5 AcceptedRunSnapshot as AcceptedRunSnapshotCodec,
6 type RecoverWorkflowDispatchFailureResult,
7 isTerminalStatus,
8 type AcceptedRunSnapshot,
9} from "@/worker/contracts";
10import { recoverTerminalActiveRun } from "@/worker/dispatch/shared";
11import { claimWorkflowRun, executeWorkflowRun, finalizeWorkflowRun } from "./steps/index";
12import type { WorkflowRunResult } from "./types";
13import { boundedRetryStepConfig } from "./step-config";
14import { toWorkflowStartedAt } from "./execution";
15 
16const finalizeFailedWorkflowExecution = async (
17 env: Env,
18 snapshot: AcceptedRunSnapshot,
19 startedAt: UnixTimestampMsType,
20 error: unknown,
21): Promise<WorkflowRunResult> => {
22 const runMeta = await env.RUN_DO.getByName(snapshot.runId).getRunSummary(snapshot.runId);
23 await finalizeWorkflowRun(
24 env,
25 snapshot,
26 startedAt,
27 {
28 kind: "failed",
29 exitCode: 1,
30 errorMessage: error instanceof Error ? error.message : String(error),
31 },
32 runMeta?.currentStep ?? null,
33 );
34 
35 const finalizedMeta = await env.RUN_DO.getByName(snapshot.runId).getRunSummary(snapshot.runId);
36 return {
37 kind: "executed",
38 terminalStatus:
39 finalizedMeta && isTerminalStatus(finalizedMeta.status) && finalizedMeta.finishedAt !== null
40 ? finalizedMeta.status
41 : "failed",
42 };
43};
44 
45export class RunWorkflows extends WorkflowEntrypoint<Env, AcceptedRunSnapshot> {
46 async run(event: WorkflowEvent<AcceptedRunSnapshot>, step: WorkflowStep): Promise<WorkflowRunResult> {
47 const snapshot = AcceptedRunSnapshotCodec.assertDecode(event.payload);
48 let startedAt: UnixTimestampMsType | null = null;
49 
50 try {
51 const claimResult = await claimWorkflowRun(step, this.env, snapshot);
52 
53 if (claimResult.kind === "recovered") {
54 return claimResult;
55 }
56 
57 if (claimResult.kind === "stale") {
58 return claimResult;
59 }
60 
61 startedAt = claimResult.startedAt;
62 
63 return {
64 kind: "executed",
65 terminalStatus: await executeWorkflowRun(step, this.env, snapshot, claimResult.startedAt),
66 };
67 } catch (error) {
68 let executionError: unknown = error;
69 // These catch-path repairs only reconcile durable DO state, so replaying them is safe
70 // without introducing extra workflow steps.
71 
72 if (startedAt === null) {
73 const recovery = await step.do(
74 "rearm dispatch",
75 boundedRetryStepConfig(),
76 async (): Promise<RecoverWorkflowDispatchFailureResult> =>
77 await this.env.PROJECT_DO.getByName(snapshot.projectId).recoverWorkflowDispatchFailure({
78 projectId: snapshot.projectId,
79 runId: snapshot.runId,
80 errorMessage: executionError instanceof Error ? executionError.message : String(executionError),
81 }),
82 );
83 
84 if (recovery.kind === "rearmed") {
85 return {
86 kind: "stale",
87 reason: "dispatch_rearmed",
88 };
89 }
90 
91 if (recovery.kind === "already_active") {
92 const runMeta = await this.env.RUN_DO.getByName(snapshot.runId).getRunSummary(snapshot.runId);
93 if (runMeta && isTerminalStatus(runMeta.status) && runMeta.finishedAt !== null) {
94 if (await recoverTerminalActiveRun(this.env, snapshot.projectId, snapshot.runId)) {
95 return {
96 kind: "recovered",
97 };
98 }
99 
100 return {
101 kind: "stale",
102 reason: "already_terminal",
103 };
104 }
105 
106 startedAt = runMeta?.startedAt ?? toWorkflowStartedAt(Date.now());
107 try {
108 return {
109 kind: "executed",
110 terminalStatus: await executeWorkflowRun(step, this.env, snapshot, startedAt),
111 };
112 } catch (resumeError) {
113 return await finalizeFailedWorkflowExecution(this.env, snapshot, startedAt, resumeError);
114 }
115 }
116 
117 return {
118 kind: "stale",
119 reason: recovery.kind,
120 };
121 }
122 
123 return await finalizeFailedWorkflowExecution(this.env, snapshot, startedAt, executionError);
124 }
125 }
126}