diff --git a/src/mcp/channel-bridge.test.ts b/src/mcp/channel-bridge.test.ts index e09520d53874..114f10d37935 100644 --- a/src/mcp/channel-bridge.test.ts +++ b/src/mcp/channel-bridge.test.ts @@ -1,6 +1,7 @@ // Channel MCP bridge tests cover request bridging between MCP and channel APIs. import { afterEach, beforeEach, describe, expect, test, vi } from "vitest"; import { OpenClawChannelBridge } from "./channel-bridge.js"; +import type { QueueEvent, WaitFilter } from "./channel-shared.js"; const ONE_MINUTE_MS = 60 * 1_000; const ONE_HOUR_MS = 60 * ONE_MINUTE_MS; @@ -11,9 +12,15 @@ const APPROVAL_DEFAULT_TTL_MS = 30 * ONE_MINUTE_MS; // exercise. Defined as a standalone shape (not an intersection with the class) // because mixing public/private constituents collapses to `never` under tsgo. type BridgeInternals = { + queue: QueueEvent[]; pendingClaudePermissions: Map; pendingApprovals: Map; pendingSweepInterval: NodeJS.Timeout | null; + pollEvents: (filter: WaitFilter, limit?: number) => { + events: QueueEvent[]; + nextCursor: number; + }; + waitForEvent: (filter: WaitFilter, timeoutMs?: number) => Promise; handleClaudePermissionRequest: (params: { requestId: string; toolName: string; @@ -262,4 +269,47 @@ describe("OpenClawChannelBridge — pendingClaudePermissions / pendingApprovals await bridge.close(); } }); + + test("pollEvents clamps direct caller limits to the public MCP event window", async () => { + const bridge = makeBridge(); + try { + for (let cursor = 1; cursor <= 250; cursor += 1) { + bridge.queue.push({ + cursor, + type: "message", + sessionKey: "agent:main:main", + raw: { sessionKey: "agent:main:main" }, + }); + } + + const result = bridge.pollEvents({ afterCursor: 0 }, 10_000); + + expect(result.events).toHaveLength(200); + expect(result.nextCursor).toBe(200); + } finally { + await bridge.close(); + } + }); + + test("waitForEvent clamps oversized direct caller timeouts before arming timers", async () => { + const bridge = makeBridge(); + try { + let resolved = false; + const waited = bridge.waitForEvent({ afterCursor: 0 }, 3_000_000_000).then((event) => { + resolved = true; + return event; + }); + await Promise.resolve(); + + vi.advanceTimersByTime(299_999); + await Promise.resolve(); + expect(resolved).toBe(false); + + vi.advanceTimersByTime(1); + await expect(waited).resolves.toBeNull(); + expect(resolved).toBe(true); + } finally { + await bridge.close(); + } + }); }); diff --git a/src/mcp/channel-bridge.ts b/src/mcp/channel-bridge.ts index c9628a3f32df..2f299f8e527a 100644 --- a/src/mcp/channel-bridge.ts +++ b/src/mcp/channel-bridge.ts @@ -49,10 +49,21 @@ type ServerNotification = { const CLAUDE_PERMISSION_REPLY_RE = /^(yes|no)\s+([a-km-z]{5})$/i; const QUEUE_LIMIT = 1_000; +const CONVERSATIONS_LIST_LIMIT = 500; +const MESSAGES_READ_LIMIT = 200; +const EVENTS_POLL_LIMIT = 200; +const EVENTS_WAIT_TIMEOUT_LIMIT_MS = 300_000; const PENDING_CLAUDE_PERMISSION_TTL_MS = 60 * 60 * 1_000; const PENDING_APPROVAL_DEFAULT_TTL_MS = 30 * 60 * 1_000; const PENDING_SWEEP_INTERVAL_MS = 5 * 60 * 1_000; +function clampPositiveInteger(value: number | undefined, fallback: number, max: number): number { + if (typeof value !== "number" || !Number.isFinite(value)) { + return fallback; + } + return Math.min(max, Math.max(1, Math.floor(value))); +} + /** Connects the MCP server surface to a Gateway client and queues channel events for polling. */ export class OpenClawChannelBridge { private gateway: GatewayClient | null = null; @@ -212,8 +223,9 @@ export class OpenClawChannelBridge { includeLastMessage?: boolean; }): Promise { await this.waitUntilReady(); + const limit = clampPositiveInteger(params?.limit, 50, CONVERSATIONS_LIST_LIMIT); const response: SessionListResult = await this.requestGateway("sessions.list", { - limit: params?.limit ?? 50, + limit, search: params?.search, includeDerivedTitles: params?.includeDerivedTitles ?? true, includeLastMessage: params?.includeLastMessage ?? true, @@ -250,9 +262,10 @@ export class OpenClawChannelBridge { limit = 20, ): Promise> { await this.waitUntilReady(); + const requestLimit = clampPositiveInteger(limit, 20, MESSAGES_READ_LIMIT); const response: ChatHistoryResult = await this.requestGateway("sessions.get", { key: sessionKey, - limit, + limit: requestLimit, }); return response.messages ?? []; } @@ -307,7 +320,10 @@ export class OpenClawChannelBridge { /** Poll queued events after a cursor without consuming them. */ pollEvents(filter: WaitFilter, limit = 20): { events: QueueEvent[]; nextCursor: number } { - const events = this.queue.filter((event) => matchEventFilter(event, filter)).slice(0, limit); + const eventLimit = clampPositiveInteger(limit, 20, EVENTS_POLL_LIMIT); + const events = this.queue + .filter((event) => matchEventFilter(event, filter)) + .slice(0, eventLimit); const nextCursor = events.at(-1)?.cursor ?? filter.afterCursor; return { events, nextCursor }; } @@ -318,6 +334,7 @@ export class OpenClawChannelBridge { if (existing) { return existing; } + const waitTimeoutMs = clampPositiveInteger(timeoutMs, 30_000, EVENTS_WAIT_TIMEOUT_LIMIT_MS); return await new Promise((resolve) => { const waiter: PendingWaiter = { filter, @@ -327,11 +344,9 @@ export class OpenClawChannelBridge { }, timeout: null, }; - if (timeoutMs > 0) { - waiter.timeout = setTimeout(() => { - waiter.resolve(null); - }, timeoutMs); - } + waiter.timeout = setTimeout(() => { + waiter.resolve(null); + }, waitTimeoutMs); this.pendingWaiters.add(waiter); }); } diff --git a/src/mcp/channel-server.test.ts b/src/mcp/channel-server.test.ts index 92dc5c2ab103..a4f0ada8e98b 100644 --- a/src/mcp/channel-server.test.ts +++ b/src/mcp/channel-server.test.ts @@ -218,6 +218,37 @@ describe("openclaw channel mcp server", () => { ).toBe(true); }); + test("clamps direct bridge session limits to the public MCP windows", async () => { + const sessionKey = "agent:main:main"; + const gatewayRequest = vi.fn(async (method: string) => { + if (method === "sessions.list") { + return { sessions: [] }; + } + if (method === "sessions.get") { + return { messages: [] }; + } + throw new Error(`unexpected gateway method ${method}`); + }); + const bridge = new OpenClawChannelBridge({} as never, { + claudeChannelMode: "off", + verbose: false, + }); + attachReadyGateway(bridge, gatewayRequest); + + await bridge.listConversations({ limit: 10_000 }); + await bridge.readMessages(sessionKey, 10_000); + + expect(gatewayRequest).toHaveBeenNthCalledWith( + 1, + "sessions.list", + expect.objectContaining({ limit: 500 }), + ); + expect(gatewayRequest).toHaveBeenNthCalledWith(2, "sessions.get", { + key: sessionKey, + limit: 200, + }); + }); + test("serializes conversation and message payloads into MCP primary content", async () => { const mcp = await connectMcpWithoutGateway({ claudeChannelMode: "off" }); try {