fix: stabilize chat interactions and async tests (#153277)

* fix: stabilize task reads, chat navigation, and async tests

Treat verified publication supersession as an intentional nonpublication while preserving genuine errors and replacement guards. Handle rail navigation and resize in the same observation. Replace scheduler-sensitive test assertions with explicit readiness, controlled deadlines, and cleanup that joins accepted work.

* fix(ui): keep inset focus rings clear of disclosure labels

Reserve inline space for the existing inset focus ring and let dense tool summaries inherit it while retaining their block padding. Preserve the renderer-containment repair now on main.

* fix(tasks): retain publication ownership across supersession

* fix: stabilize task reads, chat navigation, and async tests

Treat verified publication supersession as an intentional nonpublication while preserving genuine errors and replacement guards. Handle rail navigation and resize in the same observation. Replace scheduler-sensitive test assertions with explicit readiness, controlled deadlines, and cleanup that joins accepted work.

* fix(ui): keep inset focus rings clear of disclosure labels

Reserve inline space for the existing inset focus ring and let dense tool summaries inherit it while retaining their block padding. Preserve the renderer-containment repair now on main.

* fix(tasks): retain publication ownership across supersession

Keep task and flow owner retirement distinct from a superseded row. Settle publication and retire pending mutations synchronously before observer microtasks wake readers.

* fix(tasks): recheck flow owner after cancellation cleanup

* fix(ui): refresh stale rail resize targets at the measured end

* test: stabilize rail resize setup and simplify Swarm diagnostics

Establish the rail clipping geometry through real wheel input before composer growth, preserving all navigation and recovery assertions. Read bounded Gateway event history instead of overriding WebSocket dispatch or adding a lint suppression.
This commit is contained in:
Peter Steinberger 2026-09-20 01:44:18 -07:00 • committed by GitHub
parent 1bb62ea57e
commit b56478b368
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
19 changed files with 1429 additions and 648 deletions

View file

@ -2,6 +2,7 @@ import fs from "node:fs/promises";
import http from "node:http"; import http from "node:http";
import os from "node:os"; import os from "node:os";
import path from "node:path"; import path from "node:path";
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
import { captureEnv, withServer } from "openclaw/plugin-sdk/test-env"; import { captureEnv, withServer } from "openclaw/plugin-sdk/test-env";
import { afterAll, beforeAll, describe, expect, it } from "vitest"; import { afterAll, beforeAll, describe, expect, it } from "vitest";
import { saveMediaStreamWithIdleTimeout } from "./media-chunk-idle.js"; import { saveMediaStreamWithIdleTimeout } from "./media-chunk-idle.js";
@ -53,16 +54,12 @@ describe("saveMediaStreamWithIdleTimeout", () => {
}); });
it("times out a stalled SDK-style HTTP stream and closes its connection", async () => { it("times out a stalled SDK-style HTTP stream and closes its connection", async () => {
let serverSawClose = false; const socketClosed = createDeferred<void>();
await withServer( await withServer(
(req, res) => { (req, res) => {
res.writeHead(200, { "content-type": "image/jpeg", "content-length": "1048576" }); res.writeHead(200, { "content-type": "image/jpeg", "content-length": "1048576" });
res.flushHeaders(); res.flushHeaders();
const markClose = () => { req.socket.once("close", () => socketClosed.resolve());
serverSawClose = true;
};
req.on("close", markClose);
res.on("close", markClose);
}, },
async (baseUrl) => { async (baseUrl) => {
const stalled = await getHttpReadable(`${baseUrl}/media`); const stalled = await getHttpReadable(`${baseUrl}/media`);
@ -73,10 +70,7 @@ describe("saveMediaStreamWithIdleTimeout", () => {
chunkTimeoutMs: 50, chunkTimeoutMs: 50,
}); });
expect(stalled.destroyed).toBe(true); expect(stalled.destroyed).toBe(true);
await new Promise<void>((resolve) => { await socketClosed.promise;
setTimeout(resolve, 20);
});
expect(serverSawClose).toBe(true);
}, },
); );
}); });

View file

