mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-09 03:22:40 +00:00
fix(mcp): cap channel bridge request limits
This commit is contained in:
@@ -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<string, unknown>;
|
||||
pendingApprovals: Map<string, unknown>;
|
||||
pendingSweepInterval: NodeJS.Timeout | null;
|
||||
pollEvents: (filter: WaitFilter, limit?: number) => {
|
||||
events: QueueEvent[];
|
||||
nextCursor: number;
|
||||
};
|
||||
waitForEvent: (filter: WaitFilter, timeoutMs?: number) => Promise<QueueEvent | null>;
|
||||
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();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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<ConversationDescriptor[]> {
|
||||
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<NonNullable<ChatHistoryResult["messages"]>> {
|
||||
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<QueueEvent | null>((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);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user