diff --git a/packages/core/src/mcp/http.ts b/packages/core/src/mcp/http.ts index 5673cd2..a81bf6d 100644 --- a/packages/core/src/mcp/http.ts +++ b/packages/core/src/mcp/http.ts @@ -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"; @@ -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 { @@ -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 | null = null; let closed = false; const stream = new ReadableStream({ @@ -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`); diff --git a/packages/core/test/mcp.test.ts b/packages/core/test/mcp.test.ts index 4a9e3fe..159c7f2 100644 --- a/packages/core/test/mcp.test.ts +++ b/packages/core/test/mcp.test.ts @@ -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"; @@ -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]!;