diff --git a/extensions/feishu/src/media-chunk-idle.test.ts b/extensions/feishu/src/media-chunk-idle.test.ts index 027a43e3dd1a..aed563973f39 100644 --- a/extensions/feishu/src/media-chunk-idle.test.ts +++ b/extensions/feishu/src/media-chunk-idle.test.ts @@ -2,6 +2,7 @@ import fs from "node:fs/promises"; import http from "node:http"; import os from "node:os"; import path from "node:path"; +import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; import { captureEnv, withServer } from "openclaw/plugin-sdk/test-env"; import { afterAll, beforeAll, describe, expect, it } from "vitest"; 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 () => { - let serverSawClose = false; + const socketClosed = createDeferred(); await withServer( (req, res) => { res.writeHead(200, { "content-type": "image/jpeg", "content-length": "1048576" }); res.flushHeaders(); - const markClose = () => { - serverSawClose = true; - }; - req.on("close", markClose); - res.on("close", markClose); + req.socket.once("close", () => socketClosed.resolve()); }, async (baseUrl) => { const stalled = await getHttpReadable(`${baseUrl}/media`); @@ -73,10 +70,7 @@ describe("saveMediaStreamWithIdleTimeout", () => { chunkTimeoutMs: 50, }); expect(stalled.destroyed).toBe(true); - await new Promise((resolve) => { - setTimeout(resolve, 20); - }); - expect(serverSawClose).toBe(true); + await socketClosed.promise; }, ); }); diff --git a/extensions/matrix/src/matrix/client.test.ts b/extensions/matrix/src/matrix/client.test.ts index ab92c6a9f5d4..33f8530b73f4 100644 --- a/extensions/matrix/src/matrix/client.test.ts +++ b/extensions/matrix/src/matrix/client.test.ts @@ -1,5 +1,6 @@ import { expectDefined } from "@openclaw/normalization-core"; import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; +import * as runtimeEnv from "openclaw/plugin-sdk/runtime-env"; // Matrix tests cover client plugin behavior. import { createRequireRecord } from "openclaw/plugin-sdk/test-fixtures"; 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 () => { + vi.useFakeTimers(); + const retryStarted = createDeferred(); + const sleepWithAbort = runtimeEnv.sleepWithAbort; + vi.spyOn(runtimeEnv, "sleepWithAbort").mockImplementation((...args) => { + const sleeping = sleepWithAbort(...args); + retryStarted.resolve(); + return sleeping; + }); matrixDoRequestMock.mockRejectedValueOnce( Object.assign(new TypeError("fetch failed"), { cause: Object.assign(new Error("read ECONNRESET"), { @@ -848,7 +857,6 @@ describe("resolveMatrixAuth", () => { }), ); const abortController = new AbortController(); - const startedAt = Date.now(); const backfillPromise = backfillMatrixAuthDeviceIdAfterStartup({ auth: { accountId: "default", @@ -860,17 +868,21 @@ describe("resolveMatrixAuth", () => { abortSignal: abortController.signal, }); - await vi.waitFor(() => { - expect(matrixDoRequestMock).toHaveBeenCalledTimes(1); - }); - abortController.abort(); + try { + await retryStarted.promise; + expect(vi.getTimerCount()).toBe(1); + abortController.abort(); - // The first retry backoff starts at 250ms; an honored abort returns long before it elapses. - await expect(backfillPromise).resolves.toBeUndefined(); - expect(Date.now() - startedAt).toBeLessThan(200); - expect(matrixDoRequestMock).toHaveBeenCalledTimes(1); - expect(repairCurrentTokenStorageMetaDeviceIdMock).not.toHaveBeenCalled(); - expect(saveBackfilledMatrixDeviceIdMock).not.toHaveBeenCalled(); + // Cancellation must settle the backoff without advancing its clock. + await expect(backfillPromise).resolves.toBeUndefined(); + expect(vi.getTimerCount()).toBe(0); + expect(matrixDoRequestMock).toHaveBeenCalledTimes(1); + expect(repairCurrentTokenStorageMetaDeviceIdMock).not.toHaveBeenCalled(); + expect(saveBackfilledMatrixDeviceIdMock).not.toHaveBeenCalled(); + } finally { + abortController.abort(); + vi.useRealTimers(); + } }); it("resolves configured accessToken SecretRefs during Matrix auth", async () => { diff --git a/extensions/qa-lab/src/scenario-terminal-ack.test.ts b/extensions/qa-lab/src/scenario-terminal-ack.test.ts new file mode 100644 index 000000000000..2a0646a594ee --- /dev/null +++ b/extensions/qa-lab/src/scenario-terminal-ack.test.ts @@ -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(); + const releaseParent = createDeferred(); + const parentSent = createDeferred(); + 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 | 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[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, + 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", + ); + } + }, + ); +}); diff --git a/src/process/spawn-broker/command-startup.test.ts b/src/process/spawn-broker/command-startup.test.ts index a51ff44e3d48..026dc1e9ce6b 100644 --- a/src/process/spawn-broker/command-startup.test.ts +++ b/src/process/spawn-broker/command-startup.test.ts @@ -1,6 +1,6 @@ import { setTimeout as delay } from "node:timers/promises"; -import { describe, expect, it, vi } from "vitest"; -import { withTestTimeout } from "../../../test/helpers/promise.js"; +import { describe, expect, it, onTestFinished, vi } from "vitest"; +import { createDeferred, withTestTimeout } from "../../../test/helpers/promise.js"; import { isPidDefinitelyDead } from "../../shared/pid-alive.js"; import { runCommandWithTimeout, runExec } from "../exec.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)", async (cooperative) => { const host = createSpawnBrokerHost(); - await host.ready(); const spawnExeca = host.spawnExeca.bind(host); let remote: ReturnType | undefined; - vi.spyOn(host, "spawnExeca").mockImplementation((...args) => { - remote = spawnExeca(...args); - 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 { + onTestFinished(async () => { + vi.useRealTimers(); try { process.kill(host.pid!, "SIGCONT"); } catch {} @@ -66,7 +26,64 @@ describe.skipIf(skipBrokerTests)("command startup cancellation", () => { } await host.close(); 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); }, ); diff --git a/src/realtime-transcription/websocket-session.test.ts b/src/realtime-transcription/websocket-session.test.ts index c4d5593340ba..f4930fbb7333 100644 --- a/src/realtime-transcription/websocket-session.test.ts +++ b/src/realtime-transcription/websocket-session.test.ts @@ -8,22 +8,36 @@ import WebSocket, { WebSocketServer } from "ws"; import { createDeferred, withTestTimeout } from "../../test/helpers/promise.js"; import { createRealtimeTranscriptionWebSocketSession, + type RealtimeTranscriptionWebSocketSessionOptions, type RealtimeTranscriptionWebSocketTransport, } from "./websocket-session.js"; let cleanup: (() => Promise) | undefined; +const sessions = new Set>(); beforeEach(() => { vi.useRealTimers(); }); afterEach(async () => { + for (const session of sessions) { + session.close(); + } + sessions.clear(); vi.restoreAllMocks(); vi.useRealTimers(); await cleanup?.(); cleanup = undefined; }); +function createSession( + options: RealtimeTranscriptionWebSocketSessionOptions, +) { + const session = createRealtimeTranscriptionWebSocketSession(options); + sessions.add(session); + return session; +} + async function createRealtimeServer(params?: { closeOnConnection?: boolean; initialEvent?: unknown; @@ -113,7 +127,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { } }, }); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: {}, url: server.url, @@ -129,7 +143,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { await framesReady.promise; expect(Buffer.concat(frames).toString()).toBe("queuedafter"); expect(session.isConnected()).toBe(true); - session.close(); }); 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", callbacks: {}, url: server.url, @@ -174,7 +187,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { encodeSequence(12_999), Buffer.from([0xaa, 0xbb, 0xcc]), ]); - session.close(); }); it.each([ @@ -189,7 +201,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { async ({ readyOnOpen, initialEvent }) => { const server = await createRealtimeServer({ initialEvent }); const onError = vi.fn(); - const session = createRealtimeTranscriptionWebSocketSession<{ type?: string }>({ + const session = createSession<{ type?: string }>({ providerId: "test", callbacks: { onError }, url: server.url, @@ -208,7 +220,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { await expect(session.connect()).rejects.toThrow("queued audio send failed"); expect(session.isConnected()).toBe(false); expect(onError).toHaveBeenCalledOnce(); - session.close(); }, ); @@ -216,7 +227,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { const server = await createRealtimeServer(); const sentFrames: string[] = []; let shouldFailSecondFrame = true; - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: {}, url: server.url, @@ -237,7 +248,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { await session.connect(); expect(sentFrames).toEqual(["first", "second"]); - session.close(); }); 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", callbacks: {}, url: server.url, @@ -290,7 +300,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { encodeSequence(12_998), encodeSequence(12_999), ]); - session.close(); }); it("discards a large queued audio tail when closed before connecting", async () => { @@ -302,7 +311,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { framesReady.resolve(); }, }); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: {}, url: server.url, @@ -322,7 +331,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { session.sendAudio(Buffer.from("live")); await framesReady.promise; expect(frames).toEqual([Buffer.from("live")]); - session.close(); }); it("keeps replacement sockets owned when retired socket callbacks arrive late", async () => { @@ -333,7 +341,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { const server = await createRealtimeServer({ onConnection: (socket) => connections.push(socket), }); - const session = createRealtimeTranscriptionWebSocketSession<{ text?: string }>({ + const session = createSession<{ text?: string }>({ providerId: "test", callbacks: { onError, onTranscript }, url: server.url, @@ -381,7 +389,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { connections[1]?.send(JSON.stringify({ text: "current transcript" })); await vi.waitFor(() => expect(onTranscript).toHaveBeenCalledWith("current transcript")); - session.close(); }); it("discards superseded asynchronous connection preparation", async () => { @@ -394,7 +401,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { resolveFirstUrl = resolve; }); let connectionAttempt = 0; - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: {}, url: async () => (++connectionAttempt === 1 ? await firstUrl : server.url), @@ -413,7 +420,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { await delay(20); expect(connections).toHaveLength(1); expect(session.isConnected()).toBe(true); - session.close(); }); it("cancels a retired reconnect delay before starting a replacement socket", async () => { @@ -421,7 +427,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { const server = await createRealtimeServer({ onConnection: (socket) => connections.push(socket), }); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: {}, url: server.url, @@ -442,7 +448,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { expect(connections).toHaveLength(2); expect(session.isConnected()).toBe(true); - session.close(); }); 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; type?: string; }>({ @@ -491,7 +496,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { connections[2]?.close(1011, "healthy connection dropped"); await vi.waitFor(() => expect(connections).toHaveLength(4), { timeout: 500 }); await vi.waitFor(() => expect(session.isConnected()).toBe(true)); - session.close(); }); 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", callbacks: { onTranscript: (text) => transcripts.push(text) }, url: server.url, @@ -542,7 +546,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { it("terminates the captured socket when graceful provider shutdown expires", async () => { const server = await createRealtimeServer(); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: {}, url: server.url, @@ -565,7 +569,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { it("terminates once when binary audio reaches the active socket buffer cap", async () => { const onError = vi.fn(); const server = await createRealtimeServer(); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: { onError }, url: server.url, @@ -600,7 +604,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { let providerTransport: RealtimeTranscriptionWebSocketTransport | undefined; const onError = vi.fn(); const server = await createRealtimeServer(); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: { onError }, url: server.url, @@ -633,7 +637,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { it("rejects connect when provider handshake frames exceed the socket buffer cap", async () => { const onError = vi.fn(); const server = await createRealtimeServer(); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: { onError }, url: server.url, @@ -670,7 +674,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { } }, }); - const session = createRealtimeTranscriptionWebSocketSession<{ type?: string }>({ + const session = createSession<{ type?: string }>({ providerId: "test", callbacks: {}, url: server.url, @@ -692,7 +696,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { { type: "session.update" }, { type: "input_audio.append", audio: Buffer.from("queued").toString("base64") }, ]); - session.close(); }); it("resolves async URLs and headers before opening the socket", async () => { @@ -702,7 +705,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { seenAuthHeaders.push(headers.authorization); }, }); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: {}, url: async () => server.url, @@ -716,13 +719,12 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { await session.connect(); expect(seenAuthHeaders).toEqual(["Bearer resolved-token"]); - session.close(); }); it("applies the connect timeout while resolving async connection details", async () => { vi.useFakeTimers(); const onError = vi.fn(); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: { onError }, url: () => new Promise(() => {}), @@ -734,23 +736,18 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { }, }); - try { - const connecting = session.connect(); - const timeoutAssertion = expect(connecting).rejects.toThrow( - "test realtime transcription connection timeout", - ); - await vi.advanceTimersByTimeAsync(10); + const connecting = session.connect(); + const timeoutAssertion = expect(connecting).rejects.toThrow( + "test realtime transcription connection timeout", + ); + await vi.advanceTimersByTimeAsync(10); - await timeoutAssertion; - expect(session.isConnected()).toBe(false); - expect(onError).toHaveBeenCalledTimes(1); - const timeoutError = requireFirstMockArg(onError, "connect timeout error"); - expect(timeoutError).toBeInstanceOf(Error); - expect(timeoutError.message).toBe("test realtime transcription connection timeout"); - } finally { - session.close(); - vi.useRealTimers(); - } + await timeoutAssertion; + expect(session.isConnected()).toBe(false); + expect(onError).toHaveBeenCalledTimes(1); + const timeoutError = requireFirstMockArg(onError, "connect timeout error"); + expect(timeoutError).toBeInstanceOf(Error); + expect(timeoutError.message).toBe("test realtime transcription connection timeout"); }); it("preserves connect failures when the error callback throws", async () => { @@ -760,7 +757,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { const onError = vi.fn((_error: Error) => { throw new Error("error observer failed"); }); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: { onError }, url: () => new Promise(() => {}), @@ -786,13 +783,11 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { expect(timeoutError).toBeInstanceOf(Error); expect(timeoutError.message).toBe("test realtime transcription connection timeout"); } finally { - session.close(); if (previousDebugProxyEnabled === undefined) { delete process.env.OPENCLAW_DEBUG_PROXY_ENABLED; } else { process.env.OPENCLAW_DEBUG_PROXY_ENABLED = previousDebugProxyEnabled; } - vi.useRealTimers(); } }); @@ -807,7 +802,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { seenAuthHeaders.push(headers.authorization); }, }); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: {}, url: () => url, @@ -830,7 +825,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { it("rejects provider setup errors before ready", async () => { const server = await createRealtimeServer({ initialEvent: { type: "error", message: "nope" } }); const onError = vi.fn(); - const session = createRealtimeTranscriptionWebSocketSession<{ + const session = createSession<{ type?: string; message?: string; }>({ @@ -859,7 +854,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { const server = await createRealtimeServer({ initialText: "{not json" }); const received = createDeferred(); const onError = vi.fn((_error: Error) => received.resolve()); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: { onError }, url: server.url, @@ -872,16 +867,12 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { }, }); - try { - await session.connect(); - await withTestTimeout(received.promise, 1_000, "Malformed JSON error not received"); - expect(onError).toHaveBeenCalledTimes(1); - const parseError = requireFirstMockArg(onError, "malformed websocket json error"); - expect(parseError).toBeInstanceOf(Error); - expect(parseError.message).toBe("Realtime transcription websocket received malformed JSON."); - } finally { - session.close(); - } + await session.connect(); + await withTestTimeout(received.promise, 1_000, "Malformed JSON error not received"); + expect(onError).toHaveBeenCalledTimes(1); + const parseError = requireFirstMockArg(onError, "malformed websocket json error"); + expect(parseError).toBeInstanceOf(Error); + expect(parseError.message).toBe("Realtime transcription websocket received malformed JSON."); }); it("keeps error callback failures inside websocket message dispatch", async () => { @@ -891,7 +882,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { received.resolve(); throw new Error("error observer failed"); }); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: { onError }, url: server.url, @@ -904,23 +895,19 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { }, }); - try { - await session.connect(); - await withTestTimeout(received.promise, 1_000, "Throwing error observer not reached"); - expect(onError).toHaveBeenCalledTimes(1); - const parseError = requireFirstMockArg(onError, "malformed websocket json error"); - expect(parseError).toBeInstanceOf(Error); - expect(parseError.message).toBe("Realtime transcription websocket received malformed JSON."); - expect(session.isConnected()).toBe(true); - } finally { - session.close(); - } + await session.connect(); + await withTestTimeout(received.promise, 1_000, "Throwing error observer not reached"); + expect(onError).toHaveBeenCalledTimes(1); + const parseError = requireFirstMockArg(onError, "malformed websocket json error"); + expect(parseError).toBeInstanceOf(Error); + expect(parseError.message).toBe("Realtime transcription websocket received malformed JSON."); + expect(session.isConnected()).toBe(true); }); it("reports pre-ready closes separately from connection timeouts", async () => { const server = await createRealtimeServer({ closeOnConnection: true }); const onError = vi.fn(); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: { onError }, url: server.url, @@ -950,7 +937,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { setTimeout(() => ws.close(1011, "flap"), 1); }, }); - const session = createRealtimeTranscriptionWebSocketSession<{ type?: string }>({ + const session = createSession<{ type?: string }>({ providerId: "test", callbacks: { onError }, url: server.url, @@ -978,7 +965,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { { timeout: 1000 }, ); expect(openCount).toBe(4); - session.close(); }); it("refreshes the reconnect budget after a stable ready connection", async () => { @@ -990,7 +976,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { initialEvent: { type: "session.created" }, onConnection: (ws) => connections.push(ws), }); - const session = createRealtimeTranscriptionWebSocketSession<{ type?: string }>({ + const session = createSession<{ type?: string }>({ providerId: "test", callbacks: { onError }, url: server.url, @@ -1025,7 +1011,6 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { ); }); expect(connections).toHaveLength(3); - session.close(); }); it("delivers a legitimate large inbound message below the payload cap", async () => { @@ -1037,7 +1022,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { }); const received = createDeferred(); const onMessage = vi.fn(() => received.resolve()); - const session = createRealtimeTranscriptionWebSocketSession<{ type?: string; text?: string }>({ + const session = createSession<{ type?: string; text?: string }>({ providerId: "test", callbacks: {}, url: server.url, @@ -1048,15 +1033,11 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { }, }); - try { - await session.connect(); - await withTestTimeout(received.promise, 1_000, "Large inbound message not received"); - expect(onMessage).toHaveBeenCalledTimes(1); - const event = requireFirstMockArg(onMessage, "large inbound message"); - expect(event).toEqual({ type: "transcript", text: largeText }); - } finally { - session.close(); - } + await session.connect(); + await withTestTimeout(received.promise, 1_000, "Large inbound message not received"); + expect(onMessage).toHaveBeenCalledTimes(1); + const event = requireFirstMockArg(onMessage, "large inbound message"); + expect(event).toEqual({ type: "transcript", text: largeText }); }); it("drops an oversized inbound message before it reaches the provider parser", async () => { @@ -1069,7 +1050,7 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { const onMessage = vi.fn(() => { throw new Error("oversized frame should not reach provider handler"); }); - const session = createRealtimeTranscriptionWebSocketSession({ + const session = createSession({ providerId: "test", callbacks: { onError }, url: server.url, @@ -1080,17 +1061,13 @@ describe("createRealtimeTranscriptionWebSocketSession", () => { }, }); - try { - await session.connect(); - await withTestTimeout(received.promise, 1_000, "Oversized message error not received"); - expect(onError).toHaveBeenCalledTimes(1); - expect(onMessage).not.toHaveBeenCalled(); - const overflowError = requireFirstMockArg(onError, "oversized inbound message error"); - expect(overflowError).toBeInstanceOf(Error); - expect(overflowError).toHaveProperty("code", "WS_ERR_UNSUPPORTED_MESSAGE_LENGTH"); - expect(overflowError.message).toMatch(/max payload/i); - } finally { - session.close(); - } + await session.connect(); + await withTestTimeout(received.promise, 1_000, "Oversized message error not received"); + expect(onError).toHaveBeenCalledTimes(1); + expect(onMessage).not.toHaveBeenCalled(); + const overflowError = requireFirstMockArg(onError, "oversized inbound message error"); + expect(overflowError).toBeInstanceOf(Error); + expect(overflowError).toHaveProperty("code", "WS_ERR_UNSUPPORTED_MESSAGE_LENGTH"); + expect(overflowError.message).toMatch(/max payload/i); }); }); diff --git a/src/tasks/task-registry-agent-events.ts b/src/tasks/task-registry-agent-events.ts index b3d1a31f42c1..935d87731f3f 100644 --- a/src/tasks/task-registry-agent-events.ts +++ b/src/tasks/task-registry-agent-events.ts @@ -398,6 +398,13 @@ async function persist(pending: PendingEvent): Promise { }, beforeObservers: async (assertCurrentPublication) => { 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); if ( pending.publication.becomesTerminal && @@ -408,16 +415,9 @@ async function persist(pending: PendingEvent): Promise { } await finishTaskMutation(context, store, flowStore, taskId, { operation: "update", - assertCurrent: () => { - assertCurrentPublication(); - if ( - getTaskRegistryStore() !== store || - getTaskFlowRegistryStore() !== flowStore - ) { - throw new Error("Task event publication owners changed"); - } - }, + assertCurrent: assertCurrentOwners, }); + assertCurrentOwners(); flowEffectsSettled = true; } }, diff --git a/src/tasks/task-registry-read.test-support.ts b/src/tasks/task-registry-read.test-support.ts new file mode 100644 index 000000000000..c3c02a59ca0c --- /dev/null +++ b/src/tasks/task-registry-read.test-support.ts @@ -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) { + 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", + }); +} diff --git a/src/tasks/task-registry-read.test.ts b/src/tasks/task-registry-read.test.ts index 6dd49e37deb2..1ac2c962838a 100644 --- a/src/tasks/task-registry-read.test.ts +++ b/src/tasks/task-registry-read.test.ts @@ -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 type { SubagentRunRecord } from "../agents/subagents/registry/subagent-registry.types.js"; import { createSubagentsTool } from "../agents/tools/subagents-tool.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 { emitAgentEvent, resetAgentEventsForTest } from "../infra/agent-events.js"; +import { emitAgentEvent } from "../infra/agent-events.js"; import { SqliteWorkerError } from "../infra/sqlite-worker-contract.js"; -import { - getActiveGatewayRootWorkCount, - getActiveGatewayRootWorkHolders, - resetGatewayWorkAdmission, -} from "../process/gateway-work-admission.js"; import { closeOpenClawStateDatabaseAsync, runOpenClawStateWriteTransaction, } from "../state/openclaw-state-db.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 { createSubagentTaskBackingDetail, resolveManagedTaskBackingDetail, } from "./task-backing-authority.js"; -import * as taskMutationEffects from "./task-executor-create.async.js"; import { createRunningTaskRunCoreWithReceiptAsync } from "./task-executor-create.async.js"; import { completeTaskRunByRunIdCore } from "./task-executor.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 * as taskRegistryListenerState from "./task-registry-listener-state.js"; import { updateTask } from "./task-registry-mutation.js"; +import { publishTaskRecordAfterAtomicStore } from "./task-registry-publication.js"; import { getTaskById, listTaskRecordPage, listFreshTasksForOwnerKey, } from "./task-registry-query.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, 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 { createTaskFixture } from "./task-registry.test-support.js"; -import { - resetTaskFlowRegistryForTests, - resetTaskRegistryForTests, -} from "./task-runtime.test-helpers.js"; +import { configureTaskFlowRegistryRuntime } from "./task-runtime.test-helpers.js"; vi.mock("node:timers/promises", { spy: true }); -afterEach(() => { - vi.restoreAllMocks(); - resetTaskRegistryForTests({ persist: false }); - resetTaskFlowRegistryForTests({ persist: false }); - resetAgentEventsForTest({ preserveListeners: true }); - resetGatewayWorkAdmission(); - subagentRuns.clear(); -}); - -async function withReadState(run: () => Promise) { - 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); - } - }); -} +afterEach(resetReadState); function emitTool(runId: string, name: string) { 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() { const entry: SubagentRunRecord = { runId: "contended-progress-child", @@ -211,117 +134,150 @@ function createReadProgressBatch() { } describe("task registry read preparation", () => { - it.each(["unchanged", "newer write", "ABA"] as const)( - "keeps registered tasks.list current after terminal publication is superseded by %s", - async (change) => { + it.each([ + "receipt", + "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 () => { - 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 = () => requestTaskList(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 task = createReadTask("registered-read-superseded"); + const flowBoundary = boundary === "flow follow-up" || boundary === "flow error"; + const cancellationBoundary = boundary === "retired flow owner during cancellation"; + const consumed = boundary === "native consumption"; + const flowFailure = new Error("Synthetic flow synchronization failure"); + if (flowBoundary || boundary === "retired flow owner" || cancellationBoundary) { + const flow = expectDefined(createTaskFlowForTask({ task }), "task flow"); + expect(linkTaskToFlowById({ taskId: task.taskId, flowId: flow.flowId })).not.toBeNull(); + } + const store = getTaskRegistryStore(); + const mutate = store.runAgentEventMutationAsync.bind(store); + const committed = createDeferred(); + const release = createDeferred(); + const fenced = createDeferred(); + const captureFence = taskAgentEventMutations.captureReadFence.bind(taskAgentEventMutations); + vi.spyOn(taskAgentEventMutations, "captureReadFence").mockImplementation((admission) => { + const result = captureFence(admission); + fenced.resolve(); + return result; }); - const captureFence = taskRegistryListenerState.captureTaskRegistryReadFence; - let captureRead = false; - vi.spyOn(taskRegistryListenerState, "captureTaskRegistryReadFence").mockImplementation( - (...args) => { - const pending = captureFence(...args); - if (captureRead) { - readCaptured.resolve(); + let eventCommitted = false; + const writes = vi + .spyOn(store, "runAgentEventMutationAsync") + .mockImplementation(async (...args) => { + if (consumed) { + committed.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 stop = onTaskRegistryChange((event) => { - if (event?.kind === "upserted" && event.task.taskId === task.taskId) { - publications.push(event.task.task); + const stop = onTaskRegistryChange(() => { + const current = tasks.get(task.taskId); + if (current) { + publications.push(current.task); } }); - const mutation = vi.spyOn(getTaskRegistryStore(), "runAgentEventMutationAsync"); - let reading: ReturnType | undefined; - let readResult: - | Promise>>[]> - | undefined; + const respond = vi.fn(); + let read: ReturnType | 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([ + emitTool(task.runId!, "accepted-tool"); + await committed.promise; + read = requestTasks(task.ownerKey, respond); + await fenced.promise; + expect(respond).not.toHaveBeenCalled(); + if (consumed) { + expect(getTaskById(task.taskId)?.toolUseCount).toBe(1); + configureTaskFlowRegistryRuntime({ store: { ...getTaskFlowRegistryStore() } }); + release.resolve(); + await read; + expect(respond.mock.calls[0]).toMatchObject([ 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 { - releaseEffects.resolve(); - await readResult; - await prepareTaskRegistryRead(); + release.resolve(); + await read?.catch(() => undefined); stop(); } }); @@ -457,6 +413,7 @@ describe("task registry read preparation", () => { ); } expect(await timer).toBe(0); + holder.release(); expect(await pending).toMatchObject( surface === "fresh owner" ? [ @@ -502,48 +459,73 @@ describe("task registry read preparation", () => { } return result; }); - const immediate = timers.setImmediate; - vi.mocked(timers.setImmediate).mockImplementationOnce(async (...args) => { - await immediate(...args); - await committed.promise; + const reader = await import("./task-registry-read.js"); + const prepare = reader.prepareTaskRegistryRead; + // Preparatory yields must not consume the scan's publication barrier. + vi.spyOn(reader, "prepareTaskRegistryRead").mockImplementationOnce(async () => { + await timers.setImmediate(); + return prepare(); }); + const immediate = timers.setImmediate; let workMs = 0; vi.spyOn(performance, "now").mockImplementation(() => workMs); let mutation: Promise | undefined; + let page: ReturnType | undefined; let selectedBeforeMutation = false; + const failures: unknown[] = []; + const recordFailure = (error: unknown) => { + if (!failures.includes(error)) { + failures.push(error); + } + }; try { - const page = await withTestTimeout( - listTaskRecordPage({ - offset: 0, - limit: 1, - prepareFilter: (batch) => { - workMs += 20; - if (!mutation) { - selectedBeforeMutation = batch.some((task) => task.taskId === selected.taskId); - mutation = createRunningTaskRunCoreWithReceiptAsync({ - runtime: selected.runtime, - runId: selected.runId!, - task: selected.task, - ownerKey: selected.ownerKey, - scopeKind: selected.scopeKind, - requesterSessionKey: selected.requesterSessionKey, - notifyPolicy: "silent", - deliveryStatus: "not_applicable", - detail: { historyGeneration: "replacement" }, - }); - } - return (task) => task.taskId === selected.taskId; - }, - }), + page = listTaskRecordPage({ + offset: 0, + limit: 1, + prepareFilter: (batch) => { + workMs += 20; + if (!mutation) { + selectedBeforeMutation = batch.some((task) => task.taskId === selected.taskId); + mutation = createRunningTaskRunCoreWithReceiptAsync({ + runtime: selected.runtime, + runId: selected.runId!, + task: selected.task, + ownerKey: selected.ownerKey, + scopeKind: selected.scopeKind, + requesterSessionKey: selected.requesterSessionKey, + notifyPolicy: "silent", + deliveryStatus: "not_applicable", + detail: { historyGeneration: "replacement" }, + }); + vi.mocked(timers.setImmediate).mockImplementationOnce(async (...args) => { + await immediate(...args); + await committed.promise; + }); + } + return (task) => task.taskId === selected.taskId; + }, + }); + const result = await withTestTimeout( + page, 5_000, "Page joined an identity-changing publication", ); expect(selectedBeforeMutation).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 { + committed.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") { createReadProgressBatch(); } - const request = () => requestTaskList(task.ownerKey); + const request = () => requestTasks(task.ownerKey); emitTool(task.runId!, "warmup"); await prepareTaskRegistryRead(); expect((await request()).mock.calls[0]?.[0]).toBe(true); @@ -605,6 +587,7 @@ describe("task registry read preparation", () => { } read = request(); expect(await timer).toBe(0); + holder.release(); expect((await read).mock.calls[0]).toMatchObject([ true, { @@ -631,115 +614,210 @@ describe("task registry read preparation", () => { }, ); - it("returns the terminal task when accepted metadata publication is superseded", async () => { - 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, - 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>>(); - const release = createDeferred(); - const writes = vi - .spyOn(store, "runAgentEventMutationAsync") - .mockImplementation(async (...args) => { + it.each(["replacement", "ABA"] as const)( + "serves registered tasks.list after a committed event publication loses to %s", + async (change) => { + await withReadState(async () => { + const task = createReadTask(`superseded-read-${change}`); + const store = getTaskRegistryStore(); + const mutate = store.runAgentEventMutationAsync.bind(store); + const snapshot = store.loadMutationSnapshotAsync.bind(store); + const entered = createDeferred(); + const release = createDeferred(); + const reading = createDeferred(); + let committed = false; + let held = false; + vi.spyOn(store, "runAgentEventMutationAsync").mockImplementation(async (...args) => { const receipt = await mutate(...args); - committed.resolve(receipt); - await release.promise; + committed = true; return receipt; }); - 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 | 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(store, "loadMutationSnapshotAsync").mockImplementation(async (...args) => { + const result = await snapshot(...args); + if (committed && !held) { + held = true; + entered.resolve(); + await release.promise; + } + return result; }); - vi.spyOn(taskAgentEventMutations, "captureReadFence").mockImplementation((...args) => { - const fence = capture(...args); - captured.resolve(); + emitTool(task.runId!, "committed-tool"); + const captureFence = taskAgentEventMutations.captureReadFence.bind(taskAgentEventMutations); + vi.spyOn(taskAgentEventMutations, "captureReadFence").mockImplementation((admission) => { + const fence = captureFence(admission); + reading.resolve(); return fence; }); - read = requestTaskList(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({ + let read: ReturnType | undefined; + try { + await entered.promise; + const newer = updateTask(task.taskId, { + task: "Newer title", + ...(change === "replacement" ? { runId: "successor" } : {}), + }); + expect(newer).not.toBeNull(); + if (change === "ABA") { + expect(updateTask(task.taskId, { task: task.task })).not.toBeNull(); + } + read = requestTasks(task.ownerKey); + await reading.promise; + release.resolve(); + expect((await read).mock.calls[0]).toMatchObject([ + true, + { + tasks: [ + { + id: task.taskId, + status: "running", + title: change === "replacement" ? "Newer title" : task.task, + toolUseCount: 1, + }, + ], + }, + ]); + expect(loadTaskRegistryStateFromSqliteReadOnly().tasks.get(task.taskId)).toMatchObject({ + status: "running", + task: change === "replacement" ? "Newer title" : task.task, + toolUseCount: 1, + }); + } 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, childSessionKey: entry.childSessionKey, - status: "succeeded", - terminalSummary: "Authoritative completion", - toolUseCount: 1, - lastToolName: "accepted-before-completion", + requesterSessionKey: entry.requesterSessionKey, + task: entry.task, + notifyPolicy: "silent", + detail: createSubagentTaskBackingDetail(entry.generation!), }); - } finally { - release.resolve(); + const store = getTaskRegistryStore(); + const mutate = store.runAgentEventMutationAsync.bind(store); + const snapshot = store.loadMutationSnapshotAsync.bind(store); + const committed = createDeferred>>(); + 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 | undefined; 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 { - stop(); + release.resolve(); + try { + await read; + } finally { + stop(); + } } - } - }); - }); + }); + }, + ); it.each([1, 8])("prepares %i readers through a fixed event fence", async (readers) => { await withReadState(async () => { @@ -854,50 +932,56 @@ describe("task registry read preparation", () => { }, ); - it.each(["publication", "unknown settlement", "undefined rejection"] as const)( - "does not acknowledge an accepted batch after %s", - async (failureKind) => { - await withReadState(async () => { - const task = createReadTask(`failed-read-${failureKind}`); - const earlier = expectDefined(await prepareTaskRegistryRead(), "earlier task read"); - const store = getTaskRegistryStore(); - const mutate = store.runAgentEventMutationAsync.bind(store); - const snapshot = store.loadMutationSnapshotAsync.bind(store); - const failure = - failureKind === "undefined rejection" + it.each([ + "publication", + "supersession message", + "unknown settlement", + "undefined rejection", + ] as const)("does not acknowledge an accepted batch after %s", async (failureKind) => { + await withReadState(async () => { + const task = createReadTask(`failed-read-${failureKind}`); + const earlier = expectDefined(await prepareTaskRegistryRead(), "earlier task read"); + const store = getTaskRegistryStore(); + 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 : new SqliteWorkerError(`Synthetic ${failureKind}`, "outcome-unknown"); - const rejection = createDeferred(); - let mutationReturned = false; - vi.spyOn(store, "runAgentEventMutationAsync").mockImplementation(async (...args) => { - const result = await mutate(...args); - mutationReturned = true; - if (failureKind !== "publication") { - rejection.reject(failure); - 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"); + const rejection = createDeferred(); + let mutationReturned = false; + vi.spyOn(store, "runAgentEventMutationAsync").mockImplementation(async (...args) => { + const result = await mutate(...args); + mutationReturned = true; + if (!publicationFailure) { + rejection.reject(failure); + return rejection.promise; } + 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 () => { await withReadState(async () => { diff --git a/src/tasks/task-registry-terminal-read.test.ts b/src/tasks/task-registry-terminal-read.test.ts new file mode 100644 index 000000000000..18b8da3adbc1 --- /dev/null +++ b/src/tasks/task-registry-terminal-read.test.ts @@ -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 | undefined; + let readResult: + | Promise>>[]> + | 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(); + } + }); + }, + ); +}); diff --git a/test/scripts/cron-mcp-cleanup-docker-client.test.ts b/test/scripts/cron-mcp-cleanup-docker-client.test.ts index 3cbfde817aca..f2be3c1966c3 100644 --- a/test/scripts/cron-mcp-cleanup-docker-client.test.ts +++ b/test/scripts/cron-mcp-cleanup-docker-client.test.ts @@ -1,13 +1,26 @@ // Cron Mcp Cleanup Docker Client tests cover cron mcp cleanup docker client script behavior. import fs from "node:fs"; -import os from "node:os"; 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 { assertCronFinishedOk, readCronMcpCleanupProbePidWaitMs, waitForProbePid, } 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()), + setTimeout: vi.fn(), +})); + +const tempDirs = useAutoCleanupTempDirTracker(afterEach); + +afterEach(() => { + vi.useRealTimers(); + vi.mocked(delay).mockReset(); +}); describe("cron MCP cleanup docker client", () => { it("rejects malformed probe pid wait limits", () => { @@ -24,31 +37,22 @@ describe("cron MCP cleanup docker client", () => { } }); - it("bounds missing probe pid waits", async () => { - const root = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-cron-mcp-client-")); - try { - const startedAt = Date.now(); - 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"); + it.each(["missing", "malformed"])("bounds %s probe pid waits", async (fixture) => { + const root = tempDirs.make("openclaw-cron-mcp-client-"); + const pidPath = path.join(root, "probe.pid"); + if (fixture === "malformed") { 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", () => { diff --git a/test/scripts/gateway-network-client.test.ts b/test/scripts/gateway-network-client.test.ts index eacd1ddd40fe..682a411b4713 100644 --- a/test/scripts/gateway-network-client.test.ts +++ b/test/scripts/gateway-network-client.test.ts @@ -299,14 +299,24 @@ describe("gateway network client", () => { }); it("rejects frame waits immediately when the socket closes", async () => { - const ws = new EventEmitter(); - const startedAt = Date.now(); - const frame = onceFrame(ws, () => false, 1000); + vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] }); + try { + 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(Date.now() - startedAt).toBeLessThan(250); + expect(rejected).toHaveBeenCalledExactlyOnceWith(new Error("closed before frame: 1006 bye")); + 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 () => { diff --git a/test/scripts/source-file-scan-cache.test.ts b/test/scripts/source-file-scan-cache.test.ts index 106f90d21f77..1e5848fec061 100644 --- a/test/scripts/source-file-scan-cache.test.ts +++ b/test/scripts/source-file-scan-cache.test.ts @@ -1,11 +1,13 @@ // 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 path from "node:path"; import { afterEach, describe, expect, it } from "vitest"; import { collectSourceFileContents } from "../../scripts/lib/source-file-scan-cache.mts"; +import { createDeferred } from "../helpers/promise.js"; const tempDirs: string[] = []; +let pendingScan: ReturnType | undefined; async function makeTempRepo() { const repoRoot = await mkdtemp(path.join(os.tmpdir(), "openclaw-source-scan-")); @@ -15,10 +17,13 @@ async function makeTempRepo() { describe("source file scan cache", () => { 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 }))); }); - 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 srcRoot = path.join(repoRoot, "src"); await mkdir(srcRoot, { recursive: true }); @@ -29,35 +34,69 @@ describe("source file scan cache", () => { }), ); - let activeReads = 0; - let maxActiveReads = 0; + let activeFiles = 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) => { - activeReads += 1; - maxActiveReads = Math.max(maxActiveReads, activeReads); - await new Promise((resolve) => { - setTimeout(resolve, 10); - }); - activeReads -= 1; + const read = reads.find((entry) => entry.name === path.basename(filePath))!; + read.started.resolve(); + await read.release.promise; + activeFiles -= 1; + read.completed.resolve(); return `content:${path.basename(filePath)}`; }; - const files = await collectSourceFileContents({ + signal.throwIfAborted(); + signal.addEventListener("abort", releaseReads, { once: true }); + const scan = (pendingScan = collectSourceFileContents({ repoRoot, scanRoots: ["src"], scanExtensions: new Set([".ts"]), ignoredDirNames: new Set(), maxConcurrentReads: 3, + statFile: (filePath) => { + activeFiles += 1; + maxActiveFiles = Math.max(maxActiveFiles, activeFiles); + return stat(filePath); + }, readFile, - }); + })); - expect(maxActiveReads).toBeGreaterThan(1); - expect(maxActiveReads).toBeLessThanOrEqual(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`), - ); + try { + for (let offset = 0; offset < reads.length; offset += 3) { + const batch = reads.slice(offset, offset + 3); + await Promise.all(batch.map((read) => read.started.promise)); + signal.throwIfAborted(); + expect(activeFiles).toBe(3); + // 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 () => { diff --git a/test/vitest/vitest.database-worker-core-paths.mjs b/test/vitest/vitest.database-worker-core-paths.mjs index dbd53171e993..54bd119f10eb 100644 --- a/test/vitest/vitest.database-worker-core-paths.mjs +++ b/test/vitest/vitest.database-worker-core-paths.mjs @@ -229,6 +229,7 @@ export const databaseWorkerCoreTestFiles = [ "src/tasks/task-registry-agent-events.test.ts", "src/tasks/task-registry-agent-events.lineage.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-lifecycle.test.ts", "src/tasks/task-registry-flow-sync.test.ts", diff --git a/ui/src/e2e/chat-position-rail-layout.e2e.test.ts b/ui/src/e2e/chat-position-rail-layout.e2e.test.ts index 2c7c66257dda..ee021ec5c2fe 100644 --- a/ui/src/e2e/chat-position-rail-layout.e2e.test.ts +++ b/ui/src/e2e/chat-position-rail-layout.e2e.test.ts @@ -1,3 +1,4 @@ +import type { Page } from "playwright"; import { expect, it } from "vitest"; import { controlUiBundledSettingsStorageKey, @@ -6,6 +7,7 @@ import { } from "../test-helpers/control-ui-e2e.ts"; import { createChatFlowE2eSuite, + captureUiProof, installMockGateway, waitForChatScrollIdle, } from "./chat-flow.test-support.ts"; @@ -13,7 +15,114 @@ import { const suite = createChatFlowE2eSuite(); 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(() => { + 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([ { count: 1, direction: "ltr" }, { count: 2, direction: "ltr" }, @@ -87,33 +196,7 @@ suite.define(() => { const markerForIndex = (index: number) => marks.locator(`[data-position-marker-id="stable-rail-${index}"]`); // Wait for the rail to reflect the visible reader before recording its anchor. - await expect - .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); + await expectPositionRailAtEnd(page); const bounds = () => track.evaluate((element) => element.getBoundingClientRect().toJSON()); const collapsed = await bounds(); diff --git a/ui/src/e2e/chat-position-rail.e2e.test.ts b/ui/src/e2e/chat-position-rail.e2e.test.ts index 27d594c9defb..08094b8cd6dd 100644 --- a/ui/src/e2e/chat-position-rail.e2e.test.ts +++ b/ui/src/e2e/chat-position-rail.e2e.test.ts @@ -720,12 +720,24 @@ suite.define(() => { ) .toBe("The shared design is ready"); 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(".chat-thread")!; const rect = element.getBoundingClientRect(); - root.scrollTop += - rect.top - root.getBoundingClientRect().top + rect.height / 2 - root.clientHeight / 2; + return ( + 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("data-visible")).toBe(""); expect( @@ -738,17 +750,7 @@ suite.define(() => { element.closest(".chat-thread")!.getBoundingClientRect().top, ), ).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"); - await composerInput.fill( - Array.from({ length: 6 }, (_, index) => `Review note ${index + 1}`).join("\n"), - ); const markerFits = () => runMarker.evaluate((element) => { const marker = element.getBoundingClientRect(); @@ -758,6 +760,26 @@ suite.define(() => { 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); const readerOffset = await thread.evaluate((element) => element.scrollTop); await thread.hover(); diff --git a/ui/src/e2e/chat-swarm-lifecycle-diagnostic.test-support.ts b/ui/src/e2e/chat-swarm-lifecycle-diagnostic.test-support.ts new file mode 100644 index 000000000000..8fa080ef4ba6 --- /dev/null +++ b/ui/src/e2e/chat-swarm-lifecycle-diagnostic.test-support.ts @@ -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( + "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 })); +} diff --git a/ui/src/e2e/chat-swarm-lifecycle.e2e.test.ts b/ui/src/e2e/chat-swarm-lifecycle.e2e.test.ts index 1726db5f60e3..e813ba950747 100644 --- a/ui/src/e2e/chat-swarm-lifecycle.e2e.test.ts +++ b/ui/src/e2e/chat-swarm-lifecycle.e2e.test.ts @@ -3,6 +3,11 @@ import { expect, it } from "vitest"; import { createControlUiE2eArtifactDir } from "../test-helpers/control-ui-e2e-artifacts.ts"; import { controlUiSessionUrl, installMockGateway } from "../test-helpers/control-ui-e2e.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"; const suite = createControlUiE2eSuite({ @@ -127,6 +132,11 @@ suite.define(() => { ) .toBe(true); const outcomeClearance = await summary.evaluate((element) => { + const pane = element.closest("openclaw-chat-pane"); + (window as SwarmDiagnosticWindow).openclawSwarmDiagnostic = { + expandedDetails: element.parentElement, + expandedEpoch: pane?.state?.connectionEpoch, + }; const outcome = element.parentElement?.querySelector(".chat-swarm__outcome"); if (!outcome) { throw new Error("Expanded Swarm outcome is missing"); @@ -174,13 +184,19 @@ suite.define(() => { agentId: "main", reason: "swarm", }); - await expect - .poll(() => - widget - .getByText("Child runs finished. Check the conversation for the final response.") - .isVisible(), - ) - .toBe(true); + try { + await expect + .poll(() => + widget + .getByText("Child runs finished. Check the conversation for the final response.") + .isVisible(), + ) + .toBe(true); + } finally { + await logSwarmDiagnostic(page, gateway, sessionKey).catch(() => { + console.info("[swarm-final-diagnostic] unavailable"); + }); + } await page.screenshot({ path: path.join(proofDir, "settled-details.png"), animations: "disabled", diff --git a/ui/src/pages/chat/components/chat-position-rail.test.ts b/ui/src/pages/chat/components/chat-position-rail.test.ts index 6db414a69366..0a4cbe2edcd5 100644 --- a/ui/src/pages/chat/components/chat-position-rail.test.ts +++ b/ui/src/pages/chat/components/chat-position-rail.test.ts @@ -31,7 +31,7 @@ describe("conversation position rail", () => { beforeEach(installTranscriptDomMocks); 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", async (scenario) => { vi.stubGlobal( @@ -44,16 +44,21 @@ describe("conversation position rail", () => { ); const transcript = createTestTranscript(); 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 = { - markers: Array.from({ length: 80 }, (_, index) => ({ + markers: Array.from({ length: count }, (_, index) => ({ id: `message-${index}`, anchorId: `message-${index}`, role: "user" as const, message: message(`message-${index}`, "user", `Checkpoint ${index}`, index + 1), })), markerIdsByMessageId: new Map( - Array.from({ length: 80 }, (_, index) => [`message-${index}`, `message-${index}`]), + Array.from({ length: count }, (_, index) => [`message-${index}`, `message-${index}`]), ), }; render( @@ -73,12 +78,13 @@ describe("conversation position rail", () => { const marks = container.querySelector(".chat-position-rail__marks")!; const marker = (index: number) => marks.querySelector(`[data-position-marker-id="message-${index}"]`)!; - let height = 597; - let marksHeight = 283; + let height = settlesAtEnd ? 668 : 597; + let scrollHeight = settlesAtEnd ? 700 : 8912; + let marksHeight = settlesAtEnd ? 60 : 283; let railOffset = 0; Object.defineProperties(root, { clientHeight: { configurable: true, get: () => height }, - scrollHeight: { configurable: true, value: 8912 }, + scrollHeight: { configurable: true, get: () => scrollHeight }, }); Object.defineProperties(marks, { clientHeight: { configurable: true, get: () => marksHeight }, @@ -86,11 +92,11 @@ describe("conversation position rail", () => { configurable: true, get: () => railOffset, 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 () => { marks.dispatchEvent(new Event("scroll")); await new Promise((resolve) => { @@ -99,9 +105,35 @@ describe("conversation position rail", () => { }; try { await flush(); - expect(marks.scrollTop).toBe(677); + expect(marks.scrollTop).toBe(startsAtTop ? 0 : 677); 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; marksHeight = 240; activeMessage.mockReturnValue("message-76"); diff --git a/ui/src/pages/chat/components/chat-position-rail.ts b/ui/src/pages/chat/components/chat-position-rail.ts index c5ca0c79bdf8..d48502bb8c69 100644 --- a/ui/src/pages/chat/components/chat-position-rail.ts +++ b/ui/src/pages/chat/components/chat-position-rail.ts @@ -337,13 +337,16 @@ class ChatPositionRailDirective extends AsyncDirective { if (this.followActive) { 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); this.resizeScrollTarget = { offset: atEnd ? maxOffset : Math.min(previous.scrollTop, maxOffset), 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; // Smooth resize compensation crosses intermediate offsets before its target. // The transcript input owner above retires it when the reader takes over.