summaryrefslogtreecommitdiff
path: root/src/serve.zig
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-20 02:33:02 -0300
committerGabriel Schneider <[email protected]>2026-09-20 02:33:02 -0300
commitba7ec40782ba7020d82a56896a5eb1b52578d6aa (patch)
treeb6a5222d124255ae1e5d568c6ffbfe2e2050ae63 /src/serve.zig
parent66f2e492c348677ab3050f5e378b9eb4c04c98ce (diff)
downloadcloud9-ba7ec40782ba7020d82a56896a5eb1b52578d6aa.tar.gz
cloud9-ba7ec40782ba7020d82a56896a5eb1b52578d6aa.zip
Add cloud9.serve: an std.Io runner around the file-server engine
Runner(Backend, Options, Limits) listens on Unix or TCP, runs a reader and a serve task per connection in one Io.Group, pushes frames into an fs.Server, and lets the backend answer now or later from any task or thread (reply, flush, wake); Tflush, greet timeout, connection limit, close and stop with cancellation are covered by tests over real sockets. The engine and the backend contract stay Io-free, so the push/step mode for freestanding targets is unchanged. Documented as the two ways to drive the engine. Co-Authored-By: Claude Fable 5.1 <[email protected]>
Diffstat (limited to 'src/serve.zig')
-rw-r--r--src/serve.zig848
1 files changed, 848 insertions, 0 deletions
diff --git a/src/serve.zig b/src/serve.zig
new file mode 100644
index 0000000..255e2ff
--- /dev/null
+++ b/src/serve.zig
@@ -0,0 +1,848 @@
+//! An `std.Io` runner for the file-server engine on hosted targets:
+//! listeners, a bounded table of connections, two tasks per connection (one
+//! reads frames, one steps the engine and writes replies) and the wake-up a
+//! backend needs when it answers later or has news for a parked read.
+//!
+//! The engine (`fs.Server`) and the backend contract are untouched: the
+//! engine stays free of OS calls and threads, and freestanding targets keep
+//! driving it with `push()`/`output()`/`wrote()` and `next()`/`retry()`/
+//! `reply()` themselves. This module is an adapter beside `http`, the way
+//! `web/main.zig` runs the gateway. See docs/design.md, "Io runner".
+const std = @import("std");
+const Io = std.Io;
+const fs = @import("fs.zig");
+const transport = @import("transport.zig");
+
+/// Comptime bounds of one runner.
+pub const Limits = struct {
+ /// The largest frame a connection negotiates: the engine's input buffer
+ /// is this big, its output buffer twice that.
+ msize: u32 = 8192,
+ /// Connections served at once; one more is accepted and closed at once.
+ connections: usize = 4,
+ /// Listeners `listen()` may add.
+ listeners: usize = 2,
+};
+
+/// A runner for `fs.Server(Backend, opts)` bounded by `limits`.
+///
+/// The backend is reached through a `Handler`: `serve` runs on the
+/// connection's task with the engine unlocked and answers with
+/// `Conn.reply()`, at once or later from any task or thread. A backend that
+/// answers `Status.again` is asked again when someone calls `Conn.wake()`
+/// (or `wakeAll()`), which is the runner's only spontaneous retry; new
+/// input, replies and closing do not retry parked requests.
+pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits: Limits) type {
+ if (limits.msize < fs.msize_min) @compileError("9P runner msize must hold one full Rwalk");
+ if (limits.connections == 0 or limits.connections > 255) @compileError("9P runner needs 1..255 connection slots");
+ if (limits.listeners == 0) @compileError("9P runner needs at least one listener slot");
+ return struct {
+ const Self = @This();
+
+ pub const Engine = fs.Server(Backend, opts);
+ pub const msize = limits.msize;
+ pub const max_connections = limits.connections;
+
+ /// How the runner reaches the backend. `serve` and `closed` run on
+ /// the connection's task, `opened` on the accepting task; none holds
+ /// the engine lock. Once `stop()` has begun, `serve` must answer
+ /// without waiting on anything the caller of `stop()` would do.
+ pub const Handler = struct {
+ ctx: ?*anyopaque = null,
+ /// One backend request. Answer it through `conn.reply()` now or
+ /// later; a `release` issued by a hangup is not waited for.
+ serve: *const fn (ctx: ?*anyopaque, conn: *Conn, req: Backend.Req) void,
+ /// A connection took the slot (before any request).
+ opened: ?*const fn (ctx: ?*anyopaque, conn: *Conn) void = null,
+ /// The connection hung up; the releases it owed have gone through
+ /// `serve`, and the slot is free once this returns.
+ closed: ?*const fn (ctx: ?*anyopaque, conn: *Conn) void = null,
+ };
+
+ pub const InitOptions = struct {
+ io: Io,
+ /// The engine's root node.
+ root: u64,
+ handler: Handler,
+ /// Salts the engines' fid indexes (`fs.Options.fid_index`).
+ seed: u32 = 0,
+ /// Milliseconds a connection may stay silent before its Tversion;
+ /// zero waits forever.
+ greet_timeout_ms: u32 = 0,
+ };
+
+ pub const ListenError = error{TooManyListeners} || Io.net.IpAddress.ListenError || Io.net.UnixAddress.ListenError || Io.net.UnixAddress.InitError || Io.ConcurrentError;
+
+ io: Io,
+ root: u64,
+ seed: u32,
+ handler: Handler,
+ greet_timeout_ms: u32,
+ conns: [limits.connections]Conn,
+ listeners: [limits.listeners]Io.net.Server,
+ nlisteners: usize,
+ /// Accept tasks and connection tasks; `stop()` cancels it.
+ group: Io.Group,
+ /// Guards `Conn.used`.
+ slots: Io.Mutex,
+ stopping: std.atomic.Value(bool),
+ active: std.atomic.Value(u32),
+
+ const input: u8 = 1;
+ const retry: u8 = 2;
+ const output: u8 = 4;
+ const eof: u8 = 8;
+ const drop: u8 = 16;
+
+ /// One connection slot. Everything but `user` belongs to the runner;
+ /// `engine` may be driven directly between `lock()` and `unlock()`,
+ /// after which `flush()` sends what that produced.
+ pub const Conn = struct {
+ runner: *Self,
+ index: u8,
+ /// The application's; the runner never touches it.
+ user: ?*anyopaque = null,
+ engine: Engine = undefined,
+ stream: Io.net.Stream = undefined,
+ mutex: Io.Mutex = .init,
+ flags: std.atomic.Value(u8) = .init(0),
+ signal: Io.Event = .unset,
+ room: Io.Event = .unset,
+ used: std.atomic.Value(bool) = .init(false),
+ in: [msize]u8 = undefined,
+ out: [2 * msize]u8 = undefined,
+ stage: [msize]u8 = undefined,
+ rbuf: [msize]u8 = undefined,
+ wbuf: [2 * msize]u8 = undefined,
+
+ /// Whether a connection holds the slot right now.
+ pub fn live(c: *const Conn) bool {
+ return c.used.load(.acquire);
+ }
+
+ /// Takes the engine; not a cancelation point.
+ pub fn lock(c: *Conn) void {
+ c.mutex.lockUncancelable(c.runner.io);
+ }
+
+ pub fn unlock(c: *Conn) void {
+ c.mutex.unlock(c.runner.io);
+ }
+
+ /// Answers a request and has the connection task carry on.
+ pub fn reply(c: *Conn, r: *const Backend.Reply, bytes: []const u8) void {
+ c.lock();
+ c.engine.reply(r, bytes);
+ c.unlock();
+ c.raise(output);
+ }
+
+ /// Sends whatever the engine has produced; for a caller that
+ /// drove the engine itself under `lock()`.
+ pub fn flush(c: *Conn) void {
+ c.raise(output);
+ }
+
+ /// Has the connection task retry every parked request.
+ pub fn wake(c: *Conn) void {
+ c.raise(retry);
+ }
+
+ /// Hangs the connection up from anywhere: the slot is free once
+ /// the backend has been paid its releases.
+ pub fn close(c: *Conn) void {
+ c.raise(drop);
+ }
+
+ fn raise(c: *Conn, bits: u8) void {
+ _ = c.flags.fetchOr(bits, .release);
+ c.signal.set(c.runner.io);
+ }
+
+ fn start(c: *Conn, stream: Io.net.Stream) void {
+ const r = c.runner;
+ c.stream = stream;
+ c.user = null;
+ c.flags = .init(0);
+ c.signal = .unset;
+ c.room = .unset;
+ // Under the lock: the slot is already visible as live.
+ c.lock();
+ c.engine = .init(.{ .in = &c.in, .out = &c.out, .root = r.root, .seed = r.seed ^ c.index });
+ c.unlock();
+ if (r.handler.opened) |f| f(r.handler.ctx, c);
+ }
+
+ fn run(c: *Conn) void {
+ const r = c.runner;
+ const io = r.io;
+ var reader = c.stream.reader(io, &c.rbuf);
+ var writer = c.stream.writer(io, &c.wbuf);
+ var readers: Io.Group = .init;
+ if (readers.concurrent(io, readLoop, .{ c, &reader.interface })) |_| {
+ c.serveLoop(&writer.interface);
+ } else |_| {}
+ readers.cancel(io);
+ c.stream.close(io);
+ c.finish();
+ r.release(c);
+ }
+
+ /// Reads whole frames and pushes them into the engine; waits for
+ /// room when a job holds the input.
+ fn readLoop(c: *Conn, reader: *Io.Reader) void {
+ const io = c.runner.io;
+ outer: while (true) {
+ const frame = transport.readFrame(reader, &c.stage, msize) catch break;
+ var off: usize = 0;
+ while (off < frame.len) {
+ c.lock();
+ const n = c.engine.push(frame[off..]);
+ const dead = c.engine.protocol.dead;
+ c.unlock();
+ if (dead) break :outer;
+ off += n;
+ if (n != 0) c.raise(input);
+ if (off < frame.len) {
+ c.room.wait(io) catch break :outer;
+ c.room.reset();
+ }
+ }
+ }
+ c.raise(eof);
+ }
+
+ fn serve(c: *Conn, req: Backend.Req) void {
+ const h = c.runner.handler;
+ h.serve(h.ctx, c, req);
+ }
+
+ /// Steps the engine on every signal, writes its output outside
+ /// the lock, and ends on EOF, a protocol error, `close()`, the
+ /// greet timeout or `stop()`.
+ fn serveLoop(c: *Conn, writer: *Io.Writer) void {
+ const r = c.runner;
+ const io = r.io;
+ const deadline: ?Io.Clock.Timestamp = if (r.greet_timeout_ms == 0) null else Io.Clock.Timestamp.now(io, .awake).addDuration(.{ .raw = .fromMilliseconds(r.greet_timeout_ms), .clock = .awake });
+ var greeted = false;
+ var flags: u8 = 0;
+ while (true) {
+ if (r.stopping.load(.acquire)) return;
+ c.lock();
+ if (flags & retry != 0) {
+ while (c.engine.retry()) |req| {
+ c.unlock();
+ c.serve(req);
+ if (r.stopping.load(.acquire)) return;
+ c.lock();
+ }
+ }
+ while (c.engine.next()) |req| {
+ c.unlock();
+ c.serve(req);
+ if (r.stopping.load(.acquire)) return;
+ c.lock();
+ }
+ const pending = c.engine.output();
+ // The writer's buffer holds a whole output buffer, so
+ // this copies and never blocks under the lock.
+ writer.writeAll(pending) catch {
+ c.unlock();
+ return;
+ };
+ c.engine.wrote(pending.len);
+ const dead = c.engine.protocol.dead;
+ greeted = greeted or c.engine.protocol.msize != 0;
+ c.unlock();
+ c.room.set(io);
+ writer.flush() catch return;
+ if (dead or flags & (eof | drop) != 0) return;
+ if (deadline != null and !greeted) {
+ c.signal.waitTimeout(io, .{ .deadline = deadline.? }) catch |err| switch (err) {
+ error.Timeout => if (c.flags.load(.acquire) & input == 0) return,
+ error.Canceled => return,
+ };
+ } else c.signal.wait(io) catch return;
+ c.signal.reset();
+ flags = c.flags.swap(0, .acquire);
+ }
+ }
+
+ /// Hangs the engine up and pays the backend the releases it owes.
+ fn finish(c: *Conn) void {
+ const r = c.runner;
+ c.lock();
+ c.engine.hangup();
+ c.unlock();
+ while (true) {
+ c.lock();
+ const req = c.engine.next();
+ c.unlock();
+ c.serve(req orelse break);
+ }
+ if (r.handler.closed) |f| f(r.handler.ctx, c);
+ }
+ };
+
+ /// Prepares `r` in place (it is large: every slot carries its
+ /// buffers). Nothing listens until `listen()`.
+ pub fn init(r: *Self, o: InitOptions) void {
+ r.io = o.io;
+ r.root = o.root;
+ r.seed = o.seed;
+ r.handler = o.handler;
+ r.greet_timeout_ms = o.greet_timeout_ms;
+ r.nlisteners = 0;
+ r.group = .init;
+ r.slots = .init;
+ r.stopping = .init(false);
+ r.active = .init(0);
+ for (&r.conns, 0..) |*c, i| c.* = .{ .runner = r, .index = @intCast(i) };
+ }
+
+ /// Binds `address` and starts accepting on it. Returns the bound
+ /// address, which names the port a TCP listener on port zero got.
+ /// Unix paths are not unlinked, chmodded or removed: that is the
+ /// caller's policy, before and after.
+ pub fn listen(r: *Self, address: transport.Address, backlog: u31) ListenError!Io.net.IpAddress {
+ if (r.nlisteners == limits.listeners) return error.TooManyListeners;
+ if (r.stopping.load(.acquire)) return error.TooManyListeners;
+ const i = r.nlisteners;
+ r.listeners[i] = try transport.listen(r.io, address, backlog);
+ errdefer r.listeners[i].deinit(r.io);
+ r.nlisteners += 1;
+ r.group.concurrent(r.io, acceptLoop, .{ r, i }) catch |err| {
+ r.nlisteners -= 1;
+ return err;
+ };
+ return r.listeners[i].socket.address;
+ }
+
+ /// Connections held right now.
+ pub fn count(r: *const Self) usize {
+ return r.active.load(.acquire);
+ }
+
+ /// `Conn.wake()` on every connection.
+ pub fn wakeAll(r: *Self) void {
+ for (&r.conns) |*c| if (c.live()) c.wake();
+ }
+
+ /// `Conn.close()` on every connection.
+ pub fn closeAll(r: *Self) void {
+ for (&r.conns) |*c| if (c.live()) c.close();
+ }
+
+ /// Stops accepting, hangs every connection up, waits for their tasks
+ /// and closes the listeners. Idempotent; the runner is spent after.
+ pub fn stop(r: *Self) void {
+ if (r.stopping.swap(true, .acq_rel)) return;
+ r.group.cancel(r.io);
+ for (r.listeners[0..r.nlisteners]) |*l| l.deinit(r.io);
+ r.nlisteners = 0;
+ }
+
+ fn acceptLoop(r: *Self, i: usize) void {
+ const io = r.io;
+ while (!r.stopping.load(.acquire)) {
+ const stream = r.listeners[i].accept(io) catch |err| switch (err) {
+ error.Canceled => return,
+ else => {
+ // Out of descriptors or memory: try again later.
+ io.sleep(.fromMilliseconds(50), .awake) catch return;
+ continue;
+ },
+ };
+ const c = r.take() orelse {
+ stream.close(io);
+ continue;
+ };
+ c.start(stream);
+ r.group.concurrent(io, Conn.run, .{c}) catch {
+ stream.close(io);
+ r.release(c);
+ };
+ }
+ }
+
+ fn take(r: *Self) ?*Conn {
+ r.slots.lockUncancelable(r.io);
+ defer r.slots.unlock(r.io);
+ for (&r.conns) |*c| {
+ if (c.live()) continue;
+ c.used.store(true, .release);
+ _ = r.active.fetchAdd(1, .acq_rel);
+ return c;
+ }
+ return null;
+ }
+
+ fn release(r: *Self, c: *Conn) void {
+ r.slots.lockUncancelable(r.io);
+ defer r.slots.unlock(r.io);
+ c.used.store(false, .release);
+ _ = r.active.fetchSub(1, .acq_rel);
+ }
+ };
+}
+
+// ---- tests ----
+
+const testing = std.testing;
+const c9 = @import("root.zig");
+const Client = c9.Client;
+
+/// A read-only tree in the style of the engine's StubFs: `/index` is a
+/// file, `/event` parks its reads until `post()` and `/dir` is a directory.
+const Stub = struct {
+ pub const Req = fs.Req;
+ pub const Reply = fs.Reply;
+
+ const root_node = 1;
+ const index_node = 2;
+ const event_node = 3;
+ const dir_node = 4;
+
+ mutex: Io.Mutex = .init,
+ event: ?[]const u8 = null,
+ calls: u32 = 0,
+ releases: u32 = 0,
+ parks: u32 = 0,
+ opened: u32 = 0,
+ closed: u32 = 0,
+
+ const Answer = struct { reply: fs.Reply, bytes: []const u8 = "" };
+
+ fn attrOf(node: u64) fs.Attr {
+ return switch (node) {
+ root_node => .{ .name = "/", .node = root_node, .dir = true, .mode = 0o500 },
+ index_node => .{ .name = "index", .node = index_node, .mode = 0o400, .size = 12 },
+ event_node => .{ .name = "event", .node = event_node, .mode = 0o400 },
+ dir_node => .{ .name = "dir", .node = dir_node, .dir = true, .mode = 0o500 },
+ else => unreachable,
+ };
+ }
+
+ fn handle(st: *Stub, req: fs.Req) Answer {
+ st.mutex.lockUncancelable(testing.io);
+ defer st.mutex.unlock(testing.io);
+ st.calls += 1;
+ const fail: Answer = .{ .reply = .fail(req.tag, fs.E.NOENT) };
+ switch (req.op) {
+ .lookup => {
+ if (std.mem.eql(u8, req.data, "..")) return .{ .reply = .{ .tag = req.tag, .attr = attrOf(root_node) } };
+ if (req.node != root_node) return fail;
+ inline for (.{ "index", "event", "dir" }, .{ index_node, event_node, dir_node }) |name, node| {
+ if (std.mem.eql(u8, req.data, name)) return .{ .reply = .{ .tag = req.tag, .attr = attrOf(node) } };
+ }
+ return fail;
+ },
+ .getattr, .setattr => return .{ .reply = .{ .tag = req.tag, .attr = attrOf(req.node) } },
+ .open => return .{ .reply = .{ .tag = req.tag, .handle = 7 } },
+ .release => {
+ st.releases += 1;
+ return .{ .reply = .{ .tag = req.tag } };
+ },
+ .readdir => return .{ .reply = .{ .tag = req.tag } },
+ .read => {
+ if (req.node == event_node) {
+ const rec = st.event orelse {
+ st.parks += 1;
+ return .{ .reply = .{ .tag = req.tag, .status = .again } };
+ };
+ st.event = null;
+ return .{ .reply = .{ .tag = req.tag }, .bytes = rec };
+ }
+ const all = "hello, index";
+ if (req.off >= all.len) return .{ .reply = .{ .tag = req.tag } };
+ const from = all[@intCast(req.off)..];
+ return .{ .reply = .{ .tag = req.tag }, .bytes = from[0..@min(from.len, req.size)] };
+ },
+ .write => return .{ .reply = .fail(req.tag, fs.E.PERM) },
+ }
+ }
+
+ fn post(st: *Stub, rec: []const u8) void {
+ st.mutex.lockUncancelable(testing.io);
+ defer st.mutex.unlock(testing.io);
+ st.event = rec;
+ }
+
+ fn of(ctx: ?*anyopaque) *Stub {
+ return @ptrCast(@alignCast(ctx.?));
+ }
+};
+
+const TestRunner = Runner(Stub, .{ .fid_capacity = 8, .slot_capacity = 4 }, .{ .msize = 4096, .connections = 2, .listeners = 2 });
+
+/// Answers at once on the connection's task.
+fn serveNow(ctx: ?*anyopaque, conn: *TestRunner.Conn, req: fs.Req) void {
+ const a = Stub.of(ctx).handle(req);
+ conn.reply(&a.reply, a.bytes);
+}
+
+fn countOpened(ctx: ?*anyopaque, conn: *TestRunner.Conn) void {
+ _ = conn;
+ const st = Stub.of(ctx);
+ st.mutex.lockUncancelable(testing.io);
+ defer st.mutex.unlock(testing.io);
+ st.opened += 1;
+}
+
+fn countClosed(ctx: ?*anyopaque, conn: *TestRunner.Conn) void {
+ _ = conn;
+ const st = Stub.of(ctx);
+ st.mutex.lockUncancelable(testing.io);
+ defer st.mutex.unlock(testing.io);
+ st.closed += 1;
+}
+
+/// A `Client` over one `std.Io` stream: submit, ship, read until the reply.
+const TestClient = struct {
+ io: Io,
+ stream: Io.net.Stream,
+ cli: Client = undefined,
+ in: [4096]u8 = undefined,
+ out: [4096]u8 = undefined,
+ rbuf: [4096]u8 = undefined,
+ wbuf: [4096]u8 = undefined,
+ stage: [4096]u8 = undefined,
+ reader: Io.net.Stream.Reader = undefined,
+ writer: Io.net.Stream.Writer = undefined,
+
+ fn open(tc: *TestClient, io: Io, address: transport.Address) !void {
+ tc.io = io;
+ tc.stream = try transport.connect(io, address);
+ tc.cli = .init(.{ .in = &tc.in, .out = &tc.out });
+ tc.reader = tc.stream.reader(io, &tc.rbuf);
+ tc.writer = tc.stream.writer(io, &tc.wbuf);
+ }
+
+ fn close(tc: *TestClient) void {
+ tc.stream.close(tc.io);
+ }
+
+ fn ship(tc: *TestClient) !void {
+ const bytes = tc.cli.output();
+ try tc.writer.interface.writeAll(bytes);
+ try tc.writer.interface.flush();
+ tc.cli.wrote(bytes.len);
+ }
+
+ fn submit(tc: *TestClient, req: Client.Request) !u16 {
+ const tag = try tc.cli.submit(req);
+ try tc.ship();
+ return tag;
+ }
+
+ fn reap(tc: *TestClient) !Client.Done {
+ while (true) {
+ if (tc.cli.take()) |done| return done;
+ if (tc.cli.dead) return error.Botch;
+ const frame = try transport.readFrame(&tc.reader.interface, &tc.stage, 4096);
+ try testing.expectEqual(frame.len, tc.cli.push(frame));
+ }
+ }
+
+ fn one(tc: *TestClient, req: Client.Request) !Client.Done {
+ const tag = try tc.submit(req);
+ const done = try tc.reap();
+ try testing.expectEqual(tag, done.tag);
+ try testing.expectEqual(std.meta.activeTag(req), done.op);
+ return done;
+ }
+
+ fn handshake(tc: *TestClient) !void {
+ const v = try tc.one(.{ .version = .{} });
+ try testing.expectEqualStrings("9P2000", v.result.version.version);
+ const a = try tc.one(.{ .attach = .{ .fid = 0, .uname = "goblin" } });
+ try testing.expectEqual(@as(u64, Stub.root_node), a.result.attach.path);
+ }
+
+ fn readIndex(tc: *TestClient, fid: u32) !void {
+ const w = try tc.one(.{ .walk = .{ .fid = 0, .newfid = fid, .names = &.{"index"} } });
+ try testing.expectEqual(@as(u16, 1), w.result.walk.nwqid);
+ try testing.expectEqual(@as(u64, Stub.index_node), w.result.walk.wqid[0].path);
+ _ = try tc.one(.{ .open = .{ .fid = fid, .mode = c9.oread } });
+ const r = try tc.one(.{ .read = .{ .fid = fid, .offset = 0, .count = 64 } });
+ try testing.expectEqualStrings("hello, index", r.result.read);
+ _ = try tc.one(.{ .clunk = .{ .fid = fid } });
+ }
+};
+
+/// The peer hung up: EOF, or the reset/EPIPE a closed socket answers with.
+fn expectHangup(result: anytype) !void {
+ if (result) |_| return error.TestUnexpectedResult else |err| {
+ const name = @errorName(err);
+ for ([_][]const u8{ "EndOfStream", "ReadFailed", "WriteFailed" }) |ok| if (std.mem.eql(u8, name, ok)) return;
+ return err;
+ }
+}
+
+const Rig = struct {
+ dir: testing.TmpDir,
+ path_buffer: [std.fs.max_path_bytes]u8 = undefined,
+ unix_buffer: [transport.sun_path_len]u8 = undefined,
+ unix: [:0]const u8 = undefined,
+ stub: Stub = .{},
+ runner: TestRunner = undefined,
+
+ fn start(rig: *Rig, handler: TestRunner.Handler, greet_ms: u32) !void {
+ const io = testing.io;
+ rig.dir = testing.tmpDir(.{});
+ errdefer rig.dir.cleanup();
+ const parent_len = try rig.dir.dir.realPath(io, &rig.path_buffer);
+ rig.unix = try std.fmt.bufPrintSentinel(&rig.unix_buffer, "{s}/9p", .{rig.path_buffer[0..parent_len]}, 0);
+ var h = handler;
+ h.ctx = &rig.stub;
+ rig.runner.init(.{ .io = io, .root = Stub.root_node, .handler = h, .greet_timeout_ms = greet_ms });
+ _ = try rig.runner.listen(.{ .unix = rig.unix }, 4);
+ }
+
+ fn end(rig: *Rig) void {
+ rig.runner.stop();
+ rig.dir.cleanup();
+ }
+
+ /// Waits until the runner holds `n` connections, briefly.
+ fn settle(rig: *Rig, n: usize) !void {
+ var tries: usize = 0;
+ while (rig.runner.count() != n) : (tries += 1) {
+ if (tries == 2000) return error.Timeout;
+ try testing.io.sleep(.fromMilliseconds(1), .awake);
+ }
+ }
+};
+
+test "serve: version, attach, walk and read over a Unix socket and TCP" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start(.{ .serve = serveNow, .opened = countOpened, .closed = countClosed }, 0);
+ defer rig.end();
+ const tcp = try rig.runner.listen(.{ .tcp = .{ .ip4 = .loopback(0) } }, 4);
+ try testing.expect(tcp.getPort() != 0);
+ try testing.expectError(error.TooManyListeners, rig.runner.listen(.{ .tcp = .{ .ip4 = .loopback(0) } }, 4));
+ for ([_]transport.Address{ .{ .unix = rig.unix }, .{ .tcp = tcp } }) |address| {
+ var tc: TestClient = .{ .io = undefined, .stream = undefined };
+ try tc.open(testing.io, address);
+ defer tc.close();
+ try tc.handshake();
+ try tc.readIndex(1);
+ const st = try tc.one(.{ .stat = .{ .fid = 0 } });
+ try testing.expectEqualStrings("/", st.result.stat.name);
+ try testing.expectEqualStrings("goblin", st.result.stat.uid);
+ const missing = try tc.one(.{ .walk = .{ .fid = 0, .newfid = 2, .names = &.{"nowhere"} } });
+ try testing.expectEqualStrings(fs.errString(fs.E.NOENT), missing.result.fail);
+ }
+ try rig.settle(0);
+ try testing.expectEqual(@as(u32, 2), rig.stub.opened);
+ try testing.expectEqual(@as(u32, 2), rig.stub.closed);
+ try testing.expectEqual(@as(u32, 2), rig.stub.releases);
+}
+
+test "serve: a parked read completes when another task posts and wakes" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start(.{ .serve = serveNow }, 0);
+ defer rig.end();
+ var tc: TestClient = .{ .io = undefined, .stream = undefined };
+ try tc.open(testing.io, .{ .unix = rig.unix });
+ defer tc.close();
+ try tc.handshake();
+ _ = try tc.one(.{ .walk = .{ .fid = 0, .newfid = 1, .names = &.{"event"} } });
+ _ = try tc.one(.{ .open = .{ .fid = 1, .mode = c9.oread } });
+ const tag = try tc.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 64 } });
+ // The read parks; nothing wakes it until the poster does.
+ var tries: usize = 0;
+ while (rig.stub.parks == 0) : (tries += 1) {
+ if (tries == 2000) return error.Timeout;
+ try testing.io.sleep(.fromMilliseconds(1), .awake);
+ }
+ try testing.io.sleep(.fromMilliseconds(20), .awake);
+ try testing.expectEqual(@as(u32, 1), rig.stub.parks);
+ const Poster = struct {
+ fn post(r: *Rig) void {
+ testing.io.sleep(.fromMilliseconds(10), .awake) catch return;
+ r.stub.post("record 1\n");
+ r.runner.wakeAll();
+ }
+ };
+ var poster = try testing.io.concurrent(Poster.post, .{&rig});
+ defer poster.cancel(testing.io);
+ const done = try tc.reap();
+ try testing.expectEqual(tag, done.tag);
+ try testing.expectEqualStrings("record 1\n", done.result.read);
+ poster.await(testing.io);
+ try testing.expect(rig.stub.parks >= 1);
+ // A second read parks and is flushed through the engine: EINTR, then Rflush.
+ const again = try tc.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 64 } });
+ tries = 0;
+ while (rig.stub.parks < 2) : (tries += 1) {
+ if (tries == 2000) return error.Timeout;
+ try testing.io.sleep(.fromMilliseconds(1), .awake);
+ }
+ const flush_tag = try tc.submit(.{ .flush = .{ .oldtag = again } });
+ const first = try tc.reap();
+ try testing.expectEqual(again, first.tag);
+ try testing.expectEqualStrings(fs.e_interrupted, first.result.fail);
+ const second = try tc.reap();
+ try testing.expectEqual(flush_tag, second.tag);
+ try testing.expectEqual(Client.Op.flush, second.op);
+ _ = try tc.one(.{ .clunk = .{ .fid = 1 } });
+}
+
+test "serve: two clients at once, a third is refused, and stop() hangs up" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start(.{ .serve = serveNow, .closed = countClosed }, 0);
+ defer rig.end();
+ const Worker = struct {
+ fn run(r: *Rig, fid: u32, failure: *?anyerror) void {
+ check(r, fid) catch |err| {
+ failure.* = err;
+ };
+ }
+ fn check(r: *Rig, fid: u32) !void {
+ var tc: TestClient = .{ .io = undefined, .stream = undefined };
+ try tc.open(testing.io, .{ .unix = r.unix });
+ defer tc.close();
+ try tc.handshake();
+ for (0..8) |_| try tc.readIndex(fid);
+ }
+ };
+ var failures: [2]?anyerror = .{ null, null };
+ var group: Io.Group = .init;
+ try group.concurrent(testing.io, Worker.run, .{ &rig, 1, &failures[0] });
+ try group.concurrent(testing.io, Worker.run, .{ &rig, 2, &failures[1] });
+ try group.await(testing.io);
+ for (failures) |f| if (f) |err| return err;
+ try rig.settle(0);
+ try testing.expectEqual(@as(u32, 16), rig.stub.releases);
+ try testing.expectEqual(@as(u32, 2), rig.stub.closed);
+
+ // Two slots: the third connection is closed without a word.
+ var a: TestClient = .{ .io = undefined, .stream = undefined };
+ try a.open(testing.io, .{ .unix = rig.unix });
+ defer a.close();
+ try a.handshake();
+ var b: TestClient = .{ .io = undefined, .stream = undefined };
+ try b.open(testing.io, .{ .unix = rig.unix });
+ defer b.close();
+ try b.handshake();
+ try rig.settle(2);
+ var c: TestClient = .{ .io = undefined, .stream = undefined };
+ try c.open(testing.io, .{ .unix = rig.unix });
+ defer c.close();
+ try expectHangup(c.one(.{ .version = .{} }));
+ try rig.settle(2);
+ // a holds an open fid: stop() pays its release and ends the connection.
+ _ = try a.one(.{ .walk = .{ .fid = 0, .newfid = 3, .names = &.{"index"} } });
+ _ = try a.one(.{ .open = .{ .fid = 3, .mode = c9.oread } });
+ const releases = rig.stub.releases;
+ rig.runner.stop();
+ try testing.expectEqual(@as(usize, 0), rig.runner.count());
+ try testing.expectEqual(releases + 1, rig.stub.releases);
+ try testing.expectEqual(@as(u32, 4), rig.stub.closed);
+ try expectHangup(a.one(.{ .clunk = .{ .fid = 3 } }));
+ try expectHangup(b.one(.{ .stat = .{ .fid = 0 } }));
+ rig.runner.stop();
+}
+
+test "serve: a silent connection is dropped at the greet timeout, Conn.close() drops one at will" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start(.{ .serve = serveNow }, 30);
+ defer rig.end();
+ var silent: TestClient = .{ .io = undefined, .stream = undefined };
+ try silent.open(testing.io, .{ .unix = rig.unix });
+ defer silent.close();
+ try rig.settle(1);
+ try expectHangup(transport.readFrame(&silent.reader.interface, &silent.stage, 4096));
+ try rig.settle(0);
+ var tc: TestClient = .{ .io = undefined, .stream = undefined };
+ try tc.open(testing.io, .{ .unix = rig.unix });
+ defer tc.close();
+ try tc.handshake();
+ try testing.io.sleep(.fromMilliseconds(40), .awake);
+ try tc.readIndex(1);
+ try rig.settle(1);
+ rig.runner.closeAll();
+ try rig.settle(0);
+ try expectHangup(tc.one(.{ .stat = .{ .fid = 0 } }));
+}
+
+/// A handler in the editor's style: the request is handed to another
+/// thread, which drives the engine itself under `lock()`.
+const Deferred = struct {
+ stub: *Stub,
+ conn: ?*TestRunner.Conn = null,
+ req: ?fs.Req = null,
+ served: u32 = 0,
+ mutex: Io.Mutex = .init,
+ posted: Io.Event = .unset,
+ answered: Io.Event = .unset,
+ stopping: bool = false,
+
+ fn serve(ctx: ?*anyopaque, conn: *TestRunner.Conn, req: fs.Req) void {
+ const d: *Deferred = @ptrCast(@alignCast(ctx.?));
+ d.mutex.lockUncancelable(testing.io);
+ if (d.stopping) {
+ d.mutex.unlock(testing.io);
+ const a = d.stub.handle(req);
+ conn.reply(&a.reply, a.bytes);
+ return;
+ }
+ d.conn = conn;
+ d.req = req;
+ d.answered.reset();
+ d.mutex.unlock(testing.io);
+ d.posted.set(testing.io);
+ d.answered.wait(testing.io) catch {};
+ }
+
+ /// The other thread: answers the posted request and everything the
+ /// engine asks after it, then lets the connection task write.
+ fn drive(d: *Deferred) !void {
+ while (true) {
+ try d.posted.wait(testing.io);
+ d.posted.reset();
+ d.mutex.lockUncancelable(testing.io);
+ const conn = d.conn.?;
+ const req = d.req.?;
+ d.req = null;
+ d.mutex.unlock(testing.io);
+ conn.lock();
+ var a = d.stub.handle(req);
+ conn.engine.reply(&a.reply, a.bytes);
+ d.served += 1;
+ while (conn.engine.next()) |more| {
+ a = d.stub.handle(more);
+ conn.engine.reply(&a.reply, a.bytes);
+ d.served += 1;
+ }
+ conn.unlock();
+ conn.flush();
+ d.answered.set(testing.io);
+ }
+ }
+};
+
+test "serve: a backend answered from another thread under lock()" {
+ var rig: Rig = .{ .dir = undefined };
+ var deferred: Deferred = .{ .stub = &rig.stub };
+ try rig.start(.{ .serve = Deferred.serve, .ctx = &deferred }, 0);
+ defer rig.end();
+ rig.runner.handler.ctx = &deferred;
+ var driver = try testing.io.concurrent(Deferred.drive, .{&deferred});
+ defer driver.cancel(testing.io) catch {};
+ var tc: TestClient = .{ .io = undefined, .stream = undefined };
+ try tc.open(testing.io, .{ .unix = rig.unix });
+ defer tc.close();
+ try tc.handshake();
+ try tc.readIndex(1);
+ const w = try tc.one(.{ .walk = .{ .fid = 0, .newfid = 2, .names = &.{ "dir", ".." } } });
+ try testing.expectEqual(@as(u16, 2), w.result.walk.nwqid);
+ // attach getattr, walk lookup, open, read, release, two more lookups.
+ try testing.expectEqual(@as(u32, 7), deferred.served);
+ deferred.mutex.lockUncancelable(testing.io);
+ deferred.stopping = true;
+ deferred.mutex.unlock(testing.io);
+ rig.runner.stop();
+ try expectHangup(tc.one(.{ .stat = .{ .fid = 0 } }));
+}