improve(gateway): reduce session roster refresh storms (#159647)

* perf(gateway): reduce session roster refresh storms

Apply held participant snapshots without refetching unfiltered rosters. Reset connection-owned ancestor delivery after successful list responses so bootstrap events cannot leave clients relying on an unadmitted revision. Preserve filtered membership and trailing refreshes for overlapping reads.

* test(gateway): complete direct request ancestor context

Provide the required ancestor-delivery reset in the shared direct request fixture used by authenticated plugin HTTP reads. Preserve all viewer privacy, merged-profile ownership, and readiness revocation assertions. The original five failures now pass with the complete context contract.
This commit is contained in:
Peter Steinberger 2026-09-27 06:57:28 -07:00 • committed by GitHub
parent c991b6bf8f
commit 7cd8ee862c
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
13 changed files with 174 additions and 3 deletions

View file

@ -125,7 +125,7 @@ count.
take precedence when present. Merge an existing
roster member's snapshot locally when the query's membership and pagination
window remain valid. The Control UI reuses lifecycle and ordinary `patch`,
`placement`, `send`, `steer`, `agent.run.started`, `agent.input.settled`, `run-capacity`, and
`participants`, `placement`, `send`, `steer`, `agent.run.started`, `agent.input.settled`, `run-capacity`, and
`chat.title` snapshots for held rows with unchanged identity, archive,
pin, owner, and parent facts and nondecreasing recency. Keyed `sessions.changed`
and `session.message` publications also carry `ancestorSessions`, an array of
@ -156,7 +156,7 @@ count.
happens to match. Missing rows, generation or revision mismatches, and uncertain
presentation ownership require the existing authoritative refresh path.
The Gateway bounds this per-connection record and sends full rows after first
delivery, reconnect, resubscribe, reset/delete, changed presentation or visibility,
delivery, a successful list read, reconnect, resubscribe, reset/delete, changed presentation or visibility,
eviction, or uncertain delivery. Unsubscribe and disconnect clear the record;
session deletion invalidates remembered ancestors.
Older web clients ignore the additive reference field. Because `ancestorSessions`

View file

@ -155,6 +155,7 @@ function createLocalGatewayRequestContext(
unsubscribeSessionEvents: (connId) => {
sessionEvents.delete(connId);
},
forgetConnectionAncestors: () => {},
subscribeSessionMessageEvents: () => undefined,
unsubscribeSessionMessageEvents: () => {},
unsubscribeAllSessionEvents: (connId) => {

View file

@ -134,6 +134,7 @@ export function createDirectChatContext(
broadcast: vi.fn(),
broadcastToConnIds: vi.fn(),
getSessionEventSubscriberConnIds: () => new Set(),
forgetConnectionAncestors: vi.fn<GatewayRequestContext["forgetConnectionAncestors"]>(),
nodeSendToSession: vi.fn(),
registerToolEventRecipient: vi.fn(),
getRuntimeConfig,

View file

@ -323,6 +323,7 @@ export function createGatewayConnectionState(params: {
};
},
clients,
forgetConnectionAncestors,
connectionWork: new GatewayConnectionWork(),
mentionInbox,
isConnectionActive,

View file

@ -121,6 +121,7 @@ export function requestContext(config: OpenClawConfig): GatewayRequestContext {
chatAbortControllers: new Map(),
getRuntimeConfig: () => config,
getSessionEventSubscriberConnIds: () => new Set(),
forgetConnectionAncestors: vi.fn(),
loadGatewayModelCatalog: async () => [],
logGateway: { debug: vi.fn() },
} as unknown as GatewayRequestContext;

View file

@ -276,6 +276,10 @@ export const sessionReadHandlers: GatewayRequestHandlers = {
diagnostics,
onResult: (result) => {
args.sessionMutationAuthorization?.assertCurrent();
// An event delivered before roster admission may not have established its ancestor rows.
if (client?.connId) {
context.forgetConnectionAncestors(client.connId);
}
respond(true, result);
},
});

View file

@ -358,6 +358,7 @@ type GatewayTransportContext = {
terminalSessions?: TerminalSessionManager;
subscribeSessionEvents: (connId: string) => void;
unsubscribeSessionEvents: (connId: string) => void;
forgetConnectionAncestors: (connId: string) => void;
subscribeSessionMessageEvents: (
connId: string,
sessionKey: string,

View file

@ -34,6 +34,7 @@ export function makeContextParams(
return {
runtime: {
getSessionRowProjection: () => undefined,
forgetConnectionAncestors: vi.fn(),
connectionWork: { track: trackAsyncWork },
deps: {} as never,
runtimeState: {

View file

@ -100,6 +100,7 @@ type GatewayRequestContextRuntime = Pick<
Pick<
GatewayCoreRuntime,
| "getSessionRowProjection"
| "forgetConnectionAncestors"
| "refreshGatewayHealthSnapshotWithRuntime"
| "hasTalkNodeConnected"
| "sharedGatewaySessionGenerationState"
@ -538,6 +539,7 @@ export function createGatewayRequestContext(
removeChatRun: runtime.removeChatRun,
subscribeSessionEvents: sessionEventSubscribers.subscribe,
unsubscribeSessionEvents: sessionEventSubscribers.unsubscribe,
forgetConnectionAncestors: runtime.forgetConnectionAncestors,
subscribeSessionMessageEvents: runtime.subscribeSessionMessageEvents,
unsubscribeSessionMessageEvents: runtime.unsubscribeSessionMessageEvents,
unsubscribeAllSessionEvents: (connId) => {

View file

@ -23,7 +23,9 @@ import {
initializeSessionReadContext,
listSessions,
requestContext,
sessionReadHandlers,
} from "./server-methods/sessions-read-cache.test-support.js";
import { sessionSubscriptionHandlers } from "./server-methods/sessions-subscriptions.js";
import type { GatewayWsClient } from "./server/ws-types.js";
import { getSessionRowProjection } from "./session-row-projection-access.js";
import { rolePolicyConfig, sharingPolicyClient } from "./session-sharing.test-utils.js";
@ -37,6 +39,111 @@ type TreeEventPayload = {
afterEach(() => vi.restoreAllMocks());
it.each(["sessions.list", "sessions.subscribe"])(
"%s restores full ancestor delivery after events overlap the roster read",
async (method) => {
await withOpenClawTestState({ scenario: "minimal" }, async () => {
const cfg = { agents: { entries: { main: {} } } };
const root = "agent:main:root";
const child = "agent:main:child";
for (const [key, parentSessionKey] of [
[root, undefined],
[child, root],
] as const) {
replaceSessionEntrySync(
{ agentId: "main", sessionKey: key },
{ sessionId: key, updatedAt: 1, visibility: "shared", parentSessionKey },
);
}
const connection = createGatewayConnectionState({
scheduler: createTestGatewayScheduler(),
bootId: "ancestor-list-recovery",
cfg,
});
const context = requestContext(cfg);
context.subscribeSessionEvents = connection.sessionEventSubscribers.subscribe;
context.forgetConnectionAncestors = connection.forgetConnectionAncestors;
const peers = ["reader", "other"].map((connId) => {
const send = vi.fn();
const client = {
connId,
usesSharedGatewayAuth: false,
connect: {
minProtocol: 1,
maxProtocol: 1,
client: {
id: "openclaw-control-ui",
version: "test",
platform: "test",
mode: "webchat",
},
role: "operator",
scopes: ["operator.admin"],
},
socket: {
readyState: WebSocket.OPEN,
bufferedAmount: 0,
send,
close: vi.fn(),
terminate: vi.fn(),
on: vi.fn(),
off: vi.fn(),
once: vi.fn(),
},
} satisfies GatewayWsClient;
connection.clients.add(client);
connection.sessionEventSubscribers.subscribe(connId);
return { client, send };
});
await initializeSessionReadContext(context);
const projection = getSessionRowProjection(context)!;
const detach = connection.attachSessionRowProjection(projection);
const publish = () =>
connection.broadcast("sessions.changed", {
sessionKey: child,
agentId: "main",
reason: "send",
});
const payloadFor = (peer: (typeof peers)[number]): TreeEventPayload =>
JSON.parse(peer.send.mock.lastCall![0]).payload;
try {
publish();
publish();
expect(payloadFor(peers[0]!).ancestorSessionRefs).toHaveLength(1);
const ensure = projection.ensureMaterialized;
vi.spyOn(projection, "ensureMaterialized").mockImplementationOnce(async () => {
publish();
await ensure();
});
const respond = vi.fn((ok: boolean) => {
expect(ok).toBe(true);
publish();
});
await (method === "sessions.list" ? sessionReadHandlers : sessionSubscriptionHandlers)[
method
]!({
req: { type: "req", id: "ancestor-list", method },
params: { agentId: "main", limit: 20 },
client: peers[0]!.client,
context,
isWebchatConnect: () => false,
respond,
});
expect(respond).toHaveBeenCalledOnce();
expect(payloadFor(peers[0]!).ancestorSessions?.map((row) => row.key)).toEqual([root]);
expect(payloadFor(peers[0]!)).not.toHaveProperty("ancestorSessionRefs");
expect(payloadFor(peers[1]!).ancestorSessionRefs).toHaveLength(1);
publish();
expect(payloadFor(peers[0]!).ancestorSessionRefs).toHaveLength(1);
} finally {
detach();
connection.mentionInbox.dispose();
projection.dispose();
}
});
},
);
it("publishes fresh ancestor rows through private intermediates with list visibility and no duplicates", async () => {
await withOpenClawTestState({ scenario: "minimal" }, async () => {
let now = 1_000_000;

View file

@ -24,7 +24,7 @@ This directory owns Control UI-specific guidance that should not live in the rep
- Session rosters apply nested Gateway row snapshots through the shared reconciler.
`lib/sessions/session-list-query.ts` owns whether a snapshot preserves a held
window: lifecycle, placement, patch/send/steer, run-start/settlement/capacity, and title
window: lifecycle, participants, placement, patch/send/steer, run-start/settlement/capacity, and title
updates can avoid list reads when membership, lineage, and pin/owner/archive
facts stay unchanged and recency does not move backwards. Tree events require
the Gateway's complete, access-scoped `ancestorSessions` snapshots plus any

View file

@ -1,4 +1,5 @@
// @vitest-environment node
import { asOptionalRecord } from "@openclaw/normalization-core/record-coerce";
import { describe, expect, it, vi } from "vitest";
import { createDeferred } from "../../../../test/helpers/promise.js";
import type { GatewaySessionRow, SessionsListResult } from "../../api/types.ts";
@ -81,6 +82,56 @@ function treeHarness(rows = [child, parent, grandparent]) {
}
describe("tree row snapshots", () => {
it("applies participant snapshots locally while refreshing involvement-filtered membership", async () => {
vi.useFakeTimers();
const participants = [{ identity: { type: "profile" as const, id: "viewer" } }];
const updated = { ...settledChild, participants, participantCount: 1 };
let participated = false;
const request = vi.fn(async (_method: string, params?: unknown) =>
sessionsResult(
asOptionalRecord(params)?.involvingMe
? participated
? [updated]
: []
: [child, parent, grandparent],
100,
),
);
const gateway = createGatewayHarness(createTestGatewayClient(request));
const sessions = createTestSessionCapability(gateway.gateway);
const query = { agentId: "main", involvingMe: true };
const stop = sessions.subscribeList(query, () => {});
try {
await sessions.refresh({ agentId: "main", force: true });
await sessions.refreshList({ ...query, force: true });
request.mockClear();
participated = true;
gateway.emitEvent({
type: "event",
event: "sessions.changed",
payload: {
agentId: "main",
reason: "participants",
session: updated,
ancestorSessions: [settledParent, settledGrandparent],
ts: 101,
},
});
expect(sessions.state.result?.sessions[0]).toEqual(updated);
expect(sessions.listSnapshot(query).result?.sessions).toEqual([]);
await vi.advanceTimersByTimeAsync(5_000);
expect(request).toHaveBeenCalledExactlyOnceWith(
"sessions.list",
expect.objectContaining({ involvingMe: true }),
);
expect(sessions.listSnapshot(query).result?.sessions).toEqual([updated]);
} finally {
stop();
sessions.dispose();
vi.useRealTimers();
}
});
it.each(["sessions.changed", "session.message"])(
"keeps %s ancestor references equivalent to full snapshots without roster or descriptor reads",
async (event) => {

View file

@ -30,6 +30,7 @@ import {
const ROW_SNAPSHOT_REASONS = new Set([
"patch",
"participants",
"placement",
"send",
"steer",