diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-19 21:26:05 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-09-19 21:26:05 -0300 |
| commit | b05abcba3ea09ea106ad28364c6e40a3ec31b890 (patch) | |
| tree | 9170fac5e7e5d8bde108de34a182aaa9d6844117 /9player/src/nine.zig | |
| parent | ae310a207534b33b7321dd2b9f423a73b1969159 (diff) | |
| download | cloud9-b05abcba3ea09ea106ad28364c6e40a3ec31b890.tar.gz cloud9-b05abcba3ea09ea106ad28364c6e40a3ec31b890.zip | |
Add 9player and introspect as programs beside the library
9player/: FUSE mount CLI that mounts a 9P2000 tree into a fresh user+mount
namespace and runs a program in it (no root, no libfuse, no libc).
introspect/: the 9P debug/introspection library (freestanding core, value
renderers, Linux probe with threads/stacks/memory/breakpoints/panics) and
its demo server. Each has its own build fragment; the root build.zig wires
them behind -D9player/-Dintrospect with namespaced steps (9player-itest,
introspect-check-freestanding, programs-test, ...) and exports the
introspect module for dependents. This is the layout for related programs.
Co-Authored-By: Claude Fable 5.1 <[email protected]>
Diffstat (limited to '9player/src/nine.zig')
| -rw-r--r-- | 9player/src/nine.zig | 756 |
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 } })); +} |
