openclaw / src /process /command-queue.deadline.test.ts
SaylorTwift's picture
SaylorTwift HF Staff
Add files using upload-large-folder tool
3144483 verified
Raw
History Blame Contribute Delete
5.56 kB
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { createDeferred } from "../../test/helpers/promise.js";
import { enqueueCommandInLane, getCommandLaneSnapshot } from "./command-queue.js";
import { resetCommandQueueStateForTest } from "./command-queue.test-support.js";
import type { CommandQueueTaskDeadline } from "./command-queue.types.js";
const lane = "runtime-deadline-test";
const finishers: Array<() => void> = [];
function enqueueOwnedTask(initialDeadline: CommandQueueTaskDeadline, initiallyAborted = false) {
const finish = createDeferred();
finishers.push(finish.resolve);
const abort = new AbortController();
const release = new AbortController();
if (initiallyAborted) {
abort.abort();
}
const unsubscribe = vi.fn();
let progressAtMs = Date.now();
let publish: (deadline: CommandQueueTaskDeadline | undefined) => void = () => {
throw new Error("deadline subscription is not active");
};
const task = enqueueCommandInLane(lane, () => finish.promise, {
taskTimeoutMs: 25,
taskTimeoutProgressAtMs: () => progressAtMs,
taskTimeoutAbortSignal: abort.signal,
taskTimeoutAbortGraceMs: 5,
taskTimeoutReleaseSignal: release.signal,
taskTimeoutSubscribe: (onDeadline) => {
publish = onDeadline;
onDeadline(initialDeadline);
return unsubscribe;
},
});
const outcome = task.then(
() => ({ status: "completed" as const }),
(error: unknown) => ({ status: "failed" as const, error }),
);
return {
abort,
release,
finish: finish.resolve,
outcome,
unsubscribe,
publish: (deadline: CommandQueueTaskDeadline | undefined) => publish(deadline),
progress: () => {
progressAtMs = Date.now();
},
};
}
beforeEach(() => {
resetCommandQueueStateForTest();
vi.useFakeTimers();
vi.setSystemTime(Date.parse("2026-08-20T12:00:00Z"));
});
afterEach(async () => {
for (const finish of finishers.splice(0)) {
finish();
}
await vi.advanceTimersByTimeAsync(0);
resetCommandQueueStateForTest();
vi.useRealTimers();
});
describe("command lane owner deadlines", () => {
it("replaces idle timing with an absolute deadline that progress cannot extend", async () => {
const owner = enqueueOwnedTask({ kind: "bounded", deadlineAtMs: Date.now() + 100 });
const next = vi.fn(async () => "next");
const queued = enqueueCommandInLane(lane, next);
await vi.advanceTimersByTimeAsync(50);
owner.progress();
await vi.advanceTimersByTimeAsync(49);
expect(next).not.toHaveBeenCalled();
expect(getCommandLaneSnapshot(lane)).toMatchObject({ activeCount: 1, queuedCount: 1 });
await vi.advanceTimersByTimeAsync(1);
await expect(owner.outcome).resolves.toMatchObject({
status: "failed",
error: {
name: "CommandLaneTaskTimeoutError",
message: expect.stringContaining("owner deadline"),
},
});
await expect(queued).resolves.toBe("next");
expect(owner.unsubscribe).toHaveBeenCalledOnce();
});
it("lets unlimited execution hand off to bounded terminal settlement", async () => {
const owner = enqueueOwnedTask({ kind: "unlimited" });
await vi.advanceTimersByTimeAsync(49 * 60 * 60 * 1000);
expect(getCommandLaneSnapshot(lane).activeCount).toBe(1);
owner.publish({ kind: "bounded", deadlineAtMs: Date.now() + 120_000 });
await vi.advanceTimersByTimeAsync(119_999);
expect(getCommandLaneSnapshot(lane).activeCount).toBe(1);
owner.finish();
await expect(owner.outcome).resolves.toEqual({ status: "completed" });
owner.publish({ kind: "bounded", deadlineAtMs: Date.now() });
await vi.advanceTimersByTimeAsync(1);
expect(getCommandLaneSnapshot(lane).activeCount).toBe(0);
expect(vi.getTimerCount()).toBe(0);
});
it("restores ordinary idle recovery after the runtime releases deadline ownership", async () => {
const owner = enqueueOwnedTask({ kind: "unlimited" });
await vi.advanceTimersByTimeAsync(1_000);
owner.progress();
owner.publish(undefined);
await vi.advanceTimersByTimeAsync(24);
expect(getCommandLaneSnapshot(lane).activeCount).toBe(1);
await vi.advanceTimersByTimeAsync(1);
await expect(owner.outcome).resolves.toMatchObject({
status: "failed",
error: {
name: "CommandLaneTaskTimeoutError",
message: expect.stringContaining("no progress for 25ms"),
},
});
});
it("does not let a late deadline update replace an accepted abort grace", async () => {
const owner = enqueueOwnedTask({ kind: "unlimited" });
owner.abort.abort();
owner.publish({ kind: "unlimited" });
await vi.advanceTimersByTimeAsync(4);
expect(getCommandLaneSnapshot(lane).activeCount).toBe(1);
await vi.advanceTimersByTimeAsync(1);
await expect(owner.outcome).resolves.toMatchObject({
status: "failed",
error: {
name: "CommandLaneTaskTimeoutError",
message: expect.stringContaining("abort grace 5ms"),
},
});
});
it.each([false, true])(
"honors immediate release when initially aborted=%s",
async (initiallyAborted) => {
const owner = enqueueOwnedTask({ kind: "unlimited" }, initiallyAborted);
owner.release.abort();
await expect(owner.outcome).resolves.toMatchObject({
status: "failed",
error: {
name: "CommandLaneTaskTimeoutError",
message: expect.stringContaining("lane release requested"),
},
});
expect(getCommandLaneSnapshot(lane).activeCount).toBe(0);
},
);
});