summaryrefslogtreecommitdiff
path: root/9player/src/nine.zig
diff options
context:
space:
mode:
Diffstat (limited to '9player/src/nine.zig')
-rw-r--r--9player/src/nine.zig756
1 files changed, 756 insertions, 0 deletions
diff --git a/9player/src/nine.zig b/9player/src/nine.zig
new file mode 100644
index 0000000..70633e6
--- /dev/null
+++ b/9player/src/nine.zig
@@ -0,0 +1,756 @@
+//! Synchronous 9P2000 session over a blocking file descriptor.
+//!
+//! A thin RPC layer over `cloud9.Client` (push/take, allocation-free). One request
+//! is outstanding at a time: the FUSE loop that drives this is single-threaded, so
+//! every call here blocks until its reply (or the connection's death) arrives.
+//! Fids are handed out from a free list; fid 0 is reserved for the root.
+const std = @import("std");
+const cloud9 = @import("cloud9");
+const linux = std.os.linux;
+
+pub const Address = union(enum) {
+ unix: []const u8,
+ tcp: struct { host: []const u8, port: u16 },
+ fd: i32,
+};
+
+/// A Stat whose every field means "leave unchanged" in a Twstat.
+pub const dontcare = cloud9.Stat{
+ .type = 0xFFFF,
+ .dev = 0xFFFF_FFFF,
+ .qid = .{ .type = 0xFF, .version = 0xFFFF_FFFF, .path = 0xFFFF_FFFF_FFFF_FFFF },
+ .mode = 0xFFFF_FFFF,
+ .atime = 0xFFFF_FFFF,
+ .mtime = 0xFFFF_FFFF,
+ .length = 0xFFFF_FFFF_FFFF_FFFF,
+ .name = "",
+ .uid = "",
+ .gid = "",
+ .muid = "",
+};
+
+pub const Session = struct {
+ pub const Error = error{ Nine, Protocol, Io, Closed, Stopped, TooLarge, OutOfMemory };
+
+ pub const Walk = struct { nwqid: u16, wqid: [cloud9.max_welem]cloud9.Qid };
+ pub const Open = struct { qid: cloud9.Qid, iounit: u32 };
+
+ gpa: std.mem.Allocator,
+ fd: i32,
+ client: cloud9.Client,
+ in_buf: []u8,
+ out_buf: []u8,
+ /// After `error.Nine`, the server's Rerror text (copied, bounded).
+ ename: [256]u8 = undefined,
+ ename_len: usize = 0,
+ /// Negotiated maximum message size.
+ msize: u32,
+ next_fid: u32 = 1,
+ free_fids: std.ArrayList(u32) = .empty,
+ /// Per-fid iounit learned from open/create (0 = none); used to chunk read/write.
+ iounits: std.AutoHashMapUnmanaged(u32, u32) = .empty,
+ /// Optional descriptor watched while waiting for a reply: when it becomes
+ /// readable (the bridge's "child exited" pipe) the pending rpc fails with
+ /// `error.Stopped` instead of blocking on a server that never answers.
+ stop_fd: i32 = -1,
+
+ /// Connect to `address`, then negotiate the protocol version.
+ /// `msize` is the maximum message size to ask for (0 = the buffers' size).
+ pub fn connect(gpa: std.mem.Allocator, address: Address, msize: u32) !Session {
+ const want: u32 = if (msize == 0) 8192 else @max(msize, 24);
+ const fd = try openTransport(address);
+ errdefer if (address != .fd) {
+ _ = linux.close(fd);
+ };
+
+ const in_buf = try gpa.alloc(u8, want);
+ errdefer gpa.free(in_buf);
+ const out_buf = try gpa.alloc(u8, want);
+ errdefer gpa.free(out_buf);
+
+ var s: Session = .{
+ .gpa = gpa,
+ .fd = fd,
+ .client = .init(.{ .in = in_buf, .out = out_buf }),
+ .in_buf = in_buf,
+ .out_buf = out_buf,
+ .msize = want,
+ };
+ const r = try s.rpc(.{ .version = .{ .msize = want } });
+ if (!std.mem.eql(u8, r.version.version, "9P2000")) return error.Protocol;
+ s.msize = r.version.msize;
+ return s;
+ }
+
+ /// Closes the descriptor and frees the buffers. Fids are not clunked.
+ pub fn deinit(s: *Session) void {
+ _ = linux.close(s.fd);
+ s.free_fids.deinit(s.gpa);
+ s.iounits.deinit(s.gpa);
+ s.gpa.free(s.in_buf);
+ s.gpa.free(s.out_buf);
+ s.* = undefined;
+ }
+
+ pub fn attach(s: *Session, fid: u32, uname: []const u8, aname: []const u8) Error!cloud9.Qid {
+ const r = try s.rpc(.{ .attach = .{ .fid = fid, .uname = uname, .aname = aname } });
+ return r.attach;
+ }
+
+ /// Fid 0 is never handed out: it belongs to the root attach.
+ pub fn allocFid(s: *Session) u32 {
+ if (s.free_fids.pop()) |fid| return fid;
+ const fid = s.next_fid;
+ s.next_fid += 1;
+ return fid;
+ }
+
+ /// Fids currently bound (excluding fid 0); a debugging aid for leak hunting.
+ pub fn fidsInUse(s: *const Session) usize {
+ return (s.next_fid - 1) - s.free_fids.items.len;
+ }
+
+ pub fn freeFid(s: *Session, fid: u32) void {
+ _ = s.iounits.remove(fid);
+ // If the free list cannot grow the fid is simply leaked; the counter keeps going.
+ s.free_fids.append(s.gpa, fid) catch {};
+ }
+
+ /// Generic RPC. Result slices borrow the input buffer until the next call.
+ pub fn rpc(s: *Session, req: cloud9.Client.Request) Error!cloud9.Client.Result {
+ s.ename_len = 0;
+ _ = s.client.submit(req) catch |e| switch (e) {
+ error.NoTags, error.Handshake, error.Dead => return error.Protocol,
+ error.NoSpace, error.TooLarge => return error.TooLarge,
+ error.BadRequest => {
+ s.setEname("bad request");
+ return error.Nine;
+ },
+ };
+ try s.flush();
+ var tmp: [64 * 1024]u8 = undefined;
+ while (true) {
+ if (s.client.take()) |done| {
+ switch (done.result) {
+ .fail => |ename| {
+ s.setEname(ename);
+ return error.Nine;
+ },
+ else => return done.result,
+ }
+ }
+ if (s.client.dead) return error.Protocol;
+ // After take() returned null the previous frame is gone, so the free
+ // space is at least what the pending frame still needs.
+ const room = s.client.in.len - s.client.in_len;
+ if (room == 0) return error.Protocol;
+ const n = try readSome(s.fd, s.stop_fd, tmp[0..@min(room, tmp.len)]);
+ if (n == 0) return error.Closed;
+ const pushed = s.client.push(tmp[0..n]);
+ if (pushed != n) return error.Protocol;
+ }
+ }
+
+ /// Walk `names` from `fid` to `newfid`. A partial walk leaves `newfid` unbound
+ /// (9P semantics) and reports `error.Nine` with ename "file does not exist".
+ pub fn walk(s: *Session, fid: u32, newfid: u32, names: []const []const u8) Error!Walk {
+ const r = try s.rpc(.{ .walk = .{ .fid = fid, .newfid = newfid, .names = names } });
+ if (r.walk.nwqid < names.len) {
+ s.setEname("file does not exist");
+ return error.Nine;
+ }
+ return .{ .nwqid = r.walk.nwqid, .wqid = r.walk.wqid };
+ }
+
+ /// allocFid + zero-element walk. The fid is released again on failure.
+ pub fn clone(s: *Session, fid: u32) Error!u32 {
+ const newfid = s.allocFid();
+ errdefer s.freeFid(newfid);
+ _ = try s.walk(fid, newfid, &.{});
+ return newfid;
+ }
+
+ pub fn open(s: *Session, fid: u32, mode: u8) Error!Open {
+ const r = try s.rpc(.{ .open = .{ .fid = fid, .mode = mode } });
+ s.noteIounit(fid, r.open.iounit);
+ return .{ .qid = r.open.qid, .iounit = r.open.iounit };
+ }
+
+ pub fn create(s: *Session, fid: u32, name: []const u8, perm: u32, mode: u8) Error!Open {
+ const r = try s.rpc(.{ .create = .{ .fid = fid, .name = name, .perm = perm, .mode = mode } });
+ s.noteIounit(fid, r.create.iounit);
+ return .{ .qid = r.create.qid, .iounit = r.create.iounit };
+ }
+
+ /// Reads into `buf`, chunking by min(maxRead, iounit) and stopping at the first
+ /// short read. Returns the number of bytes read (0 at end of file).
+ pub fn read(s: *Session, fid: u32, offset: u64, buf: []u8) Error!usize {
+ return readWith(s, rpc, fid, offset, buf, s.chunk(fid));
+ }
+
+ /// Writes `data`, chunking like `read` and stopping at the first short write.
+ pub fn write(s: *Session, fid: u32, offset: u64, data: []const u8) Error!usize {
+ return writeWith(s, rpc, fid, offset, data, s.chunkWrite(fid));
+ }
+
+ /// The returned Stat's strings (name/uid/gid/muid) borrow the session's input
+ /// buffer: they are valid only until the next rpc. Copy what must outlive it.
+ pub fn stat(s: *Session, fid: u32) Error!cloud9.Stat {
+ const r = try s.rpc(.{ .stat = .{ .fid = fid } });
+ return r.stat;
+ }
+
+ pub fn wstat(s: *Session, fid: u32, st: cloud9.Stat) Error!void {
+ _ = try s.rpc(.{ .wstat = .{ .fid = fid, .stat = st } });
+ }
+
+ /// Frees the fid locally even when the server reports an error.
+ pub fn clunk(s: *Session, fid: u32) Error!void {
+ defer s.freeFid(fid);
+ _ = try s.rpc(.{ .clunk = .{ .fid = fid } });
+ }
+
+ /// Frees the fid locally even when the server reports an error.
+ pub fn remove(s: *Session, fid: u32) Error!void {
+ defer s.freeFid(fid);
+ _ = try s.rpc(.{ .remove = .{ .fid = fid } });
+ }
+
+ /// Maps the last Rerror text to an errno (case-insensitive substring match).
+ pub fn errno(s: *const Session) linux.E {
+ return enameToErrno(s.ename[0..s.ename_len]);
+ }
+
+ // -- internals --------------------------------------------------------------
+
+ fn setEname(s: *Session, text: []const u8) void {
+ const n = @min(text.len, 255);
+ @memcpy(s.ename[0..n], text[0..n]);
+ s.ename_len = n;
+ }
+
+ fn noteIounit(s: *Session, fid: u32, iounit: u32) void {
+ if (iounit == 0) {
+ _ = s.iounits.remove(fid);
+ } else {
+ s.iounits.put(s.gpa, fid, iounit) catch {};
+ }
+ }
+
+ fn chunk(s: *Session, fid: u32) u32 {
+ return chunkSize(s.client.maxRead(), s.iounits.get(fid) orelse 0);
+ }
+
+ fn chunkWrite(s: *Session, fid: u32) u32 {
+ return chunkSize(s.client.maxWrite(), s.iounits.get(fid) orelse 0);
+ }
+
+ /// Writes everything in the client's output buffer to the socket.
+ fn flush(s: *Session) Error!void {
+ while (s.client.output().len != 0) {
+ const out = s.client.output();
+ const rc = linux.write(s.fd, out.ptr, out.len);
+ switch (linux.errno(rc)) {
+ .SUCCESS => {
+ if (rc == 0) return error.Closed;
+ s.client.wrote(rc);
+ },
+ .INTR, .AGAIN => continue,
+ .PIPE, .CONNRESET => return error.Closed,
+ else => return error.Io,
+ }
+ }
+ }
+};
+
+fn chunkSize(max: u32, iounit: u32) u32 {
+ if (iounit != 0 and iounit < max) return iounit;
+ return max;
+}
+
+/// Chunked read over any rpc-shaped function (injected so the loop is testable).
+fn readWith(
+ s: anytype,
+ comptime rpcFn: anytype,
+ fid: u32,
+ offset: u64,
+ buf: []u8,
+ max_chunk: u32,
+) Session.Error!usize {
+ if (max_chunk == 0) return error.Protocol;
+ var done: usize = 0;
+ while (done < buf.len) {
+ const want: u32 = @intCast(@min(buf.len - done, max_chunk));
+ const r = try rpcFn(s, .{ .read = .{ .fid = fid, .offset = offset + done, .count = want } });
+ const data = r.read;
+ @memcpy(buf[done..][0..data.len], data);
+ done += data.len;
+ if (data.len < want) break;
+ }
+ return done;
+}
+
+/// Chunked write over any rpc-shaped function.
+fn writeWith(
+ s: anytype,
+ comptime rpcFn: anytype,
+ fid: u32,
+ offset: u64,
+ data: []const u8,
+ max_chunk: u32,
+) Session.Error!usize {
+ if (max_chunk == 0) return error.Protocol;
+ var done: usize = 0;
+ while (done < data.len) {
+ const want: usize = @min(data.len - done, max_chunk);
+ const r = try rpcFn(s, .{ .write = .{ .fid = fid, .offset = offset + done, .data = data[done..][0..want] } });
+ done += r.write;
+ if (r.write < want) break;
+ }
+ return done;
+}
+
+fn readSome(fd: i32, stop_fd: i32, buf: []u8) Session.Error!usize {
+ while (true) {
+ if (stop_fd >= 0) {
+ var pfds = [_]linux.pollfd{
+ .{ .fd = fd, .events = linux.POLL.IN, .revents = 0 },
+ .{ .fd = stop_fd, .events = linux.POLL.IN, .revents = 0 },
+ };
+ const prc = linux.poll(&pfds, pfds.len, -1);
+ switch (linux.errno(prc)) {
+ .SUCCESS => {},
+ .INTR, .AGAIN => continue,
+ else => return error.Io,
+ }
+ if (pfds[1].revents != 0 and pfds[0].revents == 0) return error.Stopped;
+ }
+ const rc = linux.read(fd, buf.ptr, buf.len);
+ switch (linux.errno(rc)) {
+ .SUCCESS => return rc,
+ .INTR, .AGAIN => continue,
+ .CONNRESET => return error.Closed,
+ else => return error.Io,
+ }
+ }
+}
+
+/// Rerror text → errno, per docs/DESIGN.md (first match wins).
+pub fn enameToErrno(ename: []const u8) linux.E {
+ const Rule = struct { needle: []const u8, err: linux.E };
+ const rules = [_]Rule{
+ .{ .needle = "not exist", .err = .NOENT },
+ .{ .needle = "not found", .err = .NOENT },
+ .{ .needle = "no such", .err = .NOENT },
+ .{ .needle = "exists", .err = .EXIST },
+ .{ .needle = "not empty", .err = .NOTEMPTY },
+ .{ .needle = "not a dir", .err = .NOTDIR },
+ .{ .needle = "is a dir", .err = .ISDIR },
+ .{ .needle = "permission", .err = .ACCES },
+ .{ .needle = "denied", .err = .ACCES },
+ .{ .needle = "read-only", .err = .ROFS },
+ .{ .needle = "read only", .err = .ROFS },
+ .{ .needle = "readonly", .err = .ROFS },
+ .{ .needle = "no space", .err = .NOSPC },
+ .{ .needle = "not allowed", .err = .PERM },
+ .{ .needle = "not permitted", .err = .PERM },
+ .{ .needle = "cannot", .err = .PERM },
+ .{ .needle = "fid", .err = .BADF },
+ .{ .needle = "bad offset", .err = .INVAL },
+ .{ .needle = "invalid", .err = .INVAL },
+ .{ .needle = "bad ", .err = .INVAL },
+ .{ .needle = "busy", .err = .BUSY },
+ .{ .needle = "in use", .err = .BUSY },
+ .{ .needle = "too long", .err = .NAMETOOLONG },
+ .{ .needle = "not supported", .err = .OPNOTSUPP },
+ .{ .needle = "unsupported", .err = .OPNOTSUPP },
+ };
+ for (rules) |rule| {
+ if (std.ascii.findIgnoreCase(ename, rule.needle) != null) return rule.err;
+ }
+ return .IO;
+}
+
+// -- transport ------------------------------------------------------------------
+
+fn openTransport(address: Address) !i32 {
+ switch (address) {
+ .fd => |fd| return fd,
+ .unix => |path| {
+ if (path.len == 0 or path.len >= 108) return error.NameTooLong;
+ var sa: linux.sockaddr.un = .{ .path = @splat(0) };
+ @memcpy(sa.path[0..path.len], path);
+ const fd = try newSocket(linux.AF.UNIX, 0);
+ errdefer _ = linux.close(fd);
+ try doConnect(fd, @ptrCast(&sa), @sizeOf(linux.sockaddr.un));
+ return fd;
+ },
+ .tcp => |t| {
+ const ip = std.Io.net.IpAddress.parse(t.host, t.port) catch return error.InvalidAddress;
+ switch (ip) {
+ .ip4 => |a| {
+ const sa: linux.sockaddr.in = .{
+ .port = std.mem.nativeToBig(u16, t.port),
+ .addr = @bitCast(a.bytes),
+ };
+ const fd = try newSocket(linux.AF.INET, linux.IPPROTO.TCP);
+ errdefer _ = linux.close(fd);
+ setNodelay(fd);
+ try doConnect(fd, @ptrCast(&sa), @sizeOf(linux.sockaddr.in));
+ return fd;
+ },
+ .ip6 => |a| {
+ const sa: linux.sockaddr.in6 = .{
+ .port = std.mem.nativeToBig(u16, t.port),
+ .flowinfo = 0,
+ .addr = a.bytes,
+ .scope_id = 0,
+ };
+ const fd = try newSocket(linux.AF.INET6, linux.IPPROTO.TCP);
+ errdefer _ = linux.close(fd);
+ setNodelay(fd);
+ try doConnect(fd, @ptrCast(&sa), @sizeOf(linux.sockaddr.in6));
+ return fd;
+ },
+ }
+ },
+ }
+}
+
+fn newSocket(domain: u32, protocol: u32) !i32 {
+ const rc = linux.socket(domain, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, protocol);
+ switch (linux.errno(rc)) {
+ .SUCCESS => return @intCast(rc),
+ .MFILE, .NFILE => return error.ProcessFdQuotaExceeded,
+ .AFNOSUPPORT, .PROTONOSUPPORT => return error.AddressFamilyNotSupported,
+ .ACCES => return error.AccessDenied,
+ .NOMEM, .NOBUFS => return error.SystemResources,
+ else => return error.Unexpected,
+ }
+}
+
+fn setNodelay(fd: i32) void {
+ const one: u32 = 1;
+ _ = linux.setsockopt(fd, linux.IPPROTO.TCP, linux.TCP.NODELAY, @ptrCast(&one), @sizeOf(u32));
+}
+
+fn doConnect(fd: i32, addr: *const linux.sockaddr, len: linux.socklen_t) !void {
+ while (true) {
+ const rc = linux.connect(fd, addr, len);
+ switch (linux.errno(rc)) {
+ .SUCCESS => return,
+ .INTR => continue,
+ .CONNREFUSED => return error.ConnectionRefused,
+ .NOENT, .NOTDIR => return error.FileNotFound,
+ .ACCES, .PERM => return error.AccessDenied,
+ .TIMEDOUT => return error.ConnectionTimedOut,
+ .NETUNREACH, .HOSTUNREACH => return error.NetworkUnreachable,
+ .ADDRNOTAVAIL => return error.AddressNotAvailable,
+ .AGAIN, .INPROGRESS => return error.WouldBlock,
+ else => return error.Unexpected,
+ }
+ }
+}
+
+// -- tests ----------------------------------------------------------------------
+
+const testing = std.testing;
+
+test {
+ testing.refAllDecls(@This());
+}
+
+test "ename → errno mapping" {
+ try testing.expectEqual(linux.E.NOENT, enameToErrno("file does not exist"));
+ try testing.expectEqual(linux.E.NOENT, enameToErrno("No Such File"));
+ try testing.expectEqual(linux.E.NOENT, enameToErrno("directory entry not found"));
+ try testing.expectEqual(linux.E.EXIST, enameToErrno("file already exists"));
+ try testing.expectEqual(linux.E.NOTEMPTY, enameToErrno("directory not empty"));
+ try testing.expectEqual(linux.E.NOTDIR, enameToErrno("not a directory"));
+ try testing.expectEqual(linux.E.ISDIR, enameToErrno("is a directory"));
+ try testing.expectEqual(linux.E.ACCES, enameToErrno("permission denied"));
+ try testing.expectEqual(linux.E.ACCES, enameToErrno("access denied"));
+ try testing.expectEqual(linux.E.ROFS, enameToErrno("read-only file system"));
+ try testing.expectEqual(linux.E.NOSPC, enameToErrno("no space left"));
+ try testing.expectEqual(linux.E.PERM, enameToErrno("operation not permitted"));
+ try testing.expectEqual(linux.E.PERM, enameToErrno("cannot remove root"));
+ try testing.expectEqual(linux.E.BADF, enameToErrno("unknown fid"));
+ try testing.expectEqual(linux.E.BADF, enameToErrno("fid in use")); // "fid" precedes "in use"
+ try testing.expectEqual(linux.E.INVAL, enameToErrno("bad offset"));
+ try testing.expectEqual(linux.E.INVAL, enameToErrno("invalid argument"));
+ try testing.expectEqual(linux.E.INVAL, enameToErrno("bad request"));
+ try testing.expectEqual(linux.E.BUSY, enameToErrno("device busy"));
+ try testing.expectEqual(linux.E.NAMETOOLONG, enameToErrno("name too long"));
+ try testing.expectEqual(linux.E.OPNOTSUPP, enameToErrno("operation not supported"));
+ try testing.expectEqual(linux.E.IO, enameToErrno("something odd happened"));
+ try testing.expectEqual(linux.E.IO, enameToErrno(""));
+}
+
+test "fid allocator recycles and never hands out 0" {
+ var s: Session = undefined;
+ s.gpa = testing.allocator;
+ s.next_fid = 1;
+ s.free_fids = .empty;
+ s.iounits = .empty;
+ defer s.free_fids.deinit(s.gpa);
+ defer s.iounits.deinit(s.gpa);
+
+ const a = s.allocFid();
+ const b = s.allocFid();
+ const c = s.allocFid();
+ try testing.expectEqual(@as(u32, 1), a);
+ try testing.expectEqual(@as(u32, 2), b);
+ try testing.expectEqual(@as(u32, 3), c);
+ s.freeFid(b);
+ try testing.expectEqual(b, s.allocFid());
+ s.freeFid(a);
+ s.freeFid(c);
+ const x = s.allocFid();
+ const y = s.allocFid();
+ try testing.expect((x == a and y == c) or (x == c and y == a));
+ try testing.expectEqual(@as(u32, 4), s.allocFid());
+ try testing.expect(a != 0 and b != 0 and c != 0);
+}
+
+test "chunkSize honours iounit only when smaller" {
+ try testing.expectEqual(@as(u32, 100), chunkSize(100, 0));
+ try testing.expectEqual(@as(u32, 40), chunkSize(100, 40));
+ try testing.expectEqual(@as(u32, 100), chunkSize(100, 400));
+}
+
+/// Fake rpc for the chunked read/write loops: a file of `len` bytes where byte i == i & 0xff.
+const FakeFile = struct {
+ len: usize,
+ calls: usize = 0,
+ max_count: u32 = 0,
+ short_write_at: ?usize = null,
+ scratch: [4096]u8 = undefined,
+
+ fn rpc(f: *FakeFile, req: cloud9.Client.Request) Session.Error!cloud9.Client.Result {
+ f.calls += 1;
+ switch (req) {
+ .read => |r| {
+ f.max_count = @max(f.max_count, r.count);
+ if (r.offset >= f.len) return .{ .read = "" };
+ const n: usize = @min(@as(usize, r.count), f.len - @as(usize, @intCast(r.offset)));
+ for (f.scratch[0..n], 0..) |*b, i| b.* = @truncate(r.offset + i);
+ return .{ .read = f.scratch[0..n] };
+ },
+ .write => |w| {
+ f.max_count = @max(f.max_count, @as(u32, @intCast(w.data.len)));
+ if (f.short_write_at) |at| {
+ if (w.offset + w.data.len > at) {
+ const n: usize = if (w.offset >= at) 0 else @intCast(at - w.offset);
+ return .{ .write = @intCast(n) };
+ }
+ }
+ return .{ .write = @intCast(w.data.len) };
+ },
+ else => unreachable,
+ }
+ }
+};
+
+test "read chunks by max_chunk and stops at a short read" {
+ var f: FakeFile = .{ .len = 2500 };
+ var buf: [4000]u8 = undefined;
+ const n = try readWith(&f, FakeFile.rpc, 7, 0, &buf, 1000);
+ try testing.expectEqual(@as(usize, 2500), n);
+ try testing.expectEqual(@as(usize, 3), f.calls); // 1000, 1000, 500 (short → stop)
+ try testing.expectEqual(@as(u32, 1000), f.max_count);
+ for (buf[0..n], 0..) |b, i| try testing.expectEqual(@as(u8, @truncate(i)), b);
+
+ // Reading exactly up to a chunk boundary uses one call per chunk and no more.
+ f = .{ .len = 2000 };
+ try testing.expectEqual(@as(usize, 2000), try readWith(&f, FakeFile.rpc, 7, 0, buf[0..2000], 1000));
+ try testing.expectEqual(@as(usize, 2), f.calls);
+
+ // Offset past EOF → 0.
+ f = .{ .len = 10 };
+ try testing.expectEqual(@as(usize, 0), try readWith(&f, FakeFile.rpc, 7, 50, &buf, 1000));
+}
+
+test "write chunks and stops at a short write" {
+ var f: FakeFile = .{ .len = 0 };
+ var data: [2500]u8 = undefined;
+ for (&data, 0..) |*b, i| b.* = @truncate(i);
+ try testing.expectEqual(@as(usize, 2500), try writeWith(&f, FakeFile.rpc, 7, 0, &data, 1000));
+ try testing.expectEqual(@as(usize, 3), f.calls);
+ try testing.expectEqual(@as(u32, 1000), f.max_count);
+
+ f = .{ .len = 0, .short_write_at = 1500 };
+ try testing.expectEqual(@as(usize, 1500), try writeWith(&f, FakeFile.rpc, 7, 0, &data, 1000));
+ try testing.expectEqual(@as(usize, 2), f.calls);
+}
+
+// -- in-process server test ---------------------------------------------------------
+
+/// A tiny 9P2000 backend on a cloud9.Server: answers version/attach/walk/stat/open/
+/// read/clunk/remove with canned data. Runs in its own thread over a socketpair.
+const FakeServer = struct {
+ fd: i32,
+ msize: u32,
+ max_read_count: u32 = 0,
+ file_len: usize,
+
+ const file_qid: cloud9.Qid = .{ .type = 0, .version = 3, .path = 0x1234 };
+ const dir_qid: cloud9.Qid = .{ .type = cloud9.qtdir, .version = 1, .path = 0x1 };
+
+ fn run(fs: *FakeServer) void {
+ fs.loop() catch |e| std.debug.print("fake server: {s}\n", .{@errorName(e)});
+ _ = linux.close(fs.fd);
+ }
+
+ fn loop(fs: *FakeServer) !void {
+ const gpa = testing.allocator;
+ const in = try gpa.alloc(u8, fs.msize);
+ defer gpa.free(in);
+ const out = try gpa.alloc(u8, fs.msize * 2);
+ defer gpa.free(out);
+ var srv: cloud9.Server = .init(.{ .in = in, .out = out });
+ var tmp: [4096]u8 = undefined;
+ var data: [8192]u8 = undefined;
+ while (true) {
+ while (try srv.receive()) |req| {
+ const tag = req.tag;
+ switch (req.msg) {
+ .tversion => |m| try srv.negotiate(m.msize, m.version),
+ .tattach => try srv.reply(tag, .{ .rattach = .{ .qid = dir_qid } }),
+ .twalk => |m| {
+ var wq: [cloud9.max_welem]cloud9.Qid = @splat(dir_qid);
+ var n: u16 = 0;
+ for (m.wname[0..m.nwname]) |name| {
+ if (std.mem.eql(u8, name, "file")) {
+ wq[n] = file_qid;
+ } else if (std.mem.eql(u8, name, "dir")) {
+ wq[n] = dir_qid;
+ } else break;
+ n += 1;
+ }
+ if (n == 0 and m.nwname != 0) {
+ try srv.reply(tag, .{ .rerror = .{ .ename = "file does not exist" } });
+ } else {
+ try srv.reply(tag, .{ .rwalk = .{ .nwqid = n, .wqid = wq } });
+ }
+ },
+ .tstat => try srv.reply(tag, .{ .rstat = .{ .stat = .{
+ .type = 0,
+ .dev = 0,
+ .qid = file_qid,
+ .mode = 0o644,
+ .atime = 1,
+ .mtime = 2,
+ .length = fs.file_len,
+ .name = "file",
+ .uid = "u",
+ .gid = "g",
+ .muid = "u",
+ } } }),
+ .topen => |m| try srv.reply(tag, .{ .ropen = .{ .qid = file_qid, .iounit = if (m.mode == cloud9.owrite) 700 else 0 } }),
+ .tread => |m| {
+ fs.max_read_count = @max(fs.max_read_count, m.count);
+ var n: usize = 0;
+ if (m.offset < fs.file_len) n = @min(@as(usize, m.count), fs.file_len - @as(usize, @intCast(m.offset)));
+ n = @min(n, data.len);
+ for (data[0..n], 0..) |*b, i| b.* = @truncate(m.offset + i);
+ try srv.reply(tag, .{ .rread = .{ .data = data[0..n] } });
+ },
+ .twrite => |m| try srv.reply(tag, .{ .rwrite = .{ .count = @intCast(m.data.len) } }),
+ .tclunk => try srv.reply(tag, .rclunk),
+ .tremove => try srv.reply(tag, .{ .rerror = .{ .ename = "permission denied" } }),
+ .twstat => try srv.reply(tag, .rwstat),
+ // A flush is the test's "hang up now" signal.
+ .tflush => return,
+ else => try srv.reply(tag, .{ .rerror = .{ .ename = "not supported" } }),
+ }
+ srv.release();
+ }
+ while (srv.output().len != 0) {
+ const o = srv.output();
+ const rc = linux.write(fs.fd, o.ptr, o.len);
+ if (linux.errno(rc) != .SUCCESS) return error.Write;
+ srv.wrote(rc);
+ }
+ const rc = linux.read(fs.fd, &tmp, tmp.len);
+ if (linux.errno(rc) != .SUCCESS) return error.Read;
+ if (rc == 0) return;
+ if (srv.push(tmp[0..rc]) != rc) return error.Overflow;
+ }
+ }
+};
+
+test "session against an in-process cloud9.Server" {
+ var fds: [2]i32 = undefined;
+ try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.socketpair(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0, &fds)));
+
+ var fs: FakeServer = .{ .fd = fds[1], .msize = 8192, .file_len = 20_000 };
+ const th = try std.Thread.spawn(.{}, FakeServer.run, .{&fs});
+
+ var s = try Session.connect(testing.allocator, .{ .fd = fds[0] }, 8192);
+ defer {
+ s.deinit();
+ th.join();
+ }
+ try testing.expectEqual(@as(u32, 8192), s.msize);
+
+ const root = try s.attach(0, "me", "");
+ try testing.expectEqual(FakeServer.dir_qid.path, root.path);
+
+ // Plain rpc + stat borrowing the input buffer.
+ const fid = s.allocFid();
+ const w = try s.walk(0, fid, &.{"file"});
+ try testing.expectEqual(@as(u16, 1), w.nwqid);
+ try testing.expectEqual(FakeServer.file_qid.path, w.wqid[0].path);
+ const st = try s.stat(fid);
+ try testing.expectEqualStrings("file", st.name);
+ try testing.expectEqual(@as(u64, 20_000), st.length);
+
+ // Chunked read: 20000 bytes at maxRead = msize - 11 = 8181 per chunk.
+ _ = try s.open(fid, cloud9.oread);
+ const buf = try testing.allocator.alloc(u8, 30_000);
+ defer testing.allocator.free(buf);
+ const n = try s.read(fid, 0, buf);
+ try testing.expectEqual(@as(usize, 20_000), n);
+ for (buf[0..n], 0..) |b, i| try testing.expectEqual(@as(u8, @truncate(i)), b);
+ try testing.expectEqual(@as(u32, 8181), fs.max_read_count);
+ try testing.expectEqual(@as(usize, 0), try s.read(fid, 20_000, buf));
+
+ // iounit from open bounds the chunk.
+ const wfid = try s.clone(fid);
+ _ = try s.open(wfid, cloud9.owrite);
+ fs.max_read_count = 0;
+ _ = try s.read(wfid, 0, buf[0..3000]);
+ try testing.expectEqual(@as(u32, 700), fs.max_read_count);
+ try testing.expectEqual(@as(usize, 3000), try s.write(wfid, 0, buf[0..3000]));
+
+ // Partial walk → error.Nine with a "not exist" ename → ENOENT.
+ const pfid = s.allocFid();
+ try testing.expectError(error.Nine, s.walk(0, pfid, &.{ "dir", "nope" }));
+ try testing.expectEqual(linux.E.NOENT, s.errno());
+ try testing.expectEqualStrings("file does not exist", s.ename[0..s.ename_len]);
+ s.freeFid(pfid);
+
+ // Server Rerror → error.Nine, ename copied, fid freed by remove even on error.
+ try testing.expectError(error.Nine, s.remove(wfid));
+ try testing.expectEqual(linux.E.ACCES, s.errno());
+ try testing.expectEqual(wfid, s.allocFid()); // recycled
+ s.freeFid(wfid);
+
+ // Unsupported op → "not supported" → ENOTSUP; a plain wstat succeeds.
+ try testing.expectError(error.Nine, s.rpc(.{ .auth = .{ .afid = 5, .uname = "me" } }));
+ try testing.expectEqual(linux.E.OPNOTSUPP, s.errno());
+ try s.wstat(fid, dontcare);
+ try s.clunk(fid);
+ try testing.expectEqual(fid, s.allocFid());
+ s.freeFid(fid);
+
+ // A clone bound to a fid that then fails to walk must release the fid.
+ const before = s.next_fid;
+ const cfid = s.allocFid();
+ s.freeFid(cfid);
+ try testing.expectError(error.Nine, s.walk(0, cfid, &.{"nope"}));
+ try testing.expectEqual(before, s.next_fid);
+
+ // The server hanging up makes the pending rpc fail with error.Closed.
+ try testing.expectError(error.Closed, s.rpc(.{ .flush = .{ .oldtag = 0 } }));
+}