Skip to content
Merged

fix #90

Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 24 additions & 4 deletions packages/core/src/mcp/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,9 @@
* Authenticated with per-run bearer tokens (mcp/tokens.ts) or, on `/mcp` only, the key of a connected app
* (connect/connectors.ts) — never with the user's access token.
* `tools/call` answers over SSE when the client accepts it, with keepalives, so long calls
* (agent_delegate waiting for a peer) survive idle timeouts; everything else answers with JSON.
* (agent_delegate waiting for a peer) survive idle timeouts: progress notifications when the call carries a
* progressToken (Claude Code aborts a tool that sends neither a response nor progress for 5 minutes), comments
* otherwise. Everything else answers with JSON.
*/
import type { Context, Hono } from "hono";
import { VERSION } from "../config";
Expand All @@ -23,7 +25,7 @@ import { SSH_INSTRUCTIONS, UnknownSshToolError, callSshTool, listSshTools } from
const log = logger("mcp");

const DEFAULT_PROTOCOL_VERSION = "2025-06-18";
const KEEPALIVE_MS = 15_000;
let keepaliveMs = 15_000;
const SLOW_TOOL_MS = 10_000;

function toolResultText(result: Record<string, unknown>): string {
Expand Down Expand Up @@ -189,9 +191,21 @@ export function disableIdleTimeout(c: Context) {
}
}

/** Answer one tools/call over SSE: headers go out immediately, keepalive comments until the result is ready. */
export function __setKeepaliveForTests(ms: number | null) {
keepaliveMs = ms ?? 15_000;
}

function progressTokenOf(msg: unknown): string | number | null {
const meta = isObj(msg) && isObj(msg.params) && isObj(msg.params._meta) ? msg.params._meta : null;
const token = meta?.progressToken;
return typeof token === "string" || typeof token === "number" ? token : null;
}

/** Answer one tools/call over SSE: headers go out immediately, keepalives until the result is ready. */
function sseCall(ctx: RunContext, msg: unknown, server: McpServerDef): Response {
const encoder = new TextEncoder();
const progressToken = progressTokenOf(msg);
const started = Date.now();
let keepalive: ReturnType<typeof setInterval> | null = null;
let closed = false;
const stream = new ReadableStream<Uint8Array>({
Expand All @@ -205,7 +219,13 @@ function sseCall(ctx: RunContext, msg: unknown, server: McpServerDef): Response
}
};
send(": godmode\n\n");
keepalive = setInterval(() => send(": keepalive\n\n"), KEEPALIVE_MS);
let beats = 0;
keepalive = setInterval(() => {
if (progressToken === null) return send(": keepalive\n\n");
const seconds = Math.round((Date.now() - started) / 1000);
const params = { progressToken, progress: ++beats, message: `Still working (${seconds}s)` };
send(`event: message\ndata: ${JSON.stringify({ jsonrpc: "2.0", method: "notifications/progress", params })}\n\n`);
}, keepaliveMs);
try {
const response = await handleRpc(ctx, msg, server);
if (response) send(`event: message\ndata: ${JSON.stringify(response)}\n\n`);
Expand Down
23 changes: 23 additions & 0 deletions packages/core/test/mcp.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import { createRoutine, listRoutines } from "../src/services/routines";
import { updateSettings } from "../src/services/settings";
import { createWorkspace } from "../src/services/workspaces";
import { createMcpServer } from "../src/integrations/mcpServers";
import { __setKeepaliveForTests } from "../src/mcp/http";

const PASSPHRASE = "correct horse battery staple";
const PASSWORD = "s3cret-Pass-9876";
Expand Down Expand Up @@ -394,6 +395,28 @@ describe("agents + delegation", () => {
expect(child.prompt).toContain("Your final answer goes back to Delegator.]\n\nSay hello");
});

test("a long agent_delegate sends progress for the caller's progressToken, so Claude Code's idle timeout doesn't fire", async () => {
__setKeepaliveForTests(50);
try {
const res = await post(tokens[delegator.id]!, {
jsonrpc: "2.0",
id: 11,
method: "tools/call",
params: { name: "agent_delegate", arguments: { agentId: worker.id, task: "SLOW_STREAM", timeoutSeconds: 60 }, _meta: { progressToken: 11 } },
});
const events = (await res.text()).split("\n").filter((l) => l.startsWith("data: ")).map((l) => JSON.parse(l.slice(6)));
const progress = events.filter((e) => e.method === "notifications/progress");
expect(progress.length).toBeGreaterThan(1);
expect(progress.every((e) => e.params.progressToken === 11)).toBe(true);
expect(progress.map((e) => e.params.progress)).toEqual(progress.map((_, i) => i + 1));
const last = events.at(-1);
expect(last.id).toBe(11);
expect(last.result.content[0].text).toContain("Worker finished the task");
} finally {
__setKeepaliveForTests(null);
}
});

test("agent_delegate without waiting + delegation_status", async () => {
const r = await call(tokens[delegator.id]!, "agent_delegate", { agentId: worker.id, task: "Background job", wait: false });
const runId = /run (run_[A-Za-z0-9]+)/.exec(r.content[0]!.text)![1]!;
Expand Down
Loading