Skip to content

Commit 5a31666

Browse files
cxbAsDevsteipete
andauthored
fix(gateway): catch lazy import rejections in runtime event subscriptions (#100401)
* fix(gateway): catch lazy import rejections in runtime event subscriptions * fix(gateway): consolidate event dispatch failures * docs(changelog): note gateway event dispatch fix --------- Co-authored-by: Peter Steinberger <[email protected]>
1 parent 95f1217 commit 5a31666

4 files changed

Lines changed: 160 additions & 3 deletions

File tree

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ Docs: https://docs.openclaw.ai
1919
### Fixes
2020

2121
- **Source build portability:** keep tsdown configuration self-contained so builds do not depend on resolving the tsdown package from unrun's temporary module directory.
22+
- **Gateway event dispatch:** catch and log lazy subscriber setup and handler failures instead of leaking unhandled promise rejections. (#100401) Thanks @cxbAsDev.
2223
- **Diffs rendering:** render viewer and image output from one SSR preload, preserve language-pack highlighting through hydration, normalize language hints case-insensitively, skip identical before/after inputs with an explicit `changed` result, report truthful file-render and input errors, cache hash-pinned viewer runtimes, and prefer canonical file settings over stale aliases. (#100487)
2324
- **Remote browser reliability:** bound persistent Playwright tab enumeration by the existing remote CDP timeout budget and retire timed-out connection attempts so late completions cannot restore a stuck connection. (#80147, #58968) Thanks @HemantSudarshan and @KeaneYan.
2425
- **Managed browser cookie persistence:** initialize new isolated macOS headless profiles with a non-interactive encryption key while preserving existing profile keys, and close Chromium through CDP before bounded signal fallback so persistent logins survive graceful browser and Gateway restarts. (#96704, #98284) Thanks @TurboTheTurtle.
Lines changed: 120 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,120 @@
1+
// Tests for gateway runtime subscription wiring.
2+
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
3+
import { emitAgentEvent, resetAgentEventsForTest } from "../infra/agent-events.js";
4+
import type { SubsystemLogger } from "../logging/subsystem.js";
5+
import { emitSessionLifecycleEvent } from "../sessions/session-lifecycle-events.js";
6+
import {
7+
emitInternalSessionTranscriptUpdate,
8+
type InternalSessionTranscriptUpdate,
9+
} from "../sessions/transcript-events.js";
10+
import {
11+
createChatRunState,
12+
createSessionEventSubscriberRegistry,
13+
createSessionMessageSubscriberRegistry,
14+
createToolEventRecipientRegistry,
15+
} from "./server-chat-state.js";
16+
17+
const warn = vi.fn();
18+
const mockLog: SubsystemLogger = {
19+
subsystem: "gateway-test",
20+
isEnabled: () => true,
21+
trace: vi.fn(),
22+
debug: vi.fn(),
23+
info: vi.fn(),
24+
warn,
25+
error: vi.fn(),
26+
fatal: vi.fn(),
27+
raw: vi.fn(),
28+
child: () => mockLog,
29+
};
30+
31+
vi.mock("./server-chat.js", () => {
32+
throw new Error("server-chat lazy load failure");
33+
});
34+
35+
vi.mock("./server-session-key.js", () => ({
36+
resolveSessionKeyForRun: () => "agent:main:main",
37+
}));
38+
39+
vi.mock("./server-session-events.js", () => ({
40+
createTranscriptUpdateBroadcastHandler: () => () => {
41+
throw new Error("transcript handler failure");
42+
},
43+
createLifecycleEventBroadcastHandler: () => () => {
44+
throw new Error("lifecycle handler failure");
45+
},
46+
}));
47+
48+
const { startGatewayEventSubscriptions } = await import("./server-runtime-subscriptions.js");
49+
type SubscriptionParams = Parameters<typeof startGatewayEventSubscriptions>[0];
50+
51+
function createParams(): SubscriptionParams {
52+
return {
53+
log: mockLog,
54+
broadcast: vi.fn(),
55+
broadcastToConnIds: vi.fn(),
56+
nodeSendToSession: vi.fn(),
57+
agentRunSeq: new Map(),
58+
chatRunState: createChatRunState(),
59+
toolEventRecipients: createToolEventRecipientRegistry(),
60+
sessionEventSubscribers: createSessionEventSubscriberRegistry(),
61+
sessionMessageSubscribers: createSessionMessageSubscriberRegistry(),
62+
chatAbortControllers: new Map(),
63+
restartRecoveryCandidates: new Map(),
64+
};
65+
}
66+
67+
describe("startGatewayEventSubscriptions", () => {
68+
let unsubs: ReturnType<typeof startGatewayEventSubscriptions> | undefined;
69+
70+
beforeEach(() => {
71+
vi.clearAllMocks();
72+
});
73+
74+
afterEach(() => {
75+
unsubs?.agentUnsub();
76+
unsubs?.heartbeatUnsub();
77+
unsubs?.transcriptUnsub();
78+
unsubs?.lifecycleUnsub();
79+
resetAgentEventsForTest();
80+
});
81+
82+
it("logs lazy agent event module failures", async () => {
83+
unsubs = startGatewayEventSubscriptions(createParams());
84+
85+
emitAgentEvent({ runId: "run-1", stream: "lifecycle", data: { phase: "start" } });
86+
87+
await vi.waitFor(() => expect(warn).toHaveBeenCalledTimes(1));
88+
expect(warn).toHaveBeenCalledWith(
89+
"Agent event dispatch failed",
90+
expect.objectContaining({ runId: "run-1", stream: "lifecycle" }),
91+
);
92+
});
93+
94+
it("logs transcript handler failures", async () => {
95+
unsubs = startGatewayEventSubscriptions(createParams());
96+
97+
emitInternalSessionTranscriptUpdate({
98+
sessionFile: "/tmp/sess.jsonl",
99+
sessionKey: "agent:main:main",
100+
} as InternalSessionTranscriptUpdate);
101+
102+
await vi.waitFor(() => expect(warn).toHaveBeenCalledTimes(1));
103+
expect(warn).toHaveBeenCalledWith(
104+
"Transcript update dispatch failed",
105+
expect.objectContaining({ sessionKey: "agent:main:main" }),
106+
);
107+
});
108+
109+
it("logs lifecycle handler failures", async () => {
110+
unsubs = startGatewayEventSubscriptions(createParams());
111+
112+
emitSessionLifecycleEvent({ sessionKey: "agent:main:main", reason: "created" });
113+
114+
await vi.waitFor(() => expect(warn).toHaveBeenCalledTimes(1));
115+
expect(warn).toHaveBeenCalledWith(
116+
"Lifecycle event dispatch failed",
117+
expect.objectContaining({ sessionKey: "agent:main:main" }),
118+
);
119+
});
120+
});

src/gateway/server-runtime-subscriptions.ts

Lines changed: 38 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
// Gateway event subscription wiring for agent, heartbeat, transcript, and lifecycle broadcasts.
22
import { clearAgentRunContext, onAgentEvent } from "../infra/agent-events.js";
33
import { onHeartbeatEvent } from "../infra/heartbeat-events.js";
4+
import type { SubsystemLogger } from "../logging/subsystem.js";
45
import { onSessionLifecycleEvent } from "../sessions/session-lifecycle-events.js";
56
import { onInternalSessionTranscriptUpdate } from "../sessions/transcript-events.js";
67
import { createLazyPromise } from "../shared/lazy-runtime.js";
@@ -16,8 +17,24 @@ import type {
1617
ToolEventRecipientRegistry,
1718
} from "./server-chat-state.js";
1819

20+
function dispatchEventHandler<TEvent>(params: {
21+
loadHandler: () => Promise<(event: TEvent) => unknown>;
22+
event: TEvent;
23+
log: SubsystemLogger;
24+
failureMessage: string;
25+
context: Record<string, unknown>;
26+
}) {
27+
void params
28+
.loadHandler()
29+
.then((handler) => handler(params.event))
30+
.catch((error: unknown) => {
31+
params.log.warn(params.failureMessage, { ...params.context, error });
32+
});
33+
}
34+
1935
/** Register gateway runtime event subscriptions and return unsubscribe handles. */
2036
export function startGatewayEventSubscriptions(params: {
37+
log: SubsystemLogger;
2138
broadcast: (event: string, payload: unknown, opts?: { dropIfSlow?: boolean }) => void;
2239
broadcastToConnIds: (
2340
event: string,
@@ -234,19 +251,37 @@ export function startGatewayEventSubscriptions(params: {
234251
}
235252
}
236253
}
237-
void getAgentEventHandler().then((handler) => handler(evt));
254+
dispatchEventHandler({
255+
loadHandler: getAgentEventHandler,
256+
event: evt,
257+
log: params.log,
258+
failureMessage: "Agent event dispatch failed",
259+
context: { runId: evt.runId, stream: evt.stream },
260+
});
238261
});
239262

240263
const heartbeatUnsub = onHeartbeatEvent((evt) => {
241264
params.broadcast("heartbeat", evt, { dropIfSlow: true });
242265
});
243266

244267
const transcriptUnsub = onInternalSessionTranscriptUpdate((evt) => {
245-
void getTranscriptUpdateHandler().then((handler) => handler(evt));
268+
dispatchEventHandler({
269+
loadHandler: getTranscriptUpdateHandler,
270+
event: evt,
271+
log: params.log,
272+
failureMessage: "Transcript update dispatch failed",
273+
context: { sessionKey: evt.sessionKey },
274+
});
246275
});
247276

248277
const lifecycleUnsub = onSessionLifecycleEvent((evt) => {
249-
void getLifecycleEventHandler().then((handler) => handler(evt));
278+
dispatchEventHandler({
279+
loadHandler: getLifecycleEventHandler,
280+
event: evt,
281+
log: params.log,
282+
failureMessage: "Lifecycle event dispatch failed",
283+
context: { sessionKey: evt.sessionKey },
284+
});
250285
});
251286

252287
return {

src/gateway/server.impl.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1178,6 +1178,7 @@ export async function startGatewayServer(
11781178
);
11791179
const runtimeSubscriptions = await startupTrace.measure("runtime.subscriptions", () =>
11801180
startGatewayEventSubscriptions({
1181+
log,
11811182
broadcast,
11821183
broadcastToConnIds,
11831184
nodeSendToSession,

0 commit comments

Comments
 (0)