Skip to content
File

Blob: tests/worker/project-do/watchdog-dispatch.test.ts

typescript723 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 { DEFAULT_DISPATCH_MODE, DEFAULT_EXECUTION_RUNTIME, UnixTimestampMs } from "@/contracts";
7import { expectTrusted, PositiveInteger } from "@/worker/contracts";
8import { DISPATCH_RETRY_DELAYS_MS, HEARTBEAT_STALE_AFTER_MS } from "@/worker/durable/project-do/constants";
9import { createD1Db } from "@/worker/db/d1";
10import * as d1Schema from "@/worker/db/d1/schema";
11import { ProjectDO } from "@/worker/durable";
12import {
13 dispatchExecutableRun,
14 reconcileAcceptedRunD1Sync,
15 reconcileActiveRunWatchdog,
16 reconcileTerminalRunD1Sync,
17} from "@/worker/durable/project-do/reconciliation";
18import { recoverWorkflowDispatchFailure } from "@/worker/durable/project-do/commands";
19import {
20 getDispatchRetryAt,
21 getSandboxCleanupRetryState,
22 setHeartbeatAt,
23} from "@/worker/durable/project-do/sidecar-state";
24import type { ProjectDoContext } from "@/worker/durable/project-do/types";
25 
26import {
27 acceptManualRunWithoutAlarm,
28 claimRunWorkWithoutAlarm,
29 createTestProjectDoContext,
30 expectAcceptedManualRun,
31} from "../../helpers/project-do";
32import { registerWorkerRuntimeHooks } from "../../helpers/worker-hooks";
33import { readProjectDoRows, seedProject, seedUser } from "../../helpers/runtime";
34 
35describe("ProjectDO watchdog and dispatch recovery", () => {
36 registerWorkerRuntimeHooks();
37 
38 it("fails stale active runs as runner_lost during watchdog recovery", async () => {
39 const baseTime = 1_710_000_000_000;
40 const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => baseTime);
41 
42 try {
43 const user = await seedUser({
44 email: "watchdog@example.com",
45 slug: "watchdog-user",
46 });
47 const project = await seedProject(user, {
48 projectSlug: "watchdog-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 const claim = await claimRunWorkWithoutAlarm(projectStub, {
59 projectId: project.id,
60 runId: accepted.runId,
61 });
62 expect(claim.kind).toBe("execute");
63 
64 nowSpy.mockImplementation(() => baseTime + HEARTBEAT_STALE_AFTER_MS + 1);
65 await runInDurableObject(env.PROJECT_DO.getByName(project.id), async (instance: ProjectDO) => {
66 await instance.alarm();
67 });
68 
69 const rows = await readProjectDoRows(project.id);
70 expect(rows.state?.activeRunId).toBeNull();
71 expect(rows.runs[0]?.status).toBe("failed");
72 expect(rows.runs[0]?.lastError).toBe("runner_lost");
73 
74 const runMeta = await env.RUN_DO.getByName(accepted.runId).getRunSummary(accepted.runId);
75 expect(runMeta?.status).toBe("failed");
76 expect(runMeta?.errorMessage).toBe("runner_lost");
77 expect(runMeta?.startedAt).toBeNull();
78 
79 const db = createD1Db(env.DB);
80 const d1Row = await db.select().from(d1Schema.runIndex).where(eq(d1Schema.runIndex.id, accepted.runId)).limit(1);
81 expect(d1Row[0]?.status).toBe("failed");
82 expect(d1Row[0]?.startedAt).toBeNull();
83 } finally {
84 nowSpy.mockRestore();
85 }
86 });
87 
88 it("repairs the active step before failing a stale active run during watchdog recovery", async () => {
89 const baseTime = 1_712_000_000_000;
90 const startedAt = expectTrusted(UnixTimestampMs, baseTime, "UnixTimestampMs");
91 const user = await seedUser({
92 email: "watchdog-step-repair@example.com",
93 slug: "watchdog-step-repair-user",
94 });
95 const project = await seedProject(user, {
96 projectSlug: "watchdog-step-repair-project",
97 });
98 const projectStub = env.PROJECT_DO.getByName(project.id);
99 
100 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
101 projectId: project.id,
102 triggeredByUserId: user.id,
103 branch: project.defaultBranch,
104 });
105 
106 const claim = await claimRunWorkWithoutAlarm(projectStub, {
107 projectId: project.id,
108 runId: accepted.runId,
109 });
110 expect(claim.kind).toBe("execute");
111 
112 const stepPosition = expectTrusted(PositiveInteger, 1, "PositiveInteger");
113 const runStub = env.RUN_DO.getByName(accepted.runId);
114 await runStub.replaceSteps({
115 runId: accepted.runId,
116 steps: [
117 {
118 position: stepPosition,
119 name: "Install",
120 command: "npm ci",
121 },
122 ],
123 });
124 await runStub.updateRunState({
125 runId: accepted.runId,
126 status: "starting",
127 currentStep: null,
128 startedAt,
129 finishedAt: null,
130 exitCode: null,
131 errorMessage: null,
132 });
133 await runStub.updateStepState({
134 runId: accepted.runId,
135 position: stepPosition,
136 status: "running",
137 startedAt,
138 finishedAt: null,
139 exitCode: null,
140 });
141 await runStub.updateRunState({
142 runId: accepted.runId,
143 status: "running",
144 currentStep: stepPosition,
145 startedAt,
146 finishedAt: null,
147 exitCode: null,
148 errorMessage: null,
149 });
150 
151 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
152 const context = createTestProjectDoContext(instance);
153 await setHeartbeatAt(context, accepted.runId, baseTime);
154 await expect(reconcileActiveRunWatchdog(context, project.id)).resolves.toBe(accepted.runId);
155 });
156 
157 const detail = await runStub.getRunDetail(accepted.runId);
158 expect(detail.meta?.status).toBe("failed");
159 expect(detail.meta?.currentStep).toBeNull();
160 expect(detail.meta?.errorMessage).toBe("runner_lost");
161 expect(detail.steps[0]?.status).toBe("failed");
162 expect(detail.steps[0]?.finishedAt).not.toBeNull();
163 });
164 
165 it("persists sandbox cleanup retry state when watchdog cleanup fails", async () => {
166 const baseTime = 1_715_000_000_000;
167 const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => baseTime);
168 
169 try {
170 const user = await seedUser({
171 email: "watchdog-cleanup@example.com",
172 slug: "watchdog-cleanup-user",
173 });
174 const project = await seedProject(user, {
175 projectSlug: "watchdog-cleanup-project",
176 });
177 const projectStub = env.PROJECT_DO.getByName(project.id);
178 
179 const accepted = expectAcceptedManualRun(
180 await projectStub.acceptManualRun({
181 projectId: project.id,
182 triggeredByUserId: user.id,
183 branch: project.defaultBranch,
184 }),
185 );
186 
187 const claim = await projectStub.claimRunWork({
188 projectId: project.id,
189 runId: accepted.runId,
190 });
191 expect(claim.kind).toBe("execute");
192 
193 nowSpy.mockImplementation(() => baseTime + HEARTBEAT_STALE_AFTER_MS + 1);
194 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
195 const baseContext = createTestProjectDoContext(instance);
196 const context: ProjectDoContext = {
197 ...baseContext,
198 env: Object.assign(Object.create(baseContext.env), {
199 Sandbox: Object.assign(Object.create(baseContext.env.Sandbox), {
200 getByName: (sandboxId: string) => {
201 const stub = baseContext.env.Sandbox.getByName(sandboxId);
202 if (sandboxId !== accepted.runId) {
203 return stub;
204 }
205 
206 return Object.assign(Object.create(stub), {
207 setKeepAlive: vi.fn(async () => {}),
208 destroy: vi.fn(async () => {
209 throw new Error("sandbox still reachable");
210 }),
211 });
212 },
213 }),
214 }) as Env,
215 };
216 
217 await expect(reconcileActiveRunWatchdog(context, project.id)).resolves.toBe(accepted.runId);
218 
219 const retryState = await getSandboxCleanupRetryState(baseContext, accepted.runId);
220 expect(retryState?.attempt).toBe(1);
221 expect(retryState?.nextAt ?? 0).toBeGreaterThan(Date.now());
222 });
223 } finally {
224 nowSpy.mockRestore();
225 }
226 });
227 
228 it("marks a run failed after repeated dispatch enqueue failures", async () => {
229 const baseTime = 1_720_000_000_000;
230 const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => baseTime);
231 
232 try {
233 const user = await seedUser({
234 email: "dispatch@example.com",
235 slug: "dispatch-user",
236 });
237 const project = await seedProject(user, {
238 projectSlug: "dispatch-project",
239 dispatchMode: "queue",
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 
248 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
249 const baseContext = createTestProjectDoContext(instance);
250 const context: ProjectDoContext = {
251 ...baseContext,
252 env: Object.assign(Object.create(baseContext.env), {
253 RUN_QUEUE: {
254 send: async () => {
255 throw new Error("queue unavailable");
256 },
257 },
258 }) as Env,
259 };
260 await baseContext.ctx.storage.deleteAlarm();
261 
262 await reconcileAcceptedRunD1Sync(context, project.id);
263 
264 for (let attempt = 0; attempt <= DISPATCH_RETRY_DELAYS_MS.length; attempt += 1) {
265 await dispatchExecutableRun(context, project.id);
266 if (attempt < DISPATCH_RETRY_DELAYS_MS.length) {
267 nowSpy.mockImplementation(
268 () =>
269 baseTime + DISPATCH_RETRY_DELAYS_MS.slice(0, attempt + 1).reduce((sum, delay) => sum + delay + 1, 0),
270 );
271 }
272 }
273 
274 await reconcileTerminalRunD1Sync(context, project.id);
275 });
276 
277 const rows = await readProjectDoRows(project.id);
278 expect(rows.runs[0]?.status).toBe("failed");
279 expect(rows.runs[0]?.dispatchStatus).toBe("terminal");
280 expect(rows.runs[0]?.lastError).toBe("dispatch_failed");
281 
282 const runMeta = await env.RUN_DO.getByName(accepted.runId).getRunSummary(accepted.runId);
283 expect(runMeta?.status).toBe("failed");
284 expect(runMeta?.errorMessage).toBe("dispatch_failed");
285 expect(runMeta?.startedAt).toBeNull();
286 
287 const db = createD1Db(env.DB);
288 const d1Row = await db.select().from(d1Schema.runIndex).where(eq(d1Schema.runIndex.id, accepted.runId)).limit(1);
289 expect(d1Row[0]?.status).toBe("failed");
290 expect(d1Row[0]?.startedAt).toBeNull();
291 } finally {
292 nowSpy.mockRestore();
293 }
294 });
295 
296 it("starts a Workflows instance when the run dispatch mode is workflows", async () => {
297 const user = await seedUser({
298 email: "workflow-dispatch@example.com",
299 slug: "workflow-dispatch-user",
300 });
301 const project = await seedProject(user, {
302 projectSlug: "workflow-dispatch-project",
303 dispatchMode: "workflows",
304 });
305 const projectStub = env.PROJECT_DO.getByName(project.id);
306 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
307 projectId: project.id,
308 triggeredByUserId: user.id,
309 branch: project.defaultBranch,
310 });
311 
312 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
313 const baseContext = createTestProjectDoContext(instance);
314 const createBatch = vi.fn(async () => [{ id: accepted.runId }]);
315 const context: ProjectDoContext = {
316 ...baseContext,
317 env: Object.assign(Object.create(baseContext.env), {
318 RUN_WORKFLOWS: {
319 createBatch,
320 get: vi.fn(),
321 },
322 }) as Env,
323 };
324 await baseContext.ctx.storage.deleteAlarm();
325 
326 await dispatchExecutableRun(context, project.id);
327 
328 expect(createBatch).toHaveBeenCalledWith([
329 expect.objectContaining({
330 id: accepted.runId,
331 params: expect.objectContaining({
332 runId: accepted.runId,
333 projectId: project.id,
334 dispatchMode: "workflows",
335 executionRuntime: DEFAULT_EXECUTION_RUNTIME,
336 }),
337 }),
338 ]);
339 });
340 
341 const rows = await readProjectDoRows(project.id);
342 expect(rows.runs[0]?.dispatchStatus).toBe("queued");
343 });
344 
345 it.each([
346 {
347 workflowStatus: "errored" as const,
348 expectedError: "reached terminal status errored",
349 },
350 {
351 workflowStatus: "terminated" as const,
352 expectedError: "reached terminal status terminated",
353 },
354 ])(
355 "rearms a queued workflow dispatch when the retained instance is $workflowStatus",
356 async ({ workflowStatus, expectedError }) => {
357 const user = await seedUser({
358 email: `workflow-${workflowStatus}-queued@example.com`,
359 slug: `workflow-${workflowStatus}-queued-user`,
360 });
361 const project = await seedProject(user, {
362 projectSlug: `workflow-${workflowStatus}-queued-project`,
363 dispatchMode: "workflows",
364 });
365 const projectStub = env.PROJECT_DO.getByName(project.id);
366 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
367 projectId: project.id,
368 triggeredByUserId: user.id,
369 branch: project.defaultBranch,
370 });
371 
372 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
373 const baseContext = createTestProjectDoContext(instance);
374 const status = vi.fn(async () => ({
375 status: workflowStatus,
376 }));
377 const context: ProjectDoContext = {
378 ...baseContext,
379 env: Object.assign(Object.create(baseContext.env), {
380 RUN_WORKFLOWS: {
381 createBatch: vi.fn(async () => [{ id: accepted.runId }]),
382 get: vi.fn(async () => ({
383 status,
384 restart: vi.fn(async () => {}),
385 })),
386 },
387 }) as Env,
388 };
389 await baseContext.ctx.storage.deleteAlarm();
390 
391 await dispatchExecutableRun(context, project.id);
392 await expect(dispatchExecutableRun(context, project.id)).resolves.toBeNull();
393 
394 const retryAt = await getDispatchRetryAt(baseContext, accepted.runId);
395 expect(retryAt ?? 0).toBeGreaterThan(Date.now());
396 });
397 
398 const rows = await readProjectDoRows(project.id);
399 expect(rows.runs[0]?.status).toBe("executable");
400 expect(rows.runs[0]?.dispatchStatus).toBe("pending");
401 expect(rows.runs[0]?.dispatchAttempts).toBe(1);
402 expect(rows.runs[0]?.lastError).toContain(expectedError);
403 },
404 );
405 
406 it("rearms a queued workflow dispatch when workflow status inspection fails", async () => {
407 const user = await seedUser({
408 email: "workflow-queued-status-throw@example.com",
409 slug: "workflow-queued-status-throw-user",
410 });
411 const project = await seedProject(user, {
412 projectSlug: "workflow-queued-status-throw-project",
413 dispatchMode: "workflows",
414 });
415 const projectStub = env.PROJECT_DO.getByName(project.id);
416 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
417 projectId: project.id,
418 triggeredByUserId: user.id,
419 branch: project.defaultBranch,
420 });
421 
422 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
423 const baseContext = createTestProjectDoContext(instance);
424 const context: ProjectDoContext = {
425 ...baseContext,
426 env: Object.assign(Object.create(baseContext.env), {
427 RUN_WORKFLOWS: {
428 createBatch: vi.fn(async () => [{ id: accepted.runId }]),
429 get: vi.fn(async () => ({
430 status: vi.fn(async () => {
431 throw new Error("workflow status unavailable");
432 }),
433 restart: vi.fn(async () => {}),
434 })),
435 },
436 }) as Env,
437 };
438 await baseContext.ctx.storage.deleteAlarm();
439 
440 await dispatchExecutableRun(context, project.id);
441 await expect(dispatchExecutableRun(context, project.id)).resolves.toBeNull();
442 
443 const retryAt = await getDispatchRetryAt(baseContext, accepted.runId);
444 expect(retryAt ?? 0).toBeGreaterThan(Date.now());
445 });
446 
447 const rows = await readProjectDoRows(project.id);
448 expect(rows.runs[0]?.status).toBe("executable");
449 expect(rows.runs[0]?.dispatchStatus).toBe("pending");
450 expect(rows.runs[0]?.dispatchAttempts).toBe(1);
451 expect(rows.runs[0]?.lastError).toContain("workflow status unavailable");
452 });
453 
454 it("preserves a workflow dispatch retry when the workflow rearms itself before dispatch returns", async () => {
455 const baseTime = 1_730_500_000_000;
456 const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => baseTime);
457 
458 try {
459 const user = await seedUser({
460 email: "workflow-fast-failure@example.com",
461 slug: "workflow-fast-failure-user",
462 });
463 const project = await seedProject(user, {
464 projectSlug: "workflow-fast-failure-project",
465 dispatchMode: "workflows",
466 });
467 const projectStub = env.PROJECT_DO.getByName(project.id);
468 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
469 projectId: project.id,
470 triggeredByUserId: user.id,
471 branch: project.defaultBranch,
472 });
473 
474 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
475 const baseContext = createTestProjectDoContext(instance);
476 let context: ProjectDoContext;
477 context = {
478 ...baseContext,
479 env: Object.assign(Object.create(baseContext.env), {
480 RUN_WORKFLOWS: {
481 createBatch: vi.fn(async () => {
482 const recovery = await recoverWorkflowDispatchFailure(baseContext, {
483 projectId: project.id,
484 runId: accepted.runId,
485 errorMessage: "claim exploded",
486 });
487 expect(recovery).toEqual({
488 kind: "rearmed",
489 });
490 return [{ id: accepted.runId }];
491 }),
492 get: vi.fn(),
493 },
494 }) as Env,
495 };
496 await baseContext.ctx.storage.deleteAlarm();
497 
498 await dispatchExecutableRun(context, project.id);
499 
500 const retryAt = await getDispatchRetryAt(baseContext, accepted.runId);
501 expect(retryAt ?? 0).toBeGreaterThan(Date.now());
502 });
503 
504 const rows = await readProjectDoRows(project.id);
505 expect(rows.runs[0]?.status).toBe("executable");
506 expect(rows.runs[0]?.dispatchStatus).toBe("pending");
507 expect(rows.runs[0]?.dispatchAttempts).toBe(1);
508 expect(rows.runs[0]?.lastError).toBe("claim exploded");
509 } finally {
510 nowSpy.mockRestore();
511 }
512 });
513 
514 it("restarts an existing Workflows instance after a rearmed pre-active failure", async () => {
515 const baseTime = 1_730_000_000_000;
516 const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => baseTime);
517 try {
518 const user = await seedUser({
519 email: "workflow-restart@example.com",
520 slug: "workflow-restart-user",
521 });
522 const project = await seedProject(user, {
523 projectSlug: "workflow-restart-project",
524 dispatchMode: "workflows",
525 });
526 const projectStub = env.PROJECT_DO.getByName(project.id);
527 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
528 projectId: project.id,
529 triggeredByUserId: user.id,
530 branch: project.defaultBranch,
531 });
532 
533 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
534 const baseContext = createTestProjectDoContext(instance);
535 const restart = vi.fn(async () => {});
536 const context: ProjectDoContext = {
537 ...baseContext,
538 env: Object.assign(Object.create(baseContext.env), {
539 RUN_WORKFLOWS: {
540 createBatch: vi.fn(async () => [{ id: accepted.runId }]),
541 get: vi.fn(async () => ({
542 restart,
543 })),
544 },
545 }) as Env,
546 };
547 await baseContext.ctx.storage.deleteAlarm();
548 
549 await dispatchExecutableRun(context, project.id);
550 const recovery = await recoverWorkflowDispatchFailure(context, {
551 projectId: project.id,
552 runId: accepted.runId,
553 errorMessage: "claim exploded",
554 });
555 expect(recovery).toEqual({
556 kind: "rearmed",
557 });
558 nowSpy.mockImplementation(() => baseTime + DISPATCH_RETRY_DELAYS_MS[0] + 1);
559 
560 await dispatchExecutableRun(
561 {
562 ...context,
563 env: Object.assign(Object.create(context.env), {
564 RUN_WORKFLOWS: {
565 createBatch: vi.fn(async () => []),
566 get: vi.fn(async () => ({
567 status: vi.fn(async () => ({
568 status: "complete" as const,
569 })),
570 restart,
571 })),
572 },
573 }) as Env,
574 },
575 project.id,
576 );
577 
578 expect(restart).toHaveBeenCalledTimes(1);
579 });
580 } finally {
581 nowSpy.mockRestore();
582 }
583 });
584 
585 it("does not restart an already-running Workflows instance when dispatch is replayed", async () => {
586 const user = await seedUser({
587 email: "workflow-running@example.com",
588 slug: "workflow-running-user",
589 });
590 const project = await seedProject(user, {
591 projectSlug: "workflow-running-project",
592 dispatchMode: "workflows",
593 });
594 const projectStub = env.PROJECT_DO.getByName(project.id);
595 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
596 projectId: project.id,
597 triggeredByUserId: user.id,
598 branch: project.defaultBranch,
599 });
600 
601 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
602 const baseContext = createTestProjectDoContext(instance);
603 const restart = vi.fn(async () => {});
604 const status = vi.fn(async () => ({
605 status: "running" as const,
606 }));
607 const context: ProjectDoContext = {
608 ...baseContext,
609 env: Object.assign(Object.create(baseContext.env), {
610 RUN_WORKFLOWS: {
611 createBatch: vi.fn(async () => []),
612 get: vi.fn(async () => ({
613 status,
614 restart,
615 })),
616 },
617 }) as Env,
618 };
619 await baseContext.ctx.storage.deleteAlarm();
620 
621 await dispatchExecutableRun(context, project.id);
622 
623 expect(status).toHaveBeenCalledTimes(1);
624 expect(restart).not.toHaveBeenCalled();
625 });
626 
627 const rows = await readProjectDoRows(project.id);
628 expect(rows.runs[0]?.dispatchStatus).toBe("queued");
629 });
630 
631 it.each(["paused", "waitingForPause"] as const)(
632 "does not restart a retained Workflows instance when dispatch is replayed in %s state",
633 async (workflowStatus) => {
634 const user = await seedUser({
635 email: `workflow-${workflowStatus.toLowerCase()}@example.com`,
636 slug: `workflow-${workflowStatus.toLowerCase()}-user`,
637 });
638 const project = await seedProject(user, {
639 projectSlug: `workflow-${workflowStatus.toLowerCase()}-project`,
640 dispatchMode: "workflows",
641 });
642 const projectStub = env.PROJECT_DO.getByName(project.id);
643 const accepted = await acceptManualRunWithoutAlarm(projectStub, {
644 projectId: project.id,
645 triggeredByUserId: user.id,
646 branch: project.defaultBranch,
647 });
648 
649 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
650 const baseContext = createTestProjectDoContext(instance);
651 const restart = vi.fn(async () => {});
652 const status = vi.fn(async () => ({
653 status: workflowStatus,
654 }));
655 const context: ProjectDoContext = {
656 ...baseContext,
657 env: Object.assign(Object.create(baseContext.env), {
658 RUN_WORKFLOWS: {
659 createBatch: vi.fn(async () => []),
660 get: vi.fn(async () => ({
661 status,
662 restart,
663 })),
664 },
665 }) as Env,
666 };
667 await baseContext.ctx.storage.deleteAlarm();
668 
669 await dispatchExecutableRun(context, project.id);
670 
671 expect(status).toHaveBeenCalledTimes(1);
672 expect(restart).not.toHaveBeenCalled();
673 });
674 
675 const rows = await readProjectDoRows(project.id);
676 expect(rows.runs[0]?.dispatchStatus).toBe("queued");
677 },
678 );
679 
680 it("keeps workflow dispatch pending when an existing instance has an unsupported status", async () => {
681 const user = await seedUser({
682 email: "workflow-unknown@example.com",
683 slug: "workflow-unknown-user",
684 });
685 const project = await seedProject(user, {
686 projectSlug: "workflow-unknown-project",
687 dispatchMode: "workflows",
688 });
689 const projectStub = env.PROJECT_DO.getByName(project.id);
690 await acceptManualRunWithoutAlarm(projectStub, {
691 projectId: project.id,
692 triggeredByUserId: user.id,
693 branch: project.defaultBranch,
694 });
695 
696 await runInDurableObject(projectStub, async (instance: ProjectDO) => {
697 const baseContext = createTestProjectDoContext(instance);
698 const context: ProjectDoContext = {
699 ...baseContext,
700 env: Object.assign(Object.create(baseContext.env), {
701 RUN_WORKFLOWS: {
702 createBatch: vi.fn(async () => []),
703 get: vi.fn(async () => ({
704 status: vi.fn(async () => ({
705 status: "unknown" as const,
706 })),
707 restart: vi.fn(async () => {}),
708 })),
709 },
710 }) as Env,
711 };
712 await baseContext.ctx.storage.deleteAlarm();
713 
714 await expect(dispatchExecutableRun(context, project.id)).resolves.toBeNull();
715 });
716 
717 const rows = await readProjectDoRows(project.id);
718 expect(rows.runs[0]?.dispatchStatus).toBe("pending");
719 expect(rows.runs[0]?.dispatchAttempts).toBe(1);
720 expect(rows.runs[0]?.lastError).toContain("unsupported status unknown");
721 });
722});