@ -1,5 +1,6 @@
import { expectDefined } from "@openclaw/normalization-core"; import { expectDefined } from "@openclaw/normalization-core";
import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
import * as runtimeEnv from "openclaw/plugin-sdk/runtime-env";
// Matrix tests cover client plugin behavior. // Matrix tests cover client plugin behavior.
import { createRequireRecord } from "openclaw/plugin-sdk/test-fixtures"; import { createRequireRecord } from "openclaw/plugin-sdk/test-fixtures";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
@ -840,6 +841,14 @@ describe("resolveMatrixAuth", () => {
}); });
it("stops waiting on whoami retry backoff when startup backfill is aborted", async () => { it("stops waiting on whoami retry backoff when startup backfill is aborted", async () => {
vi.useFakeTimers();
const retryStarted = createDeferred<void>();
const sleepWithAbort = runtimeEnv.sleepWithAbort;
vi.spyOn(runtimeEnv, "sleepWithAbort").mockImplementation((...args) => {
const sleeping = sleepWithAbort(...args);
retryStarted.resolve();
return sleeping;
});
matrixDoRequestMock.mockRejectedValueOnce( matrixDoRequestMock.mockRejectedValueOnce(
Object.assign(new TypeError("fetch failed"), { Object.assign(new TypeError("fetch failed"), {
cause: Object.assign(new Error("read ECONNRESET"), { cause: Object.assign(new Error("read ECONNRESET"), {
@ -848,7 +857,6 @@ describe("resolveMatrixAuth", () => {
}), }),
); );
const abortController = new AbortController(); const abortController = new AbortController();
const startedAt = Date.now();
const backfillPromise = backfillMatrixAuthDeviceIdAfterStartup({ const backfillPromise = backfillMatrixAuthDeviceIdAfterStartup({
auth: { auth: {
accountId: "default", accountId: "default",
@ -860,17 +868,21 @@ describe("resolveMatrixAuth", () => {
abortSignal: abortController.signal, abortSignal: abortController.signal,
}); });
await vi.waitFor(() => { try {
expect(matrixDoRequestMock).toHaveBeenCalledTimes(1); await retryStarted.promise;
}); expect(vi.getTimerCount()).toBe(1);
abortController.abort(); abortController.abort();
// The first retry backoff starts at 250ms; an honored abort returns long before it elapses. // Cancellation must settle the backoff without advancing its clock.
await expect(backfillPromise).resolves.toBeUndefined(); await expect(backfillPromise).resolves.toBeUndefined();
expect(Date.now() - startedAt).toBeLessThan(200); expect(vi.getTimerCount()).toBe(0);
expect(matrixDoRequestMock).toHaveBeenCalledTimes(1); expect(matrixDoRequestMock).toHaveBeenCalledTimes(1);
expect(repairCurrentTokenStorageMetaDeviceIdMock).not.toHaveBeenCalled(); expect(repairCurrentTokenStorageMetaDeviceIdMock).not.toHaveBeenCalled();
expect(saveBackfilledMatrixDeviceIdMock).not.toHaveBeenCalled(); expect(saveBackfilledMatrixDeviceIdMock).not.toHaveBeenCalled();
} finally {
abortController.abort();
vi.useRealTimers();
}
}); });
it("resolves configured accessToken SecretRefs during Matrix auth", async () => { it("resolves configured accessToken SecretRefs during Matrix auth", async () => {

View file

@ -0,0 +1,194 @@
import fs from "node:fs/promises";
import path from "node:path";
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
import { isRecord } from "openclaw/plugin-sdk/string-coerce-runtime";
import { afterEach, describe, expect, it } from "vitest";
import { createQaBusState } from "./bus-state.js";
import { readQaScenarioById } from "./scenario-catalog.js";
import { runLoadedScenarioFlow } from "./scenario-flow-runner.test-support.js";
import { createTempDirHarness } from "./temp-dir.test-helper.js";
const temporary = createTempDirHarness();
afterEach(() => temporary.cleanup());
async function startPublicCase(acknowledgments: number) {
const scenario = readQaScenarioById("subagent-completion-direct-fallback");
const step = scenario.execution.flow!.steps[0]!;
const guarded = step.actions.find((action) => {
if (!isRecord(action)) {
throw new Error("invalid terminal flow action");
}
return "try" in action;
});
if (!isRecord(guarded) || !isRecord(guarded.try) || !Array.isArray(guarded.try.actions)) {
throw new Error("expected guarded terminal flow actions");
}
const actions: unknown[] = guarded.try.actions;
const publicIndex = actions.findIndex((action) => {
if (!isRecord(action)) {
throw new Error("invalid guarded terminal flow action");
}
return "forEach" in action;
});
if (publicIndex < 0) {
throw new Error("expected public terminal cases");
}
const state = createQaBusState();
const observed = createDeferred<unknown>();
const releaseParent = createDeferred<void>();
const parentSent = createDeferred<void>();
const marker = "QA-SUBAGENT-TERMINAL-FALLBACK-OK";
const task = {
taskId: "child-task",
title: "qa-terminal-fallback",
status: "completed",
deliveryStatus: "delivered",
sessionKey: "parent",
childSessionKey: "child",
runId: "run",
};
const requests = [{ plannedToolName: "sessions_spawn", plannedToolArgs: { label: task.title } }];
let parentSend: Promise<void> | undefined;
const result = runLoadedScenarioFlow(scenario.id, {
state,
// Execute the shipped public-case actions and assertions, stopping before
// the independent private/restart scenarios rather than emulating them.
flow: {
steps: [
{
name: step.name,
actions: [
...step.actions.slice(0, step.actions.indexOf(guarded)),
...actions.slice(0, publicIndex + 1),
],
},
],
},
api: {
fs,
path,
config: {
...scenario.execution.config,
cases: [{ name: "fallback", marker, expectedSendCount: 1 }],
},
env: {
providerMode: "mock-openai",
outputDir: await temporary.makeTempDir("terminal-ack-"),
mock: { baseUrl: "http://mock.invalid" },
gateway: {
call: async (method: string) => {
if (method === "tasks.list") {
return { tasks: [task] };
}
if (method === "chat.history") {
return {
messages: [
{
role: "assistant",
provider: "openclaw",
model: "delivery-mirror",
__openclaw: { idempotencyKey: "announce:v1:child:run:text-direct" },
content: [{ type: "text", text: marker }],
},
],
};
}
throw new Error(`unexpected RPC ${method}`);
},
},
},
transport: {
sendInbound: async (input: Parameters<typeof state.addInboundMessage>[0]) => {
const message = state.addInboundMessage(input);
const send = (text: string) =>
state.addOutboundMessage({
accountId: "default",
to: `dm:${input.conversation.id}`,
text,
});
send(marker);
parentSend = releaseParent.promise.then(() => {
for (let i = 0; i < acknowledgments; i++) {
send("Worker started.");
}
parentSent.resolve();
});
return message;
},
},
fetchJson: async (url: string) => (url.endsWith("request-cursor") ? { cursor: 0 } : requests),
recentOutboundSummary: () => "synthetic terminal messages",
waitForCondition: async (
check: () => Promise<unknown>,
timeout: number,
interval: number,
) => {
expect([timeout, interval]).toEqual([60000, 250]);
const early = await check();
observed.resolve(early);
if (early !== undefined) {
return early;
}
await parentSent.promise;
const settled = await check();
if (settled === undefined) {
throw new Error("terminal observation deadline: parent acknowledgment missing");
}
return settled;
},
},
});
// Observe failures immediately while the test coordinates the held send.
const outcome = result.then(
(value) => ({ value }),
(error: unknown) => ({ error }),
);
return {
observed: Promise.race([
observed.promise,
outcome.then((settled) => {
if ("error" in settled) {
throw settled.error;
}
throw new Error("terminal flow ended before observing child delivery");
}),
]),
outcome,
async release() {
releaseParent.resolve();
await parentSend;
await outcome;
},
};
}
describe("terminal completion scenario parent acknowledgment", () => {
it("waits for the parent send after child delivery and its receipt settle", async () => {
const run = await startPublicCase(1);
try {
expect(await run.observed).toBeUndefined();
} finally {
await run.release();
}
expect(await run.outcome).toMatchObject({ value: { status: "pass" } });
});
it.each([0, 2])(
"rejects %i parent acknowledgments despite settled child delivery",
async (count) => {
const run = await startPublicCase(count);
try {
await run.observed;
} finally {
await run.release();
}
const outcome = await run.outcome;
expect(outcome).toHaveProperty("error");
if ("error" in outcome) {
expect(String(outcome.error)).toContain(
count === 0 ? "parent acknowledgment missing" : "spawning parent did not acknowledge",
);
}
},
);
});

View file

@ -1,6 +1,6 @@
import { setTimeout as delay } from "node:timers/promises"; import { setTimeout as delay } from "node:timers/promises";
import { describe, expect, it, vi } from "vitest"; import { describe, expect, it, onTestFinished, vi } from "vitest";
import { withTestTimeout } from "../../../test/helpers/promise.js"; import { createDeferred, withTestTimeout } from "../../../test/helpers/promise.js";
import { isPidDefinitelyDead } from "../../shared/pid-alive.js"; import { isPidDefinitelyDead } from "../../shared/pid-alive.js";
import { runCommandWithTimeout, runExec } from "../exec.js"; import { runCommandWithTimeout, runExec } from "../exec.js";
import { runWithSpawnBroker } from "./context.js"; import { runWithSpawnBroker } from "./context.js";
@ -13,50 +13,10 @@ describe.skipIf(skipBrokerTests)("command startup cancellation", () => {
"keeps one runExec deadline across admission and execution (cooperative exit: %s)", "keeps one runExec deadline across admission and execution (cooperative exit: %s)",
async (cooperative) => { async (cooperative) => {
const host = createSpawnBrokerHost(); const host = createSpawnBrokerHost();
await host.ready();
const spawnExeca = host.spawnExeca.bind(host); const spawnExeca = host.spawnExeca.bind(host);
let remote: ReturnType<SpawnBrokerHost["spawnExeca"]> | undefined; let remote: ReturnType<SpawnBrokerHost["spawnExeca"]> | undefined;
vi.spyOn(host, "spawnExeca").mockImplementation((...args) => { onTestFinished(async () => {
remote = spawnExeca(...args); vi.useRealTimers();
return remote;
});
const source = `
${cooperative ? "process.on('SIGTERM',()=>{process.stdout.write('-stopped');process.exit(0)});" : ""}
process.stdout.write('started');process.stderr.write('diagnostic');
setTimeout(()=>process.stdout.write('-finished'),1400);
`;
process.kill(host.pid!, "SIGSTOP");
try {
const command = runWithSpawnBroker(host, () =>
runExec(process.execPath, ["-e", source], {
timeoutMs: 2000,
logOutput: false,
}),
);
const outcome = command.then(
(value) => ({ value }),
(error: unknown) => ({ error }),
);
await delay(900);
process.kill(host.pid!, "SIGCONT");
await withTestTimeout(
remote!.child.ready(),
1000,
"command did not start within its remaining budget",
);
expect(await outcome).toMatchObject({
error: {
timedOut: true,
message: "Command timed out",
shortMessage: "Command timed out",
stdout: cooperative ? "started-stopped" : "started",
stderr: "diagnostic",
...(cooperative ? { exitCode: 0 } : { signal: "SIGTERM" }),
},
});
await remote!.child.waitForClose();
expect(isPidDefinitelyDead(remote!.child.pid!)).toBe(true);
} finally {
try { try {
process.kill(host.pid!, "SIGCONT"); process.kill(host.pid!, "SIGCONT");
} catch {} } catch {}
@ -66,7 +26,64 @@ describe.skipIf(skipBrokerTests)("command startup cancellation", () => {
} }
await host.close(); await host.close();
vi.restoreAllMocks(); vi.restoreAllMocks();
} });
await host.ready();
vi.spyOn(host, "spawnExeca").mockImplementation((argv, options) => {
// This case owns the parent deadline; Execa's independent execution
// timeout has parity coverage and must not rescue a broken parent clock.
remote = spawnExeca(argv, { ...options, timeout: undefined });
return remote;
});
const source = `
${cooperative ? "process.on('SIGTERM',()=>{process.stdout.write('-stopped');process.exit(0)});" : ""}
process.stdout.write('started');process.stderr.write('diagnostic');
setInterval(()=>{},1000);
`;
const outputReady = createDeferred();
const output = { stdout: "", stderr: "" };
process.kill(host.pid!, "SIGSTOP");
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] });
const command = runWithSpawnBroker(host, () =>
runExec(process.execPath, ["-e", source], {
timeoutMs: 2000,
logOutput: false,
onOutputChunk: (chunk, stream) => {
output[stream] += chunk.toString();
if (output.stdout === "started" && output.stderr === "diagnostic") {
outputReady.resolve();
}
},
}),
);
const outcome = command.then(
(value) => ({ value }),
(error: unknown) => ({ error }),
);
await vi.advanceTimersByTimeAsync(900);
expect(remote!.child.pid).toBeUndefined();
process.kill(host.pid!, "SIGCONT");
await Promise.race([
outputReady.promise,
outcome.then(() => {
throw new Error("command ended before output readiness");
}),
]);
await vi.advanceTimersByTimeAsync(1099);
expect(remote!.child.killed).toBe(false);
await vi.advanceTimersByTimeAsync(1);
expect(remote!.child.killed).toBe(true);
expect(await outcome).toMatchObject({
error: {
timedOut: true,
message: "Command timed out",
shortMessage: "Command timed out",
stdout: cooperative ? "started-stopped" : "started",
stderr: "diagnostic",
...(cooperative ? { exitCode: 0 } : { signal: "SIGTERM" }),
},
});
await remote!.child.waitForClose();
expect(isPidDefinitelyDead(remote!.child.pid!)).toBe(true);
}, },
); );

View file

@ -8,22 +8,36 @@ import WebSocket, { WebSocketServer } from "ws";
import { createDeferred, withTestTimeout } from "../../test/helpers/promise.js"; import { createDeferred, withTestTimeout } from "../../test/helpers/promise.js";
import { import {
createRealtimeTranscriptionWebSocketSession, createRealtimeTranscriptionWebSocketSession,
type RealtimeTranscriptionWebSocketSessionOptions,
type RealtimeTranscriptionWebSocketTransport, type RealtimeTranscriptionWebSocketTransport,
} from "./websocket-session.js"; } from "./websocket-session.js";
let cleanup: (() => Promise<void>) | undefined; let cleanup: (() => Promise<void>) | undefined;
const sessions = new Set<ReturnType<typeof createRealtimeTranscriptionWebSocketSession>>();
beforeEach(() => { beforeEach(() => {
vi.useRealTimers(); vi.useRealTimers();
}); });
afterEach(async () => { afterEach(async () => {
for (const session of sessions) {
session.close();
}
sessions.clear();
vi.restoreAllMocks(); vi.restoreAllMocks();
vi.useRealTimers(); vi.useRealTimers();
await cleanup?.(); await cleanup?.();
cleanup = undefined; cleanup = undefined;
}); });
function createSession<Event = unknown>(
options: RealtimeTranscriptionWebSocketSessionOptions<Event>,
) {
const session = createRealtimeTranscriptionWebSocketSession(options);
sessions.add(session);
return session;
}
async function createRealtimeServer(params?: { async function createRealtimeServer(params?: {
closeOnConnection?: boolean; closeOnConnection?: boolean;
initialEvent?: unknown; initialEvent?: unknown;
@ -113,7 +127,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
} }
}, },
}); });
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: server.url, url: server.url,
@ -129,7 +143,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
await framesReady.promise; await framesReady.promise;
expect(Buffer.concat(frames).toString()).toBe("queuedafter"); expect(Buffer.concat(frames).toString()).toBe("queuedafter");
expect(session.isConnected()).toBe(true); expect(session.isConnected()).toBe(true);
session.close();
}); });
it("drops the oldest queued audio by bytes and flushes the retained tail in order", async () => { it("drops the oldest queued audio by bytes and flushes the retained tail in order", async () => {
@ -143,7 +156,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
} }
}, },
}); });
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: server.url, url: server.url,
@ -174,7 +187,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
encodeSequence(12_999), encodeSequence(12_999),
Buffer.from([0xaa, 0xbb, 0xcc]), Buffer.from([0xaa, 0xbb, 0xcc]),
]); ]);
session.close();
}); });
it.each([ it.each([
@ -189,7 +201,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
async ({ readyOnOpen, initialEvent }) => { async ({ readyOnOpen, initialEvent }) => {
const server = await createRealtimeServer({ initialEvent }); const server = await createRealtimeServer({ initialEvent });
const onError = vi.fn(); const onError = vi.fn();
const session = createRealtimeTranscriptionWebSocketSession<{ type?: string }>({ const session = createSession<{ type?: string }>({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: server.url, url: server.url,
@ -208,7 +220,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
await expect(session.connect()).rejects.toThrow("queued audio send failed"); await expect(session.connect()).rejects.toThrow("queued audio send failed");
expect(session.isConnected()).toBe(false); expect(session.isConnected()).toBe(false);
expect(onError).toHaveBeenCalledOnce(); expect(onError).toHaveBeenCalledOnce();
session.close();
}, },
); );
@ -216,7 +227,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
const server = await createRealtimeServer(); const server = await createRealtimeServer();
const sentFrames: string[] = []; const sentFrames: string[] = [];
let shouldFailSecondFrame = true; let shouldFailSecondFrame = true;
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: server.url, url: server.url,
@ -237,7 +248,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
await session.connect(); await session.connect();
expect(sentFrames).toEqual(["first", "second"]); expect(sentFrames).toEqual(["first", "second"]);
session.close();
}); });
it("flushes a large retained audio tail in order after reconnect", async () => { it("flushes a large retained audio tail in order after reconnect", async () => {
@ -253,7 +263,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
} }
}, },
}); });
const session = createRealtimeTranscriptionWebSocketSession<{ type?: string }>({ const session = createSession<{ type?: string }>({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: server.url, url: server.url,
@ -290,7 +300,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
encodeSequence(12_998), encodeSequence(12_998),
encodeSequence(12_999), encodeSequence(12_999),
]); ]);
session.close();
}); });
it("discards a large queued audio tail when closed before connecting", async () => { it("discards a large queued audio tail when closed before connecting", async () => {
@ -302,7 +311,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
framesReady.resolve(); framesReady.resolve();
}, },
}); });
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: server.url, url: server.url,
@ -322,7 +331,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
session.sendAudio(Buffer.from("live")); session.sendAudio(Buffer.from("live"));
await framesReady.promise; await framesReady.promise;
expect(frames).toEqual([Buffer.from("live")]); expect(frames).toEqual([Buffer.from("live")]);
session.close();
}); });
it("keeps replacement sockets owned when retired socket callbacks arrive late", async () => { it("keeps replacement sockets owned when retired socket callbacks arrive late", async () => {
@ -333,7 +341,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
const server = await createRealtimeServer({ const server = await createRealtimeServer({
onConnection: (socket) => connections.push(socket), onConnection: (socket) => connections.push(socket),
}); });
const session = createRealtimeTranscriptionWebSocketSession<{ text?: string }>({ const session = createSession<{ text?: string }>({
providerId: "test", providerId: "test",
callbacks: { onError, onTranscript }, callbacks: { onError, onTranscript },
url: server.url, url: server.url,
@ -381,7 +389,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
connections[1]?.send(JSON.stringify({ text: "current transcript" })); connections[1]?.send(JSON.stringify({ text: "current transcript" }));
await vi.waitFor(() => expect(onTranscript).toHaveBeenCalledWith("current transcript")); await vi.waitFor(() => expect(onTranscript).toHaveBeenCalledWith("current transcript"));
session.close();
}); });
it("discards superseded asynchronous connection preparation", async () => { it("discards superseded asynchronous connection preparation", async () => {
@ -394,7 +401,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
resolveFirstUrl = resolve; resolveFirstUrl = resolve;
}); });
let connectionAttempt = 0; let connectionAttempt = 0;
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: async () => (++connectionAttempt === 1 ? await firstUrl : server.url), url: async () => (++connectionAttempt === 1 ? await firstUrl : server.url),
@ -413,7 +420,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
await delay(20); await delay(20);
expect(connections).toHaveLength(1); expect(connections).toHaveLength(1);
expect(session.isConnected()).toBe(true); expect(session.isConnected()).toBe(true);
session.close();
}); });
it("cancels a retired reconnect delay before starting a replacement socket", async () => { it("cancels a retired reconnect delay before starting a replacement socket", async () => {
@ -421,7 +427,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
const server = await createRealtimeServer({ const server = await createRealtimeServer({
onConnection: (socket) => connections.push(socket), onConnection: (socket) => connections.push(socket),
}); });
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: server.url, url: server.url,
@ -442,7 +448,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
expect(connections).toHaveLength(2); expect(connections).toHaveLength(2);
expect(session.isConnected()).toBe(true); expect(session.isConnected()).toBe(true);
session.close();
}); });
it("reconnects after a healthy successor closes following a failed connection", async () => { it("reconnects after a healthy successor closes following a failed connection", async () => {
@ -460,7 +465,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
); );
}, },
}); });
const session = createRealtimeTranscriptionWebSocketSession<{ const session = createSession<{
message?: string; message?: string;
type?: string; type?: string;
}>({ }>({
@ -491,7 +496,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
connections[2]?.close(1011, "healthy connection dropped"); connections[2]?.close(1011, "healthy connection dropped");
await vi.waitFor(() => expect(connections).toHaveLength(4), { timeout: 500 }); await vi.waitFor(() => expect(connections).toHaveLength(4), { timeout: 500 });
await vi.waitFor(() => expect(session.isConnected()).toBe(true)); await vi.waitFor(() => expect(session.isConnected()).toBe(true));
session.close();
}); });
it("delivers graceful provider finals before natural close and finalizes only once", async () => { it("delivers graceful provider finals before natural close and finalizes only once", async () => {
@ -512,7 +516,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
} }
}, },
}); });
const session = createRealtimeTranscriptionWebSocketSession<{ text?: string }>({ const session = createSession<{ text?: string }>({
providerId: "test", providerId: "test",
callbacks: { onTranscript: (text) => transcripts.push(text) }, callbacks: { onTranscript: (text) => transcripts.push(text) },
url: server.url, url: server.url,
@ -542,7 +546,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
it("terminates the captured socket when graceful provider shutdown expires", async () => { it("terminates the captured socket when graceful provider shutdown expires", async () => {
const server = await createRealtimeServer(); const server = await createRealtimeServer();
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: server.url, url: server.url,
@ -565,7 +569,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
it("terminates once when binary audio reaches the active socket buffer cap", async () => { it("terminates once when binary audio reaches the active socket buffer cap", async () => {
const onError = vi.fn(); const onError = vi.fn();
const server = await createRealtimeServer(); const server = await createRealtimeServer();
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: server.url, url: server.url,
@ -600,7 +604,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
let providerTransport: RealtimeTranscriptionWebSocketTransport | undefined; let providerTransport: RealtimeTranscriptionWebSocketTransport | undefined;
const onError = vi.fn(); const onError = vi.fn();
const server = await createRealtimeServer(); const server = await createRealtimeServer();
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: server.url, url: server.url,
@ -633,7 +637,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
it("rejects connect when provider handshake frames exceed the socket buffer cap", async () => { it("rejects connect when provider handshake frames exceed the socket buffer cap", async () => {
const onError = vi.fn(); const onError = vi.fn();
const server = await createRealtimeServer(); const server = await createRealtimeServer();
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: server.url, url: server.url,
@ -670,7 +674,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
} }
}, },
}); });
const session = createRealtimeTranscriptionWebSocketSession<{ type?: string }>({ const session = createSession<{ type?: string }>({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: server.url, url: server.url,
@ -692,7 +696,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
{ type: "session.update" }, { type: "session.update" },
{ type: "input_audio.append", audio: Buffer.from("queued").toString("base64") }, { type: "input_audio.append", audio: Buffer.from("queued").toString("base64") },
]); ]);
session.close();
}); });
it("resolves async URLs and headers before opening the socket", async () => { it("resolves async URLs and headers before opening the socket", async () => {
@ -702,7 +705,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
seenAuthHeaders.push(headers.authorization); seenAuthHeaders.push(headers.authorization);
}, },
}); });
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: async () => server.url, url: async () => server.url,
@ -716,13 +719,12 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
await session.connect(); await session.connect();
expect(seenAuthHeaders).toEqual(["Bearer resolved-token"]); expect(seenAuthHeaders).toEqual(["Bearer resolved-token"]);
session.close();
}); });
it("applies the connect timeout while resolving async connection details", async () => { it("applies the connect timeout while resolving async connection details", async () => {
vi.useFakeTimers(); vi.useFakeTimers();
const onError = vi.fn(); const onError = vi.fn();
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: () => new Promise<string>(() => {}), url: () => new Promise<string>(() => {}),
@ -734,23 +736,18 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
}, },
}); });
try { const connecting = session.connect();
const connecting = session.connect(); const timeoutAssertion = expect(connecting).rejects.toThrow(
const timeoutAssertion = expect(connecting).rejects.toThrow( "test realtime transcription connection timeout",
"test realtime transcription connection timeout", );
); await vi.advanceTimersByTimeAsync(10);
await vi.advanceTimersByTimeAsync(10);
await timeoutAssertion; await timeoutAssertion;
expect(session.isConnected()).toBe(false); expect(session.isConnected()).toBe(false);
expect(onError).toHaveBeenCalledTimes(1); expect(onError).toHaveBeenCalledTimes(1);
const timeoutError = requireFirstMockArg(onError, "connect timeout error"); const timeoutError = requireFirstMockArg(onError, "connect timeout error");
expect(timeoutError).toBeInstanceOf(Error); expect(timeoutError).toBeInstanceOf(Error);
expect(timeoutError.message).toBe("test realtime transcription connection timeout"); expect(timeoutError.message).toBe("test realtime transcription connection timeout");
} finally {
session.close();
vi.useRealTimers();
}
}); });
it("preserves connect failures when the error callback throws", async () => { it("preserves connect failures when the error callback throws", async () => {
@ -760,7 +757,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
const onError = vi.fn((_error: Error) => { const onError = vi.fn((_error: Error) => {
throw new Error("error observer failed"); throw new Error("error observer failed");
}); });
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: () => new Promise<string>(() => {}), url: () => new Promise<string>(() => {}),
@ -786,13 +783,11 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
expect(timeoutError).toBeInstanceOf(Error); expect(timeoutError).toBeInstanceOf(Error);
expect(timeoutError.message).toBe("test realtime transcription connection timeout"); expect(timeoutError.message).toBe("test realtime transcription connection timeout");
} finally { } finally {
session.close();
if (previousDebugProxyEnabled === undefined) { if (previousDebugProxyEnabled === undefined) {
delete process.env.OPENCLAW_DEBUG_PROXY_ENABLED; delete process.env.OPENCLAW_DEBUG_PROXY_ENABLED;
} else { } else {
process.env.OPENCLAW_DEBUG_PROXY_ENABLED = previousDebugProxyEnabled; process.env.OPENCLAW_DEBUG_PROXY_ENABLED = previousDebugProxyEnabled;
} }
vi.useRealTimers();
} }
}); });
@ -807,7 +802,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
seenAuthHeaders.push(headers.authorization); seenAuthHeaders.push(headers.authorization);
}, },
}); });
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: () => url, url: () => url,
@ -830,7 +825,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
it("rejects provider setup errors before ready", async () => { it("rejects provider setup errors before ready", async () => {
const server = await createRealtimeServer({ initialEvent: { type: "error", message: "nope" } }); const server = await createRealtimeServer({ initialEvent: { type: "error", message: "nope" } });
const onError = vi.fn(); const onError = vi.fn();
const session = createRealtimeTranscriptionWebSocketSession<{ const session = createSession<{
type?: string; type?: string;
message?: string; message?: string;
}>({ }>({
@ -859,7 +854,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
const server = await createRealtimeServer({ initialText: "{not json" }); const server = await createRealtimeServer({ initialText: "{not json" });
const received = createDeferred(); const received = createDeferred();
const onError = vi.fn((_error: Error) => received.resolve()); const onError = vi.fn((_error: Error) => received.resolve());
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: server.url, url: server.url,
@ -872,16 +867,12 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
}, },
}); });
try { await session.connect();
await session.connect(); await withTestTimeout(received.promise, 1_000, "Malformed JSON error not received");
await withTestTimeout(received.promise, 1_000, "Malformed JSON error not received"); expect(onError).toHaveBeenCalledTimes(1);
expect(onError).toHaveBeenCalledTimes(1); const parseError = requireFirstMockArg(onError, "malformed websocket json error");
const parseError = requireFirstMockArg(onError, "malformed websocket json error"); expect(parseError).toBeInstanceOf(Error);
expect(parseError).toBeInstanceOf(Error); expect(parseError.message).toBe("Realtime transcription websocket received malformed JSON.");
expect(parseError.message).toBe("Realtime transcription websocket received malformed JSON.");
} finally {
session.close();
}
}); });
it("keeps error callback failures inside websocket message dispatch", async () => { it("keeps error callback failures inside websocket message dispatch", async () => {
@ -891,7 +882,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
received.resolve(); received.resolve();
throw new Error("error observer failed"); throw new Error("error observer failed");
}); });
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: server.url, url: server.url,
@ -904,23 +895,19 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
}, },
}); });
try { await session.connect();
await session.connect(); await withTestTimeout(received.promise, 1_000, "Throwing error observer not reached");
await withTestTimeout(received.promise, 1_000, "Throwing error observer not reached"); expect(onError).toHaveBeenCalledTimes(1);
expect(onError).toHaveBeenCalledTimes(1); const parseError = requireFirstMockArg(onError, "malformed websocket json error");
const parseError = requireFirstMockArg(onError, "malformed websocket json error"); expect(parseError).toBeInstanceOf(Error);
expect(parseError).toBeInstanceOf(Error); expect(parseError.message).toBe("Realtime transcription websocket received malformed JSON.");
expect(parseError.message).toBe("Realtime transcription websocket received malformed JSON."); expect(session.isConnected()).toBe(true);
expect(session.isConnected()).toBe(true);
} finally {
session.close();
}
}); });
it("reports pre-ready closes separately from connection timeouts", async () => { it("reports pre-ready closes separately from connection timeouts", async () => {
const server = await createRealtimeServer({ closeOnConnection: true }); const server = await createRealtimeServer({ closeOnConnection: true });
const onError = vi.fn(); const onError = vi.fn();
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: server.url, url: server.url,
@ -950,7 +937,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
setTimeout(() => ws.close(1011, "flap"), 1); setTimeout(() => ws.close(1011, "flap"), 1);
}, },
}); });
const session = createRealtimeTranscriptionWebSocketSession<{ type?: string }>({ const session = createSession<{ type?: string }>({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: server.url, url: server.url,
@ -978,7 +965,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
{ timeout: 1000 }, { timeout: 1000 },
); );
expect(openCount).toBe(4); expect(openCount).toBe(4);
session.close();
}); });
it("refreshes the reconnect budget after a stable ready connection", async () => { it("refreshes the reconnect budget after a stable ready connection", async () => {
@ -990,7 +976,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
initialEvent: { type: "session.created" }, initialEvent: { type: "session.created" },
onConnection: (ws) => connections.push(ws), onConnection: (ws) => connections.push(ws),
}); });
const session = createRealtimeTranscriptionWebSocketSession<{ type?: string }>({ const session = createSession<{ type?: string }>({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: server.url, url: server.url,
@ -1025,7 +1011,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
); );
}); });
expect(connections).toHaveLength(3); expect(connections).toHaveLength(3);
session.close();
}); });
it("delivers a legitimate large inbound message below the payload cap", async () => { it("delivers a legitimate large inbound message below the payload cap", async () => {
@ -1037,7 +1022,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
}); });
const received = createDeferred(); const received = createDeferred();
const onMessage = vi.fn(() => received.resolve()); const onMessage = vi.fn(() => received.resolve());
const session = createRealtimeTranscriptionWebSocketSession<{ type?: string; text?: string }>({ const session = createSession<{ type?: string; text?: string }>({
providerId: "test", providerId: "test",
callbacks: {}, callbacks: {},
url: server.url, url: server.url,
@ -1048,15 +1033,11 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
}, },
}); });
try { await session.connect();
await session.connect(); await withTestTimeout(received.promise, 1_000, "Large inbound message not received");
await withTestTimeout(received.promise, 1_000, "Large inbound message not received"); expect(onMessage).toHaveBeenCalledTimes(1);
expect(onMessage).toHaveBeenCalledTimes(1); const event = requireFirstMockArg(onMessage, "large inbound message");
const event = requireFirstMockArg(onMessage, "large inbound message"); expect(event).toEqual({ type: "transcript", text: largeText });
expect(event).toEqual({ type: "transcript", text: largeText });
} finally {
session.close();
}
}); });
it("drops an oversized inbound message before it reaches the provider parser", async () => { it("drops an oversized inbound message before it reaches the provider parser", async () => {
@ -1069,7 +1050,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
const onMessage = vi.fn(() => { const onMessage = vi.fn(() => {
throw new Error("oversized frame should not reach provider handler"); throw new Error("oversized frame should not reach provider handler");
}); });
const session = createRealtimeTranscriptionWebSocketSession({ const session = createSession({
providerId: "test", providerId: "test",
callbacks: { onError }, callbacks: { onError },
url: server.url, url: server.url,
@ -1080,17 +1061,13 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
}, },
}); });
try { await session.connect();
await session.connect(); await withTestTimeout(received.promise, 1_000, "Oversized message error not received");
await withTestTimeout(received.promise, 1_000, "Oversized message error not received"); expect(onError).toHaveBeenCalledTimes(1);
expect(onError).toHaveBeenCalledTimes(1); expect(onMessage).not.toHaveBeenCalled();
expect(onMessage).not.toHaveBeenCalled(); const overflowError = requireFirstMockArg(onError, "oversized inbound message error");
const overflowError = requireFirstMockArg(onError, "oversized inbound message error"); expect(overflowError).toBeInstanceOf(Error);
expect(overflowError).toBeInstanceOf(Error); expect(overflowError).toHaveProperty("code", "WS_ERR_UNSUPPORTED_MESSAGE_LENGTH");
expect(overflowError).toHaveProperty("code", "WS_ERR_UNSUPPORTED_MESSAGE_LENGTH"); expect(overflowError.message).toMatch(/max payload/i);
expect(overflowError.message).toMatch(/max payload/i);
} finally {
session.close();
}
}); });
}); });

