diff --git a/src/cron/service.failure-alert-notifications.test.ts b/src/cron/service.failure-alert-notifications.test.ts index d06b212c14f5..24115a23931c 100644 --- a/src/cron/service.failure-alert-notifications.test.ts +++ b/src/cron/service.failure-alert-notifications.test.ts @@ -7,7 +7,7 @@ import { resolveHeartbeatRunPrompt, } from "../infra/heartbeat-runner-prompt.js"; import { startHeartbeatRunner } from "../infra/heartbeat-runner-scheduler.js"; -import { requestHeartbeat as requestHeartbeatWake } from "../infra/heartbeat-wake.js"; +import { requestHeartbeatAndWait } from "../infra/heartbeat-wake.js"; import { drainSystemEvents, enqueueSystemEvent, @@ -76,6 +76,7 @@ describe("CronService failure notification delivery", () => { threadId: 77, }; const observed: Array<{ prompt: string; deliveryContext?: DeliveryContext }> = []; + const pendingWakes: Array> = []; const runOnce = vi.fn(async (options: HeartbeatRunOptions) => { const pendingEventEntries = peekSystemEventEntries(options.sessionKey ?? ""); const preflight = await resolveHeartbeatPreflight({ @@ -125,13 +126,17 @@ describe("CronService failure notification delivery", () => { contextKey: options?.contextKey, deliveryContext: options?.deliveryContext, }), - requestHeartbeat: (wake) => - requestHeartbeatWake({ - ...wake, - sessionKey: - wake.sessionKey ?? resolveAgentMainSessionKey({ cfg, agentId: wake.agentId ?? "main" }), - coalesceMs: 0, - }), + requestHeartbeat: (wake) => { + pendingWakes.push( + requestHeartbeatAndWait({ + ...wake, + sessionKey: + wake.sessionKey ?? + resolveAgentMainSessionKey({ cfg, agentId: wake.agentId ?? "main" }), + coalesceMs: 0, + }), + ); + }, sendCronFailureAlert, runIsolatedAgentJob: async () => ({ status: "error", @@ -156,6 +161,7 @@ describe("CronService failure notification delivery", () => { "failure alert channel unavailable", ); await vi.advanceTimersByTimeAsync(1); + await Promise.all(pendingWakes); expect(peekSystemEventEntries(testCase.sessionKey)).toHaveLength(1); expect(runOnce).toHaveBeenCalledTimes(testCase.wakesNow ? 1 : 0); @@ -182,8 +188,14 @@ describe("CronService failure notification delivery", () => { } } finally { cron.stop(); - runner.stop(); - drainSystemEvents(testCase.sessionKey); + try { + // Stopping the runner retains unfinished notifications for its successor. + await vi.advanceTimersByTimeAsync(1); + await Promise.all(pendingWakes); + } finally { + runner.stop(); + drainSystemEvents(testCase.sessionKey); + } } }); }); diff --git a/src/cron/service.heartbeat-busy-poll.test.ts b/src/cron/service.heartbeat-busy-poll.test.ts index 8fedf7035f38..d3a76ea3fccc 100644 --- a/src/cron/service.heartbeat-busy-poll.test.ts +++ b/src/cron/service.heartbeat-busy-poll.test.ts @@ -289,9 +289,12 @@ describe("native heartbeat busy poll settlement", () => { await vi.advanceTimersByTimeAsync(firstTick - Date.now()); await waitForRequest(1); expect(request).toHaveBeenCalledOnce(); + await vi.advanceTimersByTimeAsync(250); + // Scratch preflight uses a real worker, which fake-clock advancement cannot join. + await runOnce.mock.results[0]?.value; // Observe the full original watchdog window on both versions. The // unfixed scheduler records a timeout; the fixed poll settled promptly. - await vi.advanceTimersByTimeAsync(600_001); + await vi.advanceTimersByTimeAsync(600_001 - 250); await waitForFinished(1); expect(finished()).toHaveLength(1); // Finished precedes schedule maintenance and release of the active marker. @@ -423,6 +426,7 @@ describe("native heartbeat busy poll settlement", () => { monitor, sessionKey, reply, + runOnce, request, waitForRequest, finished, @@ -443,6 +447,7 @@ describe("native heartbeat busy poll settlement", () => { await waitForRequest(1); expect(request).toHaveBeenCalledOnce(); await vi.advanceTimersByTimeAsync(250); + await runOnce.mock.results[0]?.value; expect(finished()).toHaveLength(0); expect(peekSystemEventEntries(sessionKey).map((entry) => entry.text)).toContain(text); expect(reply).not.toHaveBeenCalled(); @@ -544,12 +549,14 @@ describe("native heartbeat busy poll settlement", () => { await waitForRequest(1); expect(request).toHaveBeenCalledOnce(); await vi.advanceTimersByTimeAsync(250); + await runOnce.mock.results[0]?.value; expect(deps.isReplyRunActive).toHaveBeenCalledTimes(2); expect(finished()).toHaveLength(0); expect(reply).not.toHaveBeenCalled(); const releaseMain = await holdLane(CommandLane.Main); await vi.advanceTimersByTimeAsync(60_000); expect(runOnce).toHaveBeenCalledTimes(2); + await runOnce.mock.results[1]?.value; expect(finished()).toHaveLength(0); await releaseMain(); await vi.advanceTimersByTimeAsync(60_000); diff --git a/src/infra/update-candidate-canary-progress.test-support.ts b/src/infra/update-candidate-canary-progress.test-support.ts index 52ea7f8e5785..968db3bbcbe0 100644 --- a/src/infra/update-candidate-canary-progress.test-support.ts +++ b/src/infra/update-candidate-canary-progress.test-support.ts @@ -18,6 +18,7 @@ import { SqliteWorkerError } from "./sqlite-worker-contract.js"; import { validateUpdateCandidateCanary } from "./update-candidate-canary.js"; import { stubHealthyGateway } from "./update-candidate-canary.test-support.js"; import * as candidateIo from "./update-candidate-io.js"; +import { UpdateRequesterRevokedError } from "./update-requester-authority.js"; import { createUpdateRun, getUpdateRunAsync } from "./update-run-ledger.js"; import { updateRunStepsFromResultStep } from "./update-run-step.js"; @@ -129,6 +130,7 @@ export function registerCanaryProgressWorkerTests( } } }; + const onStepComplete = vi.fn(); const pending = Promise.resolve().then(() => validateUpdateCandidateWithProgress( { @@ -138,7 +140,7 @@ export function registerCanaryProgressWorkerTests( assertCurrent: guards.assertCurrent, writeOptions, }, - { opts: { json: true }, progress: {} }, + { opts: { json: true }, progress: { onStepComplete } }, run, ), ); @@ -160,48 +162,38 @@ export function registerCanaryProgressWorkerTests( expect(mocks.spawn).not.toHaveBeenCalled(); return; } - if (outcome === "interrupted-after-acceptance") { - const observed = await pending.then( - (result) => ({ result }), - (error: unknown) => ({ error }), - ); + if (outcome === "revoked-at-commit" || outcome === "interrupted-after-acceptance") { + await expect(pending).rejects.toBeInstanceOf(UpdateRequesterRevokedError); expect(checkedCommit).toBe(true); const saved = await getUpdateRunAsync(run.runId, { env }); - expect(saved?.steps).toContainEqual( + if (outcome === "revoked-at-commit") { + expect(saved).toEqual(created); + } else { + expect(saved?.steps).toContainEqual( + expect.objectContaining({ + step: "candidate-state-snapshot", + status: "in_progress", + detail: expect.stringContaining("completed, attempt 1, 920445/920445 pages"), + }), + ); + } + expect(onStepComplete).toHaveBeenCalledWith( expect.objectContaining({ - step: "candidate-state-snapshot", - status: "in_progress", - detail: expect.stringContaining("completed, attempt 1, 920445/920445 pages"), - }), - ); - expect(observed).toMatchObject({ - result: { - status: "error", - phase: "snapshot", - steps: expect.arrayContaining([ + name: "candidate-state-snapshot", + exitCode: 1, + failureFacts: expect.arrayContaining([ expect.objectContaining({ - exitCode: 1, - failureFacts: expect.arrayContaining([ - expect.objectContaining({ - message: expect.stringContaining("requester-revoked"), - }), - ]), + message: expect.stringContaining("requester-revoked"), }), ]), - }, - }); + }), + ); expect(mocks.spawn).not.toHaveBeenCalled(); return; } const result = await pending; expect(checkedCommit).toBe(true); const saved = await getUpdateRunAsync(run.runId, { env }); - if (outcome === "revoked-at-commit") { - expect(result).toMatchObject({ status: "error", phase: "snapshot" }); - expect(saved).toEqual(created); - expect(mocks.spawn).not.toHaveBeenCalled(); - return; - } expect(result.status).toBe("ok"); expect(saved?.steps).toContainEqual( expect.objectContaining({