diff options
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; |
