summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--build.zig.zon4
-rw-r--r--docs/cloud9.md11
-rw-r--r--docs/fs.md10
-rw-r--r--src/9p_io.zig565
-rw-r--r--src/detached/server.zig56
-rw-r--r--src/gui/gui.zig6
-rw-r--r--src/macos.zig4
-rw-r--r--src/tty/tty.zig4
8 files changed, 391 insertions, 269 deletions
diff --git a/build.zig.zon b/build.zig.zon
index 5bf0d5dc..cc53c227 100644
--- a/build.zig.zon
+++ b/build.zig.zon
@@ -9,8 +9,8 @@
// read-only HTTPS URL is what a manifest can carry. Re-pin with
// `zig fetch --save=cloud9 git+https://git.sr.ht/~gbrls/cloud9#<commit>`.
.cloud9 = .{
- .url = "git+https://git.sr.ht/~gbrls/cloud9#65209217b5b68f56bc0bd5bc6c4dce33911ded59",
- .hash = "cloud9-0.1.0-yt86qs3eEQDqC8duRUGBEWOLj82-iCpAb_afpu0KPjPB",
+ .url = "git+https://git.sr.ht/~gbrls/cloud9#ba7ec40782ba7020d82a56896a5eb1b52578d6aa",
+ .hash = "cloud9-0.1.0-yt86qiEeEwBu1DUzSB0tFDCD2R4CS5gSc_Lm1IXmlBwH",
},
// ZLS as a LIBRARY, not a language server: src/lsp_zls.zig imports the
// `zls` module its build.zig publishes and calls the analyser in
diff --git a/docs/cloud9.md b/docs/cloud9.md
index d31d2cb3..d5acc37c 100644
--- a/docs/cloud9.md
+++ b/docs/cloud9.md
@@ -13,8 +13,15 @@ The file-server engine (fids, jobs, parking, flush, hangup) is cloud9's
editor's reply payload locator). `src/9p.zig` names the editor's and the
board's `fs.Options` and re-exports the wire names the transports use.
Mounting, Unix namespace discovery and permissions, the editor event loop,
-connection limits, and exported tree policy remain here. `src/9p_quic.zig`
-selects the existing `pardes-9p` ALPN for cloud9's optional OpenSSL transport.
+connection limits, and exported tree policy remain here. The Unix and TCP
+listeners run on cloud9's `serve.Runner` (`std.Io`: an accept task per
+listener, a reader and a writer task per connection, four slots); its handler
+queues every backend request for the editor's thread, which answers them all
+in `Listener.tick` after each editor update, retries parked reads there, and
+has the runner ship the replies. `src/9p_quic.zig` selects the existing
+`pardes-9p` ALPN for cloud9's optional OpenSSL transport; QUIC still runs on
+the poll loop in `src/9p_io.zig`, since cloud9's QUIC adapter is
+nonblocking-descriptor based rather than `std.Io` based.
The standalone protocol/GPIO tests import the same module. The separate
`05-zig-p4` build also supplies cloud9 for the GPIO firmware entry.
diff --git a/docs/fs.md b/docs/fs.md
index a9293393..55ffdc9c 100644
--- a/docs/fs.md
+++ b/docs/fs.md
@@ -40,8 +40,9 @@ IPv4/IPv6 addresses, not DNS names. Listener port zero chooses a free port;
All connections have session access, including `os`. TCP is unencrypted.
QUIC uses an ephemeral TLS identity without peer verification or login.
It carries 9P2000 on one bidirectional stream with ALPN `pardes-9p`.
-The four application connection slots are shared across transports;
-OpenSSL's internal buffers are separate, dynamically allocated memory.
+Unix and TCP connections share four slots served by cloud9's `std.Io`
+runner; QUIC has four of its own on the editor's poll loop. OpenSSL's
+internal buffers are separate, dynamically allocated memory.
[Plan9port's client](https://9fans.github.io/plan9port/man/man1/9p.html) can
drive Unix or TCP without a kernel mount, and `9ns` mounts the tree in a
@@ -141,8 +142,9 @@ and the editor's reply payload over cloud9's backend contract), `pane.zig`
(pane files), `ctl.zig`, `addr.zig`, `pty.zig`, `events.zig` (event and log
streams), `screen.zig` and `sources.zig`. The protocol engine is cloud9's
`fs.Server`, configured in `src/9p.zig` (the editor's and the board's
-capacities); the transports are `src/9p_io.zig`; `src/fs.zig` keeps host
-access, mounts, resolution, find and grep.
+capacities); the transports are `src/9p_io.zig` (cloud9's `serve.Runner`
+for Unix and TCP, a poll loop for QUIC, and the 9P client for mounts);
+`src/fs.zig` keeps host access, mounts, resolution, find and grep.
`zig build fs-test` drives real sessions using the independent Python client
in `test/ninep.py`; `zig build fs-discovery-test` checks that browsing creates
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;