From 2f97e8c9ebdcfceb5d7cdee7c39c78dd9ba31f5b Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sun, 9 Aug 2026 18:50:08 -0700 Subject: [PATCH] fix(gateway): deliver scoped global session events (#121166) Centralize global session subscription keys by agent, preserve the default-agent raw global compatibility path, and let authoritative observer audiences reach their selected sockets without changing the public gateway broadcast contract. --- src/gateway/chat-abort.test.ts | 2 + src/gateway/chat-abort.ts | 21 ++- src/gateway/server-broadcast.board.test.ts | 173 +++++++++++++++++- src/gateway/server-broadcast.ts | 4 + src/gateway/server-chat.ts | 68 ++++--- src/gateway/server-methods/chat-broadcast.ts | 43 +++-- src/gateway/server-methods/chat.ts | 5 +- .../server-methods/sessions-subscriptions.ts | 6 +- .../server-methods/sessions-suggestions.ts | 5 +- .../server-methods/sessions-typing.test.ts | 48 ++++- src/gateway/server-session-events.ts | 29 ++- src/gateway/session-message-events.test.ts | 86 +++++++++ src/gateway/session-observer-audience.ts | 25 ++- src/gateway/session-observer-model.ts | 8 +- src/gateway/session-observer.terminal.test.ts | 7 +- src/gateway/session-observer.test.ts | 12 +- src/gateway/session-observer.ts | 31 ++-- src/gateway/session-subscription-keys.ts | 23 +++ 18 files changed, 465 insertions(+), 131 deletions(-) create mode 100644 src/gateway/session-subscription-keys.ts diff --git a/src/gateway/chat-abort.test.ts b/src/gateway/chat-abort.test.ts index 0ac55aad1fff..007f00e2fbb0 100644 --- a/src/gateway/chat-abort.test.ts +++ b/src/gateway/chat-abort.test.ts @@ -606,6 +606,8 @@ describe("abortChatRunById", () => { expect(result).toEqual({ aborted: true }); const payload = firstBroadcastPayload(ops) as ChatAbortPayload; expect(payload.agentId).toBe("main"); + const delivery = { sessionKeys: ["agent:main:global", "global"] }; + expect(ops.broadcast).toHaveBeenCalledWith("chat", payload, delivery); expect(ops.nodeSendToSession).toHaveBeenCalledWith("agent:main:global", "chat", payload); expect(ops.nodeSendToSession).toHaveBeenCalledWith("global", "chat", payload); }); diff --git a/src/gateway/chat-abort.ts b/src/gateway/chat-abort.ts index 513906494062..35d89fbe06f5 100644 --- a/src/gateway/chat-abort.ts +++ b/src/gateway/chat-abort.ts @@ -25,6 +25,10 @@ import { type ChatRunPlanSnapshot, type ChatRunState, } from "./server-chat-state.js"; +import { + resolveSessionSubscriptionKey, + resolveSessionSubscriptionKeys, +} from "./session-subscription-keys.js"; const DEFAULT_CHAT_RUN_ABORT_GRACE_MS = 60_000; @@ -477,20 +481,19 @@ function resolveChatAbortDeliverySessionKeys( sessionKey: string, agentId: string | undefined, ): string[] { - if (sessionKey !== "global") { - return [sessionKey]; - } const scopedAgentId = normalizeActiveAgentId(agentId); if (!scopedAgentId) { return [sessionKey]; } - const keys = [`agent:${scopedAgentId}:global`]; - const cfg = ops.getRuntimeConfig?.(); - const defaultAgentId = cfg ? resolveDefaultAgentId(cfg) : undefined; - if (defaultAgentId && scopedAgentId === defaultAgentId) { - keys.push(sessionKey); + const canonicalKey = resolveSessionSubscriptionKey(sessionKey, scopedAgentId); + if (canonicalKey === sessionKey) { + return [canonicalKey]; } - return keys; + return resolveSessionSubscriptionKeys( + sessionKey, + scopedAgentId, + resolveDefaultGlobalAgentId(ops), + ); } function broadcastChatAborted( diff --git a/src/gateway/server-broadcast.board.test.ts b/src/gateway/server-broadcast.board.test.ts index 31c29b9b72bb..a67c71bd1462 100644 --- a/src/gateway/server-broadcast.board.test.ts +++ b/src/gateway/server-broadcast.board.test.ts @@ -1,8 +1,16 @@ import { describe, expect, it, vi } from "vitest"; -import { GATEWAY_CLIENT_CAPS } from "../../packages/gateway-protocol/src/client-info.js"; +import { + GATEWAY_CLIENT_CAPS, + GATEWAY_CLIENT_IDS, +} from "../../packages/gateway-protocol/src/client-info.js"; import { createGatewayBroadcaster } from "./server-broadcast.js"; -import { createSessionMessageSubscriberRegistry } from "./server-chat-state.js"; +import { + createSessionEventSubscriberRegistry, + createSessionMessageSubscriberRegistry, +} from "./server-chat-state.js"; import type { GatewayWsClient } from "./server/ws-types.js"; +import { createSessionObserverAudience } from "./session-observer-audience.js"; +import { resolveSessionSubscriptionKeys } from "./session-subscription-keys.js"; type RecordingSocket = { bufferedAmount: number; @@ -132,6 +140,167 @@ describe("collaboration event scope guards", () => { expect(unsubscribed.socket.events).toEqual([]); }); + it("delivers prepared global observer audiences exactly once", () => { + const main = makeClient("main", "operator", ["operator.read"]); + const legacy = makeClient("legacy", "operator", ["operator.read"]); + const both = makeClient("both", "operator", ["operator.read"]); + const work = makeClient("work", "operator", ["operator.read"]); + const workRaw = makeClient("work-raw", "operator", ["operator.read"]); + for (const entry of [main, legacy, both, work, workRaw]) { + entry.client.connect.caps = [GATEWAY_CLIENT_CAPS.SESSION_SCOPED_EVENTS]; + } + const subscribers = createSessionMessageSubscriberRegistry(); + subscribers.subscribe(main.client.connId, "agent:main:global"); + subscribers.subscribe(legacy.client.connId, "global"); + subscribers.subscribe(both.client.connId, "agent:main:global"); + subscribers.subscribe(both.client.connId, "global"); + subscribers.subscribe(work.client.connId, "agent:work:global"); + subscribers.subscribe(workRaw.client.connId, "global"); + const audience = createSessionObserverAudience({ + subscribers, + isVisible: () => true, + getDefaultAgentId: () => "main", + }); + const { broadcastToConnIds } = createGatewayBroadcaster({ + clients: new Set([main.client, legacy.client, both.client, work.client, workRaw.client]), + sessionMessageSubscribers: subscribers, + }); + + for (const agentId of ["main", "work"]) { + const sessionKeys = resolveSessionSubscriptionKeys(" GLOBAL ", agentId, "MAIN"); + const recipients = audience.recipients("global", agentId); + broadcastToConnIds("session.observer", { sessionKey: "global", agentId }, recipients, { + sessionKeys, + agentId, + }); + } + + expect(main.socket.events).toEqual(["session.observer"]); + expect(legacy.socket.events).toEqual(["session.observer"]); + expect(both.socket.events).toEqual(["session.observer"]); + expect(work.socket.events).toEqual(["session.observer"]); + expect(workRaw.socket.events).toEqual(["session.observer"]); + }); + + it("preserves event-only recipients selected by the critical observer audience", () => { + const message = makeClient("message", "operator", ["operator.read"]); + const eventOnly = makeClient("event-only", "operator", ["operator.read"]); + const unrelated = makeClient("unrelated", "operator", ["operator.read"]); + for (const entry of [message, eventOnly, unrelated]) { + entry.client.connect.caps = [GATEWAY_CLIENT_CAPS.SESSION_SCOPED_EVENTS]; + } + const subscribers = createSessionMessageSubscriberRegistry(); + const sessionEventSubscribers = createSessionEventSubscriberRegistry(); + subscribers.subscribe(message.client.connId, "agent:main:global"); + sessionEventSubscribers.subscribe(eventOnly.client.connId); + const audience = createSessionObserverAudience({ + subscribers, + sessionEventSubscribers, + isVisible: () => true, + getDefaultAgentId: () => "main", + }); + const { broadcastToConnIds } = createGatewayBroadcaster({ + clients: new Set([message.client, eventOnly.client, unrelated.client]), + sessionMessageSubscribers: subscribers, + }); + + const recipients = audience.criticalRecipients("global", "main"); + broadcastToConnIds( + "session.observer", + { sessionKey: "global", agentId: "main" }, + recipients, + audience.deliveryOptions("global", "main"), + ); + + expect(message.socket.events).toEqual(["session.observer"]); + expect(eventOnly.socket.events).toEqual(["session.observer"]); + expect(unrelated.socket.events).toEqual([]); + }); + + it.each(["agent", "chat", "chat.side_result"])( + "keeps global %s events scoped to their owning agent", + (event) => { + const work = makeClient("work", "operator", ["operator.read"]); + const main = makeClient("main", "operator", ["operator.read"]); + const bareGlobal = makeClient("bare-global", "operator", ["operator.read"]); + for (const entry of [work, main, bareGlobal]) { + entry.client.connect.caps = [GATEWAY_CLIENT_CAPS.SESSION_SCOPED_EVENTS]; + } + const subscribers = createSessionMessageSubscriberRegistry(); + subscribers.subscribe(work.client.connId, "agent:work:global"); + subscribers.subscribe(main.client.connId, "agent:main:global"); + subscribers.subscribe(bareGlobal.client.connId, "global"); + const { broadcast } = createGatewayBroadcaster({ + clients: new Set([work.client, main.client, bareGlobal.client]), + sessionMessageSubscribers: subscribers, + }); + + broadcast( + event, + { sessionKey: "global", agentId: "work" }, + { + sessionKeys: resolveSessionSubscriptionKeys("global", "work", "main"), + agentId: "work", + }, + ); + + expect(work.socket.events).toEqual([event]); + expect(main.socket.events).toEqual([]); + expect(bareGlobal.socket.events).toEqual([]); + }, + ); + + it("subscription-gates Browser Copilot without relying on the capability bit", () => { + const subscribed = makeClient("subscribed", "operator", ["operator.read"]); + const unrelated = makeClient("unrelated", "operator", ["operator.read"]); + for (const entry of [subscribed, unrelated]) { + entry.client.connect.client = { id: GATEWAY_CLIENT_IDS.BROWSER_COPILOT } as never; + entry.client.connect.caps = []; + } + const subscribers = createSessionMessageSubscriberRegistry(); + subscribers.subscribe(subscribed.client.connId, "agent:work:global"); + subscribers.subscribe(unrelated.client.connId, "agent:other:global"); + const { broadcast } = createGatewayBroadcaster({ + clients: new Set([subscribed.client, unrelated.client]), + sessionMessageSubscribers: subscribers, + }); + + broadcast( + "chat", + { sessionKey: "global", agentId: "work" }, + { + sessionKeys: ["agent:work:global"], + agentId: "work", + }, + ); + + expect(subscribed.socket.events).toEqual(["chat"]); + expect(unrelated.socket.events).toEqual([]); + }); + + it.each([ + { sessionKey: "agent:work:global", agentId: "work" }, + { sessionKey: "agent:work:other", agentId: "work" }, + { sessionKey: "global", agentId: undefined }, + ])("preserves exact subscription keys for $sessionKey", ({ sessionKey, agentId }) => { + const subscribed = makeClient("subscribed", "operator", ["operator.read"]); + const unrelated = makeClient("unrelated", "operator", ["operator.read"]); + subscribed.client.connect.caps = [GATEWAY_CLIENT_CAPS.SESSION_SCOPED_EVENTS]; + unrelated.client.connect.caps = [GATEWAY_CLIENT_CAPS.SESSION_SCOPED_EVENTS]; + const subscribers = createSessionMessageSubscriberRegistry(); + subscribers.subscribe(subscribed.client.connId, sessionKey); + subscribers.subscribe(unrelated.client.connId, "agent:other:global"); + const { broadcast } = createGatewayBroadcaster({ + clients: new Set([subscribed.client, unrelated.client]), + sessionMessageSubscribers: subscribers, + }); + + broadcast("chat", { sessionKey, ...(agentId ? { agentId } : {}) }); + + expect(subscribed.socket.events).toEqual(["chat"]); + expect(unrelated.socket.events).toEqual([]); + }); + it("guards suggestion and typing events and forwards payloads to visibility filtering", () => { const pairing = makeClient("pairing", "operator", ["operator.pairing"]); const reader = makeClient("reader", "operator", ["operator.read"]); diff --git a/src/gateway/server-broadcast.ts b/src/gateway/server-broadcast.ts index 59a9d336a675..9679606f64d3 100644 --- a/src/gateway/server-broadcast.ts +++ b/src/gateway/server-broadcast.ts @@ -256,6 +256,9 @@ export function createGatewayBroadcaster(params: { } return frameBase; }; + const sessionSubscriptionVerified = + (opts as { sessionSubscriptionVerified?: boolean } | undefined) + ?.sessionSubscriptionVerified === true; for (const c of params.clients) { if (c.invalidated === true) { continue; @@ -280,6 +283,7 @@ export function createGatewayBroadcaster(params: { SESSION_SUBSCRIPTION_EVENTS.has(event)); if ( requiresSessionSubscription && + !(isTargeted && sessionSubscriptionVerified) && (!sessionKeys.length || !sessionKeys.some((sessionKey) => params.sessionMessageSubscribers?.get(sessionKey).has(c.connId), diff --git a/src/gateway/server-chat.ts b/src/gateway/server-chat.ts index 2053d6c06299..0539ef1028c8 100644 --- a/src/gateway/server-chat.ts +++ b/src/gateway/server-chat.ts @@ -60,6 +60,10 @@ import { isRestartRecoveryLifecycleEvent, isStaleLifecycleEventForSession, } from "./session-lifecycle-state.js"; +import { + resolveSessionSubscriptionKey, + resolveSessionSubscriptionKeys, +} from "./session-subscription-keys.js"; import { loadSessionEntryReadOnly } from "./session-utils.js"; import { formatForLog } from "./ws-log.js"; @@ -631,24 +635,13 @@ export function createAgentEventHandler({ }; }; - const resolveSessionDeliveryKey = (sessionKey: string, agentId?: string) => { - if (sessionKey !== "global") { - return sessionKey; - } - const scopedAgentId = agentId ?? resolveDefaultAgentId(getRuntimeConfig()); - return `agent:${scopedAgentId}:global`; - }; - const resolveNodeSessionDeliveryKeys = (sessionKey: string, agentId?: string) => { - if (sessionKey !== "global") { - return [sessionKey]; + const resolveSessionDeliveryKeys = (sessionKey: string, agentId?: string) => { + const canonicalKey = resolveSessionSubscriptionKey(sessionKey, agentId ?? ""); + if (canonicalKey === sessionKey) { + return [canonicalKey]; } const defaultAgentId = resolveDefaultAgentId(getRuntimeConfig()); - const scopedAgentId = agentId ?? defaultAgentId; - const keys = [`agent:${scopedAgentId}:global`]; - if (scopedAgentId === defaultAgentId) { - keys.push("global"); - } - return keys; + return resolveSessionSubscriptionKeys(sessionKey, agentId ?? defaultAgentId, defaultAgentId); }; const sendNodeSessionPayloadForAgent = ( sessionKey: string, @@ -656,7 +649,7 @@ export function createAgentEventHandler({ payload: unknown, agentId?: string, ) => { - for (const deliverySessionKey of resolveNodeSessionDeliveryKeys(sessionKey, agentId)) { + for (const deliverySessionKey of resolveSessionDeliveryKeys(sessionKey, agentId)) { nodeSendToSession(deliverySessionKey, event, payload); } }; @@ -725,9 +718,9 @@ export function createAgentEventHandler({ isChatAbortMarkerCurrent(chatRunState.runs.get(clientRunId)?.abortMarker, chatLink) || isChatAbortMarkerCurrent(chatRunState.runs.get(evt.runId)?.abortMarker, chatLink); const lifecycleAborted = evt.data?.aborted === true; - const deliverySessionKey = sessionKey - ? resolveSessionDeliveryKey(sessionKey, sessionAgentId) - : undefined; + const deliverySessionKeys = sessionKey + ? resolveSessionDeliveryKeys(sessionKey, sessionAgentId) + : []; const restartRecoveryState = opts?.restartRecoveryState ?? (restartRecoverySessionKey @@ -760,7 +753,9 @@ export function createAgentEventHandler({ !suppressRestartRecoveryProjection && sessionKey && (isControlUiVisible || - (deliverySessionKey ? sessionMessageSubscribers.get(deliverySessionKey).size > 0 : false)) + deliverySessionKeys.some( + (deliverySessionKey) => sessionMessageSubscribers.get(deliverySessionKey).size > 0, + )) ) { if (!isAborted) { const finished = chatLink ? chatRunState.registry.shift(evt.runId) : undefined; @@ -1083,20 +1078,22 @@ export function createAgentEventHandler({ payload: unknown, opts?: { agentId?: string; controlUiVisible?: boolean; dropIfSlow?: boolean }, ) => { - const deliverySessionKey = resolveSessionDeliveryKey(sessionKey, opts?.agentId); + const deliverySessionKeys = resolveSessionDeliveryKeys(sessionKey, opts?.agentId); if (opts?.controlUiVisible ?? true) { broadcast("chat", payload, { dropIfSlow: opts?.dropIfSlow, - sessionKeys: [deliverySessionKey], + sessionKeys: deliverySessionKeys, }); sendNodeSessionPayloadForAgent(sessionKey, "chat", payload, opts?.agentId); return; } - const recipients = sessionMessageSubscribers.get(deliverySessionKey); + const recipients = new Set( + deliverySessionKeys.flatMap((deliveryKey) => [...sessionMessageSubscribers.get(deliveryKey)]), + ); if (recipients.size > 0) { broadcastToConnIds("chat", payload, recipients, { dropIfSlow: opts?.dropIfSlow, - sessionKeys: [deliverySessionKey], + sessionKeys: deliverySessionKeys, }); } }; @@ -1174,9 +1171,7 @@ export function createAgentEventHandler({ ) => { if (opts?.controlUiVisible ?? true) { broadcast("agent", payload, { - sessionKeys: sessionKey - ? [resolveSessionDeliveryKey(sessionKey, opts?.agentId)] - : undefined, + sessionKeys: sessionKey ? resolveSessionDeliveryKeys(sessionKey, opts?.agentId) : undefined, }); if (sessionKey) { sendNodeSessionPayloadForAgent(sessionKey, "agent", payload, opts?.agentId); @@ -1186,12 +1181,14 @@ export function createAgentEventHandler({ if (!sessionKey) { return; } - const deliverySessionKey = resolveSessionDeliveryKey(sessionKey, opts?.agentId); - const recipients = sessionMessageSubscribers.get(deliverySessionKey); + const deliverySessionKeys = resolveSessionDeliveryKeys(sessionKey, opts?.agentId); + const recipients = new Set( + deliverySessionKeys.flatMap((deliveryKey) => [...sessionMessageSubscribers.get(deliveryKey)]), + ); if (recipients.size > 0) { broadcastToConnIds("agent", payload, recipients, { dropIfSlow: opts?.dropIfSlow, - sessionKeys: [deliverySessionKey], + sessionKeys: deliverySessionKeys, }); } }; @@ -1414,8 +1411,9 @@ export function createAgentEventHandler({ ...(isHeartbeat !== undefined && { isHeartbeat }), }; const hasSessionMessageSubscribers = sessionKey - ? sessionMessageSubscribers.get(resolveSessionDeliveryKey(sessionKey, sessionAgentId)).size > - 0 + ? resolveSessionDeliveryKeys(sessionKey, sessionAgentId).some( + (deliverySessionKey) => sessionMessageSubscribers.get(deliverySessionKey).size > 0, + ) : false; const last = agentRunSeq.get(evt.runId) ?? 0; const isToolEvent = evt.stream === "tool"; @@ -1454,7 +1452,7 @@ export function createAgentEventHandler({ }, { sessionKeys: sessionKey - ? [resolveSessionDeliveryKey(sessionKey, sessionAgentId)] + ? resolveSessionDeliveryKeys(sessionKey, sessionAgentId) : undefined, }, ); @@ -1557,7 +1555,7 @@ export function createAgentEventHandler({ runToolRecipients, { sessionKeys: sessionKey - ? [resolveSessionDeliveryKey(sessionKey, sessionAgentId)] + ? resolveSessionDeliveryKeys(sessionKey, sessionAgentId) : undefined, }, ); diff --git a/src/gateway/server-methods/chat-broadcast.ts b/src/gateway/server-methods/chat-broadcast.ts index 7e5bad0540f8..50f0f0c9737f 100644 --- a/src/gateway/server-methods/chat-broadcast.ts +++ b/src/gateway/server-methods/chat-broadcast.ts @@ -1,7 +1,10 @@ import { resolveDefaultAgentId } from "../../agents/agent-scope.js"; import { getReplyPayloadMetadata, type ReplyPayload } from "../../auto-reply/reply-payload.js"; -import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { projectChatDisplayMessage } from "../chat-display-projection.js"; +import { + resolveSessionSubscriptionKey, + resolveSessionSubscriptionKeys, +} from "../session-subscription-keys.js"; import type { GatewayRequestContext } from "./types.js"; type ChatBroadcastContext = Pick< @@ -27,21 +30,21 @@ function nextChatSeq(context: { agentRunSeq: Map }, runId: strin return next; } -function resolveGlobalAwareNodeChatDeliveryKeys(params: { - cfg: OpenClawConfig; +function resolveChatSessionKeys(params: { + context: Partial>; sessionKey: string; agentId?: string; }): string[] { - if (params.sessionKey !== "global") { - return [params.sessionKey]; + const canonicalKey = resolveSessionSubscriptionKey(params.sessionKey, params.agentId ?? ""); + if (canonicalKey === params.sessionKey) { + return [canonicalKey]; } - const defaultAgentId = resolveDefaultAgentId(params.cfg); - const scopedAgentId = params.agentId ?? defaultAgentId; - const keys = [`agent:${scopedAgentId}:global`]; - if (scopedAgentId === defaultAgentId) { - keys.push("global"); - } - return keys; + const defaultAgentId = resolveDefaultAgentId(params.context.getRuntimeConfig?.() ?? {}); + return resolveSessionSubscriptionKeys( + params.sessionKey, + params.agentId ?? defaultAgentId, + defaultAgentId, + ); } export function sendGlobalAwareNodeChatPayload(params: { @@ -52,8 +55,8 @@ export function sendGlobalAwareNodeChatPayload(params: { event: string; payload: unknown; }): void { - const deliveryKeys = resolveGlobalAwareNodeChatDeliveryKeys({ - cfg: params.context.getRuntimeConfig?.() ?? ({} as OpenClawConfig), + const deliveryKeys = resolveChatSessionKeys({ + context: params.context, sessionKey: params.sessionKey, agentId: params.agentId, }); @@ -80,8 +83,8 @@ export function broadcastChatFinal(params: { message: projectChatDisplayMessage(params.message), }; params.context.broadcast("chat", payload, { - sessionKeys: resolveGlobalAwareNodeChatDeliveryKeys({ - cfg: params.context.getRuntimeConfig?.() ?? ({} as OpenClawConfig), + sessionKeys: resolveChatSessionKeys({ + context: params.context, sessionKey: params.sessionKey, agentId: payloadAgentId, }), @@ -121,8 +124,8 @@ export function broadcastSideResult(params: { seq, }; params.context.broadcast("chat.side_result", payload, { - sessionKeys: resolveGlobalAwareNodeChatDeliveryKeys({ - cfg: params.context.getRuntimeConfig?.() ?? ({} as OpenClawConfig), + sessionKeys: resolveChatSessionKeys({ + context: params.context, sessionKey: params.payload.sessionKey, agentId: payloadAgentId, }), @@ -154,8 +157,8 @@ export function broadcastChatError(params: { errorMessage: params.errorMessage, }; params.context.broadcast("chat", payload, { - sessionKeys: resolveGlobalAwareNodeChatDeliveryKeys({ - cfg: params.context.getRuntimeConfig?.() ?? ({} as OpenClawConfig), + sessionKeys: resolveChatSessionKeys({ + context: params.context, sessionKey: params.sessionKey, agentId: payloadAgentId, }), diff --git a/src/gateway/server-methods/chat.ts b/src/gateway/server-methods/chat.ts index 6fcb083c950a..487e92132d86 100644 --- a/src/gateway/server-methods/chat.ts +++ b/src/gateway/server-methods/chat.ts @@ -5,7 +5,7 @@ import { validateChatInjectParams, validateChatToolTitlesParams, } from "../../../packages/gateway-protocol/src/index.js"; -import { resolveSessionAgentId } from "../../agents/agent-scope.js"; +import { resolveDefaultAgentId, resolveSessionAgentId } from "../../agents/agent-scope.js"; import { resolveSessionWorkStartError } from "../../config/sessions.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { beginSessionWorkAdmission } from "../../sessions/session-lifecycle-admission.js"; @@ -14,6 +14,7 @@ import { projectChatDisplayMessage, resolveEffectiveChatHistoryMaxChars, } from "../chat-display-projection.js"; +import { resolveSessionSubscriptionKeys } from "../session-subscription-keys.js"; import { loadSessionEntry, loadSessionEntryReadOnly, @@ -214,7 +215,7 @@ export const chatHandlers: GatewayRequestHandlers = { message, }; context.broadcast("chat", chatPayload, { - sessionKeys: sessionKey === "global" && agentId ? [`agent:${agentId}:global`] : [sessionKey], + sessionKeys: resolveSessionSubscriptionKeys(sessionKey, agentId, resolveDefaultAgentId(cfg)), }); sendGlobalAwareNodeChatPayload({ context, diff --git a/src/gateway/server-methods/sessions-subscriptions.ts b/src/gateway/server-methods/sessions-subscriptions.ts index ce39d1cd51e5..1ab6dff360fe 100644 --- a/src/gateway/server-methods/sessions-subscriptions.ts +++ b/src/gateway/server-methods/sessions-subscriptions.ts @@ -9,8 +9,8 @@ import { import { resolveDefaultAgentId } from "../../agents/agent-scope.js"; import { canReviewOperatorApproval } from "../operator-approval-authorization.js"; import { APPROVALS_SCOPE } from "../operator-scopes.js"; -import { sessionObserverScopeKey } from "../session-observer-model.js"; import { resolveRequestedSessionAgentId as resolveRequestedGlobalAgentId } from "../session-request-agent.js"; +import { resolveSessionSubscriptionKey } from "../session-subscription-keys.js"; import { resolveSessionStoreKey } from "../session-utils.js"; import { requireSessionKey } from "./sessions-shared.js"; import type { GatewayRequestHandlers } from "./types.js"; @@ -109,7 +109,7 @@ export const sessionSubscriptionHandlers: GatewayRequestHandlers = { sessionKey: key, ...(requestedAgentId ? { storeAgentId: requestedAgentId } : {}), }); - const subscriptionKey = sessionObserverScopeKey( + const subscriptionKey = resolveSessionSubscriptionKey( canonicalKey, requestedAgentId ?? resolveDefaultAgentId(cfg), ); @@ -194,7 +194,7 @@ export const sessionSubscriptionHandlers: GatewayRequestHandlers = { sessionKey: key, ...(requestedAgentId ? { storeAgentId: requestedAgentId } : {}), }); - const subscriptionKey = sessionObserverScopeKey( + const subscriptionKey = resolveSessionSubscriptionKey( canonicalKey, requestedAgentId ?? resolveDefaultAgentId(cfg), ); diff --git a/src/gateway/server-methods/sessions-suggestions.ts b/src/gateway/server-methods/sessions-suggestions.ts index 79c9217948b9..e150056d21fb 100644 --- a/src/gateway/server-methods/sessions-suggestions.ts +++ b/src/gateway/server-methods/sessions-suggestions.ts @@ -11,6 +11,7 @@ import { type SessionSharingIdentity, type SessionTypingEvent, } from "../../../packages/gateway-protocol/src/index.js"; +import { resolveDefaultAgentId } from "../../agents/agent-scope.js"; import { addSessionSuggestion, claimSessionSuggestionDispatch, @@ -30,6 +31,7 @@ import { resolveSessionSharingTarget, resolveSessionVisibility, } from "../session-sharing.js"; +import { resolveSessionSubscriptionKeys as subscriptionKeys } from "../session-subscription-keys.js"; import { handleChatSend } from "./chat-send-handler.js"; import { gatewayClientSessionCreator } from "./gateway-client-identity.js"; import { resolveVisibleActiveSessionRunState } from "./session-active-runs.js"; @@ -695,6 +697,7 @@ export const sessionSuggestionHandlers: GatewayRequestHandlers = { if (liveIdentities.size < 2 || !liveIdentities.has(actor.id)) { return false; } + const defaultAgentId = resolveDefaultAgentId(context.getRuntimeConfig()); const event: SessionTypingEvent = { sessionKey: target.canonicalKey, sessionId: current.entry.sessionId, @@ -704,7 +707,7 @@ export const sessionSuggestionHandlers: GatewayRequestHandlers = { ts: Date.now(), }; context.broadcast("session.typing", event, { - sessionKeys: [...sessionKeys].toSorted(), + sessionKeys: subscriptionKeys(current.canonicalKey, current.agentId, defaultAgentId), agentId: target.agentId, dropIfSlow: true, }); diff --git a/src/gateway/server-methods/sessions-typing.test.ts b/src/gateway/server-methods/sessions-typing.test.ts index d9915183e9f6..b3d819e678d0 100644 --- a/src/gateway/server-methods/sessions-typing.test.ts +++ b/src/gateway/server-methods/sessions-typing.test.ts @@ -1,5 +1,6 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { upsertSessionEntry } from "../../config/sessions/session-accessor.js"; +import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { closeOpenClawAgentDatabasesForTest } from "../../state/openclaw-agent-db.js"; import { withOpenClawTestState } from "../../test-utils/openclaw-test-state.js"; import { sessionSuggestionHandlers } from "./sessions-suggestions.js"; @@ -42,9 +43,9 @@ function client(profileId: string, connId: string): GatewayClient { }; } -function context(broadcast = vi.fn()): GatewayRequestContext { +function context(broadcast = vi.fn(), cfg: OpenClawConfig = {}): GatewayRequestContext { return { - getRuntimeConfig: () => ({}), + getRuntimeConfig: () => cfg, broadcast, broadcastToConnIds: vi.fn(), chatAbortControllers: new Map(), @@ -56,6 +57,7 @@ async function callTyping(params: { sessionKey: string; sessionId: string; typing: boolean; + agentId?: string; client: GatewayClient; context: GatewayRequestContext; }) { @@ -63,6 +65,7 @@ async function callTyping(params: { const requestParams = { sessionKey: params.sessionKey, sessionId: params.sessionId, + ...(params.agentId ? { agentId: params.agentId } : {}), typing: params.typing, }; await sessionSuggestionHandlers["session.typing"]?.({ @@ -87,6 +90,47 @@ afterEach(() => { }); describe("session typing handler", () => { + it.each([ + { agentId: "main", expected: ["agent:main:global", "global"] }, + { agentId: "work", expected: ["agent:work:global"] }, + ])("uses the canonical global subscription keys for $agentId", async ({ agentId, expected }) => { + await withOpenClawTestState({ scenario: "minimal" }, async () => { + const cfg = { + agents: { list: [{ id: "main", default: true }, { id: "work" }] }, + } satisfies OpenClawConfig; + await upsertSessionEntry( + { agentId, sessionKey: "global" }, + { + sessionId: `session-${agentId}`, + updatedAt: 1, + createdActor: { type: "human", id: "owner" }, + visibility: "shared", + }, + ); + mocks.presence = [ + { user: { id: "alice" }, watchedSessions: ["global"] }, + { user: { id: "owner" }, watchedSessions: ["global"] }, + ]; + const broadcast = vi.fn(); + + expect( + await callTyping({ + sessionKey: "global", + sessionId: `session-${agentId}`, + agentId, + typing: true, + client: client("alice", `alice-${agentId}`), + context: context(broadcast, cfg), + }), + ).toEqual({ ok: true, broadcast: true }); + expect(broadcast).toHaveBeenCalledWith( + "session.typing", + expect.objectContaining({ agentId, sessionKey: "global" }), + expect.objectContaining({ agentId, sessionKeys: expected }), + ); + }); + }); + it("keeps an identity typing until its last active connection stops", async () => { await withOpenClawTestState({ scenario: "minimal" }, async () => { vi.useFakeTimers(); diff --git a/src/gateway/server-session-events.ts b/src/gateway/server-session-events.ts index 036c40543526..febec53791ff 100644 --- a/src/gateway/server-session-events.ts +++ b/src/gateway/server-session-events.ts @@ -28,6 +28,7 @@ import { buildGatewaySessionEventFields, buildGatewaySessionEventRow, } from "./session-event-payload.js"; +import { resolveSessionSubscriptionKeys } from "./session-subscription-keys.js"; import { attachOpenClawTranscriptMeta, readSessionMessageCountAsync, @@ -91,23 +92,6 @@ function readTranscriptUpdateLifecycleOwner( return lifecycleRevision ? { lifecycleRevision } : {}; } -function resolveSessionMessageBroadcastKeys(sessionKey: string, agentId?: string): string[] { - // Global sessions can be subscribed through either the raw global key or the - // default-agent scoped key; non-default agent global sessions stay scoped. - const normalizedAgentId = normalizeOptionalString(agentId); - if (sessionKey === "global") { - const defaultAgentId = normalizeAgentId(resolveDefaultAgentId(getRuntimeConfig())); - if (normalizedAgentId) { - const scopedKey = `agent:${normalizeAgentId(normalizedAgentId)}:global`; - return normalizeAgentId(normalizedAgentId) === defaultAgentId - ? [scopedKey, sessionKey] - : [scopedKey]; - } - return [`agent:${defaultAgentId}:global`, sessionKey]; - } - return [sessionKey]; -} - function buildGatewaySessionSnapshot(params: { sessionRow: GatewaySessionRow | null | undefined; agentId?: string; @@ -280,7 +264,16 @@ async function handleTranscriptUpdateBroadcast( for (const connId of params.sessionEventSubscribers.getAll()) { connIds.add(connId); } - for (const broadcastKey of resolveSessionMessageBroadcastKeys(sessionKey, routingAgentId)) { + let broadcastKeys = [sessionKey]; + if (sessionKey === "global") { + const defaultAgentId = resolveDefaultAgentId(getRuntimeConfig()); + broadcastKeys = resolveSessionSubscriptionKeys( + sessionKey, + routingAgentId ?? defaultAgentId, + defaultAgentId, + ); + } + for (const broadcastKey of broadcastKeys) { for (const connId of params.sessionMessageSubscribers.get(broadcastKey)) { connIds.add(connId); } diff --git a/src/gateway/session-message-events.test.ts b/src/gateway/session-message-events.test.ts index 2411ab19c27b..38da25b8f169 100644 --- a/src/gateway/session-message-events.test.ts +++ b/src/gateway/session-message-events.test.ts @@ -123,6 +123,21 @@ function waitForSessionMessageEvent( ); } +function waitForSessionObserverEvent( + ws: Awaited>["openWs"]>>, + runId: string, + timeoutMs?: number, +) { + return onceMessage( + ws, + (message) => + message.type === "event" && + message.event === "session.observer" && + (message.payload as { runId?: string } | undefined)?.runId === runId, + timeoutMs, + ); +} + function waitForSessionsChangedMessagePhase( ws: Awaited>["openWs"]>>, sessionKey: string, @@ -1998,6 +2013,77 @@ describe("session.message websocket events", () => { } }); + test("routes a subscribed global observer event through the real gateway socket once", async () => { + const storePath = await createSessionStoreFile(); + testState.agentsConfig = { list: [{ id: "main", default: true }, { id: "work" }] }; + await writeSessionStore({ + entries: { global: { sessionId: "sess-work-observer", updatedAt: Date.now() } }, + storePath, + agentId: "work", + }); + const workWs = await harness.openWs(); + const mainWs = await harness.openWs(); + const runId = "run-work-global-observer"; + const workEvents: unknown[] = []; + const mainEvents: unknown[] = []; + const collect = (target: unknown[]) => (data: RawData) => { + const message = JSON.parse(rawDataToString(data)) as { event?: string; payload?: unknown }; + if (message.event === "session.observer") { + target.push(message.payload); + } + }; + const collectWork = collect(workEvents); + const collectMain = collect(mainEvents); + workWs.on("message", collectWork); + mainWs.on("message", collectMain); + try { + const caps = [GATEWAY_CLIENT_CAPS.SESSION_SCOPED_EVENTS]; + await connectOk(workWs, { scopes: ["operator.read"], caps }); + await connectOk(mainWs, { scopes: ["operator.read"], caps }); + expect( + await rpcReq(workWs, "sessions.messages.subscribe", { + key: " GLOBAL ", + agentId: " WORK ", + }), + ).toMatchObject({ ok: true, payload: { key: "global", subscribed: true } }); + await rpcReq(mainWs, "sessions.messages.subscribe", { key: "global", agentId: "main" }); + await rpcReq(workWs, "sessions.observer.visibility", { visible: true }); + await rpcReq(mainWs, "sessions.observer.visibility", { visible: true }); + + const workEvent = waitForSessionObserverEvent(workWs, runId); + const noMainEvent = expectNoMessageWithin({ + watch: (timeoutMs) => waitForSessionObserverEvent(mainWs, runId, timeoutMs), + timeoutMs: 250, + }); + emitAgentEvent({ + runId, + sessionKey: "global", + agentId: "work", + stream: "item", + data: { + kind: "preamble", + phase: "update", + progressText: "Inspecting the work session", + }, + }); + + await workEvent; + await noMainEvent; + await new Promise((resolve) => { + setTimeout(resolve, 50); + }); + expect(workEvents).toHaveLength(1); + expect(mainEvents).toHaveLength(0); + } finally { + workWs.off("message", collectWork); + mainWs.off("message", collectMain); + workWs.close(); + mainWs.close(); + testState.agentsConfig = undefined; + testState.sessionStorePath = undefined; + } + }); + test("routes unscoped global transcript events to default-agent global subscribers", async () => { const storePath = await createSessionStoreFile(); const transcriptPath = path.join(path.dirname(storePath), "sess-default-global.jsonl"); diff --git a/src/gateway/session-observer-audience.ts b/src/gateway/session-observer-audience.ts index 0e46a723a7cf..e7ceff78dd4b 100644 --- a/src/gateway/session-observer-audience.ts +++ b/src/gateway/session-observer-audience.ts @@ -1,9 +1,8 @@ -import { normalizeAgentId } from "../routing/session-key.js"; import type { SessionEventSubscriberRegistry, SessionMessageSubscriberRegistry, } from "./server-chat-state.js"; -import { sessionObserverScopeKey } from "./session-observer-model.js"; +import { resolveSessionSubscriptionKeys } from "./session-subscription-keys.js"; export function createSessionObserverAudience(params: { subscribers: SessionMessageSubscriberRegistry; @@ -12,18 +11,7 @@ export function createSessionObserverAudience(params: { getDefaultAgentId: () => string; }) { const messageSubscriberKeys = (sessionKey: string, agentId: string): string[] => { - // sessions.messages.subscribe canonicalizes selected-agent global aliases - // to this same qualified key before registering the connection. - const scopedKey = sessionObserverScopeKey(sessionKey, agentId); - if ( - sessionKey === "global" && - normalizeAgentId(agentId) === normalizeAgentId(params.getDefaultAgentId()) - ) { - // Keep legacy default-agent global subscribers while non-default global - // sessions remain confined to their agent-qualified stream. - return [scopedKey, sessionKey]; - } - return [scopedKey]; + return resolveSessionSubscriptionKeys(sessionKey, agentId, params.getDefaultAgentId()); }; const messageRecipients = (sessionKey: string, agentId: string): Set => { @@ -37,6 +25,15 @@ export function createSessionObserverAudience(params: { }; return { + deliveryOptions(sessionKey: string, agentId: string) { + return { + agentId, + dropIfSlow: true, + sessionKeys: messageSubscriberKeys(sessionKey, agentId), + sessionSubscriptionVerified: true, + }; + }, + has(sessionKey: string, agentId: string): boolean { for (const connId of messageRecipients(sessionKey, agentId)) { if (params.isVisible(connId)) { diff --git a/src/gateway/session-observer-model.ts b/src/gateway/session-observer-model.ts index 20eb87d23dea..987f25e4d233 100644 --- a/src/gateway/session-observer-model.ts +++ b/src/gateway/session-observer-model.ts @@ -24,11 +24,11 @@ import type { OpenClawConfig } from "../config/types.openclaw.js"; import type { AgentEventPayload } from "../infra/agent-events.js"; import { pruneMapToMaxSize } from "../infra/map-size.js"; import { redactToolPayloadText } from "../logging/redact.js"; -import { normalizeAgentId } from "../routing/session-key.js"; import type { SessionEventSubscriberRegistry, SessionMessageSubscriberRegistry, } from "./server-chat-state.js"; +import { resolveSessionSubscriptionKey } from "./session-subscription-keys.js"; const HEADLINE_MAX_CHARS = 120; const ASSESSMENT_MAX_CHARS = 320; @@ -38,10 +38,6 @@ const MAX_DORMANT_RUNS = 256; const MAX_DISABLED_RUNS = 512; export const SESSION_OBSERVER_MODEL_MAX_TOKENS = 300; - -export function sessionObserverScopeKey(sessionKey: string, agentId: string): string { - return sessionKey === "global" ? `agent:${normalizeAgentId(agentId)}:global` : sessionKey; -} type PrepareModel = typeof prepareSimpleCompletionModelForAgent; type CompleteModel = typeof completeWithPreparedSimpleCompletionModel; type PreparedModel = Awaited>; @@ -125,7 +121,7 @@ export function rememberSessionObserverDormantRun( // map so a later resume cannot restart below an already broadcast revision. rememberSessionObserverRevisionFloor( floors, - sessionObserverScopeKey(evicted.sessionKey, evicted.agentId), + resolveSessionSubscriptionKey(evicted.sessionKey, evicted.agentId), { revision: evicted.revision, previousDigest: evicted.previousDigest, diff --git a/src/gateway/session-observer.terminal.test.ts b/src/gateway/session-observer.terminal.test.ts index c73b70eea792..f307f0a44530 100644 --- a/src/gateway/session-observer.terminal.test.ts +++ b/src/gateway/session-observer.terminal.test.ts @@ -114,7 +114,12 @@ describe("session observer terminal, persistence, synthesis, and races", () => { "session.observer", expect.objectContaining({ health: "done" }), harness.subscribers.get("agent:main:session-1"), - { dropIfSlow: true }, + expect.objectContaining({ + agentId: "main", + dropIfSlow: true, + sessionKeys: ["agent:main:session-1"], + sessionSubscriptionVerified: true, + }), ); }); diff --git a/src/gateway/session-observer.test.ts b/src/gateway/session-observer.test.ts index db609908e7c1..0e24ce9326ed 100644 --- a/src/gateway/session-observer.test.ts +++ b/src/gateway/session-observer.test.ts @@ -97,7 +97,7 @@ describe("session observer", () => { "session.observer", expect.objectContaining({ headline: "Inspecting mobile rows", health: "on-track" }), new Set(["conn-1"]), - { dropIfSlow: true }, + expect.objectContaining({ dropIfSlow: true }), ); expect(harness.persistDigest).toHaveBeenCalledOnce(); harness.observer.dispose(); @@ -108,9 +108,13 @@ describe("session observer", () => { vi.setSystemTime(1_000); const harness = createHarness({ subscribe: false }); harness.subscribers.subscribe("conn-main", "agent:main:global")?.commit(); + harness.subscribers.subscribe("conn-legacy", "global")?.commit(); harness.subscribers.subscribe("conn-work", "agent:work:global")?.commit(); + harness.subscribers.subscribe("conn-work-raw", "global")?.commit(); declareObserverVisibility(harness.observer, "conn-main"); + declareObserverVisibility(harness.observer, "conn-legacy"); declareObserverVisibility(harness.observer, "conn-work"); + declareObserverVisibility(harness.observer, "conn-work-raw"); harness.observer.handleEvent( event({ @@ -136,14 +140,14 @@ describe("session observer", () => { [ "session.observer", expect.objectContaining({ agentId: "main", revision: 1, sessionKey: "global" }), - new Set(["conn-main"]), - { dropIfSlow: true }, + new Set(["conn-main", "conn-legacy", "conn-work-raw"]), + expect.objectContaining({ sessionKeys: ["agent:main:global", "global"] }), ], [ "session.observer", expect.objectContaining({ agentId: "work", revision: 1, sessionKey: "global" }), new Set(["conn-work"]), - { dropIfSlow: true }, + expect.objectContaining({ sessionKeys: ["agent:work:global"] }), ], ]); harness.observer.dispose(); diff --git a/src/gateway/session-observer.ts b/src/gateway/session-observer.ts index 80eddd08fcc0..4ae8c62d2013 100644 --- a/src/gateway/session-observer.ts +++ b/src/gateway/session-observer.ts @@ -29,7 +29,6 @@ import { rememberSessionObserverDisabledRun, rememberSessionObserverDormantRun, rememberSessionObserverRevisionFloor, - sessionObserverScopeKey, synthesizeSessionObserverTerminalDigest, } from "./session-observer-model.js"; import type { @@ -41,6 +40,7 @@ import type { import { createSessionObserverDigestPersister } from "./session-observer-persistence.js"; import { createSessionObserverPreamblePublisher } from "./session-observer-preamble.js"; import { resolveStoredSessionKeyForAgentStore as resolveStoreKey } from "./session-store-key.js"; +import { resolveSessionSubscriptionKey } from "./session-subscription-keys.js"; const observerLog = createSubsystemLogger("gateway/session-observer"); @@ -76,7 +76,7 @@ export function createSessionObserver(deps: SessionObserverDeps): SessionObserve const cfg = deps.getConfig(); const agentId = resolveSessionAgentId({ sessionKey, config: cfg }); const canonicalSessionKey = resolveStoreKey({ cfg, agentId, sessionKey }); - const state = states.get(sessionObserverScopeKey(canonicalSessionKey, agentId)); + const state = states.get(resolveSessionSubscriptionKey(canonicalSessionKey, agentId)); if (state) { flushSessionActivityAssistantNote(state); return { @@ -105,7 +105,8 @@ export function createSessionObserver(deps: SessionObserverDeps): SessionObserve const runStillCurrent = (runId: string, sessionKey: string, agentId: string) => () => !disposed && !supersededRuns.has(runId) && - (states.get(sessionObserverScopeKey(sessionKey, agentId))?.runId ?? runId) === runId; + (states.get(resolveSessionSubscriptionKey(sessionKey, agentId))?.runId ?? runId) === runId; + const persistAcceptedDigest = createSessionObserverDigestPersister({ now, persistDigest, @@ -132,7 +133,7 @@ export function createSessionObserver(deps: SessionObserverDeps): SessionObserve "session.observer", digest, audience.recipients(state.sessionKey, state.agentId), - { dropIfSlow: true }, + audience.deliveryOptions(state.sessionKey, state.agentId), ); void persistAcceptedDigest(state, digest, false, "preamble"); }, @@ -174,9 +175,7 @@ export function createSessionObserver(deps: SessionObserverDeps): SessionObserve "session.observer", digest, audience.recipients(digest.sessionKey, agentId), - { - dropIfSlow: true, - }, + audience.deliveryOptions(digest.sessionKey, agentId), ); } } catch (error) { @@ -185,7 +184,7 @@ export function createSessionObserver(deps: SessionObserverDeps): SessionObserve } const stateIsTracked = (state: SessionObserverState): boolean => - states.get(sessionObserverScopeKey(state.sessionKey, state.agentId)) === state; + states.get(resolveSessionSubscriptionKey(state.sessionKey, state.agentId)) === state; const dropState = (state: SessionObserverState) => { preamblePublisher.clear(state); @@ -193,9 +192,8 @@ export function createSessionObserver(deps: SessionObserverDeps): SessionObserve clearTimeoutFn(state.timer); } modelSlots.invalidateRequest(state); - const scopeKey = sessionObserverScopeKey(state.sessionKey, state.agentId); if (stateIsTracked(state)) { - states.delete(scopeKey); + states.delete(resolveSessionSubscriptionKey(state.sessionKey, state.agentId)); } }; @@ -417,7 +415,12 @@ export function createSessionObserver(deps: SessionObserverDeps): SessionObserve const recipients = criticalTransition ? audience.criticalRecipients(state.sessionKey, state.agentId) : audience.recipients(state.sessionKey, state.agentId); - deps.broadcastToConnIds("session.observer", digest, recipients, { dropIfSlow: true }); + deps.broadcastToConnIds( + "session.observer", + digest, + recipients, + audience.deliveryOptions(state.sessionKey, state.agentId), + ); await persistAcceptedDigest(state, digest, final); if (final) { dormantRuns.delete(state.runId); @@ -479,7 +482,7 @@ export function createSessionObserver(deps: SessionObserverDeps): SessionObserve if (!agentId || !audience.has(sessionKey, agentId)) { return undefined; } - const scopeKey = sessionObserverScopeKey(sessionKey, agentId); + const scopeKey = resolveSessionSubscriptionKey(sessionKey, agentId); const cfg = deps.getConfig(); if (cfg.gateway?.controlUi?.sessionObserver === false) { return undefined; @@ -592,7 +595,7 @@ export function createSessionObserver(deps: SessionObserverDeps): SessionObserve } return; } - const scopeKey = sessionObserverScopeKey(sessionKey, agentId); + const scopeKey = resolveSessionSubscriptionKey(sessionKey, agentId); if (terminal && audience.recipients(sessionKey, agentId).size === 0) { void synthesizeTerminalDigest({ event, state: states.get(scopeKey) }); dormantRuns.delete(event.runId); @@ -621,7 +624,7 @@ export function createSessionObserver(deps: SessionObserverDeps): SessionObserve const superseded = [...dormantRuns.values()] .filter( (run) => - sessionObserverScopeKey(run.sessionKey, run.agentId) === scopeKey && + resolveSessionSubscriptionKey(run.sessionKey, run.agentId) === scopeKey && run.runId !== event.runId, ) .toSorted( diff --git a/src/gateway/session-subscription-keys.ts b/src/gateway/session-subscription-keys.ts new file mode 100644 index 000000000000..685531dd7c98 --- /dev/null +++ b/src/gateway/session-subscription-keys.ts @@ -0,0 +1,23 @@ +import { normalizeLowercaseStringOrEmpty } from "@openclaw/normalization-core/string-coerce"; +import { normalizeAgentId } from "../routing/session-key.js"; + +export function resolveSessionSubscriptionKey(sessionKey: string, agentId: string): string { + return normalizeLowercaseStringOrEmpty(sessionKey) === "global" + ? `agent:${normalizeAgentId(agentId)}:global` + : sessionKey; +} + +export function resolveSessionSubscriptionKeys( + sessionKey: string, + agentId: string, + defaultAgentId?: string, +): string[] { + const canonicalKey = resolveSessionSubscriptionKey(sessionKey, agentId); + // Raw global is a legacy default-agent stream. Non-default agents must stay + // on their qualified key even when callers supplied the global alias. + return defaultAgentId && + normalizeLowercaseStringOrEmpty(sessionKey) === "global" && + normalizeAgentId(agentId) === normalizeAgentId(defaultAgentId) + ? [canonicalKey, "global"] + : [canonicalKey]; +}