Skip to content
File

Blob: tests/worker/dispatch/workflows/run-workflows.test.ts

typescript361 lines
1import type { WorkflowStep } from "cloudflare:workers";
2import { describe, expect, it, vi } from "vitest";
3 
4import { UnixTimestampMs } from "@/contracts";
5import { AcceptedRunSnapshot, PositiveInteger, expectTrusted } from "@/worker/contracts";
6import * as workflowSteps from "@/worker/dispatch/workflows/steps/index";
7import { RunWorkflows } from "@/worker/dispatch/workflows";
8 
9const WORKFLOW_NAME = "anvil-run-workflows";
10 
11vi.mock("@/worker/dispatch/workflows/steps/index", async () => {
12 const actual = await vi.importActual<typeof import("@/worker/dispatch/workflows/steps/index")>(
13 "@/worker/dispatch/workflows/steps/index",
14 );
15 
16 return {
17 ...actual,
18 claimWorkflowRun: vi.fn(actual.claimWorkflowRun),
19 executeWorkflowRun: vi.fn(actual.executeWorkflowRun),
20 finalizeWorkflowRun: vi.fn(actual.finalizeWorkflowRun),
21 };
22});
23 
24const createSnapshot = (payload: Parameters<typeof AcceptedRunSnapshot.assertDecode>[0]) =>
25 AcceptedRunSnapshot.assertDecode(payload);
26 
27const createWorkflowEvent = (payload: Parameters<typeof AcceptedRunSnapshot.assertDecode>[0]) => ({
28 workflowName: WORKFLOW_NAME,
29 payload: createSnapshot(payload),
30 timestamp: new Date("2026-03-24T12:00:00.000Z"),
31 instanceId: "run_0000000000000000000000",
32});
33 
34const createWorkflowStep = (
35 handler?: (
36 name: string,
37 callback: (ctx: { attempt: number }) => Promise<unknown>,
38 config: unknown,
39 ) => Promise<unknown>,
40) => {
41 const doMock = vi.fn(async (...args: unknown[]) => {
42 const name = args[0] as string;
43 const config = typeof args[1] === "function" ? undefined : args[1];
44 const callback = typeof args[1] === "function" ? args[1] : args[2];
45 if (handler) {
46 return await handler(name, callback as (ctx: { attempt: number }) => Promise<unknown>, config);
47 }
48 return await (callback as (ctx: { attempt: number }) => Promise<unknown>)({ attempt: 1 });
49 });
50 
51 return {
52 step: {
53 do: doMock,
54 } as unknown as WorkflowStep,
55 doMock,
56 };
57};
58 
59const createWorkflow = (env: Env): RunWorkflows =>
60 Object.assign(Object.create(RunWorkflows.prototype), {
61 env,
62 ctx: {},
63 }) as RunWorkflows;
64 
65const toEnv = (env: unknown): Env => env as Env;
66 
67const WORKFLOW_SNAPSHOT_INPUT = {
68 runId: "run_0000000000000000000000",
69 projectId: "prj_0000000000000000000000",
70 triggerType: "manual" as const,
71 triggeredByUserId: null,
72 branch: "main",
73 commitSha: null,
74 repoUrl: "https://github.com/example/anvil",
75 configPath: ".anvil.yml",
76 dispatchMode: "workflows" as const,
77 executionRuntime: "cloudflare_sandbox" as const,
78 queuedAt: 1_711_111_111_111,
79};
80 
81describe("run workflows", () => {
82 it("returns stale when ProjectDO rejects the run claim", async () => {
83 const projectStub = {
84 claimRunWork: vi.fn(async () => ({
85 kind: "stale" as const,
86 reason: "run_missing" as const,
87 })),
88 };
89 const workflow = createWorkflow(
90 toEnv({
91 PROJECT_DO: {
92 getByName: () => projectStub,
93 },
94 RUN_DO: {
95 getByName: () => ({}),
96 },
97 }),
98 );
99 
100 const result = await workflow.run(createWorkflowEvent(WORKFLOW_SNAPSHOT_INPUT), createWorkflowStep().step);
101 
102 expect(result).toEqual({
103 kind: "stale",
104 reason: "run_missing",
105 });
106 });
107 
108 it("recovers an already-terminal active run during claim", async () => {
109 const snapshot = createSnapshot(WORKFLOW_SNAPSHOT_INPUT);
110 const finalizeRunExecution = vi.fn(async () => ({
111 snapshot,
112 }));
113 const workflow = createWorkflow(
114 toEnv({
115 PROJECT_DO: {
116 getByName: () => ({
117 claimRunWork: vi.fn(async () => ({
118 kind: "stale" as const,
119 reason: "run_active" as const,
120 })),
121 finalizeRunExecution,
122 kickReconciliation: vi.fn(async () => {}),
123 }),
124 },
125 RUN_DO: {
126 getByName: () => ({
127 getRunSummary: vi.fn(async () => ({
128 status: "passed" as const,
129 currentStep: null,
130 startedAt: snapshot.queuedAt,
131 finishedAt: snapshot.queuedAt + 1_000,
132 exitCode: 0,
133 errorMessage: null,
134 })),
135 }),
136 },
137 }),
138 );
139 
140 const result = await workflow.run(createWorkflowEvent(WORKFLOW_SNAPSHOT_INPUT), createWorkflowStep().step);
141 
142 expect(result).toEqual({
143 kind: "recovered",
144 });
145 expect(finalizeRunExecution).toHaveBeenCalledWith({
146 projectId: snapshot.projectId,
147 runId: snapshot.runId,
148 terminalStatus: "passed",
149 lastError: null,
150 sandboxDestroyed: false,
151 });
152 });
153 
154 it("rearms dispatch when claim fails before the run becomes active", async () => {
155 const recoverWorkflowDispatchFailure = vi.fn(async () => ({
156 kind: "rearmed" as const,
157 }));
158 const { step, doMock } = createWorkflowStep(async (name, callback) => {
159 if (name === "claim run") {
160 throw new Error("claim failed");
161 }
162 
163 return await callback({ attempt: 1 });
164 });
165 const workflow = createWorkflow(
166 toEnv({
167 PROJECT_DO: {
168 getByName: () => ({
169 recoverWorkflowDispatchFailure,
170 }),
171 },
172 RUN_DO: {
173 getByName: () => ({}),
174 },
175 }),
176 );
177 
178 const result = await workflow.run(createWorkflowEvent(WORKFLOW_SNAPSHOT_INPUT), step);
179 
180 expect(result).toEqual({
181 kind: "stale",
182 reason: "dispatch_rearmed",
183 });
184 expect(recoverWorkflowDispatchFailure).toHaveBeenCalledWith({
185 projectId: WORKFLOW_SNAPSHOT_INPUT.projectId,
186 runId: WORKFLOW_SNAPSHOT_INPUT.runId,
187 errorMessage: "claim failed",
188 });
189 await expect(recoverWorkflowDispatchFailure.mock.results[0]?.value).resolves.toEqual({
190 kind: "rearmed",
191 });
192 expect(doMock.mock.calls.map(([name]) => name)).toEqual(["claim run", "rearm dispatch"]);
193 expect((doMock.mock.calls[0]?.[1] as { retries?: { limit: number } }).retries?.limit).toBe(3);
194 expect((doMock.mock.calls[1]?.[1] as { retries?: { limit: number } }).retries?.limit).toBe(3);
195 });
196 
197 it("continues execution when pre-start recovery finds the run is already active", async () => {
198 const snapshot = createSnapshot(WORKFLOW_SNAPSHOT_INPUT);
199 const recoverWorkflowDispatchFailure = vi.fn(async () => ({
200 kind: "already_active" as const,
201 }));
202 const { step, doMock } = createWorkflowStep(async (name, callback) => {
203 if (name === "claim run") {
204 throw new Error("claim failed after activation");
205 }
206 
207 if (name === "execute run") {
208 return "passed";
209 }
210 
211 return await callback({ attempt: 1 });
212 });
213 const workflow = createWorkflow(
214 toEnv({
215 PROJECT_DO: {
216 getByName: () => ({
217 recoverWorkflowDispatchFailure,
218 }),
219 },
220 RUN_DO: {
221 getByName: () => ({
222 getRunSummary: vi.fn(async () => ({
223 status: "running" as const,
224 currentStep: null,
225 startedAt: snapshot.queuedAt,
226 finishedAt: null,
227 exitCode: null,
228 errorMessage: null,
229 })),
230 }),
231 },
232 }),
233 );
234 
235 const result = await workflow.run(createWorkflowEvent(WORKFLOW_SNAPSHOT_INPUT), step);
236 
237 expect(result).toEqual({
238 kind: "executed",
239 terminalStatus: "passed",
240 });
241 expect(doMock.mock.calls.map(([name]) => name)).toEqual(["claim run", "rearm dispatch", "execute run"]);
242 });
243 
244 it("finalizes a failed resumed execution after pre-start recovery finds the run already active", async () => {
245 const snapshot = createSnapshot(WORKFLOW_SNAPSHOT_INPUT);
246 const recoverWorkflowDispatchFailure = vi.fn(async () => ({
247 kind: "already_active" as const,
248 }));
249 const finalizeWorkflowRun = vi.mocked(workflowSteps.finalizeWorkflowRun).mockResolvedValueOnce();
250 const executeWorkflowRun = vi
251 .mocked(workflowSteps.executeWorkflowRun)
252 .mockRejectedValueOnce(new Error("resume failed"));
253 const runningMeta = {
254 status: "running" as const,
255 currentStep: expectTrusted(PositiveInteger, 2, "PositiveInteger"),
256 startedAt: snapshot.queuedAt,
257 finishedAt: null,
258 exitCode: null,
259 errorMessage: null,
260 };
261 const finalizedMeta = {
262 status: "failed" as const,
263 currentStep: null,
264 startedAt: snapshot.queuedAt,
265 finishedAt: expectTrusted(UnixTimestampMs, snapshot.queuedAt + 1_000, "UnixTimestampMs"),
266 exitCode: 1,
267 errorMessage: "resume failed",
268 };
269 const getRunSummary = vi
270 .fn(async (): Promise<typeof runningMeta | typeof finalizedMeta> => runningMeta)
271 .mockResolvedValueOnce(runningMeta)
272 .mockResolvedValueOnce(runningMeta)
273 .mockResolvedValueOnce(finalizedMeta);
274 const { step, doMock } = createWorkflowStep(async (name, callback) => {
275 if (name === "claim run") {
276 throw new Error("claim failed after activation");
277 }
278 
279 return await callback({ attempt: 1 });
280 });
281 const workflow = createWorkflow(
282 toEnv({
283 PROJECT_DO: {
284 getByName: () => ({
285 recoverWorkflowDispatchFailure,
286 }),
287 },
288 RUN_DO: {
289 getByName: () => ({
290 getRunSummary,
291 }),
292 },
293 }),
294 );
295 
296 const result = await workflow.run(createWorkflowEvent(WORKFLOW_SNAPSHOT_INPUT), step);
297 
298 expect(result).toEqual({
299 kind: "executed",
300 terminalStatus: "failed",
301 });
302 expect(executeWorkflowRun).toHaveBeenCalledWith(step, expect.anything(), snapshot, snapshot.queuedAt);
303 expect(finalizeWorkflowRun).toHaveBeenCalledWith(
304 expect.anything(),
305 snapshot,
306 snapshot.queuedAt,
307 {
308 kind: "failed",
309 exitCode: 1,
310 errorMessage: "resume failed",
311 },
312 runningMeta.currentStep,
313 );
314 expect(doMock.mock.calls.map(([name]) => name)).toEqual(["claim run", "rearm dispatch"]);
315 });
316 
317 it("returns the terminal status from the execute step on the happy path", async () => {
318 const snapshot = createSnapshot(WORKFLOW_SNAPSHOT_INPUT);
319 const { step, doMock } = createWorkflowStep(async (name, callback) => {
320 if (name === "execute run") {
321 return "failed";
322 }
323 
324 return await callback({ attempt: 1 });
325 });
326 const workflow = createWorkflow(
327 toEnv({
328 PROJECT_DO: {
329 getByName: () => ({
330 claimRunWork: vi.fn(async () => ({
331 kind: "execute" as const,
332 snapshot,
333 })),
334 }),
335 },
336 RUN_DO: {
337 getByName: () => ({
338 ensureInitialized: vi.fn(async () => {}),
339 }),
340 },
341 }),
342 );
343 
344 const result = await workflow.run(
345 {
346 workflowName: WORKFLOW_NAME,
347 payload: snapshot,
348 timestamp: new Date("2026-03-24T12:00:00.000Z"),
349 instanceId: snapshot.runId,
350 },
351 step,
352 );
353 
354 expect(result).toEqual({
355 kind: "executed",
356 terminalStatus: "failed",
357 });
358 expect(doMock.mock.calls.map(([name]) => name)).toEqual(["claim run", "execute run"]);
359 });
360});