diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-20 02:33:02 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-09-20 02:33:02 -0300 |
| commit | ba7ec40782ba7020d82a56896a5eb1b52578d6aa (patch) | |
| tree | b6a5222d124255ae1e5d568c6ffbfe2e2050ae63 /src/serve.zig | |
| parent | 66f2e492c348677ab3050f5e378b9eb4c04c98ce (diff) | |
| download | cloud9-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.zig | 848 |
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 } })); +} |
