diff --git a/apps/shared/OpenClawKit/Sources/OpenClawProtocol/GatewayModels.swift b/apps/shared/OpenClawKit/Sources/OpenClawProtocol/GatewayModels.swift index 3cbed9d9af26..628986baf2bb 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawProtocol/GatewayModels.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawProtocol/GatewayModels.swift @@ -13551,6 +13551,7 @@ public struct SessionCatalogHost: Codable, Sendable { public let label: String public let kind: AnyCodable public let connected: Bool + public let pending: Bool? public let nodeid: String? public let canstartterminal: Bool? public let sessions: [SessionCatalogSession] @@ -13562,6 +13563,7 @@ public struct SessionCatalogHost: Codable, Sendable { label: String, kind: AnyCodable, connected: Bool, + pending: Bool? = nil, nodeid: String? = nil, canstartterminal: Bool? = nil, sessions: [SessionCatalogSession], @@ -13572,6 +13574,7 @@ public struct SessionCatalogHost: Codable, Sendable { self.label = label self.kind = kind self.connected = connected + self.pending = pending self.nodeid = nodeid self.canstartterminal = canstartterminal self.sessions = sessions @@ -13584,6 +13587,7 @@ public struct SessionCatalogHost: Codable, Sendable { case label case kind case connected + case pending case nodeid = "nodeId" case canstartterminal = "canStartTerminal" case sessions @@ -15973,6 +15977,7 @@ public struct SessionsCatalogListParams: Codable, Sendable { public let cursors: [String: AnyCodable]? public let agentid: String? public let progressid: String? + public let allowpartialresults: Bool? public let search: String? public let limitperhost: Int? public let hostids: [String]? @@ -15983,6 +15988,7 @@ public struct SessionsCatalogListParams: Codable, Sendable { cursors: [String: AnyCodable]? = nil, agentid: String? = nil, progressid: String? = nil, + allowpartialresults: Bool? = nil, search: String? = nil, limitperhost: Int? = nil, hostids: [String]? = nil) @@ -15992,6 +15998,7 @@ public struct SessionsCatalogListParams: Codable, Sendable { self.cursors = cursors self.agentid = agentid self.progressid = progressid + self.allowpartialresults = allowpartialresults self.search = search self.limitperhost = limitperhost self.hostids = hostids @@ -16003,6 +16010,7 @@ public struct SessionsCatalogListParams: Codable, Sendable { case cursors case agentid = "agentId" case progressid = "progressId" + case allowpartialresults = "allowPartialResults" case search case limitperhost = "limitPerHost" case hostids = "hostIds" diff --git a/docs/plugins/codex-supervision.md b/docs/plugins/codex-supervision.md index 0b76adf6e6b5..aafa9e5b3381 100644 --- a/docs/plugins/codex-supervision.md +++ b/docs/plugins/codex-supervision.md @@ -197,9 +197,14 @@ hide results from healthy hosts. All queries of the same local home share one resident index. Initial native hydration uses the existing source failure backoff; completed rows remain in memory and in the reconstructible SQLite snapshot. Normal list requests never -restart discovery after a TTL. Paired nodes retain their separate eight-second -foreground response deadline; upgrade their catalog reader to obtain resident -listing on those hosts too. +restart discovery after a TTL. Progressive sidebar lists reuse the last published +paired-node page for the same query and refresh it in the background. A node +without a matching page gets up to 250 ms to answer; after that its host is marked +pending, preserving visible rows until the existing host update event arrives. +Node disconnects, reconnects, configuration changes, and newer publications +invalidate retained pages. Older refreshes cannot replace a newer publication. +One-shot lists, host-specific lookups, and pagination still await fresh node data +under the existing eight-second response deadline. The sidebar hides the Codex group when it has no visible sessions, including when discovery fails. Normal discovery refreshes continue, so the group appears diff --git a/docs/plugins/sdk-entrypoints/define-plugin-entry.md b/docs/plugins/sdk-entrypoints/define-plugin-entry.md index b935691ff113..f9476be68667 100644 --- a/docs/plugins/sdk-entrypoints/define-plugin-entry.md +++ b/docs/plugins/sdk-entrypoints/define-plugin-entry.md @@ -59,6 +59,21 @@ export default definePluginEntry({ `onHost(host)` callback as each host settles; the returned host array remains required as the final compatibility snapshot. + The optional `allowPartialResults` flag is true only when a connected caller + explicitly opts in while receiving host progress on a list without host selection + or cursors. When true, a provider may return retained host snapshots or mark a + still-loading host `pending: true`, then publish its completed snapshot through + `onHost` and `waitUntil`. Pending hosts preserve existing client rows and cursors; + omitted hosts are removed. Clear `pending` on a completed host. + Each `onHost` publication must be authoritative for that host: the Gateway + includes the latest publication for each retained host in the aggregate response + even when another provider or visibility projection delays delivery. The final + host set is authoritative: an omitted host is withdrawn, not restored from an + earlier publication. Preserve the last known + rows while refreshing; do not publish an empty host to represent pending work. + When the flag is absent or false, return the complete compatibility snapshot. + Targeted host lookups and pagination retain that complete-response contract. + If a host can finish after `list` returns a fail-soft snapshot, register its bounded completion with the optional `waitUntil(completion: Promise)` hook before `list` settles. Include host mapping and the `onHost` call in that @@ -77,7 +92,7 @@ export default definePluginEntry({ grant new authority, or permit starting work after the owner retires. Providers remain responsible for bounded work that settles after cancellation. - Keep `onHost`, `waitUntil`, and `signal` separate from validated catalog query + Keep `allowPartialResults`, `onHost`, `waitUntil`, and `signal` separate from validated catalog query objects and node command payloads. The request-owned `sessionEntries` snapshot and `listNodes` hook must be released when `list` settles, or when the optional list operation below closes. Prepare the facts needed by late host mapping diff --git a/extensions/anthropic/session-catalog-watch.test-support.ts b/extensions/anthropic/session-catalog-watch.test-support.ts new file mode 100644 index 000000000000..9dc56a1baaaf --- /dev/null +++ b/extensions/anthropic/session-catalog-watch.test-support.ts @@ -0,0 +1,71 @@ +import { EventEmitter } from "node:events"; +import fs from "node:fs"; +import path from "node:path"; +import { vi } from "vitest"; + +/** Controls OS event delivery while exercising the real watcher and filesystem cache owners. */ +export function createClaudeCatalogWatchDriver(home: string) { + const watchers = new Map(); + class Watcher extends EventEmitter { + constructor(private readonly root: string) { + super(); + } + close() { + if (watchers.get(this.root)?.watcher === this) { + watchers.delete(this.root); + } + this.removeAllListeners(); + } + ref() { + return this; + } + unref() { + return this; + } + } + const nativeWatch = fs.watch.bind(fs); + vi.spyOn(fs, "watch").mockImplementation( + ( + target: fs.PathLike, + options: fs.WatchOptionsWithStringEncoding | fs.WatchListener, + listener?: fs.WatchListener, + ) => { + const root = String(target); + if (!root.startsWith(`${home}${path.sep}`)) { + return typeof options === "function" + ? nativeWatch(target, options) + : nativeWatch(target, options, listener); + } + const watcher = new Watcher(root); + const callback = listener ?? (typeof options === "function" ? options : undefined); + if (callback) { + watcher.on("change", callback); + } + watchers.set(root, { + watcher, + recursive: typeof options === "object" && options.recursive === true, + }); + return watcher; + }, + ); + let monotonicNow = performance.now(); + vi.spyOn(performance, "now").mockImplementation(() => monotonicNow); + return { + arm: () => { + monotonicNow += 250; + }, + change: (file: string, event: fs.WatchEventType = "change") => { + const absolute = path.resolve(home, file); + const selected = [...watchers] + .filter(([root, { recursive }]) => + recursive ? absolute.startsWith(`${root}${path.sep}`) : path.dirname(absolute) === root, + ) + .toSorted(([left], [right]) => right.length - left.length)[0]; + if (!selected) { + throw new Error(`No catalog watcher covers ${absolute}`); + } + const [root, { watcher }] = selected; + watcher.emit("change", event, path.relative(root, absolute)); + }, + }; +} diff --git a/extensions/anthropic/session-catalog.test.ts b/extensions/anthropic/session-catalog.test.ts index 2184e3290b7d..3f253667dbf6 100644 --- a/extensions/anthropic/session-catalog.test.ts +++ b/extensions/anthropic/session-catalog.test.ts @@ -16,6 +16,7 @@ import { createClaudeSessionNodeInvokePolicies, registerClaudeSessionDiscovery, } from "./session-catalog-registration.js"; +import { createClaudeCatalogWatchDriver } from "./session-catalog-watch.test-support.js"; import { CLAUDE_CLI_NODE_RUN_COMMAND, CLAUDE_SESSIONS_LIST_COMMAND, @@ -1854,6 +1855,7 @@ describe("Claude session catalog", () => { it("does not revive an earlier color when a clear or invalid value is appended", async () => { const home = await createHome(); + const watches = createClaudeCatalogWatchDriver(home); const sessionId = "cleared-color"; await writeProject({ home, @@ -1867,6 +1869,9 @@ describe("Claude session catalog", () => { "-workspace", `${sessionId}.jsonl`, ); + await listLocalClaudeSessionPage({}, home); + watches.arm(); + await listLocalClaudeSessionPage({}, home); for (const agentColor of [ "default", "reset", @@ -1882,24 +1887,23 @@ describe("Claude session catalog", () => { transcriptPath, `${JSON.stringify({ type: "agent-color", agentColor: "red", sessionId })}\n`, ); - await expectClaudeCatalogEventually(home, (page) => - expect(page.sessions[0]?.color).toBe("red"), - ); + watches.change(transcriptPath); + expect((await listLocalClaudeSessionPage({}, home)).sessions[0]?.color).toBe("red"); await fs.appendFile( transcriptPath, `${JSON.stringify({ type: "agent-color", agentColor, sessionId })}\n`, ); + watches.change(transcriptPath); // The plugin passes strings through; the Gateway's palette seam removes invalid names. - await expectClaudeCatalogEventually(home, (page) => - expect(page.sessions[0]?.color).toBe( - typeof agentColor === "string" && agentColor ? agentColor : undefined, - ), + expect((await listLocalClaudeSessionPage({}, home)).sessions[0]?.color).toBe( + typeof agentColor === "string" && agentColor ? agentColor : undefined, ); } }); it("reads appended metadata beyond the prefix budget and refreshes the cached tail", async () => { const home = await createHome(); + const watches = createClaudeCatalogWatchDriver(home); const sessionId = "large-colored-session"; let now = Date.now(); vi.spyOn(Date, "now").mockImplementation(() => now); @@ -1920,6 +1924,8 @@ describe("Claude session catalog", () => { const openSpy = vi.spyOn(fs, "open"); const first = await listLocalClaudeSessionPage({}, home); expect(first.sessions[0]).toMatchObject({ name: "Tail rename", color: "green" }); + watches.arm(); + expect(await listLocalClaudeSessionPage({}, home)).toEqual(first); openSpy.mockClear(); expect(await listLocalClaudeSessionPage({}, home)).toEqual(first); expect(openSpy).not.toHaveBeenCalled(); @@ -1940,24 +1946,21 @@ describe("Claude session catalog", () => { .map((row) => JSON.stringify(row)) .join("\n"), ); - await expectClaudeCatalogEventually(home, (page) => - expect(page.sessions[0]).toMatchObject({ - name: "New tail rename", - color: "orange", - }), - ); + watches.change(transcriptPath); + expect((await listLocalClaudeSessionPage({}, home)).sessions[0]).toMatchObject({ + name: "New tail rename", + color: "orange", + }); await writeDesktopMetadata(home, "colored-cli", { cliSessionId: sessionId, title: "Desktop title", }); now += 60_001; - await expectClaudeCatalogEventually(home, (page) => - expect(page.sessions[0]).toMatchObject({ - name: "Desktop title", - source: "claude-desktop", - color: undefined, - }), - ); + expect((await listLocalClaudeSessionPage({}, home)).sessions[0]).toMatchObject({ + name: "Desktop title", + source: "claude-desktop", + color: undefined, + }); }); it("serves an unchanged assembled scan without reparsing transcript files", async () => { @@ -1998,6 +2001,7 @@ describe("Claude session catalog", () => { it("re-stats only the changed project directory on the next poll", async () => { const home = await createHome(); + const watches = createClaudeCatalogWatchDriver(home); for (const project of ["changed", "untouched"]) { await writeProject({ home, @@ -2009,35 +2013,24 @@ describe("Claude session catalog", () => { }); } const first = await listLocalClaudeSessionPage({}, home); - const armingSpies = (["lstat", "readdir", "open"] as const).map((method) => - vi.spyOn(fs, method), - ); - await expectClaudeCatalogQuiescent( - home, - armingSpies, - (spy) => - spy.mock.calls.filter(([target]) => typeof target === "string" && target.startsWith(home)), - first, - ); - for (const spy of armingSpies) { - spy.mockRestore(); - } + watches.arm(); + expect(await listLocalClaudeSessionPage({}, home)).toEqual(first); const changedDir = path.join(home, ".claude", "projects", "changed"); const changedFile = path.join(changedDir, "changed.jsonl"); const readdir = vi.spyOn(fs, "readdir"); const lstat = vi.spyOn(fs, "lstat"); const open = vi.spyOn(fs, "open"); + expect(await listLocalClaudeSessionPage({}, home)).toEqual(first); await fs.appendFile( changedFile, `${JSON.stringify({ type: "custom-title", sessionId: "changed", customTitle: "Updated title" })}\n`, ); - await expectClaudeCatalogEventually(home, (page) => - expect(page.sessions).toEqual( - expect.arrayContaining([ - expect.objectContaining({ threadId: "changed", name: "Updated title" }), - expect.objectContaining({ threadId: "untouched", name: "untouched" }), - ]), - ), + watches.change(changedFile); + expect((await listLocalClaudeSessionPage({}, home)).sessions).toEqual( + expect.arrayContaining([ + expect.objectContaining({ threadId: "changed", name: "Updated title" }), + expect.objectContaining({ threadId: "untouched", name: "untouched" }), + ]), ); expect(readdir.mock.calls.map(([target]) => target)).toEqual([changedDir]); expect( @@ -2054,6 +2047,7 @@ describe("Claude session catalog", () => { it("keeps the CLI records when only the Desktop store changes", async () => { const home = await createHome(); + const watches = createClaudeCatalogWatchDriver(home); let now = Date.now(); vi.spyOn(Date, "now").mockImplementation(() => now); await writeProject({ @@ -2068,39 +2062,29 @@ describe("Claude session catalog", () => { title: "Desktop before", }); const first = await listLocalClaudeSessionPage({}, home); - const armingSpies = (["stat", "lstat", "readdir", "open"] as const).map((method) => - vi.spyOn(fs, method), - ); - await expectClaudeCatalogQuiescent( - home, - armingSpies, - (spy) => - spy.mock.calls.filter(([target]) => typeof target === "string" && target.startsWith(home)), - first, - ); - for (const spy of armingSpies) { - spy.mockRestore(); - } + watches.arm(); + expect(await listLocalClaudeSessionPage({}, home)).toEqual(first); const readdir = vi.spyOn(fs, "readdir"); const transcriptIo = (["stat", "lstat", "open"] as const).map((method) => vi.spyOn(fs, method)); + expect(await listLocalClaudeSessionPage({}, home)).toEqual(first); await writeDesktopMetadata(home, "overlay", { cliSessionId: "desktop", title: "Desktop after", }); // Desktop is macOS-owned; synthetic stores elsewhere refresh through the same TTL backstop. - if (process.platform !== "darwin") { + if (process.platform === "darwin") { + watches.change( + "Library/Application Support/Claude/claude-code-sessions/account/workspace/local_overlay.json", + ); + } else { now += 60_001; } - await expectClaudeCatalogEventually(home, (page) => - expect(page.sessions[0]).toMatchObject({ - name: "Desktop after", - source: "claude-desktop", - }), - ); + expect((await listLocalClaudeSessionPage({}, home)).sessions[0]).toMatchObject({ + name: "Desktop after", + source: "claude-desktop", + }); now += 60_001; - await expectClaudeCatalogEventually(home, (page) => - expect(page.sessions[0]?.name).toBe("Desktop after"), - ); + expect((await listLocalClaudeSessionPage({}, home)).sessions[0]?.name).toBe("Desktop after"); for (const spy of transcriptIo) { expect( spy.mock.calls.filter( @@ -2647,6 +2631,7 @@ describe("Claude session catalog", () => { it("evicts a deleted transcript after a complete scan", async () => { const home = await createHome(); + const watches = createClaudeCatalogWatchDriver(home); const projectDir = path.join(home, ".claude", "projects", "-workspace"); const sessionId = "deleted-session"; const transcriptPath = path.join(projectDir, `${sessionId}.jsonl`); @@ -2660,10 +2645,13 @@ describe("Claude session catalog", () => { await fs.utimes(projectDir, fixedTime, fixedTime); const originalStat = await fs.stat(transcriptPath); await listLocalClaudeSessionPage({}, home); + watches.arm(); + await listLocalClaudeSessionPage({}, home); await fs.rm(transcriptPath); await fs.utimes(projectDir, fixedTime, fixedTime); - await expectClaudeCatalogEventually(home, (page) => expect(page.sessions).toEqual([])); + watches.change(transcriptPath, "rename"); + expect((await listLocalClaudeSessionPage({}, home)).sessions).toEqual([]); await fs.writeFile(transcriptPath, `${JSON.stringify(sdkCliMessage(sessionId, "Bravo"))}\n`); await fs.utimes(transcriptPath, fixedTime, fixedTime); await fs.utimes(projectDir, fixedTime, fixedTime); @@ -2674,11 +2662,10 @@ describe("Claude session catalog", () => { }); const openSpy = vi.spyOn(fs, "open"); - await expectClaudeCatalogEventually(home, (page) => - expect(page.sessions).toEqual([ - expect.objectContaining({ threadId: sessionId, name: "Bravo" }), - ]), - ); + watches.change(transcriptPath, "rename"); + expect((await listLocalClaudeSessionPage({}, home)).sessions).toEqual([ + expect.objectContaining({ threadId: sessionId, name: "Bravo" }), + ]); expect(openSpy).toHaveBeenCalledTimes(1); }); @@ -3360,7 +3347,7 @@ describe("Claude session catalog", () => { nodes: { list: async () => ({ nodes: [] }) }, } as unknown as PluginRuntime); - await expect(provider.list({})).resolves.toMatchObject([ + await expect(provider.list({ allowPartialResults: true })).resolves.toMatchObject([ { sessions: [{ threadId: sessionId, canOpenTerminal: false }] }, ]); await expect( diff --git a/extensions/anthropic/session-catalog.ts b/extensions/anthropic/session-catalog.ts index 901e28bdd61a..f91fda833be8 100644 --- a/extensions/anthropic/session-catalog.ts +++ b/extensions/anthropic/session-catalog.ts @@ -176,6 +176,7 @@ export function createClaudeSessionCatalogRuntime( const localCliAvailable = catalogTerminal.isClaudeCliAvailable(); const { allowProcessHomeFallback, + allowPartialResults: _allowPartialResults, agentId: _agentId, listNodes, onHost, diff --git a/extensions/codex/src/session-catalog-list-operation.test-support.ts b/extensions/codex/src/session-catalog-list-operation.test-support.ts new file mode 100644 index 000000000000..501638be5836 --- /dev/null +++ b/extensions/codex/src/session-catalog-list-operation.test-support.ts @@ -0,0 +1,148 @@ +import type { SessionCatalogProvider } from "openclaw/plugin-sdk/session-catalog"; +import { vi } from "vitest"; +import type { + CodexSessionCatalogPage, + CodexSessionCatalogPageParams, +} from "./session-catalog-types.js"; +import { + CODEX_APP_SERVER_THREADS_LIST_COMMAND, + config, + createCodexSessionCatalogControlFactory, + createCodexTestBindingStore, + createControl, + createGatewayApi, + createRuntime, + registerCodexSessionCatalog, +} from "./session-catalog.test-helpers.js"; + +export function page(ids: string[], nextCursor?: string): CodexSessionCatalogPage { + return { + sessions: ids.map((threadId) => ({ + threadId, + name: threadId, + status: "idle", + source: "cli", + archived: false, + })), + ...(nextCursor ? { nextCursor } : {}), + }; +} + +export function observe(promise: Promise) { + const state = { settled: false }; + const done = promise + .then( + (value) => ({ status: "fulfilled" as const, value }), + (reason: unknown) => ({ status: "rejected" as const, reason }), + ) + .finally(() => { + state.settled = true; + }); + return { state, done }; +} + +export async function fixture(homeCount = 1) { + const { runtime } = createRuntime(); + const base = createCodexSessionCatalogControlFactory({ + getPluginConfig: () => ({ supervision: { enabled: true } }), + getRuntimeConfig: () => config, + }); + const primary = (await base.homesForAgent("main"))[0]!; + const homes = Array.from({ length: homeCount }, (_, index) => ({ + ...primary, + sourceHomeId: `home-${index}`, + hostId: index === 0 ? "gateway:local" : `gateway:local:home-${index}`, + label: `Home ${index}`, + })); + const listPage = + vi.fn< + (homeId: string, params: CodexSessionCatalogPageParams) => Promise + >(); + listPage.mockResolvedValue(page(["visible"])); + const snapshot = vi.fn( + async () => new Map(homes.map((home) => [home.sourceHomeId, new Set(["managed"])])), + ); + const bindingStore = Object.assign(createCodexTestBindingStore(), { + managedThreads: { has: vi.fn(async () => false), mark: vi.fn(async () => true), snapshot }, + }); + const { api, getProvider } = createGatewayApi(runtime, config); + let runtimeConfig = config; + registerCodexSessionCatalog({ + api, + bindingStore, + control: { + ...base, + homesForAgent: async () => homes, + forRequest: (_agentId, source) => + createControl({ + listPage: (params) => listPage(source!.sourceHomeId, params), + }), + }, + getRuntimeConfig: () => runtimeConfig, + }); + const controller = new AbortController(); + const onHost = vi.fn(); + const publications: Promise[] = []; + const start = (params: Partial[0]> = {}) => { + const provider = getProvider()!; + if (!provider.createListOperation) { + throw new Error("Codex list operation is unavailable"); + } + return provider.createListOperation({ + agentId: "main", + limitPerHost: 1, + hostIds: homes.map((home) => home.hostId), + signal: controller.signal, + onHost, + waitUntil: (completion) => { + publications.push(completion); + }, + ...params, + }); + }; + return { + runtime, + homes, + listPage, + snapshot, + controller, + onHost, + publications, + start, + replaceConfig: () => { + runtimeConfig = structuredClone(config); + }, + }; +} + +export async function nodeFixture() { + const f = await fixture(); + const node = { + nodeId: "remote", + connected: true, + connectedAtMs: 1, + commands: [CODEX_APP_SERVER_THREADS_LIST_COMMAND], + }; + const listNodes = vi.fn(async () => ({ nodes: [node] })); + const invoke = vi.mocked(f.runtime.nodes.invoke); + invoke.mockResolvedValue({ payloadJSON: JSON.stringify(page(["original"])) }); + const read = async (params: Partial[0]> = {}) => { + const operation = f.start({ + hostIds: undefined, + allowPartialResults: true, + listNodes, + ...params, + }); + try { + for (;;) { + const result = await operation.next(); + if (result.done) { + return result.hosts; + } + } + } finally { + operation.close(); + } + }; + return { ...f, node, listNodes, invoke, read }; +} diff --git a/extensions/codex/src/session-catalog-list-operation.test.ts b/extensions/codex/src/session-catalog-list-operation.test.ts index 2b1814529c38..e19f555051d7 100644 --- a/extensions/codex/src/session-catalog-list-operation.test.ts +++ b/extensions/codex/src/session-catalog-list-operation.test.ts @@ -1,109 +1,15 @@ import { setImmediate as nextTurn } from "node:timers/promises"; import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; -import type { SessionCatalogProvider } from "openclaw/plugin-sdk/session-catalog"; import { describe, expect, it, vi } from "vitest"; -import { CODEX_TERMINAL_START_COMMAND } from "./session-catalog-terminal.js"; -import type { - CodexSessionCatalogPage, - CodexSessionCatalogPageParams, -} from "./session-catalog-types.js"; import { - CODEX_APP_SERVER_THREADS_LIST_COMMAND, - config, - createCodexSessionCatalogControlFactory, - createCodexTestBindingStore, - createControl, - createGatewayApi, - createRuntime, - registerCodexSessionCatalog, -} from "./session-catalog.test-helpers.js"; - -function page(ids: string[], nextCursor?: string): CodexSessionCatalogPage { - return { - sessions: ids.map((threadId) => ({ - threadId, - name: threadId, - status: "idle", - source: "cli", - archived: false, - })), - ...(nextCursor ? { nextCursor } : {}), - }; -} - -function observe(promise: Promise) { - const state = { settled: false }; - const done = promise - .then( - (value) => ({ status: "fulfilled" as const, value }), - (reason: unknown) => ({ status: "rejected" as const, reason }), - ) - .finally(() => { - state.settled = true; - }); - return { state, done }; -} - -async function fixture(homeCount = 1) { - const { runtime } = createRuntime(); - const base = createCodexSessionCatalogControlFactory({ - getPluginConfig: () => ({ supervision: { enabled: true } }), - getRuntimeConfig: () => config, - }); - const primary = (await base.homesForAgent("main"))[0]!; - const homes = Array.from({ length: homeCount }, (_, index) => ({ - ...primary, - sourceHomeId: `home-${index}`, - hostId: index === 0 ? "gateway:local" : `gateway:local:home-${index}`, - label: `Home ${index}`, - })); - const listPage = - vi.fn< - (homeId: string, params: CodexSessionCatalogPageParams) => Promise - >(); - listPage.mockResolvedValue(page(["visible"])); - const snapshot = vi.fn( - async () => new Map(homes.map((home) => [home.sourceHomeId, new Set(["managed"])])), - ); - const bindingStore = Object.assign(createCodexTestBindingStore(), { - managedThreads: { has: vi.fn(async () => false), mark: vi.fn(async () => true), snapshot }, - }); - const { api, getProvider } = createGatewayApi(runtime, config); - registerCodexSessionCatalog({ - api, - bindingStore, - control: { - ...base, - homesForAgent: async () => homes, - forRequest: (_agentId, source) => - createControl({ - listPage: (params) => listPage(source!.sourceHomeId, params), - }), - }, - getRuntimeConfig: () => config, - }); - const controller = new AbortController(); - const onHost = vi.fn(); - const publications: Promise[] = []; - const start = (params: Partial[0]> = {}) => { - const provider = getProvider()!; - if (!provider.createListOperation) { - throw new Error("Codex list operation is unavailable"); - } - return provider.createListOperation({ - agentId: "main", - limitPerHost: 1, - hostIds: homes.map((home) => home.hostId), - signal: controller.signal, - onHost, - waitUntil: (completion) => { - publications.push(completion); - }, - ...params, - }); - }; - return { runtime, homes, listPage, snapshot, controller, onHost, publications, start }; -} + fixture, + nodeFixture, + observe, + page, +} from "./session-catalog-list-operation.test-support.js"; +import { CODEX_APP_SERVER_THREADS_LIST_COMMAND } from "./session-catalog-parsing.js"; +import { CODEX_TERMINAL_START_COMMAND } from "./session-catalog-terminal.js"; +import type { CodexSessionCatalogPage } from "./session-catalog-types.js"; describe("Codex catalog list operation", () => { it("yields an inert exclusion checkpoint and retains filled rows, limits and cursors", async () => { @@ -550,7 +456,7 @@ describe("Codex catalog list operation", () => { }, ); - it("returns a complete local list at the node response deadline while its publication stays owned", async () => { + it("returns local hosts within 250 ms without publishing an empty cold node", async () => { const f = await fixture(); vi.useFakeTimers(); const invoked = createDeferred(); @@ -561,6 +467,7 @@ describe("Codex catalog list operation", () => { }); const operation = f.start({ hostIds: undefined, + allowPartialResults: true, listNodes: async () => ({ nodes: [ { nodeId: "remote", connected: true, commands: [CODEX_APP_SERVER_THREADS_LIST_COMMAND] }, @@ -570,14 +477,15 @@ describe("Codex catalog list operation", () => { const advancing = observe(operation.next()); try { await invoked.promise; - await vi.advanceTimersByTimeAsync(8_000); + await vi.advanceTimersByTimeAsync(250); + expect(advancing.state.settled).toBe(true); await expect(advancing.done).resolves.toMatchObject({ status: "fulfilled", value: { done: true, hosts: [ { sessions: [{ threadId: "visible" }] }, - { hostId: "node:remote", error: { code: "NODE_INVOKE_FAILED" } }, + { hostId: "node:remote", pending: true, sessions: [] }, ], }, }); @@ -601,6 +509,302 @@ describe("Codex catalog list operation", () => { } }); + it("serves the retained node immediately and rejects an older refresh after a newer publication", async () => { + const f = await nodeFixture(); + await f.read(); + const older = createDeferred(); + const newer = createDeferred(); + f.invoke.mockReturnValueOnce(older.promise).mockReturnValueOnce(newer.promise); + const first = observe(f.read()); + const second = observe(f.read()); + try { + await nextTurn(); + expect(first.state.settled).toBe(true); + expect(second.state.settled).toBe(true); + for (const result of [first, second]) { + await expect(result.done).resolves.toMatchObject({ + status: "fulfilled", + value: [ + { hostId: "gateway:local" }, + { hostId: "node:remote", sessions: [{ threadId: "original" }] }, + ], + }); + } + newer.resolve({ payloadJSON: JSON.stringify(page(["newer"])) }); + await nextTurn(); + older.resolve({ payloadJSON: JSON.stringify(page(["older"])) }); + await Promise.all(f.publications); + const published = f.onHost.mock.calls.flatMap(([host]) => + host.hostId === "node:remote" + ? host.sessions.map((row: { threadId: string }) => row.threadId) + : [], + ); + expect(published).toContain("newer"); + expect(published).not.toContain("older"); + const held = createDeferred(); + f.invoke.mockReturnValueOnce(held.promise); + try { + await expect(f.read()).resolves.toMatchObject([ + { hostId: "gateway:local" }, + { hostId: "node:remote", sessions: [{ threadId: "newer" }] }, + ]); + } finally { + held.resolve({ payloadJSON: JSON.stringify(page(["final"])) }); + } + } finally { + older.resolve({ payloadJSON: JSON.stringify(page([])) }); + newer.resolve({ payloadJSON: JSON.stringify(page([])) }); + await Promise.allSettled([first.done, second.done, ...f.publications]); + } + }); + + it("replays a newer shared snapshot before returning a list held by a local host", async () => { + const f = await nodeFixture(); + await f.read(); + const local = createDeferred(); + const obsolete = createDeferred(); + const cachedFinalizer = createDeferred(); + f.listPage.mockReturnValueOnce(local.promise); + f.invoke.mockReturnValueOnce(obsolete.promise); + const publications: Promise[] = []; + const onHost = vi.fn((host: { hostId: string }) => + host.hostId === "node:remote" ? cachedFinalizer.promise : undefined, + ); + const pending = observe( + f.read({ + // oxlint-disable-next-line typescript/no-misused-promises -- Retain asynchronous JS callbacks even though the SDK declares void. + onHost, + waitUntil: (completion) => { + publications.push(completion); + }, + }), + ); + try { + await nextTurn(); + f.invoke.mockResolvedValueOnce({ payloadJSON: JSON.stringify(page(["newer"])) }); + await f.read({ hostIds: ["node:remote"], allowPartialResults: false }); + local.resolve(page(["visible"])); + await expect(pending.done).resolves.toMatchObject({ + status: "fulfilled", + value: [ + { hostId: "gateway:local" }, + { hostId: "node:remote", sessions: [{ threadId: "newer" }] }, + ], + }); + expect(onHost).toHaveBeenLastCalledWith( + expect.objectContaining({ + hostId: "node:remote", + sessions: [expect.objectContaining({ threadId: "newer" })], + }), + ); + obsolete.resolve({ payloadJSON: JSON.stringify(page(["obsolete"])) }); + const tails = observe(Promise.all(publications)); + await nextTurn(); + expect(tails.state.settled).toBe(false); + cachedFinalizer.resolve(); + await tails.done; + } finally { + cachedFinalizer.resolve(); + local.resolve(page([])); + obsolete.resolve({ payloadJSON: JSON.stringify(page([])) }); + await Promise.allSettled([pending.done, ...publications, ...f.publications]); + } + }); + + it.each([false, true])( + "delivers cold nodes across overlapping queries (newer query first: %s)", + async (newerFirst) => { + const f = await nodeFixture(); + const older = createDeferred(); + const newer = createDeferred(); + f.invoke.mockReturnValueOnce(older.promise).mockReturnValueOnce(newer.promise); + const firstHost = vi.fn(); + const secondHost = vi.fn(); + vi.useFakeTimers(); + const first = observe(f.read({ onHost: firstHost })); + const second = observe( + f.read({ onHost: secondHost, ...(newerFirst ? { limitPerHost: 2 } : {}) }), + ); + try { + await nextTurn(); + await vi.advanceTimersByTimeAsync(250); + expect(first.state.settled).toBe(true); + expect(second.state.settled).toBe(true); + if (newerFirst) { + newer.resolve({ payloadJSON: JSON.stringify(page(["newer-answer"])) }); + await nextTurn(); + } + older.resolve({ payloadJSON: JSON.stringify(page(["first-answer"])) }); + await nextTurn(); + expect(firstHost).toHaveBeenCalledWith( + expect.objectContaining({ + hostId: "node:remote", + sessions: [expect.objectContaining({ threadId: "first-answer" })], + }), + ); + newer.resolve({ payloadJSON: JSON.stringify(page(["newer-answer"])) }); + await Promise.all(f.publications); + expect(secondHost).toHaveBeenCalledWith( + expect.objectContaining({ + hostId: "node:remote", + sessions: [expect.objectContaining({ threadId: "newer-answer" })], + }), + ); + } finally { + older.resolve({ payloadJSON: JSON.stringify(page([])) }); + newer.resolve({ payloadJSON: JSON.stringify(page([])) }); + await Promise.allSettled([first.done, second.done, ...f.publications]); + vi.useRealTimers(); + } + }, + ); + + it.each([{ search: "matching" }, { limitPerHost: 2 }])( + "does not reuse another query's node page: %j", + async (query) => { + const f = await nodeFixture(); + await f.read(); + const held = createDeferred(); + f.invoke.mockReturnValueOnce(held.promise); + vi.useFakeTimers(); + const pending = observe(f.read(query)); + try { + await nextTurn(); + await vi.advanceTimersByTimeAsync(250); + expect(pending.state.settled).toBe(true); + await expect(pending.done).resolves.toMatchObject({ + status: "fulfilled", + value: [ + { hostId: "gateway:local" }, + { hostId: "node:remote", pending: true, sessions: [] }, + ], + }); + const result = await pending.done; + if (result.status === "fulfilled") { + expect(result.value).toHaveLength(2); + } + } finally { + held.resolve({ payloadJSON: JSON.stringify(page(["matching"])) }); + await pending.done; + await Promise.allSettled(f.publications); + vi.useRealTimers(); + } + }, + ); + + it.each(["offline", "removed", "reconnected", "config", "aborted"] as const)( + "invalidates retained node pages and delayed publications when %s", + async (change) => { + const f = await nodeFixture(); + await f.read(); + const old = createDeferred(); + f.invoke.mockReturnValueOnce(old.promise); + await f.read(); + if (change === "offline") { + f.listNodes.mockResolvedValue({ nodes: [{ ...f.node, connected: false }] }); + } + if (change === "removed") { + f.listNodes.mockResolvedValue({ nodes: [] }); + } + if (change === "reconnected") { + f.listNodes.mockResolvedValue({ nodes: [{ ...f.node, connectedAtMs: 2 }] }); + } + if (change === "config") { + f.replaceConfig(); + } + if (change === "aborted") { + f.controller.abort(); + } + const held = createDeferred(); + f.invoke.mockReturnValueOnce(held.promise); + vi.useFakeTimers(); + const pending = observe(f.read({ signal: new AbortController().signal })); + try { + await nextTurn(); + await vi.advanceTimersByTimeAsync(250); + old.resolve({ payloadJSON: JSON.stringify(page(["obsolete"])) }); + await nextTurn(); + expect(pending.state.settled).toBe(true); + expect( + f.onHost.mock.calls.flatMap(([host]) => + host.sessions.map((row: { threadId: string }) => row.threadId), + ), + ).not.toContain("obsolete"); + const result = await pending.done; + expect(result.status).toBe("fulfilled"); + if (result.status === "fulfilled" && change !== "aborted") { + expect( + result.value.flatMap((host) => host.sessions.map((row) => row.threadId)), + ).not.toContain("original"); + if (change === "offline") { + expect(result.value[1]?.error?.code).toBe("NODE_OFFLINE"); + } + } + } finally { + old.resolve({ payloadJSON: JSON.stringify(page([])) }); + held.resolve({ payloadJSON: JSON.stringify(page([])) }); + await pending.done; + await Promise.allSettled(f.publications); + vi.useRealTimers(); + } + }, + ); + + it.each(["connection", "config"] as const)( + "does not bind a delayed inventory to a newer %s", + async (replacement) => { + const f = await nodeFixture(); + const inventory = createDeferred>>(); + const inventoryStarted = createDeferred(); + const newer = createDeferred(); + f.invoke + .mockReturnValueOnce(newer.promise) + .mockResolvedValueOnce({ payloadJSON: JSON.stringify(page(["obsolete"])) }); + const first = observe( + f.read({ + listNodes: () => { + inventoryStarted.resolve(); + return inventory.promise; + }, + }), + ); + await inventoryStarted.promise; + if (replacement === "config") { + f.replaceConfig(); + } else { + f.listNodes.mockResolvedValue({ nodes: [{ ...f.node, connectedAtMs: 2 }] }); + } + vi.useFakeTimers(); + const second = observe(f.read()); + try { + await nextTurn(); + inventory.resolve({ nodes: [f.node] }); + await nextTurn(); + await vi.advanceTimersByTimeAsync(250); + expect(first.state.settled).toBe(true); + expect(second.state.settled).toBe(true); + expect( + f.onHost.mock.calls.flatMap(([host]) => + host.sessions.map((row: { threadId: string }) => row.threadId), + ), + ).not.toContain("obsolete"); + newer.resolve({ payloadJSON: JSON.stringify(page(["current"])) }); + await Promise.all(f.publications); + expect(f.onHost).toHaveBeenCalledWith( + expect.objectContaining({ + hostId: "node:remote", + sessions: [expect.objectContaining({ threadId: "current" })], + }), + ); + } finally { + inventory.resolve({ nodes: [f.node] }); + newer.resolve({ payloadJSON: JSON.stringify(page([])) }); + await Promise.allSettled([first.done, second.done, ...f.publications]); + vi.useRealTimers(); + } + }, + ); + it("joins a started node sibling before rejecting a fatal publication failure", async () => { const f = await fixture(); const held = createDeferred(); diff --git a/extensions/codex/src/session-catalog-list-operation.ts b/extensions/codex/src/session-catalog-list-operation.ts index c8392a4b65ba..65faba99d550 100644 --- a/extensions/codex/src/session-catalog-list-operation.ts +++ b/extensions/codex/src/session-catalog-list-operation.ts @@ -11,6 +11,11 @@ import type { CodexAppServerBindingStore } from "./app-server/session-binding.js import { CodexCatalogLoadingError } from "./session-catalog-availability.js"; import { currentCodexCatalogListDiagnostics } from "./session-catalog-diagnostics.js"; import type { CodexCatalogHome } from "./session-catalog-homes.js"; +import type { CatalogNode } from "./session-catalog-node-continue.js"; +import { + CodexCatalogNodeSnapshots, + createNodeHostPublication, +} from "./session-catalog-node-snapshot.js"; import { catalogError, CODEX_APP_SERVER_THREADS_LIST_COMMAND, @@ -40,6 +45,8 @@ type ListParams = { waitUntil?: (completion: Promise) => void; signal?: AbortSignal; sessionEntries?: SessionCatalogEntrySnapshot; + allowPartialResults?: boolean; + nodeSnapshots?: CodexCatalogNodeSnapshots; includeLocal?: boolean; localHomes?: CodexCatalogHome[]; }; @@ -63,6 +70,39 @@ type PreparedList = { requestedHostIds?: Set; }; +async function boundedNodeHost( + pending: Promise, +): Promise { + let timer: ReturnType | undefined; + try { + return await Promise.race([ + pending, + new Promise((resolve) => { + timer = setTimeout(() => resolve(undefined), 250); + }), + ]); + } finally { + clearTimeout(timer); + } +} + +function measureNodeHost( + host: Promise, + diagnostics: ReturnType, + started: number, +): Promise { + // The node can outlive the list; retain diagnostics without the request's lexical context. + return diagnostics + ? host.finally(() => { + if (!diagnostics.closed) { + diagnostics.fields.pairedNodeSettled = (diagnostics.fields.pairedNodeSettled ?? 0) + 1; + diagnostics.fields.nodeWaitSumMs = + (diagnostics.fields.nodeWaitSumMs ?? 0) + performance.now() - started; + } + }) + : host; +} + function hostFailure( source: CodexCatalogHome | undefined, error: unknown, @@ -179,6 +219,9 @@ class CodexCatalogListDriver { private locals: LocalHost[] = []; private nodeHosts: CodexSessionCatalogHost[] | undefined; private nodeActive = false; + private nodeResults: Array<() => CodexSessionCatalogHost | undefined> = []; + private readonly nodeSnapshots: CodexCatalogNodeSnapshots; + private readonly nodeGeneration: number; private nodeDiscoveryFailed = false; private readonly nodePublications = { pending: 0 }; private nodesStarted = false; @@ -191,6 +234,8 @@ class CodexCatalogListDriver { constructor(params: ListParams) { this.params = params; + this.nodeSnapshots = params.nodeSnapshots ?? new CodexCatalogNodeSnapshots(); + this.nodeGeneration = this.nodeSnapshots.start(params.config); } private request(): ListParams { @@ -334,10 +379,12 @@ class CodexCatalogListDriver { if (diagnostics) { diagnostics.fields.nodeRegistryCalls = 1; } - let nodes: Awaited>["nodes"]; + let nodes: CatalogNode[]; + let inventory: CatalogNode[]; try { try { - nodes = (await (params.listNodes?.() ?? params.runtime.nodes.list())).nodes + inventory = (await (params.listNodes?.() ?? params.runtime.nodes.list())).nodes; + nodes = inventory .filter( (node) => node.gatewayLocal !== true && @@ -366,8 +413,10 @@ class CodexCatalogListDriver { return [host]; } params.signal?.throwIfAborted(); - const { listNodeAdoptedSessionEntries } = await import("./session-catalog-node-adoption.js"); - const { compareNodeLabels, listPairedNode } = + this.nodeSnapshots.observe(this.nodeGeneration, inventory); + const { listNodeAdoptedSessionEntries, nodeAdoptedSourceKey } = + await import("./session-catalog-node-adoption.js"); + const { compareNodeLabels, listPairedNode, nodeLabel } = await import("./session-catalog-node-continue.js"); params.signal?.throwIfAborted(); const adopted = listNodeAdoptedSessionEntries({ @@ -381,7 +430,43 @@ class CodexCatalogListDriver { diagnostics.fields.pairedNodeSettled = 0; } const trackPublication = createNodePublicationTracker(this.nodePublications, params.waitUntil); + const partial = + params.allowPartialResults === true && Boolean(params.onHost && params.waitUntil); const pendingHosts = nodes.toSorted(compareNodeLabels).map((node) => { + const key = JSON.stringify([ + agentId, + query.limitPerHost, + query.search, + query.cursors?.[`node:${node.nodeId}`], + node.displayName, + node.remoteIp, + node.caps, + node.commands, + node.invocableCommands, + ]); + const publication = this.nodeSnapshots.forNode(node, this.nodeGeneration, key); + const { project, publish, publishCached, readPublished } = createNodeHostPublication( + publication, + adopted, + nodeAdoptedSourceKey, + params.onHost, + params.signal, + partial, + trackPublication, + ); + const cached = partial ? publication.read() : undefined; + let result: CodexSessionCatalogHost | undefined; + this.nodeResults.push(() => { + const latest = partial ? publication.read() : undefined; + if (latest) { + publishCached(latest); + return project(latest.host); + } + return !partial ? result : publication.valid() ? (readPublished() ?? result) : undefined; + }); + if (cached) { + publishCached(cached); + } const nodeStarted = diagnostics ? performance.now() : 0; if (diagnostics && !diagnostics.closed) { diagnostics.fields.pairedNodeCalls = (diagnostics.fields.pairedNodeCalls ?? 0) + 1; @@ -391,25 +476,36 @@ class CodexCatalogListDriver { runtime: params.runtime, node, query, - adoptedSessions: adopted, terminalCapabilities: codexNodeTerminalCapability(node), waitUntil: trackPublication, signal: params.signal, - ...(params.onHost ? { onHost: params.onHost } : {}), + onHost: publish, + }); + const completion = measureNodeHost(host, diagnostics, nodeStarted); + if (cached) { + void completion.catch(() => undefined); + return Promise.resolve(project(cached.host)); + } + return (partial ? boundedNodeHost(completion) : completion).then((value) => { + result = value + ? project(value) + : { + hostId: `node:${node.nodeId}`, + label: nodeLabel(node), + kind: "node", + nodeId: node.nodeId, + connected: true, + pending: true, + ...codexNodeTerminalCapability(node), + sessions: [], + }; + return result; }); - return diagnostics - ? host.finally(() => { - if (!diagnostics.closed) { - diagnostics.fields.pairedNodeSettled = - (diagnostics.fields.pairedNodeSettled ?? 0) + 1; - diagnostics.fields.nodeWaitSumMs = - (diagnostics.fields.nodeWaitSumMs ?? 0) + performance.now() - nodeStarted; - } - }) - : host; }); try { - return await Promise.all(pendingHosts); + return (await Promise.all(pendingHosts)).filter( + (host): host is CodexSessionCatalogHost => host !== undefined, + ); } catch (error) { // A fatal callback still owns every started node's fail-soft foreground result. await Promise.allSettled(pendingHosts); @@ -446,14 +542,24 @@ class CodexCatalogListDriver { this.step.reject(this.failure.error); } } else if (this.nodeHosts && this.locals.every((host) => host.value !== undefined)) { - this.complete = true; - this.step.resolve({ - done: true, - hosts: [ - ...this.locals.flatMap((host) => (host.value ? [host.value] : [])), - ...this.nodeHosts, - ], - }); + try { + this.complete = true; + this.step.resolve({ + done: true, + hosts: [ + ...this.locals.flatMap((host) => (host.value ? [host.value] : [])), + ...(this.nodeResults.length + ? this.nodeResults.flatMap((read) => { + const host = read(); + return host ? [host] : []; + }) + : this.nodeHosts), + ], + }); + } catch (error) { + this.failure ??= { error }; + this.step.reject(error); + } } else if (this.canPause()) { this.step.resolve({ done: false }); } @@ -504,6 +610,7 @@ class CodexCatalogListDriver { this.locals = []; this.prepared = undefined; this.nodeHosts = undefined; + this.nodeResults = []; this.failure = undefined; } } diff --git a/extensions/codex/src/session-catalog-node-continue.ts b/extensions/codex/src/session-catalog-node-continue.ts index a12ea78f1d0d..86d8f653354d 100644 --- a/extensions/codex/src/session-catalog-node-continue.ts +++ b/extensions/codex/src/session-catalog-node-continue.ts @@ -98,7 +98,6 @@ export async function listPairedNode(params: { runtime: PluginRuntime; node: CatalogNode; query: CodexSessionCatalogParams; - adoptedSessions: ReadonlyMap; terminalCapabilities: Pick; onHost?: (host: CodexSessionCatalogHost) => void; waitUntil?: (completion: Promise) => void; @@ -154,19 +153,9 @@ export async function listPairedNode(params: { ...page, canContinueCodex: common.canContinueCodex && page.canContinueCodex === true && Boolean(page.sourceHomeId), - sessions: page.sessions.map((session) => { - const adopted = page.sourceHomeId - ? params.adoptedSessions.get( - nodeAdoptedSourceKey(hostId, session.threadId, page.sourceHomeId), - ) - : undefined; - return Object.assign( - {}, - session, - page.sourceHomeId ? { sourceHomeId: page.sourceHomeId } : {}, - adopted ? { sessionKey: adopted.key } : {}, - ); - }), + sessions: page.sessions.map((session) => + Object.assign({}, session, page.sourceHomeId ? { sourceHomeId: page.sourceHomeId } : {}), + ), }; }) .catch((error: unknown) => ({ diff --git a/extensions/codex/src/session-catalog-node-listing.test.ts b/extensions/codex/src/session-catalog-node-listing.test.ts index 9a214b64d4ff..382f6054731f 100644 --- a/extensions/codex/src/session-catalog-node-listing.test.ts +++ b/extensions/codex/src/session-catalog-node-listing.test.ts @@ -445,7 +445,6 @@ describe("Codex supervision catalog", () => { commands: [CODEX_APP_SERVER_THREADS_LIST_COMMAND], }, query: { limitPerHost: 40 }, - adoptedSessions: new Map(), terminalCapabilities: { canStartTerminal: true, canOpenTerminalCodex: true }, }); diff --git a/extensions/codex/src/session-catalog-node-performance.test.ts b/extensions/codex/src/session-catalog-node-performance.test.ts new file mode 100644 index 000000000000..a46b3df39834 --- /dev/null +++ b/extensions/codex/src/session-catalog-node-performance.test.ts @@ -0,0 +1,51 @@ +import { it, vi } from "vitest"; +import { fixture, page } from "./session-catalog-list-operation.test-support.js"; +import { CODEX_APP_SERVER_THREADS_LIST_COMMAND } from "./session-catalog-parsing.js"; + +it.runIf(process.env.OPENCLAW_CATALOG_NODE_BENCH === "1")( + "measures a one-second paired node", + async () => { + const f = await fixture(); + vi.mocked(f.runtime.nodes.invoke).mockImplementation(async () => { + await new Promise((resolve) => { + setTimeout(resolve, 1_000); + }); + return { payloadJSON: JSON.stringify(page(["node-row"])) }; + }); + const elapsed: number[] = []; + const local: number[] = []; + for (let i = 0; i < 6; i++) { + const started = performance.now(); + const operation = f.start({ + hostIds: undefined, + allowPartialResults: true, + listNodes: async () => ({ + nodes: [ + { + nodeId: "remote", + connected: true, + commands: [CODEX_APP_SERVER_THREADS_LIST_COMMAND], + }, + ], + }), + onHost: (host) => { + if (host.kind === "gateway") { + local.push(performance.now() - started); + } + }, + }); + try { + let step = await operation.next(); + while (!step.done) { + step = await operation.next(); + } + elapsed.push(performance.now() - started); + } finally { + operation.close(); + await Promise.all(f.publications); + } + } + console.info("paired node benchmark", JSON.stringify({ elapsedMs: elapsed, localMs: local })); + }, + 15_000, +); diff --git a/extensions/codex/src/session-catalog-node-snapshot.ts b/extensions/codex/src/session-catalog-node-snapshot.ts new file mode 100644 index 000000000000..81944fe527f5 --- /dev/null +++ b/extensions/codex/src/session-catalog-node-snapshot.ts @@ -0,0 +1,143 @@ +import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; +import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; +import type { AdoptedSessionEntry, nodeAdoptedSourceKey } from "./session-catalog-node-adoption.js"; +import type { CatalogNode } from "./session-catalog-node-continue.js"; +import type { CodexSessionCatalogHost } from "./session-catalog-types.js"; + +type NodeSnapshot = { key: string; generation: number; host: CodexSessionCatalogHost }; +type NodePublication = { + connection: number | undefined; + generation: number; + snapshot?: NodeSnapshot; +}; + +/** One query-compatible native page per node, owned by the registered catalog provider. */ +export class CodexCatalogNodeSnapshots { + private config: OpenClawConfig | undefined; + private generation = 0; + private inventoryGeneration = 0; + private configGeneration = 0; + private readonly nodes = new Map(); + + start(config: OpenClawConfig | undefined): number { + if (this.config !== config) { + this.nodes.clear(); + this.config = config; + this.configGeneration = this.generation + 1; + this.inventoryGeneration = this.configGeneration; + } + return ++this.generation; + } + + observe(generation: number, nodes: readonly CatalogNode[]): void { + if (generation < this.inventoryGeneration) { + return; + } + this.inventoryGeneration = generation; + const connected = new Map( + nodes.filter((node) => node.connected).map((node) => [node.nodeId, node]), + ); + for (const [nodeId, publication] of this.nodes) { + const node = connected.get(nodeId); + if (!node || node.connectedAtMs !== publication.connection) { + this.nodes.delete(nodeId); + } + } + for (const node of connected.values()) { + if (!this.nodes.has(node.nodeId)) { + this.nodes.set(node.nodeId, { connection: node.connectedAtMs, generation: 0 }); + } + } + } + + forNode(node: CatalogNode, generation: number, key: string) { + let publication = this.nodes.get(node.nodeId); + if (!publication) { + publication = { connection: node.connectedAtMs, generation: 0 }; + if (generation >= this.inventoryGeneration) { + this.nodes.set(node.nodeId, publication); + } + } + const valid = () => + generation >= this.configGeneration && + this.nodes.get(node.nodeId) === publication && + publication.connection === node.connectedAtMs; + return { + read: () => (valid() && publication.snapshot?.key === key ? publication.snapshot : undefined), + publish: (host: CodexSessionCatalogHost) => { + if (!valid() || generation < publication.generation) { + return false; + } + publication.generation = generation; + publication.snapshot = + node.connected && host.connected ? { key, generation, host } : undefined; + return true; + }, + isCurrent: (snapshot: NodeSnapshot) => valid() && publication.snapshot === snapshot, + valid, + }; + } +} + +export function createNodeHostPublication( + publication: ReturnType, + adopted: ReadonlyMap, + sourceKey: typeof nodeAdoptedSourceKey, + onHost: ((host: CodexSessionCatalogHost) => void) | undefined, + signal: AbortSignal | undefined, + partial: boolean, + waitUntil: ((completion: Promise) => void) | undefined, +) { + let publishedSnapshot: NodeSnapshot | undefined; + let publishedHost: CodexSessionCatalogHost | undefined; + // Late callbacks retain prepared adoption facts, never the request's entries or node inventory. + const project = (host: CodexSessionCatalogHost): CodexSessionCatalogHost => ({ + ...host, + sessions: host.sessions.map((session) => { + const entry = session.sourceHomeId + ? adopted.get(sourceKey(host.hostId, session.threadId, session.sourceHomeId)) + : undefined; + return entry ? { ...session, sessionKey: entry.key } : session; + }), + }); + const emit = (host: CodexSessionCatalogHost) => { + publishedHost = project(host); + return onHost?.(publishedHost); + }; + return { + project, + readPublished: () => publishedHost, + publish: (host: CodexSessionCatalogHost) => { + if (signal?.aborted) { + // Complete callers still own cancellation callbacks; never retain their failed result. + return partial ? undefined : onHost?.(project(host)); + } + if (!publication.valid()) { + return; + } + if (publication.publish(host)) { + publishedSnapshot = publication.read(); + return emit(host); + } + // A newer different query may own the cache, but cannot discard this caller's answer. + const latest = publication.read(); + return emit(latest?.host ?? host); + }, + publishCached: (snapshot: NodeSnapshot) => { + if (signal?.aborted || snapshot === publishedSnapshot || !publication.isCurrent(snapshot)) { + return; + } + const completion = createDeferred(); + void completion.promise.catch(() => undefined); + try { + waitUntil?.(completion.promise); + // Keep replay atomic with the final array; a deferred older frame could overwrite it. + completion.resolve(emit(snapshot.host)); + publishedSnapshot = snapshot; + } catch (error) { + completion.reject(error); + throw error; + } + }, + }; +} diff --git a/extensions/codex/src/session-catalog-types.ts b/extensions/codex/src/session-catalog-types.ts index 8e2082be6349..3d3035c90de6 100644 --- a/extensions/codex/src/session-catalog-types.ts +++ b/extensions/codex/src/session-catalog-types.ts @@ -122,6 +122,7 @@ export type CodexSessionCatalogError = { }; export type CodexSessionCatalogHost = { + pending?: boolean; hostId: string; label: string; kind: "gateway" | "node"; diff --git a/extensions/codex/src/session-catalog.ts b/extensions/codex/src/session-catalog.ts index 8dcaaa1fb13b..835ddcb0b4f1 100644 --- a/extensions/codex/src/session-catalog.ts +++ b/extensions/codex/src/session-catalog.ts @@ -21,6 +21,7 @@ import { runCatalogListInline, } from "./session-catalog-list-operation.js"; import { readCodexSessionTranscript } from "./session-catalog-listing.js"; +import { CodexCatalogNodeSnapshots } from "./session-catalog-node-snapshot.js"; import { CatalogParamsError, CODEX_APP_SERVER_THREADS_LIST_COMMAND, @@ -98,6 +99,7 @@ function toGenericCatalogHost( label: host.label, kind: host.kind, connected: host.connected, + ...(host.pending ? { pending: true } : {}), ...(host.nodeId ? { nodeId: host.nodeId } : {}), sessions: host.sessions.map((session) => { const continuableStatus = @@ -311,11 +313,13 @@ function registerCodexSessionCatalog(params: { return { ...bound, source: bound.source }; }; const checkUpstreamActivity = upstream.createChecker(params); + const nodeSnapshots = new CodexCatalogNodeSnapshots(); const createListOperation: NonNullable = (query) => withCatalogListScope(async () => { const { agentId: requestedAgentId, allowProcessHomeFallback, + allowPartialResults, listNodes, onHost, waitUntil, @@ -346,6 +350,8 @@ function registerCodexSessionCatalog(params: { signal, sessionEntries, localHomes, + allowPartialResults, + nodeSnapshots, ...(onHost ? { onHost: mappedHostPublisher(onHost, mapHost) } : {}), }), mapHost, diff --git a/packages/gateway-protocol/src/schema/sessions-catalog.test.ts b/packages/gateway-protocol/src/schema/sessions-catalog.test.ts index 1cf07ae74afe..3dc579427042 100644 --- a/packages/gateway-protocol/src/schema/sessions-catalog.test.ts +++ b/packages/gateway-protocol/src/schema/sessions-catalog.test.ts @@ -136,6 +136,7 @@ describe("SessionsCatalogListParamsSchema", () => { Value.Check(SessionsCatalogListParamsSchema, { agentId: "main", progressId: "progress-1", + allowPartialResults: true, }), ).toBe(true); }); @@ -184,6 +185,12 @@ describe("SessionsCatalogHostEventSchema", () => { }; expect(Value.Check(SessionsCatalogHostEventSchema, event)).toBe(true); + expect( + Value.Check(SessionsCatalogHostEventSchema, { + ...event, + catalog: { ...event.catalog, hosts: [{ ...event.catalog.hosts[0], pending: true }] }, + }), + ).toBe(true); expect(Value.Check(SessionsCatalogHostEventSchema, { ...event, unexpected: true })).toBe(false); expect( Value.Check(SessionsCatalogHostEventSchema, { diff --git a/packages/gateway-protocol/src/schema/sessions-catalog.ts b/packages/gateway-protocol/src/schema/sessions-catalog.ts index 70866b2dd15b..3b61579ab0d3 100644 --- a/packages/gateway-protocol/src/schema/sessions-catalog.ts +++ b/packages/gateway-protocol/src/schema/sessions-catalog.ts @@ -90,6 +90,8 @@ export const SessionCatalogHostSchema = closedObject({ label: NonEmptyString, kind: Type.Union([Type.Literal("gateway"), Type.Literal("node")]), connected: Type.Boolean(), + /** First snapshot is still loading; retain prior rows until the host publication arrives. */ + pending: Type.Optional(Type.Boolean()), nodeId: Type.Optional(NonEmptyString), canStartTerminal: Type.Optional(Type.Boolean()), sessions: Type.Array(SessionCatalogSessionSchema), @@ -109,6 +111,8 @@ export const SessionCatalogSchema = closedObject({ const SessionsCatalogListCommonProperties = { agentId: Type.Optional(NonEmptyString), progressId: Type.Optional(Type.String({ minLength: 1, maxLength: 128 })), + /** Opt into pending hosts completed by incremental publications for this progressId. */ + allowPartialResults: Type.Optional(Type.Boolean()), search: Type.Optional(Type.String()), limitPerHost: Type.Optional(Type.Integer({ minimum: 1 })), hostIds: Type.Optional(Type.Array(NonEmptyString)), diff --git a/src/gateway/server-methods/session-catalog-list-operations.ts b/src/gateway/server-methods/session-catalog-list-operations.ts index 1dad245b5782..97752f4ad5ba 100644 --- a/src/gateway/server-methods/session-catalog-list-operations.ts +++ b/src/gateway/server-methods/session-catalog-list-operations.ts @@ -1,12 +1,17 @@ -import type { SessionCatalog } from "../../../packages/gateway-protocol/src/index.js"; +import type { + SessionCatalog, + SessionsCatalogListParams, +} from "../../../packages/gateway-protocol/src/index.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import type { SessionCatalogInstances } from "./session-catalog-entry-snapshot.js"; import type { SessionCatalogListLifetime } from "./session-catalog-list-lifetime.js"; import type { CatalogRegistrationSnapshot } from "./session-catalog-provider-access.js"; +import type { GatewayClient } from "./types.js"; export type CatalogListEnumeration = { catalogs: SessionCatalog[]; instances: SessionCatalogInstances; + publishedHosts?: Map>; }; type CatalogListOperation = { @@ -21,6 +26,62 @@ type CatalogListOperations = { }; const catalogListsByConfig = new WeakMap(); +const catalogCallerIds = new WeakMap(); +let nextCatalogCallerId = 0; + +export function sessionCatalogListKey(params: { + agentId: string; + client: GatewayClient | null; + request: SessionsCatalogListParams; + allowPartialResults: boolean; + search?: string; + allowProcessHomeFallback: boolean; + visibilityKey: string; +}): string { + // Providers inherit this exact caller through Gateway async scope, including node APIs. + // A matching profile alone cannot make another connection's enumeration reusable. + let callerId = params.client ? catalogCallerIds.get(params.client) : 0; + if (params.client && callerId === undefined) { + callerId = ++nextCatalogCallerId; + catalogCallerIds.set(params.client, callerId); + } + const cursors = params.request.cursors + ? Object.entries(params.request.cursors).toSorted(([left], [right]) => + left.localeCompare(right), + ) + : null; + return JSON.stringify([ + params.agentId, + params.request.catalogId ?? null, + params.allowPartialResults, + params.search ?? null, + params.request.limitPerHost ?? null, + params.request.hostIds ?? null, + cursors, + params.allowProcessHomeFallback, + params.visibilityKey, + callerId, + params.client?.connect?.scopes?.toSorted() ?? [], + params.client?.connect?.role ?? null, + params.client?.connect?.device?.id ?? null, + ]); +} + +export function resolvePublishedSessionCatalogs(result: CatalogListEnumeration): SessionCatalog[] { + return result.publishedHosts + ? result.catalogs.map((catalog) => { + const published = result.publishedHosts?.get(catalog.id); + if (!published?.size || catalog.error) { + return catalog; + } + // The final host set owns withdrawal; pending hosts already occupy their final positions. + return { + ...catalog, + hosts: catalog.hosts.map((host) => published.get(host.hostId) ?? host), + }; + }) + : result.catalogs; +} export function getSessionCatalogListOperations( config: OpenClawConfig, diff --git a/src/gateway/server-methods/session-catalog-progress.test.ts b/src/gateway/server-methods/session-catalog-progress.test.ts index 93c4fe5f5e13..f6c3c26f3255 100644 --- a/src/gateway/server-methods/session-catalog-progress.test.ts +++ b/src/gateway/server-methods/session-catalog-progress.test.ts @@ -171,6 +171,168 @@ describe("session catalog progress ownership", () => { } }); + it.each([ + { request: {}, connected: true, partial: false }, + { request: { progressId: "legacy" }, connected: true, partial: false }, + { request: { allowPartialResults: true, progressId: "live" }, connected: true, partial: true }, + { + request: { allowPartialResults: true, progressId: "live" }, + connected: false, + partial: false, + }, + { + request: { allowPartialResults: true, progressId: "live", hostIds: ["node:fast"] }, + connected: true, + partial: false, + }, + { + request: { allowPartialResults: true, progressId: "live", cursors: { "node:fast": "page" } }, + connected: true, + partial: false, + }, + ])( + "negotiates partial catalog results for $request (connected=$connected)", + async ({ request, connected, partial }) => { + const list = vi.fn(async () => []); + hoisted.activeRegistry.sessionCatalogs = [{ provider: provider("fixture", { list }) }]; + await call( + "sessions.catalog.list", + { catalogId: "fixture", ...request }, + {}, + { connId: "requester" }, + { isConnectionActive: () => connected }, + ); + expect(list).toHaveBeenCalledWith(expect.objectContaining({ allowPartialResults: partial })); + }, + ); + + it("does not share a partial list with a caller awaiting a complete response", async () => { + const release = createDeferredCore(); + const list = vi.fn(async () => { + await release.promise; + return []; + }); + hoisted.activeRegistry.sessionCatalogs = [{ provider: provider("fixture", { list }) }]; + const config = {}; + const client = { connId: "requester" }; + const progressive = startCall( + "sessions.catalog.list", + { allowPartialResults: true, progressId: "live" }, + config, + client, + ); + const complete = startCall("sessions.catalog.list", {}, config, client); + try { + await vi.waitFor(() => expect(list).toHaveBeenCalledTimes(2)); + } finally { + release.resolve(); + await Promise.allSettled([progressive.completion, complete.completion]); + } + }); + + it.each([false, true])( + "keeps newer host publications in the aggregate while another provider waits (cold=%s)", + async (cold) => { + const sibling = createDeferredCore(); + const publish = createDeferredCore(); + const local = { + hostId: "gateway:local", + label: "Local", + kind: "gateway" as const, + connected: true, + sessions: [], + }; + const cached = { + hostId: "node:slow", + label: "Cached", + kind: "node" as const, + connected: true, + sessions: [], + }; + const fresh = { ...cached, label: "Fresh" }; + let publication: Promise | undefined; + hoisted.activeRegistry.sessionCatalogs = [ + { + provider: provider("fixture", { + list: async ({ onHost, waitUntil }) => { + publication = publish.promise.then(() => onHost?.(fresh)); + waitUntil?.(publication); + return cold ? [local, { ...cached, pending: true }] : [local, cached]; + }, + }), + }, + { + provider: provider("sibling", { + list: async () => { + await sibling.promise; + return []; + }, + }), + }, + ]; + const broadcastToConnIds = vi.fn(); + const pending = startCall( + "sessions.catalog.list", + { allowPartialResults: true, progressId: "live" }, + {}, + { connId: "requester" }, + { broadcastToConnIds }, + ); + try { + await vi.waitFor(() => expect(publication).toBeDefined()); + publish.resolve(); + await publication; + expect(broadcastToConnIds).toHaveBeenCalledOnce(); + expect(pending.respond).not.toHaveBeenCalled(); + sibling.resolve(); + await pending.completion; + expect(pending.respond).toHaveBeenCalledWith(true, { + catalogs: [ + expect.objectContaining({ id: "fixture", hosts: [local, fresh] }), + expect.objectContaining({ id: "sibling", hosts: [] }), + ], + }); + } finally { + publish.resolve(); + sibling.resolve(); + await Promise.allSettled([pending.completion, publication]); + } + }, + ); + + it("does not restore a host withdrawn from the provider's final snapshot", async () => { + const broadcastToConnIds = vi.fn(); + const cached = { + hostId: "node:removed", + label: "Removed node", + kind: "node" as const, + connected: true, + sessions: [], + }; + hoisted.activeRegistry.sessionCatalogs = [ + { + provider: provider("fixture", { + list: async ({ onHost }) => { + onHost?.(cached); + return []; + }, + }), + }, + ]; + const respond = await call( + "sessions.catalog.list", + { progressId: "withdrawn", allowPartialResults: true }, + {}, + { connId: "requester" }, + { broadcastToConnIds }, + ); + expect(broadcastToConnIds).toHaveBeenCalledOnce(); + expect(respond.mock.calls[0]?.[1]?.catalogs[0]?.error).toBeUndefined(); + expect(respond).toHaveBeenCalledWith(true, { + catalogs: [expect.objectContaining({ id: "fixture", hosts: [] })], + }); + }); + it.each([0, 128])( "keeps an active list shared after %i distinct lists settle", async (completedQueries) => { diff --git a/src/gateway/server-methods/session-catalog.ts b/src/gateway/server-methods/session-catalog.ts index e8d4447b51f5..8ad0f0f6d84b 100644 --- a/src/gateway/server-methods/session-catalog.ts +++ b/src/gateway/server-methods/session-catalog.ts @@ -41,7 +41,9 @@ import { } from "./session-catalog-list-lifetime.js"; import { getSessionCatalogListOperations, + resolvePublishedSessionCatalogs, retireSessionCatalogLists, + sessionCatalogListKey, type CatalogListEnumeration, } from "./session-catalog-list-operations.js"; import { @@ -105,8 +107,6 @@ const providerCreateTargetsByConfig = new WeakMap< >(); type CatalogListResult = { catalogs: SessionCatalog[] }; -const catalogCallerIds = new WeakMap(); -let nextCatalogCallerId = 0; function providerCreateTargetCache( config: OpenClawConfig, @@ -177,42 +177,6 @@ export function resolveRegisteredCatalogCreateTarget( : resolved; } -function sessionCatalogListKey(params: { - agentId: string; - client: GatewayClient | null; - request: SessionsCatalogListParams; - search?: string; - allowProcessHomeFallback: boolean; - visibilityKey: string; -}): string { - // Providers inherit this exact caller through Gateway async scope, including node APIs. - // A matching profile alone cannot make another connection's enumeration reusable. - let callerId = params.client ? catalogCallerIds.get(params.client) : 0; - if (params.client && callerId === undefined) { - callerId = ++nextCatalogCallerId; - catalogCallerIds.set(params.client, callerId); - } - const cursors = params.request.cursors - ? Object.entries(params.request.cursors).toSorted(([left], [right]) => - left.localeCompare(right), - ) - : null; - return JSON.stringify([ - params.agentId, - params.request.catalogId ?? null, - params.search ?? null, - params.request.limitPerHost ?? null, - params.request.hostIds ?? null, - cursors, - params.allowProcessHomeFallback, - params.visibilityKey, - callerId, - params.client?.connect?.scopes?.toSorted() ?? [], - params.client?.connect?.role ?? null, - params.client?.connect?.device?.id ?? null, - ]); -} - function providerOrRespond( catalogId: string, respond: RespondFn, @@ -362,39 +326,48 @@ export const sessionCatalogHandlers: GatewayRequestHandlers = { // Shared provider enumeration is not permission. Each synchronous delivery gets current // caller facts and one canonical index, never the provider's pre-await planning snapshot. const projectResult = (result: CatalogListEnumeration): CatalogListResult => { + const catalogs = resolvePublishedSessionCatalogs(result); const currentConfig = context.getRuntimeConfig(); const visibility = resolveSessionCatalogVisibility(client, currentConfig); const requestEntries = createSessionCatalogRequestEntrySnapshot({ cfg: currentConfig, fallbackAgentId: resolvedAgent.agentId, projection, - sessionKeys: result.catalogs + sessionKeys: catalogs .flatMap((catalog) => catalog.hosts) .flatMap((host) => host.sessions) .flatMap(({ sessionKey }) => (sessionKey ? [sessionKey] : [])), }); return { - catalogs: result.catalogs.map((catalog) => ({ - ...catalog, - hosts: catalog.hosts.map((host) => - filterSessionCatalogHost( - requestEntries.projectHostSessions( - host, - result.instances, - providerAudiences.get(catalog.id), + catalogs: catalogs.map((catalog) => + Object.assign({}, catalog, { + hosts: catalog.hosts.map((host) => + filterSessionCatalogHost( + requestEntries.projectHostSessions( + host, + result.instances, + providerAudiences.get(catalog.id), + ), + visibility, + { + audience: providerAudiences.get(catalog.id), + requestEntries, + }, ), - visibility, - { - audience: providerAudiences.get(catalog.id), - requestEntries, - }, ), - ), - })), + }), + ), }; }; const progressId = request.progressId; const progressConnId = progressId && client?.connId ? client.connId : undefined; + const isProgressCurrent = () => + progressConnId !== undefined && + client?.invalidated !== true && + context.isConnectionActive?.(progressConnId) !== false && + (!client?.internal?.agentRuntimeIdentity || + context.validateAgentRuntimeApprovalAuthority?.(client.internal.agentRuntimeIdentity) === + true); const subscriber: CatalogListProgressSubscriber | undefined = progressConnId && progressId ? (catalog, instances) => @@ -409,18 +382,21 @@ export const sessionCatalogHandlers: GatewayRequestHandlers = { { dropIfSlow: true }, ) : undefined; + const allowPartialResults = Boolean( + request.allowPartialResults === true && + subscriber && + isProgressCurrent() && + !client?.connectionSignal?.aborted && + !signal?.aborted && + request.hostIds === undefined && + request.cursors === undefined, + ); const subscribe = (progress: SessionCatalogListLifetime) => { if (subscriber && progressConnId) { progress.subscribe( `${progressConnId}\0${progressId}`, subscriber, - () => - client?.invalidated !== true && - context.isConnectionActive?.(progressConnId) !== false && - (!client?.internal?.agentRuntimeIdentity || - context.validateAgentRuntimeApprovalAuthority?.( - client.internal.agentRuntimeIdentity, - ) === true), + isProgressCurrent, client?.connectionSignal ?? signal, () => (projection.needsMaterialization ? projection.ensureMaterialized() : undefined), ); @@ -430,6 +406,7 @@ export const sessionCatalogHandlers: GatewayRequestHandlers = { agentId: resolvedAgent.agentId, client, request, + allowPartialResults, search, allowProcessHomeFallback: allowHomeFallback, visibilityKey: resolveSessionCatalogVisibility(client, config).cacheKey, @@ -482,6 +459,10 @@ export const sessionCatalogHandlers: GatewayRequestHandlers = { : undefined; requestEntries?.freeze(); const instances: SessionCatalogInstances = new Map(); + // Partial lists can publish a newer host while another provider or projection still waits. + const publishedHosts: CatalogListEnumeration["publishedHosts"] = allowPartialResults + ? new Map() + : undefined; const listNodes = createSessionCatalogRequestNodeSnapshot(); const catalogList = await Promise.all( selected.map(async (provider): Promise => { @@ -489,16 +470,21 @@ export const sessionCatalogHandlers: GatewayRequestHandlers = { const resolution = resolveProviderCreateTarget(provider, resolvedAgent.agentId, config); const createTarget = resolution.ok ? resolution.target : undefined; const onHost = (host: SessionCatalog["hosts"][number]) => { + if (publishedHosts) { + const hosts = publishedHosts.get(provider.id) ?? new Map(); + hosts.set(host.hostId, host); + publishedHosts.set(provider.id, hosts); + } requestEntries?.captureHostInstances(host, instances); const catalog = catalogResult(provider, shareRoute, [host], undefined, createTarget); - // Progressive frames are an optimization. The final RPC response remains - // authoritative when a slow client drops an intermediate host update. + // The final response also reconciles these snapshots if a slow client drops a frame. progress.publish(catalog, instances); }; try { const hosts = await progress.runProvider(onHost, (lifetime) => { const providerParams = { agentId: resolvedAgent.agentId, + allowPartialResults, allowProcessHomeFallback: allowHomeFallback, search, limitPerHost: request.limitPerHost, @@ -519,7 +505,7 @@ export const sessionCatalogHandlers: GatewayRequestHandlers = { } }), ); - return { catalogs: catalogList, instances }; + return { catalogs: catalogList, instances, publishedHosts }; })(); const entry = { progress, result: operation }; // Coalesce only concurrent requests; each subsequent list sees current provider rows. diff --git a/src/plugins/session-catalog.ts b/src/plugins/session-catalog.ts index 9d932c92571d..0b1d588eb7bf 100644 --- a/src/plugins/session-catalog.ts +++ b/src/plugins/session-catalog.ts @@ -29,6 +29,8 @@ export type SessionCatalogListProviderParams = { listNodes?: () => ReturnType; /** Publishes completed hosts without waiting for slower machines in the same list. */ onHost?: (host: SessionCatalogHost) => void; + /** True when the caller accepts retained/pending hosts and later authoritative onHost updates. */ + allowPartialResults?: boolean; /** Register host publication before the logical list settles; includes the onHost callback. */ waitUntil?: (completion: Promise) => void; /** Catalog owner retirement, independent of the requesting connection's lifetime. */ diff --git a/ui/src/components/app-sidebar-session-catalog-live.test.ts b/ui/src/components/app-sidebar-session-catalog-live.test.ts index ce62e3e49342..142cc406370f 100644 --- a/ui/src/components/app-sidebar-session-catalog-live.test.ts +++ b/ui/src/components/app-sidebar-session-catalog-live.test.ts @@ -1,7 +1,12 @@ // @vitest-environment node -import { describe, expect, it } from "vitest"; -import type { SessionCatalog } from "../../../packages/gateway-protocol/src/index.ts"; +import { describe, expect, it, vi } from "vitest"; +import type { + SessionCatalog, + SessionCatalogHost, +} from "../../../packages/gateway-protocol/src/index.ts"; +import type { GatewayBrowserClient } from "../api/gateway.ts"; import { SessionCatalogLiveState } from "./app-sidebar-session-catalog-live.ts"; +import { refetchExpandedSessionCatalogPages } from "./app-sidebar-session-catalog-state.ts"; import { sessionCatalogHostKey } from "./app-sidebar-session-types.ts"; function catalog(id: string, hostCount: number): SessionCatalog { @@ -29,6 +34,86 @@ function catalog(id: string, hostCount: number): SessionCatalog { } describe("SessionCatalogLiveState", () => { + it.each(["final", "incremental"] as const)( + "retains rows and cursors for a pending host in %s publications", + async (publication) => { + const live = new SessionCatalogLiveState(); + const { progressId } = live.beginRequest(1); + const current = catalog("codex", 2); + current.hosts[0]!.nextCursor = "next-page"; + const pending = { ...current.hosts[0]!, pending: true, sessions: [], nextCursor: undefined }; + const incoming = { ...current, hosts: [pending] }; + const catalogs = + publication === "final" + ? live.mergeFinal([incoming], [current]) + : live.applyHost({ + payload: { progressId, agentId: "main", catalog: incoming }, + agentId: "main", + catalogs: [current], + pageDepths: new Map(), + })!.catalogs; + expect(catalogs[0]?.hosts[0]).toMatchObject({ + pending: true, + sessions: current.hosts[0]!.sessions, + nextCursor: "next-page", + }); + // A final response still removes genuinely absent hosts. + expect(catalogs[0]?.hosts).toHaveLength(publication === "final" ? 1 : 2); + const request = vi.fn(); + expect( + await refetchExpandedSessionCatalogPages({ + catalogs, + previousCatalogs: [current], + client: { request } as unknown as GatewayBrowserClient, + agentId: "main", + pageDepths: new Map([[sessionCatalogHostKey("codex", pending.hostId), 1]]), + isCurrent: () => true, + canRequestPage: () => true, + }), + ).toEqual(catalogs); + expect(request).not.toHaveBeenCalled(); + }, + ); + + it.each([ + { name: "fresh", details: {} }, + { name: "error", details: { error: { code: "UNAVAILABLE", message: "Node unavailable" } } }, + { name: "offline", details: { connected: false } }, + ])("clears pending when an expanded host publishes $name data", ({ details }) => { + const live = new SessionCatalogLiveState(); + const { progressId } = live.beginRequest(1); + const current = catalog("codex", 1); + current.hosts[0]!.pending = true; + const { pending: _pending, ...fresh } = current.hosts[0]!; + const result = live.applyHost({ + payload: { + progressId, + agentId: "main", + catalog: { ...current, hosts: [{ ...fresh, ...details }] }, + }, + agentId: "main", + catalogs: [current], + pageDepths: new Map([[sessionCatalogHostKey("codex", fresh.hostId), 1]]), + }); + expect(result?.catalogs[0]?.hosts[0]?.pending).toBeUndefined(); + }); + + it("does not replace a settled progressive host with its pending final response", () => { + const live = new SessionCatalogLiveState(); + const { progressId } = live.beginRequest(1); + const current = catalog("codex", 1); + const pending: SessionCatalogHost = { ...current.hosts[0]!, pending: true, sessions: [] }; + const published = live.applyHost({ + payload: { progressId, agentId: "main", catalog: current }, + agentId: "main", + catalogs: [{ ...current, hosts: [pending] }], + pageDepths: new Map(), + })!; + expect(live.mergeFinal([{ ...current, hosts: [pending] }], published.catalogs)).toEqual([ + current, + ]); + }); + it("clears unavailable native readiness without losing another host or expanded rows", () => { const live = new SessionCatalogLiveState(); const { progressId } = live.beginRequest(1); diff --git a/ui/src/components/app-sidebar-session-catalog-live.ts b/ui/src/components/app-sidebar-session-catalog-live.ts index ec6530962a14..83b53720a6ff 100644 --- a/ui/src/components/app-sidebar-session-catalog-live.ts +++ b/ui/src/components/app-sidebar-session-catalog-live.ts @@ -151,7 +151,7 @@ export class SessionCatalogLiveState { const key = sessionCatalogHostKey(catalog.id, host.hostId); currentKeys.add(key); const discovery = this.discoveryPages.get(key); - if (!discovery || host.error || catalog.error) { + if (!discovery || host.pending || host.error || catalog.error) { return host; } // Recheck the head on each refresh. A changed anchor or newly visible row @@ -184,6 +184,13 @@ export class SessionCatalogLiveState { hosts: catalog.hosts.map((host) => { const hostKey = sessionCatalogHostKey(catalog.id, host.hostId); const progressiveHost = currentHosts.get(hostKey); + if (host.pending) { + return this.requestChangedHostKeys.has(hostKey) && + progressiveHost && + !progressiveHost.pending + ? progressiveHost + : preserveExpandedCatalogHost(host, progressiveHost); + } return host.error && this.requestChangedHostKeys.has(hostKey) && progressiveHost && @@ -317,13 +324,16 @@ export class SessionCatalogLiveState { const discovery = this.discoveryPages.get(hostKey); if ( discovery && + !freshHost.pending && !freshHost.error && (freshHost.sessions.length > 0 || freshHost.nextCursor !== discovery.headCursor) ) { this.discoveryPages.delete(hostKey); } const mergedHost = - (params.pageDepths.get(hostKey) ?? 0) > 0 || this.discoveryPages.has(hostKey) + freshHost.pending || + (params.pageDepths.get(hostKey) ?? 0) > 0 || + this.discoveryPages.has(hostKey) ? preserveExpandedCatalogHost(freshHost, currentHost) : freshHost; const hosts = currentHost @@ -435,6 +445,7 @@ export async function refreshSessionCatalogsLive(params: { agentId: params.agentId, limitPerHost: 40, progressId, + allowPartialResults: true, }); if (!requestIsCurrent() || !result?.catalogs) { return; diff --git a/ui/src/components/app-sidebar-session-catalog-state.ts b/ui/src/components/app-sidebar-session-catalog-state.ts index ecf0edc850b1..3ef534891fbb 100644 --- a/ui/src/components/app-sidebar-session-catalog-state.ts +++ b/ui/src/components/app-sidebar-session-catalog-state.ts @@ -66,7 +66,7 @@ export function preserveExpandedCatalogHost( return freshHost; } const { sessions: _freshSessions, nextCursor: _freshNextCursor, ...freshDetails } = freshHost; - const { nextCursor, ...previousDetails } = previous; + const { nextCursor, pending: _pending, ...previousDetails } = previous; return { ...previousDetails, ...freshDetails, @@ -112,7 +112,12 @@ export function mergeSessionCatalogPage(params: { } else { advancedHostIds.push(host.hostId); } - const { nextCursor: _currentCursor, error: _currentError, ...currentHost } = host; + const { + nextCursor: _currentCursor, + error: _currentError, + pending: _pending, + ...currentHost + } = host; return { ...currentHost, ...pageHostDetails, @@ -154,7 +159,7 @@ export async function refetchExpandedSessionCatalogPages(params: { catalog.hosts.map(async (host) => { const pageDepth = params.pageDepths.get(sessionCatalogHostKey(catalog.id, host.hostId)) ?? 0; - if (pageDepth === 0) { + if (pageDepth === 0 || host.pending) { return host; } const previous = previousHosts.get(host.hostId); diff --git a/ui/src/components/app-sidebar.catalog-discovery.test.ts b/ui/src/components/app-sidebar.catalog-discovery.test.ts index 59a45406e5e5..347e525da0b4 100644 --- a/ui/src/components/app-sidebar.catalog-discovery.test.ts +++ b/ui/src/components/app-sidebar.catalog-discovery.test.ts @@ -100,6 +100,7 @@ describe("AppSidebar hidden catalog discovery", () => { agentId: scope === "agent" ? "research" : "main", limitPerHost: 40, progressId: expect.any(String), + allowPartialResults: true, }); retiredPage.resolve( diff --git a/ui/src/components/app-sidebar.catalog-hidden-pages.test.ts b/ui/src/components/app-sidebar.catalog-hidden-pages.test.ts index 123bdc90c041..97a13d32da67 100644 --- a/ui/src/components/app-sidebar.catalog-hidden-pages.test.ts +++ b/ui/src/components/app-sidebar.catalog-hidden-pages.test.ts @@ -73,6 +73,42 @@ describe("AppSidebar expanded catalog refresh visibility", () => { vi.useRealTimers(); }); + it("holds expanded pending rows until an explicit fresh page is requested", async () => { + const pendingPage = catalogPage([]); + pendingPage.catalogs[0]!.hosts[0]!.pending = true; + const request = vi + .fn() + .mockResolvedValueOnce(page(1)) + .mockResolvedValueOnce(page(2)) + .mockResolvedValueOnce(page(3)) + .mockResolvedValueOnce(pendingPage) + .mockResolvedValue(page(4, "Fresh", "")); + const { sidebar } = await mountExpanded(request); + await sidebar.sessionData.refreshSessionCatalogs(); + await settle(sidebar); + expect(request).toHaveBeenCalledTimes(4); + expect(sidebar.sessionData.sessionCatalogs[0]?.hosts[0]).toMatchObject({ + pending: true, + nextCursor: "page-4", + sessions: [ + expect.objectContaining({ threadId: "thread-1" }), + expect.objectContaining({ threadId: "thread-2" }), + expect.objectContaining({ threadId: "thread-3" }), + ], + }); + await loadMore(sidebar); + expect(request).toHaveBeenCalledTimes(5); + expect(request).toHaveBeenLastCalledWith("sessions.catalog.list", { + agentId: "main", + catalogId: "codex", + hostIds: ["gateway:local"], + cursors: { "gateway:local": "page-4" }, + }); + expect(sidebar.sessionData.sessionCatalogs[0]?.hosts[0]?.pending).toBeUndefined(); + expect(sidebar.textContent).toContain("Original 3"); + expect(sidebar.textContent).toContain("Fresh 4"); + }); + it.each(["base", "expanded"] as const)( "stops new automatic pages after hiding during the %s response and catches up once", async (heldStage) => { diff --git a/ui/src/components/session-data-controller-catalog.ts b/ui/src/components/session-data-controller-catalog.ts index f9d5379b2016..8865bee5df17 100644 --- a/ui/src/components/session-data-controller-catalog.ts +++ b/ui/src/components/session-data-controller-catalog.ts @@ -301,7 +301,7 @@ function hiddenSessionCatalogPages(owner: SessionCatalogDataOwner) { return []; } const hostIds = catalog.hosts - .filter((host) => host.nextCursor && !host.error) + .filter((host) => host.nextCursor && !host.pending && !host.error) .map((host) => host.hostId); return hostIds.length > 0 ? [{ catalogId: catalog.id, hostIds }] : []; }); diff --git a/ui/src/e2e/session-catalog-pending-host.e2e.test.ts b/ui/src/e2e/session-catalog-pending-host.e2e.test.ts new file mode 100644 index 000000000000..7bcab882e9ec --- /dev/null +++ b/ui/src/e2e/session-catalog-pending-host.e2e.test.ts @@ -0,0 +1,88 @@ +import path from "node:path"; +import { expect, it } from "vitest"; +import type { SessionCatalog } from "../../../packages/gateway-protocol/src/index.ts"; +import type { AppSidebarSessionNavigationElement } from "../components/app-sidebar-session-navigation.ts"; +import { createControlUiE2eArtifactDir } from "../test-helpers/control-ui-e2e-artifacts.ts"; +import { installMockGateway } from "../test-helpers/control-ui-e2e.ts"; +import { createControlUiE2eSuite } from "./control-ui-e2e-suite.test-support.ts"; + +const suite = createControlUiE2eSuite({ + name: "Pending paired-node catalog", + startServerBeforeBrowser: true, +}); + +suite.define(() => { + it("keeps the paired-node row through pending refresh and applies its later publication", async () => { + const artifactDir = createControlUiE2eArtifactDir("session-catalog-pending-host"); + await suite.withPage({ viewport: { width: 1280, height: 900 } }, async ({ page }) => { + const catalog: SessionCatalog = { + id: "codex", + label: "Codex", + capabilities: { continueSession: false, archive: false }, + hosts: [ + { + hostId: "node:devbox", + label: "Dev Box", + kind: "node", + connected: true, + sessions: [ + { + threadId: "release-review", + name: "Paired node release review", + status: "stored", + archived: false, + canContinue: false, + canArchive: false, + }, + ], + }, + ], + }; + const gateway = await installMockGateway(page, { + featureMethods: ["chat.metadata", "chat.startup", "sessions.catalog.list"], + methodResponses: { "sessions.catalog.list": { catalogs: [catalog] } }, + }); + await page.goto(`${suite.server.baseUrl}chat`); + const sidebar = page.locator("openclaw-app-sidebar"); + const heldRow = sidebar.getByText("Paired node release review", { exact: true }); + await heldRow.waitFor(); + const pendingCatalog = { + ...catalog, + hosts: [{ ...catalog.hosts[0]!, sessions: [], pending: true }], + }; + await gateway.setMethodResponse("sessions.catalog.list", { catalogs: [pendingCatalog] }); + await sidebar.evaluate(async (element) => { + const sidebarElement = element as AppSidebarSessionNavigationElement; + await sidebarElement.sessionData.refreshSessionCatalogs(); + await sidebarElement.updateComplete; + }); + const request = (await gateway.getRequests("sessions.catalog.list")).at(-1)!; + // Capture before the assertion so the original row-clearing regression has visual evidence. + await page.screenshot({ path: path.join(artifactDir, "pending-node.png") }); + expect(await heldRow.count()).toBe(1); + expect(request.params).toMatchObject({ + allowPartialResults: true, + progressId: expect.any(String), + }); + + await gateway.emitGatewayEvent("sessions.catalog.host", { + progressId: (request.params as { progressId: string }).progressId, + agentId: "main", + catalog: { + ...catalog, + hosts: [ + { + ...catalog.hosts[0]!, + sessions: [ + { ...catalog.hosts[0]!.sessions[0]!, name: "Paired node review refreshed" }, + ], + }, + ], + }, + }); + await sidebar.getByText("Paired node review refreshed", { exact: true }).waitFor(); + expect(await heldRow.count()).toBe(0); + await page.screenshot({ path: path.join(artifactDir, "refreshed-node.png") }); + }); + }); +}); diff --git a/ui/src/test-helpers/app-sidebar-cases/catalog-compat.ts b/ui/src/test-helpers/app-sidebar-cases/catalog-compat.ts index d8c082000bb9..9aa4dee72df8 100644 --- a/ui/src/test-helpers/app-sidebar-cases/catalog-compat.ts +++ b/ui/src/test-helpers/app-sidebar-cases/catalog-compat.ts @@ -35,6 +35,7 @@ describe("AppSidebar session catalog pagination", () => { agentId: "main", limitPerHost: 40, progressId: expect.any(String), + allowPartialResults: true, }); const selection = context.agentSelection.state as { @@ -51,6 +52,7 @@ describe("AppSidebar session catalog pagination", () => { agentId: "research", limitPerHost: 40, progressId: expect.any(String), + allowPartialResults: true, }); } finally { vi.useRealTimers(); diff --git a/ui/src/test-helpers/app-sidebar-cases/catalog-live.ts b/ui/src/test-helpers/app-sidebar-cases/catalog-live.ts index a7fd48e5cda2..a58a7720295a 100644 --- a/ui/src/test-helpers/app-sidebar-cases/catalog-live.ts +++ b/ui/src/test-helpers/app-sidebar-cases/catalog-live.ts @@ -64,6 +64,7 @@ describe("AppSidebar session catalog pagination", () => { agentId: "main", limitPerHost: 40, progressId: expect.any(String), + allowPartialResults: true, }); } finally { vi.useRealTimers(); diff --git a/ui/src/test-helpers/app-sidebar-cases/catalog-page-hosts.ts b/ui/src/test-helpers/app-sidebar-cases/catalog-page-hosts.ts index d331537604f9..1755ae03bf4e 100644 --- a/ui/src/test-helpers/app-sidebar-cases/catalog-page-hosts.ts +++ b/ui/src/test-helpers/app-sidebar-cases/catalog-page-hosts.ts @@ -63,6 +63,7 @@ export function registerCatalogPageHostTests() { agentId: "main", limitPerHost: 40, progressId: expect.any(String), + allowPartialResults: true, }); expect(catalogRows()).toHaveLength(2); loadMore()?.click(); @@ -87,6 +88,7 @@ export function registerCatalogPageHostTests() { agentId: "main", limitPerHost: 40, progressId: expect.any(String), + allowPartialResults: true, }); expect(request).toHaveBeenNthCalledWith(4, "sessions.catalog.list", { agentId: "main", diff --git a/ui/src/test-helpers/app-sidebar-cases/catalog-reconnect.ts b/ui/src/test-helpers/app-sidebar-cases/catalog-reconnect.ts index d07485e7c235..ded24e1ba31a 100644 --- a/ui/src/test-helpers/app-sidebar-cases/catalog-reconnect.ts +++ b/ui/src/test-helpers/app-sidebar-cases/catalog-reconnect.ts @@ -45,6 +45,7 @@ describe("AppSidebar catalog reconnect", () => { agentId: "main", limitPerHost: 40, progressId: expect.any(String), + allowPartialResults: true, }); } finally { vi.useRealTimers();