summaryrefslogtreecommitdiff
path: root/src/9p_io.zig
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-20 02:34:12 -0300
committerGabriel Schneider <[email protected]>2026-10-01 00:12:14 -0300
commitfd50bd971d6dea7eaaa4ee5436ba16b95fa25b30 (patch)
tree5c57378525ea774fee7dc3a698a3d4fd9011712b /src/9p_io.zig
parent2ce8956c8e9843eba7407ab12593871ee451b991 (diff)
downloadpardes-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/9p_io.zig')
-rw-r--r--src/9p_io.zig565
1 files changed, 348 insertions, 217 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;