//! 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 = "", }; /// How the owner of a session (the FUSE bridge) gets a say while an rpc waits /// for its reply. `watch` names a descriptor to poll alongside the socket, or /// -1 to poll nothing extra right now; when it becomes readable `onReadable` /// consumes whatever is there and returns true if the request in flight should /// be cancelled with a Tflush. `armed` reports whether such a cancellation was /// requested earlier for the operation in progress; the chunked read/write /// loops stop between chunks when it is set (a reply that raced the flush still /// leaves the caller wanting out). pub const Interrupt = struct { ctx: *anyopaque, watch: *const fn (ctx: *anyopaque) i32, onReadable: *const fn (ctx: *anyopaque) Session.Error!bool, armed: *const fn (ctx: *anyopaque) bool, }; pub const Session = struct { pub const Error = error{ Nine, Protocol, Io, Closed, Stopped, Interrupted, 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, /// Optional interrupt source (the bridge's FUSE descriptor) consulted while /// a reply is outstanding; see `Interrupt`. interrupt: ?Interrupt = null, /// While set, a cancellation ends the rpc at once with `error.Interrupted` /// and wedges the session, with no Tflush: for the handshake (version, /// attach, the root stat), where there is no session yet to flush a /// request out of, and the honest answer to "stop waiting" is to hang up. abort_on_cancel: bool = false, /// Milliseconds a Tflush may go unanswered before the server is declared /// wedged: the rpc fails with `error.Interrupted` and the session with it. /// The protocol says a client waits for the Rflush; a server that has not /// managed one in this long is not going to, and the process behind the /// interrupt is unkillable until we stop waiting. 0 waits forever. flush_grace_ms: i32 = 3000, /// The server is gone as far as this session is concerned (see /// `abort_on_cancel`, `flush_grace_ms`); every rpc answers `error.Closed`. wedged: bool = false, /// 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 { return connectWatched(gpa, address, msize, -1); } /// `connect`, with `stop_fd` watched for the whole handshake (the version /// rpc included). A server that accepts the connection and then never /// answers the Tversion would otherwise pin the caller in a blocking read /// with no way out; -1 disables the watch. The field stays set on the /// returned session. An `.fd` address is the caller's to close on failure. pub fn connectWatched(gpa: std.mem.Allocator, address: Address, msize: u32, stop_fd: i32) !Session { var s = try dial(gpa, address, msize, stop_fd); s.version() catch |e| { s.freeBuffers(); if (address != .fd) _ = linux.close(s.fd); return e; }; return s; } /// The transport and the buffers, no handshake: for a caller that wants /// its interrupt source in place before `version()` (9ns's mntgen /// worker, so a walk interrupted mid-dial can abandon the dial). Owns /// the descriptor from here: `deinit` closes it. pub fn dial(gpa: std.mem.Allocator, address: Address, msize: u32, stop_fd: i32) !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); return .{ .gpa = gpa, .fd = fd, .client = .init(.{ .in = in_buf, .out = out_buf }), .in_buf = in_buf, .out_buf = out_buf, .msize = want, .stop_fd = stop_fd, }; } /// Negotiates the protocol version with the msize `dial` was given. pub fn version(s: *Session) Error!void { const r = try s.rpc(.{ .version = .{ .msize = s.msize } }); if (!std.mem.eql(u8, r.version.version, "9P2000")) return error.Protocol; s.msize = r.version.msize; } fn freeBuffers(s: *Session) void { s.free_fids.deinit(s.gpa); s.iounits.deinit(s.gpa); s.gpa.free(s.in_buf); s.gpa.free(s.out_buf); } /// Closes the descriptor and frees the buffers. Fids are not clunked. pub fn deinit(s: *Session) void { _ = linux.close(s.fd); s.freeBuffers(); 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. /// /// While the reply is outstanding the socket is polled together with /// `stop_fd` (→ `error.Stopped`) and the interrupt source's descriptor. When /// the latter asks for a cancellation a Tflush for the request's tag goes out /// and the wait continues until either the original reply arrives (the flush /// lost the race; the result is returned as if nothing happened and the /// Rflush is swallowed by a later call) or the Rflush does (→ /// `error.Interrupted`; the server has dropped the request), or neither /// within `flush_grace_ms` (→ `error.Interrupted`, and the session is /// wedged: the server stopped talking). Under `abort_on_cancel` the /// cancellation itself wedges the session, with nothing sent. pub fn rpc(s: *Session, req: cloud9.Client.Request) Error!cloud9.Client.Result { if (s.wedged) return error.Closed; s.ename_len = 0; const tag = 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 flush_tag: ?u16 = null; // Monotonic ms by which the Tflush must have been answered. var flush_deadline: ?i64 = null; var tmp: [64 * 1024]u8 = undefined; while (true) { while (s.client.take()) |done| { if (done.tag == tag) { switch (done.result) { .fail => |ename| { s.setEname(ename); return error.Nine; }, else => return done.result, } } if (flush_tag != null and done.tag == flush_tag.?) return error.Interrupted; // An Rflush for a flush whose original reply won the race in an // earlier call: the client has released both tags; nothing to do. if (done.op == .flush) continue; return error.Protocol; } 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 grace: i32 = if (flush_deadline) |d| @intCast(@max(d - nowMs(), 0)) else -1; switch (try s.wait(grace)) { .socket => { const n = try readSocket(s.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; }, .cancel => { if (s.abort_on_cancel) { s.wedged = true; return error.Interrupted; } if (flush_tag == null) { flush_tag = s.client.submit(.{ .flush = .{ .oldtag = tag } }) catch return error.Protocol; try s.flush(); if (s.flush_grace_ms != 0) flush_deadline = nowMs() + s.flush_grace_ms; } }, .timeout => { s.wedged = true; return error.Interrupted; }, } } } /// True when the interrupt source has asked for the operation in progress to /// stop; consulted between the chunks of a read or write. pub fn interruptArmed(s: *const Session) bool { const i = s.interrupt orelse return false; return i.armed(i.ctx); } /// 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); } const Ready = enum { socket, cancel, timeout }; /// Blocks until the socket is readable (`.socket`), the interrupt source /// wants the request in flight cancelled (`.cancel`), `timeout_ms` passes /// with neither (`.timeout`; -1 waits forever), or `stop_fd` fires /// (`error.Stopped`). Anything the interrupt source consumes without asking /// for a cancellation simply resumes the wait. fn wait(s: *Session, timeout_ms: i32) Error!Ready { while (true) { var pfds: [3]linux.pollfd = undefined; var n: usize = 0; pfds[n] = .{ .fd = s.fd, .events = linux.POLL.IN, .revents = 0 }; n += 1; const stop_at: ?usize = if (s.stop_fd >= 0) n else null; if (stop_at != null) { pfds[n] = .{ .fd = s.stop_fd, .events = linux.POLL.IN, .revents = 0 }; n += 1; } const ifd: i32 = if (s.interrupt) |i| i.watch(i.ctx) else -1; const int_at: ?usize = if (ifd >= 0) n else null; if (int_at != null) { pfds[n] = .{ .fd = ifd, .events = linux.POLL.IN, .revents = 0 }; n += 1; } if (n == 1 and timeout_ms < 0) return .socket; const prc = linux.poll(&pfds, @intCast(n), timeout_ms); switch (linux.errno(prc)) { .SUCCESS => {}, .INTR, .AGAIN => continue, else => return error.Io, } if (prc == 0) return .timeout; // A reply that is already there wins over everything else. if (pfds[0].revents != 0) return .socket; if (stop_at) |i| { if (pfds[i].revents != 0) return error.Stopped; } if (int_at) |i| { if (pfds[i].revents != 0) { const src = s.interrupt.?; if (try src.onReadable(src.ctx)) return .cancel; } } } } /// 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, } } } }; /// The monotonic clock in milliseconds: deadlines, not timestamps. fn nowMs() i64 { var ts: linux.timespec = undefined; _ = linux.clock_gettime(.MONOTONIC, &ts); return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000); } 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) { // Like read(2): an interruption after some data arrived is a short read. if (done != 0 and s.interruptArmed()) break; const want: u32 = @intCast(@min(buf.len - done, max_chunk)); const r = rpcFn(s, .{ .read = .{ .fid = fid, .offset = offset + done, .count = want } }) catch |e| { if (e == error.Interrupted and done != 0) break; return e; }; 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) { if (done != 0 and s.interruptArmed()) break; const want: usize = @min(data.len - done, max_chunk); const r = rpcFn(s, .{ .write = .{ .fid = fid, .offset = offset + done, .data = data[done..][0..want] } }) catch |e| { if (e == error.Interrupted and done != 0) break; return e; }; done += r.write; if (r.write < want) break; } return done; } fn readSocket(fd: i32, buf: []u8) Session.Error!usize { while (true) { 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{ // A server that answers a flushed request with an error (Pardes says // "Interrupted system call") should look like a flush to the caller. .{ .needle = "interrupt", .err = .INTR }, .{ .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.INTR, enameToErrno("Interrupted system call")); try testing.expectEqual(linux.E.INTR, enameToErrno("read interrupted")); 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, /// Fail the call with `error.Interrupted` once this many calls were made. interrupt_at: ?usize = null, /// Report an armed interrupt once this many calls were made. armed_at: ?usize = null, scratch: [4096]u8 = undefined, fn interruptArmed(f: *const FakeFile) bool { return if (f.armed_at) |at| f.calls >= at else false; } fn rpc(f: *FakeFile, req: cloud9.Client.Request) Session.Error!cloud9.Client.Result { f.calls += 1; if (f.interrupt_at) |at| if (f.calls > at) return error.Interrupted; 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); } test "interrupted chunk loops: partial count if data moved, Interrupted otherwise" { var buf: [4000]u8 = undefined; // The second chunk's rpc is interrupted: the first chunk is returned. var f: FakeFile = .{ .len = 2500, .interrupt_at = 1 }; try testing.expectEqual(@as(usize, 1000), try readWith(&f, FakeFile.rpc, 7, 0, &buf, 1000)); // The first chunk's rpc is interrupted: nothing was transferred. f = .{ .len = 2500, .interrupt_at = 0 }; try testing.expectError(error.Interrupted, readWith(&f, FakeFile.rpc, 7, 0, &buf, 1000)); // A reply that raced the flush arms the interrupt: stop before the next chunk. f = .{ .len = 2500, .armed_at = 2 }; try testing.expectEqual(@as(usize, 2000), try readWith(&f, FakeFile.rpc, 7, 0, &buf, 1000)); try testing.expectEqual(@as(usize, 2), f.calls); // Same for writes. var data: [2500]u8 = undefined; for (&data, 0..) |*b, i| b.* = @truncate(i); f = .{ .len = 0, .interrupt_at = 2 }; try testing.expectEqual(@as(usize, 2000), try writeWith(&f, FakeFile.rpc, 7, 0, &data, 1000)); f = .{ .len = 0, .interrupt_at = 0 }; try testing.expectError(error.Interrupted, writeWith(&f, FakeFile.rpc, 7, 0, &data, 1000)); f = .{ .len = 0, .armed_at = 1 }; try testing.expectEqual(@as(usize, 1000), try writeWith(&f, FakeFile.rpc, 7, 0, &data, 1000)); } // -- in-process server test --------------------------------------------------------- test "flush grace: a Tflush the server never answers wedges the session after the grace" { var pi = try PipeInterrupt.init(7); defer pi.deinit(); var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .ignore, .on_hang_inject = &pi }; var pair = try fs.start(); defer pair.close(); const s = &pair.session; s.interrupt = pi.interface(); s.flush_grace_ms = 100; // The read at offset 0 hangs; the injected INTERRUPT sends a Tflush; the // server ignores it; the grace runs out. var buf: [100]u8 = undefined; try testing.expectError(error.Interrupted, s.read(1, 0, &buf)); try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst)); try testing.expect(s.wedged); // From here on the session is closed for business, without another // byte to the server. try testing.expectError(error.Closed, s.read(1, 10, &buf)); try testing.expectError(error.Closed, s.clunk(1)); try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst)); // The socket is still the session's to close: the server's thread ends // when `close` drops it. } test "flush grace: zero waits for the Rflush, however late" { var pi = try PipeInterrupt.init(7); defer pi.deinit(); var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush, .on_hang_inject = &pi, .read_delay_ns = 0 }; var pair = try fs.start(); defer pair.close(); const s = &pair.session; s.interrupt = pi.interface(); s.flush_grace_ms = 0; var buf: [100]u8 = undefined; try testing.expectError(error.Interrupted, s.read(1, 0, &buf)); try testing.expect(!s.wedged); // The session lives: the Rflush released the tag and reads go on. try testing.expectEqual(@as(usize, 40), try s.read(1, 10, buf[0..40])); } test "abort_on_cancel: a cancellation ends the rpc at once, sends no Tflush, wedges the session" { var pi = try PipeInterrupt.init(7); defer pi.deinit(); var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush, .on_hang_inject = &pi }; var pair = try fs.start(); defer pair.close(); const s = &pair.session; s.interrupt = pi.interface(); s.abort_on_cancel = true; var buf: [100]u8 = undefined; try testing.expectError(error.Interrupted, s.read(1, 0, &buf)); try testing.expect(s.wedged); try testing.expectEqual(@as(u32, 0), fs.flushes.load(.seq_cst)); try testing.expectError(error.Closed, s.read(1, 10, &buf)); } /// 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. /// Test support only (bridge.zig's tests use it too). pub const FakeServer = struct { fd: i32, msize: u32, max_read_count: u32 = 0, file_len: usize, /// A Tread at this offset is never answered (a blocked stream read); the /// server keeps serving whatever else arrives, notably a Tflush. hang_offset: ?u64 = null, /// What a Tflush gets: the connection dropped (`.hangup`), an Rflush, /// first the Rread the flush was aimed at and then the Rflush (the /// race), or nothing at all (`.ignore`: counted, never answered, the /// read stays hung — a wedged server). on_flush: enum { hangup, rflush, reply_then_rflush, ignore } = .hangup, /// Observed by the test thread: number of Tflush seen and the last oldtag. flushes: std.atomic.Value(u32) = .init(0), flush_oldtag: std.atomic.Value(u32) = .init(0xFFFF), /// The tag of the hung Tread, for the test to compare with `flush_oldtag`. hung_tag: std.atomic.Value(u32) = .init(0xFFFF), /// When set, the server injects that source's INTERRUPT the moment a read /// hangs, so the cancellation provably arrives while the wait is on. on_hang_inject: ?*PipeInterrupt = null, /// Delay before every Rread, so a test can be sure the client is waiting. read_delay_ns: u64 = 0, pub const file_qid: cloud9.Qid = .{ .type = 0, .version = 3, .path = 0x1234 }; pub const dir_qid: cloud9.Qid = .{ .type = cloud9.qtdir, .version = 1, .path = 0x1 }; pub 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; var hung: ?struct { tag: u16, offset: u64, count: u32 } = null; 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); if (fs.hang_offset != null and fs.hang_offset.? == m.offset) { hung = .{ .tag = tag, .offset = m.offset, .count = m.count }; fs.hung_tag.store(tag, .seq_cst); if (fs.on_hang_inject) |p| try p.inject(p.unique); } else { if (fs.read_delay_ns != 0) { const ts: linux.timespec = .{ .sec = @intCast(fs.read_delay_ns / std.time.ns_per_s), .nsec = @intCast(fs.read_delay_ns % std.time.ns_per_s) }; _ = linux.nanosleep(&ts, null); } try srv.reply(tag, .{ .rread = .{ .data = fs.fill(&data, m.offset, m.count) } }); } }, .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), .tflush => |m| { fs.flush_oldtag.store(m.oldtag, .seq_cst); _ = fs.flushes.fetchAdd(1, .seq_cst); switch (fs.on_flush) { // The test's "hang up now" signal. .hangup => return, .ignore => {}, .rflush => { if (hung != null and hung.?.tag == m.oldtag) hung = null; try srv.reply(tag, .rflush); }, .reply_then_rflush => { if (hung) |h| if (h.tag == m.oldtag) { try srv.reply(h.tag, .{ .rread = .{ .data = fs.fill(&data, h.offset, h.count) } }); hung = null; }; try srv.reply(tag, .rflush); }, } }, 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; } } /// File contents: byte i == i & 0xff, `file_len` bytes long. fn fill(fs: *const FakeServer, data: []u8, offset: u64, count: u32) []const u8 { var n: usize = 0; if (offset < fs.file_len) n = @min(@as(usize, count), fs.file_len - @as(usize, @intCast(offset))); n = @min(n, data.len); for (data[0..n], 0..) |*b, i| b.* = @truncate(offset + i); return data[0..n]; } /// A connected session (fid 0 attached, fid 1 walked to "file" and opened /// for reading) plus the server thread; `close` when done. pub const Pair = struct { server: *FakeServer, session: Session, thread: std.Thread, pub fn close(p: *Pair) void { p.session.deinit(); p.thread.join(); } }; pub fn start(fs: *FakeServer) !Pair { var fds: [2]i32 = undefined; if (linux.errno(linux.socketpair(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0, &fds)) != .SUCCESS) return error.Io; fs.fd = fds[1]; const th = try std.Thread.spawn(.{}, FakeServer.run, .{fs}); var s = try Session.connect(testing.allocator, .{ .fd = fds[0] }, fs.msize); errdefer s.deinit(); _ = try s.attach(0, "me", ""); const fid = s.allocFid(); _ = try s.walk(0, fid, &.{"file"}); _ = try s.open(fid, cloud9.oread); return .{ .server = fs, .session = s, .thread = th }; } }; /// A fake interrupt source for the rpc wait loop: a SOCK_SEQPACKET pair stands /// in for the FUSE descriptor (one datagram per request, like /dev/fuse /// delivers one request per read). Mirrors the bridge's rules: an INTERRUPT for /// `unique` arms and cancels, anything else is consumed and ignored. pub const PipeInterrupt = struct { read_end: i32, write_end: i32, unique: u64, armed_flag: bool = false, /// Requests consumed that were not the matching INTERRUPT. ignored: usize = 0, const fuse_interrupt_opcode: u32 = 36; pub fn init(unique: u64) !PipeInterrupt { var fds: [2]i32 = undefined; const flags = linux.SOCK.SEQPACKET | linux.SOCK.CLOEXEC | linux.SOCK.NONBLOCK; if (linux.errno(linux.socketpair(linux.AF.UNIX, flags, 0, &fds)) != .SUCCESS) return error.Io; return .{ .read_end = fds[0], .write_end = fds[1], .unique = unique }; } pub fn deinit(p: *PipeInterrupt) void { _ = linux.close(p.read_end); _ = linux.close(p.write_end); } pub fn interface(p: *PipeInterrupt) Interrupt { return .{ .ctx = p, .watch = watch, .onReadable = onReadable, .armed = armed }; } /// Writes a FUSE_INTERRUPT request (InHeader + InterruptIn) naming `target`. pub fn inject(p: *PipeInterrupt, target: u64) !void { var wire: [48]u8 = undefined; std.mem.writeInt(u32, wire[0..4], 48, .little); // len std.mem.writeInt(u32, wire[4..8], fuse_interrupt_opcode, .little); // opcode std.mem.writeInt(u64, wire[8..16], 0x8000_0000_0000_0001, .little); // the interrupt's own unique @memset(wire[16..40], 0); // nodeid, uid, gid, pid, extlen, padding std.mem.writeInt(u64, wire[40..48], target, .little); // InterruptIn.unique if (linux.write(p.write_end, &wire, wire.len) != wire.len) return error.Io; } fn watch(ctx: *anyopaque) i32 { const p: *PipeInterrupt = @ptrCast(@alignCast(ctx)); return p.read_end; } fn onReadable(ctx: *anyopaque) Session.Error!bool { const p: *PipeInterrupt = @ptrCast(@alignCast(ctx)); var buf: [4096]u8 = undefined; const rc = linux.read(p.read_end, &buf, buf.len); if (linux.errno(rc) != .SUCCESS or rc != 48) return error.Io; const opcode = std.mem.readInt(u32, buf[4..8], .little); const target = std.mem.readInt(u64, buf[40..48], .little); if (opcode == fuse_interrupt_opcode and target == p.unique) { p.armed_flag = true; return true; } p.ignored += 1; return false; } fn armed(ctx: *anyopaque) bool { const p: *PipeInterrupt = @ptrCast(@alignCast(ctx)); return p.armed_flag; } }; test "rpc wait loop: INTERRUPT → Tflush → Rflush → error.Interrupted; session still usable" { var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush }; var pair = try fs.start(); defer pair.close(); const s = &pair.session; var pi = try PipeInterrupt.init(77); defer pi.deinit(); s.interrupt = pi.interface(); defer s.interrupt = null; // An INTERRUPT for some other request is consumed and ignored: the wait // goes on, and the reply (a read past the hang offset) arrives normally. var buf: [100]u8 = undefined; try pi.inject(78); try testing.expectEqual(@as(usize, 40), try s.read(1, 10, &buf)); try testing.expect(!pi.armed_flag); // The read at offset 0 hangs; the INTERRUPT for our request cancels it // (the stray one above is consumed along the way if the reply beat it). try pi.inject(77); try testing.expectError(error.Interrupted, s.read(1, 0, &buf)); try testing.expect(pi.armed_flag); try testing.expectEqual(@as(usize, 1), pi.ignored); try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst)); try testing.expectEqual(fs.hung_tag.load(.seq_cst), fs.flush_oldtag.load(.seq_cst)); // Both tags are free again: further rpcs work. const st = try s.stat(1); try testing.expectEqual(@as(u64, 50), st.length); try testing.expectEqual(@as(usize, 40), try s.read(1, 10, &buf)); try testing.expectEqual(@as(usize, 0), s.client.pending()); } test "rpc wait loop: the reply beats the Rflush → data returned, stray Rflush swallowed" { var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .reply_then_rflush }; var pair = try fs.start(); defer pair.close(); const s = &pair.session; var pi = try PipeInterrupt.init(5); defer pi.deinit(); s.interrupt = pi.interface(); defer s.interrupt = null; var buf: [40]u8 = undefined; try pi.inject(5); try testing.expectEqual(@as(usize, 30), try s.read(1, 0, buf[0..30])); for (buf[0..30], 0..) |b, i| try testing.expectEqual(@as(u8, @truncate(i)), b); try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst)); try testing.expectEqual(fs.hung_tag.load(.seq_cst), fs.flush_oldtag.load(.seq_cst)); // The Rflush is still in flight (or already buffered): the next rpcs must // step over it, and afterwards nothing is pending in the client. pi.armed_flag = false; const st = try s.stat(1); try testing.expectEqual(@as(u64, 50), st.length); try testing.expectEqual(@as(usize, 40), try s.read(1, 10, &buf)); try testing.expectEqual(@as(usize, 0), s.client.pending()); } test "rpc wait loop: a read interrupted after some data is a short read" { // Chunks of maxRead = 1013 (msize 1024); the third chunk (offset 2026) hangs // and the server fires the INTERRUPT at that moment. var pi = try PipeInterrupt.init(9); defer pi.deinit(); var fs: FakeServer = .{ .fd = -1, .msize = 1024, .file_len = 5000, .hang_offset = 2026, .on_flush = .rflush, .on_hang_inject = &pi }; var pair = try fs.start(); defer pair.close(); const s = &pair.session; s.interrupt = pi.interface(); defer s.interrupt = null; var buf: [4000]u8 = undefined; try testing.expectEqual(@as(usize, 2026), try s.read(1, 0, &buf)); for (buf[0..2026], 0..) |b, i| try testing.expectEqual(@as(u8, @truncate(i)), b); try testing.expect(pi.armed_flag); try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst)); try testing.expectEqual(@as(usize, 0), s.client.pending()); // An INTERRUPT already waiting when the read starts: the first chunk's // reply races the flush and wins, the armed flag then stops the loop. pi.armed_flag = false; try pi.inject(9); try testing.expectEqual(@as(usize, 1013), try s.read(1, 0, &buf)); // The stray Rflush is consumed by the next call. _ = try s.stat(1); try testing.expectEqual(@as(usize, 0), s.client.pending()); } test "connectWatched: a silent server cannot pin the handshake past stop_fd" { // A server that accepts and then never answers: the version handshake has // nothing to read. With a readable stop_fd the connect must come back with // error.Stopped instead of blocking in readSocket (the fd is blocking), and // the caller's descriptor must survive: `Address.fd` is not ours to close. var sv: [2]i32 = undefined; try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.socketpair(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0, &sv))); defer _ = linux.close(sv[0]); defer _ = linux.close(sv[1]); var p: [2]i32 = undefined; try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.pipe2(&p, .{ .CLOEXEC = true, .NONBLOCK = true }))); defer _ = linux.close(p[0]); defer _ = linux.close(p[1]); try testing.expectEqual(@as(usize, 1), linux.write(p[1], "x", 1)); const Probe = struct { const Self = @This(); done: std.atomic.Value(bool) = .init(false), stopped: std.atomic.Value(bool) = .init(false), fd_open: std.atomic.Value(bool) = .init(false), fn run(w: *Self, client: i32, stop: i32) void { if (Session.connectWatched(testing.allocator, .{ .fd = client }, 8192, stop)) |session| { var s = session; s.deinit(); } else |e| w.stopped.store(e == error.Stopped, .release); w.fd_open.store(linux.errno(linux.fcntl(client, linux.F.GETFD, 0)) == .SUCCESS, .release); w.done.store(true, .release); } }; var w: Probe = .{}; const th = try std.Thread.spawn(.{}, Probe.run, .{ &w, sv[0], p[0] }); var waited_ms: usize = 0; while (!w.done.load(.acquire) and waited_ms < 3000) : (waited_ms += 10) { const ts: linux.timespec = .{ .sec = 0, .nsec = 10 * std.time.ns_per_ms }; _ = linux.nanosleep(&ts, null); } try testing.expect(w.done.load(.acquire)); try testing.expect(w.stopped.load(.acquire)); try testing.expect(w.fd_open.load(.acquire)); th.join(); } 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 } })); }