diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-20 02:34:12 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-10-01 00:12:14 -0300 |
| commit | fd50bd971d6dea7eaaa4ee5436ba16b95fa25b30 (patch) | |
| tree | 5c57378525ea774fee7dc3a698a3d4fd9011712b /src/detached/server.zig | |
| parent | 2ce8956c8e9843eba7407ab12593871ee451b991 (diff) | |
| download | pardes-fd50bd971d6dea7eaaa4ee5436ba16b95fa25b30.tar.gz pardes-fd50bd971d6dea7eaaa4ee5436ba16b95fa25b30.zip | |
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 <[email protected]>
Diffstat (limited to 'src/detached/server.zig')
| -rw-r--r-- | src/detached/server.zig | 56 |
1 files changed, 19 insertions, 37 deletions
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; |
