diff options
Diffstat (limited to 'src')
| -rw-r--r-- | src/9p_io.zig | 565 | ||||
| -rw-r--r-- | src/detached/server.zig | 56 | ||||
| -rw-r--r-- | src/gui/gui.zig | 6 | ||||
| -rw-r--r-- | src/macos.zig | 4 | ||||
| -rw-r--r-- | src/tty/tty.zig | 4 |
5 files changed, 374 insertions, 261 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig index 57ce67b5..090eb838 100644 --- a/src/9p_io.zig +++ b/src/9p_io.zig @@ -2,7 +2,8 @@ const std = @import("std"); const libc = std.c; const builtin = @import("builtin"); const ninep = @import("9p.zig"); -const transport = @import("cloud9").transport; +const cloud9 = @import("cloud9"); +const transport = cloud9.transport; const pardes = @import("pardes.zig"); const limits = @import("memory.zig").limits; pub const quic_enabled = @import("9p_options").quic; @@ -140,49 +141,112 @@ fn localIp(address: std.Io.net.IpAddress) bool { const Srv = ninep.Server(pardes.ctlfs, ninep.editor); +/// The Unix and TCP listeners run on cloud9's `std.Io` runner: it accepts, +/// reads frames, steps each connection's engine and writes replies on its +/// own tasks, and hands every backend request to the editor's thread +/// through `Slot` (see `tick`). QUIC keeps the poll loop below. +const Runner = if (supported) cloud9.serve.Runner(pardes.ctlfs, ninep.editor, .{ + .msize = msize, + .connections = max_conns, + .listeners = 2, +}) else void; + +/// Answers one backend request on the editor's thread. +fn step(core: *pardes.Pardes, srv: *Srv, req: pardes.ctlfs.Req) void { + core.update(.{ .fs_req = req }); + var answered = false; + while (core.nextEffect()) |effect| { + if (effect == .fs_reply) { + const reply = effect.fs_reply; + if (reply.tag == req.tag) answered = true; + srv.reply(&reply, core.fsPayload(reply)); + } else core.perform(effect); + } + if (!answered) { + const reply = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO); + srv.reply(&reply, ""); + } +} + +/// Pays the backend a request whose connection is gone; replies go nowhere. +fn pay(core: *pardes.Pardes, req: pardes.ctlfs.Req) void { + core.update(.{ .fs_req = req }); + while (core.nextEffect()) |effect| if (effect != .fs_reply) core.perform(effect); +} + +/// One runner slot's mailbox: the requests its connection task handed over, +/// answered by the editor's thread in `tick`. A hangup's releases (at most +/// one per fid) and one outstanding request are the most it ever holds, +/// because a closing connection waits for its entries to be paid before +/// the slot is reused. +const Slot = struct { + const capacity = ninep.editor.fid_capacity + 2; + const Entry = struct { gen: u32, req: pardes.ctlfs.Req }; + + mutex: std.Io.Mutex = .init, + queue: [capacity]Entry = undefined, + head: usize = 0, + len: usize = 0, + /// Bumped when a connection takes the slot; entries carry the value. + gen: u32 = 0, + closing: bool = false, + drained: std.Io.Event = .unset, + /// Requests refused for want of room; never expected. + refused: u32 = 0, + + fn push(s: *Slot, req: pardes.ctlfs.Req) bool { + if (s.len == s.queue.len) { + s.refused += 1; + return false; + } + s.queue[(s.head + s.len) % capacity] = .{ .gen = s.gen, .req = req }; + s.len += 1; + return true; + } + + fn pop(s: *Slot) ?Entry { + if (s.len == 0) return null; + const e = s.queue[s.head]; + s.head = (s.head + 1) % capacity; + s.len -= 1; + return e; + } +}; + +const quic_slots = if (quic_enabled) max_conns else 0; + +/// A QUIC connection on the poll path. const Conn = struct { - fd: c_int = -1, quic: if (quic_enabled) ?quic.Connection else void = if (quic_enabled) null else {}, draining: bool = false, accepted_ms: i64 = 0, srv: Srv = undefined, in: [msize]u8 = undefined, out: [2 * msize]u8 = undefined, - - fn step(c: *Conn, core: *pardes.Pardes, req: pardes.ctlfs.Req) void { - core.update(.{ .fs_req = req }); - var answered = false; - while (core.nextEffect()) |effect| { - if (effect == .fs_reply) { - const reply = effect.fs_reply; - if (reply.tag == req.tag) answered = true; - c.srv.reply(&reply, core.fsPayload(reply)); - } else core.perform(effect); - } - if (!answered) { - const reply = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO); - c.srv.reply(&reply, ""); - } - } }; const accept_pause_ms: i64 = 100; pub const Listener = struct { - fd: c_int = -1, - tcp_fd: c_int = -1, + io: std.Io, + runner: Runner = undefined, + slots: [max_conns]Slot = @splat(.{}), + stopping: std.atomic.Value(bool) = .init(false), + /// While `reset` runs: the generation each slot is paying off; anything + /// newer waits for the replacement core. + resetting: ?[max_conns]u32 = null, tcp_address: ?std.Io.net.IpAddress = null, quic: if (quic_enabled) ?quic.Listener else void = if (quic_enabled) null else {}, quic_address: ?std.Io.net.IpAddress = null, paused_ms: i64 = 0, path_buf: [sun_path_len]u8 = undefined, path_len: usize = 0, - conns: [max_conns]Conn = @splat(.{}), + conns: [quic_slots]Conn = @splat(.{}), control: [2]c_int = .{ -1, -1 }, watcher: ?std.Thread = null, - stopping: std.atomic.Value(bool) = .init(false), + watch_stop: std.atomic.Value(bool) = .init(false), watch_lock: std.atomic.Mutex = .unlocked, - watch_fds: [max_conns + 3 + @as(usize, @intFromBool(quic_enabled))]libc.pollfd = undefined, + watch_fds: [2]libc.pollfd = undefined, watch_len: usize = 0, watch_timeout: c_int = -1, wake_ctx: ?*anyopaque = null, @@ -192,104 +256,103 @@ pub const Listener = struct { return l.path_buf[0..l.path_len]; } - fn listenTcp(l: *Listener, address: std.Io.net.IpAddress) !void { - const fd = try transport.listenFd(.{ .tcp = address }, max_conns); - errdefer transport.close(fd); - var addr: libc.sockaddr.storage = undefined; - var len: libc.socklen_t = @sizeOf(@TypeOf(addr)); - if (libc.getsockname(fd, @ptrCast(&addr), &len) != 0) return error.SocketAddressFailed; - l.tcp_address = sockaddrIp(@ptrCast(&addr)) orelse return error.SocketAddressFailed; - l.tcp_fd = fd; - log.info("serving 9P2000 over TCP on {f}", .{l.tcp_address.?}); + fn of(ctx: ?*anyopaque) *Listener { + return @ptrCast(@alignCast(ctx.?)); } - pub fn accept(l: *Listener) void { - if (comptime !supported) return; - for ([_]c_int{ l.fd, l.tcp_fd }) |listener_fd| { - if (listener_fd < 0) continue; - for (0..max_conns + 1) |_| { - const fd = (transport.acceptFd(listener_fd, listener_fd == l.tcp_fd) catch { - l.paused_ms = nowMs() +| accept_pause_ms; - log.warn("accept failed; pausing the listener for {d} ms", .{accept_pause_ms}); - return; - }) orelse break; - const c = for (&l.conns, 0..) |*cand, i| { - if (!l.live(@intCast(i)) and !cand.draining) break cand; - } else { - log.debug("refusing a connection, all {d} slots busy", .{max_conns}); - _ = libc.close(fd); - continue; - }; - c.fd = fd; - c.draining = false; - c.accepted_ms = nowMs(); - c.srv = .init(.{ - .in = &c.in, - .out = &c.out, - .root = pardes.ctlfs.root, - }); - } + /// The runner's handler: on the connection's task. + fn onOpened(ctx: ?*anyopaque, conn: *Runner.Conn) void { + const l = of(ctx); + const s = &l.slots[conn.index]; + s.mutex.lockUncancelable(l.io); + s.gen +%= 1; + s.mutex.unlock(l.io); + } + + fn onServe(ctx: ?*anyopaque, conn: *Runner.Conn, req: pardes.ctlfs.Req) void { + const l = of(ctx); + const s = &l.slots[conn.index]; + s.mutex.lockUncancelable(l.io); + const queued = s.push(req); + s.mutex.unlock(l.io); + if (!queued) { + const reply = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO); + conn.reply(&reply, ""); + return; } - if (comptime quic_enabled) { - if (l.quic) |*listener| for (0..max_conns + 1) |_| { - var connection = (listener.accept() catch |err| { - log.warn("QUIC accept failed: {s}", .{@errorName(err)}); - return; - }) orelse break; - const c = for (&l.conns, 0..) |*cand, i| { - if (!l.live(@intCast(i)) and !cand.draining) break cand; - } else { - connection.deinit(); - continue; - }; - c.fd = -1; - c.quic = connection; - c.draining = false; - c.accepted_ms = nowMs(); - c.srv = .init(.{ .in = &c.in, .out = &c.out, .root = pardes.ctlfs.root }); - }; + l.kick(); + } + + /// Holds the slot until the editor's thread has paid what the + /// connection still owes, so a successor never inherits its entries. + fn onClosed(ctx: ?*anyopaque, conn: *Runner.Conn) void { + const l = of(ctx); + const s = &l.slots[conn.index]; + if (l.stopping.load(.acquire)) return; + s.mutex.lockUncancelable(l.io); + if (s.len == 0) { + s.mutex.unlock(l.io); + return; } + s.closing = true; + s.drained.reset(); + s.mutex.unlock(l.io); + l.kick(); + s.drained.wait(l.io) catch {}; + } + + fn kick(l: *Listener) void { + if (l.stopping.load(.acquire)) return; + if (l.wake) |f| f(l.wake_ctx); + } + + /// Accepts and reads QUIC connections (the runner does its own). + pub fn accept(l: *Listener) void { + if (comptime !quic_enabled) return; + if (l.quic) |*listener| for (0..max_conns + 1) |_| { + var connection = (listener.accept() catch |err| { + log.warn("QUIC accept failed: {s}", .{@errorName(err)}); + return; + }) orelse break; + const c = for (&l.conns, 0..) |*cand, i| { + if (!l.live(@intCast(i)) and !cand.draining) break cand; + } else { + connection.deinit(); + continue; + }; + c.quic = connection; + c.draining = false; + c.accepted_ms = nowMs(); + c.srv = .init(.{ .in = &c.in, .out = &c.out, .root = pardes.ctlfs.root }); + }; } pub const greet_deadline_ms: i64 = 5000; pub fn expire(l: *Listener) void { - if (comptime !supported) return; + if (comptime !quic_enabled) return; const now = nowMs(); if (now == 0) return; for (&l.conns, 0..) |*c, i| { if (!l.live(@intCast(i)) or c.srv.protocol.msize != 0) continue; if (now - c.accepted_ms < greet_deadline_ms) continue; - log.debug("slot {d} never sent Tversion; taking it back", .{i}); + log.debug("QUIC slot {d} never sent Tversion; taking it back", .{i}); l.drop(@intCast(i)); } } - pub fn accepting(l: *const Listener) bool { - if (comptime !supported) return false; - if (l.fd < 0 and l.tcp_fd < 0) return false; - if (l.paused_ms == 0) return true; - const now = nowMs(); - return now == 0 or now >= l.paused_ms; - } - pub fn nextDue(l: *const Listener) ?i32 { - if (comptime !supported) return null; + if (comptime !quic_enabled) return null; const now = nowMs(); if (now == 0) return null; var due: ?i64 = null; - if (l.paused_ms > now) due = l.paused_ms; - if (comptime quic_enabled) { - if (l.quic) |*listener| if (listener.nextDue()) |ms| { - const at = now + ms; - due = if (due) |d| @min(d, at) else at; - }; - } + if (l.quic) |*listener| if (listener.nextDue()) |ms| { + const at = now + ms; + due = if (due) |d| @min(d, at) else at; + }; for (&l.conns, 0..) |*c, i| { if (!l.live(@intCast(i))) continue; - if (comptime quic_enabled) { - if (c.quic) |*connection| if (connection.pending()) return 0; - } + if (c.quic) |*connection| if (connection.pending()) return 0; if (c.srv.protocol.msize != 0) continue; const at = c.accepted_ms + greet_deadline_ms; due = if (due) |d| @min(d, at) else at; @@ -301,16 +364,13 @@ pub const Listener = struct { const nowMs = transport.nowMs; pub fn fill(l: *Listener, i: u8) void { - if (comptime !supported) return; + if (comptime !quic_enabled) return; const c = &l.conns[i]; if (c.srv.protocol.dead) return l.drop(i); const room = c.srv.protocol.in.len - c.srv.protocol.in_len; if (room == 0) return; var buf: [msize]u8 = undefined; - const got = if (quic_enabled and c.quic != null) - (c.quic.?.read(buf[0..@min(room, buf.len)]) catch return l.drop(i)) orelse return - else - (transport.read(c.fd, buf[0..@min(room, buf.len)]) catch return l.drop(i)) orelse return; + const got = (c.quic.?.read(buf[0..@min(room, buf.len)]) catch return l.drop(i)) orelse return; if (got == 0) return l.drop(i); const n = c.srv.push(buf[0..@intCast(got)]); if (c.srv.protocol.dead) return l.drop(i); @@ -318,40 +378,29 @@ pub const Listener = struct { } pub fn flush(l: *Listener, i: u8) void { - if (comptime !supported) return; + if (comptime !quic_enabled) return; const c = &l.conns[i]; if (!l.live(i)) return; while (true) { const bytes = c.srv.output(); if (bytes.len == 0) return; - const n = if (quic_enabled and c.quic != null) - c.quic.?.write(bytes) catch return l.drop(i) - else - (transport.write(c.fd, bytes) catch return l.drop(i)) orelse return; + const n = c.quic.?.write(bytes) catch return l.drop(i); if (n == 0) return; c.srv.wrote(@intCast(n)); } } - pub fn owes(l: *const Listener, i: u8) bool { - return l.conns[i].srv.output().len != 0; - } - pub fn live(l: *const Listener, i: u8) bool { - return l.conns[i].fd >= 0 or (quic_enabled and l.conns[i].quic != null); + if (comptime !quic_enabled) return false; + return l.conns[i].quic != null; } pub fn drop(l: *Listener, i: u8) void { + if (comptime !quic_enabled) return; const c = &l.conns[i]; - if (c.fd >= 0) { - _ = libc.close(c.fd); - c.fd = -1; - } - if (comptime quic_enabled) { - if (c.quic) |*connection| { - connection.deinit(); - c.quic = null; - } + if (c.quic) |*connection| { + connection.deinit(); + c.quic = null; } if (c.draining or c.accepted_ms == 0) return; c.srv.hangup(); @@ -360,38 +409,80 @@ pub const Listener = struct { pub const Drained = struct { count: usize = 0, pending: bool = false }; + /// Answers what every connection has asked since the last call, retries + /// their parked reads against the editor's current state and has the + /// runner send the replies. Call it after each editor update, on the + /// editor's thread; `pending` says QUIC has more to do right away. pub fn drain(l: *Listener, core: *pardes.Pardes) Drained { core.fs.socket_path = l.path(); core.fs.tcp_address = l.tcp_address; core.fs.quic_address = l.quic_address; l.expire(); var result: Drained = .{}; - for (&l.conns, 0..) |*conn, i| { - if (!l.live(@intCast(i)) and !conn.draining) continue; - var count: usize = 0; - while (conn.srv.retry()) |req| { - conn.step(core, req); - l.collectOs(core); - count += 1; - } - while (count < 64) { - const req = conn.srv.next() orelse break; - conn.step(core, req); - l.collectOs(core); - count += 1; + if (comptime supported) { + for (&l.slots, 0..) |*s, i| { + const conn = &l.runner.conns[i]; + s.mutex.lockUncancelable(l.io); + const only: ?u32 = if (l.resetting) |gens| gens[i] else null; + var count: usize = 0; + if (s.len != 0 or conn.live()) { + conn.lock(); + while (s.len != 0) { + if (only) |gen| if (s.queue[s.head].gen != gen) break; + const e = s.pop().?; + if (e.gen == s.gen) step(core, &conn.engine, e.req) else if (e.req.op == .release) pay(core, e.req); + count += 1; + } + if (conn.live() and (only == null or only.? == s.gen)) { + while (conn.engine.retry()) |req| { + step(core, &conn.engine, req); + count += 1; + } + while (conn.engine.next()) |req| { + step(core, &conn.engine, req); + count += 1; + } + } + const owed = conn.engine.output().len != 0; + conn.unlock(); + if (owed) conn.flush(); + } + if (s.closing and s.len == 0) { + s.closing = false; + s.drained.set(l.io); + } + s.mutex.unlock(l.io); + result.count += count; + if (count != 0) l.collectOs(core); } - result.count += count; - result.pending = result.pending or count >= 64; - if (conn.draining) { - if (count == 0) conn.draining = false else result.pending = true; - } else l.flush(@intCast(i)); - if (comptime quic_enabled) { + } + if (comptime quic_enabled) { + for (&l.conns, 0..) |*conn, i| { + if (!l.live(@intCast(i)) and !conn.draining) continue; + var count: usize = 0; + while (conn.srv.retry()) |req| { + step(core, &conn.srv, req); + l.collectOs(core); + count += 1; + } + while (count < 64) { + const req = conn.srv.next() orelse break; + step(core, &conn.srv, req); + l.collectOs(core); + count += 1; + } + result.count += count; + result.pending = result.pending or count >= 64; + if (conn.draining) { + if (count == 0) conn.draining = false else result.pending = true; + } else l.flush(@intCast(i)); if (conn.quic) |*connection| result.pending = result.pending or connection.pending(); } } return result; } + /// `drain` plus the QUIC listener's events, accepts and reads. pub fn tick(l: *Listener, core: *pardes.Pardes) Drained { if (comptime quic_enabled) { if (l.quic) |*listener| listener.events() catch |err| { @@ -401,17 +492,44 @@ pub const Listener = struct { l.quic = null; l.quic_address = null; }; + l.accept(); + for (0..quic_slots) |i| if (l.live(@intCast(i))) l.fill(@intCast(i)); } - if (l.accepting()) l.accept(); - for (0..max_conns) |i| if (l.live(@intCast(i))) l.fill(@intCast(i)); const result = l.drain(core); l.arm(); return result; } + /// Hangs every connection up and pays `core` what they owed, so that a + /// replacement core starts with no client holding anything. A client + /// that connects meanwhile is not answered until the caller has put + /// the replacement in and calls `tick` again. pub fn reset(l: *Listener, core: *pardes.Pardes) void { - for (0..max_conns) |i| l.drop(@intCast(i)); - while (l.drain(core).pending) {} + for (0..quic_slots) |i| l.drop(@intCast(i)); + if (comptime supported) { + var gens: [max_conns]u32 = undefined; + for (&l.slots, &gens) |*s, *gen| { + s.mutex.lockUncancelable(l.io); + gen.* = s.gen; + s.mutex.unlock(l.io); + } + l.resetting = gens; + defer l.resetting = null; + l.runner.closeAll(); + const deadline = nowMs() +| 2000; + while (true) { + const drained = l.drain(core); + var idle = !drained.pending; + for (&l.slots, &gens, 0..) |*s, gen, i| { + s.mutex.lockUncancelable(l.io); + if (s.len != 0 and s.queue[s.head].gen == gen) idle = false; + if (s.gen == gen and l.runner.conns[i].live()) idle = false; + s.mutex.unlock(l.io); + } + if (idle or nowMs() >= deadline) break; + Client.nap(1); + } + } else while (l.drain(core).pending) {} l.collectOs(core); for (&l.conns) |*conn| conn.accepted_ms = 0; l.arm(); @@ -421,15 +539,7 @@ pub const Listener = struct { var i: usize = 0; while (i < core.fs.os_paths.items.len) { const entry = core.fs.os_paths.items[i]; - var held = false; - for (&l.conns, 0..) |*conn, j| { - if (!l.live(@intCast(j)) and !conn.draining) continue; - if (conn.srv.references(entry.node)) { - held = true; - break; - } - } - if (held) { + if (l.references(entry.node)) { i += 1; } else { core.gpa.free(entry.path); @@ -438,8 +548,42 @@ pub const Listener = struct { } } - pub fn wakeThread(l: *Listener, ctx: ?*anyopaque, wake: *const fn (?*anyopaque) void) !void { - if (l.watcher != null) return; + /// Whether any connection, or a release still to be paid, names `node`. + fn references(l: *Listener, node: u64) bool { + if (comptime supported) { + for (&l.slots, 0..) |*s, i| { + const conn = &l.runner.conns[i]; + s.mutex.lockUncancelable(l.io); + defer s.mutex.unlock(l.io); + var k: usize = 0; + while (k < s.len) : (k += 1) { + if (s.queue[(s.head + k) % Slot.capacity].req.node == node) return true; + } + if (!conn.live()) continue; + conn.lock(); + defer conn.unlock(); + if (conn.engine.references(node)) return true; + } + } + for (&l.conns, 0..) |*conn, j| { + if (!l.live(@intCast(j)) and !conn.draining) continue; + if (conn.srv.references(node)) return true; + } + return false; + } + + /// Where the runner's tasks report that a request waits for `tick`. + pub fn setWake(l: *Listener, ctx: ?*anyopaque, wake: *const fn (?*anyopaque) void) void { + l.wake_ctx = ctx; + l.wake = wake; + } + + /// `setWake`, and for QUIC a thread that watches its socket and calls + /// `wake` when `tick` has something to read. + pub fn watch(l: *Listener, ctx: ?*anyopaque, wake: *const fn (?*anyopaque) void) !void { + l.setWake(ctx, wake); + if (comptime !quic_enabled) return; + if (l.quic == null or l.watcher != null) return; if (libc.pipe(&l.control) != 0) return error.PipeFailed; errdefer { _ = libc.close(l.control[0]); @@ -450,37 +594,18 @@ pub const Listener = struct { setCloexec(fd); setNonblock(fd); } - l.wake_ctx = ctx; - l.wake = wake; l.arm(); - l.watcher = try std.Thread.spawn(.{}, watch, .{l}); + l.watcher = try std.Thread.spawn(.{}, watchQuic, .{l}); } fn arm(l: *Listener) void { + if (comptime !quic_enabled) return; if (l.control[1] < 0) return; while (!l.watch_lock.tryLock()) std.atomic.spinLoopHint(); l.watch_fds[0] = .{ .fd = l.control[0], .events = @intCast(libc.POLL.IN), .revents = 0 }; l.watch_len = 1; - if (l.accepting()) { - for ([_]c_int{ l.fd, l.tcp_fd }) |fd| { - if (fd < 0) continue; - l.watch_fds[l.watch_len] = .{ .fd = fd, .events = @intCast(libc.POLL.IN), .revents = 0 }; - l.watch_len += 1; - } - } - if (comptime quic_enabled) { - if (l.quic) |*listener| { - l.watch_fds[l.watch_len] = listener.poll(); - l.watch_len += 1; - } - } - for (0..max_conns) |i| { - if (l.conns[i].fd < 0) continue; - l.watch_fds[l.watch_len] = .{ - .fd = l.conns[i].fd, - .events = @as(i16, @intCast(libc.POLL.IN)) | if (l.owes(@intCast(i))) @as(i16, @intCast(libc.POLL.OUT)) else 0, - .revents = 0, - }; + if (l.quic) |*listener| { + l.watch_fds[l.watch_len] = listener.poll(); l.watch_len += 1; } l.watch_timeout = l.nextDue() orelse -1; @@ -488,10 +613,10 @@ pub const Listener = struct { _ = libc.write(l.control[1], "w", 1); } - fn watch(l: *Listener) void { + fn watchQuic(l: *Listener) void { var notified = false; - while (!l.stopping.load(.acquire)) { - var fds: [max_conns + 3 + @as(usize, @intFromBool(quic_enabled))]libc.pollfd = undefined; + while (!l.watch_stop.load(.acquire)) { + var fds: [2]libc.pollfd = undefined; while (!l.watch_lock.tryLock()) std.atomic.spinLoopHint(); const len = if (notified) 1 else l.watch_len; @memcpy(fds[0..len], l.watch_fds[0..len]); @@ -505,26 +630,25 @@ pub const Listener = struct { notified = false; continue; } - if (!l.stopping.load(.acquire)) l.wake.?(l.wake_ctx); + if (!l.watch_stop.load(.acquire)) l.wake.?(l.wake_ctx); notified = true; } } pub fn deinit(l: *Listener, gpa: std.mem.Allocator) void { + l.stopping.store(true, .release); if (l.watcher) |thread| { - l.stopping.store(true, .release); + l.watch_stop.store(true, .release); _ = libc.write(l.control[1], "q", 1); thread.join(); for (l.control) |fd| _ = libc.close(fd); } - for (0..max_conns) |i| l.drop(@intCast(i)); + for (0..quic_slots) |i| l.drop(@intCast(i)); if (comptime quic_enabled) { if (l.quic) |*listener| listener.deinit(); } - if (l.tcp_fd >= 0) _ = libc.close(l.tcp_fd); - if (l.fd >= 0) { - _ = libc.close(l.fd); - l.fd = -1; + if (comptime supported) l.runner.stop(); + if (l.path_len != 0) { var z: [sun_path_len:0]u8 = undefined; @memcpy(z[0..l.path_len], l.path_buf[0..l.path_len]); z[l.path_len] = 0; @@ -540,7 +664,7 @@ pub fn socketPath(buf: *[sun_path_len]u8, dir: []const u8, name: []const u8) ?[: return std.fmt.bufPrintSentinel(buf, "{s}/" ++ prefix ++ "{s}.sock", .{ dir, name }, 0) catch null; } -pub fn listen(gpa: std.mem.Allocator, named: []const u8, fallback: []const u8, tcp_dial: ?[]const u8, quic_dial: ?[]const u8) ?*Listener { +pub fn listen(io: std.Io, gpa: std.mem.Allocator, named: []const u8, fallback: []const u8, tcp_dial: ?[]const u8, quic_dial: ?[]const u8) ?*Listener { if (comptime !supported) return null; var dir_buf: [sun_path_len:0]u8 = undefined; const dir = socketDir(&dir_buf) orelse { @@ -549,36 +673,38 @@ pub fn listen(gpa: std.mem.Allocator, named: []const u8, fallback: []const u8, t }; if (!ensureSocketDir(dir)) return null; const l = gpa.create(Listener) catch return null; - l.* = .{}; + l.* = .{ .io = io }; + l.runner.init(.{ + .io = io, + .root = pardes.ctlfs.root, + .handler = .{ .ctx = l, .serve = Listener.onServe, .opened = Listener.onOpened, .closed = Listener.onClosed }, + .greet_timeout_ms = Listener.greet_deadline_ms, + }); const p = socketPath(&l.path_buf, dir, if (named.len != 0) named else fallback) orelse { gpa.destroy(l); return null; }; - const fd = transport.listenFd(.{ .unix = p }, max_conns) catch |err| retry: { - const io = std.Io.Threaded.global_single_threaded.io(); + _ = l.runner.listen(.{ .unix = p }, max_conns) catch |err| retry: { const existing = std.Io.Dir.cwd().statFile(io, p, .{ .follow_symlinks = false }) catch null; - if (err != error.Bind or existing == null or existing.?.kind != .unix_domain_socket or alive(p)) { + if (err != error.AddressInUse or existing == null or existing.?.kind != .unix_domain_socket or alive(p)) { log.warn("something is already listening on {s}", .{p}); - gpa.destroy(l); + l.deinit(gpa); return null; } if (libc.unlink(p) != 0) { - gpa.destroy(l); + l.deinit(gpa); return null; } - break :retry transport.listenFd(.{ .unix = p }, max_conns) catch { - gpa.destroy(l); + break :retry l.runner.listen(.{ .unix = p }, max_conns) catch { + l.deinit(gpa); return null; }; }; + l.path_len = p.len; if (libc.chmod(p, 0o600) != 0) { - transport.close(fd); - _ = libc.unlink(p); - gpa.destroy(l); + l.deinit(gpa); return null; } - l.fd = fd; - l.path_len = p.len; if (tcp_dial) |dial| { if (!std.mem.startsWith(u8, dial, "tcp!")) { l.deinit(gpa); @@ -589,11 +715,13 @@ pub fn listen(gpa: std.mem.Allocator, named: []const u8, fallback: []const u8, t l.deinit(gpa); return null; }; - l.listenTcp(address) catch |err| { + const bound = l.runner.listen(.{ .tcp = address }, max_conns) catch |err| { log.warn("cannot listen on {s}: {s}", .{ dial, @errorName(err) }); l.deinit(gpa); return null; }; + l.tcp_address = canonicalIp(bound); + log.info("serving 9P2000 over TCP on {f}", .{l.tcp_address.?}); } if (quic_dial) |dial| { if (comptime quic_enabled) { @@ -651,9 +779,14 @@ test "a name that is not one path component is no address at all" { test "one connection's buffers are sized from the one msize constant" { try testing.expect(msize >= ninep.min_msize); - const c: Conn = .{}; - try testing.expectEqual(@as(usize, msize), c.in.len); - try testing.expectEqual(@as(usize, 2 * msize), c.out.len); + try testing.expectEqual(@as(u32, msize), Runner.msize); + try testing.expectEqual(@as(usize, msize), @typeInfo(@FieldType(Runner.Conn, "in")).array.len); + try testing.expectEqual(@as(usize, 2 * msize), @typeInfo(@FieldType(Runner.Conn, "out")).array.len); + try testing.expectEqual(@as(usize, max_conns), Runner.max_connections); + if (quic_enabled) { + const c: Conn = .{}; + try testing.expectEqual(@as(usize, msize), c.in.len); + } } test "TCP addresses are numeric and normalize mapped IPv4" { @@ -776,14 +909,14 @@ test "Unix TCP and QUIC share one listener through reads writes reconnects and r var bind_buf: [64]u8 = undefined; const bind = try std.fmt.bufPrint(&bind_buf, "{s}!{s}!0", .{ protocol, host }); const tcp = std.mem.eql(u8, protocol, "tcp"); - const l = listen(gpa, "roundtrip", "", if (tcp) bind else null, if (tcp) null else bind) orelse return error.ListenFailed; + const l = listen(testing.io, gpa, "roundtrip", "", if (tcp) bind else null, if (tcp) null else bind) orelse return error.ListenFailed; defer { l.reset(p); l.deinit(gpa); } - try testing.expect(l.fd >= 0); - try testing.expectEqual(tcp, l.tcp_fd >= 0); - try testing.expectEqual(@as(usize, 4), l.conns.len); + try testing.expect(l.path().len != 0); + try testing.expectEqual(tcp, l.tcp_address != null); + try testing.expectEqual(@as(usize, 4), l.runner.conns.len); try testing.expect(l.watcher == null); const port = (if (tcp) l.tcp_address else l.quic_address).?.getPort(); try testing.expect(port != 0); @@ -802,11 +935,9 @@ test "Unix TCP and QUIC share one listener through reads writes reconnects and r } try testing.expect(worker.done.load(.acquire)); if (worker.failure) |err| return err; - const unix_fd = l.fd; - const tcp_fd = l.tcp_fd; l.reset(p); - try testing.expectEqual(unix_fd, l.fd); - try testing.expectEqual(tcp_fd, l.tcp_fd); + try testing.expectEqual(@as(usize, 0), l.runner.count()); + for (&l.slots) |*slot| try testing.expectEqual(@as(usize, 0), slot.len); for (&l.conns, 0..) |conn, i| try testing.expect(!l.live(@intCast(i)) and !conn.draining); for (p.fs.snapshots) |snapshot| try testing.expect(snapshot.node == 0); } @@ -816,10 +947,10 @@ test "Unix TCP and QUIC share one listener through reads writes reconnects and r extern "c" fn setenv(name: [*:0]const u8, value: [*:0]const u8, overwrite: c_int) c_int; extern "c" fn unsetenv(name: [*:0]const u8) c_int; -pub fn start(gpa: std.mem.Allocator, core: *pardes.Pardes) ?*Listener { +pub fn start(io: std.Io, gpa: std.mem.Allocator, core: *pardes.Pardes) ?*Listener { var name: [16]u8 = undefined; const fallback = std.fmt.bufPrint(&name, "{d}", .{@as(u32, @intCast(libc.getpid()))}) catch unreachable; - const listener = listen(gpa, core.opts.ninep_name, fallback, core.opts.ninep_tcp, core.opts.ninep_quic) orelse { + const listener = listen(io, gpa, core.opts.ninep_name, fallback, core.opts.ninep_tcp, core.opts.ninep_quic) orelse { core.reportError(0, "9p listener", error.ListenFailed); return null; }; @@ -869,7 +1000,7 @@ test "9P shell environment preserves identity when nested Look forwarding is dis const listener = try testing.allocator.create(Listener); defer testing.allocator.destroy(listener); - listener.* = .{}; + listener.* = .{ .io = testing.io }; const path = "/tmp/pardes-example.sock"; @memcpy(listener.path_buf[0..path.len], path); listener.path_len = path.len; 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; diff --git a/src/gui/gui.zig b/src/gui/gui.zig index dbad10d4..27f3346c 100644 --- a/src/gui/gui.zig +++ b/src/gui/gui.zig @@ -2382,7 +2382,7 @@ fn localSession( inotify_fd = -1; }; var watches: file_watch.Table = @splat(null); - var fs = ninep_io.start(gpa, core); + var fs = ninep_io.start(io, gpa, core); defer if (fs) |f| f.deinit(gpa); var shell: Shell = .{ @@ -2412,7 +2412,7 @@ fn localSession( watch_stop = try stopPipe(); watch_reader = try std.Thread.spawn(.{}, watchThread, .{ inotify_fd, watch_stop[0], &queue }); } - if (fs) |f| try f.wakeThread(&queue, wakeFs); + if (fs) |f| try f.watch(&queue, wakeFs); _ = c.SDL_StartTextInput(g.window); @@ -2680,7 +2680,7 @@ fn runGrid(init: std.process.Init, opts_in: pardes.Options) !void { var pipe_tasks: PipeTasks = .{}; defer pipe_tasks.cancelAll(io); var watches: file_watch.Table = @splat(null); - const fs = ninep_io.start(gpa, core); + const fs = ninep_io.start(io, gpa, core); defer if (fs) |f| f.deinit(gpa); var shell: Shell = .{ .core = core, diff --git a/src/macos.zig b/src/macos.zig index d0d379de..79036c3e 100644 --- a/src/macos.zig +++ b/src/macos.zig @@ -745,7 +745,7 @@ fn initCore(runtime: ?*const Runtime, cols_arg: u16, rows_arg: u16) !void { .runtime = if (runtime) |r| r.* else .{}, }; const st = &state.?; - st.ninep = ninep_io.start(gpa, core); + st.ninep = ninep_io.start(io, gpa, core); core.host = hostFor(st); const cols = @max(1, cols_arg); @@ -754,7 +754,7 @@ fn initCore(runtime: ?*const Runtime, cols_arg: u16, rows_arg: u16) !void { while (core.nextEffect()) |effect| core.perform(effect); st.started = true; - if (st.ninep) |listener| listener.wakeThread(st, wakeNinep) catch |err| core.reportError(0, "9p wake", err); + if (st.ninep) |listener| listener.watch(st, wakeNinep) catch |err| core.reportError(0, "9p wake", err); for (&st.ptys, 0..) |*slot, id| if (slot.*) |*pt| startReader(st, pt, @intCast(id)); pardes.lsp.setStatusSink(st, lspStatusSink); diff --git a/src/tty/tty.zig b/src/tty/tty.zig index 1927577b..23b7c09b 100644 --- a/src/tty/tty.zig +++ b/src/tty/tty.zig @@ -653,7 +653,7 @@ fn localSession( var frame_arena: std.heap.ArenaAllocator = .init(allocs.frame); defer frame_arena.deinit(); - var fs = ninep_io.start(gpa, core); + var fs = ninep_io.start(io, gpa, core); defer if (fs) |f| f.deinit(gpa); var sh: Shell = .{ @@ -717,7 +717,7 @@ fn localSession( try startInput(loop, input_cache); defer if (attached.* == null) stopInput(loop); (try std.Thread.spawn(.{}, winchWatch, .{ loop, vx, tty })).detach(); - if (fs) |f| try f.wakeThread(loop, wakeFs); + if (fs) |f| try f.watch(loop, wakeFs); try vx.queryTerminalSend(tty.writer()); sh.threads_ok = true; |
