Skip to content
File

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

typescript550 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, UnixTimestampMs } from "@/contracts";
7import { expectTrusted } from "@/worker/contracts";
8import { createD1Db } from "@/worker/db/d1";
9import * as d1Repositories from "@/worker/db/d1/repositories";
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 {
15 reconcileAcceptedRunD1Sync,
16 reconcileRunMetadataD1Sync,
17 reconcileTerminalRunD1Sync,
18} from "@/worker/durable/project-do/reconciliation";
19import * as sidecarState from "@/worker/durable/project-do/sidecar-state";
20 
21import {
22 acceptManualRunWithoutAlarm,
23 claimRunWorkWithoutAlarm,
24 createTestProjectDoContext,
25 expectAcceptedManualRun,
26} from "../../helpers/project-do";
27import { registerWorkerRuntimeHooks } from "../../helpers/worker-hooks";
28import { readProjectDoRows, seedProject, seedUser } from "../../helpers/runtime";
29 
30describe("ProjectDO D1 synchronization", () => {
31 registerWorkerRuntimeHooks();
32 
33 describe("resolved commit metadata backfill", () => {
34 it("backfills a manual run commit SHA into ProjectDO, RunDO, and D1", async () => {
35 const user = await seedUser({
36 email: "commit-backfill@example.com",
37 slug: "commit-backfill-user",
38 });
39 const project = await seedProject(user, {
40 projectSlug: "commit-backfill-project",
41 });
42 const stub = env.PROJECT_DO.getByName(project.id);
43 
44 const accepted = await acceptManualRunWithoutAlarm(stub, {
45 projectId: project.id,
46 triggeredByUserId: user.id,
47 branch: project.defaultBranch,
48 });
49 const claim = await claimRunWorkWithoutAlarm(stub, {
50 projectId: project.id,
51 runId: accepted.runId,
52 });
53 expect(claim.kind).toBe("execute");
54 const commitSha = CommitSha.assertDecode("0123456789abcdef0123456789abcdef01234567");
55 
56 await expect(
57 stub.recordRunResolvedCommit({
58 projectId: project.id,
59 runId: accepted.runId,
60 commitSha,
61 }),
62 ).resolves.toEqual({
63 kind: "applied",
64 });
65 
66 const rows = await readProjectDoRows(project.id);
67 expect(rows.runs[0]?.commitSha).toBe(commitSha);
68 
69 const runMeta = await env.RUN_DO.getByName(accepted.runId).getRunSummary(accepted.runId);
70 expect(runMeta?.commitSha).toBe(commitSha);
71 
72 const db = createD1Db(env.DB);
73 const d1Row = await db.select().from(d1Schema.runIndex).where(eq(d1Schema.runIndex.id, accepted.runId)).limit(1);
74 expect(d1Row[0]?.commitSha).toBe(commitSha);
75 });
76 
77 it("keeps resolved-commit backfill non-fatal when RunDO fan-out fails and repairs it on retry", async () => {
78 const user = await seedUser({
79 email: "commit-backfill-retry@example.com",
80 slug: "commit-backfill-retry-user",
81 });
82 const project = await seedProject(user, {
83 projectSlug: "commit-backfill-retry-project",
84 });
85 const projectStub = env.PROJECT_DO.getByName(project.id);
86 
87 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
88 projectId: project.id,
89 triggeredByUserId: user.id,
90 branch: project.defaultBranch,
91 });
92 const claim = await claimRunWorkWithoutAlarm(projectStub, {
93 projectId: project.id,
94 runId: accepted.runId,
95 });
96 expect(claim.kind).toBe("execute");
97 
98 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
99 const context = createTestProjectDoContext(instance);
100 
101 await reconcileAcceptedRunD1Sync(context, project.id);
102 });
103 
104 const commitSha = CommitSha.assertDecode("89abcdef0123456789abcdef0123456789abcdef");
105 let remainingFailures = 2;
106 
107 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
108 const baseContext = createTestProjectDoContext(instance);
109 const failingEnv = Object.assign(Object.create(baseContext.env), {
110 RUN_DO: Object.assign(Object.create(baseContext.env.RUN_DO), {
111 getByName: (runId: string) => {
112 const stub = baseContext.env.RUN_DO.getByName(runId);
113 if (runId !== accepted.runId) {
114 return stub;
115 }
116 
117 return Object.assign(Object.create(stub), {
118 ensureInitialized: async (payload: { commitSha: string | null }) => {
119 if (payload.commitSha === commitSha && remainingFailures > 0) {
120 remainingFailures -= 1;
121 throw new Error("transient rundo reset");
122 }
123 
124 await stub.ensureInitialized(payload as Parameters<typeof stub.ensureInitialized>[0]);
125 },
126 });
127 },
128 }),
129 }) as Env;
130 const context = createTestProjectDoContext(instance, failingEnv);
131 
132 await expect(
133 recordRunResolvedCommitCommand(context, {
134 projectId: project.id,
135 runId: accepted.runId,
136 commitSha,
137 }),
138 ).resolves.toEqual({
139 kind: "applied",
140 });
141 });
142 
143 const rows = await readProjectDoRows(project.id);
144 expect(rows.runs[0]?.commitSha).toBe(commitSha);
145 expect(rows.runs[0]?.d1SyncStatus).toBe("needs_update");
146 
147 const staleRunMeta = await env.RUN_DO.getByName(accepted.runId).getRunSummary(accepted.runId);
148 expect(staleRunMeta?.commitSha).toBeNull();
149 
150 const staleDb = createD1Db(env.DB);
151 const staleD1Row = await staleDb
152 .select()
153 .from(d1Schema.runIndex)
154 .where(eq(d1Schema.runIndex.id, accepted.runId))
155 .limit(1);
156 expect(staleD1Row[0]).toMatchObject({
157 commitSha,
158 status: "queued",
159 startedAt: null,
160 finishedAt: null,
161 exitCode: null,
162 });
163 
164 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
165 const context = createTestProjectDoContext(instance);
166 
167 await sidecarState.setD1RetryState(context, accepted.runId, null);
168 await reconcileRunMetadataD1Sync(context, project.id);
169 });
170 
171 const repairedRows = await readProjectDoRows(project.id);
172 expect(repairedRows.runs[0]?.commitSha).toBe(commitSha);
173 expect(repairedRows.runs[0]?.d1SyncStatus).toBe("current");
174 
175 const repairedRunMeta = await env.RUN_DO.getByName(accepted.runId).getRunSummary(accepted.runId);
176 expect(repairedRunMeta?.commitSha).toBe(commitSha);
177 
178 const repairedDb = createD1Db(env.DB);
179 const repairedD1Row = await repairedDb
180 .select()
181 .from(d1Schema.runIndex)
182 .where(eq(d1Schema.runIndex.id, accepted.runId))
183 .limit(1);
184 expect(repairedD1Row[0]?.commitSha).toBe(commitSha);
185 });
186 });
187 
188 describe("initial create synchronization", () => {
189 it("creates the initial D1 row even when RunDO is unavailable during create sync", async () => {
190 const user = await seedUser({
191 email: "create-sync-rundo-unavailable@example.com",
192 slug: "create-sync-rundo-unavailable-user",
193 });
194 const project = await seedProject(user, {
195 projectSlug: "create-sync-rundo-unavailable-project",
196 });
197 const projectStub = env.PROJECT_DO.getByName(project.id);
198 
199 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
200 projectId: project.id,
201 triggeredByUserId: user.id,
202 branch: project.defaultBranch,
203 });
204 
205 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
206 const baseContext = createTestProjectDoContext(instance);
207 const failingEnv = Object.assign(Object.create(baseContext.env), {
208 RUN_DO: Object.assign(Object.create(baseContext.env.RUN_DO), {
209 getByName: () => {
210 throw new Error("RunDO unavailable during create sync");
211 },
212 }),
213 }) as Env;
214 const context = createTestProjectDoContext(instance, failingEnv);
215 
216 await expect(reconcileAcceptedRunD1Sync(context, project.id)).resolves.toBe(accepted.runId);
217 });
218 
219 const rows = await readProjectDoRows(project.id);
220 expect(rows.runs[0]?.d1SyncStatus).toBe("current");
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]).toMatchObject({
225 status: "queued",
226 commitSha: null,
227 startedAt: null,
228 finishedAt: null,
229 exitCode: null,
230 });
231 });
232 
233 it("does not let create sync overwrite fresher resolved-commit metadata", async () => {
234 const user = await seedUser({
235 email: "create-sync-race@example.com",
236 slug: "create-sync-race-user",
237 });
238 const project = await seedProject(user, {
239 projectSlug: "create-sync-race-project",
240 });
241 const projectStub = env.PROJECT_DO.getByName(project.id);
242 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
243 projectId: project.id,
244 triggeredByUserId: user.id,
245 branch: project.defaultBranch,
246 });
247 const claim = await claimRunWorkWithoutAlarm(projectStub, {
248 projectId: project.id,
249 runId: accepted.runId,
250 });
251 expect(claim.kind).toBe("execute");
252 if (claim.kind !== "execute") {
253 throw new Error("Expected claimRunWork to return an executable snapshot.");
254 }
255 
256 const runStub = env.RUN_DO.getByName(accepted.runId);
257 await runStub.ensureInitialized(claim.snapshot);
258 await runStub.updateRunState({
259 runId: accepted.runId,
260 status: "starting",
261 currentStep: null,
262 startedAt: accepted.queuedAt,
263 finishedAt: null,
264 exitCode: null,
265 errorMessage: null,
266 });
267 
268 const commitSha = CommitSha.assertDecode("fedcba9876543210fedcba9876543210fedcba98");
269 let injectedResolvedCommitSync = false;
270 
271 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
272 const baseContext = createTestProjectDoContext(instance);
273 const originalUpsertRunIndex = d1Repositories.upsertRunIndex;
274 const upsertSpy = vi
275 .spyOn(d1Repositories, "upsertRunIndex")
276 .mockImplementation(async (...args: Parameters<typeof d1Repositories.upsertRunIndex>) => {
277 const [, row] = args;
278 
279 if (!injectedResolvedCommitSync && row.id === accepted.runId) {
280 injectedResolvedCommitSync = true;
281 await recordRunResolvedCommitCommand(baseContext, {
282 projectId: project.id,
283 runId: accepted.runId,
284 commitSha,
285 });
286 }
287 
288 return await originalUpsertRunIndex(...args);
289 });
290 
291 try {
292 await expect(reconcileAcceptedRunD1Sync(baseContext, project.id)).resolves.toBe(accepted.runId);
293 } finally {
294 upsertSpy.mockRestore();
295 }
296 });
297 
298 const rows = await readProjectDoRows(project.id);
299 expect(rows.runs[0]?.commitSha).toBe(commitSha);
300 expect(rows.runs[0]?.d1SyncStatus).toBe("current");
301 
302 const runMeta = await runStub.getRunSummary(accepted.runId);
303 expect(runMeta).toMatchObject({
304 commitSha,
305 status: "starting",
306 startedAt: accepted.queuedAt,
307 });
308 
309 const db = createD1Db(env.DB);
310 const d1Row = await db.select().from(d1Schema.runIndex).where(eq(d1Schema.runIndex.id, accepted.runId)).limit(1);
311 expect(d1Row[0]).toMatchObject({
312 commitSha,
313 status: "queued",
314 startedAt: null,
315 finishedAt: null,
316 exitCode: null,
317 });
318 });
319 });
320 
321 describe("terminal synchronization guards", () => {
322 it("does not publish terminal D1 state before ProjectDO finalizes the run", async () => {
323 const user = await seedUser({
324 email: "terminal-before-finalize@example.com",
325 slug: "terminal-before-finalize-user",
326 });
327 const project = await seedProject(user, {
328 projectSlug: "terminal-before-finalize-project",
329 });
330 const projectStub = env.PROJECT_DO.getByName(project.id);
331 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
332 projectId: project.id,
333 triggeredByUserId: user.id,
334 branch: project.defaultBranch,
335 });
336 const claim = await claimRunWorkWithoutAlarm(projectStub, {
337 projectId: project.id,
338 runId: accepted.runId,
339 });
340 expect(claim.kind).toBe("execute");
341 if (claim.kind !== "execute") {
342 throw new Error("Expected claimRunWork to return an executable snapshot.");
343 }
344 
345 const runStub = env.RUN_DO.getByName(accepted.runId);
346 const finishedAt = expectTrusted(UnixTimestampMs, accepted.queuedAt + 5_000, "UnixTimestampMs");
347 
348 await runStub.ensureInitialized(claim.snapshot);
349 await runStub.updateRunState({
350 runId: accepted.runId,
351 status: "starting",
352 currentStep: null,
353 startedAt: accepted.queuedAt,
354 finishedAt: null,
355 exitCode: null,
356 errorMessage: null,
357 });
358 await runStub.updateRunState({
359 runId: accepted.runId,
360 status: "running",
361 currentStep: null,
362 startedAt: accepted.queuedAt,
363 finishedAt: null,
364 exitCode: 0,
365 errorMessage: null,
366 });
367 await runStub.updateRunState({
368 runId: accepted.runId,
369 status: "failed",
370 currentStep: null,
371 startedAt: accepted.queuedAt,
372 finishedAt,
373 exitCode: 1,
374 errorMessage: "checkout_failed",
375 });
376 
377 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
378 const context = createTestProjectDoContext(instance);
379 
380 await expect(reconcileAcceptedRunD1Sync(context, project.id)).resolves.toBe(accepted.runId);
381 });
382 
383 const preFinalizeRows = await readProjectDoRows(project.id);
384 expect(preFinalizeRows.state?.activeRunId).toBe(accepted.runId);
385 expect(preFinalizeRows.runs[0]?.status).toBe("active");
386 expect(preFinalizeRows.runs[0]?.d1SyncStatus).toBe("current");
387 
388 const db = createD1Db(env.DB);
389 const staleD1Row = await db
390 .select()
391 .from(d1Schema.runIndex)
392 .where(eq(d1Schema.runIndex.id, accepted.runId))
393 .limit(1);
394 expect(staleD1Row[0]).toMatchObject({
395 status: "queued",
396 startedAt: null,
397 finishedAt: null,
398 exitCode: null,
399 });
400 
401 await projectStub.finalizeRunExecution({
402 projectId: project.id,
403 runId: accepted.runId,
404 terminalStatus: "failed",
405 lastError: "checkout_failed",
406 sandboxDestroyed: true,
407 });
408 
409 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
410 const context = createTestProjectDoContext(instance);
411 
412 await expect(reconcileTerminalRunD1Sync(context, project.id)).resolves.toBe(accepted.runId);
413 });
414 
415 const finalizedRows = await readProjectDoRows(project.id);
416 expect(finalizedRows.state?.activeRunId).toBeNull();
417 expect(finalizedRows.runs[0]?.status).toBe("failed");
418 expect(finalizedRows.runs[0]?.d1SyncStatus).toBe("done");
419 
420 const terminalD1Row = await db
421 .select()
422 .from(d1Schema.runIndex)
423 .where(eq(d1Schema.runIndex.id, accepted.runId))
424 .limit(1);
425 expect(terminalD1Row[0]).toMatchObject({
426 status: "failed",
427 startedAt: accepted.queuedAt,
428 finishedAt,
429 exitCode: 1,
430 });
431 });
432 
433 it("does not let resolved-commit metadata sync regress an already-terminal D1 row", async () => {
434 const user = await seedUser({
435 email: "commit-backfill-terminal@example.com",
436 slug: "commit-backfill-terminal-user",
437 });
438 const project = await seedProject(user, {
439 projectSlug: "commit-backfill-terminal-project",
440 });
441 const projectStub = env.PROJECT_DO.getByName(project.id);
442 
443 const accepted = expectAcceptedManualRun(
444 await projectStub.acceptManualRun({
445 projectId: project.id,
446 triggeredByUserId: user.id,
447 branch: project.defaultBranch,
448 }),
449 );
450 const claim = await projectStub.claimRunWork({
451 projectId: project.id,
452 runId: accepted.runId,
453 });
454 expect(claim.kind).toBe("execute");
455 
456 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
457 const context = createTestProjectDoContext(instance);
458 
459 await reconcileAcceptedRunD1Sync(context, project.id);
460 });
461 
462 await env.RUN_DO.getByName(accepted.runId).updateRunState({
463 runId: accepted.runId,
464 status: "starting",
465 currentStep: null,
466 startedAt: accepted.queuedAt,
467 finishedAt: null,
468 exitCode: null,
469 errorMessage: null,
470 });
471 
472 const commitSha = CommitSha.assertDecode("76543210fedcba9876543210fedcba9876543210");
473 const finishedAt = accepted.queuedAt + 1;
474 let injectedTerminalState = false;
475 
476 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
477 const baseContext = createTestProjectDoContext(instance);
478 const originalUpsertRunIndex = d1Repositories.upsertRunIndex;
479 const upsertSpy = vi
480 .spyOn(d1Repositories, "upsertRunIndex")
481 .mockImplementation(async (...args: Parameters<typeof d1Repositories.upsertRunIndex>) => {
482 const [, row] = args;
483 
484 if (!injectedTerminalState && row.id === accepted.runId) {
485 injectedTerminalState = true;
486 await baseContext.db
487 .update(projectDoSchema.projectRuns)
488 .set({
489 status: "failed",
490 position: null,
491 dispatchStatus: "terminal",
492 d1SyncStatus: "done",
493 lastError: "runner_lost",
494 })
495 .where(eq(projectDoSchema.projectRuns.runId, accepted.runId));
496 await baseContext.db
497 .update(projectDoSchema.projectState)
498 .set({
499 activeRunId: null,
500 updatedAt: finishedAt,
501 })
502 .where(eq(projectDoSchema.projectState.projectId, project.id));
503 
504 const db = createD1Db(baseContext.env.DB);
505 await db
506 .update(d1Schema.runIndex)
507 .set({
508 status: "failed",
509 startedAt: accepted.queuedAt,
510 finishedAt,
511 exitCode: 1,
512 })
513 .where(eq(d1Schema.runIndex.id, accepted.runId));
514 }
515 
516 return await originalUpsertRunIndex(...args);
517 });
518 
519 try {
520 await expect(
521 recordRunResolvedCommitCommand(baseContext, {
522 projectId: project.id,
523 runId: accepted.runId,
524 commitSha,
525 }),
526 ).resolves.toEqual({
527 kind: "applied",
528 });
529 } finally {
530 upsertSpy.mockRestore();
531 }
532 });
533 
534 const db = createD1Db(env.DB);
535 const d1Row = await db.select().from(d1Schema.runIndex).where(eq(d1Schema.runIndex.id, accepted.runId)).limit(1);
536 expect(d1Row[0]).toMatchObject({
537 commitSha,
538 status: "failed",
539 startedAt: accepted.queuedAt,
540 finishedAt,
541 exitCode: 1,
542 });
543 
544 const rows = await readProjectDoRows(project.id);
545 expect(rows.runs[0]?.status).toBe("failed");
546 expect(rows.runs[0]?.d1SyncStatus).toBe("done");
547 });
548 });
549});