Skip to content

Commit a39e548

Browse files
committed
fix(mcp): cap channel bridge request limits
1 parent 328a446 commit a39e548

3 files changed

Lines changed: 104 additions & 8 deletions

File tree

src/mcp/channel-bridge.test.ts

Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
// Channel MCP bridge tests cover request bridging between MCP and channel APIs.
22
import { afterEach, beforeEach, describe, expect, test, vi } from "vitest";
33
import { OpenClawChannelBridge } from "./channel-bridge.js";
4+
import type { QueueEvent, WaitFilter } from "./channel-shared.js";
45

56
const ONE_MINUTE_MS = 60 * 1_000;
67
const ONE_HOUR_MS = 60 * ONE_MINUTE_MS;
@@ -11,9 +12,15 @@ const APPROVAL_DEFAULT_TTL_MS = 30 * ONE_MINUTE_MS;
1112
// exercise. Defined as a standalone shape (not an intersection with the class)
1213
// because mixing public/private constituents collapses to `never` under tsgo.
1314
type BridgeInternals = {
15+
queue: QueueEvent[];
1416
pendingClaudePermissions: Map<string, unknown>;
1517
pendingApprovals: Map<string, unknown>;
1618
pendingSweepInterval: NodeJS.Timeout | null;
19+
pollEvents: (filter: WaitFilter, limit?: number) => {
20+
events: QueueEvent[];
21+
nextCursor: number;
22+
};
23+
waitForEvent: (filter: WaitFilter, timeoutMs?: number) => Promise<QueueEvent | null>;
1724
handleClaudePermissionRequest: (params: {
1825
requestId: string;
1926
toolName: string;
@@ -262,4 +269,47 @@ describe("OpenClawChannelBridge — pendingClaudePermissions / pendingApprovals
262269
await bridge.close();
263270
}
264271
});
272+
273+
test("pollEvents clamps direct caller limits to the public MCP event window", async () => {
274+
const bridge = makeBridge();
275+
try {
276+
for (let cursor = 1; cursor <= 250; cursor += 1) {
277+
bridge.queue.push({
278+
cursor,
279+
type: "message",
280+
sessionKey: "agent:main:main",
281+
raw: { sessionKey: "agent:main:main" },
282+
});
283+
}
284+
285+
const result = bridge.pollEvents({ afterCursor: 0 }, 10_000);
286+
287+
expect(result.events).toHaveLength(200);
288+
expect(result.nextCursor).toBe(200);
289+
} finally {
290+
await bridge.close();
291+
}
292+
});
293+
294+
test("waitForEvent clamps oversized direct caller timeouts before arming timers", async () => {
295+
const bridge = makeBridge();
296+
try {
297+
let resolved = false;
298+
const waited = bridge.waitForEvent({ afterCursor: 0 }, 3_000_000_000).then((event) => {
299+
resolved = true;
300+
return event;
301+
});
302+
await Promise.resolve();
303+
304+
vi.advanceTimersByTime(299_999);
305+
await Promise.resolve();
306+
expect(resolved).toBe(false);
307+
308+
vi.advanceTimersByTime(1);
309+
await expect(waited).resolves.toBeNull();
310+
expect(resolved).toBe(true);
311+
} finally {
312+
await bridge.close();
313+
}
314+
});
265315
});

src/mcp/channel-bridge.ts

