test(ci): await heartbeat completion and match candidate authority checks (#161455)

Co-authored-by: Sarah Fortune <sarah.fortune@gmail.com>
This commit is contained in:
Sarah Fortune 2026-09-29 18:06:57 -07:00 • committed by GitHub
parent 188e031c25
commit 5f76cc437d
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 53 additions and 42 deletions

View file

@ -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<ReturnType<typeof requestHeartbeatAndWait>> = [];
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);
}
}
});
});

View file

@ -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);

View file

@ -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({