mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
perf(gateway): bound session catalog provider waits (#159457)
Return completed catalogs within a one-second provider budget and retain marked stale pages for slow or unavailable providers. Keep late refreshes owned and cached adoption identities isolated through delivery. Related: #127743
This commit is contained in:
parent
af4a236330
commit
4a2115c767
6 changed files with 558 additions and 63 deletions
|
|
@ -8,6 +8,21 @@ title: "Node session catalogs"
|
|||
sidebarTitle: "Session catalogs"
|
||||
---
|
||||
|
||||
Catalog listing waits up to one second per provider, concurrently. Providers that
|
||||
finish within that budget return normally. A slow provider returns
|
||||
`catalog.error.code: "catalog_pending"`; when available, its last successful page
|
||||
for the same caller and query is included with hosts marked `pending: true` and
|
||||
a stale-results message. A failed refresh can likewise return the last page with
|
||||
`catalog.error.code: "catalog_stale"` and the underlying error in its message.
|
||||
Other providers remain usable. Late results refresh the page for the next list;
|
||||
existing host progress updates remain supported.
|
||||
|
||||
These fallback pages are bounded in memory and invalidated by configuration or
|
||||
provider registration changes and catalog archive operations. Every delivery
|
||||
rechecks current visibility and local session identity. The one-second budget
|
||||
covers provider discovery and queueing, not session-projection preparation or
|
||||
time spent waiting for the Gateway event loop.
|
||||
|
||||
## Codex sessions and transcripts
|
||||
|
||||
The official `codex` plugin can expose non-archived Codex sessions on a
|
||||
|
|
|
|||
328
src/gateway/server-methods/session-catalog-budget.test.ts
Normal file
328
src/gateway/server-methods/session-catalog-budget.test.ts
Normal file
|
|
@ -0,0 +1,328 @@
|
|||
import { afterEach, beforeEach, expect, it, vi } from "vitest";
|
||||
import {
|
||||
getActiveGatewayRootWorkCount,
|
||||
tryBeginGatewayRootWorkAdmission,
|
||||
} from "../../process/gateway-work-admission.js";
|
||||
import { createDeferredCore } from "../../shared/deferred.js";
|
||||
import {
|
||||
call,
|
||||
hoisted,
|
||||
markPluginRegistryActive,
|
||||
provider,
|
||||
resetSessionCatalogTestState,
|
||||
setSessionCatalogEntries,
|
||||
startCall,
|
||||
type SessionCatalogProvider,
|
||||
} from "./session-catalog.test-helpers.js";
|
||||
|
||||
const { getActivePluginRegistry } = await import("../../plugins/runtime.js");
|
||||
|
||||
beforeEach(() => {
|
||||
resetSessionCatalogTestState();
|
||||
vi.useFakeTimers();
|
||||
});
|
||||
afterEach(() => vi.useRealTimers());
|
||||
|
||||
it("returns three fast catalogs within one second while an eight-second provider finishes later", async () => {
|
||||
const host = {
|
||||
hostId: "gateway:local",
|
||||
label: "Local",
|
||||
kind: "gateway" as const,
|
||||
connected: true,
|
||||
sessions: [],
|
||||
};
|
||||
hoisted.activeRegistry.sessionCatalogs = [
|
||||
...["fast-a", "fast-b", "fast-c"].map((id) => ({
|
||||
provider: provider(id, { list: async () => [host] }),
|
||||
})),
|
||||
{
|
||||
provider: provider("slow", {
|
||||
list: async () => {
|
||||
await new Promise<void>((resolve) => {
|
||||
setTimeout(resolve, 8_000);
|
||||
});
|
||||
return [host];
|
||||
},
|
||||
}),
|
||||
},
|
||||
];
|
||||
const startedAt = Date.now();
|
||||
const rootsBefore = getActiveGatewayRootWorkCount();
|
||||
const root = tryBeginGatewayRootWorkAdmission("catalog-budget-fixture");
|
||||
expect(root).not.toBeNull();
|
||||
const pending = await root!.run(async () => startCall("sessions.catalog.list", {}));
|
||||
let elapsed: number | undefined;
|
||||
void pending.completion.then(() => {
|
||||
elapsed = Date.now() - startedAt;
|
||||
});
|
||||
try {
|
||||
await vi.advanceTimersByTimeAsync(1_000);
|
||||
expect(pending.respond).toHaveBeenCalledWith(true, {
|
||||
catalogs: [
|
||||
...["fast-a", "fast-b", "fast-c"].map((id) =>
|
||||
expect.objectContaining({ id, hosts: [host] }),
|
||||
),
|
||||
expect.objectContaining({
|
||||
id: "slow",
|
||||
error: expect.objectContaining({ code: "catalog_pending" }),
|
||||
}),
|
||||
],
|
||||
});
|
||||
expect(elapsed).toBeLessThanOrEqual(1_000);
|
||||
root!.release();
|
||||
expect(getActiveGatewayRootWorkCount()).toBe(rootsBefore + 1);
|
||||
} finally {
|
||||
await vi.advanceTimersByTimeAsync(7_000);
|
||||
await pending.completion;
|
||||
root!.release();
|
||||
expect(getActiveGatewayRootWorkCount()).toBe(rootsBefore);
|
||||
console.info("catalog slow-provider fixture", {
|
||||
providers: 4,
|
||||
providerDelayMs: 8_000,
|
||||
responseMs: elapsed,
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
it("reuses a pending provider, serves its stale page, and retains its late refresh after errors", async () => {
|
||||
const oldHost = {
|
||||
hostId: "node:slow",
|
||||
label: "Old page",
|
||||
kind: "node" as const,
|
||||
connected: true,
|
||||
sessions: [],
|
||||
};
|
||||
const freshHost = { ...oldHost, label: "Late page" };
|
||||
const late = createDeferredCore<(typeof freshHost)[]>();
|
||||
const list = vi
|
||||
.fn<SessionCatalogProvider["list"]>()
|
||||
.mockResolvedValueOnce([oldHost])
|
||||
.mockImplementationOnce(() => late.promise)
|
||||
.mockRejectedValue(Object.assign(new Error("node did not respond"), { code: "TIMEOUT" }));
|
||||
hoisted.activeRegistry.sessionCatalogs = [{ provider: provider("fixture", { list }) }];
|
||||
const config = {};
|
||||
const client = { connId: "owner" };
|
||||
await call("sessions.catalog.list", {}, config, client);
|
||||
try {
|
||||
for (let index = 0; index < 3; index++) {
|
||||
const waiting = startCall("sessions.catalog.list", {}, config, client);
|
||||
await vi.advanceTimersByTimeAsync(1_000);
|
||||
await waiting.completion;
|
||||
expect(waiting.respond).toHaveBeenCalledWith(true, {
|
||||
catalogs: [
|
||||
expect.objectContaining({
|
||||
hosts: [{ ...oldHost, pending: true }],
|
||||
error: expect.objectContaining({
|
||||
code: "catalog_pending",
|
||||
message: expect.stringContaining("stale"),
|
||||
}),
|
||||
}),
|
||||
],
|
||||
});
|
||||
}
|
||||
expect(list).toHaveBeenCalledTimes(2);
|
||||
late.resolve([freshHost]);
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
const failed = await call("sessions.catalog.list", {}, config, client);
|
||||
expect(failed).toHaveBeenCalledWith(true, {
|
||||
catalogs: [
|
||||
expect.objectContaining({
|
||||
hosts: [freshHost],
|
||||
error: expect.objectContaining({
|
||||
code: "catalog_stale",
|
||||
message: expect.stringContaining("TIMEOUT"),
|
||||
}),
|
||||
}),
|
||||
],
|
||||
});
|
||||
list.mockResolvedValueOnce([
|
||||
{ ...freshHost, sessions: [], error: { code: "UNAVAILABLE", message: "node offline" } },
|
||||
]);
|
||||
const offline = await call("sessions.catalog.list", {}, config, client);
|
||||
expect(offline).toHaveBeenCalledWith(true, {
|
||||
catalogs: [
|
||||
expect.objectContaining({
|
||||
hosts: [freshHost],
|
||||
error: expect.objectContaining({
|
||||
code: "catalog_stale",
|
||||
message: expect.stringContaining("UNAVAILABLE"),
|
||||
}),
|
||||
}),
|
||||
],
|
||||
});
|
||||
} finally {
|
||||
late.resolve([freshHost]);
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
}
|
||||
});
|
||||
|
||||
it.each(["caller", "query", "config", "registration", "epoch", "gateway", "archive"] as const)(
|
||||
"does not reuse a cached page after a change to %s",
|
||||
async (change) => {
|
||||
const host = {
|
||||
hostId: "gateway:local",
|
||||
label: "Private page",
|
||||
kind: "gateway" as const,
|
||||
connected: true,
|
||||
sessions: [],
|
||||
};
|
||||
const list = vi
|
||||
.fn<SessionCatalogProvider["list"]>()
|
||||
.mockResolvedValueOnce([host])
|
||||
.mockRejectedValue(Object.assign(new Error("unavailable"), { code: "UNAVAILABLE" }));
|
||||
const fixture = provider("fixture", { list, archive: async () => ({ ok: true }) });
|
||||
hoisted.activeRegistry.sessionCatalogs = [{ provider: fixture }];
|
||||
let config = {};
|
||||
let client = { connId: "first" };
|
||||
let gateway = new AbortController();
|
||||
await call("sessions.catalog.list", {}, config, client, {
|
||||
requestEntryLifetime: { signal: gateway.signal },
|
||||
});
|
||||
if (change === "caller") {
|
||||
client = { connId: "second" };
|
||||
}
|
||||
if (change === "config") {
|
||||
config = {};
|
||||
}
|
||||
if (change === "registration") {
|
||||
hoisted.activeRegistry.sessionCatalogs = [{ provider: fixture }];
|
||||
}
|
||||
if (change === "epoch") {
|
||||
markPluginRegistryActive(getActivePluginRegistry());
|
||||
}
|
||||
if (change === "gateway") {
|
||||
gateway.abort();
|
||||
gateway = new AbortController();
|
||||
}
|
||||
if (change === "archive") {
|
||||
const archived = await call(
|
||||
"sessions.catalog.archive",
|
||||
{
|
||||
catalogId: "fixture",
|
||||
hostId: host.hostId,
|
||||
threadId: "archived",
|
||||
confirmNoOtherRunner: true,
|
||||
},
|
||||
config,
|
||||
client,
|
||||
);
|
||||
expect(archived).toHaveBeenCalledWith(true, { ok: true });
|
||||
}
|
||||
const refreshed = await call(
|
||||
"sessions.catalog.list",
|
||||
change === "query" ? { search: "different" } : {},
|
||||
config,
|
||||
client,
|
||||
{ requestEntryLifetime: { signal: gateway.signal } },
|
||||
);
|
||||
expect(refreshed).toHaveBeenCalledWith(true, {
|
||||
catalogs: [
|
||||
expect.objectContaining({
|
||||
hosts: [],
|
||||
error: expect.objectContaining({ code: "UNAVAILABLE" }),
|
||||
}),
|
||||
],
|
||||
});
|
||||
},
|
||||
);
|
||||
|
||||
it("never projects a cached adoption onto a replacement session identity", async () => {
|
||||
const sessionKey = "agent:main:adopted";
|
||||
setSessionCatalogEntries([
|
||||
{ sessionKey, entry: { sessionId: "original", pluginOwnerId: "fixture" } },
|
||||
]);
|
||||
const list = vi
|
||||
.fn<SessionCatalogProvider["list"]>()
|
||||
.mockResolvedValueOnce([
|
||||
{
|
||||
hostId: "gateway:local",
|
||||
label: "Local",
|
||||
kind: "gateway",
|
||||
connected: true,
|
||||
sessions: [
|
||||
{
|
||||
threadId: "old-thread",
|
||||
sessionKey,
|
||||
status: "stored",
|
||||
archived: false,
|
||||
canContinue: true,
|
||||
canArchive: true,
|
||||
},
|
||||
],
|
||||
},
|
||||
])
|
||||
.mockRejectedValue(new Error("offline"));
|
||||
hoisted.activeRegistry.sessionCatalogs = [{ provider: provider("fixture", { list }) }];
|
||||
const config = {};
|
||||
const original = await call("sessions.catalog.list", {}, config);
|
||||
expect(original.mock.calls[0]?.[1]?.catalogs[0]?.hosts[0]?.sessions[0]?.sessionKey).toBe(
|
||||
sessionKey,
|
||||
);
|
||||
setSessionCatalogEntries([
|
||||
{ sessionKey, entry: { sessionId: "replacement", pluginOwnerId: "fixture" } },
|
||||
]);
|
||||
const stale = await call("sessions.catalog.list", {}, config);
|
||||
expect(stale.mock.calls[0]?.[1]?.catalogs[0]?.hosts[0]?.sessions[0]).toMatchObject({
|
||||
threadId: "old-thread",
|
||||
});
|
||||
expect(stale.mock.calls[0]?.[1]?.catalogs[0]?.hosts[0]?.sessions[0]).not.toHaveProperty(
|
||||
"sessionKey",
|
||||
);
|
||||
});
|
||||
|
||||
it.each([false, true])(
|
||||
"keeps fresh sibling adoption identities separate from stale pages (reused key=%s)",
|
||||
async (reusedKey) => {
|
||||
const keys = ["agent:main:first", reusedKey ? "agent:main:first" : "agent:main:second"];
|
||||
const setEntries = (secondId: string) =>
|
||||
setSessionCatalogEntries(
|
||||
keys.map((sessionKey, index) => ({
|
||||
sessionKey,
|
||||
entry: { sessionId: index === 0 ? "first" : secondId, pluginOwnerId: `fixture-${index}` },
|
||||
})),
|
||||
);
|
||||
const hosts = keys.map((sessionKey, index) => ({
|
||||
hostId: "gateway:local",
|
||||
label: `Provider ${index}`,
|
||||
kind: "gateway" as const,
|
||||
connected: true,
|
||||
sessions: [
|
||||
{
|
||||
threadId: `thread-${index}`,
|
||||
sessionKey,
|
||||
status: "stored",
|
||||
archived: false,
|
||||
canContinue: true,
|
||||
canArchive: true,
|
||||
},
|
||||
],
|
||||
}));
|
||||
const late = createDeferredCore<typeof hosts>();
|
||||
const slow = vi
|
||||
.fn<SessionCatalogProvider["list"]>()
|
||||
.mockResolvedValueOnce([hosts[0]!])
|
||||
.mockImplementationOnce(() => late.promise);
|
||||
hoisted.activeRegistry.sessionCatalogs = [
|
||||
{ provider: provider("fixture-0", { list: slow }) },
|
||||
{ provider: provider("fixture-1", { list: async () => [hosts[1]!] }) },
|
||||
];
|
||||
const config = {};
|
||||
setEntries("original");
|
||||
await call("sessions.catalog.list", {}, config);
|
||||
setEntries("replacement");
|
||||
const next = startCall("sessions.catalog.list", {}, config);
|
||||
try {
|
||||
await vi.advanceTimersByTimeAsync(1_000);
|
||||
await next.completion;
|
||||
expect(next.respond.mock.calls[0]?.[1]?.catalogs[1]?.hosts[0]?.sessions[0]?.sessionKey).toBe(
|
||||
keys[1],
|
||||
);
|
||||
expect(next.respond.mock.calls[0]?.[1]?.catalogs[0]?.hosts[0]?.sessions[0]?.sessionKey).toBe(
|
||||
reusedKey ? undefined : keys[0],
|
||||
);
|
||||
} finally {
|
||||
late.resolve([hosts[0]!]);
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
}
|
||||
},
|
||||
);
|
||||
|
|
@ -214,6 +214,7 @@ export class SessionCatalogListLifetime {
|
|||
this.publishers.delete(releasePublisher);
|
||||
};
|
||||
this.publishers.add(releasePublisher);
|
||||
this.pending += 1;
|
||||
const settle = () => {
|
||||
pending -= 1;
|
||||
this.pending -= 1;
|
||||
|
|
@ -227,33 +228,37 @@ export class SessionCatalogListLifetime {
|
|||
// Completion callbacks can arrive from a different async context; both owners
|
||||
// belong to this listing, and finishListing releases zero-background lists.
|
||||
this.releaseRoot ??= retainGatewayRootWorkAdmissionContinuation() ?? undefined;
|
||||
return await run({
|
||||
signal,
|
||||
onHost: (host) => {
|
||||
if (this.active()) {
|
||||
publish?.(host);
|
||||
}
|
||||
},
|
||||
waitUntil: (completion) => {
|
||||
if (!listing) {
|
||||
throw new Error("Session catalog completion registration is closed");
|
||||
}
|
||||
// Retirement closes delivery, not accounting for work already started.
|
||||
// Join the publication finalizer before the Gateway releases its dependencies.
|
||||
pending += 1;
|
||||
this.pending += 1;
|
||||
void trackWork(() => completion.then(settle, settle));
|
||||
},
|
||||
});
|
||||
return await trackWork(() =>
|
||||
run({
|
||||
signal,
|
||||
onHost: (host) => {
|
||||
if (this.active()) {
|
||||
publish?.(host);
|
||||
}
|
||||
},
|
||||
waitUntil: (completion) => {
|
||||
if (!listing) {
|
||||
throw new Error("Session catalog completion registration is closed");
|
||||
}
|
||||
// Retirement closes delivery, not accounting for work already started.
|
||||
// Join the publication finalizer before the Gateway releases its dependencies.
|
||||
pending += 1;
|
||||
this.pending += 1;
|
||||
void trackWork(() => completion.then(settle, settle));
|
||||
},
|
||||
}),
|
||||
);
|
||||
} catch (error) {
|
||||
releasePublisher();
|
||||
controller.abort(error);
|
||||
throw error;
|
||||
} finally {
|
||||
listing = false;
|
||||
this.pending -= 1;
|
||||
if (pending === 0) {
|
||||
releasePublisher();
|
||||
}
|
||||
this.finish();
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@ import type {
|
|||
SessionsCatalogListParams,
|
||||
} from "../../../packages/gateway-protocol/src/index.js";
|
||||
import type { OpenClawConfig } from "../../config/types.openclaw.js";
|
||||
import { capturePluginRegistryLifecycleEpoch } from "../../plugins/registry-lifecycle.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";
|
||||
|
|
@ -10,7 +11,7 @@ import type { GatewayClient } from "./types.js";
|
|||
|
||||
export type CatalogListEnumeration = {
|
||||
catalogs: SessionCatalog[];
|
||||
instances: SessionCatalogInstances;
|
||||
instancesByCatalog: Map<string, SessionCatalogInstances>;
|
||||
publishedHosts?: Map<string, Map<string, SessionCatalog["hosts"][number]>>;
|
||||
};
|
||||
|
||||
|
|
@ -21,7 +22,11 @@ type CatalogListOperation = {
|
|||
|
||||
type CatalogListOperations = {
|
||||
registrations: CatalogRegistrationSnapshot;
|
||||
epoch: ReturnType<typeof capturePluginRegistryLifecycleEpoch>;
|
||||
gatewaySignal?: AbortSignal;
|
||||
pending: Map<string, CatalogListOperation>;
|
||||
providers: Map<string, CatalogListOperation>;
|
||||
pages: Map<string, CatalogListEnumeration>;
|
||||
retirement: AbortController;
|
||||
};
|
||||
|
||||
|
|
@ -86,11 +91,28 @@ export function resolvePublishedSessionCatalogs(result: CatalogListEnumeration):
|
|||
export function getSessionCatalogListOperations(
|
||||
config: OpenClawConfig,
|
||||
registrations: CatalogRegistrationSnapshot,
|
||||
gatewaySignal?: AbortSignal,
|
||||
): CatalogListOperations {
|
||||
let state = catalogListsByConfig.get(config);
|
||||
if (!state || state.registrations !== registrations) {
|
||||
const epoch = registrations.registry
|
||||
? capturePluginRegistryLifecycleEpoch(registrations.registry)
|
||||
: undefined;
|
||||
if (
|
||||
!state ||
|
||||
state.registrations !== registrations ||
|
||||
state.epoch !== epoch ||
|
||||
state.gatewaySignal !== gatewaySignal
|
||||
) {
|
||||
state?.retirement.abort();
|
||||
state = { registrations, pending: new Map(), retirement: new AbortController() };
|
||||
state = {
|
||||
registrations,
|
||||
epoch,
|
||||
gatewaySignal,
|
||||
pending: new Map(),
|
||||
providers: new Map(),
|
||||
pages: new Map(),
|
||||
retirement: new AbortController(),
|
||||
};
|
||||
catalogListsByConfig.set(config, state);
|
||||
}
|
||||
return state;
|
||||
|
|
@ -105,4 +127,89 @@ export function retireSessionCatalogLists(config: OpenClawConfig): void {
|
|||
operations.retirement.abort();
|
||||
operations.retirement = new AbortController();
|
||||
operations.pending.clear();
|
||||
operations.providers.clear();
|
||||
operations.pages.clear();
|
||||
}
|
||||
|
||||
export async function listSessionCatalogWithinBudget(
|
||||
operations: CatalogListOperations,
|
||||
key: string,
|
||||
progress: SessionCatalogListLifetime,
|
||||
subscribe: (progress: SessionCatalogListLifetime) => void,
|
||||
empty: SessionCatalog,
|
||||
run: () => Promise<CatalogListEnumeration>,
|
||||
): Promise<CatalogListEnumeration> {
|
||||
const retirement = operations.retirement.signal;
|
||||
let active = operations.providers.get(key);
|
||||
if (!active) {
|
||||
const result = run().then((page) => {
|
||||
const catalog = page.catalogs[0]!;
|
||||
if (
|
||||
!retirement.aborted &&
|
||||
!operations.gatewaySignal?.aborted &&
|
||||
!catalog.error &&
|
||||
catalog.hosts.every((host) => !host.error && !host.pending)
|
||||
) {
|
||||
operations.pages.delete(key);
|
||||
operations.pages.set(key, page);
|
||||
if (operations.pages.size > 128) {
|
||||
operations.pages.delete(operations.pages.keys().next().value!);
|
||||
}
|
||||
}
|
||||
return page;
|
||||
});
|
||||
active = { progress, result };
|
||||
operations.providers.set(key, active);
|
||||
const entry = active;
|
||||
void result
|
||||
.finally(() => {
|
||||
if (operations.providers.get(key) === entry) {
|
||||
operations.providers.delete(key);
|
||||
}
|
||||
})
|
||||
.catch(() => undefined);
|
||||
}
|
||||
if (active.progress !== progress) {
|
||||
subscribe(active.progress);
|
||||
}
|
||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||
let result: CatalogListEnumeration | undefined;
|
||||
try {
|
||||
result = await Promise.race([
|
||||
active.result,
|
||||
new Promise<undefined>((resolve) => {
|
||||
timer = setTimeout(() => resolve(undefined), 1_000);
|
||||
timer.unref();
|
||||
}),
|
||||
]);
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
const catalog = result?.catalogs[0];
|
||||
const error = catalog?.error ?? catalog?.hosts.find((host) => host.error)?.error;
|
||||
if (result && !error) {
|
||||
return result;
|
||||
}
|
||||
const cached =
|
||||
retirement.aborted || operations.gatewaySignal?.aborted ? undefined : operations.pages.get(key);
|
||||
if (result && !cached) {
|
||||
return result;
|
||||
}
|
||||
const previous = cached?.catalogs[0] ?? empty;
|
||||
return {
|
||||
catalogs: [
|
||||
{
|
||||
...previous,
|
||||
hosts: result
|
||||
? previous.hosts
|
||||
: previous.hosts.map((host) => Object.assign({}, host, { pending: true })),
|
||||
error: {
|
||||
code: result ? "catalog_stale" : "catalog_pending",
|
||||
message: `${cached ? "Showing stale results. " : ""}${error ? `Refresh failed: [${error.code}] ${error.message}` : "Catalog refresh is still pending; retry shortly."}`,
|
||||
},
|
||||
},
|
||||
],
|
||||
// Preserve the original adoption identity; delivery rechecks current authority.
|
||||
instancesByCatalog: cached?.instancesByCatalog ?? new Map([[empty.id, new Map()]]),
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -27,6 +27,7 @@ import {
|
|||
} from "./session-catalog-list-lifetime.js";
|
||||
import {
|
||||
getSessionCatalogListOperations,
|
||||
listSessionCatalogWithinBudget,
|
||||
resolvePublishedSessionCatalogs,
|
||||
sessionCatalogListKey,
|
||||
type CatalogListEnumeration,
|
||||
|
|
@ -174,7 +175,7 @@ export const listSessionCatalogHandler: GatewayRequestHandlers["sessions.catalog
|
|||
filterSessionCatalogHost(
|
||||
requestEntries.projectHostSessions(
|
||||
host,
|
||||
result.instances,
|
||||
result.instancesByCatalog.get(catalog.id)!,
|
||||
providerAudiences.get(catalog.id),
|
||||
),
|
||||
visibility,
|
||||
|
|
@ -213,7 +214,10 @@ export const listSessionCatalogHandler: GatewayRequestHandlers["sessions.catalog
|
|||
{
|
||||
progressId,
|
||||
agentId: resolvedAgent.agentId,
|
||||
catalog: projectResult({ catalogs: [catalog], instances }).catalogs[0],
|
||||
catalog: projectResult({
|
||||
catalogs: [catalog],
|
||||
instancesByCatalog: new Map([[catalog.id, instances]]),
|
||||
}).catalogs[0],
|
||||
},
|
||||
new Set([progressConnId]),
|
||||
{ dropIfSlow: true },
|
||||
|
|
@ -248,7 +252,11 @@ export const listSessionCatalogHandler: GatewayRequestHandlers["sessions.catalog
|
|||
allowProcessHomeFallback: allowHomeFallback,
|
||||
visibilityKey: resolveSessionCatalogVisibility(client, config).cacheKey,
|
||||
});
|
||||
const operations = getSessionCatalogListOperations(config, catalogRegistrations);
|
||||
const operations = getSessionCatalogListOperations(
|
||||
config,
|
||||
catalogRegistrations,
|
||||
context.requestEntryLifetime?.signal,
|
||||
);
|
||||
const pending = operations.pending.get(listKey);
|
||||
if (pending) {
|
||||
// progressId is connection-owned and excluded from the work key.
|
||||
|
|
@ -313,7 +321,7 @@ export const listSessionCatalogHandler: GatewayRequestHandlers["sessions.catalog
|
|||
} finally {
|
||||
finishPlanning?.();
|
||||
}
|
||||
const instances: SessionCatalogInstances = new Map();
|
||||
const instancesByCatalog = new Map<string, SessionCatalogInstances>();
|
||||
// Partial lists can publish a newer host while another provider or projection still waits.
|
||||
const publishedHosts: CatalogListEnumeration["publishedHosts"] = allowPartialResults
|
||||
? new Map()
|
||||
|
|
@ -327,46 +335,77 @@ export const listSessionCatalogHandler: GatewayRequestHandlers["sessions.catalog
|
|||
const shareRoute = catalogRegistrations.shareRoutes.get(provider);
|
||||
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);
|
||||
// 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,
|
||||
hostIds: request.hostIds,
|
||||
...(request.cursors !== undefined ? { cursors: request.cursors } : {}),
|
||||
sessionEntries: requestEntries?.sessionEntries,
|
||||
listNodes,
|
||||
...lifetime,
|
||||
const page = await listSessionCatalogWithinBudget(
|
||||
operations,
|
||||
JSON.stringify([listKey, provider.id]),
|
||||
progress,
|
||||
subscribe,
|
||||
catalogResult(provider, shareRoute, [], undefined, createTarget),
|
||||
async () => {
|
||||
const providerInstances: SessionCatalogInstances = new Map();
|
||||
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, providerInstances);
|
||||
const catalog = catalogResult(
|
||||
provider,
|
||||
shareRoute,
|
||||
[host],
|
||||
undefined,
|
||||
createTarget,
|
||||
);
|
||||
// The final response also reconciles these snapshots if a slow client drops a frame.
|
||||
progress.publish(catalog, providerInstances);
|
||||
};
|
||||
return listSessionCatalogProvider(provider, providerParams, progress.assertCurrent);
|
||||
});
|
||||
for (const host of hosts) {
|
||||
requestEntries?.captureHostInstances(host, instances);
|
||||
}
|
||||
return catalogResult(provider, shareRoute, hosts, undefined, createTarget);
|
||||
} catch (error) {
|
||||
return catalogResult(provider, shareRoute, [], catalogError(error), createTarget);
|
||||
}
|
||||
try {
|
||||
const hosts = await progress.runProvider(onHost, (lifetime) => {
|
||||
const providerParams = {
|
||||
agentId: resolvedAgent.agentId,
|
||||
allowPartialResults,
|
||||
allowProcessHomeFallback: allowHomeFallback,
|
||||
search,
|
||||
limitPerHost: request.limitPerHost,
|
||||
hostIds: request.hostIds,
|
||||
...(request.cursors !== undefined ? { cursors: request.cursors } : {}),
|
||||
sessionEntries: requestEntries?.sessionEntries,
|
||||
listNodes,
|
||||
...lifetime,
|
||||
};
|
||||
return listSessionCatalogProvider(
|
||||
provider,
|
||||
providerParams,
|
||||
progress.assertCurrent,
|
||||
);
|
||||
});
|
||||
for (const host of hosts) {
|
||||
requestEntries?.captureHostInstances(host, providerInstances);
|
||||
}
|
||||
return {
|
||||
catalogs: [catalogResult(provider, shareRoute, hosts, undefined, createTarget)],
|
||||
instancesByCatalog: new Map([[provider.id, providerInstances]]),
|
||||
publishedHosts,
|
||||
};
|
||||
} catch (error) {
|
||||
return {
|
||||
catalogs: [
|
||||
catalogResult(provider, shareRoute, [], catalogError(error), createTarget),
|
||||
],
|
||||
instancesByCatalog: new Map([[provider.id, providerInstances]]),
|
||||
};
|
||||
}
|
||||
},
|
||||
);
|
||||
instancesByCatalog.set(provider.id, page.instancesByCatalog.get(provider.id)!);
|
||||
return resolvePublishedSessionCatalogs(page)[0]!;
|
||||
}),
|
||||
);
|
||||
} finally {
|
||||
finishProvider?.();
|
||||
}
|
||||
return { catalogs: catalogList, instances, publishedHosts };
|
||||
return { catalogs: catalogList, instancesByCatalog, publishedHosts };
|
||||
})();
|
||||
const entry = { progress, result: operation };
|
||||
// Coalesce only concurrent requests; each subsequent list sees current provider rows.
|
||||
|
|
|
|||
|
|
@ -88,21 +88,22 @@ it("bounds six Gateway connections and filters the shared refresh at each delive
|
|||
elapsed.push(Date.now() - started);
|
||||
}),
|
||||
);
|
||||
await vi.advanceTimersByTimeAsync(5_000);
|
||||
await vi.advanceTimersByTimeAsync(1_000);
|
||||
expect(elapsed).toHaveLength(6);
|
||||
expect(Math.max(...elapsed)).toBeLessThanOrEqual(5_000);
|
||||
expect(Math.max(...elapsed)).toBeLessThanOrEqual(1_000);
|
||||
expect(invokeNode).toHaveBeenCalledTimes(1);
|
||||
for (const call of calls) {
|
||||
expect(call.respond).toHaveBeenCalledWith(true, {
|
||||
catalogs: [
|
||||
expect.objectContaining({
|
||||
hosts: [expect.objectContaining({ pending: true, sessions: [] })],
|
||||
hosts: [],
|
||||
error: expect.objectContaining({ code: "catalog_pending" }),
|
||||
}),
|
||||
],
|
||||
});
|
||||
}
|
||||
clients[5]!.connect.scopes = ["operator.read"];
|
||||
await vi.advanceTimersByTimeAsync(25_000);
|
||||
await vi.advanceTimersByTimeAsync(29_000);
|
||||
await done;
|
||||
for (const [index, broadcast] of broadcasts.entries()) {
|
||||
expect(broadcast).toHaveBeenLastCalledWith(
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue