test(state,shared,integration): remove low-value tests (batch d150) (#163514)

* test(state): deslop s978 tests

* test(shared): deslop s962 tests

* test(trajectory): deslop s998 tests

* test(state): deslop s983 tests

* test(status): deslop s986 tests

* test(transcripts): deslop s999 tests

* test(integration): deslop s1005 tests

* test(integration): deslop s1008 tests

* test(core): deslop s1003 tests

* test(scripts): deslop s1014 tests

* test(state): await mapped read failures

Handle synchronous throws and asynchronous rejections while retaining the exact mapped-error identity assertion. Fixes the type-aware no-floating-promises CI failure from the batch replay.
This commit is contained in:
Peter Steinberger 2026-10-02 08:00:00 -05:00 • committed by GitHub
parent 6e63f2812a
commit 889f54b266
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
43 changed files with 889 additions and 1901 deletions

View file

@ -5,32 +5,6 @@ import { AsyncWorkScope, getAsyncWorkSignal, trackAsyncWork } from "./async-work
import { createDeferredCore } from "./deferred.js";
describe("async work resources", () => {
it("returns the default logical result while its owner still joins cleanup", async () => {
const owner = new AsyncWorkScope();
const cleanupStarted = createDeferredCore();
const finishCleanup = createDeferredCore();
const release = vi.fn(async () => {
cleanupStarted.resolve();
await finishCleanup.promise;
});
const result = owner.run(() =>
runWithAsyncWorkResources(async (onAcquired) => {
onAcquired({ release });
return "accepted";
}),
);
try {
expect(await result).toBe("accepted");
await cleanupStarted.promise;
expect(owner.hasPendingWork).toBe(true);
expect(release).toHaveBeenCalledOnce();
} finally {
finishCleanup.resolve();
await owner.drain();
}
expect(release).toHaveBeenCalledOnce();
});
it("releases opted-in idle resources exactly once before returning the result", async () => {
const owner = new AsyncWorkScope();
const cleanupStarted = createDeferredCore();
@ -64,28 +38,6 @@ describe("async work resources", () => {
expect(release).toHaveBeenCalledOnce();
});
it("returns an opted-in result while accepted work still owns its resources", async () => {
const owner = new AsyncWorkScope();
const finishAdmission = createDeferredCore();
const release = vi.fn();
const result = owner.run(() =>
runWithAsyncWorkResources(async (onAcquired) => {
onAcquired({ release, releaseBeforeResultWhenIdle: true });
void trackAsyncWork(() => finishAdmission.promise);
return "enqueued";
}),
);
try {
expect(await result).toBe("enqueued");
expect(owner.hasPendingWork).toBe(true);
expect(release).not.toHaveBeenCalled();
} finally {
finishAdmission.resolve();
await owner.drain();
}
expect(release).toHaveBeenCalledOnce();
});
it("rejects opted-in idle cleanup failures without attempting release twice", async () => {
const owner = new AsyncWorkScope();
const failure = new Error("resource cleanup failed");
@ -105,8 +57,10 @@ describe("async work resources", () => {
it("does not replace a default logical result when later cleanup fails", async () => {
const owner = new AsyncWorkScope();
const cleanupStarted = createDeferredCore();
const finishCleanup = createDeferredCore();
const release = vi.fn(async () => {
cleanupStarted.resolve();
await finishCleanup.promise;
throw new Error("resource cleanup failed");
});
@ -118,6 +72,9 @@ describe("async work resources", () => {
);
try {
expect(await result).toBe("accepted");
await cleanupStarted.promise;
expect(owner.hasPendingWork).toBe(true);
expect(release).toHaveBeenCalledOnce();
} finally {
finishCleanup.resolve();
await owner.drain();
@ -141,6 +98,8 @@ describe("async work resources", () => {
);
try {
expect(await result).toBe("enqueued");
expect(owner.hasPendingWork).toBe(true);
expect(release).not.toHaveBeenCalled();
expect(operationSignal?.aborted).toBe(false);
const reason = new Error("requester cancelled");
owner.beginClose(reason);

View file

@ -15,98 +15,73 @@ afterEach(() => {
vi.clearAllMocks();
});
const activeAuthority = () => {};
describe("channel read completion ownership", () => {
it("keeps nested output until the outer read accepts it without retaining the inner callable", async () => {
const settle = vi.fn(async (_accepted: boolean) => {});
let innerAssertion: (() => void) | undefined;
await withChannelReadAuthority(
() => {},
async () => {
await withChannelReadAuthority(
() => {},
async () => {
innerAssertion = captureChannelReadAuthority();
captureChannelReadScope()!.registerResource({ key: "created-media", settle });
},
);
expect(() => innerAssertion!()).toThrow("no longer active");
expect(settle).not.toHaveBeenCalled();
expect(() => captureChannelReadAuthority()!()).not.toThrow();
},
);
await withChannelReadAuthority(activeAuthority, async () => {
await withChannelReadAuthority(activeAuthority, async () => {
innerAssertion = captureChannelReadAuthority();
captureChannelReadScope()!.registerResource({ key: "created-media", settle });
});
expect(() => innerAssertion!()).toThrow("no longer active");
expect(settle).not.toHaveBeenCalled();
expect(() => captureChannelReadAuthority()!()).not.toThrow();
});
expect(settle).toHaveBeenCalledExactlyOnceWith(true);
});
it("retains the creating provider's live check after its inner scope returns", async () => {
const settle = vi.fn(async (_accepted: boolean) => {});
let providerActive = true;
await expect(
withChannelReadAuthority(
() => {},
async () => {
it.each(["provider", "source signal"])(
"retains child %s revocation until the outer read accepts output",
async (authority) => {
const settle = vi.fn(async (_accepted: boolean) => {});
const source = new AbortController();
const revoked = new Error("child read revoked");
let providerActive = true;
await expect(
withChannelReadAuthority(activeAuthority, async () => {
await withChannelReadAuthority(
() => {
if (!providerActive) {
throw new Error("provider revoked");
throw revoked;
}
},
async () => {
captureChannelReadScope()!.registerResource({ key: "created-media", settle });
},
authority === "source signal" ? source.signal : undefined,
);
providerActive = false;
},
),
).rejects.toThrow("provider revoked");
expect(settle).toHaveBeenCalledExactlyOnceWith(false);
});
if (authority === "provider") {
providerActive = false;
} else {
source.abort(revoked);
}
}),
).rejects.toBe(revoked);
expect(settle).toHaveBeenCalledExactlyOnceWith(false);
},
);
it("discards a failed inner read even when its parent catches the error", async () => {
const outer = vi.fn(async (_accepted: boolean) => {});
const inner = vi.fn(async (_accepted: boolean) => {});
await withChannelReadAuthority(
() => {},
async () => {
captureChannelReadScope()!.registerResource({ key: "outer-media", settle: outer });
await expect(
withChannelReadAuthority(
() => {},
async () => {
captureChannelReadScope()!.registerResource({ key: "inner-media", settle: inner });
throw new Error("download failed");
},
),
).rejects.toThrow("download failed");
expect(inner).toHaveBeenCalledExactlyOnceWith(false);
expect(outer).not.toHaveBeenCalled();
},
);
await withChannelReadAuthority(activeAuthority, async () => {
captureChannelReadScope()!.registerResource({ key: "outer-media", settle: outer });
await expect(
withChannelReadAuthority(activeAuthority, async () => {
captureChannelReadScope()!.registerResource({ key: "inner-media", settle: inner });
throw new Error("download failed");
}),
).rejects.toThrow("download failed");
expect(inner).toHaveBeenCalledExactlyOnceWith(false);
expect(outer).not.toHaveBeenCalled();
});
expect(outer).toHaveBeenCalledExactlyOnceWith(true);
expect(inner).toHaveBeenCalledOnce();
});
it("retains a child's independent source signal until the enclosing read accepts its output", async () => {
const source = new AbortController();
const aborted = new Error("child source canceled");
const settle = vi.fn(async (_accepted: boolean) => {});
await expect(
withChannelReadAuthority(
() => {},
async () => {
await withChannelReadAuthority(
() => {},
async () => {
captureChannelReadScope()!.registerResource({ key: "created-media", settle });
},
source.signal,
);
source.abort(aborted);
},
),
).rejects.toBe(aborted);
expect(settle).toHaveBeenCalledExactlyOnceWith(false);
});
it("publishes checked acceptance before asynchronous resource teardown", async () => {
const closing = createDeferred();
const release = createDeferred();
@ -168,16 +143,13 @@ describe("channel read completion ownership", () => {
const otherChunk = await import("./channel-read-authority.js");
const settle = vi.fn(async (_accepted: boolean) => {});
let retained: (() => void) | undefined;
await withChannelReadAuthority(
() => {},
async () => {
retained = otherChunk.captureChannelReadAuthority();
expect(typeof retained).toBe("function");
expect(retained).toBe(captureChannelReadAuthority());
retained!();
otherChunk.captureChannelReadScope()!.registerResource({ key: "created-media", settle });
},
);
await withChannelReadAuthority(activeAuthority, async () => {
retained = otherChunk.captureChannelReadAuthority();
expect(typeof retained).toBe("function");
expect(retained).toBe(captureChannelReadAuthority());
retained!();
otherChunk.captureChannelReadScope()!.registerResource({ key: "created-media", settle });
});
expect(settle).toHaveBeenCalledExactlyOnceWith(true);
expect(() => retained!()).toThrow("no longer active");
});

View file

@ -41,10 +41,9 @@ function fixture() {
return 0;
},
);
const library = { func: vi.fn(() => sysctl) };
const load = vi.fn(() => library);
const load = vi.fn(() => ({ func: () => sysctl }));
loadNative.mockReturnValue({ load });
return { boot, proc, parts, sysctl, library, load };
return { proc, parts, sysctl, load };
}
async function read(pid = 42) {
@ -56,12 +55,11 @@ async function read(pid = 42) {
it.each(["x64", "arm64"] as const)("reads exact kernel microseconds on %s", async (arch) => {
vi.spyOn(process, "arch", "get").mockReturnValue(arch);
const { sysctl, load, library } = fixture();
const { sysctl, load } = fixture();
expect(await read()).toBe(5_200_002);
expect(await read()).toBe(5_200_002);
expect(loadNative).toHaveBeenCalledTimes(1);
expect(load).toHaveBeenCalledExactlyOnceWith(null);
expect(library.func).toHaveBeenCalledTimes(1);
expect(
sysctl.mock.calls.slice(0, 3).map(([mib, count, output]) => [mib, count, output.length]),
).toEqual([
@ -85,53 +83,26 @@ it("keeps the identity stable when the boot offset and wall-clock start move tog
expect(await read()).toBe(5_200_002);
});
it.each([
"changed boot bracket",
"wrong struct size",
"unknown layout",
"wrong PID",
"negative seconds",
"negative microseconds",
"microsecond overflow",
"negative start",
"unsafe start",
])("refuses %s without an identity approximation", async (variant) => {
const { proc, parts } = fixture();
if (variant === "changed boot bracket") {
parts[2]!.writeBigInt64LE(700_002n, 8);
}
if (variant === "wrong struct size") {
proc.writeInt32LE(1080);
}
if (variant === "unknown layout") {
proc.writeInt32LE(1, 4);
}
if (variant === "wrong PID") {
proc.writeInt32LE(43, 72);
}
if (variant === "negative seconds") {
proc.writeBigInt64LE(-1n, 336);
}
if (variant === "negative microseconds") {
proc.writeBigInt64LE(-1n, 344);
}
if (variant === "microsecond overflow") {
proc.writeBigInt64LE(1_000_000n, 344);
}
if (variant === "negative start") {
proc.writeBigInt64LE(0n, 336);
}
if (variant === "unsafe start") {
proc.writeBigInt64LE(BigInt(Number.MAX_SAFE_INTEGER), 336);
}
it.each<[string, (data: ReturnType<typeof fixture>) => void]>([
["changed boot bracket", ({ parts }) => parts[2]!.writeBigInt64LE(700_002n, 8)],
["wrong struct size", ({ proc }) => proc.writeInt32LE(1080)],
["unknown layout", ({ proc }) => proc.writeInt32LE(1, 4)],
["wrong PID", ({ proc }) => proc.writeInt32LE(43, 72)],
["negative seconds", ({ proc }) => proc.writeBigInt64LE(-1n, 336)],
["negative microseconds", ({ proc }) => proc.writeBigInt64LE(-1n, 344)],
["microsecond overflow", ({ proc }) => proc.writeBigInt64LE(1_000_000n, 344)],
["negative start", ({ proc }) => proc.writeBigInt64LE(0n, 336)],
["unsafe start", ({ proc }) => proc.writeBigInt64LE(BigInt(Number.MAX_SAFE_INTEGER), 336)],
])("refuses %s without an identity approximation", async (_name, corrupt) => {
corrupt(fixture());
expect(await read()).toBeNull();
});
it.each(
[0, 1, 2].flatMap((failedCall) =>
["short", "long", "status"].map((failure) => ({ failedCall, failure })),
),
)("refuses $failure output from sysctl call $failedCall", async ({ failedCall, failure }) => {
it.each([
{ failedCall: 0, failure: "short" },
{ failedCall: 1, failure: "long" },
{ failedCall: 2, failure: "status" },
])("refuses $failure output from sysctl call $failedCall", async ({ failedCall, failure }) => {
const { sysctl } = fixture();
const successful = sysctl.getMockImplementation()!;
let call = 0;
@ -171,14 +142,11 @@ it("retries failed native loading and fails closed on native errors", async () =
expect(await read()).toBeNull();
});
it.each([0, -1, 1.5, Number.NaN, Number.POSITIVE_INFINITY, 0x80000000])(
"rejects invalid native PID %s before loading",
async (pid) => {
fixture();
expect(await read(pid)).toBeNull();
expect(loadNative).not.toHaveBeenCalled();
},
);
it.each([0, 1.5, 0x80000000])("rejects invalid native PID %s before loading", async (pid) => {
fixture();
expect(await read(pid)).toBeNull();
expect(loadNative).not.toHaveBeenCalled();
});
it("does not load for other operating systems, architectures or byte orders", async () => {
fixture();

View file

@ -17,6 +17,14 @@ let bytes: Buffer;
const query = vi.fn();
const load = vi.fn();
function mockShellIdentity() {
return vi
.spyOn(childProcess, "execFileSync")
.mockImplementation((_file, args) =>
args?.[1] === "lstart=" ? "Thu Sep 24 00:00:00 2026\n" : "42 7 Thu Sep 24 00:00:00 2026\n",
);
}
beforeEach(() => {
vi.resetModules();
vi.spyOn(process, "platform", "get").mockReturnValue("darwin");
@ -108,49 +116,30 @@ it("keeps the bounded shell path on x64 without loading Koffi", async () => {
expect(nativeKoffi).not.toHaveBeenCalled();
});
it.each([
"short read",
"wrong PID",
"invalid parent",
"zero start",
"unsafe start",
"invalid microseconds",
"query error",
"missing native package",
])("uses the bounded shell fallback after %s", async (failure) => {
if (failure === "short read") {
query.mockReturnValue(135);
}
if (failure === "wrong PID") {
bytes.writeUInt32LE(43, 12);
}
if (failure === "invalid parent") {
bytes.writeUInt32LE(0xffffffff, 16);
}
if (failure === "zero start") {
bytes.writeBigUInt64LE(0n, 120);
}
if (failure === "unsafe start") {
bytes.writeBigUInt64LE(2n ** 53n, 120);
}
if (failure === "invalid microseconds") {
bytes.writeBigUInt64LE(1_000_000n, 128);
}
if (failure === "query error") {
query.mockImplementation(() => {
throw new Error("denied");
});
}
if (failure === "missing native package") {
nativeKoffi.mockImplementation(() => {
throw new Error("unavailable");
});
}
const shell = vi
.spyOn(childProcess, "execFileSync")
.mockImplementation((_file, args) =>
args?.[1] === "lstart=" ? "Thu Sep 24 00:00:00 2026\n" : "42 7 Thu Sep 24 00:00:00 2026\n",
);
it.each<[string, () => void]>([
["short read", () => query.mockReturnValue(135)],
["wrong PID", () => bytes.writeUInt32LE(43, 12)],
["invalid parent", () => bytes.writeUInt32LE(0xffffffff, 16)],
["zero start", () => bytes.writeBigUInt64LE(0n, 120)],
["unsafe start", () => bytes.writeBigUInt64LE(2n ** 53n, 120)],
["invalid microseconds", () => bytes.writeBigUInt64LE(1_000_000n, 128)],
[
"query error",
() =>
query.mockImplementation(() => {
throw new Error("denied");
}),
],
[
"missing native package",
() =>
nativeKoffi.mockImplementation(() => {
throw new Error("unavailable");
}),
],
])("uses the bounded shell fallback after %s", async (_name, failNative) => {
failNative();
const shell = mockShellIdentity();
const { getFileLockProcessStartTime, getProcessInstanceStartTime, readDarwinProcessIdentity } =
await import("./pid-alive.js");
const expected = Date.UTC(2026, 8, 24) / 1000;
@ -203,11 +192,7 @@ it("preserves default shell recovery after slow native loading fails", async ()
now += 1500;
throw new Error("native unavailable");
});
const shell = vi
.spyOn(childProcess, "execFileSync")
.mockImplementation((_file, args) =>
args?.[1] === "lstart=" ? "Thu Sep 24 00:00:00 2026\n" : "42 7 Thu Sep 24 00:00:00 2026\n",
);
const shell = mockShellIdentity();
const { getFileLockProcessStartTime, readDarwinProcessIdentity } = await import("./pid-alive.js");
const expected = Date.UTC(2026, 8, 24) / 1000;
expect(getFileLockProcessStartTime(42)).toBe(expected);
@ -225,11 +210,7 @@ it.each([600, 1000])("charges %sms native loading to an explicit deadline", asyn
now += elapsed;
throw new Error("native unavailable");
});
const shell = vi
.spyOn(childProcess, "execFileSync")
.mockImplementation((_file, args) =>
args?.[1] === "lstart=" ? "Thu Sep 24 00:00:00 2026\n" : "42 7 Thu Sep 24 00:00:00 2026\n",
);
const shell = mockShellIdentity();
const { getFileLockProcessStartTime, readDarwinProcessIdentity } = await import("./pid-alive.js");
expect(getFileLockProcessStartTime(42, process.env, 1000)).toBe(
elapsed === 1000 ? null : Date.UTC(2026, 8, 24) / 1000,
@ -247,15 +228,11 @@ it.each([600, 1000])("charges %sms native loading to an explicit deadline", asyn
}
});
it.each([0, -1, 1.5, Number.NaN, Infinity])(
"rejects invalid PID %s before native conversion",
async (pid) => {
const shell = vi.spyOn(childProcess, "execFileSync");
const { getFileLockProcessStartTime, readDarwinProcessIdentity } =
await import("./pid-alive.js");
expect(getFileLockProcessStartTime(pid)).toBeNull();
expect(readDarwinProcessIdentity(pid)).toBeNull();
expect(nativeKoffi).not.toHaveBeenCalled();
expect(shell).not.toHaveBeenCalled();
},
);
it.each([0, 1.5])("rejects invalid PID %s before native conversion", async (pid) => {
const shell = vi.spyOn(childProcess, "execFileSync");
const { getFileLockProcessStartTime, readDarwinProcessIdentity } = await import("./pid-alive.js");
expect(getFileLockProcessStartTime(pid)).toBeNull();
expect(readDarwinProcessIdentity(pid)).toBeNull();
expect(nativeKoffi).not.toHaveBeenCalled();
expect(shell).not.toHaveBeenCalled();
});

View file

@ -4,25 +4,14 @@ import { isCronSessionDisplayKey, isSystemCreatedSessionRow } from "./session-li
describe("session display visibility", () => {
it.each([
["cron:", true],
[" CRON:nightly ", true],
["cron", false],
["agent:main:cron:nightly", true],
[" AGENT:MAIN:CRON:nightly ", true],
["agent::main::cron::nightly", true],
["agent: :cron:nightly", true],
["agent:\n:cron:nightly", true],
["agent:main:cron: :child", true],
["agent:main:cron:nightly:child", true],
["agent:main:cron:", false],
["agent:main:cron: ", false],
["agent:main::cron::", false],
["agent::cron:nightly", false],
["agent:main:other:cron:nightly", false],
["agent:main:cronicle:nightly", false],
["agent:main:main", false],
[":agent:main:cron:nightly", false],
["prefix:cron:nightly", false],
["", false],
] as const)("classifies %j as automation=%s", (key, expected) => {
expect(isCronSessionDisplayKey(key)).toBe(expected);
});

View file

@ -3,7 +3,7 @@ import { sortAndLimitBy, sortAndLimitByWork } from "./sort-and-limit.js";
import { runSynchronousWork } from "./synchronous-work.js";
describe("sortAndLimitBy", () => {
it.each([1, 5, 50, 100, 200, 201, 1000, 1024, 1025, 2050, 4096, undefined])(
it.each([1, 200, 201, 1025, 4096, undefined])(
"matches stable full ordering without mutating input for limit %s",
(limit) => {
const entries = Array.from({ length: 2051 }, (_, id) => ({ id, rank: (id * 37) % 23 }));

View file

@ -3,7 +3,6 @@ import path from "node:path";
import { afterEach, beforeEach, expect, it, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import * as sqliteRuntime from "../infra/bun-sqlite-library.js";
import { createRetainedOperation, type RetainedOperation } from "../infra/retained-operation.js";
import { createOwnedWorkerTaskPoolMock } from "../infra/worker-task-pool.mock.test-support.js";
import type { OwnedWorkerTask, RetainedWorkerTask } from "../infra/worker-task-pool.types.js";
import { PluginBlobStoreError } from "../plugin-state/plugin-blob-store.types.js";
@ -13,6 +12,7 @@ import {
closeOpenClawStateDatabaseByPathAsync,
} from "./openclaw-state-db-cache.js";
import { executeExistingOpenClawStateRead } from "./openclaw-state-db-readonly.js";
import { observeAsyncFixture } from "./openclaw-state-db-readonly.test-support.js";
import { withExistingOpenClawStateSchema } from "./openclaw-state-db-schema-policy.js";
import type {
OpenClawStateReadPhase,
@ -22,17 +22,6 @@ import type {
import { captureOpenClawStateReadWorkerContext } from "./openclaw-state-worker-context.js";
import { encodeOpenClawStateWorkerError } from "./openclaw-state-worker-error.js";
// These awaited fixtures observe Promise settlement; they do not prove blocked-host progress.
function observeAsyncFixture<T>(run: () => Promise<T>): RetainedOperation<T> {
const completion = createRetainedOperation<T>(() => undefined);
try {
void run().then(completion.resolve, completion.reject);
} catch (error) {
completion.reject(error);
}
return completion.operation;
}
const mock = vi.hoisted(() => ({
run: vi.fn<() => Promise<OpenClawStateReadReply>>(),
close: vi.fn<OwnedWorkerTask<OpenClawStateReadReply>["close"]>(),
@ -128,15 +117,17 @@ it.each(["retired", "different-source"] as const)(
await closeOpenClawStateDatabaseByPathAsync(options.path);
}
const { mapped, mapError } = mapper();
await expect(
Promise.resolve().then(() =>
executeExistingOpenClawStateRead(
kind === "different-source" ? source() : options,
{ type: "fleet.list" },
{ context, mapError },
),
),
).rejects.toBe(mapped);
let failure: unknown;
try {
await executeExistingOpenClawStateRead(
kind === "different-source" ? source() : options,
{ type: "fleet.list" },
{ context, mapError },
);
} catch (error) {
failure = error;
}
expect(failure).toBe(mapped);
expect(mapError).toHaveBeenCalledExactlyOnceWith(
expect.objectContaining(
kind === "retired"
@ -165,29 +156,20 @@ it("maps synchronous read admission refusal once before read work", () => {
expect(mock.close).not.toHaveBeenCalled();
});
it.each(["read admission", "schema scope"] as const)(
"maps captured %s retirement once before read work",
async (kind) => {
const options = source();
const context =
kind === "schema scope"
? withExistingOpenClawStateSchema({ path: options.path }, () =>
captureOpenClawStateReadWorkerContext(options),
)
: captureOpenClawStateReadWorkerContext(options);
if (kind === "read admission") {
await closeOpenClawStateDatabaseByPathAsync(options.path);
}
const { mapped, mapError } = mapper();
expect(() =>
executeExistingOpenClawStateRead(options, { type: "fleet.list" }, { context, mapError }),
).toThrow(mapped);
expect(mapError).toHaveBeenCalledOnce();
expect(mapError.mock.calls[0]?.[1]).toBe("before-read");
expect(mock.run).not.toHaveBeenCalled();
expect(mock.close).not.toHaveBeenCalled();
},
);
it("maps captured schema scope retirement once before read work", () => {
const options = source();
const context = withExistingOpenClawStateSchema({ path: options.path }, () =>
captureOpenClawStateReadWorkerContext(options),
);
const { mapped, mapError } = mapper();
expect(() =>
executeExistingOpenClawStateRead(options, { type: "fleet.list" }, { context, mapError }),
).toThrow(mapped);
expect(mapError).toHaveBeenCalledOnce();
expect(mapError.mock.calls[0]?.[1]).toBe("before-read");
expect(mock.run).not.toHaveBeenCalled();
expect(mock.close).not.toHaveBeenCalled();
});
it("maps an authoritative pre-read error after cleanup", async () => {
const original = new Error("source admission failed");

View file

@ -33,7 +33,6 @@ beforeEach(() => {
});
it.each([
{ stage: "read", schemaAdmission: false },
{ stage: "schema", schemaAdmission: false },
{ stage: "read", schemaAdmission: true },
{ stage: "schema", schemaAdmission: true },

View file

@ -3,7 +3,6 @@ import fs from "node:fs";
import { join } from "node:path";
import { afterAll, afterEach, beforeAll, beforeEach, expect, it, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import { createRetainedOperation, type RetainedOperation } from "../infra/retained-operation.js";
import {
adoptPreparedLocation,
cleanupSnapshotOperations,
@ -19,6 +18,10 @@ import {
withArtifactPreservingStateReads,
withOpenClawStateDatabaseReadSnapshot,
} from "./openclaw-state-db-readonly.js";
import {
observeAsyncFixture,
retainFixturePreparation,
} from "./openclaw-state-db-readonly.test-support.js";
import type {
OpenClawStateReadAuthority,
OpenClawStateReadLocation,
@ -26,31 +29,6 @@ import type {
} from "./openclaw-state-read.types.js";
import type { OpenClawStateWorkerContext } from "./openclaw-state-worker-context.types.js";
// These awaited fixtures observe Promise settlement; they do not prove blocked-host progress.
function observeAsyncFixture<T>(run: () => Promise<T>): RetainedOperation<T> {
const completion = createRetainedOperation<T>(() => undefined);
try {
void run().then(completion.resolve, completion.reject);
} catch (error) {
completion.reject(error);
}
return completion.operation;
}
// This fixture's preparation owns no resource beyond the separately cleaned prepared location.
function retainFixturePreparation<T>(preparation: RetainedOperation<T>) {
const close = createRetainedOperation<void>(() => {
if (preparation.read().status !== "pending") {
close.resolve(undefined);
}
});
void preparation.result.then(
() => close.operation.service(),
() => close.operation.service(),
);
return { ...preparation, startClose: () => close.operation };
}
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
const mocks = vi.hoisted(() => ({
@ -182,27 +160,9 @@ vi.mock("./openclaw-state-db-read-connection.js", () => ({
vi.mock("./openclaw-state-db-schema-version.js", () => ({
assertSupportedStateSchemaVersion: mocks.forbidden,
}));
vi.mock("./openclaw-state-read-worker.js", () => {
const progress = new Set<() => void>();
return {
captureOpenClawStateReadSource: () => ({
createTransport: () => ({
startRead: (source: OpenClawStateReadLocation, authority: OpenClawStateReadAuthority) =>
observeAsyncFixture(() => mocks.read(source, authority)),
startValidateFresh: () => observeAsyncFixture(async () => {}),
startClose: () => observeAsyncFixture(mocks.close),
}),
own(service: () => void) {
progress.add(service);
return () => progress.delete(service);
},
service() {
for (const service of Array.from(progress)) {
service();
}
},
}),
};
vi.mock("./openclaw-state-read-worker.js", async () => {
const { createReadWorkerFixture } = await import("./openclaw-state-db-readonly.test-support.js");
return createReadWorkerFixture(mocks.read, mocks.close);
});
let exitCleanup: (() => void) | undefined;

View file

@ -2,27 +2,16 @@ import { AsyncLocalStorage } from "node:async_hooks";
import fs from "node:fs";
import path from "node:path";
import { beforeEach, expect, it, vi } from "vitest";
import { createRetainedOperation, type RetainedOperation } from "../infra/retained-operation.js";
import { getAsyncWorkSignal } from "../shared/async-work-scope.js";
import { createDeferredCore } from "../shared/deferred.js";
import { withTempDir } from "../test-utils/temp-dir.js";
import { createReadWorkerFixture } from "./openclaw-state-db-readonly.test-support.js";
import type {
OpenClawStateReadAuthority,
OpenClawStateReadLocation,
OpenClawStateReadOutcome,
} from "./openclaw-state-read.types.js";
// These awaited fixtures observe Promise settlement; they do not prove blocked-host progress.
function observeAsyncFixture<T>(run: () => Promise<T>): RetainedOperation<T> {
const completion = createRetainedOperation<T>(() => undefined);
try {
void run().then(completion.resolve, completion.reject);
} catch (error) {
completion.reject(error);
}
return completion.operation;
}
const mocks = vi.hoisted(() => ({
capture: vi.fn(),
assertCurrent: vi.fn<() => void>(),
@ -71,28 +60,9 @@ vi.mock("./openclaw-state-db-read-connection.js", () => ({
withOpenClawStateReadOnlyLocation: mocks.forbiddenNative,
}));
vi.mock("./openclaw-state-read-worker.js", () => {
const progress = new Set<() => void>();
return {
captureOpenClawStateReadSource: () => ({
createTransport: () => ({
startRead: (source: OpenClawStateReadLocation, authority: OpenClawStateReadAuthority) =>
observeAsyncFixture(() => mocks.read(source, authority)),
startValidateFresh: () => observeAsyncFixture(async () => {}),
startClose: () => observeAsyncFixture(async () => {}),
}),
own(service: () => void) {
progress.add(service);
return () => progress.delete(service);
},
service() {
for (const service of Array.from(progress)) {
service();
}
},
}),
};
});
vi.mock("./openclaw-state-read-worker.js", () =>
createReadWorkerFixture(mocks.read, async () => {}),
);
import { isStateDatabaseReadAdmissionInvalidatedError } from "./openclaw-state-db-async-lifecycle.js";
import {
@ -209,50 +179,29 @@ const rejectedAdmissions = {
worker: true,
};
it.each(["closing", "cleanup-failed", "closed"] as const)(
"rejects escaped snapshot reads while %s instead of reopening a live source",
async (phase) => {
await withTempDir("openclaw-retired-snapshot-", async (root) => {
const source = path.join(root, "source");
fs.writeFileSync(source, "mock source; never opened as SQLite");
const enteredCleanup = createDeferredCore();
const finishCleanup = createDeferredCore<boolean>();
let escape!: ReturnType<typeof AsyncLocalStorage.snapshot>;
mocks.cleanup.mockImplementationOnce(async () => {
enteredCleanup.resolve();
return phase === "closing" ? await finishCleanup.promise : phase === "closed";
});
const closing = withOpenClawStateDatabaseReadSnapshot(
it("rejects escaped snapshot reads after failed cleanup instead of reopening a live source", async () => {
await withTempDir("openclaw-retired-snapshot-", async (root) => {
const source = path.join(root, "source");
fs.writeFileSync(source, "mock source; never opened as SQLite");
let escape!: ReturnType<typeof AsyncLocalStorage.snapshot>;
mocks.cleanup.mockResolvedValueOnce(false);
await expect(
withOpenClawStateDatabaseReadSnapshot(
async () => {
escape = AsyncLocalStorage.snapshot();
expect(getActiveOpenClawStateDatabaseReadSnapshot({ path: source })).toBeDefined();
},
{ path: source },
).then(
() => undefined,
(error: unknown) => error,
);
await enteredCleanup.promise;
if (phase !== "closing") {
await closing;
}
try {
const observed = await escape(() => probeRetiredAdmission(source));
expect(observed).toEqual(rejectedAdmissions);
expect(
escape(() =>
getActiveOpenClawStateDatabaseReadSnapshot({ path: path.join(root, "other") }),
),
).toBeUndefined();
expect(mocks.forbiddenNative).not.toHaveBeenCalled();
expect(mocks.read).not.toHaveBeenCalled();
} finally {
finishCleanup.resolve(true);
await closing;
}
});
},
);
),
).rejects.toThrow("snapshot cleanup failed");
expect(await escape(() => probeRetiredAdmission(source))).toEqual(rejectedAdmissions);
expect(
escape(() => getActiveOpenClawStateDatabaseReadSnapshot({ path: path.join(root, "other") })),
).toBeUndefined();
expect(mocks.forbiddenNative).not.toHaveBeenCalled();
expect(mocks.read).not.toHaveBeenCalled();
});
});
it.each(["snapshot", "disposable"] as const)(
"drains an admitted %s reader while rejecting escaped new reads, then rejects the closed scope",

View file

@ -0,0 +1,60 @@
import { createRetainedOperation, type RetainedOperation } from "../infra/retained-operation.js";
import type {
OpenClawStateReadAuthority,
OpenClawStateReadLocation,
OpenClawStateReadOutcome,
} from "./openclaw-state-read.types.js";
// These awaited fixtures observe Promise settlement; they do not prove blocked-host progress.
export function observeAsyncFixture<T>(run: () => Promise<T>): RetainedOperation<T> {
const completion = createRetainedOperation<T>(() => undefined);
try {
void run().then(completion.resolve, completion.reject);
} catch (error) {
completion.reject(error);
}
return completion.operation;
}
// Preparation owns no resource beyond the separately cleaned prepared location.
export function retainFixturePreparation<T>(preparation: RetainedOperation<T>) {
const close = createRetainedOperation<void>(() => {
if (preparation.read().status !== "pending") {
close.resolve(undefined);
}
});
void preparation.result.then(
() => close.operation.service(),
() => close.operation.service(),
);
return { ...preparation, startClose: () => close.operation };
}
export function createReadWorkerFixture(
read: (
source: OpenClawStateReadLocation,
authority: OpenClawStateReadAuthority,
) => Promise<OpenClawStateReadOutcome>,
close: () => Promise<void>,
) {
const progress = new Set<() => void>();
return {
captureOpenClawStateReadSource: () => ({
createTransport: () => ({
startRead: (source: OpenClawStateReadLocation, authority: OpenClawStateReadAuthority) =>
observeAsyncFixture(() => read(source, authority)),
startValidateFresh: () => observeAsyncFixture(async () => {}),
startClose: () => observeAsyncFixture(close),
}),
own(service: () => void) {
progress.add(service);
return () => progress.delete(service);
},
service() {
for (const service of Array.from(progress)) {
service();
}
},
}),
};
}

View file

@ -2,9 +2,13 @@ import fs from "node:fs";
import path from "node:path";
import { afterEach, beforeEach, expect, it, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import { createRetainedOperation, type RetainedOperation } from "../infra/retained-operation.js";
import type { AsyncPreparedSqliteReadOnlyLocation } from "../infra/sqlite-readonly-location.types.js";
import { createDeferredCore } from "../shared/deferred.js";
import {
createReadWorkerFixture,
observeAsyncFixture,
retainFixturePreparation,
} from "./openclaw-state-db-readonly.test-support.js";
import type {
OpenClawStateReadAuthority,
OpenClawStateReadLocation,
@ -16,31 +20,6 @@ vi.hoisted(() => {
vi.resetModules();
});
// These awaited fixtures observe Promise settlement; they do not prove blocked-host progress.
function observeAsyncFixture<T>(run: () => Promise<T>): RetainedOperation<T> {
const completion = createRetainedOperation<T>(() => undefined);
try {
void run().then(completion.resolve, completion.reject);
} catch (error) {
completion.reject(error);
}
return completion.operation;
}
// This fixture's preparation owns no resource beyond the separately cleaned prepared location.
function retainFixturePreparation<T>(preparation: RetainedOperation<T>) {
const close = createRetainedOperation<void>(() => {
if (preparation.read().status !== "pending") {
close.resolve(undefined);
}
});
void preparation.result.then(
() => close.operation.service(),
() => close.operation.service(),
);
return { ...preparation, startClose: () => close.operation };
}
const mock = vi.hoisted(() => ({
close: vi.fn<() => Promise<void>>(),
read: vi.fn<
@ -67,28 +46,7 @@ vi.mock("../infra/sqlite-readonly-location.js", async (importOriginal) => ({
prepareSqliteReadOnlyLocationFromOwnedDatabase: mock.prepareNative,
}));
let finishProducer: (() => void) | undefined;
vi.mock("./openclaw-state-read-worker.js", () => {
const progress = new Set<() => void>();
return {
captureOpenClawStateReadSource: () => ({
createTransport: () => ({
startRead: (location: OpenClawStateReadLocation, authority: OpenClawStateReadAuthority) =>
observeAsyncFixture(() => mock.read(location, authority)),
startValidateFresh: () => observeAsyncFixture(async () => {}),
startClose: () => observeAsyncFixture(mock.close),
}),
own(service: () => void) {
progress.add(service);
return () => progress.delete(service);
},
service() {
for (const service of Array.from(progress)) {
service();
}
},
}),
};
});
vi.mock("./openclaw-state-read-worker.js", () => createReadWorkerFixture(mock.read, mock.close));
vi.mock("../infra/sqlite-snapshot-source.js", async (importOriginal) => ({
...(await importOriginal<typeof import("../infra/sqlite-snapshot-source.js")>()),
prepareSqliteReadOnlyLocation: mock.prepareSource,
@ -196,30 +154,6 @@ it("reads independently when native snapshot borrowing refuses a transaction", a
expect(mock.prepareSourceAsync).not.toHaveBeenCalled();
});
it("prepares artifact reads from the retained native source with its cleanup owner", async () => {
const options = source();
const database = { db: {} };
const observe = vi.fn();
const release = vi.fn();
mock.borrow.mockReturnValue({ database, assertCurrent() {}, observe, release });
await expect(
withArtifactPreservingStateReads(() =>
executeExistingOpenClawStateRead(options, { type: "fleet.list" }),
),
).resolves.toEqual({ ok: true, type: "fleet.list", sourceAdmitted: true, cells: [] });
expect(mock.prepareNative).toHaveBeenCalledWith(
database.db,
expect.any(Function),
expect.any(AbortSignal),
"async",
);
expect(mock.prepareSource).not.toHaveBeenCalled();
expect(mock.independent).not.toHaveBeenCalled();
expect(observe).toHaveBeenCalledOnce();
expect(release).toHaveBeenCalledOnce();
expect(mock.cleanup).toHaveBeenCalledOnce();
});
it("retains the borrowed source through pending preparation and failed published cleanup", async () => {
const options = source();
const started = createDeferredCore();
@ -257,131 +191,33 @@ it("retains the borrowed source through pending preparation and failed published
expect(release).toHaveBeenCalledOnce();
});
it("joins the cold snapshot query transport before producer cleanup", async () => {
it("retains ordered cleanup for canonical retry after a read failure", async () => {
const options = source();
const { identity } = captureOpenClawStateDatabaseReadAdmission(options.path);
const stopping = createDeferredCore();
const stopped = createDeferredCore();
const events: string[] = [];
mock.read.mockImplementation(async () => {
events.push("query completed");
return { value: { ok: true, type: "fleet.list", sourceAdmitted: true, cells: [] } };
});
mock.close.mockImplementationOnce(async () => {
events.push("transport closing");
stopping.resolve();
await stopped.promise;
events.push("transport stopped");
});
mock.cleanup.mockImplementation(async () => {
events.push("producer snapshot cleaned");
return true;
});
const result = withArtifactPreservingStateReads(() =>
executeExistingOpenClawStateRead(options, { type: "fleet.list" }),
);
try {
await stopping.promise;
expect(mock.prepareSourceAsync).toHaveBeenCalledExactlyOnceWith(options.path, {
expectedSourceIdentity: { key: identity.key },
preserveSourceArtifacts: true,
signal: expect.any(AbortSignal),
});
expect(mock.prepareSource).not.toHaveBeenCalled();
expect(mock.prepareNative).not.toHaveBeenCalled();
expect(mock.independent).not.toHaveBeenCalled();
expect(mock.read).toHaveBeenCalledExactlyOnceWith(
expect.objectContaining({
location: "/fixture/prepared.sqlite",
snapshotRoot: "/fixture/prepared",
checkFreshAdmission: true,
expectedIdentity: undefined,
}),
expect.objectContaining({
assertCurrent: expect.any(Function),
signal: expect.any(AbortSignal),
}),
);
expect(events).toEqual(["query completed", "transport closing"]);
expect(mock.cleanup).not.toHaveBeenCalled();
} finally {
stopped.resolve();
}
await expect(result).resolves.toEqual({
ok: true,
type: "fleet.list",
sourceAdmitted: true,
cells: [],
});
expect(events).toEqual([
"query completed",
"transport closing",
"transport stopped",
"producer snapshot cleaned",
]);
expect(mock.close).toHaveBeenCalledOnce();
expect(mock.cleanup).toHaveBeenCalledOnce();
});
const readFailure = new Error("read failed");
const stopFailure = new Error("transport stop not acknowledged");
mock.read.mockRejectedValue(readFailure);
mock.close.mockRejectedValueOnce(stopFailure);
mock.cleanup.mockResolvedValueOnce(false).mockResolvedValue(true);
it("preserves native transaction refusal for artifact reads without independent fallback", async () => {
const options = source();
const failure = new Error(
"Asynchronous shared-state reads cannot run inside a native transaction",
);
mock.borrow.mockImplementation(() => {
throw failure;
});
mock.independent.mockReturnValue({ assertCurrent() {}, observe() {}, release() {} });
await expect(
withArtifactPreservingStateReads(() =>
executeExistingOpenClawStateRead(options, { type: "fleet.list" }),
),
).rejects.toBe(failure);
expect(mock.independent).not.toHaveBeenCalled();
expect(mock.prepareNative).not.toHaveBeenCalled();
expect(mock.prepareSource).not.toHaveBeenCalled();
expect(mock.prepareSourceAsync).not.toHaveBeenCalled();
expect(mock.read).not.toHaveBeenCalled();
).rejects.toMatchObject({ cause: readFailure, errors: [readFailure, stopFailure] });
expect(mock.cleanup).not.toHaveBeenCalled();
await expect(closeOpenClawStateDatabaseByPathAsync(options.path)).rejects.toThrow(
"snapshot cleanup failed",
);
expect(mock.close).toHaveBeenCalledTimes(2);
expect(() => captureOpenClawStateDatabaseReadAdmission(options.path)).toThrow(/closed/);
await expect(closeOpenClawStateDatabaseByPathAsync(options.path)).resolves.toBe(false);
expect(mock.close).toHaveBeenCalledTimes(2);
expect(mock.cleanup).toHaveBeenCalledTimes(2);
expect(() =>
captureOpenClawStateDatabaseReadAdmission(options.path).assertCurrent(),
).not.toThrow();
});
it.each([false, true])(
"retains ordered cleanup for canonical retry after read failure=%s",
async (readFails) => {
const options = source();
const readFailure = new Error("read failed");
const stopFailure = new Error("transport stop not acknowledged");
if (readFails) {
mock.read.mockRejectedValue(readFailure);
}
mock.close.mockRejectedValueOnce(stopFailure);
mock.cleanup.mockResolvedValueOnce(false).mockResolvedValue(true);
const result = withArtifactPreservingStateReads(() =>
executeExistingOpenClawStateRead(options, { type: "fleet.list" }),
);
if (readFails) {
await expect(result).rejects.toMatchObject({
cause: readFailure,
errors: [readFailure, stopFailure],
});
} else {
await expect(result).rejects.toBe(stopFailure);
}
expect(mock.cleanup).not.toHaveBeenCalled();
await expect(closeOpenClawStateDatabaseByPathAsync(options.path)).rejects.toThrow(
"snapshot cleanup failed",
);
expect(mock.close).toHaveBeenCalledTimes(2);
expect(() => captureOpenClawStateDatabaseReadAdmission(options.path)).toThrow(/closed/);
await expect(closeOpenClawStateDatabaseByPathAsync(options.path)).resolves.toBe(false);
expect(mock.close).toHaveBeenCalledTimes(2);
expect(mock.cleanup).toHaveBeenCalledTimes(2);
expect(() =>
captureOpenClawStateDatabaseReadAdmission(options.path).assertCurrent(),
).not.toThrow();
},
);
it("retains an outer snapshot until its child transport acknowledges cleanup", async () => {
const options = source();
const failure = new Error("child stop not acknowledged");

View file

@ -88,17 +88,6 @@ it("moves owned-native artifact token SQL off the caller while retaining its sou
const changedOffsets = [...shmBefore.keys()].filter(
(offset) => shmBefore[offset] !== shmAfter[offset],
);
const observation = {
tokenSql,
before,
beforeBackup,
afterBackup,
after,
changedOffsets,
readMarksBefore: shmBefore.subarray(100, 120).toString("hex"),
readMarksAfter: shmAfter.subarray(100, 120).toString("hex"),
};
console.info("owned-native-token observation", JSON.stringify(observation));
expect(beforeBackup).toEqual(before);
expect(after).toEqual(afterBackup);
expect(after.filter((entry) => entry.suffix !== "-shm")).toEqual(

View file

@ -56,7 +56,7 @@ beforeEach(() => {
});
});
it.each(["success", "query-error", "schema-error"] as const)(
it.each(["query-error", "schema-error"] as const)(
"reports source admission at the schema-validated callback for %s",
(outcome) => {
const failure = new Error("controlled reader failure");
@ -64,19 +64,17 @@ it.each(["success", "query-error", "schema-error"] as const)(
mock.admit.mockImplementation(() => {
throw failure;
});
} else if (outcome === "query-error") {
} else {
mock.query.mockImplementation(() => {
throw failure;
});
}
const reply = mock.handler(request);
if (reply.ok) {
expect(reply).toEqual({ ok: true, type: "fleet.list", sourceAdmitted: true, cells: [] });
} else {
expect(reply.message).toBe(failure.message);
expect(reply.sourceAdmitted).toBe(outcome === "query-error" ? true : undefined);
}
expect(reply.ok).toBe(outcome === "success");
expect(reply).toMatchObject({
ok: false,
message: failure.message,
sourceAdmitted: outcome === "query-error" ? true : undefined,
});
expect(mock.query).toHaveBeenCalledTimes(outcome === "schema-error" ? 0 : 1);
},
);

View file

@ -1,3 +1,4 @@
import assert from "node:assert/strict";
import { expect, it } from "vitest";
import { findSourceImportBackedges } from "../../test/helpers/source-import-closure.js";
import {
@ -7,43 +8,33 @@ import {
} from "./openclaw-state-worker-error.js";
import { SessionMetadataUnavailableError } from "./session-metadata-unavailable-error.js";
it.each([false, true])(
"preserves canonical metadata unavailability and SQLite causes with aggregate=%s",
(aggregate) => {
const cause = Object.assign(new Error("synthetic SQLite read failure"), {
code: "ERR_SQLITE_ERROR",
errcode: 1,
});
const unavailable = new SessionMetadataUnavailableError("table-missing", { cause }, [
"session_nodes",
]);
const cleanup = new Error("synthetic database close failure");
const payload = encodeOpenClawStateWorkerError(
aggregate
? new AggregateError([unavailable, cleanup], "read and close failed", { cause: cleanup })
: unavailable,
);
expect(payload).toBeDefined();
const carrier = new Error("remote error");
retainOpenClawStateWorkerErrorPayload(carrier, structuredClone(payload));
const decoded = hydrateOpenClawStateWorkerError(carrier);
const restored = decoded instanceof AggregateError ? decoded.errors[0] : decoded;
expect(restored).toBeInstanceOf(SessionMetadataUnavailableError);
expect(restored).toMatchObject({
reason: "table-missing",
missingTables: ["session_nodes"],
cause: { message: cause.message, code: "ERR_SQLITE_ERROR", errcode: 1 },
});
if (aggregate) {
expect(decoded).toBeInstanceOf(AggregateError);
if (!(decoded instanceof AggregateError)) {
throw new Error("Expected retained cleanup error");
}
expect(decoded.errors[1]).toBe(decoded.cause);
expect(decoded.cause).toMatchObject({ message: cleanup.message });
}
},
);
it("preserves canonical metadata unavailability and SQLite causes through cleanup", () => {
const cause = Object.assign(new Error("synthetic SQLite read failure"), {
code: "ERR_SQLITE_ERROR",
errcode: 1,
});
const unavailable = new SessionMetadataUnavailableError("table-missing", { cause }, [
"session_nodes",
]);
const cleanup = new Error("synthetic database close failure");
const payload = encodeOpenClawStateWorkerError(
new AggregateError([unavailable, cleanup], "read and close failed", { cause: cleanup }),
);
expect(payload).toBeDefined();
const carrier = new Error("remote error");
retainOpenClawStateWorkerErrorPayload(carrier, structuredClone(payload));
const decoded = hydrateOpenClawStateWorkerError(carrier);
assert(decoded instanceof AggregateError);
const restored = decoded.errors[0];
expect(restored).toBeInstanceOf(SessionMetadataUnavailableError);
expect(restored).toMatchObject({
reason: "table-missing",
missingTables: ["session_nodes"],
cause: { message: cause.message, code: "ERR_SQLITE_ERROR", errcode: 1 },
});
expect(decoded.errors[1]).toBe(decoded.cause);
expect(decoded.cause).toMatchObject({ message: cleanup.message });
});
it("keeps shared error transport independent of agent schema classification", () => {
expect(

View file

@ -29,10 +29,7 @@ import {
markOpenClawStateDatabaseFailure,
} from "./openclaw-state-db-failure.js";
import { OpenClawStateDatabaseSchemaMigrationRequiredError } from "./openclaw-state-db-schema-migration-required.js";
import {
OpenClawStateLeaseError,
toOpenClawStateLeaseVerificationError,
} from "./openclaw-state-lease-error.js";
import { toOpenClawStateLeaseVerificationError } from "./openclaw-state-lease-error.js";
import {
encodeOpenClawStateWorkerError,
hydrateOpenClawStateWorkerError,
@ -99,8 +96,6 @@ describe("shared-state worker error transport", () => {
);
it.each([
{ code: "ERR_SQLITE_ERROR", errcode: 5 },
{ code: "ERR_SQLITE_ERROR", errcode: 6 },
{ code: "ERR_SQLITE_ERROR", errcode: 517 },
{ code: "ERR_SQLITE_ERROR", errcode: 262 },
{ code: "SQLITE_BUSY" },
@ -111,66 +106,35 @@ describe("shared-state worker error transport", () => {
});
it.each(
[RangeError, SyntaxError, TypeError, SkillUploadRequestError].flatMap((ErrorType) =>
[false, true].map((aggregate) => ({ ErrorType, name: ErrorType.name, aggregate })),
),
)("preserves $name identity with aggregate=$aggregate", ({ ErrorType, aggregate }) => {
[RangeError, SyntaxError, TypeError, SkillUploadRequestError].map((ErrorType) => ({
ErrorType,
name: ErrorType.name,
})),
)("preserves $name identity and cause", ({ ErrorType }) => {
const original = Object.assign(new ErrorType("Synthetic invalid request"), {
code: "ERR_OUT_OF_RANGE",
cause: new Error("Synthetic decoding cause"),
});
const root = aggregate
? new AggregateError([original, original], "Read and cleanup", { cause: original })
: original;
const decoded = roundTrip(root);
const restored = aggregate ? decoded.cause : decoded;
expect(restored).toBeInstanceOf(ErrorType);
expect(restored).toMatchObject({
const decoded = roundTrip(original);
expect(decoded).toBeInstanceOf(ErrorType);
expect(decoded).toMatchObject({
name: original.name,
message: original.message,
code: original.code,
cause: { message: "Synthetic decoding cause" },
});
if (aggregate) {
assert(decoded instanceof AggregateError);
expect(decoded.errors).toHaveLength(2);
expect(decoded.errors[0]).toBe(restored);
expect(decoded.errors[1]).toBe(restored);
}
expect(hydrateOpenClawStateWorkerError(decoded)).toBe(decoded);
});
it.each([
["PLUGIN_BLOB_OPEN_FAILED", "open"],
["PLUGIN_BLOB_WRITE_FAILED", "register"],
["PLUGIN_BLOB_READ_FAILED", "lookup"],
["PLUGIN_BLOB_CORRUPT", "entries"],
["PLUGIN_BLOB_LIMIT_EXCEEDED", "register"],
["PLUGIN_BLOB_INVALID_INPUT", "sweep"],
] as const)("preserves Blob error identity for %s", (code, operation) => {
const error = new PluginBlobStoreError("Synthetic blob refusal", {
code,
operation,
path: "/fixture/blob.sqlite",
cause: Object.assign(new Error("Synthetic SQLite cause"), {
code: "ERR_SQLITE_ERROR",
errcode: 1,
}),
});
const decoded = roundTrip(error);
expect(decoded).toMatchObject({
code,
operation,
path: "/fixture/blob.sqlite",
cause: { message: "Synthetic SQLite cause", code: "ERR_SQLITE_ERROR", errcode: 1 },
});
});
it("retains Blob primary failure and shared references in cleanup aggregates", () => {
const primary = new PluginBlobStoreError("Synthetic read failure", {
code: "PLUGIN_BLOB_READ_FAILED",
operation: "lookup",
path: "/fixture/blob.sqlite",
cause: Object.assign(new Error("Synthetic SQLite cause"), {
code: "ERR_SQLITE_ERROR",
errcode: 1,
}),
});
const cleanup = new Error("Synthetic cleanup failure");
const combined = new AggregateError([primary, cleanup, primary], "read and cleanup", {
@ -180,28 +144,18 @@ describe("shared-state worker error transport", () => {
const decoded = roundTrip(combined);
assert(decoded instanceof AggregateError);
expect(decoded.errors[0]).toBeInstanceOf(PluginBlobStoreError);
expect(decoded.errors[0]).toMatchObject({
code: "PLUGIN_BLOB_READ_FAILED",
operation: "lookup",
path: "/fixture/blob.sqlite",
cause: { message: "Synthetic SQLite cause", code: "ERR_SQLITE_ERROR", errcode: 1 },
});
expect(decoded.errors[0]).toBe(decoded.errors[2]);
expect(decoded.cause).toBe(decoded.errors[0]);
expect(decoded.errors[1]).toMatchObject({ message: cleanup.message });
expect(decoded.errors[3]).toBe(decoded);
});
it.each([
"OPENCLAW_STATE_LEASE_INVALID_INPUT",
"OPENCLAW_STATE_LEASE_HELD",
"OPENCLAW_STATE_LEASE_ABORTED",
"OPENCLAW_STATE_LEASE_LOST",
"OPENCLAW_STATE_LEASE_STORAGE_FAILED",
] as const)("preserves lease classification and cause for %s", (code) => {
const error = new OpenClawStateLeaseError("Synthetic lease refusal", {
code,
cause: new Error("Synthetic verification cause"),
});
const decoded = roundTrip(error);
expect(decoded).toMatchObject({ code });
expect(decoded.cause).toBeInstanceOf(Error);
expect(decoded.cause).toMatchObject({ message: "Synthetic verification cause" });
});
it("preserves canonical verification wrapping through worker transport", () => {
const cause = new Error("Synthetic read failure");
const identity = { scope: "test", key: "read", leaseLabel: "test lease" };
@ -209,6 +163,7 @@ describe("shared-state worker error transport", () => {
expect(wrapped.cause).toBe(cause);
expect(toOpenClawStateLeaseVerificationError(identity, wrapped)).toBe(wrapped);
const decoded = roundTrip(wrapped);
expect(decoded.cause).toBeInstanceOf(Error);
expect(decoded).toMatchObject({
code: "OPENCLAW_STATE_LEASE_STORAGE_FAILED",
message: "failed to verify test lease test/read",
@ -350,22 +305,6 @@ describe("shared-state worker error transport", () => {
expect(combined.errors).toEqual([first, second]);
});
it("encodes and hydrates ordinary error graphs only with an explicit opt-in", () => {
const cause = Object.assign(new Error("native failure"), { code: "SQLITE_IOERR" });
const original = new AggregateError([cause], "load and cleanup", { cause });
cause.cause = original;
expect(encodeOpenClawStateWorkerError(original)).toBeUndefined();
const payload = encodeOpenClawStateWorkerError(original, { includeOrdinary: true });
expect(payload).toBeDefined();
const retained = remoteError(structuredClone(payload));
expect(hydrateOpenClawStateWorkerError(retained)).toBe(retained);
const decoded = hydrateOpenClawStateWorkerError(retained, { includeOrdinary: true });
expect(decoded.cause).toMatchObject({ message: "native failure", code: "SQLITE_IOERR" });
assert(decoded instanceof AggregateError && decoded.cause instanceof Error);
expect(decoded.errors[0]).toBe(decoded.cause);
expect(decoded.cause.cause).toBe(decoded);
});
it("leaves ordinary and already-current error graphs identical", () => {
const local = new Error("caller rejected");
const ordinary = new AggregateError([local], "caller and cleanup", { cause: local });
@ -401,10 +340,6 @@ describe("shared-state worker error transport", () => {
error: new StartupMaintenanceRequiredError("state-migrations", "state migration"),
fields: { kind: "state-migrations", reason: "state migration" },
},
{
error: new StartupMaintenanceRequiredError("legacy-session-store", "session migration"),
fields: { kind: "legacy-session-store", reason: "session store migration" },
},
{
error: new SqliteSchemaVersionError("newer schema"),
fields: { kind: "newer-schema", reason: "a newer OpenClaw build" },
@ -554,6 +489,7 @@ describe("shared-state worker error transport", () => {
});
original.errors.push(original);
const options = { includeOrdinary: true };
expect(encodeOpenClawStateWorkerError(original)).toBeUndefined();
const payload = encodeOpenClawStateWorkerError(original, options);
expect(payload).toBeDefined();
expect(JSON.stringify(payload)).not.toContain("fixture-not-for-transport");

View file

@ -153,70 +153,65 @@ function createLeaseFixture() {
return { maintenance, context };
}
it.each(["execute", "runOperation"] as const)(
"closes fresh %s admission while accepted callbacks and terminal writes finish",
async (entry) => {
const { maintenance, context } = createLeaseFixture();
const acceptedCommand = {
type: "capture.endSession" as const,
input: { sessionId: "accepted", endedAt: 1 },
};
const terminalCommand = {
type: "capture.endSession" as const,
input: { sessionId: "terminal", endedAt: 2 },
};
const forbiddenCommand = {
type: "capture.endSession" as const,
input: { sessionId: "not-admitted", endedAt: 3 },
};
const finishAccepted = createDeferredCore();
const acceptedFinished = createDeferredCore();
const pendingCommand = createDeferredCore();
const lease = maintenance.run(() =>
createOpenClawStateWorkerLease(context, async (terminal) => {
physical.events.push("finalize-start");
finishAccepted.resolve();
await acceptedFinished.promise;
await terminal.execute(terminalCommand);
physical.events.push("finalize-end");
}),
);
await lease.ready;
const accepted = lease.runOperation(async (operation) => {
try {
await finishAccepted.promise;
await operation.execute(acceptedCommand);
physical.events.push("accepted-complete");
} finally {
acceptedFinished.resolve();
}
});
void maintenance.track(pendingCommand.promise);
const closed = maintenance.close();
const lateCallback = vi.fn(async (operation: DomainScope) =>
operation.execute(forbiddenCommand),
);
it("closes fresh callback admission while accepted callbacks and terminal writes finish", async () => {
const { maintenance, context } = createLeaseFixture();
const acceptedCommand = {
type: "capture.endSession" as const,
input: { sessionId: "accepted", endedAt: 1 },
};
const terminalCommand = {
type: "capture.endSession" as const,
input: { sessionId: "terminal", endedAt: 2 },
};
const forbiddenCommand = {
type: "capture.endSession" as const,
input: { sessionId: "not-admitted", endedAt: 3 },
};
const finishAccepted = createDeferredCore();
const acceptedFinished = createDeferredCore();
const pendingCommand = createDeferredCore();
const lease = maintenance.run(() =>
createOpenClawStateWorkerLease(context, async (terminal) => {
physical.events.push("finalize-start");
finishAccepted.resolve();
await acceptedFinished.promise;
await terminal.execute(terminalCommand);
physical.events.push("finalize-end");
}),
);
await lease.ready;
const accepted = lease.runOperation(async (operation) => {
try {
const late =
entry === "execute" ? lease.execute(forbiddenCommand) : lease.runOperation(lateCallback);
await expect(late).rejects.toThrow("Database maintenance resource admission is closed");
expect(lateCallback).not.toHaveBeenCalled();
expect(physical.events).toEqual([]);
await finishAccepted.promise;
await operation.execute(acceptedCommand);
physical.events.push("accepted-complete");
} finally {
pendingCommand.resolve();
await closed;
await accepted;
acceptedFinished.resolve();
}
expect(physical.events).toEqual([
"finalize-start",
acceptedCommand,
"accepted-complete",
terminalCommand,
"finalize-end",
"physical-close",
]);
},
);
});
void maintenance.track(pendingCommand.promise);
const closed = maintenance.close();
const lateCallback = vi.fn(async (operation: DomainScope) => operation.execute(forbiddenCommand));
try {
await expect(lease.runOperation(lateCallback)).rejects.toThrow(
"Database maintenance resource admission is closed",
);
expect(lateCallback).not.toHaveBeenCalled();
expect(physical.events).toEqual([]);
} finally {
pendingCommand.resolve();
await closed;
await accepted;
}
expect(physical.events).toEqual([
"finalize-start",
acceptedCommand,
"accepted-complete",
terminalCommand,
"finalize-end",
"physical-close",
]);
});
it("drains a command accepted before cold acquisition after its callback returns", async () => {
const openGate = createDeferredCore();

View file

@ -105,7 +105,6 @@ describe("status model authentication and endpoint", () => {
it.each([
["api_key", "api-key (codex)", "https://api.openai.com/v1"],
["oauth", "oauth (codex)", "https://chatgpt.com/backend-api/codex"],
["token", "token (codex)", "https://chatgpt.com/backend-api/codex"],
] as const)(
"renders the prepared %s mode and selected route without a host credential",
async (mode, authLabel, endpoint) => {
@ -227,7 +226,6 @@ describe("status model authentication and endpoint", () => {
it.each([
["apiKey", "api_key", "api-key (codex)"],
["chatgpt", "oauth", "oauth (codex)"],
["chatgpt", "token", "token (codex)"],
] as const)(
"renders %s discovery with its observed %s mode",
@ -250,18 +248,6 @@ describe("status model authentication and endpoint", () => {
).toEqual({ authLabel: "native (codex)", endpoint: undefined });
});
it("rejects a retired discovery observation together with its mode", async () => {
expect(
await statusAuth(
{ source: "native", mode: "api_key" },
{
nativeDiscovery: { accountType: "apiKey", authMode: "api_key" },
current: () => false,
},
)(selection),
).toEqual({ authLabel: "unknown" });
});
it("does not substitute native login for an unavailable explicit profile", async () => {
const sessionEntry: SessionEntry = {
sessionId: "status-pin",

View file

@ -25,7 +25,6 @@ describe("model endpoint display", () => {
});
it.each([
["https://api.openai.com/v1/responses", "https://api.openai.com/v1/responses"],
[
"wss://chatgpt.com/backend-api/codex/responses",
"wss://chatgpt.com/backend-api/codex/responses",
@ -35,10 +34,8 @@ describe("model endpoint display", () => {
"https://example.test:9443/v1/responses",
],
["http://127.0.0.1:1234/private-token/v1/responses", "http://127.0.0.1:1234/[path hidden]"],
["https://example.test/tenant/secret", "https://example.test/[path hidden]"],
["https://example.test/v1/", "https://example.test/v1"],
["https://example.test", "https://example.test"],
["https://example.test/[path hidden]", "https://example.test/[path hidden]"],
["file:///tmp/credentials", undefined],
["not a URL", undefined],
["https://example.test/" + "x".repeat(8192), undefined],

View file

@ -1,7 +1,6 @@
import { withTempHome } from "openclaw/plugin-sdk/test-env";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { testing as cliBackendsTesting } from "../agents/cli-backends.test-support.js";
import { resolveDefaultModelForAgent } from "../agents/model-selection-config.js";
import { resolveSessionStorePathCore } from "../config/sessions/paths.js";
import {
appendTranscriptMessageSync,
@ -48,9 +47,20 @@ describe("buildStatusText prepared context windows", () => {
state = await createOpenClawTestState({ label: "status-model" });
});
afterEach(async () => {
vi.restoreAllMocks();
cliBackendsTesting.resetDepsForTest();
await state.cleanup();
});
const tokenUsage = {
totalTokens: 45_000,
totalTokensFresh: true,
totalTokensVersion: 1,
} satisfies Partial<InternalSessionEntry>;
const fallbackNotice = {
kind: "active",
selectedModel: "deepseek/deepseek-v4-flash",
activeModel: "fallback/small-model",
} satisfies NonNullable<InternalSessionEntry["fallbackNotice"]>;
const catalog = [
{
provider: "deepseek",
@ -72,17 +82,13 @@ describe("buildStatusText prepared context windows", () => {
},
];
async function renderPreparedStatus(
overrides: Partial<Parameters<typeof buildStatusReplyParts>[0]> = {},
) {
async function renderPreparedStatus(overrides: Partial<StatusTextParams> = {}) {
return await buildStatusReplyParts({
cfg: {},
sessionEntry: {
sessionId: "prepared-context",
updatedAt: 0,
totalTokens: 45_000,
totalTokensFresh: true,
totalTokensVersion: 1,
...tokenUsage,
},
sessionKey: "agent:main:main",
statusChannel: "mobilechat",
@ -103,40 +109,31 @@ describe("buildStatusText prepared context windows", () => {
});
}
it.each([
{ agentThinking: undefined, agentDefault: undefined, expected: "high" },
{ agentThinking: false, agentDefault: undefined, expected: "off" },
{ agentThinking: "high", agentDefault: "minimal", expected: "minimal" },
] as const)(
"renders configured thinking precedence (model=$agentThinking, agent=$agentDefault)",
async ({ agentThinking, agentDefault, expected }) => {
const parts = await renderPreparedStatus({
cfg: {
agents: {
defaults: {
thinkingDefault: "low",
it("renders the agent thinking default ahead of model and global defaults", async () => {
const parts = await renderPreparedStatus({
cfg: {
agents: {
defaults: {
thinkingDefault: "low",
models: { "fixture/reasoning-model": { params: { thinking: "high" } } },
},
entries: {
main: {
thinkingDefault: "minimal",
models: { "fixture/reasoning-model": { params: { thinking: "high" } } },
},
entries: {
main: {
thinkingDefault: agentDefault,
models: { "fixture/reasoning-model": { params: { thinking: agentThinking } } },
},
},
},
},
provider: "fixture",
model: "reasoning-model",
thinkingCatalog: [{ provider: "fixture", id: "reasoning-model", reasoning: true }],
});
},
provider: "fixture",
model: "reasoning-model",
thinkingCatalog: [{ provider: "fixture", id: "reasoning-model", reasoning: true }],
});
expect(parts.text).toContain(`think ${expected}`);
},
);
expect(parts.text).toContain("think minimal");
});
it.each([
{ name: "selected model on", configured: true, expected: "on" },
{ name: "selected model off", configured: false, expected: "off" },
{
name: "selected model auto with a different active fallback cutoff",
configured: "auto",
@ -149,12 +146,6 @@ describe("buildStatusText prepared context windows", () => {
sessionFast: false,
expected: "off",
},
{
name: "session on overrides model off",
configured: false,
sessionFast: true,
expected: "on",
},
{
name: "prepared off overrides model on",
configured: true,
@ -209,7 +200,7 @@ describe("buildStatusText prepared context windows", () => {
entry?: Partial<InternalSessionEntry>;
message?: Record<string, unknown>;
laterMessage?: Record<string, unknown>;
status?: Partial<Parameters<typeof buildStatusReplyParts>[0]>;
status?: Partial<StatusTextParams>;
} = {},
) {
return await withTempHome(async () => {
@ -229,15 +220,8 @@ describe("buildStatusText prepared context windows", () => {
agentHarnessId: "openclaw",
contextTokens: 1_000_000,
contextTokensSource: "runtime",
totalTokens: 45_000,
totalTokensFresh: true,
totalTokensVersion: 1,
fallbackNotice: {
kind: "active",
selectedModel: "deepseek/deepseek-v4-flash",
activeModel: "fallback/small-model",
reason: "provider unavailable",
},
...tokenUsage,
fallbackNotice: { ...fallbackNotice, reason: "provider unavailable" },
...params.entry,
};
replaceSessionEntrySync(scope, entry);
@ -275,29 +259,14 @@ describe("buildStatusText prepared context windows", () => {
it.each([
["stale runtime telemetry", {}],
["stale resolved context", { entry: { contextTokensSource: "resolved-v1" } }],
["absent runtime model", { entry: { modelProvider: undefined, model: undefined } }],
[
"padded selected notice",
{
entry: {
fallbackNotice: {
kind: "active",
selectedModel: " deepseek/deepseek-v4-flash ",
activeModel: "fallback/small-model",
},
},
},
],
[
"literal provider-local model",
{
status: { provider: "MiXeD", model: "Vendor/Model:opaque" },
entry: {
fallbackNotice: {
kind: "active",
...fallbackNotice,
selectedModel: "mixed/Vendor/Model:opaque",
activeModel: "fallback/small-model",
},
},
},
@ -308,9 +277,8 @@ describe("buildStatusText prepared context windows", () => {
entry: {
modelOverride: "MiXeD/Model:Case",
fallbackNotice: {
kind: "active",
...fallbackNotice,
selectedModel: "MiXeD/Model:Case",
activeModel: "fallback/small-model",
},
},
},
@ -332,17 +300,12 @@ describe("buildStatusText prepared context windows", () => {
it.each([
["running session", { entry: { status: "running" } }],
["failed session", { entry: { status: "failed" } }],
["killed session", { entry: { status: "killed" } }],
["timed-out session", { entry: { status: "timeout" } }],
["missing run", { entry: { lastRunId: undefined } }],
["other run", { entry: { lastRunId: "other-run" } }],
["failed assistant", { message: { stopReason: "error" } }],
["hidden assistant", { message: { content: [] } }],
["undisplayed assistant", { message: { display: false } }],
["oversized tail", { message: { content: [{ type: "text", text: "x".repeat(300_000) }] } }],
["later user", { laterMessage: { role: "user", content: "New turn" } }],
["later tool", { laterMessage: { role: "toolResult", toolCallId: "later", content: [] } }],
[
"later run",
{
@ -363,8 +326,7 @@ describe("buildStatusText prepared context windows", () => {
{
entry: {
fallbackNotice: {
kind: "active",
selectedModel: "deepseek/deepseek-v4-flash",
...fallbackNotice,
activeModel: "fallback/other-model",
},
},
@ -381,22 +343,17 @@ describe("buildStatusText prepared context windows", () => {
it("skips terminal transcript access for a stale selected notice", async () => {
const readTail = vi.spyOn(transcriptTail, "readSessionTranscriptBoundedMessageTailPage");
try {
const parts = await renderTerminalFallback({
entry: {
fallbackNotice: {
kind: "active",
selectedModel: "deepseek/older-model",
activeModel: "fallback/small-model",
},
const parts = await renderTerminalFallback({
entry: {
fallbackNotice: {
...fallbackNotice,
selectedModel: "deepseek/older-model",
},
});
expect(parts.text).not.toContain("Fallback: fallback/small-model");
expect(parts.text).toContain("Context: 45k/1.0m");
expect(readTail).not.toHaveBeenCalled();
} finally {
readTail.mockRestore();
}
},
});
expect(parts.text).not.toContain("Fallback: fallback/small-model");
expect(parts.text).toContain("Context: 45k/1.0m");
expect(readTail).not.toHaveBeenCalled();
});
it("retains the incoming prepared cap when it already belongs to the terminal pair", async () => {
@ -429,25 +386,20 @@ describe("buildStatusText prepared context windows", () => {
inputTokens: 10,
outputTokens: 2,
});
try {
const parts = await renderTerminalFallback({
entry: {
modelProvider: undefined,
model: undefined,
fallbackNotice: {
kind: "active",
selectedModel: "deepseek/deepseek-v4-flash",
activeModel: notice,
reason: "provider unavailable",
},
const parts = await renderTerminalFallback({
entry: {
modelProvider: undefined,
model: undefined,
fallbackNotice: {
...fallbackNotice,
activeModel: notice,
reason: "provider unavailable",
},
status: { includeTranscriptUsage: true },
});
expect(readUsage).toHaveBeenCalled();
expect(parts.text).toContain(`Fallback: ${notice}`);
} finally {
readUsage.mockRestore();
}
},
status: { includeTranscriptUsage: true },
});
expect(readUsage).toHaveBeenCalled();
expect(parts.text).toContain(`Fallback: ${notice}`);
});
const budget: SessionContextBudgetStatus = {
@ -471,19 +423,15 @@ describe("buildStatusText prepared context windows", () => {
sessionId: "terminal-fallback",
};
it.each([
["matching budget", {}, false, true],
["selected model budget", { provider: "deepseek", model: "deepseek-v4-flash" }, false, false],
["other session budget", { sessionId: "other-session" }, false, false],
["other cap budget", { contextTokenBudget: 1_000_000 }, false, false],
["pending switch budget", {}, true, false],
] satisfies Array<[string, Partial<SessionContextBudgetStatus>, boolean, boolean]>)(
["matching budget", {}, true],
["selected model budget", { provider: "deepseek", model: "deepseek-v4-flash" }, false],
] satisfies Array<[string, Partial<SessionContextBudgetStatus>, boolean]>)(
"projects only a %s owned by the terminal model",
async (_name, patch, pending, expected) => {
async (_name, patch, expected) => {
const parts = await renderTerminalFallback({
entry: {
totalTokens: undefined,
contextBudgetStatus: { ...budget, ...patch },
liveModelSwitchPending: pending,
},
});
expect(parts.text).toContain("Fallback: fallback/small-model");
@ -493,7 +441,6 @@ describe("buildStatusText prepared context windows", () => {
);
it.each([
["prepared alias", "candidate", "middle", {}, "candidate/middle"],
[
"opaque empty provider",
"",
@ -508,13 +455,6 @@ describe("buildStatusText prepared context windows", () => {
{ providerOverride: "candidate", modelOverride: "middle" },
"candidate/middle",
],
[
"legacy explicit override",
"candidate",
"entry",
{ modelOverride: "fallback/small-model" },
"fallback/small-model",
],
] satisfies Array<[string, string, string, Partial<InternalSessionEntry>, string]>)(
"preserves typed selection for %s through the status owner",
async (_name, provider, model, patch, expected) => {
@ -552,105 +492,25 @@ describe("buildStatusText prepared context windows", () => {
expect(parts.text).toContain("Model: deepseek/deepseek-v4-flash");
});
it.each([
{
name: "configured CLI default with absent session model fields",
cfg: { agents: { defaults: { model: "claude-cli/opus" } } },
expectedModel: "claude-cli/opus",
},
{
name: "canonical default with absent agent configuration",
cfg: {},
expectedModel: "openai/gpt-6-astra",
},
{
name: "literal self-provider prefix in a prepared model ID",
cfg: { agents: { defaults: { model: "deepseek/deepseek-v4-flash" } } },
input: { provider: "custom", model: "custom/model" },
expectedModel: "custom/custom/model",
absent: ["Model: custom/model"],
},
{
name: "manual selection that remains pending on an active session",
cfg: { agents: { defaults: { model: "deepseek/deepseek-v4-flash" } } },
input: { provider: "deepseek", model: "deepseek-v4-flash" },
entry: {
status: "running",
providerOverride: "fallback",
modelOverride: "small-model",
modelOverrideSource: "user",
modelProvider: "deepseek",
model: "deepseek-v4-flash",
liveModelSwitchPending: true,
},
expectedModel: "fallback/small-model",
expected: ["live switch pending", "pinned session"],
absent: ["Fallback:", "auto fallback"],
},
{
name: "configured subagent selection with matching automatic origin",
cfg: {
agents: {
ownership: "explicit",
defaults: { model: "deepseek/deepseek-v4-flash" },
entries: { worker: { subagents: { model: "fallback/small-model" } } },
},
},
agentId: "worker",
sessionKey: "agent:worker:subagent:configured",
input: { provider: "deepseek", model: "deepseek-v4-flash" },
entry: {
providerOverride: "fallback",
modelOverride: "small-model",
modelOverrideSource: "auto",
modelOverrideFallbackOriginProvider: "fallback",
modelOverrideFallbackOriginModel: "small-model",
},
expectedModel: "fallback/small-model",
absent: ["auto fallback", "check provider", "pinned session"],
},
] satisfies Array<{
name: string;
cfg: StatusTextParams["cfg"];
agentId?: string;
sessionKey?: string;
input?: Pick<StatusTextParams, "provider" | "model">;
entry?: Partial<InternalSessionEntry>;
expectedModel: string;
expected?: string[];
absent?: string[];
}>)("preserves $name through the actual status owner", async (control) => {
const agentId = control.agentId ?? "main";
it("preserves a literal self-provider prefix in a prepared model ID", async () => {
const sessionEntry: InternalSessionEntry = {
sessionId: "selection-owner-control",
updatedAt: 1,
...control.entry,
};
const original = structuredClone(sessionEntry);
// Defaults use the real production selector; other rows carry prepared caller facts.
const selection =
control.input ??
resolveDefaultModelForAgent({ cfg: control.cfg, agentId, allowPluginNormalization: false });
const parts = await renderPreparedStatus({
cfg: control.cfg,
agentId,
sessionKey: control.sessionKey ?? "agent:main:main",
provider: selection.provider,
model: selection.model,
cfg: { agents: { defaults: { model: "deepseek/deepseek-v4-flash" } } },
provider: "custom",
model: "custom/model",
sessionEntry,
resolvedThinkLevel: "off",
});
expect(parts.text).toContain(`Model: ${control.expectedModel}`);
for (const expected of control.expected ?? []) {
expect(parts.text).toContain(expected);
}
for (const absent of control.absent ?? []) {
expect(parts.text).not.toContain(absent);
}
expect(parts.text).toContain("Model: custom/custom/model");
expect(parts.text).not.toContain("Model: custom/model");
const table = parts.presentation.blocks.find((block) => block.type === "table");
expect(table?.type === "table" ? table.rows : []).toContainEqual([
"🧠 Model",
expect.stringContaining(control.expectedModel),
expect.stringContaining("custom/custom/model"),
]);
expect(sessionEntry).toEqual(original);
});
@ -665,28 +525,20 @@ describe("buildStatusText prepared context windows", () => {
.mockImplementation(() => {
throw error;
});
try {
const sessionEntry: InternalSessionEntry = {
sessionId: "projection",
updatedAt: 1,
status: "done",
lastRunId: "settled-run",
fallbackNotice: {
kind: "active",
selectedModel: "deepseek/deepseek-v4-flash",
activeModel: "fallback/small-model",
},
};
const result = renderPreparedStatus({ sessionEntry });
if (unavailable) {
expect((await result).text).not.toContain("Fallback:");
} else {
await expect(result).rejects.toBe(error);
}
expect(readTail).toHaveBeenCalledOnce();
} finally {
readTail.mockRestore();
const sessionEntry: InternalSessionEntry = {
sessionId: "projection",
updatedAt: 1,
status: "done",
lastRunId: "settled-run",
fallbackNotice,
};
const result = renderPreparedStatus({ sessionEntry });
if (unavailable) {
expect((await result).text).not.toContain("Fallback:");
} else {
await expect(result).rejects.toBe(error);
}
expect(readTail).toHaveBeenCalledOnce();
});
it("renders a cold-cache prepared window in plain and rich status", async () => {
@ -711,9 +563,7 @@ describe("buildStatusText prepared context windows", () => {
modelOverrideSource: "user",
modelProvider: "fallback",
model: "small-model",
totalTokens: 45_000,
totalTokensFresh: true,
totalTokensVersion: 1,
...tokenUsage,
},
});
@ -730,15 +580,8 @@ describe("buildStatusText prepared context windows", () => {
modelOverride: "deepseek-v4-flash",
modelProvider: "fallback",
model: "small-model",
fallbackNotice: {
kind: "active",
selectedModel: "deepseek/deepseek-v4-flash",
activeModel: "fallback/small-model",
reason: "provider unavailable",
},
totalTokens: 45_000,
totalTokensFresh: true,
totalTokensVersion: 1,
fallbackNotice: { ...fallbackNotice, reason: "provider unavailable" },
...tokenUsage,
},
});
@ -771,9 +614,7 @@ describe("buildStatusText prepared context windows", () => {
agentHarnessId: "claude-cli",
contextTokens: 256_000,
contextTokensSource: "resolved",
totalTokens: 45_000,
totalTokensFresh: true,
totalTokensVersion: 1,
...tokenUsage,
},
thinkingCatalog: [
{

View file

@ -8,7 +8,7 @@ import {
} from "./command-export.js";
import { resolveTrajectoryFilePath } from "./paths.js";
import { createTrajectoryRuntimeRecorder } from "./runtime.js";
import type { TrajectoryBundleManifest, TrajectoryEvent } from "./types.js";
import type { TrajectoryBundleManifest } from "./types.js";
const tool = { name: "synthetic_read", description: "Read a synthetic document." };
const complete = { systemPrompt: "Synthetic system prompt.", tools: [tool] };
@ -19,14 +19,12 @@ const cases: Array<{
contextFiles: string[];
outputPath?: string;
}> = [
{ name: "no context", contexts: [], contextFiles: [], outputPath: "nested/inventory" },
{
name: "complete context",
contexts: [complete],
contextFiles: ["system-prompt.txt", "tools.json"],
outputPath: "prefix/../~/inventory",
},
{ name: "oversized prompt", contexts: [oversized], contextFiles: ["tools.json"] },
{
name: "empty prompt",
contexts: [{ systemPrompt: "", tools: [] }],
@ -105,13 +103,6 @@ describe("trajectory command export inventory", () => {
}
recorder.recordEvent("model.completed", { stopReason: "stop" });
await recorder.flush();
expect(writes).toHaveLength(contexts.length + 3);
const recorded = writes.map((line) => JSON.parse(line) as TrajectoryEvent);
const latest = recorded.findLast((event) => event.type === "context.compiled")?.data;
expect(typeof latest?.systemPrompt === "string" && latest.systemPrompt.length > 0).toBe(
contextFiles.includes("system-prompt.txt"),
);
expect(Array.isArray(latest?.tools)).toBe(contextFiles.includes("tools.json"));
const runtimeBytes = writes.join("");
await fs.writeFile(runtimeFile, runtimeBytes);
@ -156,7 +147,7 @@ describe("trajectory command export inventory", () => {
},
);
it.each([".openclaw", ".openclaw/trajectory-exports", ".openclaw/trajectory-exports/alias"])(
it.each([".openclaw", ".openclaw/trajectory-exports/alias"])(
"rejects an escaping %s directory without writing outside the workspace",
async (relativeLink) => {
await withTempDir("openclaw-trajectory-boundary-", async (root) => {

View file

@ -11,6 +11,10 @@ function arrayReturning(value: unknown): unknown[] {
});
}
function projectParameters(parameters: unknown) {
return toTrajectoryToolDefinitions([{ name: "sample", parameters }])[0]?.parameters;
}
describe("trajectory tool definition preparation", () => {
afterEach(() => {
vi.restoreAllMocks();
@ -51,71 +55,49 @@ describe("trajectory tool definition preparation", () => {
recorder?.recordEvent("context.compiled", {
tools: toTrajectoryToolDefinitions([{ name: "sample", parameters: { description: text } }]),
});
try {
record(description);
record(description);
registerSecretValueForRedaction(description);
record(description);
record("Authorization: Bearer synthetic-changed-value");
record("policy-fixture-value");
applyLoggingConfig({ redactPatterns: ["policy-fixture-value"] });
record("policy-fixture-value");
applyLoggingConfig(undefined);
record("policy-fixture-value");
record(description);
record(description);
registerSecretValueForRedaction(description);
record(description);
record("Authorization: Bearer synthetic-changed-value");
record("policy-fixture-value");
applyLoggingConfig({ redactPatterns: ["policy-fixture-value"] });
record("policy-fixture-value");
applyLoggingConfig(undefined);
record("policy-fixture-value");
expect(writes).toHaveLength(7);
expect(writes[0]).toContain(description);
expect(writes[1]).toContain(description);
expect(writes[2]).not.toContain(description);
expect(writes[3]).not.toContain("synthetic-changed-value");
expect(JSON.parse(writes[3]!).data.tools[0].parameters.description).toContain("redacted");
expect(writes[4]).toContain("policy-fixture-value");
expect(writes[5]).not.toContain("policy-fixture-value");
expect(writes[6]).toContain("policy-fixture-value");
} finally {
resetSecretRedactionRegistryForTest();
}
expect(writes).toHaveLength(7);
expect(writes[0]).toContain(description);
expect(writes[1]).toContain(description);
expect(writes[2]).not.toContain(description);
expect(writes[3]).not.toContain("synthetic-changed-value");
expect(JSON.parse(writes[3]!).data.tools[0].parameters.description).toContain("redacted");
expect(writes[4]).toContain("policy-fixture-value");
expect(writes[5]).not.toContain("policy-fixture-value");
expect(writes[6]).toContain("policy-fixture-value");
});
it("does not invoke custom serialization hooks or cache opaque array-operation results", () => {
const toJSON = vi.fn(() => ({ unexpected: true }));
const opaque = Object.defineProperty({ ordinary: "first" }, "toJSON", { value: toJSON });
const parameters = arrayReturning(opaque);
expect(toTrajectoryToolDefinitions([{ name: "sample", parameters }])[0]?.parameters).toEqual({
ordinary: "first",
});
expect(projectParameters(parameters)).toEqual({ ordinary: "first" });
opaque.ordinary = "second";
expect(toTrajectoryToolDefinitions([{ name: "sample", parameters }])[0]?.parameters).toEqual({
ordinary: "second",
});
expect(projectParameters(parameters)).toEqual({ ordinary: "second" });
expect(toJSON).not.toHaveBeenCalled();
});
it("preserves truncation metadata without probing its records as native headers", () => {
const has = vi.spyOn(Headers.prototype, "has");
try {
expect(
toTrajectoryToolDefinitions([
{ name: "sample", parameters: { description: "x".repeat(32_769) } },
]),
).toEqual([
{
name: "sample",
description: undefined,
parameters: {
description: {
truncated: true,
reason: "trajectory-field-size-limit",
originalChars: 32_769,
limitChars: 32_768,
},
},
},
]);
expect(has).not.toHaveBeenCalled();
} finally {
has.mockRestore();
}
expect(projectParameters({ description: "x".repeat(32_769) })).toEqual({
description: {
truncated: true,
reason: "trajectory-field-size-limit",
originalChars: 32_769,
limitChars: 32_768,
},
});
expect(has).not.toHaveBeenCalled();
});
it("preserves native retry metadata returned by custom array operations", () => {
@ -124,20 +106,10 @@ describe("trajectory tool definition preparation", () => {
authorization: "Bearer synthetic-credential",
});
expect(
toTrajectoryToolDefinitions([
{
name: "sample",
parameters: { ordinary: "kept", nested: arrayReturning(headers) },
},
]),
).toEqual([
{
name: "sample",
description: undefined,
parameters: { ordinary: "kept", nested: { "retry-after-ms": 7_000 } },
},
]);
expect(projectParameters({ ordinary: "kept", nested: arrayReturning(headers) })).toEqual({
ordinary: "kept",
nested: { "retry-after-ms": 7_000 },
});
});
it("preserves prototype traps when a copied field changes the prepared record prototype", () => {
@ -148,9 +120,7 @@ describe("trajectory tool definition preparation", () => {
ordinary: "kept",
};
expect(toTrajectoryToolDefinitions([{ name: "sample", parameters }])).toEqual([
{ name: "sample", description: undefined, parameters: { ordinary: "kept" } },
]);
expect(projectParameters(parameters)).toEqual({ ordinary: "kept" });
expect(getPrototypeOf).toHaveBeenCalled();
expect(Object.getPrototypeOf(parameters)).toBe(Object.prototype);
});
@ -163,13 +133,7 @@ describe("trajectory tool definition preparation", () => {
get: readNested,
});
expect(toTrajectoryToolDefinitions([{ name: "sample", parameters }])).toEqual([
{
name: "sample",
description: undefined,
parameters: { nested: { ordinary: "kept", toJSON } },
},
]);
expect(projectParameters(parameters)).toEqual({ nested: { ordinary: "kept", toJSON } });
expect(readNested).toHaveBeenCalledOnce();
expect(toJSON).not.toHaveBeenCalled();
});

View file

@ -11,16 +11,20 @@ import {
const fixture = useTranscriptStatusFixture();
async function startCapture(sessionId: string) {
const f = fixture();
const started = createDeferred<TranscriptStartRequest>();
f.provider.start = async (request) => {
started.resolve(request);
return { ok: true, session: request.session };
};
await f.start({ ...room, sessionId });
return { f, request: await started.promise };
}
describe("transcript capture accepted append drainage", () => {
it("persists concurrently accepted speech in order before terminal metadata and notes", async () => {
const f = fixture();
const started = createDeferred<TranscriptStartRequest>();
f.provider.start = async (request) => {
started.resolve(request);
return { ok: true, session: request.session };
};
await f.start({ ...room, sessionId: "terminal-drain" });
const request = await started.promise;
const { f, request } = await startCapture("terminal-drain");
const appendEntered = createDeferred();
const releaseAppend = createDeferred();
const events: string[] = [];
@ -266,14 +270,7 @@ describe("transcript capture accepted append drainage", () => {
]);
});
it("settles metadata rejection promptly while preserving earlier accepted speech ordering", async () => {
const f = fixture();
const started = createDeferred<TranscriptStartRequest>();
f.provider.start = async (request) => {
started.resolve(request);
return { ok: true, session: request.session };
};
await f.start({ ...room, sessionId: "serialization-rejection" });
const request = await started.promise;
const { f, request } = await startCapture("serialization-rejection");
const appendEntered = createDeferred();
const releaseAppend = createDeferred();
const append = f.store.appendUtteranceForSession.bind(f.store);
@ -322,14 +319,7 @@ describe("transcript capture accepted append drainage", () => {
});
it("drains an accepted append when its metadata serialization ends the capture", async () => {
const f = fixture();
const started = createDeferred<TranscriptStartRequest>();
f.provider.start = async (request) => {
started.resolve(request);
return { ok: true, session: request.session };
};
await f.start({ ...room, sessionId: "serialization-terminal" });
const request = await started.promise;
const { f, request } = await startCapture("serialization-terminal");
const terminal = createDeferred();
const toJSON = vi.fn(() => {
// Provider metadata can synchronously end capture while this speech is being prepared.

View file

@ -1,4 +1,3 @@
import path from "node:path";
import { setImmediate } from "node:timers/promises";
import { afterEach, describe, expect, it, vi } from "vitest";
import { createDeferred } from "../../test/helpers/promise.js";
@ -10,8 +9,11 @@ import {
import { createTranscriptCaptureAppends } from "./capture-appends.js";
import { activeSessions } from "./capture-startup.js";
import { exportTranscriptLibrary, getTranscriptLibrary, listTranscriptLibrary } from "./library.js";
import type { TranscriptSessionDescriptor } from "./provider-types.js";
import { TranscriptsStore, transcriptSessionSelector } from "./store.js";
import {
createTranscriptLibraryStoreFixture,
transcriptLibrarySession as session,
} from "./library.store.test-support.js";
import { transcriptSessionSelector } from "./store.js";
import { summarizeTranscripts } from "./summary.js";
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
@ -23,25 +25,7 @@ afterEach(async () => {
});
function fixture() {
const stateDir = tempDirs.make("transcript-library-async-");
return {
store: new TranscriptsStore(path.join(stateDir, "transcripts"), {
env: { ...process.env, OPENCLAW_STATE_DIR: stateDir },
}),
};
}
function session(
sessionId: string,
overrides: Partial<TranscriptSessionDescriptor> = {},
): TranscriptSessionDescriptor {
return {
sessionId,
title: sessionId,
source: { providerId: "manual-transcript" },
startedAt: "2026-08-20T10:00:00.000Z",
...overrides,
};
return createTranscriptLibraryStoreFixture(tempDirs.make("transcript-library-async-"));
}
describe("transcript library asynchronous reads", () => {
@ -86,22 +70,15 @@ describe("transcript library asynchronous reads", () => {
(await exportTranscriptLibrary(store, { selector, format: "markdown" })).data,
"base64",
).toString("utf8");
let settled = false;
const settled = vi.fn();
const result = read();
const settlement = result.then(
() => {
settled = true;
},
() => {
settled = true;
},
);
const settlement = result.then(settled, settled);
const replacement = { ...target, title: "Replacement meeting" };
const added = { text: "Peer saved note" };
try {
await reading.promise;
await setImmediate();
expect(settled).toBe(false);
expect(settled).not.toHaveBeenCalled();
await store.writeSession(replacement);
await store.appendUtteranceForSession(replacement, added);
await store.writeSummary(
@ -123,30 +100,22 @@ describe("transcript library asynchronous reads", () => {
},
);
it.each(["get", "export"] as const)(
"propagates a failed composed %s without returning partial content",
async (kind) => {
const { store } = fixture();
const target = session("rejected-read");
await store.writeSession(target);
const failure = new Error("archive read failed");
const selector = transcriptSessionSelector(target);
if (kind === "get") {
vi.spyOn(store, "readLibraryEntry").mockRejectedValueOnce(failure);
await expect(
getTranscriptLibrary(store, { selector, includeUtterances: true, limit: 1 }),
).rejects.toBe(failure);
} else {
vi.spyOn(store, "iterateExport").mockImplementationOnce(async function* () {
yield { sequence: 0, text: "Partial content" };
throw failure;
});
await expect(exportTranscriptLibrary(store, { selector, format: "jsonl" })).rejects.toBe(
failure,
);
}
},
);
it("propagates a failed export without returning partial content", async () => {
const { store } = fixture();
const target = session("rejected-read");
await store.writeSession(target);
const failure = new Error("archive read failed");
vi.spyOn(store, "iterateExport").mockImplementationOnce(async function* () {
yield { sequence: 0, text: "Partial content" };
throw failure;
});
await expect(
exportTranscriptLibrary(store, {
selector: transcriptSessionSelector(target),
format: "jsonl",
}),
).rejects.toBe(failure);
});
it("rejects an export canceled before its completion result", async () => {
const { store } = fixture();
@ -179,19 +148,12 @@ describe("transcript library asynchronous reads", () => {
const result = listTranscriptLibrary(store, {}, () => {
throw failure;
});
let settled = false;
const settlement = result.then(
() => {
settled = true;
},
() => {
settled = true;
},
);
const settled = vi.fn();
const settlement = result.then(settled, settled);
try {
await closing.promise;
await setImmediate();
expect(settled).toBe(false);
expect(settled).not.toHaveBeenCalled();
expect(closed).toBe(false);
} finally {
cleanup.resolve();

View file

@ -49,6 +49,18 @@ function fixture() {
return createTranscriptLibraryStoreFixture(tempDirs.make("transcript-library-query-budget-"));
}
function seedExportBookkeeping(db: DatabaseSync) {
executeSqliteQuerySync(
db,
meetingTranscriptDb(db)
.updateTable("meeting_transcript_sessions")
.set({
export_manifest_json: JSON.stringify({ "retained-export.md": "x".repeat(16_384) }),
export_pending_json: JSON.stringify(["x".repeat(16_384)]),
}),
);
}
function observeArchiveReads(
store: ReturnType<typeof createTranscriptLibraryStoreFixture>["store"],
database: DatabaseSync,
@ -167,15 +179,7 @@ describe("transcript library SQLite query budgets", () => {
await store.writeSession(target);
}
const db = database();
executeSqliteQuerySync(
db,
meetingTranscriptDb(db)
.updateTable("meeting_transcript_sessions")
.set({
export_manifest_json: JSON.stringify({ "retained-export.md": "x".repeat(16_384) }),
export_pending_json: JSON.stringify(["x".repeat(16_384)]),
}),
);
seedExportBookkeeping(db);
executeSqliteQuerySync(
db,
meetingTranscriptDb(db)
@ -207,15 +211,7 @@ describe("transcript library SQLite query budgets", () => {
await store.writeSession(target);
await store.appendUtteranceForSession(target, { text: "Summarize this speech" });
const db = database();
executeSqliteQuerySync(
db,
meetingTranscriptDb(db)
.updateTable("meeting_transcript_sessions")
.set({
export_manifest_json: JSON.stringify({ "retained-export.md": "x".repeat(16_384) }),
export_pending_json: JSON.stringify(["x".repeat(16_384)]),
}),
);
seedExportBookkeeping(db);
const reads = observeArchiveReads(store, db);
expect(await store.readSummarySnapshot(target, 20)).toMatchObject({
nextSequence: 1,
@ -350,20 +346,18 @@ describe("transcript library SQLite query budgets", () => {
).rejects.toThrow(expect.objectContaining({ type: "transcript_result_too_large" }));
});
it.each(["title", "source and metadata", "last timestamp"])(
it.each(["source and metadata", "last timestamp"])(
"bounds %s in list, selector and latest descriptors without limiting full-store reads",
async (field) => {
const { store, database } = fixture();
const target = session(
"descriptor",
field === "title"
? { title: "x".repeat(TRANSCRIPTS_RESULT_MAX_BYTES + 1) }
: field === "source and metadata"
? {
source: { providerId: "manual-transcript", private: "x".repeat(600_000) },
metadata: { private: "y".repeat(600_000) },
}
: {},
field === "source and metadata"
? {
source: { providerId: "manual-transcript", private: "x".repeat(600_000) },
metadata: { private: "y".repeat(600_000) },
}
: {},
);
await store.writeSession(target);
await store.appendUtteranceForSession(target, {

View file

@ -54,6 +54,10 @@ async function settleSummaryUpdates(updates: ReturnType<typeof trackSummaryUpdat
// The capture owner rearms its timer in the continuation after detached work settles.
await vi.advanceTimersByTimeAsync(0);
}
async function advanceSummaryInterval(updates: ReturnType<typeof trackSummaryUpdates>) {
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(updates);
}
const cfg = { agents: { defaults: { utilityModel: "test/utility", model: "test/primary" } } };
const modelNotes = () => ({
text: JSON.stringify({
@ -127,6 +131,17 @@ async function saved(fixture: Awaited<ReturnType<typeof capture>>) {
return (await fixture.store.readSummary(fixture.source.session)).summary;
}
function execute(
fixture: Awaited<ReturnType<typeof capture>>,
id: string,
action: "stop" | "summarize",
) {
const tool = createTranscriptsTool(fixture.ctx);
return withPluginRuntimeRegistryScope(fixture.registry, () =>
tool.execute(id, { action, sessionId: "meeting" }),
);
}
describe("live meeting summaries", () => {
it.each(["restart-drain", "capture-policy"] as const)(
"persists final speech without loading inference after %s closes admission",
@ -145,10 +160,7 @@ describe("live meeting summaries", () => {
if (policy) {
await policy.drain();
} else {
const tool = createTranscriptsTool(fixture.ctx);
await withPluginRuntimeRegistryScope(fixture.registry, () =>
tool.execute("shutdown", { action: "stop", sessionId: "meeting" }),
);
await execute(fixture, "shutdown", "stop");
}
expect(loadRuntime).toHaveBeenCalledTimes(runtimeLoadsBeforeShutdown);
expect(complete).not.toHaveBeenCalled();
@ -176,12 +188,8 @@ describe("live meeting summaries", () => {
overview: "Final notes from the durable owner",
};
await fixture.store.writeSummary(final, stopped);
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(fixture.updates);
const tool = createTranscriptsTool(fixture.ctx);
await withPluginRuntimeRegistryScope(fixture.registry, () =>
tool.execute("manual", { action: "summarize", sessionId: "meeting" }),
);
await advanceSummaryInterval(fixture.updates);
await execute(fixture, "manual", "summarize");
expect(complete).not.toHaveBeenCalled();
expect(await saved(fixture)).toEqual(final);
});
@ -196,8 +204,7 @@ describe("live meeting summaries", () => {
expect((await store.readSession("meeting"))?.stoppedAt).toBeTruthy();
const fixture = await capture(stateDir);
await fixture.source.onUtterance({ text: "Recovered capture" });
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(fixture.updates);
await advanceSummaryInterval(fixture.updates);
expect(await saved(fixture)).toMatchObject({ source: "model", utteranceCount: 1 });
expect(complete).toHaveBeenCalledOnce();
});
@ -213,10 +220,7 @@ describe("live meeting summaries", () => {
.mockReturnValueOnce(read.promise);
await vi.advanceTimersByTimeAsync(fiveMinutes);
await vi.waitFor(() => expect(pendingRead).toHaveBeenCalledOnce());
const tool = createTranscriptsTool(fixture.ctx);
const stopped = withPluginRuntimeRegistryScope(fixture.registry, () =>
tool.execute("stop", { action: "stop", sessionId: "meeting" }),
);
const stopped = execute(fixture, "stop", "stop");
await fixture.stopStarted;
expect(fixture.stop).toHaveBeenCalledOnce();
expect(complete).not.toHaveBeenCalled();
@ -231,15 +235,13 @@ describe("live meeting summaries", () => {
const fixture = await caller.track(() => capture());
await caller.drain();
await fixture.source.onUtterance({ text: "Speech after the start request finished" });
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(fixture.updates);
await advanceSummaryInterval(fixture.updates);
expect(await saved(fixture)).toMatchObject({ source: "model", utteranceCount: 1 });
expect(fixture.ctx.logger.warn).not.toHaveBeenCalled();
});
it("publishes a captured speech prefix during continued speech and skips unchanged intervals", async () => {
const fixture = await capture();
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(fixture.updates);
await advanceSummaryInterval(fixture.updates);
expect(complete).not.toHaveBeenCalled();
await fixture.source.onUtterance({ text: "First decision" });
await vi.advanceTimersByTimeAsync(fiveMinutes - 1);
@ -261,14 +263,12 @@ describe("live meeting summaries", () => {
model: "utility",
agentId: "main",
});
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(fixture.updates);
await advanceSummaryInterval(fixture.updates);
expect(await saved(fixture)).toMatchObject({
utteranceCount: 2,
transcript: ["First decision", "Later decision"],
});
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(fixture.updates);
await advanceSummaryInterval(fixture.updates);
expect(complete).toHaveBeenCalledTimes(2);
});
@ -280,15 +280,12 @@ describe("live meeting summaries", () => {
await pending.entered;
expect(complete).toHaveBeenCalledOnce();
await fixture.source.onUtterance({ text: "Final speech" });
const tool = createTranscriptsTool(fixture.ctx);
const aborted = new Promise<void>((resolve) => {
complete.mock.calls[0]![0].abortSignal.addEventListener("abort", () => resolve(), {
once: true,
});
});
const stopped = withPluginRuntimeRegistryScope(fixture.registry, () =>
tool.execute("stop", { action: "stop", sessionId: "meeting" }),
);
const stopped = execute(fixture, "stop", "stop");
await aborted;
expect(complete.mock.calls[0]![0].abortSignal.aborted).toBe(true);
expect(complete).toHaveBeenCalledOnce();
@ -347,10 +344,7 @@ describe("live meeting summaries", () => {
expect(await terminalOutcome).toBe(appendFailure);
expect((await fixture.store.readSession("meeting"))?.stoppedAt).toBeUndefined();
expect(await saved(fixture)).toBeUndefined();
const tool = createTranscriptsTool(fixture.ctx);
await withPluginRuntimeRegistryScope(fixture.registry, () =>
tool.execute("retry-stop", { action: "stop", sessionId: "meeting" }),
);
await execute(fixture, "retry-stop", "stop");
expect(fixture.stop).not.toHaveBeenCalled();
expect(await saved(fixture)).toMatchObject({ transcript: ["Saved speech"] });
expect(activeSessions.size).toBe(0);
@ -375,8 +369,7 @@ describe("live meeting summaries", () => {
await settleSummaryUpdates(fixture.updates);
expect(gatewayWorkAdmission.getActiveGatewayRootWorkCount()).toBe(0);
expect((await saved(fixture))?.overview).toBe("Operator-edited notes");
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(fixture.updates);
await advanceSummaryInterval(fixture.updates);
expect(gatewayWorkAdmission.getActiveGatewayRootWorkCount()).toBe(0);
expect(complete).toHaveBeenCalledOnce();
expect((await saved(fixture))?.overview).toBe("Operator-edited notes");
@ -386,10 +379,7 @@ describe("live meeting summaries", () => {
await vi.advanceTimersByTimeAsync(fiveMinutes);
await next.entered;
expect(complete).toHaveBeenCalledTimes(2);
const tool = createTranscriptsTool(fixture.ctx);
const manual = withPluginRuntimeRegistryScope(fixture.registry, () =>
tool.execute("manual", { action: "summarize", sessionId: "meeting" }),
);
const manual = execute(fixture, "manual", "summarize");
await vi.advanceTimersByTimeAsync(0);
expect(complete).toHaveBeenCalledTimes(2);
await fixture.source.onUtterance({ text: "Speech while the manual summary is queued" });
@ -397,8 +387,7 @@ describe("live meeting summaries", () => {
await manual;
expect(complete).toHaveBeenCalledTimes(3);
expect(await saved(fixture)).toMatchObject({ utteranceCount: 3 });
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(fixture.updates);
await advanceSummaryInterval(fixture.updates);
expect(complete).toHaveBeenCalledTimes(3);
});
@ -432,16 +421,13 @@ describe("live meeting summaries", () => {
for (let index = 0; index < 2_001; index++) {
await fixture.source.onUtterance({ text: `Speech ${index}` });
}
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(fixture.updates);
await advanceSummaryInterval(fixture.updates);
expect((await saved(fixture))?.utteranceCount).toBe(2_000);
await fixture.source.onUtterance({ text: "Newest speech" });
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(fixture.updates);
await advanceSummaryInterval(fixture.updates);
expect(complete).toHaveBeenCalledTimes(2);
expect((await saved(fixture))?.transcript.at(-1)).toBe("Newest speech");
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(fixture.updates);
await advanceSummaryInterval(fixture.updates);
expect(complete).toHaveBeenCalledTimes(2);
});
@ -500,8 +486,7 @@ describe("live meeting summaries", () => {
{ at: new Date().toISOString(), speaker: "Ada", text: "We agreed to simplify setup." },
]);
}
await vi.advanceTimersByTimeAsync(fiveMinutes);
await settleSummaryUpdates(updates);
await advanceSummaryInterval(updates);
expect(complete).toHaveBeenCalledOnce();
expect(complete.mock.calls[0]![0]).toMatchObject({ model: "utility", agentId: "research" });
expect((await store.readSummary(descriptor.session)).summary?.source).toBe("model");

View file

@ -13,14 +13,12 @@ import { transcriptSessionSelector, TranscriptsStore } from "./store.js";
const fixture = useTranscriptStatusFixture();
describe("configured transcript shutdown cleanup", () => {
it.each(
[false, true].flatMap((whenOccupied) =>
["returned-stop", "thrown-stop", "session-write", "summary-write"].map((fault) => ({
whenOccupied,
fault,
})),
),
)(
it.each([
{ whenOccupied: false, fault: "returned-stop" },
{ whenOccupied: true, fault: "thrown-stop" },
{ whenOccupied: false, fault: "session-write" },
{ whenOccupied: true, fault: "summary-write" },
])(
"retains and drains late $fault after shutdown (occupied=$whenOccupied)",
async ({ whenOccupied, fault }) => {
vi.useFakeTimers({ toFake: ["Date", "setTimeout", "clearTimeout"] });

View file

@ -8,38 +8,28 @@ beforeEach(() => {
});
afterEach(() => vi.useRealTimers());
it("uses the selected elapsed clock when wall time changes", async () => {
let elapsed = 50;
const operation = createDeferred<string>();
const result = awaitWithinDeadline(
() => operation.promise,
60,
() => elapsed,
);
it.each([
[59.5, "in time"],
[60, ABSOLUTE_DEADLINE_EXPIRED],
] as const)(
"uses the selected clock at elapsed=%s when wall time changes",
async (settledAt, expected) => {
let elapsed = 50;
const operation = createDeferred<string>();
const result = awaitWithinDeadline(
() => operation.promise,
60,
() => elapsed,
);
vi.setSystemTime(10_000);
elapsed = 59.5;
operation.resolve("in time");
vi.setSystemTime(10_000);
elapsed = settledAt;
operation.resolve("in time");
await expect(result).resolves.toBe("in time");
expect(vi.getTimerCount()).toBe(0);
});
it("rejects a result at the selected deadline before its timer runs", async () => {
let elapsed = 50;
const operation = createDeferred<string>();
const result = awaitWithinDeadline(
() => operation.promise,
60,
() => elapsed,
);
elapsed = 60;
operation.resolve("late");
await expect(result).resolves.toBe(ABSOLUTE_DEADLINE_EXPIRED);
expect(vi.getTimerCount()).toBe(0);
});
await expect(result).resolves.toBe(expected);
expect(vi.getTimerCount()).toBe(0);
},
);
it("retains wall-clock settlement when no clock is selected", async () => {
const operation = createDeferred<string>();

View file

@ -1,7 +1,5 @@
import { mkdtempSync, rmSync } from "node:fs";
import os from "node:os";
import path from "node:path";
import { afterEach, describe, expect, it, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import {
clearRuntimeAuthProfileStoreSnapshots,
replaceRuntimeAuthProfileStoreSnapshots,
@ -23,50 +21,36 @@ vi.mock("../plugins/web-search-providers.runtime.js", () => ({
}));
describe("web search OAuth and environment auto-detection", () => {
const tempDirs: string[] = [];
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
afterEach(() => {
vi.unstubAllEnvs();
resolveProviders.mockReset();
clearRuntimeAuthProfileStoreSnapshots();
for (const dir of tempDirs.splice(0)) {
rmSync(dir, { recursive: true, force: true });
}
});
it.each([
{
name: "OAuth before environment",
oauth: true,
failGrok: false,
pin: undefined,
order: ["grok"],
},
{
name: "environment fallback after OAuth failure",
oauth: true,
failGrok: true,
pin: undefined,
order: ["grok", "tavily"],
},
{
name: "environment only",
oauth: false,
failGrok: false,
pin: undefined,
order: ["tavily"],
},
{
name: "explicit environment provider",
oauth: true,
failGrok: false,
pin: "tavily",
order: ["tavily"],
},
])("preserves $name", async ({ oauth, failGrok, pin, order }) => {
])("preserves $name", async ({ oauth, pin, order }) => {
vi.stubEnv("TAVILY_API_KEY", "tavily-synthetic-key");
const agentDir = mkdtempSync(path.join(os.tmpdir(), "openclaw-web-search-order-"));
tempDirs.push(agentDir);
const agentDir = tempDirs.make("openclaw-web-search-order-");
replaceRuntimeAuthProfileStoreSnapshots([
{
agentDir,
@ -99,10 +83,7 @@ describe("web search OAuth and environment auto-detection", () => {
parameters: {},
execute: async () => {
attempts.push("grok");
if (failGrok) {
throw new Error("grok search failed");
}
return { answer: "OAuth result" };
throw new Error("grok search failed");
},
}),
}),

View file

@ -41,58 +41,14 @@ describe("executeWebSearchCandidates", () => {
expect(fallback).not.toHaveBeenCalled();
});
it("tries every fallback and retains the first error with its cause", async () => {
const cause = new Error("first transport failure");
const firstError = new Error("first provider failed", { cause });
const attempts: string[] = [];
const result = executeWebSearchCandidates({
candidates: [
candidate("first", async () => {
attempts.push("first");
throw firstError;
}),
candidate("second", async () => {
attempts.push("second");
throw new Error("second provider failed");
}),
],
args: { query: "synthetic query" },
allowFallback: true,
});
await expect(result).rejects.toMatchObject({ provider: "first", cause: firstError });
expect(firstError.cause).toBe(cause);
expect(attempts).toEqual(["first", "second"]);
});
it.each(["first rejection", null, undefined])(
"retains the first non-Error rejection: %s",
async (firstRejection) => {
await expect(
executeWebSearchCandidates({
candidates: [
candidate(
"first",
vi.fn<WebSearchProviderToolDefinition["execute"]>().mockRejectedValue(firstRejection),
),
candidate("second", async () => {
throw new Error("second provider failed");
}),
],
args: {},
allowFallback: true,
}),
).rejects.toThrow(String(firstRejection));
},
);
it.each([true, false])(
"retains chronology across structured and thrown errors (structured first: %s)",
async (structuredFirst) => {
const thrown = new Error("ordinary provider failure");
const missingKey = async () => ({ error: "missing_fixture_api_key" });
const fail = async () => {
const missingKey = vi.fn(async () => ({ error: "missing_fixture_api_key" }));
const fail = vi.fn(async () => {
throw thrown;
};
});
const result = executeWebSearchCandidates({
candidates: [
candidate("first", structuredFirst ? missingKey : fail),
@ -108,16 +64,21 @@ describe("executeWebSearchCandidates", () => {
} else {
await expect(result).rejects.toMatchObject({ provider: "first", cause: thrown });
}
expect(missingKey).toHaveBeenCalledOnce();
expect(fail).toHaveBeenCalledOnce();
},
);
it("does not treat a thrown undefined as an unavailable factory", async () => {
it("retains an undefined rejection across an unavailable factory and later failure", async () => {
const unavailable = createWebSearchTestProvider({
id: "unavailable",
pluginId: "fixture-search",
credentialPath: "plugins.entries.fixture-search.config.unavailable",
createTool: () => null,
});
const laterFailure = vi.fn(async () => {
throw new Error("second provider failed");
});
await expect(
executeWebSearchCandidates({
candidates: [
@ -126,11 +87,13 @@ describe("executeWebSearchCandidates", () => {
vi.fn<WebSearchProviderToolDefinition["execute"]>().mockRejectedValue(undefined),
),
unavailable,
candidate("second", laterFailure),
],
args: {},
allowFallback: true,
}),
).rejects.toThrow("undefined");
expect(laterFailure).toHaveBeenCalledOnce();
});
it("gives cancellation precedence over a saved failure and later cleanup error", async () => {

View file

@ -20,14 +20,11 @@ function command(reply: ReplyPayload, index: number): string {
}
describe("login command choices", () => {
it.each([
[0, "all"],
[1, "keep"],
] as const)("returns button %s as %s once", (index, value) => {
it("accepts a delivered choice only once", () => {
const choice = createLoginChoicePrompt(prompt, new AbortController().signal, "sample");
const answer = command(choice.reply, index);
const answer = command(choice.reply, 0);
expect(choice.reply.text).toContain(answer);
expect(choice.answer(answer)).toEqual({ value });
expect(choice.answer(answer)).toEqual({ value: "all" });
expect(choice.answer(answer)).toBeUndefined();
});

View file

@ -145,11 +145,13 @@ it.each([
},
});
};
const interceptWrite: typeof targetOwner.write = async (...args) => {
const committed = await targetOwner.write(...args);
await editCredential();
return committed;
};
const interceptWrite =
(write: typeof targetOwner.write): typeof write =>
async (...args) => {
const committed = await write(...args);
await editCredential();
return committed;
};
mocks.verify.mockImplementation(async (params) => {
const route = await resolveSystemAgentConfiguredRouteFromConfig(params.config);
if (!route) {
@ -197,16 +199,12 @@ it.each([
stateDir,
required: true,
configTarget: {
write: interceptWrite,
write: interceptWrite(targetOwner.write),
read: async () => {
const read = await targetOwner.read();
return {
config: read.config,
write: async (...args) => {
const committed = await read.write(...args);
await editCredential();
return committed;
},
write: interceptWrite(read.write),
};
},
},
@ -245,7 +243,6 @@ it.each([
);
it.each([
{ pending: false, superseded: false },
{ pending: true, superseded: false },
{ pending: false, superseded: true },
])(

View file

@ -168,16 +168,12 @@ describe.each(["discord", "slack"] as const)("official %s provider read boundary
const modes = [
"allowed",
"denied",
"account",
"revoked",
"result-revoked",
"legacy",
"run-revoked-direct",
"run-result-revoked-direct",
"bundled-allowed",
"bundled-revoked",
"bundled-run-revoked-direct",
"bundled-run-result-revoked-direct",
] as const;
const cases =
channel === "discord"
@ -195,7 +191,6 @@ describe.each(["discord", "slack"] as const)("official %s provider read boundary
"tool-run-claim-revoked" as const,
"tool-run-result-revoked" as const,
"tool-run-retry-revoked" as const,
"bundled-tool-allowed" as const,
"bundled-tool-run-claim-revoked" as const,
];
it.each(cases)("routes a cross-conversation read through the provider (%s)", async (mode) => {
@ -215,11 +210,6 @@ describe.each(["discord", "slack"] as const)("official %s provider read boundary
...provider,
// Status probes have provider-specific generics and are not part of message dispatch.
status: undefined,
actions: {
...provider.actions!,
readAuthorityActions:
mode === "legacy" ? undefined : provider.actions?.readAuthorityActions,
},
};
owner.registry.plugins.push(record);
owner.createApi(record, { config: {}, registrationMode: "full" }).registerChannel({ plugin });
@ -324,7 +314,7 @@ describe.each(["discord", "slack"] as const)("official %s provider read boundary
action: "read" as const,
params: { channelId: channel === "discord" ? sibling : slackTarget, limit: 1 },
accountId: "default",
requesterAccountId: mode === "account" ? "other" : "default",
requesterAccountId: "default",
conversationReadOrigin: "delegated" as const,
assertDirectAdapterHandoff: run?.assert,
toolContext: {
@ -375,19 +365,11 @@ describe.each(["discord", "slack"] as const)("official %s provider read boundary
}),
);
expect(outcome.rejected).toBe(true);
expect(outcome.message).toContain(
mode === "legacy"
? "exact current conversation"
: mode === "account"
? "current provider and account"
: mode === "denied"
? "not allowed"
: "no longer active",
);
expect(outcome.message).toContain(mode === "denied" ? "not allowed" : "no longer active");
expect(requests.filter(isContent)).toHaveLength(
mode === "result-revoked" || (runRevocation && resultRevocation) ? 1 : 0,
);
if (mode === "legacy" || mode === "account" || mode === "tool-run-preparation-revoked") {
if (mode === "tool-run-preparation-revoked") {
expect(requests).toEqual([]);
}
if (

View file

@ -16,6 +16,7 @@ import type { RunCliAgentParams } from "../src/agents/cli-runner/types.js";
import type { ScheduledToolPolicyContext } from "../src/agents/scheduled-tool-policy.js";
import type { ChannelMessageActionAdapter } from "../src/channels/plugins/types.core.js";
import { clearRuntimeConfigSnapshot, setRuntimeConfigSnapshot } from "../src/config/config.js";
import type { DiscordConfig } from "../src/config/types.discord.js";
import type { OpenClawConfig } from "../src/config/types.openclaw.js";
import { createAgentRuntimeApprovalAuthorityValidator } from "../src/gateway/agent-runtime-approval-authority.js";
import type { CronAuthenticatedChannelRequester } from "../src/gateway/cron-creator-authority-grant.types.js";
@ -413,6 +414,17 @@ describe("CLI message authority integration", () => {
setActivePluginRegistry(owner.registry);
}
function updateDiscordConfig(update: (discord: DiscordConfig) => DiscordConfig) {
cfg = {
...cfg,
channels: {
...cfg.channels,
discord: update(expectDefined(cfg.channels?.discord, "configured Discord account")),
},
};
setRuntimeConfigSnapshot(cfg, cfg);
}
afterEach(async () => {
heldRequest?.release.resolve();
heldRequest = undefined;
@ -710,17 +722,13 @@ describe("CLI message authority integration", () => {
]) {
expectSuccess(await turn.call({ ...react, ...(target ? { target } : {}) }));
}
expect(requests.filter((request) => request.method === "PUT")).toEqual([
expect.objectContaining({
path: `${discordChannelPath}/messages/${discordMessage}/reactions/%E2%9C%85/@me`,
}),
expect.objectContaining({
path: `${discordChannelPath}/messages/${discordMessage}/reactions/%E2%9C%85/@me`,
}),
expect.objectContaining({
path: `${discordChannelPath}/messages/${discordMessage}/reactions/%E2%9C%85/@me`,
}),
]);
const reactions = requests.filter((request) => request.method === "PUT");
expect(reactions).toHaveLength(3);
for (const reaction of reactions) {
expect(reaction.path).toBe(
`${discordChannelPath}/messages/${discordMessage}/reactions/%E2%9C%85/@me`,
);
}
const before = requests.length;
expectDenied(
await turn.call(
@ -773,14 +781,7 @@ describe("CLI message authority integration", () => {
it("edits the current Discord channel with the admitted sender's permission", async () => {
const turn = await createTurn("discord");
expectSuccess(
await turn.call({
action: "channel-edit",
channel: "discord",
target: `channel:${channels.discord.current}`,
topic: "A permitted channel topic",
}),
);
expectSuccess(await turn.call(channelEdit));
expect(requests).toContainEqual(
expect.objectContaining({
method: "GET",
@ -846,7 +847,6 @@ describe("CLI message authority integration", () => {
});
it.each([
{ name: "an ordinary account job", identity: {}, namedTarget: false },
{ name: "an account job with a channel-name target", identity: {}, namedTarget: true },
{
name: "an account job with forged owner and sender arguments",
@ -881,22 +881,14 @@ describe("CLI message authority integration", () => {
] as const)(
"discovers and invokes the native editor with an $origin read origin",
async ({ ownerOrigin }) => {
const discord = expectDefined(cfg.channels?.discord, "configured Discord accounts");
cfg = {
...cfg,
channels: {
...cfg.channels,
discord: {
...discord,
defaultAccount: "default",
accounts: {
default: { actions: { channels: false } },
work: { token: discordWorkToken, actions: { channels: true } },
},
},
updateDiscordConfig((discord) => ({
...discord,
defaultAccount: "default",
accounts: {
default: { actions: { channels: false } },
work: { token: discordWorkToken, actions: { channels: true } },
},
};
setRuntimeConfigSnapshot(cfg, cfg);
}));
const turn = await createTurn("discord", {
scheduledPolicy: {
...accountScheduledPolicy,
@ -1024,7 +1016,6 @@ describe("CLI message authority integration", () => {
});
it.each([
{ change: "unchanged authority", revoke: undefined },
{ change: "immediate job revocation", revoke: "source" },
{ change: "prospective action configuration", revoke: "action" },
{ change: "prospective account configuration", revoke: "account" },
@ -1038,26 +1029,18 @@ describe("CLI message authority integration", () => {
expect(acceptedChannelEdits).toBe(0);
if (revoke === "source") {
turn.revokeScheduledPermission();
} else if (revoke) {
const discord = expectDefined(cfg.channels?.discord, "configured Discord account");
cfg = {
...cfg,
channels: {
...cfg.channels,
discord: {
...discord,
...(revoke === "action"
? { actions: { ...discord.actions, channels: false } }
: {
accounts: {
...discord.accounts,
default: { ...discord.accounts?.default, enabled: false },
},
}),
},
},
};
setRuntimeConfigSnapshot(cfg, cfg);
} else {
updateDiscordConfig((discord) => ({
...discord,
...(revoke === "action"
? { actions: { ...discord.actions, channels: false } }
: {
accounts: {
...discord.accounts,
default: { ...discord.accounts?.default, enabled: false },
},
}),
}));
}
rateLimited.release();
@ -1087,33 +1070,28 @@ describe("CLI message authority integration", () => {
}
});
it.each(["operator", "native-account"] as const)(
"settles an accepted %s edit after scheduled permission is revoked without replay",
async (kind) => {
const turn = await createTurn(
"discord",
kind === "native-account"
? { scheduledPolicy: accountScheduledPolicy, channelRequester }
: { scheduledPolicy: trustedScheduledPolicy },
);
const accepted = holdProviderRequest("PATCH");
const pending = turn.call(channelEdit);
try {
await accepted.entered(pending);
expect(acceptedChannelEdits).toBe(1);
turn.revokeScheduledPermission();
accepted.release();
it("settles an accepted native-account edit after scheduled permission is revoked without replay", async () => {
const turn = await createTurn("discord", {
scheduledPolicy: accountScheduledPolicy,
channelRequester,
});
const accepted = holdProviderRequest("PATCH");
const pending = turn.call(channelEdit);
try {
await accepted.entered(pending);
expect(acceptedChannelEdits).toBe(1);
turn.revokeScheduledPermission();
accepted.release();
expectSuccess(await pending);
expectDenied(await turn.call(channelEdit), /Scheduled source permission revoked/);
expect(requests.filter((request) => request.method === "PATCH")).toHaveLength(1);
expect(acceptedChannelEdits).toBe(1);
} finally {
accepted.release();
await pending.catch(() => undefined);
}
},
);
expectSuccess(await pending);
expectDenied(await turn.call(channelEdit), /Scheduled source permission revoked/);
expect(requests.filter((request) => request.method === "PATCH")).toHaveLength(1);
expect(acceptedChannelEdits).toBe(1);
} finally {
accepted.release();
await pending.catch(() => undefined);
}
});
it("reports provider permission denial without retrying the edit", async () => {
const turn = await createTurn("discord", { scheduledPolicy: trustedScheduledPolicy });
@ -1261,27 +1239,22 @@ describe("CLI message authority integration", () => {
messageId: discordMessage,
};
const editedContent = "A permitted scheduled message edit";
const writes = [
{ action: "edit", method: "PATCH", path: discordMessagePath },
{ action: "delete", method: "DELETE", path: discordMessagePath },
{ action: "pin", method: "PUT", path: discordPinPath },
{ action: "unpin", method: "DELETE", path: discordPinPath },
] as const;
it.each(
[
{ principal: "trusted", scheduledPolicy: trustedScheduledPolicy },
{ principal: "account", scheduledPolicy: accountScheduledPolicy },
].flatMap(({ principal, scheduledPolicy }) =>
writes.map(({ action, method, path: requestPath }) => ({
principal,
scheduledPolicy,
action,
method,
path: requestPath,
})),
),
)(
it.each([
{
principal: "trusted",
scheduledPolicy: trustedScheduledPolicy,
action: "edit",
method: "PATCH",
path: discordMessagePath,
},
{
principal: "account",
scheduledPolicy: accountScheduledPolicy,
action: "unpin",
method: "DELETE",
path: discordPinPath,
},
])(
"dispatches $action through the provider route for a $principal scheduled job",
async ({ scheduledPolicy, action, method, path: requestPath }) => {
const turn = await createTurn("discord", { scheduledPolicy });
@ -1309,15 +1282,10 @@ describe("CLI message authority integration", () => {
{ action: "pin", gate: "pins" },
] as const)("retains the $gate action gate for scheduled $action", async ({ action, gate }) => {
const turn = await createTurn("discord", { scheduledPolicy: accountScheduledPolicy });
const discord = expectDefined(cfg.channels?.discord, "configured Discord account");
cfg = {
...cfg,
channels: {
...cfg.channels,
discord: { ...discord, actions: { ...discord.actions, [gate]: false } },
},
};
setRuntimeConfigSnapshot(cfg, cfg);
updateDiscordConfig((discord) => ({
...discord,
actions: { ...discord.actions, [gate]: false },
}));
expectDenied(
await turn.call({
@ -1350,29 +1318,21 @@ describe("CLI message authority integration", () => {
reason: /Discord read target channel is not allowed/,
},
])("denies $boundary before message mutation", async ({ scheduledPolicy, args, reason }) => {
const discord = expectDefined(cfg.channels?.discord, "configured Discord account");
cfg = {
...cfg,
channels: {
...cfg.channels,
discord: {
...discord,
accounts: {
...discord.accounts,
other: { token: "synthetic-other-message-provider-token" },
},
guilds: {
[discordGuild]: {
channels: {
[channels.discord.current]: { enabled: true },
[discordSibling]: { enabled: false },
},
},
updateDiscordConfig((discord) => ({
...discord,
accounts: {
...discord.accounts,
other: { token: "synthetic-other-message-provider-token" },
},
guilds: {
[discordGuild]: {
channels: {
[channels.discord.current]: { enabled: true },
[discordSibling]: { enabled: false },
},
},
},
};
setRuntimeConfigSnapshot(cfg, cfg);
}));
const turn = await createTurn("discord", { scheduledPolicy });
expectDenied(await turn.call({ ...messageTarget, action: "delete", ...args }), reason);

View file

@ -82,35 +82,30 @@ describe("registered Codex finalizer host silence contract", () => {
return result;
}
it.each([true, false, undefined])(
"honors optional authored silence independently of empty-reply permission (%s)",
async (allowed) => {
const input = await createInput();
input.terminalBase.runParams.allowEmptyAssistantReplyAsSilent = allowed;
returnBoundedText(" NO_REPLY\n");
it("honors optional authored silence with empty replies disabled", async () => {
const input = await createInput();
input.terminalBase.runParams.allowEmptyAssistantReplyAsSilent = false;
returnBoundedText(" NO_REPLY\n");
const result = await prepareTerminalWithSettledTurnFinalization(input);
const result = await prepareTerminalWithSettledTurnFinalization(input);
expect(fixture.runBounded).toHaveBeenCalledOnce();
expect(fixture.runBounded).toHaveBeenCalledWith(
expect.objectContaining({
isolation: "private-stdio",
requireNoExternalCapabilities: true,
}),
);
expect(result.finalizationOutcome).toBe("answered");
expect(result.attempt.assistantTexts).toEqual(["NO_REPLY"]);
expect(result.prepared.payloadsWithToolMedia ?? []).toEqual([]);
expect(fixture.mirror).not.toHaveBeenCalled();
},
);
expect(fixture.runBounded).toHaveBeenCalledOnce();
expect(fixture.runBounded).toHaveBeenCalledWith(
expect.objectContaining({
isolation: "private-stdio",
requireNoExternalCapabilities: true,
}),
);
expect(result.finalizationOutcome).toBe("answered");
expect(result.attempt.assistantTexts).toEqual(["NO_REPLY"]);
expect(result.prepared.payloadsWithToolMedia ?? []).toEqual([]);
expect(fixture.mirror).not.toHaveBeenCalled();
});
it.each([
{ text: "no_reply", expectation: "optional", silent: true },
{ text: "NO_REPLY", expectation: "required", silent: false },
{ text: " ", expectation: "optional", silent: false },
{ text: " ", expectation: "required", silent: false },
] as const)("distinguishes $expectation $text output", async ({ text, expectation, silent }) => {
{ text: "NO_REPLY", expectation: "required" },
{ text: " ", expectation: "optional" },
] as const)("distinguishes $expectation $text output", async ({ text, expectation }) => {
const input = await createInput();
input.terminalBase.runParams.terminalReplyExpectation = expectation;
input.terminalBase.runParams.allowEmptyAssistantReplyAsSilent = true;
@ -118,18 +113,13 @@ describe("registered Codex finalizer host silence contract", () => {
const result = await prepareTerminalWithSettledTurnFinalization(input);
expect(fixture.runBounded).toHaveBeenCalledTimes(silent ? 1 : 2);
expect(result.finalizationOutcome).toBe(silent ? "answered" : "completed-empty");
if (silent) {
expect(result.attempt.assistantTexts).toEqual([text]);
expect(result.prepared.payloadsWithToolMedia ?? []).toEqual([]);
} else {
expect(result.prepared.payloadsWithToolMedia).toEqual([
expect.objectContaining({
text: "The tool run finished, but no final summary was produced. I did not repeat any completed actions.",
}),
]);
}
expect(fixture.runBounded).toHaveBeenCalledTimes(2);
expect(result.finalizationOutcome).toBe("completed-empty");
expect(result.prepared.payloadsWithToolMedia).toEqual([
expect.objectContaining({
text: "The tool run finished, but no final summary was produced. I did not repeat any completed actions.",
}),
]);
expect(fixture.mirror).not.toHaveBeenCalled();
});

View file

@ -23,7 +23,7 @@ describe("Codex core conversation delivery", () => {
delivered: true,
success: true,
},
...(["sent", "queued", "suppressed", "unknown"] as const).map((status) => ({
...(["sent", "queued"] as const).map((status) => ({
name: `send ${status} with a platform ID`,
toolName: "conversations_send" as const,
receipt: { ...destination, status, messageId: "outbound-1" },
@ -55,14 +55,14 @@ describe("Codex core conversation delivery", () => {
delivered: true,
success: false,
},
...(["sent", "queued", "suppressed", "unknown"] as const).map((status) => ({
...(["sent", "queued"] as const).map((status) => ({
name: `turn ${status} with a correlation error`,
toolName: "conversations_turn" as const,
receipt: {
...destination,
status,
messageId: "outbound-1",
correlationPersisted: status === "sent" || status === "queued",
correlationPersisted: true,
error: "No process-local reply waiter remains.",
},
delivered: status === "sent",

View file

@ -122,7 +122,7 @@ fs.writeFileSync('dist/marker', process.argv[2] || 'new');
);
writeShim(
"pnpm",
'if [ "$1" = --version ]; then echo 12.4.0; exit 0; fi\necho "pnpm $*" >> "$UPDATE_TEST_LOG"\nif [ "$1" = build ]; then mkdir -p dist; echo new > dist/marker; exit "${UPDATE_TEST_BUILD_EXIT:-0}"; fi',
'if [ "$1" = --version ]; then echo 12.4.0; exit 0; fi\necho "pnpm $*" >> "$UPDATE_TEST_LOG"',
);
writeShim("openclaw", 'echo "openclaw $*" >> "$UPDATE_TEST_LOG"');
});
@ -275,9 +275,7 @@ await withDistArtifactOwnership(process.cwd(), async () => {
});
it.each([
"success",
"reset-refused",
"verification-failed",
"restart-failed",
"tracked-runtime-deleted",
"tracked-runtime-changed",
@ -340,8 +338,7 @@ process.exitCode = await runUpdateGatewayBuild(...process.argv.slice(2));
"stop-ok",
[
'echo "stop:$(git rev-parse HEAD):$(cat node_modules/serving-marker)" >> "$UPDATE_TEST_LOG"',
...(mode === "verification-failed" ||
mode === "tracked-runtime-deleted" ||
...(mode === "tracked-runtime-deleted" ||
mode === "tracked-runtime-rollback" ||
mode === "restore-branch-drift"
? ['node "$UPDATE_TEST_BIN/corrupt-staged.cjs"']
@ -377,10 +374,9 @@ fs.writeFileSync(path.join(names[0],'candidate','.buildstamp'), JSON.stringify({
UPDATE_TEST_RESTART_EXIT: mode === "restart-failed" ? "29" : "0",
});
expect(result.status, result.stdout + result.stderr).toBe(
mode === "success" || mode === "tracked-runtime-changed" ? 0 : 1,
mode === "tracked-runtime-changed" ? 0 : 1,
);
const published =
mode === "success" || mode === "restart-failed" || mode === "tracked-runtime-changed";
const published = mode === "restart-failed" || mode === "tracked-runtime-changed";
const expectedHead = published
? git(path.join(scratch, "seed"), "rev-parse", "HEAD")
: original;
@ -701,14 +697,4 @@ fs.writeFileSync(path.join(names[0],'candidate','.buildstamp'), JSON.stringify({
expect(calls()).not.toContain("pnpm build");
},
);
it.each([0, 17])("keeps exact-empty restart as manual lifecycle (build exit %s)", (exit) => {
const result = runUpdater({
OPENCLAW_UPDATE_RESTART_CMD: "",
UPDATE_TEST_BUILD_EXIT: String(exit),
});
expect(result.status, result.stderr).toBe(exit);
expect(calls()).toContain("pnpm build");
expect(calls().some((call) => call.startsWith("openclaw "))).toBe(false);
});
});

View file

@ -19,6 +19,14 @@ const wsServerUrl = pathToFileURL(
).href;
const sourceSha = "1".repeat(40);
const hash = (bytes: Buffer | string) => createHash("sha256").update(bytes).digest("hex");
const writeJson = (file: string, value: unknown) => fs.writeFile(file, JSON.stringify(value));
async function readEvents(file: string) {
return (await fs.readFile(file, "utf8"))
.trim()
.split("\n")
.map((line) => JSON.parse(line));
}
async function fixture(
mode: "healthy" | "rpc-error" | "established-rpc-error" | "lingering" | "hold-health" = "healthy",
@ -28,11 +36,8 @@ async function fixture(
const installRoot = path.join(root, "install");
const packageRoot = path.join(installRoot, "node_modules", "openclaw");
await fs.mkdir(path.join(packageRoot, "dist"), { recursive: true });
await fs.writeFile(
path.join(packageRoot, "package.json"),
JSON.stringify({ name: "openclaw", version: "1.0.0" }),
);
await fs.writeFile(path.join(packageRoot, "dist", "build-info.json"), JSON.stringify({ commit }));
await writeJson(path.join(packageRoot, "package.json"), { name: "openclaw", version: "1.0.0" });
await writeJson(path.join(packageRoot, "dist", "build-info.json"), { commit });
const events = path.join(root, "events.jsonl");
await fs.writeFile(
path.join(packageRoot, "openclaw.mjs"),
@ -108,57 +113,51 @@ if (process.argv.includes("--descendant")) {
);
const tarball = path.join(root, "fixture.tgz");
await fs.writeFile(tarball, `synthetic package ${commit}; installation is fixture-owned`);
await fs.writeFile(
path.join(installRoot, "package-lock.json"),
JSON.stringify({
name: "install",
lockfileVersion: 3,
requires: true,
packages: {
"": { dependencies: { openclaw: "file:../fixture.tgz" } },
"node_modules/openclaw": {
version: "1.0.0",
resolved: "file:../fixture.tgz",
integrity: `sha512-${createHash("sha512")
.update(await fs.readFile(tarball))
.digest("base64")}`,
},
"node_modules/fixture-dependency": {
version: "1.0.0",
resolved: "https://registry.npmjs.org/fixture-dependency/-/fixture-dependency-1.0.0.tgz",
integrity: "sha512-synthetic-dependency",
optional: true,
os: ["win32"],
},
await writeJson(path.join(installRoot, "package-lock.json"), {
name: "install",
lockfileVersion: 3,
requires: true,
packages: {
"": { dependencies: { openclaw: "file:../fixture.tgz" } },
"node_modules/openclaw": {
version: "1.0.0",
resolved: "file:../fixture.tgz",
integrity: `sha512-${createHash("sha512")
.update(await fs.readFile(tarball))
.digest("base64")}`,
},
}),
);
"node_modules/fixture-dependency": {
version: "1.0.0",
resolved: "https://registry.npmjs.org/fixture-dependency/-/fixture-dependency-1.0.0.tgz",
integrity: "sha512-synthetic-dependency",
optional: true,
os: ["win32"],
},
},
});
const input = path.join(root, "input.json");
const stateRoot = path.join(root, "state");
await fs.writeFile(
input,
JSON.stringify({
sourceSha: commit,
toolingSha: sourceSha,
tarball,
installRoot,
stateRoot,
candidate: {
name: "openclaw",
packageSourceSha: commit,
version: "1.0.0",
sha256: hash(await fs.readFile(tarball)),
},
runtime: { version: process.version, sha256: hash(await fs.readFile(process.execPath)) },
artifact: {
id: 1,
runId: 1,
runAttempt: 1,
workflowSha: sourceSha,
digest: `sha256:${"2".repeat(64)}`,
},
}),
);
await writeJson(input, {
sourceSha: commit,
toolingSha: sourceSha,
tarball,
installRoot,
stateRoot,
candidate: {
name: "openclaw",
packageSourceSha: commit,
version: "1.0.0",
sha256: hash(await fs.readFile(tarball)),
},
runtime: { version: process.version, sha256: hash(await fs.readFile(process.execPath)) },
artifact: {
id: 1,
runId: 1,
runAttempt: 1,
workflowSha: sourceSha,
digest: `sha256:${"2".repeat(64)}`,
},
});
const output = path.join(root, "report.json");
return { root, packageRoot, input, output, stateRoot, events };
}
@ -166,13 +165,10 @@ if (process.argv.includes("--descendant")) {
async function comparisonFixture(mode: Parameters<typeof fixture>[0] = "healthy") {
const baseline = await fixture();
const candidate = await fixture(mode, "3".repeat(40));
await fs.writeFile(
baseline.input,
JSON.stringify({
...JSON.parse(await fs.readFile(baseline.input, "utf8")),
comparison: JSON.parse(await fs.readFile(candidate.input, "utf8")),
}),
);
await writeJson(baseline.input, {
...JSON.parse(await fs.readFile(baseline.input, "utf8")),
comparison: JSON.parse(await fs.readFile(candidate.input, "utf8")),
});
return { baseline, candidate };
}
@ -241,7 +237,7 @@ describe("installed Gateway startup benchmark entry", () => {
droppedPhases: 0,
observerErrors: [],
};
await fs.writeFile(capture.attachmentPath, JSON.stringify(attachment));
await writeJson(capture.attachmentPath, attachment);
const file = path.join(
capture.directory,
`CPU.20260918.000000.${process.pid}.0.001.cpuprofile`,
@ -287,13 +283,13 @@ describe("installed Gateway startup benchmark entry", () => {
phases: [attachment.phases[0], { ...attachment.phases[1], performanceMs: 100.5 }],
},
]) {
await fs.writeFile(capture.attachmentPath, JSON.stringify(invalid));
await writeJson(capture.attachmentPath, invalid);
await expect(collectInstalledCpuProfile(capture, process.pid)).rejects.toThrow(
/monotonic|performance clock/u,
);
}
await fs.writeFile(capture.attachmentPath, JSON.stringify(attachment));
await fs.writeFile(file, JSON.stringify({ ...profile, samples: [1, 3, 1] }));
await writeJson(capture.attachmentPath, attachment);
await writeJson(file, { ...profile, samples: [1, 3, 1] });
await expect(collectInstalledCpuProfile(capture, process.pid)).rejects.toThrow(
"CPU sample references an unknown node",
);
@ -321,10 +317,7 @@ describe("installed Gateway startup benchmark entry", () => {
"fresh",
"established",
]);
const events = (await fs.readFile(target.events, "utf8"))
.trim()
.split("\n")
.map((line) => JSON.parse(line));
const events = await readEvents(target.events);
const starts = events.filter((event) => event.type === "start");
expect(starts.map((event) => event.index)).toEqual([0, 1]);
expect(starts.map((event) => event.home)).toEqual([target.stateRoot, target.stateRoot]);
@ -435,10 +428,7 @@ describe("installed Gateway startup benchmark entry", () => {
expect(samples.map((sample: { armIndex: number }) => sample.armIndex)).toEqual([
0, 1, 2, 3, 4, 5, 6, 7, 8,
]);
const events = (await fs.readFile(target.events, "utf8"))
.trim()
.split("\n")
.map((line) => JSON.parse(line));
const events = await readEvents(target.events);
expect(events.filter((event) => event.type === "start").map((event) => event.index)).toEqual([
0, 1, 2, 3, 4, 5, 6, 7, 8,
]);
@ -473,7 +463,7 @@ describe("installed Gateway startup benchmark entry", () => {
const lock = JSON.parse(await fs.readFile(lockPath, "utf8"));
lock.packages["node_modules/fixture-dependency"][field] =
field === "optional" ? false : "changed";
await fs.writeFile(lockPath, JSON.stringify(lock));
await writeJson(lockPath, lock);
const result = await runFixture(baseline, signal);
expect(result.status, result.stderr).toBe(1);
expect(result.stderr).toContain("Installed dependency records differ");
@ -507,10 +497,7 @@ describe("installed Gateway startup benchmark entry", () => {
});
expect(report.after).toEqual(report.before);
expect(report.samples).toHaveLength(9);
const events = (await fs.readFile(target.events, "utf8"))
.trim()
.split("\n")
.map((line) => JSON.parse(line));
const events = await readEvents(target.events);
expect(events.filter((event) => event.type === "start").map((event) => event.index)).toEqual([
0, 1, 2, 3, 4, 5, 6, 7, 8,
]);
@ -554,9 +541,6 @@ describe("installed Gateway startup benchmark entry", () => {
const result = await runFixture(target, signal);
expect(result.status, result.stderr).toBe(1);
const report = JSON.parse(await fs.readFile(target.output, "utf8"));
expect(report.samples[0], JSON.stringify({ stderr: result.stderr, report })).toMatchObject({
observations: { health: { response: { ok: false } } },
});
expect(report.outcome).toBe("failed");
expect(report.samples.map((sample: { outcome: string }) => sample.outcome)).toEqual([
"failed",
@ -656,11 +640,9 @@ describe("installed Gateway startup benchmark entry", () => {
expect(report.outerSettlement.beforeCleanup).toBe("indeterminate");
expect(report.outcome).toBe("failed");
expect(report.establishedReadySummary).toBeNull();
const lingering = (await fs.readFile(target.events, "utf8"))
.trim()
.split("\n")
.map((line) => JSON.parse(line))
.filter((event) => event.type === "lingering");
const lingering = (await readEvents(target.events)).filter(
(event) => event.type === "lingering",
);
expect(lingering).toHaveLength(9);
for (const { childPid } of lingering) {
expect(isProcessAlive(childPid)).toBe(false);

View file

@ -221,15 +221,6 @@ describe("changed resolved dependencies", () => {
() => write("package.json", JSON.stringify({ ...manifest, openclaw: {} })),
"execution or resolution",
],
[
"manifest execution",
() =>
write(
"package.json",
JSON.stringify({ ...manifest, scripts: { test: "different-runner" } }),
),
"execution or resolution",
],
[
"unlocked manifest",
() =>

View file

@ -8,7 +8,7 @@ const tempDirs = useAutoCleanupTempDirTracker(afterEach);
afterEach(() => vi.restoreAllMocks());
describe("limit reporting", () => {
it.each([{}, { CI: "1" }, { GITHUB_ACTIONS: "false" }])(
it.each([{ CI: "1" }, { GITHUB_ACTIONS: "false" }])(
"keeps local checks blocking with %j",
(env) => {
const error = vi.spyOn(console, "error").mockImplementation(() => {});

View file

@ -3,7 +3,6 @@ import path from "node:path";
import { Type } from "typebox";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { createCodexDynamicToolBridge } from "../extensions/codex/test-api.js";
import { getBeforeToolCallDiagnosticOptions } from "../src/agents/before-tool-call-metadata.js";
import { asToolParamsRecord, type AnyAgentTool } from "../src/agents/tools/common.js";
import {
onDiagnosticEvent,
@ -134,11 +133,6 @@ describe("persistent skill usage through registered Codex dynamic tools", () =>
: {}),
},
});
expect(bridge.telemetry.quarantinedTools).toEqual([]);
expect(bridge.availableTools.map((tool) => tool.name)).toEqual([toolName]);
for (const tool of bridge.availableTools) {
expect(getBeforeToolCallDiagnosticOptions(tool)?.emitDiagnostics).toBe(false);
}
const call = (callId: string, filePath = skillFile) =>
bridge.handleToolCall({
threadId: "skill-usage-thread",
@ -169,6 +163,14 @@ describe("persistent skill usage through registered Codex dynamic tools", () =>
};
}
async function expectNoUsage() {
await waitForDiagnosticEventsDrained();
await unregisterUsage();
expect(usageRows()).toEqual([]);
expect(consumeRunSkillUsage(runId)).toEqual([]);
expect(trustedEvents).toEqual([]);
}
it.each([false, true])(
"counts repeated successful reads with process diagnostics=%s",
async (enabled) => {
@ -200,24 +202,17 @@ describe("persistent skill usage through registered Codex dynamic tools", () =>
},
);
it.each(["error", "failed", "blocked", "cancelled", "timed_out"])(
"does not count a structured %s read",
async (status) => {
const { call, execute } = createBridge({
execute: async () => ({
content: [{ type: "text", text: "Read did not complete" }],
details: { status },
}),
});
expect(await call("failed-read")).toMatchObject({ success: false });
await waitForDiagnosticEventsDrained();
await unregisterUsage();
expect(execute).toHaveBeenCalledOnce();
expect(consumeRunSkillUsage(runId)).toEqual([]);
expect(usageRows()).toEqual([]);
expect(trustedEvents).toEqual([]);
},
);
it.each(["error", "blocked"])("does not count a structured %s read", async (status) => {
const { call, execute } = createBridge({
execute: async () => ({
content: [{ type: "text", text: "Read did not complete" }],
details: { status },
}),
});
expect(await call("failed-read")).toMatchObject({ success: false });
expect(execute).toHaveBeenCalledOnce();
await expectNoUsage();
});
it("does not count a thrown read", async () => {
const { call } = createBridge({
@ -226,11 +221,7 @@ describe("persistent skill usage through registered Codex dynamic tools", () =>
},
});
expect(await call("thrown-read")).toMatchObject({ success: false });
await waitForDiagnosticEventsDrained();
await unregisterUsage();
expect(usageRows()).toEqual([]);
expect(consumeRunSkillUsage(runId)).toEqual([]);
expect(trustedEvents).toEqual([]);
await expectNoUsage();
});
it("does not count a read blocked before execution", async () => {
@ -244,27 +235,16 @@ describe("persistent skill usage through registered Codex dynamic tools", () =>
);
const { call, execute } = createBridge();
expect(await call("blocked-read")).toMatchObject({ success: false, executionStarted: false });
await waitForDiagnosticEventsDrained();
await unregisterUsage();
expect(execute).not.toHaveBeenCalled();
expect(usageRows()).toEqual([]);
expect(consumeRunSkillUsage(runId)).toEqual([]);
expect(trustedEvents).toEqual([]);
await expectNoUsage();
});
it.each(["skills/unknown/SKILL.md", "README.md"])(
"does not count reading %s outside the skill snapshot",
async (filePath) => {
const otherFile = await testState.writeText(filePath, "Other file\n");
const { call } = createBridge();
expect(await call("other-read", otherFile)).toMatchObject({ success: true });
await waitForDiagnosticEventsDrained();
await unregisterUsage();
expect(usageRows()).toEqual([]);
expect(consumeRunSkillUsage(runId)).toEqual([]);
expect(trustedEvents).toEqual([]);
},
);
it("does not count reading a skill outside the snapshot", async () => {
const otherFile = await testState.writeText("skills/unknown/SKILL.md", "Other file\n");
const { call } = createBridge();
expect(await call("other-read", otherFile)).toMatchObject({ success: true });
await expectNoUsage();
});
it("preserves explicit tool-dispatched skill command activation", async () => {
const { call } = createBridge({ command: true });

View file

@ -292,7 +292,6 @@ async function withDownloadFixture(
token: string;
bodyStarted: Promise<void>;
binaryRequests: () => RequestRecord[];
assertPreserved: () => Promise<void>;
}) => Promise<void>,
) {
await withOpenClawTestState(
@ -714,9 +713,6 @@ async function withDownloadFixture(
token,
binaryRequests,
bodyStarted: bodyStarted.promise,
assertPreserved: async () => {
await expect(fs.readFile(preservedPath, "utf8")).resolves.toBe(PRESERVED);
},
});
await expect(fs.readFile(preservedPath, "utf8")).resolves.toBe(PRESERVED);
expect(fixtureErrors).toEqual([]);
@ -776,7 +772,6 @@ describe("registered Slack attachment downloads", () => {
).toContain(`file=${FILE_ID}`);
expect(fixture.binaryRequests()).toHaveLength(1);
expect(fixture.binaryRequests()[0]?.authorization).toBe(`Bearer ${fixture.token}`);
await fixture.assertPreserved();
});
},
);
@ -801,7 +796,6 @@ describe("registered Slack attachment downloads", () => {
expect(
fixture.requests.every((request) => request.authorization === `Bearer ${fixture.token}`),
).toBe(true);
await fixture.assertPreserved();
},
);
});
@ -824,7 +818,6 @@ describe("registered Slack attachment downloads", () => {
options: { legacy: true },
error: /exact current conversation/i,
},
{ name: "missing fileId", options: { params: { fileId: undefined } }, error: /fileId/i },
] satisfies Array<{ name: string; options: FixtureOptions; error: RegExp }>)(
"rejects $name before provider I/O",
async ({ options, error }) => {
@ -838,34 +831,23 @@ describe("registered Slack attachment downloads", () => {
);
it.each([
{ name: "missing channel evidence", file: { channels: undefined }, params: {}, allowed: false },
{ name: "missing channel evidence", file: { channels: undefined }, params: {} },
{
name: "different shared channel",
file: { channels: ["C1111111111"] },
params: {},
allowed: false,
},
{
name: "different explicit thread",
file: { shares: { private: { [TARGET_CHANNEL]: [{ ts: THREAD }] } } },
params: { threadId: "1111111111.111111" },
allowed: false,
},
{
name: "matching explicit thread",
file: { shares: { private: { [TARGET_CHANNEL]: [{ ts: THREAD }] } } },
params: { threadId: THREAD },
allowed: true,
},
])("requires positive fresh file-share proof: $name", async ({ file, params, allowed }) => {
])("requires positive fresh file-share proof: $name", async ({ file, params }) => {
await withDownloadFixture({ file, params }, async (fixture) => {
const result = await fixture.invoke();
expect(result).toMatchObject({ details: { ok: allowed } });
expect(fixture.binaryRequests()).toHaveLength(allowed ? 1 : 0);
if (!allowed) {
expect(await filesUnder(fixture.mediaDir)).toEqual([fixture.preservedPath]);
}
await fixture.assertPreserved();
expect(result).toMatchObject({ details: { ok: false } });
expect(fixture.binaryRequests()).toHaveLength(0);
expect(await filesUnder(fixture.mediaDir)).toEqual([fixture.preservedPath]);
});
});
@ -964,7 +946,6 @@ describe("registered Slack attachment downloads", () => {
expect(await fixture.invoke()).toMatchObject({ details: { ok: false } });
expect(fixture.binaryRequests().length).toBeGreaterThan(0);
expect(await filesUnder(fixture.mediaDir)).toEqual([fixture.preservedPath]);
await fixture.assertPreserved();
});
},
);
@ -1034,16 +1015,10 @@ describe("Slack download authority through artifact completion", () => {
"rejects revoked output and retains only preexisting files after $name",
async (options) => {
await withDownloadFixture(options, async (fixture) => {
const outcome = await fixture.invoke().then(
(result) => ({ result, error: undefined }),
(error: unknown) => ({ result: undefined, error }),
);
await expect(fixture.invoke()).rejects.toBeInstanceOf(Error);
expect(fixture.phases.has(options.phase)).toBe(true);
expect(fixture.binaryRequests()).toHaveLength(options.requests);
expect(await filesUnder(fixture.mediaDir)).toEqual([fixture.preservedPath]);
expect(outcome.result).toBeUndefined();
expect(outcome.error).toBeInstanceOf(Error);
await fixture.assertPreserved();
});
},
);
@ -1077,7 +1052,6 @@ describe("Slack download authority through artifact completion", () => {
expect(remaining).toHaveLength(1);
await expect(fs.readFile(remaining[0]!, "utf8")).resolves.toBe("replacement must remain");
await expect(fs.readFile(fixture.displacedPath)).resolves.toEqual(fixture.body);
await fixture.assertPreserved();
});
});
});

View file

@ -257,33 +257,6 @@ describe("resolveSubagentCompletionOrigin", () => {
expected: { channel: "discord", accountId: "acct-1", ...expected },
spawnMode: "session" as const,
})),
{
name: "resolves bound completion delivery from the requester session, not the child session",
bindings: [
{
channel: "discord",
accountId: "bot-alpha",
targetSessionKey: "agent:worker:subagent:child",
targetKind: "subagent" as const,
conversationId: "child-window",
},
{
channel: "discord",
accountId: "acct-1",
targetSessionKey: "agent:main:main",
targetKind: "session" as const,
conversationId: "parent-main",
},
],
childSessionKey: "agent:worker:subagent:child",
requesterOrigin: {
channel: "discord",
accountId: "acct-1",
to: "channel:parent-main",
},
expected: { channel: "discord", accountId: "acct-1", to: "channel:parent-main" },
spawnMode: "session" as const,
},
{
name: "prefers requester binding when child and requester share the same channel and accountId",
bindings: [
@ -447,7 +420,6 @@ describe("completion delivery route fallback", () => {
to: "channel:requester-room",
threadId: "requester-thread",
};
const retargetedOrigin = { channel: "slack", accountId: "acct-1", to: "channel:bound-room" };
const topicOrigin = {
channel: "telegram",
accountId: "bot-1",
@ -502,18 +474,13 @@ describe("completion delivery route fallback", () => {
scenario,
),
),
...["same", "different-port", "shorter"].map((destination) => {
...["same", "different-port"].map((destination) => {
const opaqueOrigin = {
channel: "matrix",
to: "room:!example:topic:100",
threadId: "$reply",
};
const to =
destination === "same"
? opaqueOrigin.to
: destination === "shorter"
? "room:!example"
: "room:!example:topic:101";
const to = destination === "same" ? opaqueOrigin.to : "room:!example:topic:101";
return {
name: `an opaque Matrix room with ${destination} identity`,
requesterSessionOrigin: opaqueOrigin,
@ -549,16 +516,6 @@ describe("completion delivery route fallback", () => {
expected: { channel: "scoped-chat", to: "alias-upper" },
expectedChatType: "direct",
},
{
name: "a retargeted completion override",
completionDirectOrigin: retargetedOrigin,
expected: retargetedOrigin,
},
{
name: "a retargeted direct origin",
directOrigin: retargetedOrigin,
expected: retargetedOrigin,
},
{
name: "a partial completion origin",
completionDirectOrigin: { channel: "slack" },
@ -569,16 +526,6 @@ describe("completion delivery route fallback", () => {
completionDirectOrigin: { to: requesterSessionOrigin.to },
expected: requesterSessionOrigin,
},
{
name: "a to-only completion origin for a different target",
completionDirectOrigin: { to: retargetedOrigin.to },
expected: retargetedOrigin,
},
{
name: "an explicit completion thread",
completionDirectOrigin: { ...retargetedOrigin, threadId: "bound-thread" },
expected: { ...retargetedOrigin, threadId: "bound-thread" },
},
])(
"uses $name for completion and generated-media delivery",
({ name: _name, expected, expectedChatType = "channel", expectedRoute, ...override }) => {