View file

@ -398,6 +398,13 @@ async function persist(pending: PendingEvent): Promise<void> {
}, },
beforeObservers: async (assertCurrentPublication) => { beforeObservers: async (assertCurrentPublication) => {
if (pending.publication && pending.phase.kind !== "consumed") { if (pending.publication && pending.phase.kind !== "consumed") {
const assertCurrentOwners = () => {
assertCurrentPublication();
if (getTaskRegistryStore() !== store || getTaskFlowRegistryStore() !== flowStore) {
throw new Error("Task event publication owners changed");
}
};
assertCurrentOwners();
const current = tasks.get(taskId); const current = tasks.get(taskId);
if ( if (
pending.publication.becomesTerminal && pending.publication.becomesTerminal &&
@ -408,16 +415,9 @@ async function persist(pending: PendingEvent): Promise<void> {
} }
await finishTaskMutation(context, store, flowStore, taskId, { await finishTaskMutation(context, store, flowStore, taskId, {
operation: "update", operation: "update",
assertCurrent: () => { assertCurrent: assertCurrentOwners,
assertCurrentPublication();
if (
getTaskRegistryStore() !== store ||
getTaskFlowRegistryStore() !== flowStore
) {
throw new Error("Task event publication owners changed");
}
},
}); });
assertCurrentOwners();
flowEffectsSettled = true; flowEffectsSettled = true;
} }
}, },

View file

@ -0,0 +1,92 @@
import { expect, vi } from "vitest";
import { subagentRuns } from "../agents/subagents/registry/subagent-registry-memory.js";
import {
createGatewayMethodRegistry,
createCoreGatewayMethodDescriptors,
} from "../gateway/methods/registry.js";
import { handleGatewayRequest, coreGatewayHandlers } from "../gateway/server-methods.js";
import type { GatewayClient, GatewayRequestContext } from "../gateway/server-methods/types.js";
import { resetAgentEventsForTest } from "../infra/agent-events.js";
import {
getActiveGatewayRootWorkCount,
getActiveGatewayRootWorkHolders,
resetGatewayWorkAdmission,
} from "../process/gateway-work-admission.js";
import { closeOpenClawStateDatabaseAsync } from "../state/openclaw-state-db.js";
import { withOpenClawTestState } from "../test-utils/openclaw-test-state.js";
import { createTaskFixture } from "./task-registry.test-support.js";
import {
resetTaskFlowRegistryForTests,
resetTaskRegistryForTests,
} from "./task-runtime.test-helpers.js";
export async function withReadState(run: () => Promise<void>) {
await withOpenClawTestState({ layout: "state-only" }, async () => {
try {
await run();
} finally {
const holders = getActiveGatewayRootWorkHolders();
if (holders.length) {
console.info("Task read cleanup joining owners:", holders);
}
await closeOpenClawStateDatabaseAsync();
expect(
getActiveGatewayRootWorkCount(),
JSON.stringify(getActiveGatewayRootWorkHolders()),
).toBe(0);
}
});
}
export async function requestTasks(ownerKey: string, respond = vi.fn()) {
const client: GatewayClient = {
connId: "task-read-fixture",
connect: {
minProtocol: 1,
maxProtocol: 1,
client: {
id: "openclaw-control-ui",
version: "test",
platform: "test",
mode: "webchat",
},
role: "operator",
scopes: ["operator.read"],
},
};
await handleGatewayRequest({
req: {
type: "req",
id: "task-read",
method: "tasks.list",
params: { limit: 5, sessionKey: ownerKey },
},
client,
context: { getRuntimeConfig: () => ({}) } as GatewayRequestContext,
methodRegistry: createGatewayMethodRegistry(
createCoreGatewayMethodDescriptors(coreGatewayHandlers),
),
isWebchatConnect: () => false,
respond,
});
return respond;
}
export function resetReadState() {
vi.restoreAllMocks();
resetTaskRegistryForTests({ persist: false });
resetTaskFlowRegistryForTests({ persist: false });
resetAgentEventsForTest({ preserveListeners: true });
resetGatewayWorkAdmission();
subagentRuns.clear();
}
export function createReadTask(runId: string) {
return createTaskFixture("cli", {
runId,
task: "Read accepted events",
status: "running",
notifyPolicy: "silent",
deliveryStatus: "not_applicable",
});
}

View file

@ -7,133 +7,56 @@ import { subagentRuns } from "../agents/subagents/registry/subagent-registry-mem
import { settleRequesterTurnAfterSessionSpawns } from "../agents/subagents/registry/subagent-registry-requester-yield.js"; import { settleRequesterTurnAfterSessionSpawns } from "../agents/subagents/registry/subagent-registry-requester-yield.js";
import type { SubagentRunRecord } from "../agents/subagents/registry/subagent-registry.types.js"; import type { SubagentRunRecord } from "../agents/subagents/registry/subagent-registry.types.js";
import { createSubagentsTool } from "../agents/tools/subagents-tool.js"; import { createSubagentsTool } from "../agents/tools/subagents-tool.js";
import { import { emitAgentEvent } from "../infra/agent-events.js";
createGatewayMethodRegistry,
createCoreGatewayMethodDescriptors,
} from "../gateway/methods/registry.js";
import { handleGatewayRequest, coreGatewayHandlers } from "../gateway/server-methods.js";
import type { GatewayClient, GatewayRequestContext } from "../gateway/server-methods/types.js";
import { emitAgentEvent, resetAgentEventsForTest } from "../infra/agent-events.js";
import { SqliteWorkerError } from "../infra/sqlite-worker-contract.js"; import { SqliteWorkerError } from "../infra/sqlite-worker-contract.js";
import {
getActiveGatewayRootWorkCount,
getActiveGatewayRootWorkHolders,
resetGatewayWorkAdmission,
} from "../process/gateway-work-admission.js";
import { import {
closeOpenClawStateDatabaseAsync, closeOpenClawStateDatabaseAsync,
runOpenClawStateWriteTransaction, runOpenClawStateWriteTransaction,
} from "../state/openclaw-state-db.js"; } from "../state/openclaw-state-db.js";
import { captureOpenClawStateWorkerContext } from "../state/openclaw-state-worker-context.js"; import { captureOpenClawStateWorkerContext } from "../state/openclaw-state-worker-context.js";
import { withOpenClawTestState } from "../test-utils/openclaw-test-state.js";
import { holdStateDatabaseCoordinator } from "../test-utils/state-database-contention.js"; import { holdStateDatabaseCoordinator } from "../test-utils/state-database-contention.js";
import { import {
createSubagentTaskBackingDetail, createSubagentTaskBackingDetail,
resolveManagedTaskBackingDetail, resolveManagedTaskBackingDetail,
} from "./task-backing-authority.js"; } from "./task-backing-authority.js";
import * as taskMutationEffects from "./task-executor-create.async.js";
import { createRunningTaskRunCoreWithReceiptAsync } from "./task-executor-create.async.js"; import { createRunningTaskRunCoreWithReceiptAsync } from "./task-executor-create.async.js";
import { completeTaskRunByRunIdCore } from "./task-executor.js"; import { completeTaskRunByRunIdCore } from "./task-executor.js";
import { createManagedTaskFlow, createTaskFlowForTask } from "./task-flow-registry.js"; import { createManagedTaskFlow, createTaskFlowForTask } from "./task-flow-registry.js";
import { getTaskFlowRegistryStore } from "./task-flow-registry.store.js";
import { taskAgentEventMutations } from "./task-registry-agent-events.js"; import { taskAgentEventMutations } from "./task-registry-agent-events.js";
import * as taskRegistryListenerState from "./task-registry-listener-state.js";
import { updateTask } from "./task-registry-mutation.js"; import { updateTask } from "./task-registry-mutation.js";
import { publishTaskRecordAfterAtomicStore } from "./task-registry-publication.js";
import { import {
getTaskById, getTaskById,
listTaskRecordPage, listTaskRecordPage,
listFreshTasksForOwnerKey, listFreshTasksForOwnerKey,
} from "./task-registry-query.js"; } from "./task-registry-query.js";
import { prepareTaskRegistryRead } from "./task-registry-read.js"; import { prepareTaskRegistryRead } from "./task-registry-read.js";
import {
createReadTask,
requestTasks,
resetReadState,
withReadState,
} from "./task-registry-read.test-support.js";
import { linkTaskToFlowById } from "./task-registry-record-api.js"; import { linkTaskToFlowById } from "./task-registry-record-api.js";
import { tasks, taskProgressBatches } from "./task-registry-state.js"; import { tasks, taskProgressBatches } from "./task-registry-state.js";
import { getTaskRegistryStore, onTaskRegistryChange } from "./task-registry.store.js"; import {
configureTaskRegistryRuntime,
getTaskRegistryStore,
onTaskRegistryChange,
} from "./task-registry.store.js";
import { loadTaskRegistryStateFromSqliteReadOnly } from "./task-registry.store.sqlite.js"; import { loadTaskRegistryStateFromSqliteReadOnly } from "./task-registry.store.sqlite.js";
import { createTaskFixture } from "./task-registry.test-support.js"; import { createTaskFixture } from "./task-registry.test-support.js";
import { import { configureTaskFlowRegistryRuntime } from "./task-runtime.test-helpers.js";
resetTaskFlowRegistryForTests,
resetTaskRegistryForTests,
} from "./task-runtime.test-helpers.js";
vi.mock("node:timers/promises", { spy: true }); vi.mock("node:timers/promises", { spy: true });
afterEach(() => { afterEach(resetReadState);
vi.restoreAllMocks();
resetTaskRegistryForTests({ persist: false });
resetTaskFlowRegistryForTests({ persist: false });
resetAgentEventsForTest({ preserveListeners: true });
resetGatewayWorkAdmission();
subagentRuns.clear();
});
async function withReadState(run: () => Promise<void>) {
await withOpenClawTestState({ layout: "state-only" }, async () => {
try {
await run();
} finally {
const holders = getActiveGatewayRootWorkHolders();
if (holders.length) {
console.info("Task read cleanup joining owners:", holders);
}
await closeOpenClawStateDatabaseAsync();
expect(
getActiveGatewayRootWorkCount(),
JSON.stringify(getActiveGatewayRootWorkHolders()),
).toBe(0);
}
});
}
function emitTool(runId: string, name: string) { function emitTool(runId: string, name: string) {
emitAgentEvent({ runId, stream: "tool", data: { phase: "start", name } }); emitAgentEvent({ runId, stream: "tool", data: { phase: "start", name } });
} }
function createReadTask(runId: string) {
return createTaskFixture("cli", {
runId,
task: "Read accepted events",
status: "running",
notifyPolicy: "silent",
deliveryStatus: "not_applicable",
});
}
async function requestTaskList(ownerKey: string) {
const registry = createGatewayMethodRegistry(
createCoreGatewayMethodDescriptors(coreGatewayHandlers),
);
const client: GatewayClient = {
connId: "task-read-fixture",
connect: {
minProtocol: 1,
maxProtocol: 1,
client: {
id: "openclaw-control-ui",
version: "test",
platform: "test",
mode: "webchat",
},
role: "operator",
scopes: ["operator.read"],
},
};
const context = { getRuntimeConfig: () => ({}) } as GatewayRequestContext;
const respond = vi.fn();
await handleGatewayRequest({
req: {
type: "req",
id: "task-read",
method: "tasks.list",
params: { limit: 5, sessionKey: ownerKey },
},
client,
context,
methodRegistry: registry,
isWebchatConnect: () => false,
respond,
});
return respond;
}
function createReadProgressBatch() { function createReadProgressBatch() {
const entry: SubagentRunRecord = { const entry: SubagentRunRecord = {
runId: "contended-progress-child", runId: "contended-progress-child",
@ -211,117 +134,150 @@ function createReadProgressBatch() {
} }
describe("task registry read preparation", () => { describe("task registry read preparation", () => {
it.each(["unchanged", "newer write", "ABA"] as const)( it.each([
"keeps registered tasks.list current after terminal publication is superseded by %s", "receipt",
async (change) => { "flow follow-up",
"flow error",
"retired owner",
"retired flow owner",
"retired flow owner during cancellation",
"native consumption",
] as const)(
"settles a registered read after publication supersession at %s",
async (boundary) => {
await withReadState(async () => { await withReadState(async () => {
const runId = `superseded-terminal-${change}`; const task = createReadTask("registered-read-superseded");
const task = createReadTask(runId); const flowBoundary = boundary === "flow follow-up" || boundary === "flow error";
const flow = expectDefined(createTaskFlowForTask({ task }), "terminal task flow"); const cancellationBoundary = boundary === "retired flow owner during cancellation";
expect(linkTaskToFlowById({ taskId: task.taskId, flowId: flow.flowId })).not.toBeNull(); const consumed = boundary === "native consumption";
const request = () => requestTaskList(task.ownerKey); const flowFailure = new Error("Synthetic flow synchronization failure");
const terminalInstalled = createDeferred(); if (flowBoundary || boundary === "retired flow owner" || cancellationBoundary) {
const releaseEffects = createDeferred(); const flow = expectDefined(createTaskFlowForTask({ task }), "task flow");
const readCaptured = createDeferred(); expect(linkTaskToFlowById({ taskId: task.taskId, flowId: flow.flowId })).not.toBeNull();
const finish = taskMutationEffects.finishTaskMutation; }
let held = false; const store = getTaskRegistryStore();
vi.spyOn(taskMutationEffects, "finishTaskMutation").mockImplementation(async (...args) => { const mutate = store.runAgentEventMutationAsync.bind(store);
if ( const committed = createDeferred();
!held && const release = createDeferred();
args[3] === task.taskId && const fenced = createDeferred();
args[4].operation === "update" && const captureFence = taskAgentEventMutations.captureReadFence.bind(taskAgentEventMutations);
tasks.get(task.taskId)?.status === "succeeded" vi.spyOn(taskAgentEventMutations, "captureReadFence").mockImplementation((admission) => {
) { const result = captureFence(admission);
held = true; fenced.resolve();
terminalInstalled.resolve(); return result;
await releaseEffects.promise;
}
return finish(...args);
}); });
const captureFence = taskRegistryListenerState.captureTaskRegistryReadFence; let eventCommitted = false;
let captureRead = false; const writes = vi
vi.spyOn(taskRegistryListenerState, "captureTaskRegistryReadFence").mockImplementation( .spyOn(store, "runAgentEventMutationAsync")
(...args) => { .mockImplementation(async (...args) => {
const pending = captureFence(...args); if (consumed) {
if (captureRead) { committed.resolve();
readCaptured.resolve(); await release.promise;
} }
return pending; const receipt = await mutate(...args);
}, eventCommitted = true;
); if (!flowBoundary && !cancellationBoundary) {
committed.resolve();
await release.promise;
}
return receipt;
});
const syncFlow = store.syncLiveTaskFlowAsync.bind(store);
let flowHeld = false;
vi.spyOn(store, "syncLiveTaskFlowAsync").mockImplementation(async (...args) => {
const result = await syncFlow(...args);
if (flowBoundary && eventCommitted && !flowHeld) {
flowHeld = true;
committed.resolve();
await release.promise;
if (boundary === "flow error") {
throw flowFailure;
}
}
return result;
});
const initialMutation = store.runInitialMutationAsync.bind(store);
vi.spyOn(store, "runInitialMutationAsync").mockImplementation(async (...args) => {
const result = await initialMutation(...args);
if (cancellationBoundary && args[1].type === "flows.finalizeTaskCancellation") {
committed.resolve();
await release.promise;
}
return result;
});
const publications: string[] = []; const publications: string[] = [];
const stop = onTaskRegistryChange((event) => { const stop = onTaskRegistryChange(() => {
if (event?.kind === "upserted" && event.task.taskId === task.taskId) { const current = tasks.get(task.taskId);
publications.push(event.task.task); if (current) {
publications.push(current.task);
} }
}); });
const mutation = vi.spyOn(getTaskRegistryStore(), "runAgentEventMutationAsync"); const respond = vi.fn();
let reading: ReturnType<typeof request> | undefined; let read: ReturnType<typeof requestTasks> | undefined;
let readResult:
| Promise<PromiseSettledResult<Awaited<ReturnType<typeof request>>>[]>
| undefined;
try { try {
emitAgentEvent({ emitTool(task.runId!, "accepted-tool");
runId, await committed.promise;
stream: "lifecycle", read = requestTasks(task.ownerKey, respond);
data: { phase: "end", endedAt: Date.now() }, await fenced.promise;
}); expect(respond).not.toHaveBeenCalled();
await withTestTimeout( if (consumed) {
terminalInstalled.promise, expect(getTaskById(task.taskId)?.toolUseCount).toBe(1);
5_000, configureTaskFlowRegistryRuntime({ store: { ...getTaskFlowRegistryStore() } });
"Terminal event reached publication", release.resolve();
); await read;
expect(loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)?.status).toBe( expect(respond.mock.calls[0]).toMatchObject([
"succeeded",
);
captureRead = true;
reading = request();
readResult = Promise.allSettled([reading]);
await withTestTimeout(
readCaptured.promise,
5_000,
"tasks.list captured the terminal event",
);
if (change !== "unchanged") {
expect(updateTask(task.taskId, { task: "Newer task write" })).not.toBeNull();
if (change === "ABA") {
expect(updateTask(task.taskId, { task: task.task })).not.toBeNull();
}
}
const expectedTitle = change === "newer write" ? "Newer task write" : task.task;
const competingPublications = [...publications];
releaseEffects.resolve();
const [outcome] = await withTestTimeout(readResult, 5_000, "Captured tasks.list settled");
expect(loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)).toMatchObject({
task: expectedTitle,
status: "succeeded",
});
expect(mutation).toHaveBeenCalledOnce();
if (change === "unchanged") {
expect(publications).toContain(task.task);
} else {
expect(publications).toEqual(competingPublications);
}
expect((await request()).mock.calls[0]).toMatchObject([
true,
{ tasks: [{ id: task.taskId, title: expectedTitle, status: "completed" }] },
]);
expect(
outcome,
outcome?.status === "rejected" ? String(outcome.reason) : undefined,
).toMatchObject({
status: "fulfilled",
});
if (outcome?.status === "fulfilled") {
expect(outcome.value.mock.calls[0]).toMatchObject([
true, true,
{ tasks: [{ id: task.taskId, title: expectedTitle, status: "completed" }] }, { tasks: [expect.objectContaining({ id: task.taskId, toolUseCount: 1 })] },
]); ]);
expect(writes).toHaveBeenCalledOnce();
expect(publications).toEqual([task.task]);
return;
} }
const durable = loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)!;
const newer = { ...durable, task: "Newer committed task" };
store.upsertTaskWithDeliveryState({ task: newer });
publishTaskRecordAfterAtomicStore(newer);
if (
boundary === "retired owner" ||
boundary === "retired flow owner" ||
cancellationBoundary
) {
if (boundary === "retired owner") {
configureTaskRegistryRuntime({ store: { ...store } });
} else {
configureTaskFlowRegistryRuntime({ store: { ...getTaskFlowRegistryStore() } });
}
release.resolve();
await expect(read).rejects.toThrow("owner");
expect(respond).not.toHaveBeenCalled();
expect(publications).toEqual([newer.task]);
return;
}
release.resolve();
if (boundary === "flow error") {
await expect(read).rejects.toBe(flowFailure);
expect(respond).not.toHaveBeenCalled();
expect(publications).toEqual([newer.task]);
return;
}
await read;
expect(respond.mock.calls[0]).toMatchObject([
true,
{
tasks: [
expect.objectContaining({ id: task.taskId, title: newer.task, toolUseCount: 1 }),
],
},
]);
expect(writes).toHaveBeenCalledOnce();
expect(publications).toEqual([newer.task]);
expect(loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)).toMatchObject({
task: newer.task,
toolUseCount: 1,
});
} finally { } finally {
releaseEffects.resolve(); release.resolve();
await readResult; await read?.catch(() => undefined);
await prepareTaskRegistryRead();
stop(); stop();
} }
}); });
@ -457,6 +413,7 @@ describe("task registry read preparation", () => {
); );
} }
expect(await timer).toBe(0); expect(await timer).toBe(0);
holder.release();
expect(await pending).toMatchObject( expect(await pending).toMatchObject(
surface === "fresh owner" surface === "fresh owner"
? [ ? [
@ -502,48 +459,73 @@ describe("task registry read preparation", () => {
} }
return result; return result;
}); });
const immediate = timers.setImmediate; const reader = await import("./task-registry-read.js");
vi.mocked(timers.setImmediate).mockImplementationOnce(async (...args) => { const prepare = reader.prepareTaskRegistryRead;
await immediate(...args); // Preparatory yields must not consume the scan's publication barrier.
await committed.promise; vi.spyOn(reader, "prepareTaskRegistryRead").mockImplementationOnce(async () => {
await timers.setImmediate();
return prepare();
}); });
const immediate = timers.setImmediate;
let workMs = 0; let workMs = 0;
vi.spyOn(performance, "now").mockImplementation(() => workMs); vi.spyOn(performance, "now").mockImplementation(() => workMs);
let mutation: Promise<unknown> | undefined; let mutation: Promise<unknown> | undefined;
let page: ReturnType<typeof listTaskRecordPage> | undefined;
let selectedBeforeMutation = false; let selectedBeforeMutation = false;
const failures: unknown[] = [];
const recordFailure = (error: unknown) => {
if (!failures.includes(error)) {
failures.push(error);
}
};
try { try {
const page = await withTestTimeout( page = listTaskRecordPage({
listTaskRecordPage({ offset: 0,
offset: 0, limit: 1,
limit: 1, prepareFilter: (batch) => {
prepareFilter: (batch) => { workMs += 20;
workMs += 20; if (!mutation) {
if (!mutation) { selectedBeforeMutation = batch.some((task) => task.taskId === selected.taskId);
selectedBeforeMutation = batch.some((task) => task.taskId === selected.taskId); mutation = createRunningTaskRunCoreWithReceiptAsync({
mutation = createRunningTaskRunCoreWithReceiptAsync({ runtime: selected.runtime,
runtime: selected.runtime, runId: selected.runId!,
runId: selected.runId!, task: selected.task,
task: selected.task, ownerKey: selected.ownerKey,
ownerKey: selected.ownerKey, scopeKind: selected.scopeKind,
scopeKind: selected.scopeKind, requesterSessionKey: selected.requesterSessionKey,
requesterSessionKey: selected.requesterSessionKey, notifyPolicy: "silent",
notifyPolicy: "silent", deliveryStatus: "not_applicable",
deliveryStatus: "not_applicable", detail: { historyGeneration: "replacement" },
detail: { historyGeneration: "replacement" }, });
}); vi.mocked(timers.setImmediate).mockImplementationOnce(async (...args) => {
} await immediate(...args);
return (task) => task.taskId === selected.taskId; await committed.promise;
}, });
}), }
return (task) => task.taskId === selected.taskId;
},
});
const result = await withTestTimeout(
page,
5_000, 5_000,
"Page joined an identity-changing publication", "Page joined an identity-changing publication",
); );
expect(selectedBeforeMutation).toBe(true); expect(selectedBeforeMutation).toBe(true);
expect(held).toBe(true); expect(held).toBe(true);
expect(page).toEqual({ ok: false, error: "registry_changed" }); expect(result).toEqual({ ok: false, error: "registry_changed" });
} catch (error) {
recordFailure(error);
} finally { } finally {
committed.resolve();
release.resolve(); release.resolve();
await mutation; await page?.catch(recordFailure);
await mutation?.catch(recordFailure);
}
if (failures.length === 1) {
throw failures[0];
}
if (failures.length > 1) {
throw new AggregateError(failures, "Page proof and cleanup failed", { cause: failures[0] });
} }
}); });
}); });
@ -573,7 +555,7 @@ describe("task registry read preparation", () => {
if (scenario === "active progress") { if (scenario === "active progress") {
createReadProgressBatch(); createReadProgressBatch();
} }
const request = () => requestTaskList(task.ownerKey); const request = () => requestTasks(task.ownerKey);
emitTool(task.runId!, "warmup"); emitTool(task.runId!, "warmup");
await prepareTaskRegistryRead(); await prepareTaskRegistryRead();
expect((await request()).mock.calls[0]?.[0]).toBe(true); expect((await request()).mock.calls[0]?.[0]).toBe(true);
@ -605,6 +587,7 @@ describe("task registry read preparation", () => {
} }
read = request(); read = request();
expect(await timer).toBe(0); expect(await timer).toBe(0);
holder.release();
expect((await read).mock.calls[0]).toMatchObject([ expect((await read).mock.calls[0]).toMatchObject([
true, true,
{ {
@ -631,115 +614,210 @@ describe("task registry read preparation", () => {
}, },
); );
it("returns the terminal task when accepted metadata publication is superseded", async () => { it.each(["replacement", "ABA"] as const)(
await withReadState(async () => { "serves registered tasks.list after a committed event publication loses to %s",
const entry: SubagentRunRecord = { async (change) => {
runId: "metadata-terminal-overlap", await withReadState(async () => {
childSessionKey: "agent:main:subagent:metadata-terminal-overlap", const task = createReadTask(`superseded-read-${change}`);
requesterSessionKey: "agent:main:main", const store = getTaskRegistryStore();
requesterDisplayKey: "main", const mutate = store.runAgentEventMutationAsync.bind(store);
task: "Finish while task metadata settles", const snapshot = store.loadMutationSnapshotAsync.bind(store);
cleanup: "keep", const entered = createDeferred();
createdAt: Date.now(), const release = createDeferred();
generation: 1, const reading = createDeferred();
execution: { status: "running", startedAt: Date.now() }, let committed = false;
}; let held = false;
subagentRuns.set(entry.runId, entry); vi.spyOn(store, "runAgentEventMutationAsync").mockImplementation(async (...args) => {
const task = createTaskFixture("subagent", {
runId: entry.runId,
childSessionKey: entry.childSessionKey,
requesterSessionKey: entry.requesterSessionKey,
task: entry.task,
notifyPolicy: "silent",
detail: createSubagentTaskBackingDetail(entry.generation!),
});
const store = getTaskRegistryStore();
const mutate = store.runAgentEventMutationAsync.bind(store);
const committed = createDeferred<Awaited<ReturnType<typeof mutate>>>();
const release = createDeferred();
const writes = vi
.spyOn(store, "runAgentEventMutationAsync")
.mockImplementation(async (...args) => {
const receipt = await mutate(...args); const receipt = await mutate(...args);
committed.resolve(receipt); committed = true;
await release.promise;
return receipt; return receipt;
}); });
const captured = createDeferred(); vi.spyOn(store, "loadMutationSnapshotAsync").mockImplementation(async (...args) => {
const capture = taskAgentEventMutations.captureReadFence.bind(taskAgentEventMutations); const result = await snapshot(...args);
const published: string[] = []; if (committed && !held) {
const stop = onTaskRegistryChange((event) => { held = true;
if (event?.kind === "upserted" && event.task.taskId === task.taskId) { entered.resolve();
published.push(event.task.status); await release.promise;
} }
}); return result;
let read: ReturnType<typeof requestTaskList> | undefined;
try {
emitTool(entry.runId, "accepted-before-completion");
const receipt = expectDefined(
await withTestTimeout(committed.promise, 5_000, "Metadata did not commit"),
"ordinary successful metadata receipt",
);
expect(receipt.task).toMatchObject({
taskId: task.taskId,
status: "running",
toolUseCount: 1,
lastToolName: "accepted-before-completion",
}); });
vi.spyOn(taskAgentEventMutations, "captureReadFence").mockImplementation((...args) => { emitTool(task.runId!, "committed-tool");
const fence = capture(...args); const captureFence = taskAgentEventMutations.captureReadFence.bind(taskAgentEventMutations);
captured.resolve(); vi.spyOn(taskAgentEventMutations, "captureReadFence").mockImplementation((admission) => {
const fence = captureFence(admission);
reading.resolve();
return fence; return fence;
}); });
read = requestTaskList(task.ownerKey); let read: ReturnType<typeof requestTasks> | undefined;
await withTestTimeout( try {
captured.promise, await entered.promise;
5_000, const newer = updateTask(task.taskId, {
"Gateway did not capture the accepted fence", task: "Newer title",
); ...(change === "replacement" ? { runId: "successor" } : {}),
expect( });
completeTaskRunByRunIdCore({ expect(newer).not.toBeNull();
runId: entry.runId, if (change === "ABA") {
runtime: "subagent", expect(updateTask(task.taskId, { task: task.task })).not.toBeNull();
sessionKey: entry.childSessionKey, }
endedAt: Date.now(), read = requestTasks(task.ownerKey);
terminalSummary: "Authoritative completion", await reading.promise;
}), release.resolve();
).toEqual([expect.objectContaining({ taskId: task.taskId, status: "succeeded" })]); expect((await read).mock.calls[0]).toMatchObject([
release.resolve(); true,
expect((await read).mock.calls[0]).toMatchObject([ {
true, tasks: [
{ {
tasks: [ id: task.taskId,
{ status: "running",
id: task.taskId, title: change === "replacement" ? "Newer title" : task.task,
status: "completed", toolUseCount: 1,
toolUseCount: 1, },
lastToolName: "accepted-before-completion", ],
}, },
], ]);
}, expect(loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)).toMatchObject({
]); status: "running",
expect(writes).toHaveBeenCalledOnce(); task: change === "replacement" ? "Newer title" : task.task,
expect(published).toEqual(["succeeded"]); toolUseCount: 1,
expect(loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)).toMatchObject({ });
} finally {
release.resolve();
await read;
}
});
},
);
it.each(["before readback", "during readback"] as const)(
"returns the terminal task when accepted metadata publication is superseded %s",
async (timing) => {
await withReadState(async () => {
const entry: SubagentRunRecord = {
runId: "metadata-terminal-overlap",
childSessionKey: "agent:main:subagent:metadata-terminal-overlap",
requesterSessionKey: "agent:main:main",
requesterDisplayKey: "main",
task: "Finish while task metadata settles",
cleanup: "keep",
createdAt: Date.now(),
generation: 1,
execution: { status: "running", startedAt: Date.now() },
};
subagentRuns.set(entry.runId, entry);
const task = createTaskFixture("subagent", {
runId: entry.runId, runId: entry.runId,
childSessionKey: entry.childSessionKey, childSessionKey: entry.childSessionKey,
status: "succeeded", requesterSessionKey: entry.requesterSessionKey,
terminalSummary: "Authoritative completion", task: entry.task,
toolUseCount: 1, notifyPolicy: "silent",
lastToolName: "accepted-before-completion", detail: createSubagentTaskBackingDetail(entry.generation!),
}); });
} finally { const store = getTaskRegistryStore();
release.resolve(); const mutate = store.runAgentEventMutationAsync.bind(store);
const snapshot = store.loadMutationSnapshotAsync.bind(store);
const committed = createDeferred<Awaited<ReturnType<typeof mutate>>>();
const publicationPaused = createDeferred();
const release = createDeferred();
let mutationCommitted = false;
let readbackHeld = false;
const writes = vi
.spyOn(store, "runAgentEventMutationAsync")
.mockImplementation(async (...args) => {
const receipt = await mutate(...args);
mutationCommitted = true;
committed.resolve(receipt);
if (timing === "before readback") {
publicationPaused.resolve();
await release.promise;
}
return receipt;
});
vi.spyOn(store, "loadMutationSnapshotAsync").mockImplementation(async (...args) => {
const current = await snapshot(...args);
if (timing === "during readback" && mutationCommitted && !readbackHeld) {
readbackHeld = true;
publicationPaused.resolve();
await release.promise;
}
return current;
});
const captured = createDeferred();
const capture = taskAgentEventMutations.captureReadFence.bind(taskAgentEventMutations);
const published: string[] = [];
const stop = onTaskRegistryChange((event) => {
if (event?.kind === "upserted" && event.task.taskId === task.taskId) {
published.push(event.task.status);
}
});
let read: ReturnType<typeof requestTasks> | undefined;
try { try {
await read; emitTool(entry.runId, "accepted-before-completion");
const receipt = expectDefined(
await withTestTimeout(committed.promise, 5_000, "Metadata did not commit"),
"ordinary successful metadata receipt",
);
expect(receipt.task).toMatchObject({
taskId: task.taskId,
status: "running",
toolUseCount: 1,
lastToolName: "accepted-before-completion",
});
await withTestTimeout(publicationPaused.promise, 5_000, "Publication did not pause");
vi.spyOn(taskAgentEventMutations, "captureReadFence").mockImplementation((...args) => {
const fence = capture(...args);
captured.resolve();
return fence;
});
read = requestTasks(task.ownerKey);
await withTestTimeout(
captured.promise,
5_000,
"Gateway did not capture the accepted fence",
);
expect(
completeTaskRunByRunIdCore({
runId: entry.runId,
runtime: "subagent",
sessionKey: entry.childSessionKey,
endedAt: Date.now(),
terminalSummary: "Authoritative completion",
}),
).toEqual([expect.objectContaining({ taskId: task.taskId, status: "succeeded" })]);
release.resolve();
expect((await read).mock.calls[0]).toMatchObject([
true,
{
tasks: [
{
id: task.taskId,
status: "completed",
toolUseCount: 1,
lastToolName: "accepted-before-completion",
},
],
},
]);
expect(writes).toHaveBeenCalledOnce();
expect(published).toEqual(["succeeded"]);
expect(loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)).toMatchObject({
runId: entry.runId,
childSessionKey: entry.childSessionKey,
status: "succeeded",
terminalSummary: "Authoritative completion",
toolUseCount: 1,
lastToolName: "accepted-before-completion",
});
} finally { } finally {
stop(); release.resolve();
try {
await read;
} finally {
stop();
}
} }
} });
}); },
}); );
it.each([1, 8])("prepares %i readers through a fixed event fence", async (readers) => { it.each([1, 8])("prepares %i readers through a fixed event fence", async (readers) => {
await withReadState(async () => { await withReadState(async () => {
@ -854,50 +932,56 @@ describe("task registry read preparation", () => {
}, },
); );
it.each(["publication", "unknown settlement", "undefined rejection"] as const)( it.each([
"does not acknowledge an accepted batch after %s", "publication",
async (failureKind) => { "supersession message",
await withReadState(async () => { "unknown settlement",
const task = createReadTask(`failed-read-${failureKind}`); "undefined rejection",
const earlier = expectDefined(await prepareTaskRegistryRead(), "earlier task read"); ] as const)("does not acknowledge an accepted batch after %s", async (failureKind) => {
const store = getTaskRegistryStore(); await withReadState(async () => {
const mutate = store.runAgentEventMutationAsync.bind(store); const task = createReadTask(`failed-read-${failureKind}`);
const snapshot = store.loadMutationSnapshotAsync.bind(store); const earlier = expectDefined(await prepareTaskRegistryRead(), "earlier task read");
const failure = const store = getTaskRegistryStore();
failureKind === "undefined rejection" const mutate = store.runAgentEventMutationAsync.bind(store);
const snapshot = store.loadMutationSnapshotAsync.bind(store);
const publicationFailure =
failureKind === "publication" || failureKind === "supersession message";
const failure =
failureKind === "supersession message"
? new Error("Task publication was superseded by a current write")
: failureKind === "undefined rejection"
? undefined ? undefined
: new SqliteWorkerError(`Synthetic ${failureKind}`, "outcome-unknown"); : new SqliteWorkerError(`Synthetic ${failureKind}`, "outcome-unknown");
const rejection = createDeferred<never>(); const rejection = createDeferred<never>();
let mutationReturned = false; let mutationReturned = false;
vi.spyOn(store, "runAgentEventMutationAsync").mockImplementation(async (...args) => { vi.spyOn(store, "runAgentEventMutationAsync").mockImplementation(async (...args) => {
const result = await mutate(...args); const result = await mutate(...args);
mutationReturned = true; mutationReturned = true;
if (failureKind !== "publication") { if (!publicationFailure) {
rejection.reject(failure); rejection.reject(failure);
return rejection.promise; return rejection.promise;
}
return result;
});
vi.spyOn(store, "loadMutationSnapshotAsync").mockImplementation((...args) => {
if (failureKind === "publication" && mutationReturned) {
rejection.reject(failure);
return rejection.promise;
}
return snapshot(...args);
});
emitTool(task.runId!, "accepted-failure");
const read = prepareTaskRegistryRead();
await expect(read).rejects.toBe(failure);
expect(loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)).toMatchObject({
toolUseCount: 1,
lastToolName: "accepted-failure",
});
if (failureKind === "publication") {
expect(() => earlier.getTaskById(task.taskId)).toThrow("requires preparation");
} }
return result;
}); });
}, vi.spyOn(store, "loadMutationSnapshotAsync").mockImplementation((...args) => {
); if (publicationFailure && mutationReturned) {
rejection.reject(failure);
return rejection.promise;
}
return snapshot(...args);
});
emitTool(task.runId!, "accepted-failure");
const read = prepareTaskRegistryRead();
await expect(read).rejects.toBe(failure);
expect(loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)).toMatchObject({
toolUseCount: 1,
lastToolName: "accepted-failure",
});
if (publicationFailure) {
expect(() => earlier.getTaskById(task.taskId)).toThrow("requires preparation");
}
});
});
it("retires a prepared row reader with its database owner", async () => { it("retires a prepared row reader with its database owner", async () => {
await withReadState(async () => { await withReadState(async () => {

View file

@ -0,0 +1,140 @@
import { expectDefined } from "@openclaw/normalization-core";
import { afterEach, describe, expect, it, vi } from "vitest";
import { createDeferred, withTestTimeout } from "../../test/helpers/promise.js";
import { emitAgentEvent } from "../infra/agent-events.js";
import * as taskMutationEffects from "./task-executor-create.async.js";
import { createTaskFlowForTask } from "./task-flow-registry.js";
import * as taskRegistryListenerState from "./task-registry-listener-state.js";
import { updateTask } from "./task-registry-mutation.js";
import { prepareTaskRegistryRead } from "./task-registry-read.js";
import {
createReadTask,
requestTasks,
resetReadState,
withReadState,
} from "./task-registry-read.test-support.js";
import { linkTaskToFlowById } from "./task-registry-record-api.js";
import { tasks } from "./task-registry-state.js";
import { getTaskRegistryStore, onTaskRegistryChange } from "./task-registry.store.js";
import { loadTaskRegistryStateFromSqliteReadOnly } from "./task-registry.store.sqlite.js";
afterEach(resetReadState);
describe("task registry terminal read preparation", () => {
it.each(["unchanged", "newer write", "ABA"] as const)(
"keeps registered tasks.list current after terminal publication is superseded by %s",
async (change) => {
await withReadState(async () => {
const runId = `superseded-terminal-${change}`;
const task = createReadTask(runId);
const flow = expectDefined(createTaskFlowForTask({ task }), "terminal task flow");
expect(linkTaskToFlowById({ taskId: task.taskId, flowId: flow.flowId })).not.toBeNull();
const request = () => requestTasks(task.ownerKey);
const terminalInstalled = createDeferred();
const releaseEffects = createDeferred();
const readCaptured = createDeferred();
const finish = taskMutationEffects.finishTaskMutation;
let held = false;
vi.spyOn(taskMutationEffects, "finishTaskMutation").mockImplementation(async (...args) => {
if (
!held &&
args[3] === task.taskId &&
args[4].operation === "update" &&
tasks.get(task.taskId)?.status === "succeeded"
) {
held = true;
terminalInstalled.resolve();
await releaseEffects.promise;
}
return finish(...args);
});
const captureFence = taskRegistryListenerState.captureTaskRegistryReadFence;
let captureRead = false;
vi.spyOn(taskRegistryListenerState, "captureTaskRegistryReadFence").mockImplementation(
(...args) => {
const pending = captureFence(...args);
if (captureRead) {
readCaptured.resolve();
}
return pending;
},
);
const publications: string[] = [];
const stop = onTaskRegistryChange((event) => {
if (event?.kind === "upserted" && event.task.taskId === task.taskId) {
publications.push(event.task.task);
}
});
const mutation = vi.spyOn(getTaskRegistryStore(), "runAgentEventMutationAsync");
let reading: ReturnType<typeof request> | undefined;
let readResult:
| Promise<PromiseSettledResult<Awaited<ReturnType<typeof request>>>[]>
| undefined;
try {
emitAgentEvent({
runId,
stream: "lifecycle",
data: { phase: "end", endedAt: Date.now() },
});
await withTestTimeout(
terminalInstalled.promise,
5_000,
"Terminal event reached publication",
);
expect(loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)?.status).toBe(
"succeeded",
);
captureRead = true;
reading = request();
readResult = Promise.allSettled([reading]);
await withTestTimeout(
readCaptured.promise,
5_000,
"tasks.list captured the terminal event",
);
if (change !== "unchanged") {
expect(updateTask(task.taskId, { task: "Newer task write" })).not.toBeNull();
if (change === "ABA") {
expect(updateTask(task.taskId, { task: task.task })).not.toBeNull();
}
}
const expectedTitle = change === "newer write" ? "Newer task write" : task.task;
const competingPublications = [...publications];
releaseEffects.resolve();
const [outcome] = await withTestTimeout(readResult, 5_000, "Captured tasks.list settled");
expect(loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)).toMatchObject({
task: expectedTitle,
status: "succeeded",
});
expect(mutation).toHaveBeenCalledOnce();
if (change === "unchanged") {
expect(publications).toContain(task.task);
} else {
expect(publications).toEqual(competingPublications);
}
expect((await request()).mock.calls[0]).toMatchObject([
true,
{ tasks: [{ id: task.taskId, title: expectedTitle, status: "completed" }] },
]);
expect(
outcome,
outcome?.status === "rejected" ? String(outcome.reason) : undefined,
).toMatchObject({
status: "fulfilled",
});
if (outcome?.status === "fulfilled") {
expect(outcome.value.mock.calls[0]).toMatchObject([
true,
{ tasks: [{ id: task.taskId, title: expectedTitle, status: "completed" }] },
]);
}
} finally {
releaseEffects.resolve();
await readResult;
await prepareTaskRegistryRead();
stop();
}
});
},
);
});

View file

@ -1,13 +1,26 @@
// Cron Mcp Cleanup Docker Client tests cover cron mcp cleanup docker client script behavior. // Cron Mcp Cleanup Docker Client tests cover cron mcp cleanup docker client script behavior.
import fs from "node:fs"; import fs from "node:fs";
import os from "node:os";
import path from "node:path"; import path from "node:path";
import { describe, expect, it } from "vitest"; import { setTimeout as delay } from "node:timers/promises";
import { afterEach, describe, expect, it, vi } from "vitest";
import { import {
assertCronFinishedOk, assertCronFinishedOk,
readCronMcpCleanupProbePidWaitMs, readCronMcpCleanupProbePidWaitMs,
waitForProbePid, waitForProbePid,
} from "../../scripts/e2e/cron-mcp-cleanup-docker-client.ts"; } from "../../scripts/e2e/cron-mcp-cleanup-docker-client.ts";
import { useAutoCleanupTempDirTracker } from "../helpers/temp-dir.js";
vi.mock("node:timers/promises", async (importOriginal) => ({
...(await importOriginal<typeof import("node:timers/promises")>()),
setTimeout: vi.fn(),
}));
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
afterEach(() => {
vi.useRealTimers();
vi.mocked(delay).mockReset();
});
describe("cron MCP cleanup docker client", () => { describe("cron MCP cleanup docker client", () => {
it("rejects malformed probe pid wait limits", () => { it("rejects malformed probe pid wait limits", () => {
@ -24,31 +37,22 @@ describe("cron MCP cleanup docker client", () => {
} }
}); });
it("bounds missing probe pid waits", async () => { it.each(["missing", "malformed"])("bounds %s probe pid waits", async (fixture) => {
const root = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-cron-mcp-client-")); const root = tempDirs.make("openclaw-cron-mcp-client-");
try { const pidPath = path.join(root, "probe.pid");
const startedAt = Date.now(); if (fixture === "malformed") {
await expect(
waitForProbePid(path.join(root, "missing.pid"), { pollMs: 1, timeoutMs: 20 }),
).resolves.toBeUndefined();
expect(Date.now() - startedAt).toBeLessThan(1000);
} finally {
fs.rmSync(root, { force: true, recursive: true });
}
});
it("does not parse malformed probe pid prefixes", async () => {
const root = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-cron-mcp-client-"));
try {
const pidPath = path.join(root, "probe.pid");
fs.writeFileSync(pidPath, "123abc\n", "utf8"); fs.writeFileSync(pidPath, "123abc\n", "utf8");
const startedAt = Date.now();
await expect(waitForProbePid(pidPath, { pollMs: 1, timeoutMs: 20 })).resolves.toBeUndefined();
expect(Date.now() - startedAt).toBeLessThan(1000);
} finally {
fs.rmSync(root, { force: true, recursive: true });
} }
vi.useFakeTimers({ toFake: ["Date"] });
vi.setSystemTime(0);
vi.mocked(delay).mockImplementation(async (ms) => {
expect(ms).toBe(1);
expect(Date.now(), "polling must stop at the configured deadline").toBeLessThan(20);
vi.setSystemTime(Date.now() + ms!);
});
await expect(waitForProbePid(pidPath, { pollMs: 1, timeoutMs: 20 })).resolves.toBeUndefined();
expect(Date.now()).toBe(20);
}); });
it("accepts cron finished events only when the run status is ok", () => { it("accepts cron finished events only when the run status is ok", () => {

View file

@ -299,14 +299,24 @@ describe("gateway network client", () => {
}); });
it("rejects frame waits immediately when the socket closes", async () => { it("rejects frame waits immediately when the socket closes", async () => {
const ws = new EventEmitter(); vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] });
const startedAt = Date.now(); try {
const frame = onceFrame(ws, () => false, 1000); const ws = new EventEmitter();
const rejected = vi.fn();
const frame = onceFrame(ws, () => false, 1000).catch(rejected);
expect(vi.getTimerCount()).toBe(1);
ws.emit("close", 1006, Buffer.from("bye")); ws.emit("close", 1006, Buffer.from("bye"));
await vi.advanceTimersByTimeAsync(0);
await expect(frame).rejects.toThrow("closed before frame: 1006 bye"); expect(rejected).toHaveBeenCalledExactlyOnceWith(new Error("closed before frame: 1006 bye"));
expect(Date.now() - startedAt).toBeLessThan(250); expect(ws.eventNames()).toEqual([]);
expect(vi.getTimerCount()).toBe(0);
await frame;
} finally {
await vi.runAllTimersAsync();
vi.useRealTimers();
}
}); });
it("rejects frame waits immediately on socket errors", async () => { it("rejects frame waits immediately on socket errors", async () => {

View file

@ -1,11 +1,13 @@
// Source File Scan Cache tests cover source file scan cache script behavior. // Source File Scan Cache tests cover source file scan cache script behavior.
import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; import { mkdir, mkdtemp, rm, stat, writeFile } from "node:fs/promises";
import os from "node:os"; import os from "node:os";
import path from "node:path"; import path from "node:path";
import { afterEach, describe, expect, it } from "vitest"; import { afterEach, describe, expect, it } from "vitest";
import { collectSourceFileContents } from "../../scripts/lib/source-file-scan-cache.mts"; import { collectSourceFileContents } from "../../scripts/lib/source-file-scan-cache.mts";
import { createDeferred } from "../helpers/promise.js";
const tempDirs: string[] = []; const tempDirs: string[] = [];
let pendingScan: ReturnType<typeof collectSourceFileContents> | undefined;
async function makeTempRepo() { async function makeTempRepo() {
const repoRoot = await mkdtemp(path.join(os.tmpdir(), "openclaw-source-scan-")); const repoRoot = await mkdtemp(path.join(os.tmpdir(), "openclaw-source-scan-"));
@ -15,10 +17,13 @@ async function makeTempRepo() {
describe("source file scan cache", () => { describe("source file scan cache", () => {
afterEach(async () => { afterEach(async () => {
// Native test timeout releases held reads; join them before removing their files.
await Promise.allSettled(pendingScan ? [pendingScan] : []);
pendingScan = undefined;
await Promise.all(tempDirs.splice(0).map((dir) => rm(dir, { recursive: true, force: true }))); await Promise.all(tempDirs.splice(0).map((dir) => rm(dir, { recursive: true, force: true })));
}); });
it("bounds concurrent source file reads while preserving sorted output", async () => { it("bounds concurrent source file reads while preserving sorted output", async ({ signal }) => {
const repoRoot = await makeTempRepo(); const repoRoot = await makeTempRepo();
const srcRoot = path.join(repoRoot, "src"); const srcRoot = path.join(repoRoot, "src");
await mkdir(srcRoot, { recursive: true }); await mkdir(srcRoot, { recursive: true });
@ -29,35 +34,69 @@ describe("source file scan cache", () => {
}), }),
); );
let activeReads = 0; let activeFiles = 0;
let maxActiveReads = 0; let maxActiveFiles = 0;
const reads = Array.from({ length: 9 }, (_, index) => ({
name: `file-${index}.ts`,
started: createDeferred(),
release: createDeferred(),
completed: createDeferred(),
}));
const releaseReads = () => {
for (const read of reads) {
read.release.resolve();
}
};
const readFile = async (filePath: string) => { const readFile = async (filePath: string) => {
activeReads += 1; const read = reads.find((entry) => entry.name === path.basename(filePath))!;
maxActiveReads = Math.max(maxActiveReads, activeReads); read.started.resolve();
await new Promise((resolve) => { await read.release.promise;
setTimeout(resolve, 10); activeFiles -= 1;
}); read.completed.resolve();
activeReads -= 1;
return `content:${path.basename(filePath)}`; return `content:${path.basename(filePath)}`;
}; };
const files = await collectSourceFileContents({ signal.throwIfAborted();
signal.addEventListener("abort", releaseReads, { once: true });
const scan = (pendingScan = collectSourceFileContents({
repoRoot, repoRoot,
scanRoots: ["src"], scanRoots: ["src"],
scanExtensions: new Set([".ts"]), scanExtensions: new Set([".ts"]),
ignoredDirNames: new Set(), ignoredDirNames: new Set(),
maxConcurrentReads: 3, maxConcurrentReads: 3,
statFile: (filePath) => {
activeFiles += 1;
maxActiveFiles = Math.max(maxActiveFiles, activeFiles);
return stat(filePath);
},
readFile, readFile,
}); }));
expect(maxActiveReads).toBeGreaterThan(1); try {
expect(maxActiveReads).toBeLessThanOrEqual(3); for (let offset = 0; offset < reads.length; offset += 3) {
expect(files.map((file) => file.relativeFile)).toEqual( const batch = reads.slice(offset, offset + 3);
Array.from({ length: 9 }, (_, index) => `src/file-${index}.ts`), await Promise.all(batch.map((read) => read.started.promise));
); signal.throwIfAborted();
expect(files.map((file) => file.content)).toEqual( expect(activeFiles).toBe(3);
Array.from({ length: 9 }, (_, index) => `content:file-${index}.ts`), // Complete each admitted batch backwards so completion order cannot stand in for file order.
); for (const read of batch.toReversed()) {
read.release.resolve();
await read.completed.promise;
}
}
const files = await scan;
expect(maxActiveFiles).toBe(3);
expect(files.map((file) => file.relativeFile)).toEqual(
Array.from({ length: 9 }, (_, index) => `src/file-${index}.ts`),
);
expect(files.map((file) => file.content)).toEqual(
Array.from({ length: 9 }, (_, index) => `content:file-${index}.ts`),
);
} finally {
signal.removeEventListener("abort", releaseReads);
releaseReads();
await scan;
}
}); });
it("rejects oversized source files before reading them", async () => { it("rejects oversized source files before reading them", async () => {

View file

@ -229,6 +229,7 @@ export const databaseWorkerCoreTestFiles = [
"src/tasks/task-registry-agent-events.test.ts", "src/tasks/task-registry-agent-events.test.ts",
"src/tasks/task-registry-agent-events.lineage.test.ts", "src/tasks/task-registry-agent-events.lineage.test.ts",
"src/tasks/task-registry-read.test.ts", "src/tasks/task-registry-read.test.ts",
"src/tasks/task-registry-terminal-read.test.ts",
"src/tasks/task-registry-progress-runtime.test.ts", "src/tasks/task-registry-progress-runtime.test.ts",
"src/tasks/task-registry-lifecycle.test.ts", "src/tasks/task-registry-lifecycle.test.ts",
"src/tasks/task-registry-flow-sync.test.ts", "src/tasks/task-registry-flow-sync.test.ts",

View file

@ -1,3 +1,4 @@
import type { Page } from "playwright";
import { expect, it } from "vitest"; import { expect, it } from "vitest";
import { import {
controlUiBundledSettingsStorageKey, controlUiBundledSettingsStorageKey,
@ -6,6 +7,7 @@ import {
} from "../test-helpers/control-ui-e2e.ts"; } from "../test-helpers/control-ui-e2e.ts";
import { import {
createChatFlowE2eSuite, createChatFlowE2eSuite,
captureUiProof,
installMockGateway, installMockGateway,
waitForChatScrollIdle, waitForChatScrollIdle,
} from "./chat-flow.test-support.ts"; } from "./chat-flow.test-support.ts";
@ -13,7 +15,114 @@ import {
const suite = createChatFlowE2eSuite(); const suite = createChatFlowE2eSuite();
const POSITION_RAIL_MIN_TRANSCRIPT_HEIGHT = 360; const POSITION_RAIL_MIN_TRANSCRIPT_HEIGHT = 360;
function readPositionRailGeometry(page: Page) {
return page.locator(".chat-position-rail__marks").evaluate((element) => {
const thread = element.closest(".chat-thread")!;
const current = element.querySelector('[aria-current="true"]');
const activeId = current?.getAttribute("data-position-marker-id") ?? null;
const message = activeId
? thread.querySelector(`.chat-bubble[data-entry-id="${activeId}"]`)
: null;
const marker = current?.getBoundingClientRect();
const viewport = element.getBoundingClientRect();
const reader = thread.getBoundingClientRect();
const bubble = message?.getBoundingClientRect();
const distanceFromEnd = thread.scrollHeight - thread.clientHeight - thread.scrollTop;
return {
activeId,
atEnd: Math.abs(distanceFromEnd) <= 1,
currentVisible: current?.hasAttribute("data-visible") ?? false,
messageInViewport: Boolean(
bubble && bubble.bottom > reader.top && bubble.top < reader.bottom,
),
markerInViewport: Boolean(
marker &&
marker.top >= viewport.top &&
marker.bottom <= viewport.top + element.clientHeight,
),
distanceFromEnd,
transcript: {
scrollTop: thread.scrollTop,
clientHeight: thread.clientHeight,
scrollHeight: thread.scrollHeight,
},
rail: {
scrollTop: element.scrollTop,
clientHeight: element.clientHeight,
top: viewport.top,
bottom: viewport.bottom,
},
marker: marker?.toJSON() ?? null,
bubble: bubble?.toJSON() ?? null,
};
});
}
async function expectPositionRailAtEnd(page: Page) {
try {
await expect
.poll(() => readPositionRailGeometry(page))
.toMatchObject({
atEnd: true,
currentVisible: true,
messageInViewport: true,
markerInViewport: true,
});
} catch (error) {
console.error("[chat-position-rail] geometry", await readPositionRailGeometry(page));
throw error;
}
}
suite.define(() => { suite.define(() => {
it("reveals the current marker when navigation and composer resize share a frame", async () => {
await suite.withPage(
{ colorScheme: "dark", viewport: { width: 1440, height: 900 } },
async ({ page }) => {
await installMockGateway(page, {
historyMessages: Array.from({ length: 80 }, (_, index) => ({
__openclaw: { id: `resize-navigation-${index}`, seq: index + 1 },
role: index % 2 === 0 ? "user" : "assistant",
content: [
{
type: "text",
text: `Conversation checkpoint ${index + 1}: review the notes and confirm the next step.`,
},
],
})),
});
await page.addInitScript(createControlUiMockSameOriginGatewayScript());
await page.goto(`${suite.server.baseUrl}chat`);
await page.locator(".chat-position-rail__track").waitFor();
await waitForChatScrollIdle(page);
await expectPositionRailAtEnd(page);
const transcript = page.locator(".chat-thread");
await transcript.hover();
await page.mouse.wheel(0, -30000);
await expect.poll(() => transcript.evaluate((element) => element.scrollTop)).toBe(0);
await expect
.poll(() => readPositionRailGeometry(page))
.toMatchObject({
currentVisible: true,
messageInViewport: true,
markerInViewport: true,
});
await waitForChatScrollIdle(page);
// Resize through the real input handler, then navigate before observer delivery.
await page.locator(".agent-chat__composer-combobox textarea").evaluate((element) => {
const textarea = element as HTMLTextAreaElement;
textarea.value = "Keep the review notes available.\n".repeat(6);
textarea.dispatchEvent(new Event("input", { bubbles: true }));
const thread = document.querySelector(".chat-thread")!;
thread.scrollTop = thread.scrollHeight;
});
await waitForChatScrollIdle(page);
await captureUiProof(suite, page, "rail-resize-navigation", "settled.png");
await expectPositionRailAtEnd(page);
},
);
});
it.each([ it.each([
{ count: 1, direction: "ltr" }, { count: 1, direction: "ltr" },
{ count: 2, direction: "ltr" }, { count: 2, direction: "ltr" },
@ -87,33 +196,7 @@ suite.define(() => {
const markerForIndex = (index: number) => const markerForIndex = (index: number) =>
marks.locator(`[data-position-marker-id="stable-rail-${index}"]`); marks.locator(`[data-position-marker-id="stable-rail-${index}"]`);
// Wait for the rail to reflect the visible reader before recording its anchor. // Wait for the rail to reflect the visible reader before recording its anchor.
await expect await expectPositionRailAtEnd(page);
.poll(() =>
marks.evaluate((element) => {
const thread = element.closest(".chat-thread")!;
const current = element.querySelector('[aria-current="true"]');
const message = current
? thread.querySelector(
`.chat-bubble[data-entry-id="${current.getAttribute("data-position-marker-id")}"]`,
)
: null;
if (!current?.hasAttribute("data-visible") || !message) {
return false;
}
const marker = current.getBoundingClientRect();
const viewport = element.getBoundingClientRect();
const reader = thread.getBoundingClientRect();
const bubble = message.getBoundingClientRect();
return (
Math.abs(thread.scrollHeight - thread.clientHeight - thread.scrollTop) <= 1 &&
bubble.bottom > reader.top &&
bubble.top < reader.bottom &&
marker.top >= viewport.top &&
marker.bottom <= viewport.top + element.clientHeight
);
}),
)
.toBe(true);
const bounds = () => const bounds = () =>
track.evaluate((element) => element.getBoundingClientRect().toJSON()); track.evaluate((element) => element.getBoundingClientRect().toJSON());
const collapsed = await bounds(); const collapsed = await bounds();

View file

@ -720,12 +720,24 @@ suite.define(() => {
) )
.toBe("The shared design is ready"); .toBe("The shared design is ready");
const continuation = thread.locator('.chat-bubble[data-entry-id="continuation"]'); const continuation = thread.locator('.chat-bubble[data-entry-id="continuation"]');
await continuation.evaluate((element) => { await thread.hover();
const continuationDelta = await continuation.evaluate((element) => {
const root = element.closest<HTMLElement>(".chat-thread")!; const root = element.closest<HTMLElement>(".chat-thread")!;
const rect = element.getBoundingClientRect(); const rect = element.getBoundingClientRect();
root.scrollTop += return (
rect.top - root.getBoundingClientRect().top + rect.height / 2 - root.clientHeight / 2; rect.top - root.getBoundingClientRect().top + rect.height / 2 - root.clientHeight / 2
);
}); });
await page.mouse.wheel(0, continuationDelta);
await expect
.poll(() =>
continuation.evaluate((element) => {
const rect = element.getBoundingClientRect();
const viewport = element.closest(".chat-thread")!.getBoundingClientRect();
return rect.top < viewport.top && rect.bottom > viewport.bottom;
}),
)
.toBe(true);
await expect.poll(() => runMarker.getAttribute("aria-current")).toBe("true"); await expect.poll(() => runMarker.getAttribute("aria-current")).toBe("true");
await expect.poll(() => runMarker.getAttribute("data-visible")).toBe(""); await expect.poll(() => runMarker.getAttribute("data-visible")).toBe("");
expect( expect(
@ -738,17 +750,7 @@ suite.define(() => {
element.closest(".chat-thread")!.getBoundingClientRect().top, element.closest(".chat-thread")!.getBoundingClientRect().top,
), ),
).toBe(true); ).toBe(true);
expect(
await continuation.evaluate((element) => {
const rect = element.getBoundingClientRect();
const viewport = element.closest(".chat-thread")!.getBoundingClientRect();
return rect.top < viewport.top && rect.bottom > viewport.bottom;
}),
).toBe(true);
const composerInput = page.locator(".agent-chat__composer-combobox textarea"); const composerInput = page.locator(".agent-chat__composer-combobox textarea");
await composerInput.fill(
Array.from({ length: 6 }, (_, index) => `Review note ${index + 1}`).join("\n"),
);
const markerFits = () => const markerFits = () =>
runMarker.evaluate((element) => { runMarker.evaluate((element) => {
const marker = element.getBoundingClientRect(); const marker = element.getBoundingClientRect();
@ -758,6 +760,26 @@ suite.define(() => {
marker.top >= viewport.top && marker.bottom <= viewport.top + scroller.clientHeight marker.top >= viewport.top && marker.bottom <= viewport.top + scroller.clientHeight
); );
}); });
const markerBottomClearance = () =>
runMarker.evaluate((element) => {
const scroller = element.closest(".chat-position-rail__marks")!;
return (
scroller.getBoundingClientRect().top +
scroller.clientHeight -
element.getBoundingClientRect().bottom
);
});
// Resize preserves the reader's rail offset; make its clipping precondition explicit.
await composerInput.focus();
await thread.locator(".chat-position-rail__marks").hover();
await page.mouse.wheel(0, 1 - (await markerBottomClearance()));
await expect
.poll(async () => Math.abs((await markerBottomClearance()) - 1))
.toBeLessThanOrEqual(1);
await expect.poll(markerFits).toBe(true);
await composerInput.fill(
Array.from({ length: 6 }, (_, index) => `Review note ${index + 1}`).join("\n"),
);
await expect.poll(markerFits).toBe(false); await expect.poll(markerFits).toBe(false);
const readerOffset = await thread.evaluate((element) => element.scrollTop); const readerOffset = await thread.evaluate((element) => element.scrollTop);
await thread.hover(); await thread.hover();

View file

@ -0,0 +1,61 @@
import { asNullableRecord } from "@openclaw/normalization-core/record-coerce";
import type { Page } from "playwright";
import type { ApplicationContext } from "../app/context.ts";
import type { SwarmRosterHydrator } from "../lib/sessions/swarm-roster.ts";
import type { MockGatewayControls } from "../test-helpers/control-ui-e2e.ts";
export type SwarmDiagnosticPane = HTMLElement & {
state?: { sessionKey: string; connectionEpoch: number; lastError: string | null };
swarmHydrator?: SwarmRosterHydrator;
};
export type SwarmDiagnosticWindow = Window & {
openclawSwarmDiagnostic?: {
expandedDetails?: Element | null;
expandedEpoch?: number;
};
};
export async function logSwarmDiagnostic(
page: Page,
gateway: MockGatewayControls,
parentKey: string,
) {
const state = await page.evaluate((key) => {
const pane = document.querySelector<SwarmDiagnosticPane>(
"openclaw-chat-pane.chat-pane-cache__pane--active",
);
const app = document.querySelector("openclaw-app") as
| (HTMLElement & { runtime?: { context?: ApplicationContext } })
| null;
const applicationGateway = app?.runtime?.context?.gateway;
const snapshot = applicationGateway?.snapshot;
const parent = pane?.swarmHydrator?.rows.find((row) => row.key === key);
const widget = document.querySelector('[data-test-id="chat-swarm"]');
const details = widget?.querySelector("details");
const diagnostic = (window as SwarmDiagnosticWindow).openclawSwarmDiagnostic;
return {
sessionKey: pane?.state?.sessionKey,
connectionEpoch: pane?.state?.connectionEpoch,
expandedEpoch: diagnostic?.expandedEpoch,
gatewayPhase: snapshot?.phase,
lastErrorPresent: Boolean(snapshot?.lastError || pane?.state?.lastError),
parent: parent && {
status: parent.status,
hasActiveRun: parent.hasActiveRun,
updatedAt: parent.updatedAt,
},
detailsOpen: details?.open,
detailsSame: details === diagnostic?.expandedDetails,
outcome: widget?.querySelector(".chat-swarm__outcome")?.textContent?.trim(),
events: applicationGateway?.eventLog.slice(0, 40).map(({ ts, event }) => ({ ts, event })),
};
}, parentKey);
const requests = (await gateway.getRequests())
.filter((request) => ["sessions.describe", "sessions.list"].includes(request.method))
.slice(-30)
.map(({ id, method, params }) => {
const query = asNullableRecord(params);
return { id, method, key: query?.key, spawnedBy: query?.spawnedBy, limit: query?.limit };
});
console.info("[swarm-final-diagnostic] " + JSON.stringify({ ...state, requests }));
}

View file

@ -3,6 +3,11 @@ import { expect, it } from "vitest";
import { createControlUiE2eArtifactDir } from "../test-helpers/control-ui-e2e-artifacts.ts"; import { createControlUiE2eArtifactDir } from "../test-helpers/control-ui-e2e-artifacts.ts";
import { controlUiSessionUrl, installMockGateway } from "../test-helpers/control-ui-e2e.ts"; import { controlUiSessionUrl, installMockGateway } from "../test-helpers/control-ui-e2e.ts";
import { chatSessionListResponse } from "./chat-flow.test-support.ts"; import { chatSessionListResponse } from "./chat-flow.test-support.ts";
import {
logSwarmDiagnostic,
type SwarmDiagnosticPane,
type SwarmDiagnosticWindow,
} from "./chat-swarm-lifecycle-diagnostic.test-support.ts";
import { createControlUiE2eSuite } from "./control-ui-e2e-suite.test-support.ts"; import { createControlUiE2eSuite } from "./control-ui-e2e-suite.test-support.ts";
const suite = createControlUiE2eSuite({ const suite = createControlUiE2eSuite({
@ -127,6 +132,11 @@ suite.define(() => {
) )
.toBe(true); .toBe(true);
const outcomeClearance = await summary.evaluate((element) => { const outcomeClearance = await summary.evaluate((element) => {
const pane = element.closest<SwarmDiagnosticPane>("openclaw-chat-pane");
(window as SwarmDiagnosticWindow).openclawSwarmDiagnostic = {
expandedDetails: element.parentElement,
expandedEpoch: pane?.state?.connectionEpoch,
};
const outcome = element.parentElement?.querySelector(".chat-swarm__outcome"); const outcome = element.parentElement?.querySelector(".chat-swarm__outcome");
if (!outcome) { if (!outcome) {
throw new Error("Expanded Swarm outcome is missing"); throw new Error("Expanded Swarm outcome is missing");
@ -174,13 +184,19 @@ suite.define(() => {
agentId: "main", agentId: "main",
reason: "swarm", reason: "swarm",
}); });
await expect try {
.poll(() => await expect
widget .poll(() =>
.getByText("Child runs finished. Check the conversation for the final response.") widget
.isVisible(), .getByText("Child runs finished. Check the conversation for the final response.")
) .isVisible(),
.toBe(true); )
.toBe(true);
} finally {
await logSwarmDiagnostic(page, gateway, sessionKey).catch(() => {
console.info("[swarm-final-diagnostic] unavailable");
});
}
await page.screenshot({ await page.screenshot({
path: path.join(proofDir, "settled-details.png"), path: path.join(proofDir, "settled-details.png"),
animations: "disabled", animations: "disabled",

View file

@ -31,7 +31,7 @@ describe("conversation position rail", () => {
beforeEach(installTranscriptDomMocks); beforeEach(installTranscriptDomMocks);
afterEach(resetTranscriptTestDom); afterEach(resetTranscriptTestDom);
it.each(["resize", "focus", "focus-resize", "pointer", "reader"] as const)( it.each(["resize", "resize-jump", "end", "focus", "focus-resize", "pointer", "reader"] as const)(
"keeps the reader's rail position through %s updates", "keeps the reader's rail position through %s updates",
async (scenario) => { async (scenario) => {
vi.stubGlobal( vi.stubGlobal(
@ -44,16 +44,21 @@ describe("conversation position rail", () => {
); );
const transcript = createTestTranscript(); const transcript = createTestTranscript();
const container = document.body.appendChild(document.createElement("div")); const container = document.body.appendChild(document.createElement("div"));
const activeMessage = vi.fn(() => "message-79"); const settlesAtEnd = scenario === "end";
const count = settlesAtEnd ? 5 : 80;
const startsAtTop = settlesAtEnd || scenario === "resize-jump";
const activeMessage = vi.fn((): string =>
scenario === "resize-jump" ? "message-0" : settlesAtEnd ? "message-2" : "message-79",
);
const positions = { const positions = {
markers: Array.from({ length: 80 }, (_, index) => ({ markers: Array.from({ length: count }, (_, index) => ({
id: `message-${index}`, id: `message-${index}`,
anchorId: `message-${index}`, anchorId: `message-${index}`,
role: "user" as const, role: "user" as const,
message: message(`message-${index}`, "user", `Checkpoint ${index}`, index + 1), message: message(`message-${index}`, "user", `Checkpoint ${index}`, index + 1),
})), })),
markerIdsByMessageId: new Map( markerIdsByMessageId: new Map(
Array.from({ length: 80 }, (_, index) => [`message-${index}`, `message-${index}`]), Array.from({ length: count }, (_, index) => [`message-${index}`, `message-${index}`]),
), ),
}; };
render( render(
@ -73,12 +78,13 @@ describe("conversation position rail", () => {
const marks = container.querySelector<HTMLElement>(".chat-position-rail__marks")!; const marks = container.querySelector<HTMLElement>(".chat-position-rail__marks")!;
const marker = (index: number) => const marker = (index: number) =>
marks.querySelector<HTMLButtonElement>(`[data-position-marker-id="message-${index}"]`)!; marks.querySelector<HTMLButtonElement>(`[data-position-marker-id="message-${index}"]`)!;
let height = 597; let height = settlesAtEnd ? 668 : 597;
let marksHeight = 283; let scrollHeight = settlesAtEnd ? 700 : 8912;
let marksHeight = settlesAtEnd ? 60 : 283;
let railOffset = 0; let railOffset = 0;
Object.defineProperties(root, { Object.defineProperties(root, {
clientHeight: { configurable: true, get: () => height }, clientHeight: { configurable: true, get: () => height },
scrollHeight: { configurable: true, value: 8912 }, scrollHeight: { configurable: true, get: () => scrollHeight },
}); });
Object.defineProperties(marks, { Object.defineProperties(marks, {
clientHeight: { configurable: true, get: () => marksHeight }, clientHeight: { configurable: true, get: () => marksHeight },
@ -86,11 +92,11 @@ describe("conversation position rail", () => {
configurable: true, configurable: true,
get: () => railOffset, get: () => railOffset,
set: (value: number) => { set: (value: number) => {
railOffset = Math.max(0, Math.min(value, 960 - marksHeight)); railOffset = Math.max(0, Math.min(value, count * 12 - marksHeight));
}, },
}, },
}); });
root.scrollTop = 8315; root.scrollTop = startsAtTop ? 0 : 8315;
const flush = async () => { const flush = async () => {
marks.dispatchEvent(new Event("scroll")); marks.dispatchEvent(new Event("scroll"));
await new Promise<void>((resolve) => { await new Promise<void>((resolve) => {
@ -99,9 +105,35 @@ describe("conversation position rail", () => {
}; };
try { try {
await flush(); await flush();
expect(marks.scrollTop).toBe(677); expect(marks.scrollTop).toBe(startsAtTop ? 0 : 677);
expect(marks.querySelectorAll(".chat-position-rail__marker").length).toBeLessThan(50); expect(marks.querySelectorAll(".chat-position-rail__marker").length).toBeLessThan(50);
if (scenario === "resize") { if (scenario === "end") {
// Initial row measurements settle at the end before later composer growth.
height = 552;
scrollHeight = 552;
activeMessage.mockReturnValue("message-4");
await flush();
expect(marks.scrollTop).toBe(0);
height = 452;
scrollHeight = 486;
root.scrollTop = 34;
marksHeight = 47;
await flush();
expect(marker(4).getAttribute("aria-current")).toBe("true");
expect(marks.scrollTop).toBe(0);
} else if (scenario === "resize-jump") {
// Initial end navigation can share the frame that reveals the composer.
height = 554;
marksHeight = 240;
root.scrollTop = 8358;
activeMessage.mockReturnValue("message-79");
await flush();
expect(marker(79).getAttribute("aria-current")).toBe("true");
expect(Number.parseFloat(marker(79).style.top)).toBeGreaterThanOrEqual(marks.scrollTop);
expect(Number.parseFloat(marker(79).style.top) + 12).toBeLessThanOrEqual(
marks.scrollTop + marks.clientHeight,
);
} else if (scenario === "resize") {
height = 554; height = 554;
marksHeight = 240; marksHeight = 240;
activeMessage.mockReturnValue("message-76"); activeMessage.mockReturnValue("message-76");

View file

@ -337,13 +337,16 @@ class ChatPositionRailDirective extends AsyncDirective {
if (this.followActive) { if (this.followActive) {
this.scheduleLayout(); this.scheduleLayout();
} }
const atEnd = this.resizeScrollTarget?.atEnd ?? previous.anchorToEnd; // A measured end supersedes startup's non-end estimate; smooth follow may still be pending.
const atEnd = this.resizeScrollTarget?.atEnd || previous.anchorToEnd;
const maxOffset = Math.max(0, root.scrollHeight - viewport.height); const maxOffset = Math.max(0, root.scrollHeight - viewport.height);
this.resizeScrollTarget = { this.resizeScrollTarget = {
offset: atEnd ? maxOffset : Math.min(previous.scrollTop, maxOffset), offset: atEnd ? maxOffset : Math.min(previous.scrollTop, maxOffset),
atEnd, atEnd,
}; };
} else if (previous && viewport.scrollTop !== previous.scrollTop) { }
// Navigation and a resize can arrive in the same observer delivery.
if (previous && viewport.scrollTop !== previous.scrollTop) {
const target = this.resizeScrollTarget?.offset; const target = this.resizeScrollTarget?.offset;
// Smooth resize compensation crosses intermediate offsets before its target. // Smooth resize compensation crosses intermediate offsets before its target.
// The transcript input owner above retires it when the reader takes over. // The transcript input owner above retires it when the reader takes over.