diff --git a/src/main.zig b/src/main.zig index 9f689fde7..fbab31ad6 100644 --- a/src/main.zig +++ b/src/main.zig @@ -291,7 +291,7 @@ fn mcpThread(allocator: std.mem.Allocator, app: *App) void { var stdin_buf: [64 * 1024]u8 = undefined; var stdin = std.fs.File.stdin().reader(&stdin_buf); - lp.mcp.router.processRequests(mcp_server, &stdin.interface) catch |err| { + lp.mcp.router.processRequests(mcp_server, &stdin.interface, std.fs.File.stdin()) catch |err| { log.fatal(.mcp, "mcp error", .{ .err = err }); }; } diff --git a/src/mcp/Server.zig b/src/mcp/Server.zig index 4691e4d87..a2b33359c 100644 --- a/src/mcp/Server.zig +++ b/src/mcp/Server.zig @@ -66,6 +66,21 @@ pub fn deinit(self: *Self) void { self.allocator.destroy(self); } +/// Pump the browser session for one short slice while the server waits on +/// input. Page transfers live on this thread's curl multi; left unserviced +/// between tool calls they die on curl's wall-clock timeout. Returns how +/// long the caller may block before pumping again. +pub fn idle(self: *Self) u31 { + const quiet_ms = 250; + self.session.processDestroyQueues(); + var runner = self.session.runner(.{}) catch return quiet_ms; // no page yet + const result = runner.tick(.{ .ms = 25 }) catch return quiet_ms; + return switch (result) { + .done => quiet_ms, + .ok => |next_ms| @intCast(@min(next_ms, quiet_ms)), + }; +} + pub fn sendError(self: *Self, id: std.json.Value, code: protocol.ErrorCode, message: []const u8) !void { return self.transport.sendError(id, code, message); } @@ -119,7 +134,7 @@ test "MCP.Server - Integration: synchronous smoke test" { var server = try Self.init(allocator, app, &out_alloc.writer); defer server.deinit(); - try router.processRequests(server, &in_reader); + try router.processRequests(server, &in_reader, null); try testing.expectJson(.{ .jsonrpc = "2.0", .id = 1, .result = .{ .protocolVersion = "2024-11-05" } }, out_alloc.writer.buffered()); } diff --git a/src/mcp/router.zig b/src/mcp/router.zig index 64b3206ed..e865f657c 100644 --- a/src/mcp/router.zig +++ b/src/mcp/router.zig @@ -10,8 +10,10 @@ const log = lp.log; /// `handleToolList`, `handleToolCall` methods. `handleResourceList` / /// `handleResourceRead` are optional — servers that don't expose /// resources can omit them and the router returns `MethodNotFound` -/// automatically. -pub fn processRequests(server: anytype, reader: *std.io.Reader) !void { +/// automatically. When `input` is the file behind `reader`, the server must +/// also expose `idle`: the router then pumps the browser session between +/// requests instead of blocking on the read. +pub fn processRequests(server: anytype, reader: *std.io.Reader, input: ?std.fs.File) !void { var arena: std.heap.ArenaAllocator = .init(server.allocator); defer arena.deinit(); @@ -19,6 +21,8 @@ pub fn processRequests(server: anytype, reader: *std.io.Reader) !void { _ = arena.reset(.retain_capacity); const aa = arena.allocator(); + if (input) |file| idleUntilInput(server, reader, file); + const buffered_line = reader.takeDelimiter('\n') catch |err| switch (err) { error.StreamTooLong => { log.err(.mcp, "Message too long", .{}); @@ -41,6 +45,20 @@ pub fn processRequests(server: anytype, reader: *std.io.Reader) !void { } } +/// Returns once `reader` holds a delimiter or the fd is readable — +/// `takeDelimiter` may still block on a partial line, but MCP clients write +/// whole lines. A full buffer also returns so StreamTooLong surfaces. +fn idleUntilInput(server: anytype, reader: *std.io.Reader, file: std.fs.File) void { + while (std.mem.indexOfScalar(u8, reader.buffered(), '\n') == null and + reader.bufferedLen() < reader.buffer.len) + { + const wait_ms = server.idle(); + var fds = [_]std.posix.pollfd{.{ .fd = file.handle, .events = std.posix.POLL.IN, .revents = 0 }}; + const ready = std.posix.poll(&fds, wait_ms) catch return; + if (ready > 0) return; + } +} + const Method = enum { initialize, ping,