Lines changed: 23 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -49,10 +49,21 @@ type ServerNotification = {
4949

5050
const CLAUDE_PERMISSION_REPLY_RE = /^(yes|no)\s+([a-km-z]{5})$/i;
5151
const QUEUE_LIMIT = 1_000;
52+
const CONVERSATIONS_LIST_LIMIT = 500;
53+
const MESSAGES_READ_LIMIT = 200;
54+
const EVENTS_POLL_LIMIT = 200;
55+
const EVENTS_WAIT_TIMEOUT_LIMIT_MS = 300_000;
5256
const PENDING_CLAUDE_PERMISSION_TTL_MS = 60 * 60 * 1_000;
5357
const PENDING_APPROVAL_DEFAULT_TTL_MS = 30 * 60 * 1_000;
5458
const PENDING_SWEEP_INTERVAL_MS = 5 * 60 * 1_000;
5559

60+
function clampPositiveInteger(value: number | undefined, fallback: number, max: number): number {
61+
if (typeof value !== "number" || !Number.isFinite(value)) {
62+
return fallback;
63+
}
64+
return Math.min(max, Math.max(1, Math.floor(value)));
65+
}
66+
5667
/** Connects the MCP server surface to a Gateway client and queues channel events for polling. */
5768
export class OpenClawChannelBridge {
5869
private gateway: GatewayClient | null = null;
@@ -212,8 +223,9 @@ export class OpenClawChannelBridge {
212223
includeLastMessage?: boolean;
213224
}): Promise<ConversationDescriptor[]> {
214225
await this.waitUntilReady();
226+
const limit = clampPositiveInteger(params?.limit, 50, CONVERSATIONS_LIST_LIMIT);
215227
const response: SessionListResult = await this.requestGateway("sessions.list", {
216-
limit: params?.limit ?? 50,
228+
limit,
217229
search: params?.search,
218230
includeDerivedTitles: params?.includeDerivedTitles ?? true,
219231
includeLastMessage: params?.includeLastMessage ?? true,
@@ -250,9 +262,10 @@ export class OpenClawChannelBridge {
250262
limit = 20,
251263
): Promise<NonNullable<ChatHistoryResult["messages"]>> {
252264
await this.waitUntilReady();
265+
const requestLimit = clampPositiveInteger(limit, 20, MESSAGES_READ_LIMIT);
253266
const response: ChatHistoryResult = await this.requestGateway("sessions.get", {
254267
key: sessionKey,
255-
limit,
268+
limit: requestLimit,
256269
});
257270
return response.messages ?? [];
258271
}
@@ -307,7 +320,10 @@ export class OpenClawChannelBridge {
307320

308321
/** Poll queued events after a cursor without consuming them. */
309322
pollEvents(filter: WaitFilter, limit = 20): { events: QueueEvent[]; nextCursor: number } {
310-
const events = this.queue.filter((event) => matchEventFilter(event, filter)).slice(0, limit);
323+
const eventLimit = clampPositiveInteger(limit, 20, EVENTS_POLL_LIMIT);
324+
const events = this.queue
325+
.filter((event) => matchEventFilter(event, filter))
326+
.slice(0, eventLimit);
311327
const nextCursor = events.at(-1)?.cursor ?? filter.afterCursor;
312328
return { events, nextCursor };
313329
}
@@ -318,6 +334,7 @@ export class OpenClawChannelBridge {
318334
if (existing) {
319335
return existing;
320336
}
337+
const waitTimeoutMs = clampPositiveInteger(timeoutMs, 30_000, EVENTS_WAIT_TIMEOUT_LIMIT_MS);
321338
return await new Promise<QueueEvent | null>((resolve) => {
322339
const waiter: PendingWaiter = {
323340
filter,
@@ -327,11 +344,9 @@ export class OpenClawChannelBridge {
327344
},
328345
timeout: null,
329346
};
330-
if (timeoutMs > 0) {
331-
waiter.timeout = setTimeout(() => {
332-
waiter.resolve(null);
333-
}, timeoutMs);
334-
}
347+
waiter.timeout = setTimeout(() => {
348+
waiter.resolve(null);
349+
}, waitTimeoutMs);
335350
this.pendingWaiters.add(waiter);
336351
});
337352
}

src/mcp/channel-server.test.ts

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -218,6 +218,37 @@ describe("openclaw channel mcp server", () => {
218218
).toBe(true);
219219
});
220220

221+
test("clamps direct bridge session limits to the public MCP windows", async () => {
222+
const sessionKey = "agent:main:main";
223+
const gatewayRequest = vi.fn(async (method: string) => {
224+
if (method === "sessions.list") {
225+
return { sessions: [] };
226+
}
227+
if (method === "sessions.get") {
228+
return { messages: [] };
229+
}
230+
throw new Error(`unexpected gateway method ${method}`);
231+
});
232+
const bridge = new OpenClawChannelBridge({} as never, {
233+
claudeChannelMode: "off",
234+
verbose: false,
235+
});
236+
attachReadyGateway(bridge, gatewayRequest);
237+
238+
await bridge.listConversations({ limit: 10_000 });
239+
await bridge.readMessages(sessionKey, 10_000);
240+
241+
expect(gatewayRequest).toHaveBeenNthCalledWith(
242+
1,
243+
"sessions.list",
244+
expect.objectContaining({ limit: 500 }),
245+
);
246+
expect(gatewayRequest).toHaveBeenNthCalledWith(2, "sessions.get", {
247+
key: sessionKey,
248+
limit: 200,
249+
});
250+
});
251+
221252
test("serializes conversation and message payloads into MCP primary content", async () => {
222253
const mcp = await connectMcpWithoutGateway({ claudeChannelMode: "off" });
223254
try {

0 commit comments

Comments
 (0)