mcp: pump browser session during idle periods

Pumps the browser session while waiting for MCP requests to prevent
curl timeouts. Adds an idle polling loop to the request router.
This commit is contained in:
Adrià Arrufat
2026-06-11 17:48:24 +02:00
parent de7ad43bb9
commit 30c8bcbdd9
3 changed files with 37 additions and 4 deletions

View File

@@ -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 });
};
}

View File

@@ -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());
}

View File

@@ -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,