From fd50bd971d6dea7eaaa4ee5436ba16b95fa25b30 Mon Sep 17 00:00:00 2001 From: Gabriel Schneider Date: Sun, 20 Sep 2026 02:34:12 -0300 Subject: Serve Unix and TCP 9P through cloud9.serve's std.Io runner The hand-written poll loop for Unix/TCP listeners is replaced by cloud9.serve.Runner; requests are queued to the editor thread, which answers them under the connection lock on each frame and retries parked reads as before. QUIC keeps the poll path (its adapter is fd based). The detached server no longer loses a wake that lands between frames. The firmware path keeps driving the engine with push/step. cloud9 re-pinned. Co-Authored-By: Claude Fable 5.1 --- src/detached/server.zig | 56 +++++++++++++++++-------------------------------- 1 file changed, 19 insertions(+), 37 deletions(-) (limited to 'src/detached/server.zig') diff --git a/src/detached/server.zig b/src/detached/server.zig index 6b82cd7b..ac9af2ee 100644 --- a/src/detached/server.zig +++ b/src/detached/server.zig @@ -35,7 +35,7 @@ const read_chunk = 16 * 1024; const pty_chunk = 64 * 1024; -const poll_slots = 2 + max_clients + pardes.MAX_PANES + 2 + 2 + @as(usize, @intFromBool(ninep_io.quic_enabled)) + ninep_io.max_conns; +const poll_slots = 2 + max_clients + pardes.MAX_PANES + 2 + @as(usize, @intFromBool(ninep_io.quic_enabled)); const reload_retries = 4; @@ -131,9 +131,7 @@ const Source = union(enum) { client: u8, pty: u8, inotify, - ninep_listener, ninep_quic, - ninep: u8, }; pub const Session = struct { @@ -158,6 +156,8 @@ pub const Session = struct { inotify_fd: c_int = -1, ninep: ?*ninep_io.Listener = null, ninep_pending: bool = false, + /// Set by the 9P runner's tasks: a request waits for `tick`. + ninep_wake: std.atomic.Value(bool) = .init(false), watches: file_watch.Table = @splat(null), check_files: bool = false, in_loop: bool = false, @@ -587,9 +587,16 @@ pub const Session = struct { if (s.origins()) |c| s.send(c, .detach); } + /// The 9P runner has a request for `tick`: wake the poll loop. + fn wakeNinep(ctx: ?*anyopaque) void { + const s = of(ctx); + s.ninep_wake.store(true, .release); + s.mailbox.signal(); + } + fn pollFrame(ctx: ?*anyopaque) void { const s = of(ctx); - if (s.ninep) |l| s.ninep_pending = if (ninep_io.quic_enabled and l.quic != null) l.tick(s.core).pending else l.drain(s.core).pending; + if (s.ninep) |l| s.ninep_pending = l.tick(s.core).pending; for (&s.ptys, 0..) |*pt, pane| { if (pt.fd < 0) continue; var lbuf: [pardes.memory.limits.host_path_cap + 1]u8 = undefined; @@ -740,7 +747,10 @@ pub const Session = struct { fn waitInput(ctx: ?*anyopaque, timeout_ms: u32) void { const s = of(ctx); - const completed = s.drainCompletions(true); + // Both before the poll: a wake that landed since the last frame has + // had its byte drained here and must not be waited for again. + const drained = s.drainCompletions(true); + const completed = s.ninep_wake.swap(false, .acq_rel) or drained; for (&s.clients) |*c| if (c.fd >= 0) s.flush(c); const regridded = s.reconcile(); @@ -785,14 +795,8 @@ pub const Session = struct { n += 1; } if (s.ninep) |l| { - if (l.accepting()) { - for ([_]c_int{ l.fd, l.tcp_fd }) |fd| { - if (fd < 0) continue; - fds[n] = .{ .fd = fd, .events = poll_in, .revents = 0 }; - src[n] = .ninep_listener; - n += 1; - } - } + // Unix and TCP connections run on cloud9's runner, which wakes + // this loop through the mailbox; only QUIC is polled here. if (comptime ninep_io.quic_enabled) { if (l.quic) |*listener| { fds[n] = listener.poll(); @@ -800,16 +804,6 @@ pub const Session = struct { n += 1; } } - for (0..ninep_io.max_conns) |i| { - if (l.conns[i].fd < 0) continue; - fds[n] = .{ - .fd = l.conns[i].fd, - .events = if (l.owes(@intCast(i))) poll_in | poll_out else poll_in, - .revents = 0, - }; - src[n] = .{ .ninep = @intCast(i) }; - n += 1; - } } if (n == 0) return nap(if (timeout_ms == 0) 16 else timeout_ms); var timeout: c_int = if (timeout_ms == 0) -1 else @intCast(@min(timeout_ms, std.math.maxInt(c_int))); @@ -850,20 +844,7 @@ pub const Session = struct { } }, .inotify => if (pfd.revents & poll_in != 0) s.drainInotify(), - .ninep_listener => if (pfd.revents != 0) { - if (s.ninep) |l| l.accept(); - }, .ninep_quic => {}, - .ninep => |i| if (s.ninep) |l| { - if (!l.live(i)) continue; - if (pfd.revents & poll_out != 0) l.flush(i); - if (!l.live(i)) continue; - if (pfd.revents & poll_in != 0) { - l.fill(i); - } else if (pfd.revents & (poll_hup | poll_err | poll_nval) != 0) { - l.drop(i); - } - }, }; } @@ -1153,8 +1134,9 @@ pub fn run(init: std.process.Init, opts: pardes.Options, name: []const u8) !void return error.NoSocket; } - session.ninep = ninep_io.listen(gpa, opts.ninep_name, name, opts.ninep_tcp, opts.ninep_quic); + session.ninep = ninep_io.listen(init.io, gpa, opts.ninep_name, name, opts.ninep_tcp, opts.ninep_quic); if (session.ninep == null) return error.ListenFailed; + session.ninep.?.setWake(&session, Session.wakeNinep); core.fs.socket_path = session.ninep.?.path(); core.fs.tcp_address = session.ninep.?.tcp_address; core.fs.quic_address = session.ninep.?.quic_address; -- cgit v1.3