From d27f2bd6daaddf4f3dad61feea8254c3da68056d Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Tue, 29 Sep 2026 19:55:19 -0700 Subject: [PATCH] improve(ui): stop re-describing new sessions on every chat event (#161442) * improve(ui): stop re-describing new sessions on every chat event Creating a session from New Session and running two short turns made the chat pane issue 32-36 sessions.describe calls (plus identical in-flight branch/model/descriptor pairs) against a real Gateway. When the dashboard session's parent (agent:main:main) does not exist, every sessions.changed omits ancestorSessions, so descriptor observations were invalidated on every event and the pane's active-resource owner re-described ~4 ms later. The event already carries the full row; describe returns the same row and cannot certify ancestry either. - Row observations publish admitted rows immediately and deliver one paced authoritative invalidation for incomplete ancestry through the existing session event refresh coordinator (absorbed by newer reads, immediate invalidations, retirement and reconnect). - The shared describe owner treats an agent-qualified key with or without its implied agentId as one read slot; refresh still supersedes pending reads. - A session-patch metadata refresh whose chat.metadata answers before its concurrent models.list validates against that pending catalog read instead of issuing an identical replacement. - Concurrent chat branch loads share one pending sessions.branches.list. Real Gateway, create + two turns: describes 32-36 -> 8-13, identical in-flight pairs 1-3 -> 0, total requests 88-102 -> 61-72. * fix(ui): keep chat metadata publication consistently async Publication can wait for a pending catalog read, so return a Promise on every path instead of a sync-or-async union; side effects still apply synchronously when no read is pending. Await it in the metadata, command, model-catalog and model-control tests flagged by type-aware no-floating-promises. * fix(ui): publish slash commands before catalog validation settles A session-patch refresh whose chat.metadata answers before its concurrent models.list held the whole metadata result until the catalog read settled, delaying ready slash commands behind a slow model read. Publish the metadata synchronously again and let only catalog validation wait for the pending read; a changed, missing or retired catalog then invalidates and sends a follow-up catalogChanged notification while the writer is still current. Restores the synchronous publication contract, so the settlement restructure and the test awaits are no longer needed. Also waits for the held chat.startup request in the metadata-observation early-wake E2E before resolving it. --- ui/AGENTS.md | 1 + .../e2e/chat-metadata-observation.e2e.test.ts | 1 + .../chat-session-read-coalescing.e2e.test.ts | 249 ++++++++++++++++++ ui/src/lib/chat/chat-metadata-cache.ts | 3 +- ui/src/lib/chat/chat-metadata-refresh.test.ts | 87 ++++++ ui/src/lib/chat/chat-metadata-store.test.ts | 133 +++++++++- ui/src/lib/chat/chat-metadata-store.ts | 92 ++++--- ui/src/lib/model-catalog-store.ts | 98 ++++--- ui/src/lib/sessions/session-capability.ts | 2 +- ui/src/lib/sessions/session-describe.test.ts | 28 +- ui/src/lib/sessions/session-describe.ts | 84 +++--- ui/src/lib/sessions/session-reconciliation.ts | 9 +- .../sessions/session-roster-observations.ts | 58 +++- .../session-row-certification.test.ts | 180 +++++++++++++ .../pages/chat/chat-branch-freshness.test.ts | 44 +++- ui/src/pages/chat/chat-commands.ts | 30 +-- ui/src/pages/chat/chat-history-branches.ts | 64 +++-- 17 files changed, 973 insertions(+), 190 deletions(-) create mode 100644 ui/src/e2e/chat-session-read-coalescing.e2e.test.ts create mode 100644 ui/src/lib/sessions/session-row-certification.test.ts diff --git a/ui/AGENTS.md b/ui/AGENTS.md index c4845fd942d2..69a47c25f97b 100644 --- a/ui/AGENTS.md +++ b/ui/AGENTS.md @@ -38,6 +38,7 @@ This directory owns Control UI-specific guidance that should not live in the rep - Re-adopting cached lineage rows changes presentation without invalidating managed list membership. Fresh descriptor reads and Gateway events retain their authoritative invalidation paths. +- Descriptor observations apply admitted rows immediately; incomplete ancestor coverage retains one paced authoritative descriptor refresh through the coordinator instead of a read per event. - `lib/sessions/event-refresh-coordinator.ts` owns automatic refresh pacing: collect events in a four-to-five-second window sampled once when armed so browsers spread their reads and subsequent events cannot postpone them. diff --git a/ui/src/e2e/chat-metadata-observation.e2e.test.ts b/ui/src/e2e/chat-metadata-observation.e2e.test.ts index a886225da590..5b23b0394540 100644 --- a/ui/src/e2e/chat-metadata-observation.e2e.test.ts +++ b/ui/src/e2e/chat-metadata-observation.e2e.test.ts @@ -124,6 +124,7 @@ suite.define(() => { await expect.poll(() => panes.count()).toBe(paneCount); await gateway.waitForRequest("models.list"); if (wake === "early") { + await gateway.waitForRequest("chat.startup"); await setDocumentVisibility(page, "hidden"); await setDocumentVisibility(page, "visible"); await gateway.resolveDeferred("chat.startup"); diff --git a/ui/src/e2e/chat-session-read-coalescing.e2e.test.ts b/ui/src/e2e/chat-session-read-coalescing.e2e.test.ts new file mode 100644 index 000000000000..0eb0891d06f8 --- /dev/null +++ b/ui/src/e2e/chat-session-read-coalescing.e2e.test.ts @@ -0,0 +1,249 @@ +import { expect, it } from "vitest"; +import { + defaultControlUiFeatureMethods, + installMockGateway, + pauseVirtualClock, + startControlUiE2eServer, + type MockGatewayRequest, +} from "../test-helpers/control-ui-e2e.ts"; +import { createControlUiE2eSuite } from "./control-ui-e2e-suite.test-support.ts"; + +declare global { + interface Window { + readCoalescingWire: { pending: Map; duplicates: string[] }; + } +} + +const suite = createControlUiE2eSuite({ + name: "Control UI chat session read coalescing", + startServer: () => startControlUiE2eServer(), +}); + +suite.define(() => { + it("paces incomplete ancestry and validates metadata against the pending catalog", async () => { + await suite.withPage({ locale: "en-US", serviceWorkers: "block" }, async ({ page }) => { + const gateway = await installMockGateway(page, { + featureMethods: [ + ...defaultControlUiFeatureMethods, + "desktop.observe", + "browser.request", + "board.get", + "controlUi.sessionPullRequests.subscribe", + "sessions.github.publish", + ], + methodResponses: { + "question.list": { questions: [] }, + "environments.list": { environments: [] }, + "board.get": { revision: 1, tabs: [], widgets: [] }, + "sessions.github.options": { + shared: null, + personal: null, + pendingPersonal: null, + latestShared: null, + }, + "controlUi.sessionPullRequests.subscribe": { ok: true }, + }, + sessionKey: "agent:main:main", + sessions: [], + historyMessages: [], + agentModel: "fixture/echo", + models: [{ id: "echo", name: "Echo", provider: "fixture", contextWindow: 128_000 }], + }); + await page.addInitScript(() => { + const wire = (window.readCoalescingWire = { + pending: new Map(), + duplicates: [] as string[], + }); + const methods = new Set(["sessions.describe", "sessions.branches.list", "models.list"]); + const signature = ({ method, params }: MockGatewayRequest) => + JSON.stringify([ + method, + Object.entries(params ?? {}).toSorted(([a], [b]) => a.localeCompare(b)), + ]); + const instrument = (Base: typeof WebSocket) => + class extends Base { + override send(data: Parameters[0]) { + if (typeof data === "string") { + const frame = JSON.parse(data) as MockGatewayRequest & { type: string }; + if (frame.type === "req") { + const key = signature(frame); + if ( + methods.has(frame.method) && + [...wire.pending.values()].some((other) => signature(other) === key) + ) { + wire.duplicates.push(key); + } + wire.pending.set(frame.id, frame); + } + } + return super.send(data); + } + override dispatchEvent(event: Event) { + if (event instanceof MessageEvent) { + const frame = JSON.parse(String(event.data)) as { type: string; id: string }; + if (frame.type === "res") { + wire.pending.delete(frame.id); + } + } + return super.dispatchEvent(event); + } + }; + // Init-script order is unspecified; instrument the mock's replacement too. + let socket = instrument(window.WebSocket); + Object.defineProperty(window, "WebSocket", { + configurable: true, + get: () => socket, + set: (next: typeof WebSocket) => { + socket = instrument(next); + }, + }); + }); + await page.goto(`${suite.server.baseUrl}new`); + await page.clock.install(); + await pauseVirtualClock(page); + await page.evaluate(() => { + Math.random = () => 0; + }); + const count = async (method: string, match?: Record) => + (await gateway.getRequests(method, match)).length; + const settle = async (heldMethod?: string) => { + await expect + .poll(async () => { + await page.clock.runFor(1); + return page.evaluate( + (held) => + [...window.readCoalescingWire.pending.values()] + .filter(({ method }) => method !== held) + .map(({ method }) => method), + heldMethod, + ); + }) + .toEqual([]); + }; + await gateway.deferNext("sessions.create"); + const composer = page.locator(".new-session-page__message"); + await composer.fill("Say a short hello."); + await composer.press("Enter"); + const { params } = await gateway.waitForRequest("sessions.create"); + const { key } = params as { key: string }; + const runId = "fixture-run"; + await gateway.resolveDeferred("sessions.create", { + ok: true, + key, + sessionId: "fixture-session", + runId, + runStarted: true, + status: "started", + entry: { + kind: "direct", + parentSessionKey: "agent:main:main", + model: "echo", + modelProvider: "fixture", + updatedAt: await page.evaluate(() => Date.now()), + activeRunIds: [runId], + lastRunId: runId, + }, + }); + // The swarm owner connects after 250 ms; its initial read must settle too. + await expect + .poll(async () => { + await page.clock.runFor(20); + return count("sessions.list", { spawnedBy: key }); + }) + .toBeGreaterThan(0); + await settle(); + await page.locator(".agent-chat__composer-combobox textarea").waitFor(); + expect(await count("board.get")).toBeGreaterThan(0); + expect(await count("sessions.branches.list", { sessionKey: key })).toBeGreaterThan(0); + expect(await count("sessions.describe", { key: "agent:main:main" })).toBeGreaterThan(0); + const describes = () => count("sessions.describe", { key }); + const before = await describes(); + expect(before).toBeGreaterThan(0); + let row = await gateway.getSessionRow(key); + const messages: unknown[] = []; + const emit = async (event: string, extra: Record) => { + const now = await page.evaluate(() => Date.now()); + row = { ...row, snapshotAt: now, updatedAt: now }; + await gateway.setSessionsListResponse({ + sessions: [row], + count: 1, + totalCount: 1, + ts: now, + }); + await gateway.setHistoryMessages(messages); + await gateway.emitGatewayEvent(event, { + sessionKey: key, + agentId: "main", + sessionId: row.sessionId, + runId, + ts: now, + ...(event === "sessions.changed" || event === "session.message" + ? { ...row, session: row } + : {}), + ...extra, + }); + await settle(); + row = await gateway.getSessionRow(key); + }; + const changed = (reason?: string, phase?: string) => + emit("sessions.changed", { ...(reason ? { reason } : {}), ...(phase ? { phase } : {}) }); + const message = (role: string, seq: number, text: string) => ({ + role, + content: [{ type: "text", text }], + __openclaw: { id: `fixture-${seq}`, seq }, + }); + const user = message("user", 1, "Say a short hello."); + const assistant = message("assistant", 2, "Hello from the fixture."); + // Replay the captured ordering without patch invalidations that absorb certification. + // The missing parent deliberately leaves ancestorSessions absent on every event. + const burstStarted = await page.evaluate(() => Date.now()); + await changed("send"); + messages.push(user); + await emit("session.message", { + message: user, + messageId: "fixture-1", + messageSeq: 1, + senderIsOwner: true, + }); + await changed("participants"); + await changed("send"); + await emit("agent", { stream: "lifecycle", seq: 3, data: { phase: "start" } }); + await changed(undefined, "start"); + await changed("agent.run.started"); + await changed(undefined, "model"); + await emit("chat", { seq: 6, state: "delta", message: assistant }); + messages.push(assistant); + await emit("session.message", { message: assistant, messageId: "fixture-2", messageSeq: 2 }); + await changed(undefined, "model"); + row = { ...row, hasActiveRun: false, activeRunIds: [], status: "done" }; + await emit("agent", { stream: "lifecycle", seq: 10, data: { phase: "end" } }); + await emit("chat", { seq: 10, state: "final", message: assistant }); + await changed("agent.input.settled"); + await changed(undefined, "end"); + await page.getByText("Hello from the fixture.", { exact: true }).first().waitFor(); + expect.soft(await describes(), "no per-event descriptors").toBe(before); + await page.clock.runFor(4_999 - ((await page.evaluate(() => Date.now())) - burstStarted)); + expect.soft(await describes(), "no descriptor before collection expires").toBe(before); + await page.clock.runFor(1); + expect.soft(await describes(), "one shared authoritative descriptor").toBe(before + 1); + await settle(); + + const scope = { sessionKey: key }; + const modelsBefore = await count("models.list", scope); + const metadataBefore = await count("chat.metadata", scope); + await gateway.deferNext("models.list", scope); + await changed("patch"); + await page.clock.runFor(2_500); + await settle("models.list"); + expect(await count("chat.metadata", scope)).toBe(metadataBefore + 1); + expect(await count("models.list", scope)).toBe(modelsBefore + 1); + await gateway.resolveDeferred("models.list"); + await settle(); + expect + .soft(await count("models.list", scope), "no replacement catalog read") + .toBe(modelsBefore + 1); + expect.soft(await describes(), "bounded total descriptors").toBeLessThanOrEqual(8); + expect(await page.evaluate(() => window.readCoalescingWire.duplicates)).toEqual([]); + }); + }); +}); diff --git a/ui/src/lib/chat/chat-metadata-cache.ts b/ui/src/lib/chat/chat-metadata-cache.ts index d3fe9193b808..b0331a997c9b 100644 --- a/ui/src/lib/chat/chat-metadata-cache.ts +++ b/ui/src/lib/chat/chat-metadata-cache.ts @@ -11,6 +11,7 @@ import { modelCatalogKey, modelCatalogParams, type ModelCatalogInvalidation, + type ModelCatalogRead, } from "../model-catalog-cache.ts"; import { readSessionChangedEvent } from "../sessions/reconcile.ts"; import type { UiSessionDefaultsHost } from "../sessions/session-key.ts"; @@ -57,7 +58,7 @@ export type ChatMetadataEntry = { writer?: object; refreshRevision: number; refreshAfter?: number; - validateCatalog?: boolean; + validateCatalog?: ReadonlySet; catalogRevision: number; refresh?: ChatMetadataRefreshRecord; listeners: Map<(update: ChatMetadataUpdate) => void, () => boolean>; diff --git a/ui/src/lib/chat/chat-metadata-refresh.test.ts b/ui/src/lib/chat/chat-metadata-refresh.test.ts index b645f65e4b9d..d92b559a3765 100644 --- a/ui/src/lib/chat/chat-metadata-refresh.test.ts +++ b/ui/src/lib/chat/chat-metadata-refresh.test.ts @@ -1,10 +1,14 @@ import { afterEach, describe, expect, it, vi } from "vitest"; import { createDeferred } from "../../../../test/helpers/promise.js"; import { createTestGatewayClient } from "../../test-helpers/gateway-client.ts"; +import { invalidateModelCatalogCache } from "../model-catalog-cache.ts"; +import { peekModelCatalog } from "../model-catalog-store.ts"; import { invalidateChatMetadataForSessionEvent, invalidateChatMetadataStore, type ChatMetadataResult, + type ChatMetadataResponse, + type ChatMetadataUpdate, } from "./chat-metadata-cache.ts"; import { beginChatMetadataPublication, @@ -22,6 +26,89 @@ const models = [{ id: "fresh", name: "Fresh", provider: "test" }]; afterEach(() => vi.useRealTimers()); describe("automatic metadata admission", () => { + it.each(["matching", "different", "failed", "retired"] as const)( + "publishes commands immediately before validating the pending %s catalog after a session patch", + async (outcome) => { + vi.useFakeTimers(); + const metadata = createDeferred(); + const catalog = createDeferred<{ models: typeof models }>(); + const updates: ChatMetadataUpdate[] = []; + const catalogsAtChange: ReturnType[] = []; + const refreshes: ReturnType[] = []; + const request = vi.fn((method: string) => { + if (method === "chat.metadata") { + return metadata.promise; + } + return catalog.promise; + }); + const client = createTestGatewayClient(request); + const release = subscribeChatMetadata(client, scope, (update) => { + updates.push(update); + if (update.type === "result" && update.catalogChanged) { + catalogsAtChange.push(peekModelCatalog(client, scope)); + } + if (update.type === "invalidated" || update.type === "result") { + refreshes.push(loadChatMetadataRefresh(client, scope)); + } + }); + try { + invalidateChatMetadataForSessionEvent(client, { ...scope, reason: "patch" }, {}); + await vi.advanceTimersByTimeAsync(2_499); + expect(request).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(1); + expect(request.mock.calls.map(([method]) => method)).toEqual([ + "models.list", + "chat.metadata", + ]); + const commandsRead = loadChatMetadata(client, scope); + metadata.resolve({ ...commands, models }); + await vi.advanceTimersByTimeAsync(0); + expect(updates.filter((update) => update.type === "result")).toEqual([ + { type: "result", result: commands }, + ]); + await expect(commandsRead).resolves.toEqual(commands); + expect(peekChatMetadata(client, scope)).toEqual(commands); + expect(request.mock.calls.filter(([method]) => method === "models.list")).toHaveLength(1); + + request.mockImplementation((method) => + method === "chat.metadata" ? Promise.resolve(commands) : Promise.resolve({ models }), + ); + if (outcome === "retired") { + invalidateModelCatalogCache(client, scope); + } + if (outcome === "failed") { + catalog.reject(new Error("Catalog unavailable")); + } else { + catalog.resolve({ models: outcome === "different" ? [] : models }); + } + await commandsRead; + await vi.advanceTimersByTimeAsync(0); + await Promise.all(refreshes.map((refresh) => refresh.completed)); + const changed = outcome !== "matching"; + expect(updates.filter((update) => update.type === "result")).toEqual([ + { type: "result", result: commands }, + ...(changed ? [{ type: "result", result: commands, catalogChanged: true }] : []), + ]); + expect(catalogsAtChange).toEqual(changed ? [undefined] : []); + expect(updates[0]).toEqual({ + type: "invalidated", + scope: "session", + refreshSessionFacts: true, + }); + expect(request.mock.calls.filter(([method]) => method === "models.list")).toHaveLength( + changed ? 2 : 1, + ); + expect(peekModelCatalog(client, scope)).toEqual({ models }); + expect(refreshes[0]?.isCurrent()).toBe(!changed); + } finally { + metadata.resolve(commands); + catalog.resolve({ models }); + await Promise.all(refreshes.map((refresh) => refresh.completed)); + release(); + } + }, + ); + it.each(["visible", "hidden", "released", "global invalidation"])( "coalesces session patches and rechecks admission when %s", async (transition) => { diff --git a/ui/src/lib/chat/chat-metadata-store.test.ts b/ui/src/lib/chat/chat-metadata-store.test.ts index 3acb27e80dff..5235a062e773 100644 --- a/ui/src/lib/chat/chat-metadata-store.test.ts +++ b/ui/src/lib/chat/chat-metadata-store.test.ts @@ -6,6 +6,7 @@ import { afterEach, describe, expect, it, vi } from "vitest"; import { createDeferred as deferred } from "../../../../test/helpers/promise.js"; import { GatewayRequestError, type GatewayBrowserClient } from "../../api/gateway.ts"; import { gatewayHelloForMethods } from "../../test-helpers/gateway-methods.ts"; +import { invalidateModelCatalogCache } from "../model-catalog-cache.ts"; import { loadModelCatalog, peekModelCatalog } from "../model-catalog-store.ts"; import { invalidateChatMetadataForSessionEvent, @@ -46,6 +47,134 @@ afterEach(() => { }); describe("chat metadata store", () => { + it("publishes commands immediately before validating a queued catalog dispatched after the patch", async () => { + vi.useFakeTimers(); + const older = deferred<{ models: [] }>(); + const current = deferred<{ models: [] }>(); + const catalogs = vi.fn().mockReturnValueOnce(older.promise).mockReturnValue(current.promise); + const client = clientWith( + vi.fn((method: string) => + method === "models.list" + ? catalogs() + : Promise.resolve({ ...metadata("current"), models: [] }), + ), + ); + const scope = { agentId: "main", sessionKey: "agent:main:current" }; + const listener = vi.fn(); + const release = subscribeChatMetadata(client, scope, listener); + const retired = loadModelCatalog(client, scope).catch(() => undefined); + invalidateModelCatalogCache(client, scope); + const queued = loadModelCatalog(client, scope); + invalidateChatMetadataForSessionEvent(client, { ...scope, reason: "patch" }, {}); + older.reject(new Error("Old catalog unavailable")); + await vi.advanceTimersByTimeAsync(0); + expect(catalogs).toHaveBeenCalledTimes(2); + const read = loadChatMetadata(client, scope); + try { + await vi.advanceTimersByTimeAsync(0); + expect(listener.mock.calls.filter(([update]) => update.type === "result")).toEqual([ + [{ type: "result", result: metadata("current") }], + ]); + await expect(read).resolves.toEqual(metadata("current")); + current.resolve({ models: [] }); + await read; + expect(listener).toHaveBeenLastCalledWith({ type: "result", result: metadata("current") }); + expect(peekModelCatalog(client, scope)).toEqual({ models: [] }); + expect(listener.mock.calls.filter(([update]) => update.type === "result")).toHaveLength(1); + } finally { + current.resolve({ models: [] }); + await Promise.all([read, retired, queued]); + release(); + } + }); + + it.each(["older", "other-session", "other-account", "other-view"])( + "does not wait for an ineligible %s catalog during patch validation", + async (kind) => { + vi.useFakeTimers(); + const catalog = deferred<{ models: [] }>(); + const response = { ...metadata("current"), models: [] }; + const request = vi.fn((method: string) => + method === "chat.metadata" ? Promise.resolve(response) : catalog.promise, + ); + const client = clientWith(request); + const scope = { agentId: "main", sessionKey: "agent:main:current" }; + const listener = vi.fn(); + const release = subscribeChatMetadata(client, scope, listener); + const catalogScope = { + ...scope, + ...(kind === "other-session" ? { sessionKey: "agent:main:other" } : {}), + ...(kind === "other-account" ? { authProfileId: "test-account" } : {}), + ...(kind === "other-view" ? { includeDetails: true } : {}), + }; + let pending = kind === "older" ? loadModelCatalog(client, catalogScope) : undefined; + invalidateChatMetadataForSessionEvent(client, { ...scope, reason: "patch" }, {}); + pending ??= loadModelCatalog(client, catalogScope); + const read = loadChatMetadata(client, scope); + try { + await vi.advanceTimersByTimeAsync(0); + expect(listener).toHaveBeenLastCalledWith({ + type: "result", + result: metadata("current"), + catalogChanged: true, + }); + } finally { + catalog.resolve({ models: [] }); + await Promise.all([read, pending]); + release(); + } + }, + ); + + it.each(["invalidate", "release"])( + "publishes commands immediately but suppresses deferred catalog validation after %s", + async (transition) => { + vi.useFakeTimers(); + const catalog = deferred<{ models: [] }>(); + const response = { + ...metadata("obsolete"), + models: [{ id: "different", name: "Different", provider: "test" }], + }; + const client = clientWith( + vi.fn((method: string) => + method === "chat.metadata" ? Promise.resolve(response) : catalog.promise, + ), + ); + const scope = { agentId: "main", sessionKey: "agent:main:current" }; + const listener = vi.fn(); + const release = subscribeChatMetadata(client, scope, listener); + invalidateChatMetadataForSessionEvent(client, { ...scope, reason: "patch" }, {}); + const pending = loadModelCatalog(client, scope); + const read = loadChatMetadata(client, scope); + try { + await vi.advanceTimersByTimeAsync(0); + expect(listener.mock.calls.filter(([update]) => update.type === "result")).toEqual([ + [{ type: "result", result: metadata("obsolete") }], + ]); + await expect(read).resolves.toEqual(metadata("obsolete")); + if (transition === "invalidate") { + invalidateChatMetadataForSessionEvent(client, { ...scope, reason: "patch" }, {}); + } else { + release(); + } + catalog.resolve({ models: [] }); + await Promise.all([read, pending]); + await vi.advanceTimersByTimeAsync(0); + expect(peekChatMetadata(client, scope)).toEqual( + transition === "invalidate" ? undefined : metadata("obsolete"), + ); + expect(peekModelCatalog(client, scope)).toEqual({ models: [] }); + expect(listener.mock.calls.filter(([update]) => update.type === "result")).toEqual([ + [{ type: "result", result: metadata("obsolete") }], + ]); + } finally { + catalog.resolve({ models: [] }); + await Promise.all([read, pending]); + release(); + } + }, + ); + it("preserves only subscribed exact catalog scopes during session validation", async () => { const client = clientWith(vi.fn().mockResolvedValue({ models: [] })); const scope = { agentId: "main", sessionKey: "agent:main:main" }; @@ -161,7 +290,7 @@ describe("chat metadata store", () => { expect(peekChatMetadata(client, scope)).toEqual(metadata("locked")); const lateStartup = beginChatMetadataPublication(client, scope); second(); - lateStartup.publish(metadata("late")); + void lateStartup.publish(metadata("late")); expect(peekChatMetadata(client, scope)).toEqual(metadata("locked")); expect(peekChatMetadata(client, { agentId: "main" })).toEqual(metadata("neutral")); }); @@ -188,7 +317,7 @@ describe("chat metadata store", () => { invalidateChatMetadataStore(client); } expect.soft(request).toHaveBeenCalledOnce(); - startup.publish(metadata("obsolete-startup")); + void startup.publish(metadata("obsolete-startup")); if (outcome === "error") { older.reject(new Error("obsolete-read")); } else { diff --git a/ui/src/lib/chat/chat-metadata-store.ts b/ui/src/lib/chat/chat-metadata-store.ts index 6f7322457550..57040d3139b7 100644 --- a/ui/src/lib/chat/chat-metadata-store.ts +++ b/ui/src/lib/chat/chat-metadata-store.ts @@ -9,12 +9,14 @@ import type { GatewayBrowserClient } from "../../api/gateway.ts"; import type { ModelCatalogResult } from "../../api/types.ts"; import { invalidateModelCatalogCache, + getModelCatalogCache, modelCatalogKey, modelCatalogParams, } from "../model-catalog-cache.ts"; import { loadModelCatalog, peekModelCatalog, + pendingModelCatalogResult, settleModelCatalogRequests, subscribeModelCatalogCache, } from "../model-catalog-store.ts"; @@ -42,12 +44,8 @@ function notifyChatMetadataListeners(entry: ChatMetadataEntry, update: ChatMetad } } -function metadataScopeKey(scope: ChatMetadataParams): string { - return JSON.stringify([ - scope.agentId?.trim() ?? "", - scope.sessionKey ?? null, - scope.authProfileId ?? null, - ]); +function metadataScopeKey({ agentId, sessionKey, authProfileId }: ChatMetadataParams): string { + return JSON.stringify([agentId?.trim() ?? "", sessionKey ?? null, authProfileId ?? null]); } const MAX_CACHED_CHAT_METADATA = 64; @@ -104,10 +102,18 @@ function metadataEntryFor( // Retire every affected writer before subscribers can synchronously start replacements. const sessionOnly = scope?.sessionKey !== undefined && isSessionMetadataInvalidation(sessionEvent); + const catalog = sessionOnly ? getModelCatalogCache(client) : undefined; + const validation = + catalog && + new Set( + Array.from(catalog.requests.values()).flatMap((lanes) => + Array.from(lanes.values()).flatMap(({ active }) => (active ? active.read : [])), + ), + ); for (const entry of invalidated) { entry.refreshRevision += 1; entry.refreshAfter = sessionOnly ? Date.now() + SESSION_METADATA_DEBOUNCE_MS : undefined; - entry.validateCatalog = sessionOnly && entry.listeners.size > 0; + entry.validateCatalog = entry.listeners.size > 0 ? validation : undefined; entry.result = undefined; entry.writer = undefined; } @@ -120,10 +126,7 @@ function metadataEntryFor( entry.release(); } }; - cache = { - entries, - invalidate, - }; + cache = { entries, invalidate }; chatMetadataCache.set(client, cache); } const entries = cache.entries; @@ -218,25 +221,17 @@ async function requestChatMetadata( } } -function catalogProjectionKey( - models: ModelCatalogResult["models"], - accountSelection: ModelCatalogResult["accountSelection"], - modelSelectionPolicy: ModelCatalogResult["modelSelectionPolicy"], -) { +function catalogProjectionKey(projection: Partial) { // Metadata omits direct-picker policy, including on alternate runtime choices. return stableStringify([ - models.map(({ manualSelectionAllowed: _manual, runtimeChoices, ...model }) => ({ + projection.models?.map(({ manualSelectionAllowed: _manual, runtimeChoices, ...model }) => ({ ...model, - ...(runtimeChoices - ? { - runtimeChoices: runtimeChoices.map( - ({ manualSelectionAllowed: _choiceManual, ...choice }) => choice, - ), - } - : {}), + runtimeChoices: runtimeChoices?.map( + ({ manualSelectionAllowed: _choiceManual, ...choice }) => choice, + ), })), - accountSelection, - modelSelectionPolicy, + projection.accountSelection, + projection.modelSelectionPolicy, ]); } @@ -254,23 +249,36 @@ function preparePublication( const { models, accountSelection, modelSelectionPolicy, ...metadata } = result; if (isCurrent()) { let catalogChanged = false; - if (entry.validateCatalog) { - entry.validateCatalog = false; + const validateCatalog = entry.validateCatalog; + if (validateCatalog) { + entry.validateCatalog = undefined; + const catalogRevision = entry.catalogRevision; const catalog = peekModelCatalog(client, entry.scope); - // A patch can also change a session's account/runtime projection. Metadata - // is only an invalidation signal; the direct catalog remains the display owner. - if ( - !catalog || - models === undefined || - catalogProjectionKey(models, accountSelection, modelSelectionPolicy) !== - catalogProjectionKey( - catalog.models, - catalog.accountSelection, - catalog.modelSelectionPolicy, - ) - ) { + const hasCatalogChanged = (validatedCatalog: ModelCatalogResult | undefined) => + !validatedCatalog || + catalogRevision !== entry.catalogRevision || + catalogProjectionKey({ models, accountSelection, modelSelectionPolicy }) !== + catalogProjectionKey(validatedCatalog); + const pending = !catalog + ? pendingModelCatalogResult(client, entry.scope, validateCatalog) + : undefined; + if (pending) { + // Commands are ready now; only catalog validation waits for its existing producer. + void pending.then((validatedCatalog) => { + if (isCurrent() && hasCatalogChanged(validatedCatalog)) { + invalidateModelCatalogCache(client, entry.scope); + notifyChatMetadataListeners(entry, { + type: "result", + result: metadata, + catalogChanged: true, + }); + } + }); + } else { + catalogChanged = hasCatalogChanged(catalog); + } + if (catalogChanged) { invalidateModelCatalogCache(client, entry.scope); - catalogChanged = true; } } entry.result = metadata; @@ -403,7 +411,7 @@ export function subscribeChatMetadata( entry.refreshRevision += 1; entry.writer = undefined; if (entry.validateCatalog) { - entry.validateCatalog = false; + entry.validateCatalog = undefined; invalidateModelCatalogCache(client, scope); } } diff --git a/ui/src/lib/model-catalog-store.ts b/ui/src/lib/model-catalog-store.ts index 944833fb0d14..7e1271cdb3cf 100644 --- a/ui/src/lib/model-catalog-store.ts +++ b/ui/src/lib/model-catalog-store.ts @@ -57,16 +57,13 @@ export function readAgentModelCatalog( client: ModelCatalogClient | null | undefined, agentId: string | null | undefined, ): ModelCatalogPresentation { - if (!client || !agentId) { - return { models: [], hasSnapshot: false, retired: false }; - } - const scope = { agentId }; - const catalog = peekModelCatalog(client, scope, { allowStale: true }); + const catalog = + client && agentId ? peekModelCatalog(client, { agentId }, { allowStale: true }) : undefined; return { ...catalog, models: catalog?.models ?? [], hasSnapshot: catalog !== undefined, - retired: isModelCatalogRetired(client, scope), + retired: client && agentId ? isModelCatalogRetired(client, { agentId }) : false, }; } @@ -163,12 +160,33 @@ export function settleModelCatalogRequests( scope: ModelsListParams, ): Promise | undefined { const key = modelCatalogKey(modelCatalogParams(scope)); - const pending = Array.from(modelCatalogCache.get(client)?.requests.get(key)?.values() ?? []) - .map((lane) => lane.active?.transportSettled) - .filter((promise) => promise !== undefined); + const pending = Array.from( + modelCatalogCache.get(client)?.requests.get(key)?.values() ?? [], + ).flatMap(({ active }) => (active ? [active.transportSettled] : [])); return pending.length ? Promise.allSettled(pending).then(() => {}) : undefined; } +/** Observe an eligible producer without joining its cancellation or publication ownership. */ +export function pendingModelCatalogResult( + client: ModelCatalogClient, + scope: ModelsListParams, + issuedBefore: ReadonlySet, +): Promise | undefined { + const cache = modelCatalogCache.get(client); + const key = modelCatalogKey(modelCatalogParams(scope)); + const pending = Array.from(cache?.requests.get(key)?.values() ?? []).find( + ({ active }) => active && cache?.reads.has(active.read) && !issuedBefore.has(active.read), + )?.active; + return pending?.promise.then( + () => + modelCatalogCache.get(client) === cache && + cache?.entries.get(key)?.publishedRead === pending.read.order + ? peekModelCatalog(client, scope) + : undefined, + () => undefined, + ); +} + function createModelCatalogRequest(params: { client: ModelCatalogClient; scope: ModelsListParams; @@ -178,18 +196,22 @@ function createModelCatalogRequest(params: { queued: boolean; releaseLane: () => void; }): ModelCatalogRequest { - const { client, cache, lane, timeoutMs } = params; + const { client, cache, lane, timeoutMs: timeout } = params; const controller = new AbortController(); const completion = createDeferredCore(); const transportSettled = createDeferredCore(); const duration = - typeof timeoutMs === "number" && Number.isFinite(timeoutMs) - ? resolveSafeTimeoutDelayMs(timeoutMs, { minMs: 0 }) + typeof timeout === "number" && Number.isFinite(timeout) + ? resolveSafeTimeoutDelayMs(timeout, { minMs: 0 }) : undefined; const deadline = duration === undefined ? undefined : Date.now() + duration; let deadlineTimer: ReturnType | undefined; let started = false; let requestSent = false; + const rejectTimeout = (timeoutMs: number) => + pending.reject( + new GatewayProtocolRequestTimeoutError({ method: "models.list", timeoutMs, requestSent }), + ); const canRetry = () => !pending.settled && !controller.signal.aborted && @@ -198,7 +220,10 @@ function createModelCatalogRequest(params: { cache.reads.has(pending.read) && lane.active === pending && !lane.queued; - const retireCompletion = () => { + const settle = (complete: () => void) => { + if (pending.settled) { + return; + } clearTimeout(deadlineTimer); cache.reads.delete(pending.read); pending.settled = true; @@ -206,6 +231,7 @@ function createModelCatalogRequest(params: { lane.queued = undefined; params.releaseLane(); } + complete(); }; const finishTransport = () => { transportSettled.resolve(); @@ -230,18 +256,8 @@ function createModelCatalogRequest(params: { subscribers: new Set(), promise: completion.promise, transportSettled: transportSettled.promise, - resolve: (result) => { - if (!pending.settled) { - retireCompletion(); - completion.resolve(result); - } - }, - reject: (error) => { - if (!pending.settled) { - retireCompletion(); - completion.reject(error); - } - }, + resolve: (result) => settle(() => completion.resolve(result)), + reject: (error) => settle(() => completion.reject(error)), start: () => { if (started) { return; @@ -253,13 +269,7 @@ function createModelCatalogRequest(params: { } const remaining = deadline === undefined ? undefined : deadline - Date.now(); if (params.queued && duration !== undefined && remaining !== undefined && remaining <= 0) { - pending.reject( - new GatewayProtocolRequestTimeoutError({ - method: "models.list", - timeoutMs: duration, - requestSent: false, - }), - ); + rejectTimeout(duration); finishTransport(); return; } @@ -271,10 +281,10 @@ function createModelCatalogRequest(params: { try { // Only a received rejection can retry; local timeout still owns its transport. // The existing numeric deadline covers every attempt and wait in this lane. - const result = await (timeoutMs === undefined + const result = await (timeout === undefined ? client.request("models.list", requestParams) : client.request("models.list", requestParams, { - timeoutMs: duration === undefined ? timeoutMs : null, + timeoutMs: duration === undefined ? timeout : null, ...(duration === undefined ? {} : { @@ -320,13 +330,7 @@ function createModelCatalogRequest(params: { stopWatching(); } if (deadline !== undefined && duration !== undefined && deadline <= Date.now()) { - pending.reject( - new GatewayProtocolRequestTimeoutError({ - method: "models.list", - timeoutMs: duration, - requestSent, - }), - ); + rejectTimeout(duration); return; } if (!canRetry()) { @@ -344,17 +348,7 @@ function createModelCatalogRequest(params: { once: true, }); if (duration !== undefined) { - deadlineTimer = setTimeout( - () => - pending.reject( - new GatewayProtocolRequestTimeoutError({ - method: "models.list", - timeoutMs: duration, - requestSent, - }), - ), - duration, - ); + deadlineTimer = setTimeout(() => rejectTimeout(duration), duration); } return pending; } diff --git a/ui/src/lib/sessions/session-capability.ts b/ui/src/lib/sessions/session-capability.ts index 354153f8e6e2..ffde6d3eac68 100644 --- a/ui/src/lib/sessions/session-capability.ts +++ b/ui/src/lib/sessions/session-capability.ts @@ -218,7 +218,7 @@ export type SessionCapability = { captureConnectionScope: () => SessionConnectionScope | null; /** Whether a captured read-only request still belongs to the active connection. */ isConnectionScopeCurrent: (scope: SessionConnectionScope) => boolean; - /** Shares exact descriptor reads until the session changes; refresh supersedes earlier reads. */ + /** Shares descriptor reads, including agent-implied scopes, until the session changes; refresh supersedes earlier reads. */ describe: ( params: SessionsDescribeParams, options?: { refresh?: boolean; timeoutMs?: number; client?: SessionRequestClient }, diff --git a/ui/src/lib/sessions/session-describe.test.ts b/ui/src/lib/sessions/session-describe.test.ts index 364551e4040b..b57625a6670e 100644 --- a/ui/src/lib/sessions/session-describe.test.ts +++ b/ui/src/lib/sessions/session-describe.test.ts @@ -52,12 +52,38 @@ describe("session descriptor reads", () => { expect(h.read).toHaveBeenCalledTimes(1); }); + it("shares the implied agent scope while isolating a contradictory agent", async () => { + const h = harness(); + const pending = createDeferred<{ session: GatewaySessionRow }>(); + h.read.mockReturnValueOnce(pending.promise); + const reads = [ + h.sessions.describe({ key }), + h.sessions.describe({ key, agentId: "main" }), + h.sessions.describe({ key, agentId: " MAIN " }), + ]; + expect(h.read).toHaveBeenCalledTimes(1); + const other = { ...initial, sessionId: "other-agent" }; + h.read.mockResolvedValueOnce({ session: other }); + expect(await h.sessions.describe({ key, agentId: "other" })).toEqual({ session: other }); + pending.resolve({ session: initial }); + expect(await Promise.all(reads)).toEqual([ + { session: initial }, + { session: initial }, + { session: initial }, + ]); + expect(h.read).toHaveBeenCalledTimes(2); + }); + it("keeps session, agent, preview parameters and request timeout distinct", async () => { const h = harness(); for (const params of [ { key }, { key: "agent:main:other" }, { key, agentId: "work" }, + { key: "global" }, + { key: "global", agentId: "main" }, + { key: "local" }, + { key: "local", agentId: "main" }, { key, includeDerivedTitles: true }, { key, includeLastMessage: true }, ]) { @@ -65,7 +91,7 @@ describe("session descriptor reads", () => { await h.sessions.describe({ ...params }); } await h.sessions.describe({ key }, { timeoutMs: 30_000 }); - expect(h.read).toHaveBeenCalledTimes(6); + expect(h.read).toHaveBeenCalledTimes(10); }); it("preserves the runtime sample time when reusing a running descriptor", async () => { diff --git a/ui/src/lib/sessions/session-describe.ts b/ui/src/lib/sessions/session-describe.ts index 9cb215b07162..8f8947affb5e 100644 --- a/ui/src/lib/sessions/session-describe.ts +++ b/ui/src/lib/sessions/session-describe.ts @@ -73,57 +73,57 @@ export function createSessionDescribeReads(host: { } throw new Error("gateway not connected"); } + const impliedAgentId = parseAgentSessionKey(params.key)?.agentId; const key = JSON.stringify([ params.key, - params.agentId, + impliedAgentId && normalizeAgentId(impliedAgentId) === normalizeAgentId(params.agentId) + ? undefined + : params.agentId, params.includeDerivedTitles, params.includeLastMessage, options.timeoutMs, ]); - let read = reads.get(key); - if (read && (!current(read) || options.refresh || (!host.canReuse() && read.result))) { - reads.delete(key); - read = undefined; + const read = reads.get(key); + if (read && current(read) && !options.refresh && (!read.result || host.canReuse())) { + return read.promise; } - if (!read) { - const request = scope.client.request("sessions.describe", params, ...requestOptions); - const pending: Read = { - params: { ...params }, - scope, - promise: request, - issuedRow: host.currentRow(params), - reusable: host.canReuse(), - }; - pending.promise = request.then( - (response) => { - const result = - response.session?.runtimeMs === undefined - ? response - : { - ...response, - session: { ...response.session, runtimeSampledAt: Date.now() }, - }; - if (reads.get(key) === pending && host.connection.isCurrent(scope)) { - pending.result = result; - if (!pending.reusable || !host.canReuse() || !current(pending)) { - reads.delete(key); - } - trim(); - } - return result; - }, - (error: unknown) => { - if (reads.get(key) === pending) { + reads.delete(key); + const request = scope.client.request("sessions.describe", params, ...requestOptions); + const pending: Read = { + params: { ...params }, + scope, + promise: request, + issuedRow: host.currentRow(params), + reusable: host.canReuse(), + }; + pending.promise = request.then( + (response) => { + const result = + response.session?.runtimeMs === undefined + ? response + : { + ...response, + session: { ...response.session, runtimeSampledAt: Date.now() }, + }; + if (reads.get(key) === pending && host.connection.isCurrent(scope)) { + pending.result = result; + if (!pending.reusable || !host.canReuse() || !current(pending)) { reads.delete(key); } - throw error; - }, - ); - read = pending; - reads.set(key, read); - trim(); - } - return read.promise; + trim(); + } + return result; + }, + (error: unknown) => { + if (reads.get(key) === pending) { + reads.delete(key); + } + throw error; + }, + ); + reads.set(key, pending); + trim(); + return pending.promise; }; return { diff --git a/ui/src/lib/sessions/session-reconciliation.ts b/ui/src/lib/sessions/session-reconciliation.ts index 1b471451d641..ae4a7d57b636 100644 --- a/ui/src/lib/sessions/session-reconciliation.ts +++ b/ui/src/lib/sessions/session-reconciliation.ts @@ -642,7 +642,14 @@ export function createSessionReconciliation(host: Host) { eventResult: reduced, ...(!reduced.deletedKey && (!reduced.admittedRow || !Array.isArray(asNullableRecord(snapshot)?.ancestorSessions)) - ? { invalidateRevision: eventObservation.revision } + ? reduced.admittedRow && asNullableRecord(asNullableRecord(snapshot)?.session) + ? { + certification: { + revision: eventObservation.revision, + reason: invalidationReason, + }, + } + : { invalidateRevision: eventObservation.revision } : {}), }; }, diff --git a/ui/src/lib/sessions/session-roster-observations.ts b/ui/src/lib/sessions/session-roster-observations.ts index f8ff49a727cb..ebd2725afdcb 100644 --- a/ui/src/lib/sessions/session-roster-observations.ts +++ b/ui/src/lib/sessions/session-roster-observations.ts @@ -1,4 +1,5 @@ import type { GatewaySessionRow, SessionsListResult } from "../../api/types.ts"; +import { createSessionEventRefreshCoordinator } from "./event-refresh-coordinator.ts"; import { projectSessionResultRows } from "./reconcile.ts"; import type { SessionConnectionOwner, @@ -49,6 +50,8 @@ type RegisteredSessionRow = { type RowProjection = (entry: ObservedSessionRow) => { row: GatewaySessionRow | null; invalidateRevision?: number; + invalidationReason?: string; + certification?: { revision: number; reason?: string }; observationRevision?: number; eventResult?: SessionChangedRowResult; }; @@ -76,6 +79,31 @@ export function createSessionRosterObservations( const provenance = createSessionRowProvenance(); const { owner, identity, inheritRow, mergeRow, rowRevision } = provenance; const registeredRows = new Set(); + const certifications = new Map(); + const certificationRefresh = createSessionEventRefreshCoordinator({ + active: true, + refresh: async () => { + stageManagedResults( + host.connection.capture(), + (entry) => entry.snapshot.result, + ({ target, row }) => ({ + row, + invalidateRevision: certifications.get(target)?.revision, + invalidationReason: certifications.get(target)?.reason, + }), + ).notify(); + }, + }); + const clearCertification = (target: SessionRowTarget, retired = false) => { + certifications.delete(target); + if (!certifications.size) { + if (retired) { + certificationRefresh.reset(); + } else { + certificationRefresh.absorb(); + } + } + }; const registrationIsAttached = (entry: RegisteredSessionRow) => registeredRows.has(entry) && entry.scope !== null && host.connection.isCurrent(entry.scope); const registrationIsCurrent = (entry: RegisteredSessionRow) => @@ -221,6 +249,7 @@ export function createSessionRosterObservations( entry: RegisteredSessionRow; previous: RegisteredSessionRow["snapshot"]; snapshot: RegisteredSessionRow["snapshot"]; + invalidationReason?: string; }> = []; for (const [key, entry] of lists) { if (entry.connectionEpoch !== scope.epoch) { @@ -273,6 +302,22 @@ export function createSessionRosterObservations( !previous.row || entry.onInvalidate ? Math.max(previous.invalidatedRevision, projected.invalidateRevision ?? 0) : previous.invalidatedRevision; + if (projected.certification && entry.onInvalidate) { + certifications.set(entry.target, projected.certification); + certificationRefresh.schedule(); + } + const certification = certifications.get(entry.target); + if ( + certification && + (invalidatedRevision >= certification.revision || + (admitRead && + admittedRows.some( + ({ row: admitted, revision }) => + revision > certification.revision && acceptsRow(entry, admitted, revision), + ))) + ) { + clearCertification(entry.target); + } const hasObserved = previous.hasObserved || row !== null || Boolean(projected.eventResult?.deletedKey); let nextSnapshot = previous; @@ -295,6 +340,7 @@ export function createSessionRosterObservations( entry, previous, snapshot: nextSnapshot, + invalidationReason: projected.invalidationReason, }); } if ( @@ -327,6 +373,9 @@ export function createSessionRosterObservations( change.entry.snapshot === change.previous ) { change.entry.snapshot = change.snapshot; + if (change.snapshot.retired) { + clearCertification(change.entry.target, true); + } changed = true; } } @@ -339,6 +388,7 @@ export function createSessionRosterObservations( const retirement = successorRetirement(entry, publishedRows); if (retirement) { entry.snapshot = retirement.snapshot; + clearCertification(entry.target, true); rowChanges.push(retirement); } } @@ -353,7 +403,7 @@ export function createSessionRosterObservations( listener(snapshot); } } - for (const { entry, previous, snapshot } of rowChanges) { + for (const { entry, previous, snapshot, invalidationReason } of rowChanges) { if (!host.connection.isCurrent(scope)) { return; } @@ -373,7 +423,7 @@ export function createSessionRosterObservations( entry.snapshot === snapshot && snapshot.invalidatedRevision > previous.invalidatedRevision ) { - entry.onInvalidate?.(reason); + entry.onInvalidate?.(invalidationReason ?? reason); } } } @@ -440,6 +490,8 @@ export function createSessionRosterObservations( // Retained cache rows carry presentation, not evidence from the retired connection. provenance.reset(); registeredRows.clear(); + certifications.clear(); + certificationRefresh.reset(); descriptions.clear(); }, registerRow( @@ -483,6 +535,7 @@ export function createSessionRosterObservations( invalidatedRevision: Math.max(entry.snapshot.invalidatedRevision, revision), }; entry.snapshot = snapshot; + clearCertification(entry.target); return () => { if (registrationIsCurrent(entry) && entry.snapshot === snapshot) { entry.listener(null); @@ -491,6 +544,7 @@ export function createSessionRosterObservations( }, dispose: () => { registeredRows.delete(entry); + clearCertification(entry.target, true); }, }; }, diff --git a/ui/src/lib/sessions/session-row-certification.test.ts b/ui/src/lib/sessions/session-row-certification.test.ts new file mode 100644 index 000000000000..ee298891c879 --- /dev/null +++ b/ui/src/lib/sessions/session-row-certification.test.ts @@ -0,0 +1,180 @@ +// @vitest-environment node +import { beforeEach, afterEach, describe, expect, it, vi } from "vitest"; +import type { GatewaySessionRow } from "../../api/types.ts"; +import { createTestGatewayClient } from "../../test-helpers/gateway-client.ts"; +import { + createGatewayHarness, + createTestSessionCapability, + sessionsResult, +} from "./session-capability.test-support.ts"; + +const workRow: GatewaySessionRow = { + key: "global", + agentId: "work", + sessionId: "work-global-session", + kind: "global", + updatedAt: 200, + label: "Work descriptor", + archived: false, + status: "running", + hasActiveRun: true, + activeRunIds: ["work-run"], +}; + +async function descriptorOwner(onInvalidate: () => void) { + const gateway = createGatewayHarness(createTestGatewayClient(async () => sessionsResult([], 1))); + const sessions = createTestSessionCapability(gateway.gateway); + await sessions.refresh({ agentId: "main", force: true }); + const target = { key: "global", agentId: "work" }; + const changed = vi.fn(); + const observation = sessions.observeRow(target, changed, { onInvalidate }); + expect(observation.captureReconcile()(workRow)).toMatchObject({ + status: "current", + row: workRow, + }); + changed.mockClear(); + return { ...gateway, sessions, target, changed, observation }; +} + +describe("descriptor certification refresh", () => { + beforeEach(() => { + vi.useFakeTimers(); + vi.spyOn(Math, "random").mockReturnValue(0); + }); + afterEach(() => { + vi.useRealTimers(); + vi.restoreAllMocks(); + }); + it("publishes uncertified rows immediately and paces both registrations in the same tick", async () => { + const invalidated = vi.fn(() => Date.now()); + const h = await descriptorOwner(invalidated); + const otherInvalidated = vi.fn(() => Date.now()); + const otherChanged = vi.fn(); + const other = h.sessions.observeRow(h.target, otherChanged, { + onInvalidate: otherInvalidated, + }); + const staleRow = h.observation.captureReconcile(); + const staleAbsence = other.captureReconcile(); + const beforeRefresh = h.observation.captureReconcile(); + const started = Date.now(); + for (const [offset, reason] of ["patch", "send"].entries()) { + if (offset) { + await vi.advanceTimersByTimeAsync(1_000); + } + const row = { ...workRow, updatedAt: 201 + offset, label: `Admitted ${offset}` }; + h.emitEvent({ + type: "event", + event: offset ? "session.message" : "sessions.changed", + payload: { agentId: "work", reason, session: row }, + }); + expect(h.observation.row).toEqual(row); + expect(other.row).toEqual(row); + expect(h.changed).toHaveBeenLastCalledWith(row); + expect(otherChanged).toHaveBeenLastCalledWith(row); + expect(invalidated).not.toHaveBeenCalled(); + expect(otherInvalidated).not.toHaveBeenCalled(); + } + // Pacing cannot let a pre-event row or absence replace the admitted facts. + expect(staleRow(workRow)).toMatchObject({ status: "current", row: { label: "Admitted 1" } }); + expect(staleAbsence(undefined)).toMatchObject({ + status: "current", + row: { label: "Admitted 1" }, + }); + await vi.advanceTimersByTimeAsync(3_999); + expect(invalidated).not.toHaveBeenCalled(); + expect(otherInvalidated).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(1); + expect(invalidated).toHaveBeenCalledExactlyOnceWith("send"); + expect(otherInvalidated).toHaveBeenCalledExactlyOnceWith("send"); + expect(invalidated.mock.results[0]?.value).toBe(started + 5_000); + expect(otherInvalidated.mock.results[0]?.value).toBe(started + 5_000); + expect(beforeRefresh(workRow)).toEqual({ status: "invalidated" }); + }); + + it("immediately invalidates an envelope-only event and absorbs pending certification", async () => { + const invalidated = vi.fn(); + const h = await descriptorOwner(invalidated); + h.emitEvent({ + type: "event", + event: "sessions.changed", + payload: { agentId: "work", reason: "patch", session: { ...workRow, updatedAt: 201 } }, + }); + expect(invalidated).not.toHaveBeenCalled(); + const pending = h.observation.captureReconcile(); + await vi.advanceTimersByTimeAsync(4_000); + h.emitEvent({ + type: "event", + event: "sessions.changed", + payload: { sessionKey: "global", agentId: "work", reason: "swarm" }, + }); + expect(invalidated).toHaveBeenCalledExactlyOnceWith("swarm"); + expect(pending(workRow)).toEqual({ status: "invalidated" }); + await vi.advanceTimersByTimeAsync(500); + h.emitEvent({ + type: "event", + event: "sessions.changed", + payload: { agentId: "work", reason: "send", session: { ...workRow, updatedAt: 202 } }, + }); + await vi.advanceTimersByTimeAsync(500); + expect(invalidated).toHaveBeenCalledOnce(); + await vi.advanceTimersByTimeAsync(4_500); + expect(invalidated).toHaveBeenCalledTimes(2); + expect(invalidated).toHaveBeenLastCalledWith("send"); + }); + + it.each(["read", "absence", "delete", "dispose", "reset"] as const)( + "absorbs pending certification after %s", + async (action) => { + const invalidated = vi.fn(); + const h = await descriptorOwner(invalidated); + const row = { ...workRow, updatedAt: 201 }; + h.emitEvent({ + type: "event", + event: "sessions.changed", + payload: { agentId: "work", reason: "patch", session: row }, + }); + expect(invalidated).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(4_000); + if (action === "read" || action === "absence") { + expect(h.observation.captureReconcile()(action === "read" ? row : undefined)).toMatchObject( + { + status: "current", + row: action === "read" ? row : null, + }, + ); + } else if (action === "delete") { + h.emitEvent({ + type: "event", + event: "sessions.changed", + payload: { + agentId: "work", + sessionKey: "global", + sessionId: row.sessionId, + reason: "delete", + ts: 202, + }, + }); + expect(h.observation.isCurrent()).toBe(false); + } else if (action === "dispose") { + h.observation.dispose(); + } else { + h.publish(false); + } + if (action === "read") { + await vi.advanceTimersByTimeAsync(500); + h.emitEvent({ + type: "event", + event: "sessions.changed", + payload: { agentId: "work", reason: "send", session: { ...row, updatedAt: 202 } }, + }); + await vi.advanceTimersByTimeAsync(500); + expect(invalidated).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(4_500); + expect(invalidated).toHaveBeenCalledExactlyOnceWith("send"); + return; + } + await vi.advanceTimersByTimeAsync(10_000); + expect(invalidated).not.toHaveBeenCalled(); + }, + ); +}); diff --git a/ui/src/pages/chat/chat-branch-freshness.test.ts b/ui/src/pages/chat/chat-branch-freshness.test.ts index fa872694d551..198799b18107 100644 --- a/ui/src/pages/chat/chat-branch-freshness.test.ts +++ b/ui/src/pages/chat/chat-branch-freshness.test.ts @@ -2,13 +2,55 @@ import { describe, expect, it, vi } from "vitest"; import { createDeferred } from "../../../../test/helpers/promise.js"; import type { SessionBranch } from "../../api/types.ts"; import { createTestGatewayClient } from "../../test-helpers/gateway-client.ts"; -import { loadChatBranches } from "./chat-history-branches.ts"; +import { + invalidateChatBranches, + loadChatBranches, + retireChatBranchRequests, +} from "./chat-history-branches.ts"; import { loadChatHistory } from "./chat-history.ts"; import { makeChatHost } from "./chat-host.test-support.ts"; import { handlePageGatewayEvent } from "./chat-state-events.ts"; import type { ChatPageHost } from "./chat-state-host.ts"; describe("chat branch freshness", () => { + it.each([ + ["invalidation", invalidateChatBranches], + ["retirement", retireChatBranchRequests], + ] as const)("shares an initial branch read until %s retires it", async (_reason, retire) => { + const first = createDeferred<{ branches: SessionBranch[] }>(); + const replacement = createDeferred<{ branches: SessionBranch[] }>(); + const listBranches = vi + .fn() + .mockReturnValueOnce(first.promise) + .mockReturnValueOnce(replacement.promise); + const state = makeChatHost({ + sessionKey: "agent:main:main", + connectionEpoch: 1, + requestHandlers: { "sessions.branches.list": listBranches }, + }); + const initial = loadChatBranches(state); + const joined = loadChatBranches(state); + expect(listBranches).toHaveBeenCalledOnce(); + + retire(state); + const refreshed = loadChatBranches(state); + const refreshedJoin = loadChatBranches(state); + expect(listBranches).toHaveBeenCalledTimes(2); + const current = { + leafEntryId: "current", + headline: "Current reply", + messageCount: 2, + active: true, + }; + replacement.resolve({ branches: [current] }); + await Promise.all([refreshed, refreshedJoin]); + first.resolve({ branches: [] }); + await Promise.all([initial, joined]); + expect(state.chatBranches).toEqual([current]); + expect(state.chatBranchesConnectionEpoch).toBe(1); + expect(state.chatBranchesSessionKey).toBe("agent:main:main"); + }); + function createSessionEventState(overrides: Partial = {}) { const request = vi.fn().mockResolvedValue({ messages: [], diff --git a/ui/src/pages/chat/chat-commands.ts b/ui/src/pages/chat/chat-commands.ts index ec5641ae2324..209bdb010506 100644 --- a/ui/src/pages/chat/chat-commands.ts +++ b/ui/src/pages/chat/chat-commands.ts @@ -194,17 +194,6 @@ function remoteSlashCommandCacheKey(agentId: string | undefined, sessionKey?: st return JSON.stringify([agentId ?? null, sessionKey ?? null]); } -function getRemoteSlashCommandCache( - client: GatewayBrowserClient, -): Map { - let cache = remoteSlashCommandCache.get(client); - if (!cache) { - cache = new Map(); - remoteSlashCommandCache.set(client, cache); - } - return cache; -} - async function requestRemoteSlashCommands( client: GatewayBrowserClient, agentId: string | undefined, @@ -237,7 +226,11 @@ function loadRemoteSlashCommands( if (Array.isArray(metadata?.commands)) { return Promise.resolve(buildSlashCommandsFromEntries(getRemoteCommandEntries(metadata))); } - const cache = getRemoteSlashCommandCache(client); + let cache = remoteSlashCommandCache.get(client); + if (!cache) { + cache = new Map(); + remoteSlashCommandCache.set(client, cache); + } const key = remoteSlashCommandCacheKey(agentId, sessionKey); const cached = cache.get(key); const now = Date.now(); @@ -301,18 +294,13 @@ export async function refreshSlashCommands(params: { }): Promise { const seq = ++refreshSeq; const agentId = params.agentId?.trim(); - if (!params.client) { - if (seq !== refreshSeq || params.shouldApply?.() === false) { - return; - } - replaceSlashCommands(buildFallbackSlashCommands()); - return; - } - const commands = await loadRemoteSlashCommands(params.client, agentId, params.sessionKey); + const commands = params.client + ? await loadRemoteSlashCommands(params.client, agentId, params.sessionKey) + : undefined; if (seq !== refreshSeq || params.shouldApply?.() === false) { return; } - replaceSlashCommands(commands); + replaceSlashCommands(commands ?? buildFallbackSlashCommands()); } export function shouldQueueLocalSlashCommand(name: string): boolean { diff --git a/ui/src/pages/chat/chat-history-branches.ts b/ui/src/pages/chat/chat-history-branches.ts index e23abe2356e0..a8324aeb1f07 100644 --- a/ui/src/pages/chat/chat-history-branches.ts +++ b/ui/src/pages/chat/chat-history-branches.ts @@ -4,8 +4,14 @@ import { areUiSessionKeysEquivalent } from "../../lib/sessions/session-key.ts"; import { chatHistoryRequests } from "./chat-history-state.ts"; import type { ChatState } from "./chat-state-contract.ts"; +const pendingBranchLoads = new WeakMap< + ChatState, + { matches: () => boolean; promise: Promise } +>(); + export function retireChatBranchRequests(state: ChatState): void { chatHistoryRequests(state).branchVersion += 1; + pendingBranchLoads.delete(state); } export function invalidateChatBranches(state: ChatState): void { @@ -24,37 +30,47 @@ export function displayedChatSessionBranches( } export async function loadChatBranches(state: ChatState): Promise { - const sessions = state.sessions; - const client = state.client; - const sessionKey = state.sessionKey; - if (!sessions?.listBranches || !client || !state.connected) { + const { sessions, client, sessionKey, connectionEpoch } = state; + const listBranches = sessions?.listBranches; + if (!listBranches || !client || !state.connected) { return; } + const pending = pendingBranchLoads.get(state); + if (pending?.matches()) { + return pending.promise; + } const requests = chatHistoryRequests(state); const version = ++requests.branchVersion; - const connectionEpoch = state.connectionEpoch; const agentParams = scopedAgentParamsForSession(state, sessionKey); - try { - const branches = await sessions.listBranches(sessionKey, agentParams); - if ( - requests.branchVersion !== version || - state.client !== client || - !state.connected || - state.connectionEpoch !== connectionEpoch || - !visibleSessionMatches(state, sessionKey, agentParams.agentId) - ) { - return; + const isCurrent = () => + requests.branchVersion === version && + state.client === client && + state.connected && + state.connectionEpoch === connectionEpoch && + visibleSessionMatches(state, sessionKey, agentParams.agentId); + const promise = (async () => { + try { + const branches = await listBranches.call(sessions, sessionKey, agentParams); + if (isCurrent()) { + state.chatBranches = branches; + state.chatBranchesSessionKey = sessionKey; + state.chatBranchesConnectionEpoch = connectionEpoch; + } + } catch { + // Leave the success receipt unset so the next history load retries transient failures. } - state.chatBranches = branches; - state.chatBranchesSessionKey = sessionKey; - state.chatBranchesConnectionEpoch = connectionEpoch; - } catch { - // Leave chatBranchesSessionKey unset so the next history load retries; - // recording success here latched transient failures into a permanently - // hidden branch dropdown with no visible outcome. - } finally { + })().finally(() => { if (requests.branchVersion === version) { + pendingBranchLoads.delete(state); state.requestUpdate?.(); } - } + }); + pendingBranchLoads.set(state, { + matches: () => + isCurrent() && + state.sessionKey === sessionKey && + scopedAgentParamsForSession(state, sessionKey).agentId === agentParams.agentId, + promise, + }); + return promise; }