// Copyright (C) 2023-2026 Lightpanda (Selecy SAS) // // Francis Bouvier // Pierre Tachoire // // This program is free software: you can redistribute it and/or modify // it under the terms of the GNU Affero General Public License as // published by the Free Software Foundation, either version 3 of the // License, or (at your option) any later version. // // This program is distributed in the hope that it will be useful, // but WITHOUT ANY WARRANTY; without even the implied warranty of // MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the // GNU Affero General Public License for more details. // // You should have received a copy of the GNU Affero General Public License // along with this program. If not, see . const std = @import("std"); const lp = @import("lightpanda"); const App = @import("../App.zig"); const Inbox = @import("../Inbox.zig"); const ArenaPool = @import("../ArenaPool.zig"); const sys_net = @import("../sys/net.zig"); const WS = @import("WS.zig"); const CDP = @import("cdp/CDP.zig"); const Driver = @import("Driver.zig"); const posix = std.posix; const Allocator = std.mem.Allocator; const ArenaAllocator = std.heap.ArenaAllocator; // The worker's end of an upgraded connection (the loop's is Server.Worker). // Reads/framing happen on the server run loop (readAvailable → inbox); the worker // thread is the sole writer (send*). The two sides touch disjoint state // (reader+inbox vs send_arena+socket write) so no lock is needed beyond the // inbox's own. const Link = @This(); // A send that hits WouldBlock waits at most this long for the peer to take // more bytes. const SEND_TIMEOUT_MS = 5_000; // The loop reads as fast as it can. The worker _can_ be slow to process its // inbox (e.g. stuck in a syncRequest). Still, we want _some_ limit on how much // data is queued. This is 32 * the configured max message size. Which should // be plenty for a well-behaving client. const INBOX_BACKLOG_MESSAGES = 32; inbox: *Inbox, allocator: Allocator, arena_pool: *ArenaPool, socket: posix.socket_t, protocol: Driver.Protocol, reader: WS.Reader, send_arena: ArenaAllocator, // Nested serialization must not reset an outer message's storage. send_depth: u32, send_timeout_ms: i32, max_inbox_backlog: usize, pub fn init( self: *Link, app: *App, socket: posix.socket_t, protocol: Driver.Protocol, inbox: *Inbox, ) !void { // The Link owns the socket from here on errdefer sys_net.close(socket); if (lp.IS_TEST == false) { const socket_flags = try sys_net.fcntl(socket, posix.F.GETFL, 0); const nonblocking = @as(u32, @bitCast(posix.O{ .NONBLOCK = true })); lp.assert(socket_flags & nonblocking == nonblocking, "Link.init blocking", .{}); } const config = app.config; const allocator = app.allocator; self.* = .{ .inbox = inbox, .socket = socket, .protocol = protocol, .allocator = allocator, .arena_pool = &app.arena_pool, .reader = try .init(allocator, config.cdpMaxMessageSize()), .send_arena = ArenaAllocator.init(allocator), .send_depth = 0, .send_timeout_ms = SEND_TIMEOUT_MS, .max_inbox_backlog = @as(usize, config.cdpMaxMessageSize()) * INBOX_BACKLOG_MESSAGES, }; } pub fn deinit(self: *Link) void { self.reader.deinit(); self.send_arena.deinit(); sys_net.close(self.socket); } pub fn create(app: *App, socket: posix.socket_t, protocol: Driver.Protocol, inbox: *Inbox) !*Link { const link = app.allocator.create(Link) catch |err| { sys_net.close(socket); return err; }; errdefer app.allocator.destroy(link); // init immediately takes ownership of the socket try link.init(app, socket, protocol, inbox); return link; } pub fn destroy(self: *Link) void { const allocator = self.allocator; self.deinit(); allocator.destroy(self); } // Pair every call with sendDone. pub fn acquireSendArena(self: *Link) Allocator { self.send_depth += 1; return self.send_arena.allocator(); } pub fn releaseSendArena(self: *Link) void { self.send_depth -= 1; if (self.send_depth == 0) { _ = self.send_arena.reset(.{ .retain_with_limit = 1024 * 32 }); } } pub fn send(self: *Link, data: []const u8) !void { var pos: usize = 0; const socket = self.socket; while (pos < data.len) { const written = sys_net.write(socket, data[pos..]) catch |err| switch (err) { // The socket is nonblocking so loop reads never stall. Writes are // simpler if they can wait: no per-connection pending-write queue // with its own allocations. Waiting is done with poll rather than // by flipping the fd to blocking: O_NONBLOCK lives on the open // file description, so a flip would reach the loop's reads too. // Should virtually never happen. error.WouldBlock => { // The socket is nonblocking so that the main read loop doesn't // block. But we don't want to make writes truly async, because // then we'd need to allocate the message and hook that back into // the main thread, so...we'll just pull until the we can write // or we hit our send timeout var fds = [_]posix.pollfd{.{ .fd = socket, .events = posix.POLL.OUT, .revents = 0 }}; if ((try posix.poll(&fds, self.send_timeout_ms)) == 0) { return error.Timeout; } continue; }, // a signal landed mid-write; nothing was written error.Interrupted => continue, else => return err, }; if (written == 0) { return error.Closed; } pos += written; } } pub fn sendPong(self: *Link, data: []const u8) !void { if (data.len == 0) { return self.send(&WS.EMPTY_PONG); } var header_buf: [10]u8 = undefined; const header = WS.frameHeader(&header_buf, .pong, data.len); const allocator = self.acquireSendArena(); defer self.releaseSendArena(); const framed = try allocator.alloc(u8, header.len + data.len); @memcpy(framed[0..header.len], header); @memcpy(framed[header.len..], data); return self.send(framed); } // Websocket frames have a variable-length header (2-10 bytes server->client). // We serialize into a buffer whose first 10 bytes are reserved, then // backfill the header right-aligned and send the slice. pub fn sendJSON(self: *Link, message: anytype, opts: std.json.Stringify.Options) !void { const allocator = self.acquireSendArena(); defer self.releaseSendArena(); var aw = try std.Io.Writer.Allocating.initCapacity(allocator, 512); try aw.writer.writeAll(&[_]u8{0} ** 10); try std.json.Stringify.value(message, opts, &aw.writer); const framed = WS.fillHeader(aw.toArrayList()); return self.send(framed); } pub fn sendJSONRaw(self: *Link, buf: std.ArrayList(u8)) !void { // Dangerous API! Assumes the caller reserved the first 10 bytes in buf. const framed = WS.fillHeader(buf); return self.send(framed); } pub const Read = struct { // false once a close frame was consumed: stop reading, the worker // replies and disconnects itself keep: bool, // at least one frame landed in the inbox pushed: bool, }; // Server loop. The socket is readable pub fn readAvailable(self: *Link, budget: usize) !Read { if (self.inbox.queuedBytes() >= self.max_inbox_backlog) { lp.metrics.serve_inbox_backlog.incr(); return error.InboxBacklog; } var pushed = false; var remaining = budget; while (remaining > 0) { const dst = self.reader.readBuf(); if (dst.len == 0) { // a partial message already fills the buffer return error.TooLarge; } const want = dst[0..@min(dst.len, remaining)]; const n = posix.read(self.socket, want) catch |err| switch (err) { error.WouldBlock => break, else => return err, }; if (n == 0) { return error.Closed; } self.reader.len += n; if ((try self.processMessages(&pushed)) == false) { return .{ .keep = false, .pushed = pushed }; } remaining -= n; if (n < want.len) { // a short read: the socket is (very likely) drained break; } } return .{ .keep = true, .pushed = pushed }; } fn processMessages(self: *Link, pushed: *bool) !bool { var reader = &self.reader; while (true) { const msg = (try reader.next()) orelse break; const keep = switch (msg.type) { .pong => true, .ping, .text, .binary => try self.handleMessage(msg, pushed), .close => blk: { _ = try self.handleMessage(msg, pushed); break :blk false; }, }; if (msg.cleanup_fragment) { reader.cleanup(); } if (!keep) { return false; } } reader.compact(); return true; } fn handleMessage(self: *Link, msg: WS.Message, pushed: *bool) !bool { switch (msg.type) { .text, .binary => return switch (self.protocol) { .cdp => self.pushCdp(msg.data, pushed), .bidi => self.pushBiDi(msg.data, pushed), }, .ping => { const arena = try self.arena_pool.acquire(.tiny, "ws ping"); errdefer arena.release(); self.inbox.push(arena, .{ .ping = try arena.dupe(u8, msg.data) }); pushed.* = true; return true; }, .close => { const arena = try self.arena_pool.acquire(.tiny, "ws close"); self.inbox.push(arena, .close); pushed.* = true; return true; }, .pong => unreachable, // processMessages skips pong } } // Parse a CDP JSON frame on the run loop and push it already-parsed: the // consumer's allowlist works on input.method directly and the worker // doesn't re-parse. On parse failure push .disconnect(InvalidJSON) so the // worker tears down, same as a fatal framing error. fn pushCdp(self: *Link, bytes: []const u8, pushed: *bool) !bool { const arena = try self.arena_pool.acquire(bytes.len, "cdp data"); errdefer arena.release(); const raw = try arena.dupe(u8, bytes); const input = std.json.parseFromSliceLeaky( CDP.InputMessage, arena.allocator(), raw, .{ .ignore_unknown_fields = true }, ) catch { self.inbox.push(arena, .{ .disconnect = error.InvalidJSON }); pushed.* = true; return false; }; self.inbox.push(arena, .{ .cdp = .{ .raw = raw, .input = input } }); pushed.* = true; return true; } // BiDi frames are pushed raw; the worker parses them. fn pushBiDi(self: *Link, bytes: []const u8, pushed: *bool) !bool { const arena = try self.arena_pool.acquire(bytes.len, "bidi data"); errdefer arena.release(); self.inbox.push(arena, .{ .bidi = try arena.dupe(u8, bytes) }); pushed.* = true; return true; } // Server loop, closing only the read side. The worker can still send a message // (e.g. a close frame). pub fn shutdown(self: *Link) void { sys_net.shutdown(self.socket, .recv) catch {}; } const testing = @import("../testing.zig"); test "link: send gives up when the peer stops reading" { var pair: [2]posix.socket_t = undefined; if (std.c.socketpair(posix.AF.LOCAL, posix.SOCK.STREAM, 0, &pair) != 0) { return error.SocketPairFailed; } // pair[1] is the link's, closed by its deinit defer sys_net.close(pair[0]); const small = std.mem.toBytes(@as(c_int, 4096)); try posix.setsockopt(pair[0], posix.SOL.SOCKET, posix.SO.RCVBUF, &small); try posix.setsockopt(pair[1], posix.SOL.SOCKET, posix.SO.SNDBUF, &small); const nonblocking = @as(u32, @bitCast(posix.O{ .NONBLOCK = true })); const flags = try sys_net.fcntl(pair[1], posix.F.GETFL, 0); _ = try sys_net.fcntl(pair[1], posix.F.SETFL, flags | nonblocking); var inbox: Inbox = .{}; defer inbox.deinit(); var link: Link = undefined; try link.init(testing.test_app, pair[1], .cdp, &inbox); defer link.deinit(); // shorten the wait so the test doesn't sit out the real one link.send_timeout_ms = 50; const payload = try testing.allocator.alloc(u8, 1024 * 1024); defer testing.allocator.free(payload); @memset(payload, 'a'); // Nobody drains pair[0]. send() waits for writability to finish the // write; unbounded, that parks the worker forever and, because shutdown // only half-closes the read side, hangs the whole process on SIGINT. try testing.expectError(error.Timeout, link.send(payload)); // and the run loop's reads share the fd: it must still be non-blocking. // Bit test, not equality: macOS adds an internal bit to F_GETFL after a write. const after = try sys_net.fcntl(pair[1], posix.F.GETFL, 0); try testing.expect(after & nonblocking != 0); } test "link: nested serialization preserves complete frames and recovers from errors" { const Nested = struct { link: *Link, fail: bool, pub fn jsonStringify(self: @This(), w: anytype) error{WriteFailed}!void { try w.beginObject(); try w.objectField("head"); try w.write("outer head"); self.link.sendJSON(.{ .nested = "x" ** 64 }, .{}) catch return error.WriteFailed; self.link.sendPong("ping") catch return error.WriteFailed; self.link.sendJSON(.{ .nested = "y" ** 64 }, .{}) catch return error.WriteFailed; if (self.fail) return error.WriteFailed; try w.objectField("tail"); try w.write("outer tail" ** 128); try w.endObject(); } }; var ctx = try @import("cdp/testing.zig").context(); defer ctx.deinit(); const link = &ctx.cdp().link; try link.sendJSON(Nested{ .link = link, .fail = false }, .{}); try ctx.expectSent(.{ .nested = "x" ** 64 }, .{ .index = 0 }); try ctx.expectSent(.{ .nested = "y" ** 64 }, .{ .index = 1 }); try ctx.expectSent(.{ .head = "outer head", .tail = "outer tail" ** 128 }, .{ .index = 2 }); try testing.expectEqual(0, link.send_depth); try testing.expectError(error.WriteFailed, link.sendJSON(Nested{ .link = link, .fail = true }, .{})); try testing.expectEqual(0, link.send_depth); try link.sendJSON(.{ .recovered = true }, .{}); try testing.expectJson(.{ .nested = "x" ** 64 }, (try ctx.getSentMessage(3)).?); try testing.expectJson(.{ .nested = "y" ** 64 }, (try ctx.getSentMessage(4)).?); try ctx.expectSent(.{ .recovered = true }, .{ .index = 5 }); try testing.expectEqual(6, ctx.received.items.len); } test "link: stops reading once the worker's inbox backs up" { var pair: [2]posix.socket_t = undefined; if (std.c.socketpair(posix.AF.LOCAL, posix.SOCK.STREAM, 0, &pair) != 0) { return error.SocketPairFailed; } // pair[1] is the link's, closed by its deinit defer sys_net.close(pair[0]); const nonblocking = @as(u32, @bitCast(posix.O{ .NONBLOCK = true })); const flags = try sys_net.fcntl(pair[1], posix.F.GETFL, 0); _ = try sys_net.fcntl(pair[1], posix.F.SETFL, flags | nonblocking); var inbox: Inbox = .{}; defer inbox.deinit(); var link: Link = undefined; try link.init(testing.test_app, pair[1], .cdp, &inbox); defer link.deinit(); // a ceiling below one max-size message would refuse what was configured try testing.expect(link.max_inbox_backlog > testing.test_app.config.cdpMaxMessageSize()); // an empty inbox reads normally (nothing pending, so nothing pushed) const read = try link.readAvailable(1024); try testing.expectEqual(true, read.keep); try testing.expectEqual(false, read.pushed); // a worker that has fallen this far behind isn't going to catch up { const arena = try testing.test_app.arena_pool.acquire(link.max_inbox_backlog, "backlog test"); const payload = try arena.allocator().alloc(u8, link.max_inbox_backlog); inbox.push(arena, .{ .bidi = payload }); } try testing.expectEqual(link.max_inbox_backlog, inbox.queuedBytes()); try testing.expectError(error.InboxBacklog, link.readAvailable(1024)); // and it recovers once the worker drains { const msg = inbox.pop().?; defer msg.deinit(); } try testing.expectEqual(0, inbox.queuedBytes()); _ = try link.readAvailable(1024); }