//! 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"); const post = @import("post.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; pub const ListenPostedError = post.PostError || error{TooManyListeners} || 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, /// The registry socket of a `listenPosted`, zero-terminated; /// `stop()` unlinks it (unpost on stop) — but only while it is /// still this runner's entry (`posted_ino`). posted_path: [transport.sun_path_len + 1]u8 = @splat(0), posted_len: usize = 0, /// The bound registry entry's inode, the ownership proof for /// the unpost in `stop()`. posted_ino: u64 = 0, /// 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.initIn(.{ .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.posted_len = 0; r.nlisteners = 0; r.group = .init; r.slots = .init; r.stopping = .init(false); r.active = .init(0); // Field by field: a connection's engine and buffers are written // when a client takes the slot (`start`), so a runner sized for // many connections costs nothing for the ones never used. for (&r.conns, 0..) |*c, i| { c.runner = r; c.index = @intCast(i); c.user = null; c.mutex = .init; c.flags = .init(0); c.signal = .unset; c.room = .unset; c.used = .init(false); } } /// 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; var server = try transport.listen(r.io, address, backlog); errdefer server.deinit(r.io); try r.startListener(server); return server.socket.address; } /// Posts the runner on the registry socket /// `$XDG_RUNTIME_DIR/9p/` and starts accepting on it: /// `post.post` creates the 0o750 registry directory and runs the /// stale protocol (a refused entry is replaced; a live server /// owning the name is `AlreadyPosted`; a non-socket entry is /// never deleted). One posted name per runner: a second /// `listenPosted` is `AlreadyPosted` (its socket would otherwise /// be orphaned in the registry — nothing would unpost it). /// `stop()` unposts — the socket is unlinked when the runner /// stops, as long as the entry is still the runner's own. pub fn listenPosted(r: *Self, env: post.Env, name: []const u8, backlog: u31) ListenPostedError!void { if (r.nlisteners == limits.listeners) return error.TooManyListeners; if (r.stopping.load(.acquire)) return error.TooManyListeners; if (r.posted_len != 0) return error.AlreadyPosted; var p = try post.post(r.io, env, name, backlog, &r.posted_path); errdefer { post.unpost(r.io, p.path, p.inode); p.server.deinit(r.io); } try r.startListener(p.server); r.posted_len = p.path.len; r.posted_ino = p.inode; } /// Registers a bound listener and starts its accept task. fn startListener(r: *Self, server: Io.net.Server) Io.ConcurrentError!void { const i = r.nlisteners; r.listeners[i] = server; 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; }; } /// 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. A posted listener is unposted /// — its registry socket is unlinked, but only while the entry /// is still the runner's own: a name that was re-posted by /// another server (this one's socket file having been lost) /// survives the stop. 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; if (r.posted_len != 0) { post.unpost(r.io, r.posted_path[0..r.posted_len :0], r.posted_ino); r.posted_len = 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; const linux = std.os.linux; /// 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 } })); } test "serve: listenPosted serves the registry name and stop() unposts" { if (@import("builtin").os.tag != .linux) return error.SkipZigTest; var rig: Rig = .{ .dir = undefined }; const io = testing.io; // A scratch registry: XDG_RUNTIME_DIR is the rig's own temp dir. rig.dir = testing.tmpDir(.{}); errdefer rig.dir.cleanup(); var real_buf: [std.fs.max_path_bytes]u8 = undefined; const len = try rig.dir.dir.realPath(io, &real_buf); var env_buf: [std.fs.max_path_bytes]u8 = undefined; const value = try std.fmt.bufPrintZ(&env_buf, "XDG_RUNTIME_DIR={s}", .{real_buf[0..len]}); const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; const envp: post.Env = @ptrCast(&env); rig.stub = .{}; rig.runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &rig.stub, .serve = serveNow } }); try rig.runner.listenPosted(envp, "posted", 4); var path_buf: [transport.sun_path_len]u8 = undefined; const path = try post.registryPath(envp, "posted", &path_buf); // The name is posted, live, and serves a full 9P session. try testing.expect(post.probe(path) == .live); var tc: TestClient = .{ .io = undefined, .stream = undefined }; try tc.open(io, .{ .unix = path }); defer tc.close(); try tc.handshake(); try tc.readIndex(1); // A second post of the same name is refused while the runner lives. var pbuf: [transport.sun_path_len]u8 = undefined; try testing.expectError(error.AlreadyPosted, post.post(io, envp, "posted", 4, &pbuf)); // The listing sees it; stop() unposts and the entry disappears. var stage: [512]u8 = undefined; var names = try post.posted(io, envp, &stage); var seen = false; while (names.next()) |n| seen = seen or std.mem.eql(u8, n, "posted"); try testing.expect(seen); rig.runner.stop(); rig.runner.stop(); // idempotent: the second stop unposts nothing more try testing.expectError(error.FileNotFound, Io.Dir.statFile(.cwd(), io, path, .{})); names = try post.posted(io, envp, &stage); try testing.expect(names.next() == null); rig.dir.cleanup(); } test "serve: a second listenPosted is refused; stop() never unposts another's name" { if (@import("builtin").os.tag != .linux) return error.SkipZigTest; var rig: Rig = .{ .dir = undefined }; const io = testing.io; rig.dir = testing.tmpDir(.{}); errdefer rig.dir.cleanup(); var real_buf: [std.fs.max_path_bytes]u8 = undefined; const len = try rig.dir.dir.realPath(io, &real_buf); var env_buf: [std.fs.max_path_bytes]u8 = undefined; const value = try std.fmt.bufPrintZ(&env_buf, "XDG_RUNTIME_DIR={s}", .{real_buf[0..len]}); const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; const envp: post.Env = @ptrCast(&env); rig.stub = .{}; rig.runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &rig.stub, .serve = serveNow } }); try rig.runner.listenPosted(envp, "twice", 4); // One posted name per runner: a second would overwrite the first's // path and orphan its socket in the registry. try testing.expectError(error.AlreadyPosted, rig.runner.listenPosted(envp, "twice", 4)); try testing.expectError(error.AlreadyPosted, rig.runner.listenPosted(envp, "other", 4)); var path_buf: [transport.sun_path_len]u8 = undefined; const path = try post.registryPath(envp, "twice", &path_buf); // The runner's socket file is lost behind its back (rm, crash // cleanup), and another server takes the now-free name. try Io.Dir.deleteFileAbsolute(io, path); var thief = try post.post(io, envp, "twice", 4, &path_buf); try testing.expect(post.probe(thief.path) == .live); // stop() unposts only what it still owns: the thief survives. rig.runner.stop(); const st = try Io.Dir.statFile(.cwd(), io, thief.path, .{}); try testing.expect(st.kind == .unix_domain_socket); try testing.expect(post.probe(thief.path) == .live); post.unpost(io, thief.path, thief.inode); thief.server.deinit(io); rig.dir.cleanup(); } test "serve: listenPosted beside listen(): stop unposts only the registry name" { if (@import("builtin").os.tag != .linux) return error.SkipZigTest; var rig: Rig = .{ .dir = undefined }; const io = testing.io; rig.dir = testing.tmpDir(.{}); errdefer rig.dir.cleanup(); var real_buf: [std.fs.max_path_bytes]u8 = undefined; const len = try rig.dir.dir.realPath(io, &real_buf); var env_buf: [std.fs.max_path_bytes]u8 = undefined; const value = try std.fmt.bufPrintZ(&env_buf, "XDG_RUNTIME_DIR={s}", .{real_buf[0..len]}); const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; const envp: post.Env = @ptrCast(&env); rig.stub = .{}; rig.runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &rig.stub, .serve = serveNow } }); // A plain Unix listener (the application's path policy) beside the // posted name: both listener slots fill. var unix_buf: [std.fs.max_path_bytes]u8 = undefined; const unix_path = try std.fmt.bufPrintZ(&unix_buf, "{s}/plain.sock", .{real_buf[0..len]}); _ = try rig.runner.listen(.{ .unix = unix_path }, 4); try rig.runner.listenPosted(envp, "mixed", 4); var path_buf: [transport.sun_path_len]u8 = undefined; const posted_path = try post.registryPath(envp, "mixed", &path_buf); try testing.expect(post.probe(posted_path) == .live); rig.runner.stop(); // The posted name is unposted; the plain path is the application's. try testing.expectError(error.FileNotFound, Io.Dir.statFile(.cwd(), io, posted_path, .{})); const plain_st = try Io.Dir.statFile(.cwd(), io, unix_path, .{}); try testing.expect(plain_st.kind == .unix_domain_socket); try Io.Dir.deleteFileAbsolute(io, unix_path); rig.dir.cleanup(); } // ---- runner capacity and process hygiene ---- /// Counts the process's open descriptors through /proc, no allocation. fn openFdCount(io: Io) !usize { var rb: [Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined; const dir = try Io.Dir.openDirAbsolute(io, "/proc/self/fd", .{ .iterate = true }); defer Io.Dir.close(dir, io); var it = Io.Dir.Reader.init(dir, &rb); var n: usize = 0; while (try it.next(io)) |_| n += 1; return n; } /// Counts the process's OS threads through /proc, no allocation. fn taskCount(io: Io) !usize { var rb: [Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined; const dir = try Io.Dir.openDirAbsolute(io, "/proc/self/task", .{ .iterate = true }); defer Io.Dir.close(dir, io); var it = Io.Dir.Reader.init(dir, &rb); var n: usize = 0; while (try it.next(io)) |_| n += 1; return n; } fn waitCount(runner: anytype, n: usize) !void { var tries: usize = 0; while (runner.count() != n) : (tries += 1) { if (tries == 5000) return error.Timeout; try testing.io.sleep(.fromMilliseconds(1), .awake); } } const BigRunner = Runner(Stub, .{ .fid_capacity = 8, .slot_capacity = 4 }, .{ .msize = 4096, .connections = 16, .listeners = 1 }); fn serveNowBig(ctx: ?*anyopaque, conn: *BigRunner.Conn, req: fs.Req) void { const a = Stub.of(ctx).handle(req); conn.reply(&a.reply, a.bytes); } fn countClosedBig(ctx: ?*anyopaque, conn: *BigRunner.Conn) void { _ = conn; const st = Stub.of(ctx); st.mutex.lockUncancelable(testing.io); defer st.mutex.unlock(testing.io); st.closed += 1; } test "serve: sixteen connections fill the table, the seventeenth is closed, stop ends all" { const io = testing.io; var dir = testing.tmpDir(.{}); defer dir.cleanup(); var path_buffer: [std.fs.max_path_bytes]u8 = undefined; var unix_buffer: [transport.sun_path_len]u8 = undefined; const plen = try dir.dir.realPath(io, &path_buffer); const unix = try std.fmt.bufPrintSentinel(&unix_buffer, "{s}/9p", .{path_buffer[0..plen]}, 0); var stub: Stub = .{}; var runner: BigRunner = undefined; runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &stub, .serve = serveNowBig, .closed = countClosedBig } }); _ = try runner.listen(.{ .unix = unix }, 16); defer runner.stop(); var clients: [16]TestClient = undefined; for (&clients) |*tc| { tc.* = .{ .io = undefined, .stream = undefined }; try tc.open(io, .{ .unix = unix }); try tc.handshake(); } defer for (&clients) |*tc| tc.close(); try waitCount(&runner, 16); // The seventeenth connection is accepted and closed at once: it never // gets a version reply, and the sixteen keep their slots. var extra: TestClient = .{ .io = undefined, .stream = undefined }; try extra.open(io, .{ .unix = unix }); defer extra.close(); try expectHangup(extra.one(.{ .version = .{} })); try waitCount(&runner, 16); // stop() hangs every live connection up and frees every slot. const closed0 = stub.closed; runner.stop(); try testing.expectEqual(@as(usize, 0), runner.count()); try testing.expectEqual(@as(u32, 16), stub.closed - closed0); for (&clients) |*tc| try expectHangup(tc.one(.{ .stat = .{ .fid = 0 } })); } test "serve: repeated connect/disconnect leaves descriptors and threads flat" { var rig: Rig = .{ .dir = undefined }; try rig.start(.{ .serve = serveNow }, 0); defer rig.end(); const io = testing.io; const Cycle = struct { fn run(r: *Rig) !void { var tc: TestClient = .{ .io = undefined, .stream = undefined }; try tc.open(testing.io, .{ .unix = r.unix }); try tc.handshake(); try tc.readIndex(1); tc.close(); } }; // Warm the task pool to its steady state (the first connections grow // it; later ones must not). for (0..400) |_| try Cycle.run(&rig); try rig.settle(0); const fds0 = try openFdCount(io); const threads0 = try taskCount(io); for (0..1200) |_| try Cycle.run(&rig); try rig.settle(0); // Descriptors are exactly flat: no stream, buffer or address is leaked // per connection. try testing.expectEqual(fds0, try openFdCount(io)); // Threads come from the shared Io worker pool and may lazily add a // worker; they must not scale with the connection count (which would be // a per-connection thread leak). try testing.expect(try taskCount(io) <= threads0 + 2); } test "serve: a connection storm is absorbed and the runner keeps serving" { var rig: Rig = .{ .dir = undefined }; try rig.start(.{ .serve = serveNow }, 0); defer rig.end(); const io = testing.io; const fds0 = try openFdCount(io); const Storm = struct { fn run(r: *Rig, refused: *std.atomic.Value(u32)) void { for (0..32) |_| { var tc: TestClient = .{ .io = undefined, .stream = undefined }; tc.open(testing.io, .{ .unix = r.unix }) catch { _ = refused.fetchAdd(1, .monotonic); continue; }; // No handshake: the connection is torn down as soon as it // is accepted (or refused a slot by the runner). tc.close(); } } }; var group: Io.Group = .init; var refused = std.atomic.Value(u32).init(0); for (0..16) |_| try group.concurrent(io, Storm.run, .{ &rig, &refused }); try group.await(io); try rig.settle(0); // 512 connections, two slots: the extra ones are closed by the runner, // not refused by the kernel backlog. try testing.expectEqual(@as(u32, 0), refused.load(.monotonic)); try testing.expectEqual(fds0, try openFdCount(io)); // The runner still accepts and serves a full session. var tc: TestClient = .{ .io = undefined, .stream = undefined }; try tc.open(io, .{ .unix = rig.unix }); defer tc.close(); try tc.handshake(); try tc.readIndex(1); } /// The open descriptors of this process, from /proc, for the CLOEXEC audit. fn collectFds(io: Io, out: []i32) ![]i32 { var rb: [Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined; const dir = try Io.Dir.openDirAbsolute(io, "/proc/self/fd", .{ .iterate = true }); defer Io.Dir.close(dir, io); var it = Io.Dir.Reader.init(dir, &rb); var n: usize = 0; while (try it.next(io)) |e| { const fd = std.fmt.parseInt(i32, e.name, 10) catch continue; if (n < out.len) { out[n] = fd; n += 1; } } return out[0..n]; } fn hasFd(fds: []const i32, fd: i32) bool { for (fds) |f| if (f == fd) return true; return false; } test "serve: every descriptor opened by the post/serve path is close-on-exec" { if (@import("builtin").os.tag != .linux) return error.SkipZigTest; const io = testing.io; var dir = testing.tmpDir(.{}); defer dir.cleanup(); var real_buf: [std.fs.max_path_bytes]u8 = undefined; const rlen = try dir.dir.realPath(io, &real_buf); var env_buf: [std.fs.max_path_bytes]u8 = undefined; const value = try std.fmt.bufPrintZ(&env_buf, "XDG_RUNTIME_DIR={s}", .{real_buf[0..rlen]}); const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; const envp: post.Env = @ptrCast(&env); var before_buf: [64]i32 = undefined; const before = try collectFds(io, &before_buf); // Post a runner, dial its name, scan the registry and accept a client: // every descriptor this opens must not survive into an exec. var stub: Stub = .{}; var runner: TestRunner = undefined; runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &stub, .serve = serveNow } }); defer runner.stop(); try runner.listenPosted(envp, "cloexec", 4); var path_buf: [transport.sun_path_len]u8 = undefined; const posted_path = try post.registryPath(envp, "cloexec", &path_buf); var dialed = try post.dialPath(io, posted_path); defer dialed.close(io); var stage: [1024]u8 = undefined; var names = try post.posted(io, envp, &stage); try testing.expect(names.next() != null); var tc: TestClient = .{ .io = undefined, .stream = undefined }; try tc.open(io, .{ .unix = posted_path }); defer tc.close(); try tc.handshake(); var after_buf: [64]i32 = undefined; const after = try collectFds(io, &after_buf); for (after) |fd| { if (hasFd(before, fd)) continue; const flags = linux.fcntl(fd, linux.F.GETFD, 0); try testing.expect(flags & linux.FD_CLOEXEC != 0); } // Sanity: the audit saw the new descriptors (the listener, the dialed // socket, the accepted connection). try testing.expect(after.len > before.len); }