mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
perf(codex): reuse paired node snapshots for progressive catalogs (#153614)
Serve retained paired-node pages immediately and bound cold foreground waits to 250 ms. Refresh through owned incremental publications with query, configuration, and connection fences. Preserve pending sidebar rows and complete one-shot lookup behavior.
This commit is contained in:
parent
7a143dc77c
commit
3c6c393f51
32 changed files with 1481 additions and 292 deletions
|
|
@ -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"
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<void>)`
|
||||
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
|
||||
|
|
|
|||
71
extensions/anthropic/session-catalog-watch.test-support.ts
Normal file
71
extensions/anthropic/session-catalog-watch.test-support.ts
Normal file
|
|
@ -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<string, { watcher: Watcher; recursive: boolean }>();
|
||||
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<string>,
|
||||
listener?: fs.WatchListener<string>,
|
||||
) => {
|
||||
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));
|
||||
},
|
||||
};
|
||||
}
|
||||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -176,6 +176,7 @@ export function createClaudeSessionCatalogRuntime(
|
|||
const localCliAvailable = catalogTerminal.isClaudeCliAvailable();
|
||||
const {
|
||||
allowProcessHomeFallback,
|
||||
allowPartialResults: _allowPartialResults,
|
||||
agentId: _agentId,
|
||||
listNodes,
|
||||
onHost,
|
||||
|
|
|
|||
|
|
@ -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<T>(promise: Promise<T>) {
|
||||
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<CodexSessionCatalogPage>
|
||||
>();
|
||||
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<void>[] = [];
|
||||
const start = (params: Partial<Parameters<SessionCatalogProvider["list"]>[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<Parameters<SessionCatalogProvider["list"]>[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 };
|
||||
}
|
||||
|
|
@ -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<T>(promise: Promise<T>) {
|
||||
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<CodexSessionCatalogPage>
|
||||
>();
|
||||
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<void>[] = [];
|
||||
const start = (params: Partial<Parameters<SessionCatalogProvider["list"]>[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<void>();
|
||||
|
|
@ -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<unknown>();
|
||||
const newer = createDeferred<unknown>();
|
||||
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<unknown>();
|
||||
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<CodexSessionCatalogPage>();
|
||||
const obsolete = createDeferred<unknown>();
|
||||
const cachedFinalizer = createDeferred<void>();
|
||||
f.listPage.mockReturnValueOnce(local.promise);
|
||||
f.invoke.mockReturnValueOnce(obsolete.promise);
|
||||
const publications: Promise<void>[] = [];
|
||||
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<unknown>();
|
||||
const newer = createDeferred<unknown>();
|
||||
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<unknown>();
|
||||
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<unknown>();
|
||||
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<unknown>();
|
||||
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<Awaited<ReturnType<typeof f.listNodes>>>();
|
||||
const inventoryStarted = createDeferred<void>();
|
||||
const newer = createDeferred<unknown>();
|
||||
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<unknown>();
|
||||
|
|
|
|||
|
|
@ -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>) => void;
|
||||
signal?: AbortSignal;
|
||||
sessionEntries?: SessionCatalogEntrySnapshot;
|
||||
allowPartialResults?: boolean;
|
||||
nodeSnapshots?: CodexCatalogNodeSnapshots;
|
||||
includeLocal?: boolean;
|
||||
localHomes?: CodexCatalogHome[];
|
||||
};
|
||||
|
|
@ -63,6 +70,39 @@ type PreparedList = {
|
|||
requestedHostIds?: Set<string>;
|
||||
};
|
||||
|
||||
async function boundedNodeHost(
|
||||
pending: Promise<CodexSessionCatalogHost>,
|
||||
): Promise<CodexSessionCatalogHost | undefined> {
|
||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||
try {
|
||||
return await Promise.race([
|
||||
pending,
|
||||
new Promise<undefined>((resolve) => {
|
||||
timer = setTimeout(() => resolve(undefined), 250);
|
||||
}),
|
||||
]);
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
}
|
||||
|
||||
function measureNodeHost(
|
||||
host: Promise<CodexSessionCatalogHost>,
|
||||
diagnostics: ReturnType<typeof currentCodexCatalogListDiagnostics>,
|
||||
started: number,
|
||||
): Promise<CodexSessionCatalogHost> {
|
||||
// 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<ReturnType<PluginRuntime["nodes"]["list"]>>["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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -98,7 +98,6 @@ export async function listPairedNode(params: {
|
|||
runtime: PluginRuntime;
|
||||
node: CatalogNode;
|
||||
query: CodexSessionCatalogParams;
|
||||
adoptedSessions: ReadonlyMap<string, AdoptedSessionEntry>;
|
||||
terminalCapabilities: Pick<CodexSessionCatalogHost, "canOpenTerminalCodex" | "canStartTerminal">;
|
||||
onHost?: (host: CodexSessionCatalogHost) => void;
|
||||
waitUntil?: (completion: Promise<void>) => 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) => ({
|
||||
|
|
|
|||
|
|
@ -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 },
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
);
|
||||
143
extensions/codex/src/session-catalog-node-snapshot.ts
Normal file
143
extensions/codex/src/session-catalog-node-snapshot.ts
Normal file
|
|
@ -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<string, NodePublication>();
|
||||
|
||||
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<CodexCatalogNodeSnapshots["forNode"]>,
|
||||
adopted: ReadonlyMap<string, AdoptedSessionEntry>,
|
||||
sourceKey: typeof nodeAdoptedSourceKey,
|
||||
onHost: ((host: CodexSessionCatalogHost) => void) | undefined,
|
||||
signal: AbortSignal | undefined,
|
||||
partial: boolean,
|
||||
waitUntil: ((completion: Promise<void>) => 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>();
|
||||
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;
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
|
|
@ -122,6 +122,7 @@ export type CodexSessionCatalogError = {
|
|||
};
|
||||
|
||||
export type CodexSessionCatalogHost = {
|
||||
pending?: boolean;
|
||||
hostId: string;
|
||||
label: string;
|
||||
kind: "gateway" | "node";
|
||||
|
|
|
|||
|
|
@ -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<SessionCatalogProvider["createListOperation"]> = (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,
|
||||
|
|
|
|||
|
|
@ -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, {
|
||||
|
|
|
|||
|
|
@ -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)),
|
||||
|
|
|
|||
|
|
@ -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<string, Map<string, SessionCatalog["hosts"][number]>>;
|
||||
};
|
||||
|
||||
type CatalogListOperation = {
|
||||
|
|
@ -21,6 +26,62 @@ type CatalogListOperations = {
|
|||
};
|
||||
|
||||
const catalogListsByConfig = new WeakMap<OpenClawConfig, CatalogListOperations>();
|
||||
const catalogCallerIds = new WeakMap<GatewayClient, number>();
|
||||
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,
|
||||
|
|
|
|||
|
|
@ -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<SessionCatalogProvider["list"]>(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<SessionCatalogProvider["list"]>(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<void> | 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) => {
|
||||
|
|
|
|||
|
|
@ -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<GatewayClient, number>();
|
||||
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<SessionCatalog> => {
|
||||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -29,6 +29,8 @@ export type SessionCatalogListProviderParams = {
|
|||
listNodes?: () => ReturnType<PluginRuntime["nodes"]["list"]>;
|
||||
/** 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>) => void;
|
||||
/** Catalog owner retirement, independent of the requesting connection's lifetime. */
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -100,6 +100,7 @@ describe("AppSidebar hidden catalog discovery", () => {
|
|||
agentId: scope === "agent" ? "research" : "main",
|
||||
limitPerHost: 40,
|
||||
progressId: expect.any(String),
|
||||
allowPartialResults: true,
|
||||
});
|
||||
|
||||
retiredPage.resolve(
|
||||
|
|
|
|||
|
|
@ -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) => {
|
||||
|
|
|
|||
|
|
@ -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 }] : [];
|
||||
});
|
||||
|
|
|
|||
88
ui/src/e2e/session-catalog-pending-host.e2e.test.ts
Normal file
88
ui/src/e2e/session-catalog-pending-host.e2e.test.ts
Normal file
|
|
@ -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") });
|
||||
});
|
||||
});
|
||||
});
|
||||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -64,6 +64,7 @@ describe("AppSidebar session catalog pagination", () => {
|
|||
agentId: "main",
|
||||
limitPerHost: 40,
|
||||
progressId: expect.any(String),
|
||||
allowPartialResults: true,
|
||||
});
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -45,6 +45,7 @@ describe("AppSidebar catalog reconnect", () => {
|
|||
agentId: "main",
|
||||
limitPerHost: 40,
|
||||
progressId: expect.any(String),
|
||||
allowPartialResults: true,
|
||||
});
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue