Skip to content
File

Blob: tests/worker/project-do/alarms.test.ts

typescript500 lines
1import { runInDurableObject } from "cloudflare:test";
2import { env } from "cloudflare:workers";
3import { eq } from "drizzle-orm";
4import { describe, expect, it, vi } from "vitest";
5 
6import { BranchName, CommitSha, DEFAULT_DISPATCH_MODE, DEFAULT_EXECUTION_RUNTIME } from "@/contracts";
7import {
8 PROJECT_ALARM_MAX_ITERATIONS,
9 PROJECT_RECONCILIATION_LIVENESS_FALLBACK_MS,
10} from "@/worker/durable/project-do/constants";
11import * as projectDoSchema from "@/worker/db/durable/schema/project-do";
12import { ProjectDO } from "@/worker/durable";
13import {
14 acceptManualRun as acceptManualRunCommand,
15 claimRunWork as claimRunWorkCommand,
16 finalizeRunExecution as finalizeRunExecutionCommand,
17 recordRunResolvedCommit as recordRunResolvedCommitCommand,
18 requestRunCancel as requestRunCancelCommand,
19} from "@/worker/durable/project-do/commands";
20import * as reconciliationModule from "@/worker/durable/project-do/reconciliation";
21import * as sidecarState from "@/worker/durable/project-do/sidecar-state";
22import { createLogger } from "@/worker/services/logger";
23 
24import {
25 acceptManualRunWithoutAlarm,
26 claimRunWorkWithoutAlarm,
27 createAlarmWriteFailingContext,
28 createBoundStorageProxy,
29 createBoundTransactionProxy,
30 createTestProjectDoContext,
31 finalizeRunExecutionWithoutAlarm,
32 getProjectDoInternals,
33 withPatchedStorage,
34} from "../../helpers/project-do";
35import { registerWorkerRuntimeHooks } from "../../helpers/worker-hooks";
36import { readProjectDoRows, seedProject, seedUser } from "../../helpers/runtime";
37 
38describe("ProjectDO alarm behavior", () => {
39 registerWorkerRuntimeHooks();
40 
41 describe("arming fallbacks", () => {
42 it("arms a fallback alarm when non-terminal rows have no immediate reconciliation candidate", async () => {
43 const user = await seedUser({
44 email: "fallback@example.com",
45 slug: "fallback-user",
46 });
47 const project = await seedProject(user, {
48 projectSlug: "fallback-project",
49 });
50 const projectStub = env.PROJECT_DO.getByName(project.id);
51 
52 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
53 projectId: project.id,
54 triggeredByUserId: user.id,
55 branch: project.defaultBranch,
56 });
57 
58 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
59 const { ctx, env: durableEnv, db } = createTestProjectDoContext(instance);
60 
61 await db
62 .update(projectDoSchema.projectRuns)
63 .set({
64 status: "executable",
65 dispatchStatus: "queued",
66 d1SyncStatus: "current",
67 })
68 .where(eq(projectDoSchema.projectRuns.runId, accepted.runId));
69 await ctx.storage.deleteAlarm();
70 
71 const before = Date.now();
72 await sidecarState.rescheduleAlarm(
73 {
74 ctx,
75 env: durableEnv,
76 db,
77 logger: createLogger("test.project-do"),
78 cacheProjectId: () => {},
79 },
80 project.id,
81 );
82 
83 const alarmAt = await ctx.storage.getAlarm();
84 expect(alarmAt).not.toBeNull();
85 expect(alarmAt as number).toBeGreaterThanOrEqual(before + PROJECT_RECONCILIATION_LIVENESS_FALLBACK_MS - 2_000);
86 expect(alarmAt as number).toBeLessThanOrEqual(Date.now() + PROJECT_RECONCILIATION_LIVENESS_FALLBACK_MS + 2_000);
87 });
88 });
89 
90 it("falls back to a direct setAlarm when getAlarm fails during reconciliation arming", async () => {
91 const user = await seedUser({
92 email: "direct-fallback@example.com",
93 slug: "direct-fallback-user",
94 });
95 const project = await seedProject(user, {
96 projectSlug: "direct-fallback-project",
97 });
98 const projectStub = env.PROJECT_DO.getByName(project.id);
99 
100 await acceptManualRunWithoutAlarm(projectStub, {
101 projectId: project.id,
102 triggeredByUserId: user.id,
103 branch: project.defaultBranch,
104 });
105 
106 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
107 const baseContext = createTestProjectDoContext(instance);
108 const alarmWrites: number[] = [];
109 const getAlarm: DurableObjectStorage["getAlarm"] = async () => {
110 throw new Error("getAlarm failed");
111 };
112 const setAlarm: DurableObjectStorage["setAlarm"] = async (scheduledTime, options) => {
113 alarmWrites.push(Number(scheduledTime));
114 await baseContext.ctx.storage.setAlarm(scheduledTime, options);
115 };
116 const storage = createBoundStorageProxy(baseContext.ctx.storage, { getAlarm, setAlarm });
117 const context = withPatchedStorage(baseContext, storage);
118 
119 await expect(sidecarState.armReconciliation(context, project.id)).resolves.toBeUndefined();
120 expect(alarmWrites).toHaveLength(1);
121 expect(alarmWrites[0]).toBeGreaterThanOrEqual(Date.now() - 2_000);
122 expect(alarmWrites[0]).toBeLessThanOrEqual(Date.now() + 2_000);
123 });
124 });
125 
126 it("rewrites an overdue alarm when canceling an active run after restart", async () => {
127 const user = await seedUser({
128 email: "overdue-cancel-rearm@example.com",
129 slug: "overdue-cancel-rearm-user",
130 });
131 const project = await seedProject(user, {
132 projectSlug: "overdue-cancel-rearm-project",
133 dispatchMode: "workflows",
134 });
135 const projectStub = env.PROJECT_DO.getByName(project.id);
136 
137 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
138 projectId: project.id,
139 triggeredByUserId: user.id,
140 branch: project.defaultBranch,
141 });
142 const claim = await claimRunWorkWithoutAlarm(projectStub, {
143 projectId: project.id,
144 runId: accepted.runId,
145 });
146 expect(claim.kind).toBe("execute");
147 
148 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
149 const overdueAlarmAt = Date.now() - 60_000;
150 const baseContext = createTestProjectDoContext(instance);
151 const alarmWrites: number[] = [];
152 const transaction: DurableObjectStorage["transaction"] = async (closure) =>
153 baseContext.ctx.storage.transaction(async (txn) => {
154 const patchedTxn = createBoundTransactionProxy(txn, {
155 getAlarm: async () => overdueAlarmAt,
156 setAlarm: async (scheduledTime, options) => {
157 alarmWrites.push(Number(scheduledTime));
158 },
159 });
160 
161 return await closure(patchedTxn);
162 });
163 const storage = createBoundStorageProxy(baseContext.ctx.storage, { transaction });
164 const context = withPatchedStorage(baseContext, storage);
165 
166 const result = await requestRunCancelCommand(context, {
167 projectId: project.id,
168 runId: accepted.runId,
169 });
170 
171 expect(result.status).toBe("cancel_requested");
172 expect(alarmWrites).toHaveLength(1);
173 expect(alarmWrites[0]).toBeGreaterThan(overdueAlarmAt);
174 });
175 });
176 });
177 
178 describe("sandbox cleanup retries", () => {
179 it("clears retry state after a successful alarm-driven sandbox cleanup", async () => {
180 const user = await seedUser({
181 email: "sandbox-cleanup-success@example.com",
182 slug: "sandbox-cleanup-success-user",
183 });
184 const project = await seedProject(user, {
185 projectSlug: "sandbox-cleanup-success-project",
186 });
187 const projectStub = env.PROJECT_DO.getByName(project.id);
188 
189 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
190 projectId: project.id,
191 triggeredByUserId: user.id,
192 branch: project.defaultBranch,
193 });
194 const claim = await claimRunWorkWithoutAlarm(projectStub, {
195 projectId: project.id,
196 runId: accepted.runId,
197 });
198 expect(claim.kind).toBe("execute");
199 
200 await finalizeRunExecutionWithoutAlarm(projectStub, {
201 projectId: project.id,
202 runId: accepted.runId,
203 terminalStatus: "failed",
204 lastError: "cleanup_pending",
205 sandboxDestroyed: false,
206 });
207 
208 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
209 const baseContext = createTestProjectDoContext(instance);
210 await sidecarState.setSandboxCleanupRetryState(baseContext, accepted.runId, {
211 attempt: 1,
212 nextAt: Date.now() - 1,
213 });
214 
215 const setKeepAlive = vi.fn(async () => {});
216 const destroy = vi.fn(async () => {});
217 const context = createTestProjectDoContext(
218 instance,
219 Object.assign(Object.create(baseContext.env), {
220 Sandbox: Object.assign(Object.create(baseContext.env.Sandbox), {
221 getByName: (sandboxId: string) => {
222 const stub = baseContext.env.Sandbox.getByName(sandboxId);
223 if (sandboxId !== accepted.runId) {
224 return stub;
225 }
226 
227 return Object.assign(Object.create(stub), {
228 setKeepAlive,
229 destroy,
230 });
231 },
232 }),
233 }) as Env,
234 );
235 
236 await expect(reconciliationModule.reconcileSandboxCleanup(context, project.id)).resolves.toBe(accepted.runId);
237 expect(setKeepAlive).toHaveBeenCalledWith(false);
238 expect(destroy).toHaveBeenCalledTimes(1);
239 expect(await sidecarState.getSandboxCleanupRetryState(baseContext, accepted.runId)).toBeNull();
240 });
241 });
242 });
243 
244 describe("rollback on alarm persistence failure", () => {
245 it("rolls back acceptManualRun when reconciliation alarm persistence fails", async () => {
246 const user = await seedUser({
247 email: "sidecar-arm@example.com",
248 slug: "sidecar-arm-user",
249 });
250 const project = await seedProject(user, {
251 projectSlug: "sidecar-arm-project",
252 });
253 const projectStub = env.PROJECT_DO.getByName(project.id);
254 
255 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
256 const context = createAlarmWriteFailingContext(createTestProjectDoContext(instance));
257 
258 await expect(
259 acceptManualRunCommand(context, {
260 projectId: project.id,
261 triggeredByUserId: user.id,
262 branch: project.defaultBranch,
263 }),
264 ).rejects.toThrow("alarm write failed");
265 });
266 
267 const rows = await readProjectDoRows(project.id);
268 expect(rows.state).toMatchObject({
269 projectId: project.id,
270 activeRunId: null,
271 projectIndexSyncStatus: "current",
272 });
273 expect(rows.config).not.toBeNull();
274 expect(rows.runs).toEqual([]);
275 });
276 
277 it("rolls back finalizeRunExecution when promoting the next run cannot persist an alarm", async () => {
278 const user = await seedUser({
279 email: "finalize-arm-failure@example.com",
280 slug: "finalize-arm-failure-user",
281 });
282 const project = await seedProject(user, {
283 projectSlug: "finalize-arm-failure-project",
284 dispatchMode: "queue",
285 });
286 const projectStub = env.PROJECT_DO.getByName(project.id);
287 
288 const firstAccepted = await acceptManualRunWithoutAlarm(projectStub, {
289 projectId: project.id,
290 triggeredByUserId: user.id,
291 branch: project.defaultBranch,
292 });
293 const secondAccepted = await acceptManualRunWithoutAlarm(projectStub, {
294 projectId: project.id,
295 triggeredByUserId: user.id,
296 branch: BranchName.assertDecode("release"),
297 });
298 const claim = await claimRunWorkWithoutAlarm(projectStub, {
299 projectId: project.id,
300 runId: firstAccepted.runId,
301 });
302 expect(claim.kind).toBe("execute");
303 
304 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
305 const context = createAlarmWriteFailingContext(createTestProjectDoContext(instance));
306 await context.ctx.storage.deleteAlarm();
307 
308 await expect(
309 finalizeRunExecutionCommand(context, {
310 projectId: project.id,
311 runId: firstAccepted.runId,
312 terminalStatus: "passed",
313 lastError: null,
314 sandboxDestroyed: true,
315 }),
316 ).rejects.toThrow("alarm write failed");
317 });
318 
319 const rows = await readProjectDoRows(project.id);
320 expect(rows.state?.activeRunId).toBe(firstAccepted.runId);
321 expect(
322 rows.runs.map((row) => ({ runId: row.runId, status: row.status, dispatchStatus: row.dispatchStatus })),
323 ).toEqual([
324 {
325 runId: firstAccepted.runId,
326 status: "active",
327 dispatchStatus: "started",
328 },
329 {
330 runId: secondAccepted.runId,
331 status: "pending",
332 dispatchStatus: "blocked",
333 },
334 ]);
335 });
336 
337 it("rolls back requestRunCancel when promotion cannot persist a replacement alarm", async () => {
338 const user = await seedUser({
339 email: "cancel-arm-failure@example.com",
340 slug: "cancel-arm-failure-user",
341 });
342 const project = await seedProject(user, {
343 projectSlug: "cancel-arm-failure-project",
344 });
345 const projectStub = env.PROJECT_DO.getByName(project.id);
346 
347 const firstAccepted = await acceptManualRunWithoutAlarm(projectStub, {
348 projectId: project.id,
349 triggeredByUserId: user.id,
350 branch: project.defaultBranch,
351 });
352 const secondAccepted = await acceptManualRunWithoutAlarm(projectStub, {
353 projectId: project.id,
354 triggeredByUserId: user.id,
355 branch: BranchName.assertDecode("feature-x"),
356 });
357 
358 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
359 const context = createAlarmWriteFailingContext(createTestProjectDoContext(instance));
360 
361 await expect(
362 requestRunCancelCommand(context, {
363 projectId: project.id,
364 runId: firstAccepted.runId,
365 }),
366 ).rejects.toThrow("alarm write failed");
367 });
368 
369 const rows = await readProjectDoRows(project.id);
370 expect(rows.state?.activeRunId).toBeNull();
371 expect(
372 rows.runs.map((row) => ({ runId: row.runId, status: row.status, dispatchStatus: row.dispatchStatus })),
373 ).toEqual([
374 {
375 runId: firstAccepted.runId,
376 status: "executable",
377 dispatchStatus: "pending",
378 },
379 {
380 runId: secondAccepted.runId,
381 status: "pending",
382 dispatchStatus: "blocked",
383 },
384 ]);
385 });
386 
387 it("rolls back recordRunResolvedCommit when alarm persistence fails", async () => {
388 const user = await seedUser({
389 email: "resolved-commit-arm-failure@example.com",
390 slug: "resolved-commit-arm-failure-user",
391 });
392 const project = await seedProject(user, {
393 projectSlug: "resolved-commit-arm-failure-project",
394 });
395 const projectStub = env.PROJECT_DO.getByName(project.id);
396 
397 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
398 projectId: project.id,
399 triggeredByUserId: user.id,
400 branch: project.defaultBranch,
401 });
402 const claim = await claimRunWorkWithoutAlarm(projectStub, {
403 projectId: project.id,
404 runId: accepted.runId,
405 });
406 expect(claim.kind).toBe("execute");
407 const commitSha = CommitSha.assertDecode("0123456789abcdef0123456789abcdef01234567");
408 
409 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
410 const context = createAlarmWriteFailingContext(createTestProjectDoContext(instance));
411 
412 await expect(
413 recordRunResolvedCommitCommand(context, {
414 projectId: project.id,
415 runId: accepted.runId,
416 commitSha,
417 }),
418 ).rejects.toThrow("alarm write failed");
419 });
420 
421 const rows = await readProjectDoRows(project.id);
422 expect(rows.runs[0]).toMatchObject({
423 runId: accepted.runId,
424 commitSha: null,
425 d1SyncStatus: "needs_create",
426 });
427 });
428 
429 it("rolls back claimRunWork when alarm persistence fails", async () => {
430 const user = await seedUser({
431 email: "claim-arm-failure@example.com",
432 slug: "claim-arm-failure-user",
433 });
434 const project = await seedProject(user, {
435 projectSlug: "claim-arm-failure-project",
436 });
437 const projectStub = env.PROJECT_DO.getByName(project.id);
438 
439 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
440 projectId: project.id,
441 triggeredByUserId: user.id,
442 branch: project.defaultBranch,
443 });
444 
445 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
446 const context = createAlarmWriteFailingContext(createTestProjectDoContext(instance));
447 
448 await expect(
449 claimRunWorkCommand(context, {
450 projectId: project.id,
451 runId: accepted.runId,
452 }),
453 ).rejects.toThrow("alarm write failed");
454 });
455 
456 const rows = await readProjectDoRows(project.id);
457 expect(rows.state?.activeRunId).toBeNull();
458 expect(rows.runs[0]).toMatchObject({
459 runId: accepted.runId,
460 status: "executable",
461 dispatchStatus: "pending",
462 });
463 });
464 });
465 
466 describe("alarm iteration limits", () => {
467 it("caps each alarm invocation and leaves follow-up work scheduled", async () => {
468 const user = await seedUser({
469 email: "alarm-cap@example.com",
470 slug: "alarm-cap-user",
471 });
472 const project = await seedProject(user, {
473 projectSlug: "alarm-cap-project",
474 });
475 const projectStub = env.PROJECT_DO.getByName(project.id);
476 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
477 projectId: project.id,
478 triggeredByUserId: user.id,
479 branch: project.defaultBranch,
480 });
481 const cycleSpy = vi
482 .spyOn(reconciliationModule, "runAlarmCycle")
483 .mockResolvedValue({ action: "promote_pending", runId: accepted.runId });
484 
485 try {
486 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
487 const { ctx } = getProjectDoInternals(instance);
488 await ctx.storage.deleteAlarm();
489 await instance.alarm();
490 expect(await ctx.storage.getAlarm()).not.toBeNull();
491 });
492 
493 expect(cycleSpy).toHaveBeenCalledTimes(PROJECT_ALARM_MAX_ITERATIONS);
494 } finally {
495 cycleSpy.mockRestore();
496 }
497 });
498 });
499});