File
Blob: tests/worker/routes/run-history-and-log-tickets.test.ts
| 1 | import { env, exports } from "cloudflare:workers"; |
| 2 | import { describe, expect, it } from "vitest"; |
| 3 | |
| 4 | import { |
| 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"; |
| 15 | import { createD1Db } from "@/worker/db/d1"; |
| 16 | import * as d1Schema from "@/worker/db/d1/schema"; |
| 17 | import { generateDurableEntityId } from "@/worker/services"; |
| 18 | |
| 19 | import { authHeaders, fetchJson, mintCookieAuth, seedProject, seedUser } from "../../helpers/runtime"; |
| 20 | import { registerWorkerRuntimeHooks } from "../../helpers/worker-hooks"; |
| 21 | |
| 22 | const 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 | |
| 37 | const 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 | |
| 46 | const 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 | |
| 52 | const 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 | |
| 62 | describe("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 | }); |