From 71cf3a4a16bde46c2a81e4525f71d12bbf4e0a44 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Wed, 23 Sep 2026 04:49:31 -0700 Subject: [PATCH] perf(session-share): keep chat catalog startup responsive (#156400) Share concurrent Session Share node listings under service-owned authority and bound progressive chat-startup waits to five seconds. Retain compatible pending pages and publish completed refreshes through the existing catalog lifecycle, preserving current caller filtering. Targeted metadata, pagination, and non-progress clients retain complete-response semantics. A synthetic six-caller progressive burst behind a 30-second node response improves from 30000 ms p99 and six RPCs to 5000 ms and one RPC. Testbox validation passed 41 focused tests, 64 existing gateway/service tests, typechecks, lint, and architecture checks. Independent review found no remaining actionable issues. --- docs/plugins/session-share.md | 2 + .../session-share/src/session-catalog.test.ts | 278 ++++++++++++++++-- .../session-share/src/session-catalog.ts | 231 +++++++++++---- .../session-catalog.session-share.test.ts | 134 +++++++++ 4 files changed, 576 insertions(+), 69 deletions(-) create mode 100644 src/gateway/server-methods/session-catalog.session-share.test.ts diff --git a/docs/plugins/session-share.md b/docs/plugins/session-share.md index d02db26edad2..5f8b6454699c 100644 --- a/docs/plugins/session-share.md +++ b/docs/plugins/session-share.md @@ -76,6 +76,8 @@ Publication is shared with the receiver's permitted viewers, not just the named The catalog refreshes by polling, not a live transcript stream. The source node must remain connected for listings and reads. Long transcripts are paginated; individual text fields are redacted and clipped when necessary. +Progressive catalog listings used during chat startup wait at most five seconds in the foreground. Concurrent viewers share the node request and receive a pending host, retaining its last good page when available. The refreshed page arrives through the catalog's normal host update. Targeted metadata lookups, pagination, and callers without progress updates still await a complete response. Pending requests and retained pages belong to the receiver's Session Share service; configuration changes, node reconnections, and service retirement invalidate them. + Listings leave cold transcript archives untouched and use any stored title metadata. To read cold history, open the session on the source Gateway first so its normal history owner restores the archive. Each source page also bounds raw transcript reads to 8 MiB; a single larger entry returns an explicit error instead of being silently skipped. Inspect that entry on the source Gateway. ## Attribute the source node diff --git a/extensions/session-share/src/session-catalog.test.ts b/extensions/session-share/src/session-catalog.test.ts index e99c477a14fa..42dfa67d2bfa 100644 --- a/extensions/session-share/src/session-catalog.test.ts +++ b/extensions/session-share/src/session-catalog.test.ts @@ -1,5 +1,6 @@ import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; +import type { OpenClawPluginService } from "openclaw/plugin-sdk/plugin-entry"; import type { PluginRuntime } from "openclaw/plugin-sdk/plugin-runtime"; import { createTestPluginApi } from "openclaw/plugin-sdk/plugin-test-api"; import { createPluginRuntimeMock } from "openclaw/plugin-sdk/plugin-test-runtime"; @@ -43,8 +44,8 @@ const remoteIdentity = { id: "4242", }; -function catalogFixture() { - const config: OpenClawConfig = {}; +async function catalogFixture() { + let config: OpenClawConfig = {}; const list = vi.fn().mockResolvedValue({ nodes: [{ nodeId: "alpha", displayName: " Alpha ", connected: true, commands }], }); @@ -64,15 +65,248 @@ function catalogFixture() { config: { current: () => config }, nodes: { list, invoke }, }); - const catalog = createSessionShareCatalog(createTestPluginApi({ runtime })); + let service: OpenClawPluginService | undefined; + const api = createTestPluginApi({ + runtime, + registerService: (registered) => { + service = registered; + }, + }); + const catalog = createSessionShareCatalog(api); + const serviceContext = { config, logger: api.logger, stateDir: "/unused", invokeNode: invoke }; + await service?.start(serviceContext); return { catalog, list, invoke, + stop: async () => service?.stop?.(serviceContext), + configure: (next: OpenClawConfig) => { + config = next; + }, }; } describe("session-share receiver catalog", () => { + it("serves complete lookups during cache saturation without losing active publications", async () => { + vi.useFakeTimers(); + const fixture = await catalogFixture(); + const gate = createDeferred(); + fixture.list.mockResolvedValue({ + nodes: Array.from({ length: 32 }, (_, index) => ({ + nodeId: `node-${index}`, + connected: true, + commands, + })), + }); + fixture.invoke.mockImplementation(() => gate.promise); + const onHost = vi.fn(); + const publications: Promise[] = []; + try { + const listing = fixture.catalog.list({ + allowPartialResults: true, + onHost, + waitUntil: (work) => publications.push(work), + }); + await vi.advanceTimersByTimeAsync(5_000); + expect(await listing).toHaveLength(32); + const selected = fixture.catalog.list({ hostIds: ["node:node-0"], limitPerHost: 1 }); + await vi.advanceTimersByTimeAsync(0); + expect(fixture.invoke).toHaveBeenCalledTimes(33); + gate.resolve({ sessions: [nativeSession] }); + expect((await selected)[0]).toMatchObject({ sessions: [nativeSession] }); + await Promise.all(publications); + expect(onHost.mock.calls.filter(([host]) => host.sessions.length === 1)).toHaveLength(32); + } finally { + gate.resolve({ sessions: [] }); + await vi.runAllTimersAsync(); + await fixture.stop(); + vi.useRealTimers(); + } + }); + + it.each([ + { label: "cold default", query: {}, warm: false }, + { label: "warm default", query: {}, warm: true }, + { label: "explicit complete", query: { allowPartialResults: false }, warm: false }, + { label: "unsubscribed opt-in", query: { allowPartialResults: true }, warm: false }, + { label: "targeted", query: { hostIds: ["node:alpha"] }, warm: false }, + { + label: "cursor", + query: { cursors: { "node:alpha": sessionCatalogPaging.encodeCursor(20) } }, + warm: false, + }, + ])("returns complete snapshots for $label callers", async ({ query, warm }) => { + vi.useFakeTimers(); + const fixture = await catalogFixture(); + try { + if (warm) { + await fixture.catalog.list(query); + } + const refreshed = { ...nativeSession, name: "Complete refresh" }; + fixture.invoke.mockImplementation(async () => { + await new Promise((resolve) => { + setTimeout(resolve, 6_000); + }); + return { sessions: [refreshed] }; + }); + let settled = 0; + const calls = Array.from({ length: 6 }, () => + fixture.catalog.list(query).then((hosts) => { + settled++; + return hosts; + }), + ); + await vi.advanceTimersByTimeAsync(5_000); + expect(settled).toBe(0); + await vi.advanceTimersByTimeAsync(1_000); + for (const hosts of await Promise.all(calls)) { + expect(hosts[0]).toMatchObject({ sessions: [refreshed] }); + expect(hosts[0]?.pending).toBeUndefined(); + expect(hosts[0]?.error).toBeUndefined(); + } + expect(fixture.invoke).toHaveBeenCalledTimes(warm ? 2 : 1); + } finally { + await vi.runAllTimersAsync(); + await fixture.stop(); + vi.useRealTimers(); + } + }); + + it.each(["success", "failure"])( + "bounds a six-caller cold burst through a slow node %s", + async (outcome) => { + vi.useFakeTimers(); + try { + const fixture = await catalogFixture(); + fixture.invoke.mockImplementation(async () => { + await new Promise((resolve) => { + setTimeout(resolve, 30_000); + }); + if (outcome === "failure") { + throw new Error("node timeout"); + } + return { sessions: [nativeSession] }; + }); + const started = Date.now(); + const elapsed: number[] = []; + const publications: Promise[] = []; + const updates = Array.from({ length: 6 }, () => vi.fn()); + const pending = updates.map((onHost) => + fixture.catalog + .list({ + allowPartialResults: true, + onHost, + waitUntil: (work) => publications.push(work), + }) + .then((hosts) => { + elapsed.push(Date.now() - started); + return hosts; + }), + ); + await vi.advanceTimersByTimeAsync(30_000); + await Promise.all(pending); + await Promise.all(publications); + console.log( + JSON.stringify({ + p99Ms: Math.max(...elapsed), + invocations: fixture.invoke.mock.calls.length, + }), + ); + expect(Math.max(...elapsed)).toBeLessThanOrEqual(5_000); + expect(fixture.invoke).toHaveBeenCalledTimes(1); + for (const update of updates) { + expect(update).not.toHaveBeenCalledWith(expect.objectContaining({ pending: true })); + expect(update).toHaveBeenLastCalledWith( + expect.objectContaining( + outcome === "success" + ? { sessions: [nativeSession] } + : { + sessions: [], + error: { code: "NODE_INVOKE_FAILED", message: expect.any(String) }, + }, + ), + ); + } + } finally { + vi.useRealTimers(); + } + }, + ); + + it.each(["unchanged", "config", "connection", "query"])( + "retains only a compatible page during %s refresh", + async (revision) => { + const fixture = await catalogFixture(); + await fixture.catalog.list({}); + const refreshed = { ...nativeSession, name: "Refreshed" }; + const gate = createDeferred(); + fixture.invoke.mockImplementation(() => gate.promise); + if (revision === "config") { + fixture.configure({ gateway: { port: 12345 } }); + } + if (revision === "connection") { + fixture.list.mockResolvedValue({ + nodes: [{ nodeId: "alpha", connected: true, connectedAtMs: 2, commands }], + }); + } + vi.useFakeTimers(); + try { + const onHost = vi.fn(); + const publications: Promise[] = []; + const pending = fixture.catalog.list({ + allowPartialResults: true, + search: revision === "query" ? "new" : undefined, + onHost, + waitUntil: (work) => publications.push(work), + }); + await vi.advanceTimersByTimeAsync(5_000); + const hosts = await pending; + expect(hosts).toEqual([ + expect.objectContaining({ + pending: true, + sessions: revision === "unchanged" ? [nativeSession] : [], + }), + ]); + gate.resolve({ sessions: [refreshed] }); + await Promise.all(publications); + expect(onHost).toHaveBeenLastCalledWith(expect.objectContaining({ sessions: [refreshed] })); + expect(fixture.invoke).toHaveBeenCalledTimes(2); + } finally { + gate.resolve({ sessions: [] }); + await fixture.stop(); + vi.useRealTimers(); + } + }, + ); + + it("shares matching queries without reusing a caller's cancellation", async () => { + const fixture = await catalogFixture(); + const gate = createDeferred(); + fixture.invoke.mockImplementation(() => gate.promise); + const controller = new AbortController(); + const onHost = vi.fn(); + const publications: Promise[] = []; + const original = fixture.catalog.list({ signal: controller.signal, onHost }); + const other = fixture.catalog.list({ + search: "other", + waitUntil: (work) => publications.push(work), + }); + const follower = fixture.catalog.list({ + allowPartialResults: true, + onHost: vi.fn(), + waitUntil: (work) => publications.push(work), + }); + await follower; + expect(fixture.invoke).toHaveBeenCalledTimes(2); + controller.abort(new Error("caller retired")); + const rejected = expect(original).rejects.toThrow("caller retired"); + gate.resolve({ sessions: [nativeSession] }); + await rejected; + expect((await other)[0]?.sessions).toEqual([nativeSession]); + await Promise.all(publications); + expect(onHost).not.toHaveBeenCalled(); + }); + describe("slow phase diagnostics", () => { beforeEach(() => { diagnostics.enabled = true; @@ -89,7 +323,7 @@ describe("session-share receiver catalog", () => { async (nodeCommandDispatched) => { let clock = 0; vi.spyOn(performance, "now").mockImplementation(() => clock); - const fixture = catalogFixture(); + const fixture = await catalogFixture(); const list = fixture.list.getMockImplementation()!; fixture.list.mockImplementation(async (params) => { clock += 1_200; @@ -113,8 +347,16 @@ describe("session-share receiver catalog", () => { }); }); - const hosts = await fixture.catalog.list({}); - expect(hosts[0]).toMatchObject({ sessions: [], error: { code: "NODE_INVOKE_FAILED" } }); + const onHost = vi.fn(); + const publications: Promise[] = []; + await fixture.catalog.list({ onHost, waitUntil: (work) => publications.push(work) }); + await Promise.all(publications); + expect(onHost).toHaveBeenLastCalledWith( + expect.objectContaining({ + sessions: [], + error: expect.objectContaining({ code: "NODE_INVOKE_FAILED" }), + }), + ); expect(diagnostics.warn.mock.calls).toEqual([ [ "slow Session Share catalog phase", @@ -146,7 +388,7 @@ describe("session-share receiver catalog", () => { } let clock = 0; vi.spyOn(performance, "now").mockImplementation(() => clock); - const fixture = catalogFixture(); + const fixture = await catalogFixture(); const invoke = fixture.invoke.getMockImplementation()!; fixture.invoke.mockImplementation(async (params) => { clock += mode === "fast" ? 999 : 1_200; @@ -164,7 +406,7 @@ describe("session-share receiver catalog", () => { }); it("does not invoke nodes when the owner retires during discovery", async () => { - const fixture = catalogFixture(); + const fixture = await catalogFixture(); const entered = createDeferred(); const release = createDeferred(); const controller = new AbortController(); @@ -184,8 +426,8 @@ describe("session-share receiver catalog", () => { expect(fixture.invoke).not.toHaveBeenCalled(); }); - it("delivers owner retirement to an active node invocation", async () => { - const fixture = catalogFixture(); + it("delivers service retirement to an active shared node invocation", async () => { + const fixture = await catalogFixture(); const entered = createDeferred(); const release = createDeferred(); const controller = new AbortController(); @@ -207,7 +449,7 @@ describe("session-share receiver catalog", () => { const pending = fixture.catalog.list({ signal: controller.signal }); try { await entered.promise; - controller.abort(new Error("catalog owner retired")); + await fixture.stop(); expect(transportRetired).toBe(true); } finally { release.resolve(); @@ -218,7 +460,7 @@ describe("session-share receiver catalog", () => { it.each(["retirement", "publication failure"] as const)( "joins all started node work before rejecting on %s", async (failure) => { - const fixture = catalogFixture(); + const fixture = await catalogFixture(); const entered = createDeferred(); const fast = createDeferred(); const slow = createDeferred(); @@ -276,7 +518,7 @@ describe("session-share receiver catalog", () => { it.each(["openclaw", "node:alpha"])( "namespaces colliding profile claims by the invoked node, not wire domain %s", async (domain) => { - const fixture = catalogFixture(); + const fixture = await catalogFixture(); fixture.list.mockResolvedValue({ nodes: ["alpha", "beta"].map((nodeId) => ({ nodeId, commands, connected: true })), }); @@ -309,7 +551,7 @@ describe("session-share receiver catalog", () => { ); it("publishes eligible hosts progressively, preserving failures and deterministic host order", async () => { - const fixture = catalogFixture(); + const fixture = await catalogFixture(); const slow = createDeferred(); fixture.list.mockResolvedValue({ nodes: [ @@ -352,7 +594,7 @@ describe("session-share receiver catalog", () => { }); it("forwards filtered per-host pagination and uses the request-owned node snapshot", async () => { - const fixture = catalogFixture(); + const fixture = await catalogFixture(); const cursor = sessionCatalogPaging.encodeCursor(20); fixture.invoke.mockResolvedValue({ sessions: [nativeSession], @@ -396,7 +638,7 @@ describe("session-share receiver catalog", () => { { nodeId: "alpha", commands, connected: false }, { nodeId: "alpha", commands: [commands[0]!], connected: true }, ])("denies reads when the paired host is unavailable: %j", async (node) => { - const fixture = catalogFixture(); + const fixture = await catalogFixture(); fixture.list.mockResolvedValue({ nodes: [node] }); await expect( fixture.catalog.read({ hostId: "node:alpha", threadId: nativeSession.threadId }), @@ -414,7 +656,7 @@ describe("session-share receiver catalog", () => { { label: "local adoption", patch: { sessionKey: "agent:main:local" } }, { label: "write capability", patch: { canContinue: true } }, ])("rejects node rows carrying $label", async ({ patch }) => { - const fixture = catalogFixture(); + const fixture = await catalogFixture(); fixture.invoke.mockResolvedValue({ payloadJSON: JSON.stringify({ sessions: [{ ...nativeSession, ...patch }] }), }); @@ -427,7 +669,7 @@ describe("session-share receiver catalog", () => { { sender: { identity: remoteIdentity, label: "x".repeat(201) } }, { unexpected: true }, ])("rejects transcript payload outside the closed wire identity contract: %j", async (patch) => { - const fixture = catalogFixture(); + const fixture = await catalogFixture(); fixture.invoke.mockResolvedValue({ payloadJSON: JSON.stringify({ threadId: nativeSession.threadId, diff --git a/extensions/session-share/src/session-catalog.ts b/extensions/session-share/src/session-catalog.ts index f3336597c934..2cf5aad426a1 100644 --- a/extensions/session-share/src/session-catalog.ts +++ b/extensions/session-share/src/session-catalog.ts @@ -2,9 +2,13 @@ import { areDiagnosticsEnabledForProcess, createSubsystemLogger, } from "openclaw/plugin-sdk/diagnostic-runtime"; -import type { OpenClawPluginApi } from "openclaw/plugin-sdk/plugin-entry"; +import type { + OpenClawPluginApi, + OpenClawPluginServiceContext, +} from "openclaw/plugin-sdk/plugin-entry"; import type { PluginRuntime } from "openclaw/plugin-sdk/plugin-runtime"; import { + publishSessionCatalogHost, sessionCatalogPaging, type SessionCatalogHost, type SessionCatalogProvider, @@ -12,6 +16,7 @@ import { } from "openclaw/plugin-sdk/session-catalog"; import { createSessionCatalogGitHubLinker } from "openclaw/plugin-sdk/session-transcript-runtime"; import { asOptionalRecord } from "openclaw/plugin-sdk/string-coerce-runtime"; +import { withTimeout } from "openclaw/plugin-sdk/time-runtime"; import { sessionShareNodeBinding } from "./config.js"; import { SESSION_SHARE_COMMANDS, @@ -23,6 +28,13 @@ import { parseSessionSharePage, parseSessionShareTranscriptPage } from "./wire.j type CatalogNode = Awaited>["nodes"][number]; type GitHubLinker = ReturnType; type CatalogIdentity = NonNullable["identity"]>; +type CatalogPage = ReturnType; +type NodePage = { + nodeId: string; + connection: CatalogNode["connectedAtMs"]; + page?: CatalogPage; + pending?: Promise; +}; const log = createSubsystemLogger("gateway/session-catalog"); const nodeErrorCodes = new Set([ @@ -136,27 +148,34 @@ function bindSession( } export function createSessionShareCatalog(api: OpenClawPluginApi): SessionCatalogProvider { + const pages = new Map(); + let connectedNodes = new Map(); + const active = new Set>(); + let config: ReturnType | undefined; + let invokeNode: OpenClawPluginServiceContext["invokeNode"]; + let lifetime = new AbortController(); + api.registerService({ + id: "session-share-catalog", + start(ctx) { + lifetime = new AbortController(); + invokeNode = ctx.invokeNode; + }, + async stop() { + invokeNode = undefined; + lifetime.abort(); + pages.clear(); + connectedNodes.clear(); + await Promise.allSettled(active); + }, + }); const bindingFor = (nodeId: string) => sessionShareNodeBinding(api.runtime.config.current(), nodeId); - const invoke = ( - nodeId: string, - command: string, - params: Record, - signal?: AbortSignal, - ) => - api.runtime.nodes.invoke({ - nodeId, - command, - params, - timeoutMs: 30_000, - scopes: ["operator.write"], - signal, - }); - async function listNode( node: CatalogNode, query: Parameters[0], + deadline: number, ): Promise { + const { onHost, waitUntil, signal: callerSignal } = query; const hostId = `node:${node.nodeId}`; const common = { hostId, @@ -172,26 +191,16 @@ export function createSessionShareCatalog(api: OpenClawPluginApi): SessionCatalo error: { code: "NODE_OFFLINE", message: "Paired node is offline" }, }; } - try { - query.signal?.throwIfAborted(); - const cursor = query.cursors?.[hostId]; - if (cursor !== undefined) { - sessionCatalogPaging.decodeCursor(cursor); - } - const raw = await observeCatalogPhase("invoke", () => - invoke( - node.nodeId, - SESSION_SHARE_LIST_COMMAND, - { - limit: sessionCatalogPaging.boundedLimit(query.limitPerHost), - ...(query.search ? { searchTerm: query.search } : {}), - ...(cursor !== undefined ? { cursor } : {}), - }, - query.signal, - ), - ); - query.signal?.throwIfAborted(); - const page = parseSessionSharePage(raw); + const failed = (): SessionCatalogHost => ({ + ...common, + sessions: [], + error: { + code: "NODE_INVOKE_FAILED", + message: + "Cannot list OpenClaw sessions. Check the paired node's session-share configuration and connection.", + }, + }); + const project = (page: CatalogPage): SessionCatalogHost => { const binding = bindingFor(node.nodeId); const linker = binding.owner || binding.linkGitHubIdentities @@ -206,17 +215,107 @@ export function createSessionShareCatalog(api: OpenClawPluginApi): SessionCatalo bindSession(session, hostId, owner, linkParticipant), ), }; + }; + try { + query.signal?.throwIfAborted(); + const cursor = query.cursors?.[hostId]; + if (cursor !== undefined) { + sessionCatalogPaging.decodeCursor(cursor); + } + const invokeListing = invokeNode; + if (!invokeListing) { + return failed(); + } + const params = { + limit: sessionCatalogPaging.boundedLimit(query.limitPerHost), + ...(query.search ? { searchTerm: query.search } : {}), + ...(cursor !== undefined ? { cursor } : {}), + }; + const key = JSON.stringify([node.nodeId, params]); + let entry = pages.get(key); + if (!entry) { + if (pages.size >= 32) { + const evict = [...pages].find(([, candidate]) => !candidate.pending); + if (evict) { + pages.delete(evict[0]); + } + } + entry = { nodeId: node.nodeId, connection: node.connectedAtMs }; + pages.set(key, entry); + } + const publication = entry; + const revision = config; + const signal = lifetime.signal; + const current = () => + !signal.aborted && + api.runtime.config.current() === revision && + connectedNodes.has(node.nodeId) && + connectedNodes.get(node.nodeId) === publication.connection; + const follower = Boolean(publication.pending); + if (!publication.pending) { + const pending = observeCatalogPhase("invoke", () => + invokeListing({ + nodeId: node.nodeId, + command: SESSION_SHARE_LIST_COMMAND, + params, + timeoutMs: 30_000, + signal, + }), + ).then((raw) => { + const page = parseSessionSharePage(raw); + if (current()) { + publication.page = page; + } + return page; + }); + publication.pending = pending; + active.add(pending); + const release = () => { + publication.pending = undefined; + active.delete(pending); + // Bound retained pages without rejecting or retiring admitted source work. + if (pages.size > 32 && pages.get(key) === publication) { + pages.delete(key); + } + }; + void pending.then(release, release); + } + // Hydration borrows service authority, never a requesting connection's node handle. + const completed = publication.pending.then( + (page) => (current() ? project(page) : failed()), + failed, + ); + // Metadata lookups and pagination cannot consume pending host updates. + if (query.allowPartialResults !== true || !onHost || !waitUntil) { + return await completed; + } + const loading = (): SessionCatalogHost => { + publishSessionCatalogHost( + { + waitUntil, + onHost: (host) => { + if (current() && !callerSignal?.aborted) { + return onHost?.(host); + } + }, + }, + completed, + ); + return { + ...(publication.page ? project(publication.page) : { ...common, sessions: [] }), + pending: true, + }; + }; + const remaining = deadline - performance.now(); + if (follower || publication.page || remaining <= 0) { + return loading(); + } + return await withTimeout(completed, remaining, { + message: "Session Share catalog is still loading", + }).catch(loading); } catch { query.signal?.throwIfAborted(); - return { - ...common, - sessions: [], - error: { - code: "NODE_INVOKE_FAILED", - message: - "Cannot list OpenClaw sessions. Check the paired node's session-share configuration and connection.", - }, - }; + return failed(); } } @@ -226,7 +325,9 @@ export function createSessionShareCatalog(api: OpenClawPluginApi): SessionCatalo supportsProcessHomeIsolation: true, audience: "session-viewers", async list(query) { + const deadline = performance.now() + 5_000; query.signal?.throwIfAborted(); + const currentConfig = api.runtime.config.current(); let nodes: CatalogNode[]; try { nodes = ( @@ -240,6 +341,26 @@ export function createSessionShareCatalog(api: OpenClawPluginApi): SessionCatalo return []; } query.signal?.throwIfAborted(); + if (api.runtime.config.current() !== currentConfig) { + return []; + } + if (config !== currentConfig) { + config = currentConfig; + pages.clear(); + } + connectedNodes = new Map( + nodes + .filter((node) => node.connected && isSessionHost(node)) + .map((node) => [node.nodeId, node.connectedAtMs]), + ); + for (const [key, entry] of pages) { + if ( + !connectedNodes.has(entry.nodeId) || + connectedNodes.get(entry.nodeId) !== entry.connection + ) { + pages.delete(key); + } + } const requested = query.hostIds ? new Set(query.hostIds) : undefined; const eligible = nodes .filter( @@ -252,9 +373,11 @@ export function createSessionShareCatalog(api: OpenClawPluginApi): SessionCatalo ) .slice(0, 32); const pending = eligible.map(async (node) => { - const host = await listNode(node, query); + const host = await listNode(node, query, deadline); query.signal?.throwIfAborted(); - query.onHost?.(host); + if (!host.pending) { + query.onHost?.(host); + } return host; }); let hosts: SessionCatalogHost[]; @@ -281,10 +404,16 @@ export function createSessionShareCatalog(api: OpenClawPluginApi): SessionCatalo "OpenClaw session node is unavailable. Reconnect it and refresh the catalog.", ); } - const raw = await invoke(nodeId, SESSION_SHARE_READ_COMMAND, { - threadId: request.threadId, - limit: sessionCatalogPaging.boundedLimit(request.limit), - ...(request.cursor !== undefined ? { cursor: request.cursor } : {}), + const raw = await api.runtime.nodes.invoke({ + nodeId, + command: SESSION_SHARE_READ_COMMAND, + timeoutMs: 30_000, + scopes: ["operator.write"], + params: { + threadId: request.threadId, + limit: sessionCatalogPaging.boundedLimit(request.limit), + ...(request.cursor !== undefined ? { cursor: request.cursor } : {}), + }, }); const page = parseSessionShareTranscriptPage(raw, request.threadId); const linkParticipant = bindingFor(nodeId).linkGitHubIdentities diff --git a/src/gateway/server-methods/session-catalog.session-share.test.ts b/src/gateway/server-methods/session-catalog.session-share.test.ts new file mode 100644 index 000000000000..82d0ac4ee1ed --- /dev/null +++ b/src/gateway/server-methods/session-catalog.session-share.test.ts @@ -0,0 +1,134 @@ +import type { OpenClawPluginApi, OpenClawPluginService } from "openclaw/plugin-sdk/plugin-entry"; +import { createTestPluginApi } from "openclaw/plugin-sdk/plugin-test-api"; +import { createPluginRuntimeMock } from "openclaw/plugin-sdk/plugin-test-runtime"; +import { expect, it, vi } from "vitest"; +import { loadBundledPluginFacade } from "../../test-utils/bundled-plugin-public-surface.js"; +import { + bindPluginRegistryRuntime, + hoisted, + resetSessionCatalogTestState, + startCall, + type PluginRegistry, +} from "./session-catalog.test-helpers.js"; + +const { default: sessionSharePlugin } = await loadBundledPluginFacade<{ + default: { register: (api: OpenClawPluginApi) => void }; +}>({ pluginId: "session-share", artifactBasename: "index.js" }); + +it("bounds six Gateway connections and filters the shared refresh at each delivery", async () => { + resetSessionCatalogTestState(); + vi.useFakeTimers(); + const row = { + threadId: "agent:main:shared", + name: "Shared session", + status: "idle", + archived: false, + canContinue: false, + canArchive: false, + }; + const invokeNode = vi.fn(async () => { + await new Promise((resolve) => { + setTimeout(resolve, 30_000); + }); + return { sessions: [row] }; + }); + const config = {}; + const runtime = createPluginRuntimeMock({ + config: { current: () => config }, + nodes: { + list: async () => ({ + nodes: [ + { + nodeId: "source", + connected: true, + commands: ["openclaw.sessions.list.v1", "openclaw.sessions.read.v1"], + }, + ], + }), + invoke: async () => { + throw new Error("must use service authority"); + }, + }, + }); + let service: OpenClawPluginService | undefined; + const api = createTestPluginApi({ + runtime, + registerService: (registered) => { + service = registered; + }, + registerSessionCatalog: (provider) => { + hoisted.activeRegistry.sessionCatalogs = [{ provider }]; + }, + }); + sessionSharePlugin.register(api); + const context = { config, logger: api.logger, stateDir: "/unused", invokeNode }; + await service?.start(context); + bindPluginRegistryRuntime(hoisted.activeRegistry as PluginRegistry, runtime); + hoisted.hasMultipleSessionSharingIdentities.mockReturnValue(true); + const clients = Array.from({ length: 6 }, (_, index) => ({ + connId: `viewer-${index}`, + connect: { scopes: ["operator.admin"] }, + })); + const broadcasts = clients.map(() => vi.fn()); + const elapsed: number[] = []; + const started = Date.now(); + try { + const calls = clients.map((client, index) => + startCall( + "sessions.catalog.list", + { catalogId: "openclaw", progressId: `progress-${index}`, allowPartialResults: true }, + config, + client, + { broadcastToConnIds: broadcasts[index] }, + ), + ); + const done = Promise.all( + calls.map(async (call) => { + await call.completion; + elapsed.push(Date.now() - started); + }), + ); + await vi.advanceTimersByTimeAsync(5_000); + expect(elapsed).toHaveLength(6); + expect(Math.max(...elapsed)).toBeLessThanOrEqual(5_000); + expect(invokeNode).toHaveBeenCalledTimes(1); + for (const call of calls) { + expect(call.respond).toHaveBeenCalledWith(true, { + catalogs: [ + expect.objectContaining({ + hosts: [expect.objectContaining({ pending: true, sessions: [] })], + }), + ], + }); + } + clients[5]!.connect.scopes = ["operator.read"]; + await vi.advanceTimersByTimeAsync(25_000); + await done; + for (const [index, broadcast] of broadcasts.entries()) { + expect(broadcast).toHaveBeenLastCalledWith( + "sessions.catalog.host", + expect.objectContaining({ + catalog: expect.objectContaining({ + hosts: [ + expect.objectContaining({ + sessions: index === 5 ? [] : [expect.objectContaining(row)], + }), + ], + }), + }), + new Set([clients[index]!.connId]), + { dropIfSlow: true }, + ); + } + console.log( + JSON.stringify({ + gatewayP99Ms: Math.max(...elapsed), + invocations: invokeNode.mock.calls.length, + }), + ); + } finally { + await vi.runAllTimersAsync(); + await service?.stop?.(context); + vi.useRealTimers(); + } +});