Skip to content
File

Blob: tests/worker/dispatch/queue/consumer.test.ts

typescript144 lines
1import { createExecutionContext, createMessageBatch, getQueueResult, runInDurableObject } from "cloudflare:test";
2import { env } from "cloudflare:workers";
3 
4import { describe, expect, it } from "vitest";
5 
6import { ProjectId, RunId, UnixTimestampMs } from "@/contracts";
7import { expectTrusted, type RunQueueMessage } from "@/worker/contracts";
8import { ProjectDO } from "@/worker/durable";
9import { getSandboxCleanupRetryState } from "@/worker/durable/project-do/sidecar-state";
10import worker from "@/worker/index";
11 
12import { createTestProjectDoContext, expectAcceptedManualRun } from "../../../helpers/project-do";
13import { registerWorkerRuntimeHooks } from "../../../helpers/worker-hooks";
14import { readProjectDoRows, seedProject, seedUser } from "../../../helpers/runtime";
15 
16describe("queue consumer", () => {
17 const toTimestamp = (value: number) => expectTrusted(UnixTimestampMs, value, "UnixTimestampMs");
18 registerWorkerRuntimeHooks();
19 
20 it("acks queue messages for missing projects", async () => {
21 const batch = createMessageBatch("anvil-runs", [
22 {
23 id: "missing-project",
24 attempts: 0,
25 timestamp: new Date(),
26 body: {
27 projectId: ProjectId.assertDecode("prj_0000000000000000000000"),
28 runId: RunId.assertDecode("run_0000000000000000000000"),
29 } satisfies RunQueueMessage,
30 },
31 ]);
32 const ctx = createExecutionContext();
33 
34 await worker.queue!(batch as Parameters<NonNullable<typeof worker.queue>>[0], env, ctx);
35 const result = await getQueueResult(batch, ctx);
36 
37 expect(result.explicitAcks).toEqual(["missing-project"]);
38 expect(result.retryMessages).toEqual([]);
39 });
40 
41 it("retries invalid queue payloads", async () => {
42 const batch = createMessageBatch("anvil-runs", [
43 {
44 id: "invalid-message",
45 attempts: 0,
46 timestamp: new Date(),
47 body: {
48 projectId: "not-a-project-id",
49 runId: "also-invalid",
50 },
51 },
52 ]);
53 const ctx = createExecutionContext();
54 
55 await worker.queue!(batch as Parameters<NonNullable<typeof worker.queue>>[0], env, ctx);
56 const result = await getQueueResult(batch, ctx);
57 
58 expect(result.explicitAcks).toEqual([]);
59 expect(result.retryMessages).toEqual([{ msgId: "invalid-message" }]);
60 });
61 
62 it("recovers stale duplicate active deliveries from terminal RunDO state", async () => {
63 const user = await seedUser({
64 email: "queue@example.com",
65 slug: "queue-user",
66 });
67 const project = await seedProject(user, {
68 projectSlug: "queue-project",
69 dispatchMode: "queue",
70 });
71 const projectStub = env.PROJECT_DO.getByName(project.id);
72 
73 const accepted = expectAcceptedManualRun(
74 await projectStub.acceptManualRun({
75 projectId: project.id,
76 triggeredByUserId: user.id,
77 branch: project.defaultBranch,
78 }),
79 );
80 const claim = await projectStub.claimRunWork({
81 projectId: project.id,
82 runId: accepted.runId,
83 });
84 expect(claim.kind).toBe("execute");
85 
86 await env.RUN_DO.getByName(accepted.runId).updateRunState({
87 runId: accepted.runId,
88 status: "starting",
89 currentStep: null,
90 startedAt: toTimestamp(Date.now()),
91 finishedAt: null,
92 exitCode: null,
93 errorMessage: null,
94 });
95 await env.RUN_DO.getByName(accepted.runId).updateRunState({
96 runId: accepted.runId,
97 status: "running",
98 currentStep: null,
99 startedAt: toTimestamp(Date.now()),
100 finishedAt: null,
101 exitCode: null,
102 errorMessage: null,
103 });
104 await env.RUN_DO.getByName(accepted.runId).updateRunState({
105 runId: accepted.runId,
106 status: "passed",
107 currentStep: null,
108 startedAt: toTimestamp(Date.now()),
109 finishedAt: toTimestamp(Date.now() + 1_000),
110 exitCode: 0,
111 errorMessage: null,
112 });
113 
114 const batch = createMessageBatch("anvil-runs", [
115 {
116 id: "duplicate-active",
117 attempts: 0,
118 timestamp: new Date(),
119 body: {
120 projectId: project.id,
121 runId: accepted.runId,
122 } satisfies RunQueueMessage,
123 },
124 ]);
125 const ctx = createExecutionContext();
126 
127 await worker.queue!(batch as Parameters<NonNullable<typeof worker.queue>>[0], env, ctx);
128 const result = await getQueueResult(batch, ctx);
129 
130 expect(result.explicitAcks).toEqual(["duplicate-active"]);
131 expect(result.retryMessages).toEqual([]);
132 
133 const rows = await readProjectDoRows(project.id);
134 expect(rows.state?.activeRunId).toBeNull();
135 expect(rows.runs[0]?.status).toBe("passed");
136 
137 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
138 const retryState = await getSandboxCleanupRetryState(createTestProjectDoContext(instance), accepted.runId);
139 expect(retryState?.attempt).toBe(1);
140 expect(retryState?.nextAt ?? 0).toBeGreaterThan(Date.now());
141 });
142 });
143});