Skip to content
File

Blob: tests/worker/project-do/d1-terminal.test.ts

typescript522 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 { CommitSha, DEFAULT_DISPATCH_MODE, DEFAULT_EXECUTION_RUNTIME } from "@/contracts";
7import { expectTrusted, PositiveInteger } from "@/worker/contracts";
8import { D1_RETRY_DELAYS_MS } from "@/worker/durable/project-do/constants";
9import { createD1Db } from "@/worker/db/d1";
10import * as d1Schema from "@/worker/db/d1/schema";
11import * as projectDoSchema from "@/worker/db/durable/schema/project-do";
12import { ProjectDO } from "@/worker/durable";
13import { recordRunResolvedCommit as recordRunResolvedCommitCommand } from "@/worker/durable/project-do/commands";
14import { reconcileAcceptedRunD1Sync, reconcileTerminalRunD1Sync } from "@/worker/durable/project-do/reconciliation";
15import * as sidecarState from "@/worker/durable/project-do/sidecar-state";
16 
17import {
18 acceptManualRunWithoutAlarm,
19 claimRunWorkWithoutAlarm,
20 createTestProjectDoContext,
21 expectAcceptedManualRun,
22 finalizeRunExecutionWithoutAlarm,
23 requestRunCancelWithoutAlarm,
24} from "../../helpers/project-do";
25import { registerWorkerRuntimeHooks } from "../../helpers/worker-hooks";
26import { readProjectDoRows, seedProject, seedUser } from "../../helpers/runtime";
27 
28describe("ProjectDO terminal D1 synchronization", () => {
29 registerWorkerRuntimeHooks();
30 
31 describe("post-terminal metadata guards", () => {
32 it("rejects resolved-commit backfill after terminal sync", async () => {
33 const user = await seedUser({
34 email: "commit-backfill-after-done@example.com",
35 slug: "commit-backfill-after-done-user",
36 });
37 const project = await seedProject(user, {
38 projectSlug: "commit-backfill-after-done-project",
39 });
40 const projectStub = env.PROJECT_DO.getByName(project.id);
41 
42 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
43 projectId: project.id,
44 triggeredByUserId: user.id,
45 branch: project.defaultBranch,
46 });
47 const claim = await claimRunWorkWithoutAlarm(projectStub, {
48 projectId: project.id,
49 runId: accepted.runId,
50 });
51 expect(claim.kind).toBe("execute");
52 
53 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
54 const context = createTestProjectDoContext(instance);
55 
56 await reconcileAcceptedRunD1Sync(context, project.id);
57 });
58 
59 await projectStub.finalizeRunExecution({
60 projectId: project.id,
61 runId: accepted.runId,
62 terminalStatus: "failed",
63 lastError: "checkout_failed",
64 sandboxDestroyed: true,
65 });
66 
67 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
68 const context = createTestProjectDoContext(instance);
69 
70 await reconcileTerminalRunD1Sync(context, project.id);
71 });
72 
73 const commitSha = CommitSha.assertDecode("1111111111111111111111111111111111111111");
74 await expect(
75 projectStub.recordRunResolvedCommit({
76 projectId: project.id,
77 runId: accepted.runId,
78 commitSha,
79 }),
80 ).resolves.toEqual({
81 kind: "stale",
82 status: "failed",
83 });
84 
85 const rows = await readProjectDoRows(project.id);
86 expect(rows.runs[0]).toMatchObject({
87 runId: accepted.runId,
88 status: "failed",
89 commitSha: null,
90 d1SyncStatus: "done",
91 });
92 
93 const db = createD1Db(env.DB);
94 const d1Row = await db.select().from(d1Schema.runIndex).where(eq(d1Schema.runIndex.id, accepted.runId)).limit(1);
95 expect(d1Row[0]).toMatchObject({
96 commitSha: null,
97 status: "failed",
98 });
99 expect(d1Row[0]?.finishedAt).not.toBeNull();
100 expect(d1Row[0]?.exitCode).toBe(1);
101 });
102 });
103 
104 describe("retry normalization", () => {
105 it("does not let terminal D1 reconciliation wait on a stale metadata retry window", async () => {
106 const user = await seedUser({
107 email: "terminal-retry-window@example.com",
108 slug: "terminal-retry-window-user",
109 });
110 const project = await seedProject(user, {
111 projectSlug: "terminal-retry-window-project",
112 });
113 const projectStub = env.PROJECT_DO.getByName(project.id);
114 
115 const accepted = expectAcceptedManualRun(
116 await projectStub.acceptManualRun({
117 projectId: project.id,
118 triggeredByUserId: user.id,
119 branch: project.defaultBranch,
120 }),
121 );
122 const claim = await projectStub.claimRunWork({
123 projectId: project.id,
124 runId: accepted.runId,
125 });
126 expect(claim.kind).toBe("execute");
127 
128 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
129 const context = createTestProjectDoContext(instance);
130 
131 await reconcileAcceptedRunD1Sync(context, project.id);
132 });
133 
134 const commitSha = CommitSha.assertDecode("abcdef0123456789abcdef0123456789abcdef01");
135 let remainingFailures = 2;
136 
137 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
138 const baseContext = createTestProjectDoContext(instance);
139 const failingEnv = Object.assign(Object.create(baseContext.env), {
140 RUN_DO: Object.assign(Object.create(baseContext.env.RUN_DO), {
141 getByName: (runId: string) => {
142 const stub = baseContext.env.RUN_DO.getByName(runId);
143 if (runId !== accepted.runId) {
144 return stub;
145 }
146 
147 return Object.assign(Object.create(stub), {
148 ensureInitialized: async (payload: { commitSha: string | null }) => {
149 if (payload.commitSha === commitSha && remainingFailures > 0) {
150 remainingFailures -= 1;
151 throw new Error("transient rundo reset");
152 }
153 
154 await stub.ensureInitialized(payload as Parameters<typeof stub.ensureInitialized>[0]);
155 },
156 });
157 },
158 }),
159 }) as Env;
160 const context = createTestProjectDoContext(instance, failingEnv);
161 
162 await expect(
163 recordRunResolvedCommitCommand(context, {
164 projectId: project.id,
165 runId: accepted.runId,
166 commitSha,
167 }),
168 ).resolves.toEqual({
169 kind: "applied",
170 });
171 });
172 
173 let metadataRetryAt: number | null = null;
174 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
175 const context = createTestProjectDoContext(instance);
176 const retryState = await sidecarState.getD1RetryState(context, accepted.runId);
177 expect(retryState?.attempt).toBe(1);
178 expect(retryState?.phase).toBe("metadata");
179 metadataRetryAt = retryState?.nextAt ?? null;
180 expect(metadataRetryAt).not.toBeNull();
181 await context.ctx.storage.deleteAlarm();
182 });
183 
184 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
185 const context = createTestProjectDoContext(instance);
186 
187 await context.db
188 .update(projectDoSchema.projectRuns)
189 .set({
190 status: "failed",
191 position: null,
192 dispatchStatus: "terminal",
193 d1SyncStatus: "needs_terminal_update",
194 lastError: "checkout_failed",
195 })
196 .where(eq(projectDoSchema.projectRuns.runId, accepted.runId));
197 await context.db
198 .update(projectDoSchema.projectState)
199 .set({
200 activeRunId: null,
201 updatedAt: Date.now(),
202 })
203 .where(eq(projectDoSchema.projectState.projectId, project.id));
204 await context.ctx.storage.deleteAlarm();
205 await sidecarState.rescheduleAlarm(context, project.id);
206 
207 const alarmAt = await context.ctx.storage.getAlarm();
208 expect(alarmAt).not.toBeNull();
209 expect(alarmAt as number).toBeLessThan(metadataRetryAt as number);
210 expect(await reconcileTerminalRunD1Sync(context, project.id)).toBe(accepted.runId);
211 expect(await sidecarState.getD1RetryState(context, accepted.runId)).toBeNull();
212 });
213 
214 const runMeta = await env.RUN_DO.getByName(accepted.runId).getRunSummary(accepted.runId);
215 expect(runMeta?.status).toBe("failed");
216 expect(runMeta?.finishedAt).not.toBeNull();
217 
218 const rows = await readProjectDoRows(project.id);
219 expect(rows.runs[0]?.status).toBe("failed");
220 expect(rows.runs[0]?.d1SyncStatus).toBe("done");
221 
222 const db = createD1Db(env.DB);
223 const d1Row = await db.select().from(d1Schema.runIndex).where(eq(d1Schema.runIndex.id, accepted.runId)).limit(1);
224 expect(d1Row[0]?.commitSha).toBe(commitSha);
225 expect(d1Row[0]?.status).toBe("failed");
226 expect(d1Row[0]?.finishedAt).not.toBeNull();
227 expect(d1Row[0]?.exitCode).toBe(1);
228 });
229 
230 it("resets the terminal retry attempt when the stored retry came from metadata", async () => {
231 const baseTime = 1_740_000_000_000;
232 const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => baseTime);
233 
234 try {
235 const user = await seedUser({
236 email: "terminal-retry-reset@example.com",
237 slug: "terminal-retry-reset-user",
238 });
239 const project = await seedProject(user, {
240 projectSlug: "terminal-retry-reset-project",
241 });
242 const projectStub = env.PROJECT_DO.getByName(project.id);
243 
244 const accepted = expectAcceptedManualRun(
245 await projectStub.acceptManualRun({
246 projectId: project.id,
247 triggeredByUserId: user.id,
248 branch: project.defaultBranch,
249 }),
250 );
251 const claim = await projectStub.claimRunWork({
252 projectId: project.id,
253 runId: accepted.runId,
254 });
255 expect(claim.kind).toBe("execute");
256 
257 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
258 const context = createTestProjectDoContext(instance);
259 
260 await reconcileAcceptedRunD1Sync(context, project.id);
261 await context.db
262 .update(projectDoSchema.projectRuns)
263 .set({
264 status: "failed",
265 position: null,
266 dispatchStatus: "terminal",
267 d1SyncStatus: "needs_terminal_update",
268 lastError: "checkout_failed",
269 })
270 .where(eq(projectDoSchema.projectRuns.runId, accepted.runId));
271 await context.db
272 .update(projectDoSchema.projectState)
273 .set({
274 activeRunId: null,
275 updatedAt: accepted.queuedAt,
276 })
277 .where(eq(projectDoSchema.projectState.projectId, project.id));
278 await sidecarState.setD1RetryState(context, accepted.runId, {
279 attempt: 3,
280 nextAt: baseTime + 60_000,
281 phase: "metadata",
282 });
283 
284 const failingEnv = Object.assign(Object.create(context.env), {
285 RUN_DO: Object.assign(Object.create(context.env.RUN_DO), {
286 getByName: (runId: string) => {
287 const stub = context.env.RUN_DO.getByName(runId);
288 if (runId !== accepted.runId) {
289 return stub;
290 }
291 
292 return Object.assign(Object.create(stub), {
293 getRunSummary: async () => null,
294 ensureInitialized: async () => {},
295 tryUpdateRunState: async () => {
296 throw new Error("terminal write failed");
297 },
298 });
299 },
300 }),
301 }) as Env;
302 const failingContext = createTestProjectDoContext(instance, failingEnv);
303 
304 await expect(reconcileTerminalRunD1Sync(failingContext, project.id)).resolves.toBeNull();
305 
306 const retryState = await sidecarState.getD1RetryState(context, accepted.runId);
307 expect(retryState).toEqual({
308 attempt: 1,
309 nextAt: baseTime + D1_RETRY_DELAYS_MS[0],
310 phase: "terminal",
311 });
312 });
313 } finally {
314 nowSpy.mockRestore();
315 }
316 });
317 
318 it("ignores stored D1 retries that are missing a phase", async () => {
319 const user = await seedUser({
320 email: "retry-missing-phase@example.com",
321 slug: "retry-missing-phase-user",
322 });
323 const project = await seedProject(user, {
324 projectSlug: "retry-missing-phase-project",
325 });
326 const projectStub = env.PROJECT_DO.getByName(project.id);
327 
328 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
329 projectId: project.id,
330 triggeredByUserId: user.id,
331 branch: project.defaultBranch,
332 });
333 
334 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
335 const context = createTestProjectDoContext(instance);
336 const futureRetryAt = Date.now() + 60_000;
337 
338 await context.ctx.storage.put(sidecarState.d1RetryKey(accepted.runId), {
339 attempt: 4,
340 nextAt: futureRetryAt,
341 });
342 
343 expect(await sidecarState.getD1RetryState(context, accepted.runId)).toBeNull();
344 
345 await sidecarState.rescheduleAlarm(context, project.id);
346 
347 const alarmAt = await context.ctx.storage.getAlarm();
348 expect(alarmAt).not.toBeNull();
349 expect(alarmAt as number).toBeLessThan(futureRetryAt);
350 });
351 });
352 });
353 
354 describe("terminal reconciliation", () => {
355 it("reconciles passed RunDO terminal state during terminal reconciliation", async () => {
356 const user = await seedUser({
357 email: "passed-fixup@example.com",
358 slug: "passed-fixup-user",
359 });
360 const project = await seedProject(user, {
361 projectSlug: "passed-fixup-project",
362 });
363 const projectStub = env.PROJECT_DO.getByName(project.id);
364 
365 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
366 projectId: project.id,
367 triggeredByUserId: user.id,
368 branch: project.defaultBranch,
369 });
370 const claim = await claimRunWorkWithoutAlarm(projectStub, {
371 projectId: project.id,
372 runId: accepted.runId,
373 });
374 expect(claim.kind).toBe("execute");
375 
376 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
377 const context = createTestProjectDoContext(instance);
378 
379 await reconcileAcceptedRunD1Sync(context, project.id);
380 });
381 
382 const runStub = env.RUN_DO.getByName(accepted.runId);
383 await runStub.updateRunState({
384 runId: accepted.runId,
385 status: "starting",
386 currentStep: null,
387 startedAt: accepted.queuedAt,
388 finishedAt: null,
389 exitCode: null,
390 errorMessage: null,
391 });
392 await runStub.updateRunState({
393 runId: accepted.runId,
394 status: "running",
395 currentStep: null,
396 startedAt: accepted.queuedAt,
397 finishedAt: null,
398 exitCode: 0,
399 errorMessage: null,
400 });
401 
402 await finalizeRunExecutionWithoutAlarm(projectStub, {
403 projectId: project.id,
404 runId: accepted.runId,
405 terminalStatus: "passed",
406 lastError: null,
407 });
408 
409 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
410 const context = createTestProjectDoContext(instance);
411 
412 await reconcileTerminalRunD1Sync(context, project.id);
413 });
414 
415 const runMeta = await runStub.getRunSummary(accepted.runId);
416 expect(runMeta?.status).toBe("passed");
417 expect(runMeta?.finishedAt).not.toBeNull();
418 
419 const db = createD1Db(env.DB);
420 const d1Row = await db.select().from(d1Schema.runIndex).where(eq(d1Schema.runIndex.id, accepted.runId)).limit(1);
421 expect(d1Row[0]?.status).toBe("passed");
422 });
423 
424 it("repairs the active step before reconciling canceled RunDO terminal state", async () => {
425 const user = await seedUser({
426 email: "canceled-step-repair@example.com",
427 slug: "canceled-step-repair-user",
428 });
429 const project = await seedProject(user, {
430 projectSlug: "canceled-step-repair-project",
431 });
432 const projectStub = env.PROJECT_DO.getByName(project.id);
433 
434 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
435 projectId: project.id,
436 triggeredByUserId: user.id,
437 branch: project.defaultBranch,
438 });
439 const claim = await claimRunWorkWithoutAlarm(projectStub, {
440 projectId: project.id,
441 runId: accepted.runId,
442 });
443 expect(claim.kind).toBe("execute");
444 
445 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
446 const context = createTestProjectDoContext(instance);
447 
448 await reconcileAcceptedRunD1Sync(context, project.id);
449 });
450 
451 const stepPosition = expectTrusted(PositiveInteger, 1, "PositiveInteger");
452 const runStub = env.RUN_DO.getByName(accepted.runId);
453 await runStub.replaceSteps({
454 runId: accepted.runId,
455 steps: [
456 {
457 position: stepPosition,
458 name: "Install",
459 command: "npm ci",
460 },
461 ],
462 });
463 await runStub.updateRunState({
464 runId: accepted.runId,
465 status: "starting",
466 currentStep: null,
467 startedAt: accepted.queuedAt,
468 finishedAt: null,
469 exitCode: null,
470 errorMessage: null,
471 });
472 await runStub.updateStepState({
473 runId: accepted.runId,
474 position: stepPosition,
475 status: "running",
476 startedAt: accepted.queuedAt,
477 finishedAt: null,
478 exitCode: null,
479 });
480 await runStub.updateRunState({
481 runId: accepted.runId,
482 status: "running",
483 currentStep: stepPosition,
484 startedAt: accepted.queuedAt,
485 finishedAt: null,
486 exitCode: null,
487 errorMessage: null,
488 });
489 
490 const cancelResult = await requestRunCancelWithoutAlarm(projectStub, {
491 projectId: project.id,
492 runId: accepted.runId,
493 });
494 expect(cancelResult.status).toBe("cancel_requested");
495 
496 await finalizeRunExecutionWithoutAlarm(projectStub, {
497 projectId: project.id,
498 runId: accepted.runId,
499 terminalStatus: "canceled",
500 lastError: null,
501 });
502 
503 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
504 const context = createTestProjectDoContext(instance);
505 
506 await reconcileTerminalRunD1Sync(context, project.id);
507 });
508 
509 const detail = await runStub.getRunDetail(accepted.runId);
510 expect(detail.meta?.status).toBe("canceled");
511 expect(detail.meta?.currentStep).toBeNull();
512 expect(detail.steps[0]?.status).toBe("failed");
513 expect(detail.steps[0]?.finishedAt).not.toBeNull();
514 expect(detail.steps[0]?.exitCode).toBeNull();
515 
516 const db = createD1Db(env.DB);
517 const d1Row = await db.select().from(d1Schema.runIndex).where(eq(d1Schema.runIndex.id, accepted.runId)).limit(1);
518 expect(d1Row[0]?.status).toBe("canceled");
519 });
520 });
521});