Skip to content

Commit 0e97f96

Browse files
fix(mattermost): add WebSocket ping/pong keepalive (#73979)
Adds Mattermost WebSocket ping/pong liveness checks so half-open sockets terminate and the existing reconnect loop recovers. Fixes #41837. Carries forward #57621. Refs #50138, #44160, and #51104. Thanks @JasonWang1124. Co-authored-by: JasonWang1124 <[email protected]>
1 parent 2d1523e commit 0e97f96

4 files changed

Lines changed: 169 additions & 2 deletions

File tree

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -208,6 +208,7 @@ Docs: https://docs.openclaw.ai
208208
- CLI/status: fall back to a bounded local `status` RPC when loopback detail probes time out or report unknown capability, so reachable local gateways are no longer marked unreachable by slow read diagnostics. Fixes #73535; refs #48360, #62762, #51357, and #42019. Thanks @RacecarGuy, @justinschille, @DJBlackhawk, @tianyaqpzm, and @0xrsydn.
209209
- CLI/gateway: reuse cached paired-device auth during `gateway probe` and report post-connect diagnostic failures as degraded reachability, so healthy local gateways are no longer marked unreachable after loopback auth or read timeouts. Fixes #48360. Thanks @RacecarGuy.
210210
- Channels/Discord: give Discord Gateway WebSocket handshakes a 30s timeout so stalled TLS/network transitions emit an error and Carbon can continue its reconnect loop instead of leaving the bot silent until restart. Refs #50046. Thanks @codexGW.
211+
- Mattermost/WebSocket: send protocol ping/pong keepalives and terminate stale sessions when pongs stop arriving, so silent TCP drops reconnect instead of leaving monitoring idle. Fixes #41837; carries forward #57621; refs #50138, #44160, and #51104. Thanks @JasonWang1124.
211212
- Channels/Telegram: suppress standalone failed edit/write warning payloads when a user-facing assistant error reply already covers the turn, while keeping unresolved mutating failures visible behind success-looking or suppressed-error replies. Fixes #39631; refs #73750; carries forward #39636 and #39717; leaves #39406 for configurable delivery policy. Thanks @Bartok9 and @Bortlesboat.
212213
- Control UI/agents: persist the Set Default action through `agents.list[].default` instead of writing the unsupported `agents.defaultId` field, so saved default-agent changes survive config validation. Fixes #65565; carries forward #72585. Thanks @luyao618.
213214
- NVIDIA/NIM: persist the `NVIDIA_API_KEY` provider marker and mark bundled NVIDIA Chat Completions models as string-content compatible, so NIM models load from `models.json` and OpenAI-compatible subagent calls send plain text content. Fixes #73013 and #50107; refs #73014. Thanks @bautrey, @iot2edge, @ifearghal, and @futhgar.

extensions/mattermost/src/mattermost/monitor-websocket.test.ts

Lines changed: 94 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,18 +8,21 @@ import {
88

99
class FakeWebSocket implements MattermostWebSocketLike {
1010
public readonly sent: string[] = [];
11+
public pingCalls = 0;
1112
public closeCalls = 0;
1213
public terminateCalls = 0;
1314
private openListeners: Array<() => void> = [];
1415
private messageListeners: Array<(data: Buffer) => void | Promise<void>> = [];
16+
private pongListeners: Array<(data: Buffer) => void> = [];
1517
private closeListeners: Array<(code: number, reason: Buffer) => void> = [];
1618
private errorListeners: Array<(err: unknown) => void> = [];
1719

1820
on(event: "open", listener: () => void): void;
1921
on(event: "message", listener: (data: Buffer) => void | Promise<void>): void;
22+
on(event: "pong", listener: (data: Buffer) => void): void;
2023
on(event: "close", listener: (code: number, reason: Buffer) => void): void;
2124
on(event: "error", listener: (err: unknown) => void): void;
22-
on(event: "open" | "message" | "close" | "error", listener: unknown): void {
25+
on(event: "open" | "message" | "pong" | "close" | "error", listener: unknown): void {
2326
if (event === "open") {
2427
this.openListeners.push(listener as () => void);
2528
return;
@@ -28,13 +31,21 @@ class FakeWebSocket implements MattermostWebSocketLike {
2831
this.messageListeners.push(listener as (data: Buffer) => void | Promise<void>);
2932
return;
3033
}
34+
if (event === "pong") {
35+
this.pongListeners.push(listener as (data: Buffer) => void);
36+
return;
37+
}
3138
if (event === "close") {
3239
this.closeListeners.push(listener as (code: number, reason: Buffer) => void);
3340
return;
3441
}
3542
this.errorListeners.push(listener as (err: unknown) => void);
3643
}
3744

45+
ping(): void {
46+
this.pingCalls++;
47+
}
48+
3849
send(data: string): void {
3950
this.sent.push(data);
4051
}
@@ -59,6 +70,12 @@ class FakeWebSocket implements MattermostWebSocketLike {
5970
}
6071
}
6172

73+
emitPong(data = Buffer.alloc(0)): void {
74+
for (const listener of this.pongListeners) {
75+
listener(data);
76+
}
77+
}
78+
6279
emitClose(code: number, reason = ""): void {
6380
const buffer = Buffer.from(reason, "utf8");
6481
for (const listener of this.closeListeners) {
@@ -282,6 +299,82 @@ describe("mattermost websocket monitor", () => {
282299
vi.useRealTimers();
283300
});
284301

302+
it("continues protocol keepalive when Mattermost responds with pong", async () => {
303+
vi.useFakeTimers();
304+
const socket = new FakeWebSocket();
305+
const connectOnce = createMattermostConnectOnce({
306+
wsUrl: "wss://example.invalid/api/v4/websocket",
307+
botToken: "token",
308+
runtime: testRuntime(),
309+
nextSeq: () => 1,
310+
onPosted: async () => {},
311+
webSocketFactory: () => socket,
312+
pingIntervalMs: 100,
313+
pongTimeoutMs: 25,
314+
});
315+
316+
const connected = connectOnce();
317+
socket.emitOpen();
318+
319+
await vi.advanceTimersByTimeAsync(100);
320+
expect(socket.pingCalls).toBe(1);
321+
322+
socket.emitPong();
323+
await vi.advanceTimersByTimeAsync(25);
324+
expect(socket.terminateCalls).toBe(0);
325+
326+
await vi.advanceTimersByTimeAsync(75);
327+
expect(socket.pingCalls).toBe(2);
328+
329+
socket.emitClose(1000);
330+
await connected;
331+
vi.useRealTimers();
332+
});
333+
334+
it("terminates silent websocket drops when Mattermost misses pong timeout", async () => {
335+
vi.useFakeTimers();
336+
const socket = new FakeWebSocket();
337+
const runtime = testRuntime();
338+
let pollCount = 0;
339+
const connectOnce = createMattermostConnectOnce({
340+
wsUrl: "wss://example.invalid/api/v4/websocket",
341+
botToken: "token",
342+
runtime,
343+
nextSeq: () => 1,
344+
onPosted: async () => {},
345+
webSocketFactory: () => socket,
346+
getBotUpdateAt: async () => {
347+
pollCount++;
348+
return 1000;
349+
},
350+
healthCheckIntervalMs: 100,
351+
pingIntervalMs: 50,
352+
pongTimeoutMs: 25,
353+
});
354+
355+
const connected = connectOnce();
356+
socket.emitOpen();
357+
358+
await vi.advanceTimersByTimeAsync(0);
359+
expect(pollCount).toBe(1);
360+
361+
await vi.advanceTimersByTimeAsync(50);
362+
expect(socket.pingCalls).toBe(1);
363+
expect(socket.terminateCalls).toBe(0);
364+
365+
await vi.advanceTimersByTimeAsync(25);
366+
expect(socket.terminateCalls).toBe(1);
367+
expect(runtime.error).toHaveBeenCalledWith("mattermost websocket pong timeout — reconnecting");
368+
369+
await vi.advanceTimersByTimeAsync(500);
370+
expect(socket.pingCalls).toBe(1);
371+
expect(pollCount).toBe(1);
372+
373+
socket.emitClose(1006);
374+
await connected;
375+
vi.useRealTimers();
376+
});
377+
285378
it("does not terminate when getBotUpdateAt throws", async () => {
286379
vi.useFakeTimers();
287380
const socket = new FakeWebSocket();

extensions/mattermost/src/mattermost/monitor-websocket.ts

Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,8 +33,10 @@ export type MattermostEventPayload = {
3333
export type MattermostWebSocketLike = {
3434
on(event: "open", listener: () => void): void;
3535
on(event: "message", listener: (data: WebSocket.RawData) => void | Promise<void>): void;
36+
on(event: "pong", listener: (data: Buffer) => void): void;
3637
on(event: "close", listener: (code: number, reason: Buffer) => void): void;
3738
on(event: "error", listener: (err: unknown) => void): void;
39+
ping(): void;
3840
send(data: string): void;
3941
close(): void;
4042
terminate(): void;
@@ -104,6 +106,8 @@ type CreateMattermostConnectOnceOpts = {
104106
*/
105107
getBotUpdateAt?: () => Promise<number>;
106108
healthCheckIntervalMs?: number;
109+
pingIntervalMs?: number;
110+
pongTimeoutMs?: number;
107111
};
108112

109113
export const defaultMattermostWebSocketFactory: MattermostWebSocketFactory = (url) => {
@@ -144,6 +148,8 @@ export function createMattermostConnectOnce(
144148
): () => Promise<void> {
145149
const webSocketFactory = opts.webSocketFactory ?? defaultMattermostWebSocketFactory;
146150
const healthCheckIntervalMs = opts.healthCheckIntervalMs ?? 30_000;
151+
const pingIntervalMs = opts.pingIntervalMs ?? 30_000;
152+
const pongTimeoutMs = opts.pongTimeoutMs ?? 10_000;
147153
return async () => {
148154
const flowId = randomUUID();
149155
const ws = webSocketFactory(opts.wsUrl);
@@ -158,20 +164,70 @@ export function createMattermostConnectOnce(
158164
let healthCheckEnabled = getBotUpdateAt != null;
159165
let healthCheckInFlight = false;
160166
let healthCheckTimer: ReturnType<typeof setTimeout> | undefined;
167+
let protocolKeepaliveEnabled = true;
168+
let protocolPingTimer: ReturnType<typeof setTimeout> | undefined;
169+
let protocolPongTimer: ReturnType<typeof setTimeout> | undefined;
161170
let initialUpdateAt: number | undefined;
162171

163172
const clearTimers = () => {
164173
if (healthCheckTimer !== undefined) {
165174
clearTimeout(healthCheckTimer);
166175
healthCheckTimer = undefined;
167176
}
177+
if (protocolPingTimer !== undefined) {
178+
clearTimeout(protocolPingTimer);
179+
protocolPingTimer = undefined;
180+
}
181+
if (protocolPongTimer !== undefined) {
182+
clearTimeout(protocolPongTimer);
183+
protocolPongTimer = undefined;
184+
}
168185
};
169186

170187
const stopHealthChecks = () => {
171188
healthCheckEnabled = false;
189+
protocolKeepaliveEnabled = false;
172190
clearTimers();
173191
};
174192

193+
const sendProtocolPing = () => {
194+
if (!protocolKeepaliveEnabled || settled) {
195+
return;
196+
}
197+
if (protocolPongTimer !== undefined) {
198+
clearTimeout(protocolPongTimer);
199+
}
200+
protocolPongTimer = setTimeout(() => {
201+
protocolPongTimer = undefined;
202+
if (!protocolKeepaliveEnabled || settled) {
203+
return;
204+
}
205+
opts.runtime.error?.("mattermost websocket pong timeout — reconnecting");
206+
stopHealthChecks();
207+
ws.terminate();
208+
}, pongTimeoutMs);
209+
try {
210+
ws.ping();
211+
} catch (err) {
212+
if (!protocolKeepaliveEnabled || settled) {
213+
return;
214+
}
215+
opts.runtime.error?.(`mattermost websocket ping failed: ${String(err)}`);
216+
stopHealthChecks();
217+
ws.terminate();
218+
}
219+
};
220+
221+
const scheduleProtocolPing = () => {
222+
if (!protocolKeepaliveEnabled || settled || protocolPingTimer !== undefined) {
223+
return;
224+
}
225+
protocolPingTimer = setTimeout(() => {
226+
protocolPingTimer = undefined;
227+
sendProtocolPing();
228+
}, pingIntervalMs);
229+
};
230+
175231
const scheduleHealthCheck = () => {
176232
if (!getBotUpdateAt || !healthCheckEnabled || settled || healthCheckInFlight) {
177233
return;
@@ -263,6 +319,7 @@ export function createMattermostConnectOnce(
263319
meta: { subsystem: "mattermost-websocket", eventType: "authentication_challenge" },
264320
});
265321
ws.send(authPayload);
322+
scheduleProtocolPing();
266323

267324
// Periodically check if the bot account was modified (e.g. disable/enable).
268325
// After such a cycle the WebSocket silently stops delivering events even
@@ -274,6 +331,14 @@ export function createMattermostConnectOnce(
274331
}
275332
});
276333

334+
ws.on("pong", () => {
335+
if (protocolPongTimer !== undefined) {
336+
clearTimeout(protocolPongTimer);
337+
protocolPongTimer = undefined;
338+
}
339+
scheduleProtocolPing();
340+
});
341+
277342
ws.on("message", async (data) => {
278343
captureWsEvent({
279344
url: opts.wsUrl,

extensions/mattermost/src/mattermost/monitor.inbound-system-event.test.ts

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,14 +5,16 @@ class FakeWebSocket {
55
public readonly sent: string[] = [];
66
private readonly openListeners: Array<() => void> = [];
77
private readonly messageListeners: Array<(data: Buffer) => void | Promise<void>> = [];
8+
private readonly pongListeners: Array<(data: Buffer) => void> = [];
89
private readonly closeListeners: Array<(code: number, reason: Buffer) => void> = [];
910
private readonly errorListeners: Array<(err: unknown) => void> = [];
1011

1112
on(event: "open", listener: () => void): void;
1213
on(event: "message", listener: (data: Buffer) => void | Promise<void>): void;
14+
on(event: "pong", listener: (data: Buffer) => void): void;
1315
on(event: "close", listener: (code: number, reason: Buffer) => void): void;
1416
on(event: "error", listener: (err: unknown) => void): void;
15-
on(event: "open" | "message" | "close" | "error", listener: unknown): void {
17+
on(event: "open" | "message" | "pong" | "close" | "error", listener: unknown): void {
1618
if (event === "open") {
1719
this.openListeners.push(listener as () => void);
1820
return;
@@ -21,6 +23,10 @@ class FakeWebSocket {
2123
this.messageListeners.push(listener as (data: Buffer) => void | Promise<void>);
2224
return;
2325
}
26+
if (event === "pong") {
27+
this.pongListeners.push(listener as (data: Buffer) => void);
28+
return;
29+
}
2430
if (event === "close") {
2531
this.closeListeners.push(listener as (code: number, reason: Buffer) => void);
2632
return;
@@ -32,6 +38,8 @@ class FakeWebSocket {
3238
this.sent.push(data);
3339
}
3440

41+
ping(): void {}
42+
3543
close(): void {}
3644

3745
terminate(): void {}

0 commit comments

Comments
 (0)