Skip to content
File

Blob: tests/worker/dispatch/workflows/execute.test.ts

typescript307 lines
1import { type WorkflowStep } from "cloudflare:workers";
2import { beforeEach, describe, expect, it, vi } from "vitest";
3 
4import { UnixTimestampMs } from "@/contracts";
5import { AcceptedRunSnapshot, PositiveInteger, expectTrusted, type RunMetaState } from "@/worker/contracts";
6import { executeWorkflowRun } from "@/worker/dispatch/workflows/steps/execute";
7 
8import {
9 createQueueLeaseStub,
10 createQueueRunControlStub,
11 createQueueRunLogsStub,
12 createQueueRunRuntimeStub,
13 createQueueRunStoreStub,
14 createQueueScope,
15 createQueueState,
16} from "../../../helpers/dispatch/shared";
17 
18const mockedModules = vi.hoisted(() => {
19 const state = {
20 context: null as any,
21 lease: null as any,
22 createRunExecutionContext: vi.fn(() => state.context),
23 ensureRunInitialized: vi.fn(async () => {}),
24 logger: {
25 warn: vi.fn(),
26 },
27 appendFailureLogBestEffort: vi.fn(async () => {}),
28 executeRunSteps: vi.fn(async () => ({
29 kind: "passed" as const,
30 exitCode: 0,
31 })),
32 finalizeExecution: vi.fn(async () => {}),
33 mapExecutionErrorToOutcome: vi.fn(async (_context: unknown, error: unknown) => ({
34 kind: "failed" as const,
35 exitCode: 1,
36 errorMessage: error instanceof Error ? error.message : String(error),
37 })),
38 prepareExecutionEnvironment: vi.fn(async () => ({
39 repoConfig: {
40 version: 1,
41 checkout: {
42 depth: 1,
43 },
44 run: {
45 workingDirectory: ".",
46 timeoutSeconds: 60,
47 steps: [],
48 },
49 },
50 workingDirectory: "/workspace/repo",
51 })),
52 recoverTerminalActiveRun: vi.fn(async () => true),
53 RunLease: vi.fn(function MockRunLease() {
54 return state.lease;
55 }),
56 };
57 
58 return state;
59});
60 
61vi.mock("@/worker/dispatch/shared/run-execution-context", () => ({
62 createRunExecutionContext: mockedModules.createRunExecutionContext,
63 ensureRunInitialized: mockedModules.ensureRunInitialized,
64 logger: mockedModules.logger,
65}));
66 
67vi.mock("@/worker/dispatch/shared", () => ({
68 appendFailureLogBestEffort: mockedModules.appendFailureLogBestEffort,
69 executeRunSteps: mockedModules.executeRunSteps,
70 finalizeExecution: mockedModules.finalizeExecution,
71 mapExecutionErrorToOutcome: mockedModules.mapExecutionErrorToOutcome,
72 prepareExecutionEnvironment: mockedModules.prepareExecutionEnvironment,
73 recoverTerminalActiveRun: mockedModules.recoverTerminalActiveRun,
74 RunLease: mockedModules.RunLease,
75}));
76 
77const STARTED_AT = expectTrusted(UnixTimestampMs, 1_740_000_000_000, "UnixTimestampMs");
78const FINISHED_AT = expectTrusted(UnixTimestampMs, 1_740_000_001_000, "UnixTimestampMs");
79 
80const SNAPSHOT = AcceptedRunSnapshot.assertDecode({
81 runId: "run_0000000000000000000000",
82 projectId: "prj_0000000000000000000000",
83 triggerType: "manual",
84 triggeredByUserId: null,
85 branch: "main",
86 commitSha: null,
87 repoUrl: "https://github.com/example/anvil",
88 configPath: ".anvil.yml",
89 dispatchMode: "workflows",
90 executionRuntime: "cloudflare_sandbox",
91 queuedAt: STARTED_AT,
92});
93 
94const createWorkflowStep = (): WorkflowStep =>
95 ({
96 do: vi.fn(async (...args: unknown[]) => {
97 const callback = typeof args[1] === "function" ? args[1] : args[2];
98 return await (callback as () => Promise<unknown>)();
99 }),
100 }) as unknown as WorkflowStep;
101 
102const createRunMeta = (status: RunMetaState["status"]): RunMetaState => ({
103 runId: SNAPSHOT.runId,
104 projectId: SNAPSHOT.projectId,
105 status,
106 triggerType: SNAPSHOT.triggerType,
107 branch: SNAPSHOT.branch,
108 commitSha: SNAPSHOT.commitSha,
109 currentStep: null,
110 startedAt: status === "queued" ? null : STARTED_AT,
111 finishedAt: status === "passed" || status === "failed" || status === "canceled" ? FINISHED_AT : null,
112 exitCode: status === "passed" ? 0 : null,
113 errorMessage: null,
114});
115 
116const createEnv = (runSummary: RunMetaState | null = null): Env =>
117 ({
118 PROJECT_DO: {
119 getByName: () => ({
120 getProjectExecutionMaterial: vi.fn(async () => ({
121 projectId: SNAPSHOT.projectId,
122 encryptedRepoToken: null,
123 })),
124 }),
125 },
126 RUN_DO: {
127 getByName: () => ({
128 getRunSummary: vi.fn(async () => runSummary),
129 }),
130 },
131 }) as unknown as Env;
132 
133describe("workflow execute step", () => {
134 beforeEach(() => {
135 vi.clearAllMocks();
136 
137 const state = createQueueState();
138 mockedModules.context = {
139 scope: createQueueScope({
140 env: createEnv(),
141 startedAt: STARTED_AT,
142 }),
143 state,
144 runStore: createQueueRunStoreStub({
145 getMeta: vi.fn(async () => createRunMeta("queued")),
146 }),
147 logs: createQueueRunLogsStub(),
148 runtime: createQueueRunRuntimeStub(),
149 control: createQueueRunControlStub(),
150 projectControl: {
151 getFreshStub: vi.fn(),
152 recordHeartbeat: vi.fn(async () => ({
153 status: "active" as const,
154 cancelRequestedAt: null,
155 })),
156 recordResolvedCommit: vi.fn(async () => ({ kind: "applied" as const })),
157 finalizeRunExecution: vi.fn(async () => ({
158 snapshot: SNAPSHOT,
159 })),
160 kickReconciliation: vi.fn(async () => {}),
161 },
162 };
163 mockedModules.lease = {
164 start: vi.fn(),
165 ...createQueueLeaseStub(),
166 };
167 });
168 
169 it("keeps the queued path unchanged and enters starting before execution", async () => {
170 const result = await executeWorkflowRun(
171 createWorkflowStep(),
172 mockedModules.context.scope.env,
173 SNAPSHOT,
174 STARTED_AT,
175 );
176 
177 expect(result).toBe("passed");
178 expect(mockedModules.context.runStore.updateState).toHaveBeenCalledWith({
179 status: "starting",
180 startedAt: STARTED_AT,
181 currentStep: null,
182 finishedAt: null,
183 exitCode: null,
184 errorMessage: null,
185 });
186 expect(mockedModules.context.runtime.destroySandbox).toHaveBeenCalledTimes(1);
187 expect(mockedModules.prepareExecutionEnvironment).toHaveBeenCalledWith(mockedModules.context, mockedModules.lease, {
188 executionSessionId: `run-${SNAPSHOT.runId}`,
189 });
190 });
191 
192 it("configures the workflow execute step with a 30 minute timeout", async () => {
193 const doMock = vi.fn(async (...args: unknown[]) => {
194 const callback = typeof args[1] === "function" ? args[1] : args[2];
195 return await (callback as () => Promise<unknown>)();
196 });
197 
198 await executeWorkflowRun(
199 { do: doMock } as unknown as WorkflowStep,
200 mockedModules.context.scope.env,
201 SNAPSHOT,
202 STARTED_AT,
203 );
204 
205 expect(doMock).toHaveBeenCalledWith(
206 "execute run",
207 expect.objectContaining({
208 timeout: 30 * 60 * 1_000,
209 retries: expect.objectContaining({
210 limit: 1,
211 }),
212 }),
213 expect.any(Function),
214 );
215 });
216 
217 it("rebuilds the sandbox after workflow replay without re-entering starting", async () => {
218 mockedModules.context.runStore = createQueueRunStoreStub({
219 getMeta: vi.fn(async () => ({
220 ...createRunMeta("running"),
221 currentStep: expectTrusted(PositiveInteger, 2, "PositiveInteger"),
222 })),
223 });
224 
225 const result = await executeWorkflowRun(
226 createWorkflowStep(),
227 mockedModules.context.scope.env,
228 SNAPSHOT,
229 STARTED_AT,
230 );
231 
232 expect(result).toBe("passed");
233 expect(mockedModules.context.runStore.updateState).not.toHaveBeenCalled();
234 expect(mockedModules.context.state.currentStepPosition).toBe(2);
235 expect(mockedModules.context.runtime.destroySandbox).toHaveBeenCalledTimes(1);
236 expect(mockedModules.prepareExecutionEnvironment).toHaveBeenCalledWith(mockedModules.context, mockedModules.lease, {
237 executionSessionId: `run-${SNAPSHOT.runId}`,
238 });
239 expect(mockedModules.executeRunSteps).toHaveBeenCalledTimes(1);
240 });
241 
242 it("finalizes as canceled when cancellation is already requested before execution starts", async () => {
243 mockedModules.context.projectControl.recordHeartbeat = vi.fn(async () => ({
244 status: "cancel_requested" as const,
245 cancelRequestedAt: STARTED_AT,
246 }));
247 mockedModules.lease = {
248 start: vi.fn(),
249 stop: vi.fn(async () => {}),
250 throwIfOwnershipLost: vi.fn(() => {}),
251 refreshControl: vi.fn(async () => {
252 mockedModules.context.state.cancelRequestedAt = STARTED_AT;
253 }),
254 isCancellationRequested: () => mockedModules.context.state.cancelRequestedAt !== null,
255 applyCancellationIfNeeded: vi.fn(async () => {}),
256 };
257 
258 const result = await executeWorkflowRun(
259 createWorkflowStep(),
260 mockedModules.context.scope.env,
261 SNAPSHOT,
262 STARTED_AT,
263 );
264 
265 expect(result).toBe("canceled");
266 expect(mockedModules.context.runtime.destroySandbox).not.toHaveBeenCalled();
267 expect(mockedModules.context.runStore.updateState).not.toHaveBeenCalled();
268 expect(mockedModules.prepareExecutionEnvironment).not.toHaveBeenCalled();
269 expect(mockedModules.executeRunSteps).not.toHaveBeenCalled();
270 expect(mockedModules.finalizeExecution).toHaveBeenCalledWith(mockedModules.context, mockedModules.lease, {
271 kind: "canceled",
272 });
273 });
274 
275 it("returns the terminal status immediately and reconciles ProjectDO when RunDO is already terminal", async () => {
276 mockedModules.context.scope.env = createEnv({
277 ...createRunMeta("passed"),
278 finishedAt: FINISHED_AT,
279 exitCode: 0,
280 });
281 mockedModules.context.runStore = createQueueRunStoreStub({
282 getMeta: vi.fn(async () => ({
283 ...createRunMeta("passed"),
284 finishedAt: FINISHED_AT,
285 exitCode: 0,
286 })),
287 });
288 
289 const result = await executeWorkflowRun(
290 createWorkflowStep(),
291 mockedModules.context.scope.env,
292 SNAPSHOT,
293 STARTED_AT,
294 );
295 
296 expect(result).toBe("passed");
297 expect(mockedModules.recoverTerminalActiveRun).toHaveBeenCalledWith(
298 mockedModules.context.scope.env,
299 SNAPSHOT.projectId,
300 SNAPSHOT.runId,
301 );
302 expect(mockedModules.context.runtime.destroySandbox).not.toHaveBeenCalled();
303 expect(mockedModules.finalizeExecution).not.toHaveBeenCalled();
304 expect(mockedModules.prepareExecutionEnvironment).not.toHaveBeenCalled();
305 });
306});