summaryrefslogtreecommitdiff
path: root/src/detached/server.zig
diff options
context:
space:
mode:
Diffstat (limited to 'src/detached/server.zig')
-rw-r--r--src/detached/server.zig56
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;