Skip to content
File

Blob: tests/worker/dispatch/shared/finalize.test.ts

typescript477 lines
1import { env } from "cloudflare:workers";
2import { describe, expect, it, vi } from "vitest";
3 
4import { ProjectId, RunId, UnixTimestampMs } from "@/contracts";
5import { expectTrusted, PositiveInteger, type FinalizeRunExecutionInput, type RunMetaState } from "@/worker/contracts";
6import { appendFailureLogBestEffort, mapExecutionErrorToOutcome } from "@/worker/dispatch/shared/execution-errors";
7import { finalizeExecution } from "@/worker/dispatch/shared/run-finalize";
8 
9import {
10 createQueueLeaseStub,
11 createQueueProjectControlStub,
12 createQueueRunControlStub,
13 createQueueRunLogsStub,
14 createQueueRunRuntimeStub,
15 createQueueRunStoreStub,
16 createQueueScope,
17 createQueueState,
18} from "../../../helpers/dispatch/shared";
19import { registerWorkerRuntimeHooks } from "../../../helpers/worker-hooks";
20import { seedProject, seedUser } from "../../../helpers/runtime";
21 
22const toTimestamp = (value: number) => expectTrusted(UnixTimestampMs, value, "UnixTimestampMs");
23 
24type FinalizeContext = Parameters<typeof finalizeExecution>[0];
25type ErrorContext = Parameters<typeof mapExecutionErrorToOutcome>[0];
26const createLeaseStub = (cancelRequested: boolean) =>
27 createQueueLeaseStub({
28 isCancellationRequested: cancelRequested,
29 });
30 
31const buildFinalizeContext = (
32 projectId: ProjectId,
33 runId: RunId,
34 startedAt: UnixTimestampMs,
35 cancelRequestedAt: UnixTimestampMs | null,
36 destroySandboxResult: boolean,
37): {
38 context: FinalizeContext;
39 finalizeCalls: FinalizeRunExecutionInput[];
40 runStub: ReturnType<Env["RUN_DO"]["getByName"]>;
41} => {
42 const finalizeCalls: FinalizeRunExecutionInput[] = [];
43 const runStub = env.RUN_DO.getByName(runId);
44 const state = createQueueState({ cancelRequestedAt });
45 const scope = createQueueScope({
46 env,
47 projectId,
48 runId,
49 startedAt,
50 });
51 
52 const context: FinalizeContext = {
53 scope,
54 state,
55 runStore: createQueueRunStoreStub({
56 getFreshStub: () => runStub,
57 tryUpdateState: async (input: Omit<Parameters<typeof runStub.tryUpdateRunState>[0], "runId">) =>
58 await runStub.tryUpdateRunState({
59 runId,
60 ...input,
61 }),
62 repairTerminalState: async (input: Omit<Parameters<typeof runStub.repairTerminalState>[0], "runId">) => {
63 await runStub.repairTerminalState({
64 runId,
65 ...input,
66 });
67 },
68 updateStepState: async (input: Omit<Parameters<typeof runStub.updateStepState>[0], "runId">) => {
69 await runStub.updateStepState({
70 runId,
71 ...input,
72 });
73 },
74 }),
75 projectControl: createQueueProjectControlStub({
76 finalizeRunExecution: async (
77 terminalStatus: Extract<RunMetaState["status"], "passed" | "failed" | "canceled">,
78 lastError: string | null,
79 sandboxDestroyed: boolean,
80 ) => {
81 finalizeCalls.push({
82 projectId,
83 runId,
84 terminalStatus,
85 lastError,
86 sandboxDestroyed,
87 });
88 return {
89 snapshot: scope.snapshot,
90 };
91 },
92 }),
93 runtime: createQueueRunRuntimeStub({
94 destroySandbox: async () => destroySandboxResult,
95 }),
96 control: createQueueRunControlStub({
97 getRunMeta: async (): Promise<RunMetaState> => {
98 const current = await runStub.getRunSummary(runId);
99 if (!current) {
100 throw new Error(`Run ${runId} is not initialized.`);
101 }
102 
103 return current;
104 },
105 updateRunFromCurrent: async (
106 current: RunMetaState,
107 status: RunMetaState["status"],
108 overrides: Partial<
109 Pick<RunMetaState, "currentStep" | "startedAt" | "finishedAt" | "exitCode" | "errorMessage">
110 > = {},
111 ): Promise<RunMetaState> => {
112 const hasOverride = <TKey extends keyof typeof overrides>(key: TKey): boolean =>
113 Object.prototype.hasOwnProperty.call(overrides, key);
114 
115 await runStub.updateRunState({
116 runId,
117 status,
118 currentStep: hasOverride("currentStep") ? overrides.currentStep : current.currentStep,
119 startedAt: hasOverride("startedAt") ? overrides.startedAt : current.startedAt,
120 finishedAt: hasOverride("finishedAt") ? overrides.finishedAt : current.finishedAt,
121 exitCode: hasOverride("exitCode") ? overrides.exitCode : current.exitCode,
122 errorMessage: hasOverride("errorMessage") ? overrides.errorMessage : current.errorMessage,
123 });
124 
125 const updated = await runStub.getRunSummary(runId);
126 if (!updated) {
127 throw new Error(`Run ${runId} is not initialized.`);
128 }
129 
130 return updated;
131 },
132 }),
133 };
134 
135 return {
136 context,
137 finalizeCalls,
138 runStub,
139 };
140};
141 
142describe("queue finalization", () => {
143 registerWorkerRuntimeHooks();
144 
145 describe("late cancellation", () => {
146 it("marks a late-canceled run as canceled when cleanup succeeds", async () => {
147 const user = await seedUser({
148 email: "late-pass@example.com",
149 slug: "late-pass-user",
150 });
151 const project = await seedProject(user, {
152 projectSlug: "late-pass-project",
153 dispatchMode: "queue",
154 });
155 const runId = RunId.assertDecode("run_0000000000000000000001");
156 const startedAt = toTimestamp(1_740_000_000_000);
157 const cancelRequestedAt = toTimestamp(startedAt + 1_000);
158 const { context, finalizeCalls, runStub } = buildFinalizeContext(
159 project.id,
160 runId,
161 startedAt,
162 cancelRequestedAt,
163 true,
164 );
165 
166 await runStub.ensureInitialized({
167 runId,
168 projectId: project.id,
169 triggerType: "manual",
170 branch: project.defaultBranch,
171 commitSha: null,
172 });
173 await runStub.updateRunState({
174 runId,
175 status: "starting",
176 currentStep: null,
177 startedAt,
178 finishedAt: null,
179 exitCode: null,
180 errorMessage: null,
181 });
182 await runStub.updateRunState({
183 runId,
184 status: "running",
185 currentStep: null,
186 startedAt,
187 finishedAt: null,
188 exitCode: 0,
189 errorMessage: null,
190 });
191 await runStub.updateRunState({
192 runId,
193 status: "cancel_requested",
194 currentStep: null,
195 startedAt,
196 finishedAt: null,
197 exitCode: 0,
198 errorMessage: null,
199 });
200 
201 await finalizeExecution(context, createLeaseStub(true), {
202 kind: "passed",
203 exitCode: 0,
204 });
205 
206 const runMeta = await runStub.getRunSummary(runId);
207 expect(runMeta?.status).toBe("canceled");
208 expect(runMeta?.exitCode).toBeNull();
209 expect(runMeta?.errorMessage).toBeNull();
210 expect(finalizeCalls).toEqual([
211 {
212 projectId: project.id,
213 runId,
214 terminalStatus: "canceled",
215 lastError: null,
216 sandboxDestroyed: true,
217 },
218 ]);
219 });
220 
221 it("repairs the active step to failed before terminalizing the run", async () => {
222 const user = await seedUser({
223 email: "step-repair@example.com",
224 slug: "step-repair-user",
225 });
226 const project = await seedProject(user, {
227 projectSlug: "step-repair-project",
228 dispatchMode: "queue",
229 });
230 const runId = RunId.assertDecode("run_0000000000000000000003");
231 const startedAt = toTimestamp(1_740_000_030_000);
232 const { context, runStub } = buildFinalizeContext(project.id, runId, startedAt, null, true);
233 
234 await runStub.ensureInitialized({
235 runId,
236 projectId: project.id,
237 triggerType: "manual",
238 branch: project.defaultBranch,
239 commitSha: null,
240 });
241 await runStub.replaceSteps({
242 runId,
243 steps: [
244 {
245 position: expectTrusted(PositiveInteger, 1, "PositiveInteger"),
246 name: "build",
247 command: "npm run build",
248 },
249 ],
250 });
251 await runStub.updateRunState({
252 runId,
253 status: "starting",
254 currentStep: null,
255 startedAt,
256 finishedAt: null,
257 exitCode: null,
258 errorMessage: null,
259 });
260 await runStub.updateRunState({
261 runId,
262 status: "running",
263 currentStep: expectTrusted(PositiveInteger, 1, "PositiveInteger"),
264 startedAt,
265 finishedAt: null,
266 exitCode: null,
267 errorMessage: null,
268 });
269 await runStub.updateStepState({
270 runId,
271 position: expectTrusted(PositiveInteger, 1, "PositiveInteger"),
272 status: "running",
273 startedAt,
274 finishedAt: null,
275 exitCode: null,
276 });
277 context.state.currentStepPosition = 1 as RunMetaState["currentStep"];
278 
279 await finalizeExecution(context, createLeaseStub(false), {
280 kind: "failed",
281 exitCode: 1,
282 errorMessage: "step failed",
283 });
284 
285 const detail = await runStub.getRunDetail(runId);
286 expect(detail.meta?.status).toBe("failed");
287 expect(detail.steps[0]?.status).toBe("failed");
288 expect(detail.steps[0]?.finishedAt).not.toBeNull();
289 expect(detail.steps[0]?.exitCode).toBe(1);
290 });
291 
292 it("fails a late-canceled run when cancellation cleanup fails", async () => {
293 const user = await seedUser({
294 email: "late-fail@example.com",
295 slug: "late-fail-user",
296 });
297 const project = await seedProject(user, {
298 projectSlug: "late-fail-project",
299 dispatchMode: "queue",
300 });
301 const runId = RunId.assertDecode("run_0000000000000000000002");
302 const startedAt = toTimestamp(1_740_000_010_000);
303 const cancelRequestedAt = toTimestamp(startedAt + 1_000);
304 const { context, finalizeCalls, runStub } = buildFinalizeContext(
305 project.id,
306 runId,
307 startedAt,
308 cancelRequestedAt,
309 false,
310 );
311 
312 await runStub.ensureInitialized({
313 runId,
314 projectId: project.id,
315 triggerType: "manual",
316 branch: project.defaultBranch,
317 commitSha: null,
318 });
319 await runStub.updateRunState({
320 runId,
321 status: "starting",
322 currentStep: null,
323 startedAt,
324 finishedAt: null,
325 exitCode: null,
326 errorMessage: null,
327 });
328 await runStub.updateRunState({
329 runId,
330 status: "running",
331 currentStep: null,
332 startedAt,
333 finishedAt: null,
334 exitCode: 23,
335 errorMessage: null,
336 });
337 await runStub.updateRunState({
338 runId,
339 status: "cancel_requested",
340 currentStep: null,
341 startedAt,
342 finishedAt: null,
343 exitCode: 23,
344 errorMessage: null,
345 });
346 
347 await finalizeExecution(context, createLeaseStub(true), {
348 kind: "failed",
349 exitCode: 23,
350 errorMessage: "step failed",
351 });
352 
353 const runMeta = await runStub.getRunSummary(runId);
354 expect(runMeta?.status).toBe("failed");
355 expect(runMeta?.exitCode).toBe(1);
356 expect(runMeta?.errorMessage).toBe("cancel_cleanup_failed");
357 expect(finalizeCalls).toEqual([
358 {
359 projectId: project.id,
360 runId,
361 terminalStatus: "failed",
362 lastError: "cancel_cleanup_failed",
363 sandboxDestroyed: false,
364 },
365 ]);
366 });
367 
368 it("keeps the process-tree cleanup path when the session still has live processes but no current process handle", async () => {
369 const user = await seedUser({
370 email: "live-process-cleanup@example.com",
371 slug: "live-process-cleanup-user",
372 });
373 const project = await seedProject(user, {
374 projectSlug: "live-process-cleanup-project",
375 dispatchMode: "queue",
376 });
377 const runId = RunId.assertDecode("run_0000000000000000000004");
378 const startedAt = toTimestamp(1_740_000_040_000);
379 const { context, finalizeCalls, runStub } = buildFinalizeContext(project.id, runId, startedAt, null, true);
380 
381 const session = { id: "session-live" } as never;
382 context.state.session = session;
383 context.runtime.isProcessTreeAlive = vi.fn(async () => true);
384 context.runtime.waitForProcessTreeToStopSafely = vi.fn(async () => true);
385 
386 await runStub.ensureInitialized({
387 runId,
388 projectId: project.id,
389 triggerType: "manual",
390 branch: project.defaultBranch,
391 commitSha: null,
392 });
393 await runStub.updateRunState({
394 runId,
395 status: "starting",
396 currentStep: null,
397 startedAt,
398 finishedAt: null,
399 exitCode: null,
400 errorMessage: null,
401 });
402 await runStub.updateRunState({
403 runId,
404 status: "running",
405 currentStep: null,
406 startedAt,
407 finishedAt: null,
408 exitCode: null,
409 errorMessage: null,
410 });
411 
412 await finalizeExecution(context, createLeaseStub(false), {
413 kind: "passed",
414 exitCode: 0,
415 });
416 
417 expect(context.runtime.softCancelProcessTree).toHaveBeenCalledWith(session, null);
418 expect(context.runtime.deleteSession).not.toHaveBeenCalled();
419 expect(finalizeCalls).toEqual([
420 {
421 projectId: project.id,
422 runId,
423 terminalStatus: "passed",
424 lastError: null,
425 sandboxDestroyed: true,
426 },
427 ]);
428 });
429 });
430 
431 describe("failure logging", () => {
432 it("classifies errors without appending logs", async () => {
433 const runId = RunId.assertDecode("run_0000000000000000000000");
434 const projectId = expectTrusted(ProjectId, "prj_0000000000000000000000", "ProjectId");
435 const context: ErrorContext = {
436 scope: createQueueScope({
437 env,
438 projectId,
439 runId,
440 startedAt: toTimestamp(1_740_000_020_000),
441 }),
442 state: createQueueState(),
443 logs: createQueueRunLogsStub({
444 appendSystemLog: vi.fn(async () => {
445 throw new Error("append failed");
446 }),
447 }),
448 };
449 
450 await expect(mapExecutionErrorToOutcome(context, new Error("boom"))).resolves.toEqual({
451 kind: "failed",
452 exitCode: 1,
453 errorMessage: "boom",
454 });
455 expect(context.logs.appendSystemLog).not.toHaveBeenCalled();
456 });
457 
458 it("treats queue-side failure logging as best effort", async () => {
459 const context = {
460 scope: createQueueScope({
461 env,
462 projectId: expectTrusted(ProjectId, "prj_0000000000000000000000", "ProjectId"),
463 runId: RunId.assertDecode("run_0000000000000000000000"),
464 startedAt: toTimestamp(1_740_000_020_000),
465 }),
466 logs: createQueueRunLogsStub({
467 appendSystemLog: vi.fn(async () => {
468 throw new Error("append failed");
469 }),
470 }),
471 };
472 
473 await expect(appendFailureLogBestEffort(context, "boom")).resolves.toBeUndefined();
474 });
475 });
476});