Skip to content
File

Blob: tests/worker/routes/run-history-and-log-tickets.test.ts

typescript392 lines
1import { env, exports } from "cloudflare:workers";
2import { describe, expect, it } from "vitest";
3 
4import {
5 DEFAULT_DISPATCH_MODE,
6 DEFAULT_EXECUTION_RUNTIME,
7 type GetProjectRunsResponse,
8 type LogStreamTicketResponse,
9 type ProjectResponse,
10 RunId,
11 type RunDetail,
12 type RunWsMessage,
13 type TriggerRunAcceptedResponse,
14} from "@/contracts";
15import { createD1Db } from "@/worker/db/d1";
16import * as d1Schema from "@/worker/db/d1/schema";
17import { generateDurableEntityId } from "@/worker/services";
18 
19import { authHeaders, fetchJson, mintCookieAuth, seedProject, seedUser } from "../../helpers/runtime";
20import { registerWorkerRuntimeHooks } from "../../helpers/worker-hooks";
21 
22const createProjectViaRoute = async (sessionId: string, projectSlug: string) =>
23 await fetchJson<ProjectResponse>("/api/private/projects", {
24 method: "POST",
25 headers: authHeaders(sessionId, {
26 "content-type": "application/json; charset=utf-8",
27 }),
28 body: JSON.stringify({
29 projectSlug,
30 name: `Project ${projectSlug}`,
31 repoUrl: `https://github.com/example/${projectSlug}`,
32 defaultBranch: "main",
33 configPath: ".anvil.yml",
34 }),
35 });
36 
37const triggerRunViaRoute = async (sessionId: string, projectId: string, branch?: string) =>
38 await fetchJson<TriggerRunAcceptedResponse>(`/api/private/projects/${projectId}/runs`, {
39 method: "POST",
40 headers: authHeaders(sessionId, {
41 "content-type": "application/json; charset=utf-8",
42 }),
43 body: JSON.stringify(branch === undefined ? {} : { branch }),
44 });
45 
46const createLogTicketViaRoute = async (sessionId: string, runId: string) =>
47 await fetchJson<LogStreamTicketResponse>(`/api/private/runs/${runId}/log-ticket`, {
48 method: "POST",
49 headers: authHeaders(sessionId),
50 });
51 
52const openLogStream = async (runId: string, ticket: string): Promise<Response> =>
53 await exports.default.fetch(
54 `https://example.com/api/private/runs/${runId}/logs?ticket=${encodeURIComponent(ticket)}`,
55 {
56 headers: {
57 upgrade: "websocket",
58 },
59 },
60 );
61 
62describe("worker run history and log ticket routes", () => {
63 registerWorkerRuntimeHooks();
64 
65 it("uses the default branch for omitted manual runs and preserves explicit branch overrides", async () => {
66 const user = await seedUser({
67 email: "run-branches@example.com",
68 slug: "run-branches-user",
69 });
70 
71 const { sessionId } = await mintCookieAuth(user.id);
72 
73 const createdProject = await createProjectViaRoute(sessionId, "run-branches-project");
74 expect(createdProject.status).toBe(201);
75 expect(createdProject.body).not.toBeNull();
76 const projectId = createdProject.body!.project.id;
77 
78 const defaultBranchRun = await triggerRunViaRoute(sessionId, projectId);
79 expect(defaultBranchRun.status).toBe(202);
80 expect(defaultBranchRun.body).not.toBeNull();
81 
82 const overrideBranchRun = await triggerRunViaRoute(sessionId, projectId, "release");
83 expect(overrideBranchRun.status).toBe(202);
84 expect(overrideBranchRun.body).not.toBeNull();
85 
86 const defaultRunDetail = await fetchJson<RunDetail>(`/api/private/runs/${defaultBranchRun.body!.runId}`, {
87 headers: authHeaders(sessionId),
88 });
89 expect(defaultRunDetail.status).toBe(200);
90 expect(defaultRunDetail.body?.run.branch).toBe("main");
91 
92 const overrideRunDetail = await fetchJson<RunDetail>(`/api/private/runs/${overrideBranchRun.body!.runId}`, {
93 headers: authHeaders(sessionId),
94 });
95 expect(overrideRunDetail.status).toBe(200);
96 expect(overrideRunDetail.body?.run.branch).toBe("release");
97 });
98 
99 it("paginates project runs and validates limit and cursor query params", async () => {
100 const user = await seedUser({
101 email: "run-pagination@example.com",
102 slug: "run-pagination-user",
103 });
104 
105 const { sessionId } = await mintCookieAuth(user.id);
106 
107 const createdProject = await createProjectViaRoute(sessionId, "run-pagination-project");
108 expect(createdProject.status).toBe(201);
109 expect(createdProject.body).not.toBeNull();
110 const project = createdProject.body!.project;
111 
112 const db = createD1Db(env.DB);
113 const queuedAtValues = [Date.now() - 2_000, Date.now() - 1_000, Date.now()];
114 const insertedRunIds = queuedAtValues.map((queuedAt, index) => generateDurableEntityId("run", queuedAt + index));
115 
116 await db.insert(d1Schema.runIndex).values(
117 queuedAtValues.map((queuedAt, index) => ({
118 id: insertedRunIds[index],
119 projectId: project.id,
120 triggeredByUserId: user.id,
121 triggerType: "manual",
122 branch: "main",
123 commitSha: null,
124 status: "passed",
125 dispatchMode: DEFAULT_DISPATCH_MODE,
126 executionRuntime: DEFAULT_EXECUTION_RUNTIME,
127 queuedAt,
128 startedAt: queuedAt + 10,
129 finishedAt: queuedAt + 20,
130 exitCode: 0,
131 })),
132 );
133 
134 const firstPage = await fetchJson<GetProjectRunsResponse>(`/api/private/projects/${project.id}/runs?limit=2`, {
135 headers: authHeaders(sessionId),
136 });
137 expect(firstPage.status).toBe(200);
138 expect(firstPage.body).not.toBeNull();
139 expect(firstPage.body!.runs.map((run) => run.id)).toEqual([insertedRunIds[2], insertedRunIds[1]]);
140 expect(firstPage.body!.nextCursor).toEqual(expect.any(String));
141 
142 const secondPage = await fetchJson<GetProjectRunsResponse>(
143 `/api/private/projects/${project.id}/runs?limit=2&cursor=${encodeURIComponent(firstPage.body!.nextCursor!)}`,
144 {
145 headers: authHeaders(sessionId),
146 },
147 );
148 expect(secondPage.status).toBe(200);
149 expect(secondPage.body).not.toBeNull();
150 expect(secondPage.body!.runs.map((run) => run.id)).toEqual([insertedRunIds[0]]);
151 expect(secondPage.body!.nextCursor).toBeNull();
152 
153 const invalidLimit = await fetchJson(`/api/private/projects/${project.id}/runs?limit=0`, {
154 headers: authHeaders(sessionId),
155 });
156 expect(invalidLimit.status).toBe(400);
157 expect(invalidLimit.body).toMatchObject({
158 error: {
159 code: "invalid_request",
160 },
161 });
162 
163 const invalidCursor = await fetchJson(`/api/private/projects/${project.id}/runs?cursor=not-a-valid-cursor`, {
164 headers: authHeaders(sessionId),
165 });
166 expect(invalidCursor.status).toBe(400);
167 expect(invalidCursor.body).toMatchObject({
168 error: {
169 code: "invalid_cursor",
170 },
171 });
172 });
173 
174 it("requires a websocket upgrade for log streams without consuming the ticket", async () => {
175 const user = await seedUser({
176 email: "run-log-upgrade@example.com",
177 slug: "run-log-upgrade-user",
178 });
179 
180 const { sessionId } = await mintCookieAuth(user.id);
181 
182 const createdProject = await createProjectViaRoute(sessionId, "run-log-upgrade-project");
183 expect(createdProject.status).toBe(201);
184 expect(createdProject.body).not.toBeNull();
185 const projectId = createdProject.body!.project.id;
186 
187 const acceptedRun = await triggerRunViaRoute(sessionId, projectId);
188 expect(acceptedRun.status).toBe(202);
189 expect(acceptedRun.body).not.toBeNull();
190 const runId = acceptedRun.body!.runId;
191 
192 const logTicket = await createLogTicketViaRoute(sessionId, runId);
193 expect(logTicket.status).toBe(200);
194 expect(logTicket.body).not.toBeNull();
195 
196 const upgradeRequired = await fetchJson(`/api/private/runs/${runId}/logs?ticket=${logTicket.body!.ticket}`);
197 expect(upgradeRequired.status).toBe(426);
198 expect(upgradeRequired.body).toMatchObject({
199 error: {
200 code: "upgrade_required",
201 },
202 });
203 
204 const websocketResponse = await openLogStream(runId, logTicket.body!.ticket);
205 expect(websocketResponse.status).toBe(101);
206 });
207 
208 it("rejects missing, mismatched, reused, and expired log tickets", async () => {
209 const user = await seedUser({
210 email: "run-log-ticket-errors@example.com",
211 slug: "run-log-ticket-errors-user",
212 });
213 
214 const { sessionId } = await mintCookieAuth(user.id);
215 
216 const createdProject = await createProjectViaRoute(sessionId, "run-log-ticket-errors-project");
217 expect(createdProject.status).toBe(201);
218 expect(createdProject.body).not.toBeNull();
219 const projectId = createdProject.body!.project.id;
220 
221 const firstRun = await triggerRunViaRoute(sessionId, projectId);
222 expect(firstRun.status).toBe(202);
223 expect(firstRun.body).not.toBeNull();
224 
225 const secondRun = await triggerRunViaRoute(sessionId, projectId, "release");
226 expect(secondRun.status).toBe(202);
227 expect(secondRun.body).not.toBeNull();
228 
229 const missingTicket = await fetchJson(`/api/private/runs/${firstRun.body!.runId}/logs?ticket=missing-ticket`, {
230 headers: {
231 upgrade: "websocket",
232 },
233 });
234 expect(missingTicket.status).toBe(403);
235 expect(missingTicket.body).toMatchObject({
236 error: {
237 code: "invalid_log_ticket",
238 },
239 });
240 
241 const mismatchedTicket = await createLogTicketViaRoute(sessionId, firstRun.body!.runId);
242 expect(mismatchedTicket.status).toBe(200);
243 expect(mismatchedTicket.body).not.toBeNull();
244 
245 const mismatchedRun = await fetchJson(
246 `/api/private/runs/${secondRun.body!.runId}/logs?ticket=${mismatchedTicket.body!.ticket}`,
247 {
248 headers: {
249 upgrade: "websocket",
250 },
251 },
252 );
253 expect(mismatchedRun.status).toBe(403);
254 expect(mismatchedRun.body).toMatchObject({
255 error: {
256 code: "invalid_log_ticket",
257 },
258 });
259 
260 const reusableTicket = await createLogTicketViaRoute(sessionId, firstRun.body!.runId);
261 expect(reusableTicket.status).toBe(200);
262 expect(reusableTicket.body).not.toBeNull();
263 
264 const firstOpen = await openLogStream(firstRun.body!.runId, reusableTicket.body!.ticket);
265 expect(firstOpen.status).toBe(101);
266 
267 const reusedTicket = await fetchJson(
268 `/api/private/runs/${firstRun.body!.runId}/logs?ticket=${reusableTicket.body!.ticket}`,
269 {
270 headers: {
271 upgrade: "websocket",
272 },
273 },
274 );
275 expect(reusedTicket.status).toBe(403);
276 expect(reusedTicket.body).toMatchObject({
277 error: {
278 code: "invalid_log_ticket",
279 },
280 });
281 
282 const expiredTicket = await createLogTicketViaRoute(sessionId, firstRun.body!.runId);
283 expect(expiredTicket.status).toBe(200);
284 expect(expiredTicket.body).not.toBeNull();
285 
286 await env.LOG_TICKETS.put(
287 `run-log-ticket:${expiredTicket.body!.ticket}`,
288 JSON.stringify({
289 runId: firstRun.body!.runId,
290 userId: user.id,
291 expiresAt: Date.now() - 1_000,
292 }),
293 { expirationTtl: 60 },
294 );
295 
296 const expired = await fetchJson(
297 `/api/private/runs/${firstRun.body!.runId}/logs?ticket=${expiredTicket.body!.ticket}`,
298 {
299 headers: {
300 upgrade: "websocket",
301 },
302 },
303 );
304 expect(expired.status).toBe(403);
305 expect(expired.body).toMatchObject({
306 error: {
307 code: "invalid_log_ticket",
308 },
309 });
310 });
311 
312 it("sends state snapshot in envelope format on websocket connect", async () => {
313 const user = await seedUser({
314 email: "ws-envelope@example.com",
315 slug: "ws-envelope-user",
316 });
317 const project = await seedProject(user, {
318 projectSlug: "ws-envelope-project",
319 });
320 
321 const { sessionId } = await mintCookieAuth(user.id);
322 
323 const runId = RunId.assertDecode(generateDurableEntityId("run", Date.now()));
324 const runStub = env.RUN_DO.getByName(runId);
325 await runStub.ensureInitialized({
326 runId,
327 projectId: project.id,
328 triggerType: "manual",
329 branch: project.defaultBranch,
330 commitSha: null,
331 });
332 
333 const logTicket = await createLogTicketViaRoute(sessionId, runId);
334 expect(logTicket.status).toBe(200);
335 
336 // Need a D1 run-index row so the route can resolve run ownership
337 const db = createD1Db(env.DB);
338 await db.insert(d1Schema.runIndex).values({
339 id: runId,
340 projectId: project.id,
341 triggeredByUserId: user.id,
342 triggerType: "manual",
343 branch: "main",
344 commitSha: null,
345 status: "queued",
346 dispatchMode: DEFAULT_DISPATCH_MODE,
347 executionRuntime: DEFAULT_EXECUTION_RUNTIME,
348 queuedAt: Date.now(),
349 startedAt: null,
350 finishedAt: null,
351 exitCode: null,
352 });
353 
354 const wsResponse = await openLogStream(runId, logTicket.body!.ticket);
355 expect(wsResponse.status).toBe(101);
356 
357 const ws = wsResponse.webSocket!;
358 ws.accept();
359 
360 const messages: RunWsMessage[] = [];
361 const closed = new Promise<void>((resolve) => {
362 ws.addEventListener("message", (event) => {
363 messages.push(JSON.parse(event.data as string) as RunWsMessage);
364 });
365 ws.addEventListener("close", () => resolve());
366 });
367 
368 ws.close(1000, "done");
369 await closed;
370 
371 // Should receive at least an initial state snapshot
372 const stateMessages = messages.filter((m) => m.type === "state");
373 expect(stateMessages.length).toBeGreaterThanOrEqual(1);
374 
375 const initial = stateMessages[0];
376 expect(initial.type).toBe("state");
377 if (initial.type !== "state") throw new Error("unreachable");
378 
379 expect(initial.run.status).toBe("queued");
380 expect(initial.run.currentStep).toBeNull();
381 expect(initial.run.startedAt).toBeNull();
382 expect(initial.run.finishedAt).toBeNull();
383 expect(initial.run.exitCode).toBeNull();
384 expect(initial.run.errorMessage).toBeNull();
385 expect(initial.steps).toEqual([]);
386 
387 // No log messages expected for a freshly initialized run
388 const logMessages = messages.filter((m) => m.type === "log");
389 expect(logMessages).toEqual([]);
390 });
391});