mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
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:
parent
188e031c25
commit
5f76cc437d
3 changed files with 53 additions and 42 deletions
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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({
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue