From 3e9f8805f293f622bb885cf849b5ce47dc062ad1 Mon Sep 17 00:00:00 2001 From: Gabriel Schneider Date: Sat, 19 Sep 2026 23:55:47 -0300 Subject: 9ns: --name and /mnt/9p/ mounts, qid.path as inode number, interrupts as Tflush - --name NAME (default derived from the transport: socket basename, tcp-IP-PORT, spawned command, fdN) mounts at /mnt/9p/; --mount still overrides. ensureMountpoint walks down and creates missing components, shadowing the deepest unwritable ancestor. NINE_MOUNT is the only exported variable. - The inode number reported to the kernel is the 9P qid.path for every node, root included; a server handing qid.path 1 to a file (Pardes /self) no longer collides with the root. - FUSE_INTERRUPT for the request in flight becomes Tflush; a blocked read returns EINTR when the server answers the flush, chunked transfers return short counts, other requests arriving meanwhile are stashed and served next. Servers ignoring Tflush still block until they answer. - 9ns-test now covers nine/bridge/fuse; new adv_bridge_interrupt suite (28); 9ns-itest grows to 88 checks. Co-Authored-By: Claude Fable 5.1 --- 9ns/src/bridge.zig | 212 +++++++++++++++++++++++-- 9ns/src/fuse.zig | 56 +++++-- 9ns/src/main.zig | 133 +++++++++++++++- 9ns/src/nine.zig | 442 ++++++++++++++++++++++++++++++++++++++++++++++++----- 9ns/src/ns.zig | 83 ++++++++-- 5 files changed, 835 insertions(+), 91 deletions(-) (limited to '9ns/src') diff --git a/9ns/src/bridge.zig b/9ns/src/bridge.zig index 3a072ec..089c3ca 100644 --- a/9ns/src/bridge.zig +++ b/9ns/src/bridge.zig @@ -4,6 +4,12 @@ //! Everything here is single-threaded and one request at a time. State is three //! tables: inodes (nodeid → fid/qid, deduplicated by qid.path), open handles //! (fh → fid plus a cached directory listing), and the reverse qid map. +//! +//! One request at a time does not mean deaf: while a 9P reply is outstanding +//! the session polls the FUSE descriptor too (`nine.Interrupt`). A +//! FUSE_INTERRUPT for the request being served becomes a Tflush, and if the +//! server honours it the request fails with EINTR; anything else the kernel +//! sends meanwhile is parked in a one-slot stash and served next. const std = @import("std"); const cloud9 = @import("cloud9"); const fuse = @import("fuse.zig"); @@ -89,12 +95,22 @@ pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, session: *nine.Session, root_ defer b.deinit(); b.req_buf = try gpa.alignedAlloc(u8, .@"8", request_buf_len); + b.spare_buf = try gpa.alignedAlloc(u8, .@"8", request_buf_len); b.data_buf = try gpa.alloc(u8, max_write); + // Every read of the FUSE fd follows a poll; non-blocking makes sure a + // request the kernel withdrew in between cannot park us in read(2) while + // a 9P reply is due. + fuse.setNonblocking(fuse_fd) catch return error.FuseIo; + // Abandon any pending 9P reply once the child is gone (stop_fd readable), // including the initial root stat below: a silent server must not pin us. session.stop_fd = stop_fd; defer session.stop_fd = -1; + // And watch the FUSE fd meanwhile: INTERRUPTs become Tflush, other + // requests (INIT arrives during the root stat) wait in the stash. + session.interrupt = b.interruptSource(); + defer session.interrupt = null; // Node 1 is the root; its qid comes from a stat so lookups resolving back to // it (e.g. via a walk) dedupe onto node 1. @@ -115,6 +131,17 @@ pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, session: *nine.Session, root_ .{ .fd = stop_fd, .events = linux.POLL.IN, .revents = 0 }, }; while (true) { + // A request that arrived while a 9P reply was outstanding goes first. + // It lives in the spare buffer; swap so that the spare is free again + // for anything that arrives while this one is being served. + if (b.stash) |req| { + b.stash = null; + std.mem.swap([]align(8) u8, &b.req_buf, &b.spare_buf); + if (!try b.dispatch(req)) return; + continue; + } + if (b.fuse_gone) return; + if (b.fuse_fail) |e| return e; pfds[0].revents = 0; pfds[1].revents = 0; const rc = linux.poll(&pfds, pfds.len, -1); @@ -128,7 +155,8 @@ pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, session: *nine.Session, root_ return; } if (pfds[0].revents == 0) continue; - const req = (fuse.readRequest(fuse_fd, b.req_buf) catch |e| switch (e) { + const req = (fuse.readRequestOnce(fuse_fd, b.req_buf) catch |e| switch (e) { + error.Retry => continue, error.Protocol => return error.FuseProtocol, else => return error.FuseIo, }) orelse { @@ -145,7 +173,23 @@ const Bridge = struct { nine: *nine.Session, opts: Options, req_buf: []align(8) u8 = &.{}, + /// Second request buffer: what the interrupt poll reads into. Holds the + /// stashed request until `serve` swaps it in. + spare_buf: []align(8) u8 = &.{}, data_buf: []u8 = &.{}, + /// A non-INTERRUPT request read while a 9P reply was outstanding (its body + /// points into `spare_buf`). While it is set the FUSE fd is not polled + /// during waits, so a second one cannot arrive. + stash: ?fuse.Request = null, + /// `unique` of the FUSE request being served, if any: the only one an + /// INTERRUPT may cancel. + cur_unique: ?u64 = null, + /// An INTERRUPT for `cur_unique` was consumed: chunked loops stop early + /// even when the flushed reply won the race. Reset per request. + interrupted: bool = false, + /// The FUSE fd reported ENODEV / a failure while the session was waiting. + fuse_gone: bool = false, + fuse_fail: ?error{ FuseIo, FuseProtocol } = null, inodes: std.AutoHashMapUnmanaged(u64, Inode) = .empty, by_qid: std.AutoHashMapUnmanaged(u64, u64) = .empty, handles: std.AutoHashMapUnmanaged(u64, Handle) = .empty, @@ -165,9 +209,65 @@ const Bridge = struct { b.inodes.deinit(b.gpa); b.by_qid.deinit(b.gpa); if (b.req_buf.len != 0) b.gpa.free(b.req_buf); + if (b.spare_buf.len != 0) b.gpa.free(b.spare_buf); if (b.data_buf.len != 0) b.gpa.free(b.data_buf); } + // -- interrupt source (polled by nine.Session while a reply is outstanding) ----- + + fn interruptSource(b: *Bridge) nine.Interrupt { + return .{ .ctx = b, .watch = interruptWatch, .onReadable = interruptReadable, .armed = interruptArmed }; + } + + /// Poll the FUSE fd only while the stash has room: with it full a second + /// request would have nowhere to go. + fn interruptWatch(ctx: *anyopaque) i32 { + const b: *Bridge = @ptrCast(@alignCast(ctx)); + return if (b.stash == null and !b.fuse_gone and b.fuse_fail == null) b.fuse_fd else -1; + } + + /// Reads the request the kernel has ready. An INTERRUPT for the request in + /// flight asks the session to flush it; one for any other request is + /// dropped (the kernel expects no reply); anything else is stashed. + fn interruptReadable(ctx: *anyopaque) nine.Session.Error!bool { + const b: *Bridge = @ptrCast(@alignCast(ctx)); + const req = (fuse.readRequestOnce(b.fuse_fd, b.spare_buf) catch |e| switch (e) { + error.Retry => return false, + error.Protocol => { + b.fuse_fail = error.FuseProtocol; + return error.Stopped; + }, + else => { + b.fuse_fail = error.FuseIo; + return error.Stopped; + }, + }) orelse { + b.trace("fuse fd reports ENODEV while a 9P reply is outstanding", .{}); + b.fuse_gone = true; + return error.Stopped; + }; + const h = req.header; + if (h.op() == .interrupt) { + const in = fuse.body(fuse.InterruptIn, req) catch return false; + const cur = b.cur_unique orelse std.math.maxInt(u64); + if (in.unique == cur) { + b.trace("<- interrupt for unique={d} (in flight): sending Tflush", .{in.unique}); + b.interrupted = true; + return true; + } + b.trace("<- interrupt for unique={d} (not in flight; ignored)", .{in.unique}); + return false; + } + b.trace("<- {s} unique={d} nodeid={d} stashed while a 9P reply is outstanding", .{ opName(h.op()), h.unique, h.nodeid }); + b.stash = req; + return false; + } + + fn interruptArmed(ctx: *anyopaque) bool { + const b: *Bridge = @ptrCast(@alignCast(ctx)); + return b.interrupted; + } + fn trace(b: *const Bridge, comptime fmt: []const u8, args: anytype) void { if (b.opts.debug) std.debug.print("9ns: " ++ fmt ++ "\n", args); } @@ -188,7 +288,12 @@ const Bridge = struct { b.reply(h.unique, &.{}) catch {}; return false; } + b.cur_unique = h.unique; + b.interrupted = false; + defer b.cur_unique = null; b.handle(req) catch |e| { + if (e == error.Stopped and b.fuse_gone) return false; + if (e == error.Stopped and b.fuse_fail != null) return b.fuse_fail.?; const code: linux.E = switch (e) { error.Nine => b.last_err, error.BadRequest => .INVAL, @@ -201,6 +306,7 @@ const Bridge = struct { error.TooLarge => .NAMETOOLONG, error.BadDir => .IO, error.Closed, error.Protocol, error.Io, error.Stopped => .IO, + error.Interrupted => .INTR, error.FuseIo => return error.FuseIo, }; if (wants_reply) try b.replyError(h.unique, code); @@ -565,15 +671,15 @@ const Bridge = struct { while (true) { // A server that ignores the offset would otherwise feed us forever. if (offset >= max_dir_bytes) return error.BadDir; + // A flushed read whose reply still won the race: the listing is + // incomplete either way, so stop here rather than read on. + if (b.interrupted) return error.Interrupted; const n = try b.read(fid, offset, b.data_buf); if (n == 0) break; try parseDirRecords(b.gpa, b.data_buf[0..n], &list); offset += n; } - // Entries carrying the root's own qid.path get the root's ino (1), as GETATTR would report it. - for (list.entries.items[2..]) |*e| if (e.ino == b.root_path) { - e.ino = fuse.root_id; - }; + for (list.entries.items[2..]) |*e| e.ino = inoFromPath(e.ino); return list; } @@ -586,21 +692,29 @@ const Bridge = struct { } fn walkName(b: *Bridge, fid: u32, name: []const u8) nine.Session.Error!u32 { - const newfid = b.nine.allocFid(); - _ = b.nine.walk(fid, newfid, &.{name}) catch |e| { - b.nine.freeFid(newfid); - return b.nineErr("walk", fid, e); - }; + const newfid = try b.walkTo(fid, &.{name}); b.trace(" 9p walk fid={d} newfid={d} name={s} -> ok", .{ fid, newfid, name }); return newfid; } fn clone(b: *Bridge, fid: u32) nine.Session.Error!u32 { - const newfid = b.nine.clone(fid) catch |e| return b.nineErr("clone", fid, e); + const newfid = try b.walkTo(fid, &.{}); b.trace(" 9p walk fid={d} newfid={d} (clone) -> ok", .{ fid, newfid }); return newfid; } + /// allocFid + walk. On Rerror the new fid was never bound; after an + /// interruption the server may or may not have bound it (the Rflush + /// tells us only that no reply is coming), so it is clunked to be sure. + fn walkTo(b: *Bridge, fid: u32, names: []const []const u8) nine.Session.Error!u32 { + const newfid = b.nine.allocFid(); + _ = b.nine.walk(fid, newfid, names) catch |e| { + if (e == error.Interrupted) b.clunkQuiet(newfid) else b.nine.freeFid(newfid); + return b.nineErr("walk", fid, e); + }; + return newfid; + } + fn open9(b: *Bridge, fid: u32, mode: u8) nine.Session.Error!nine.Session.Open { const o = b.nine.open(fid, mode) catch |e| return b.nineErr("open", fid, e); b.trace(" 9p open fid={d} mode={d} -> iounit={d}", .{ fid, mode, o.iounit }); @@ -674,12 +788,18 @@ const Bridge = struct { // -- attrs --------------------------------------------------------------------------- + /// The inode number reported to the kernel is the 9P qid.path, for the root + /// too: FUSE only needs the root's *nodeid* to be 1, and a server may hand + /// qid.path 1 to some other file (Pardes gives it to /self), which would + /// otherwise make `find` see a directory cycle. qid.path 0 maps to a + /// sentinel because inode 0 is treated as invalid by much of userland. fn inoOf(b: *const Bridge, nodeid: u64, qid: cloud9.Qid) u64 { - return if (nodeid == fuse.root_id or qid.path == b.root_path) fuse.root_id else qid.path; + _ = b; + _ = nodeid; + return inoFromPath(qid.path); } fn inoOfNode(b: *const Bridge, nodeid: u64) u64 { - if (nodeid == fuse.root_id) return fuse.root_id; const ino = b.inodes.get(nodeid) orelse return nodeid; return b.inoOf(nodeid, ino.qid); } @@ -700,6 +820,11 @@ const Bridge = struct { // -- pure helpers (unit-tested) ------------------------------------------------------------ /// Attr from a 9P Stat: DMDIR → S_IFDIR else S_IFREG, low 9 permission bits kept. +/// qid.path → inode number; 0 becomes a sentinel (inode 0 reads as "invalid" to many tools). +pub fn inoFromPath(path: u64) u64 { + return if (path == 0) std.math.maxInt(u64) - 1 else path; +} + pub fn attrFromStat(st: cloud9.Stat, ino: u64, uid: u32, gid: u32) fuse.Attr { const ftype: u32 = if (st.mode & cloud9.dmdir != 0) fuse.S_IFDIR else fuse.S_IFREG; return .{ @@ -964,6 +1089,67 @@ test "dirent names the kernel would reject are dropped from listings" { try testing.expectEqualStrings("also", list.entries.items[1].name); } +test "interrupt source: INTERRUPT in flight → Tflush, other INTERRUPTs ignored, requests stashed" { + var pi = try nine.PipeInterrupt.init(0); // only its fake FUSE fd and inject() are used + defer pi.deinit(); + var fs: nine.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; + var b: Bridge = .{ .gpa = testing.allocator, .fuse_fd = pi.read_end, .nine = s, .opts = .{ .uid = 0, .gid = 0 } }; + defer b.deinit(); + b.spare_buf = try testing.allocator.alignedAlloc(u8, .@"8", request_buf_len); + s.interrupt = b.interruptSource(); + defer s.interrupt = null; + try testing.expectEqual(pi.read_end, Bridge.interruptWatch(&b)); + + // Serving unique 7. An INTERRUPT for 6 is already queued (ignored); the + // server fires the one for 7 (pi.unique) once the read at offset 0 hangs. + var buf: [100]u8 = undefined; + b.cur_unique = 7; + pi.unique = 7; + try pi.inject(6); + try testing.expectError(error.Interrupted, b.read(1, 0, &buf)); + try testing.expect(b.interrupted); + 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)); + try testing.expect(b.stash == null); + // The session is intact: a clunk-style cleanup rpc and a further read work. + b.interrupted = false; + b.cur_unique = 8; + try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40])); + try testing.expectEqual(@as(usize, 0), s.client.pending()); + + // A FORGET arriving during a wait is stashed, and the fd is then not watched. + var wire: [48]u8 = undefined; + const hdr = fuse.InHeader{ .len = 48, .opcode = @intFromEnum(fuse.Opcode.forget), .unique = 99, .nodeid = 5, .uid = 0, .gid = 0, .pid = 0, .total_extlen = 0, .padding = 0 }; + @memcpy(wire[0..40], std.mem.asBytes(&hdr)); + @memcpy(wire[40..48], std.mem.asBytes(&fuse.ForgetIn{ .nlookup = 1 })); + try testing.expectEqual(@as(usize, 48), linux.write(pi.write_end, &wire, wire.len)); + fs.read_delay_ns = 30 * std.time.ns_per_ms; + b.cur_unique = 9; + try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40])); + try testing.expect(!b.interrupted); + const stashed = b.stash orelse return error.TestUnexpectedResult; + try testing.expectEqual(fuse.Opcode.forget, stashed.header.op()); + try testing.expectEqual(@as(u64, 5), stashed.header.nodeid); + try testing.expectEqual(@as(u64, 1), (try fuse.body(fuse.ForgetIn, stashed)).nlookup); + try testing.expectEqual(@as(i32, -1), Bridge.interruptWatch(&b)); + // With the stash full an INTERRUPT is not even looked at. + try pi.inject(9); + try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40])); + try testing.expect(!b.interrupted); + b.stash = null; + try testing.expectEqual(pi.read_end, Bridge.interruptWatch(&b)); + // Once the stash is served the queued INTERRUPT is consumed (and ignored: + // its request is not the one in flight any more). + b.cur_unique = 10; + fs.read_delay_ns = 30 * std.time.ns_per_ms; + try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40])); + try testing.expect(!b.interrupted); + try testing.expect(b.stash == null); +} + test "DirList frees its names" { var list: DirList = .{}; try list.entries.append(testing.allocator, .{ .name = try testing.allocator.dupe(u8, "x"), .ino = 1, .dtype = fuse.DT_REG }); diff --git a/9ns/src/fuse.zig b/9ns/src/fuse.zig index 682b3d5..216e616 100644 --- a/9ns/src/fuse.zig +++ b/9ns/src/fuse.zig @@ -314,7 +314,7 @@ comptime { // Request / reply helpers // --------------------------------------------------------------------------- -pub const Error = error{ Protocol, Io, TooManyPayloads }; +pub const Error = error{ Protocol, Io, TooManyPayloads, Retry }; pub const Request = struct { header: InHeader, @@ -327,22 +327,42 @@ pub const Request = struct { /// `max_write + 4096` bytes and 8-byte aligned so `body()` can view it. pub fn readRequest(fd: i32, buf: []u8) Error!?Request { while (true) { - const rc = linux.read(fd, buf.ptr, buf.len); - switch (linux.errno(rc)) { - .SUCCESS => { - const n: usize = rc; - if (n < @sizeOf(InHeader)) return error.Protocol; - const header = std.mem.bytesToValue(InHeader, buf[0..@sizeOf(InHeader)]); - if (header.len != n) return error.Protocol; - return .{ .header = header, .body = buf[@sizeOf(InHeader)..n] }; - }, - .INTR, .AGAIN, .NOENT => continue, - .NODEV => return null, - else => return error.Io, - } + return readRequestOnce(fd, buf) catch |e| switch (e) { + error.Retry => continue, + else => return e, + }; + } +} + +/// One `read(2)` attempt: like `readRequest` but EINTR/EAGAIN/ENOENT surface as +/// `error.Retry` instead of being retried, so a caller that only reads after +/// `poll` (or on a non-blocking fd) never blocks in here. +pub fn readRequestOnce(fd: i32, buf: []u8) Error!?Request { + const rc = linux.read(fd, buf.ptr, buf.len); + switch (linux.errno(rc)) { + .SUCCESS => { + const n: usize = rc; + if (n < @sizeOf(InHeader)) return error.Protocol; + const header = std.mem.bytesToValue(InHeader, buf[0..@sizeOf(InHeader)]); + if (header.len != n) return error.Protocol; + return .{ .header = header, .body = buf[@sizeOf(InHeader)..n] }; + }, + .INTR, .AGAIN, .NOENT => return error.Retry, + .NODEV => return null, + else => return error.Io, } } +/// Sets O_NONBLOCK on `fd` so that a read after `poll` cannot block when the +/// kernel withdrew the request in between (a killed waiter, for instance). +pub fn setNonblocking(fd: i32) Error!void { + const cur = linux.fcntl(fd, linux.F.GETFL, 0); + if (linux.errno(cur) != .SUCCESS) return error.Io; + const nonblock: u32 = @bitCast(linux.O{ .NONBLOCK = true }); + const rc = linux.fcntl(fd, linux.F.SETFL, @as(usize, cur) | nonblock); + if (linux.errno(rc) != .SUCCESS) return error.Io; +} + /// Maximum number of payload slices a single `reply` can carry. pub const max_payloads = 7; @@ -650,4 +670,12 @@ test "readRequest parses one request from a pipe and rejects bad lengths" { std.mem.bytesAsValue(InHeader, bad[0..40]).len = 40; try testing.expectEqual(@as(usize, 48), linux.write(fds[1], &bad, bad.len)); try testing.expectError(error.Protocol, readRequest(fds[0], &buf)); + + // Non-blocking and empty: readRequestOnce reports Retry instead of waiting. + try setNonblocking(fds[0]); + try testing.expectError(error.Retry, readRequestOnce(fds[0], &buf)); + try testing.expectEqual(@as(usize, 48), linux.write(fds[1], &wire, wire.len)); + const again = (try readRequestOnce(fds[0], &buf)) orelse return error.Io; + try testing.expectEqual(@as(u64, 3), again.header.unique); + try testing.expectError(error.Retry, readRequestOnce(fds[0], &buf)); } diff --git a/9ns/src/main.zig b/9ns/src/main.zig index 26bd699..076aa42 100644 --- a/9ns/src/main.zig +++ b/9ns/src/main.zig @@ -21,7 +21,9 @@ const usage_text = \\ --fd N already-connected inherited descriptor \\ --spawn CMD run CMD (via /bin/sh -c) with a socketpair on its stdin/stdout \\Options: - \\ --mount PATH mountpoint inside the new namespace (default /mnt/9p) + \\ --name NAME mount name: the tree appears at /mnt/9p/NAME (one path + \\ component; default derived from the transport, see below) + \\ --mount PATH mountpoint inside the new namespace (overrides --name) \\ --uname NAME 9P user name (default $USER, else "none") \\ --aname NAME 9P tree to attach (default "") \\ --msize BYTES maximum 9P message size to request (default 131072) @@ -30,9 +32,17 @@ const usage_text = \\ --debug trace FUSE and 9P operations on stderr \\ --help, --version \\PROGRAM defaults to $SHELL (else /bin/sh). The mountpoint is exported as $NINE_MOUNT. + \\Default name: --unix PATH -> basename of PATH without .sock/.9p/.socket; + \\--tcp IP:PORT -> tcp-IP-PORT (':' becomes '-'); --spawn CMD -> basename of its + \\first word; --fd N -> fdN; 9p when nothing usable comes out of that. \\ ; +/// Where `--name NAME` mounts: `mount_root/NAME`. +const mount_root = "/mnt/9p"; +/// Name used when nothing usable can be derived from the transport. +const fallback_name = "9p"; + const own_failure: u8 = 125; /// Largest 9P message size we agree to request: the session allocates two /// buffers of this size up front, before the server negotiates it down. @@ -55,7 +65,10 @@ fn printStdout(text: []const u8) void { const Config = struct { address: ?nine.Address = null, spawn_cmd: ?[]const u8 = null, - mount: []const u8 = "/mnt/9p", + /// `--mount`: wins over `name` when set. + mount: ?[]const u8 = null, + /// `--name`: null means "derive from the transport" (see `defaultName`). + name: ?[]const u8 = null, uname: ?[]const u8 = null, aname: []const u8 = "", msize: u32 = 131072, @@ -106,7 +119,7 @@ fn parseArgs(arena: std.mem.Allocator, args: []const [:0]const u8) !ParseResult name = arg[0..eq]; inline_value = arg[eq + 1 ..]; } - const Opt = enum { unix, tcp, fd, spawn, mount, uname, aname, msize, cache, @"no-direct-io", debug, help, version, unknown }; + const Opt = enum { unix, tcp, fd, spawn, name, mount, uname, aname, msize, cache, @"no-direct-io", debug, help, version, unknown }; const opt = std.meta.stringToEnum(Opt, name[2..]) orelse .unknown; switch (opt) { .@"no-direct-io", .debug, .help, .version => if (inline_value != null) return usageError("{s} takes no value", .{name}), @@ -142,6 +155,10 @@ fn parseArgs(arena: std.mem.Allocator, args: []const [:0]const u8) !ParseResult cfg.spawn_cmd = value; transports += 1; }, + .name => { + if (!validName(value)) return usageError("--name wants a single path component (not empty, no '/', not . or ..), got '{s}'", .{value}); + cfg.name = value; + }, .mount => { if (value.len == 0) return usageError("--mount wants a path", .{}); cfg.mount = value; @@ -183,6 +200,52 @@ fn parseTcp(spec: []const u8) ?nine.Address { return .{ .tcp = .{ .host = host, .port = port } }; } +/// A mount name is one path component: non-empty, no '/', no NUL, not `.` +/// or `..`. +fn validName(name: []const u8) bool { + if (name.len == 0) return false; + if (std.mem.eql(u8, name, ".") or std.mem.eql(u8, name, "..")) return false; + for (name) |c| if (c == '/' or c == 0) return false; + return true; +} + +/// The mount name derived from the transport when `--name` is absent: +/// `--unix PATH` → basename of PATH without a trailing `.sock`/`.9p`/ +/// `.socket`; `--tcp IP:PORT` → `tcp-IP-PORT` with every ':' turned into +/// '-' (so an IPv6 literal stays one component); `--spawn CMD` → basename +/// of CMD's first word; `--fd N` → `fdN`. Anything that does not come out +/// as a valid name (empty basename, `..`, ...) becomes `9p`. The result is +/// written into `buf` (at most `buf.len` bytes; longer inputs fall back). +fn defaultName(buf: []u8, cfg: Config) []const u8 { + const raw: []const u8 = blk: { + if (cfg.spawn_cmd) |cmd| { + var words = std.mem.tokenizeAny(u8, cmd, " \t\r\n"); + break :blk std.fs.path.basename(words.next() orelse ""); + } + switch (cfg.address orelse return fallback_name) { + .unix => |path| { + const base = std.fs.path.basename(path); + inline for (.{ ".sock", ".socket", ".9p" }) |ext| { + if (base.len > ext.len and std.mem.endsWith(u8, base, ext)) break :blk base[0 .. base.len - ext.len]; + } + break :blk base; + }, + .tcp => |t| { + const text = std.fmt.bufPrint(buf, "tcp-{s}-{d}", .{ t.host, t.port }) catch return fallback_name; + std.mem.replaceScalar(u8, text, ':', '-'); + return if (validName(text)) text else fallback_name; + }, + .fd => |fd| { + const text = std.fmt.bufPrint(buf, "fd{d}", .{fd}) catch return fallback_name; + return text; + }, + } + }; + if (!validName(raw) or raw.len > buf.len) return fallback_name; + @memcpy(buf[0..raw.len], raw); + return buf[0..raw.len]; +} + /// `--spawn`: run CMD under /bin/sh with one end of a socketpair as its /// stdin/stdout; the other end is the 9P transport. const Server = struct { pid: i32, fd: i32 }; @@ -298,8 +361,15 @@ pub fn main(init: std.process.Init) !u8 { cfg.program = try arena.dupe([]const u8, &.{shell}); } const uname = cfg.uname orelse ns.getenv(envp, "USER") orelse "none"; - const mountpoint = ns.resolveMountpoint(gpa, cfg.mount) catch |err| { - std.debug.print("9ns: --mount {s}: {t}\n", .{ cfg.mount, err }); + // `--mount PATH` wins; otherwise `/mnt/9p/` with `--name` or a + // name derived from the transport. + var name_buf: [512]u8 = undefined; + const mount_arg: []const u8 = cfg.mount orelse blk: { + const name = cfg.name orelse defaultName(&name_buf, cfg); + break :blk try std.fmt.allocPrint(arena, mount_root ++ "/{s}", .{name}); + }; + const mountpoint = ns.resolveMountpoint(gpa, mount_arg) catch |err| { + std.debug.print("9ns: --mount {s}: {t}\n", .{ mount_arg, err }); return own_failure; }; defer gpa.free(mountpoint); @@ -405,7 +475,21 @@ test "parseArgs" { const r = try parseArgs(arena, &args); try std.testing.expectEqual(@as(i32, 3), r.run.address.?.fd); try std.testing.expectEqual(@as(usize, 0), r.run.program.len); - try std.testing.expectEqualStrings("/mnt/9p", r.run.mount); + try std.testing.expect(r.run.mount == null); + try std.testing.expect(r.run.name == null); + } + { + const named = [_][:0]const u8{ "9ns", "--fd", "3", "--name", "bar", "--mount=/x" }; + const r = try parseArgs(arena, &named); + try std.testing.expectEqualStrings("bar", r.run.name.?); + try std.testing.expectEqualStrings("/x", r.run.mount.?); + const eq = [_][:0]const u8{ "9ns", "--fd", "3", "--name=baz" }; + try std.testing.expectEqualStrings("baz", (try parseArgs(arena, &eq)).run.name.?); + // Invalid names: a path, empty, . and .. + for ([_][:0]const u8{ "a/b", "", ".", "..", "/" }) |bad| { + const args = [_][:0]const u8{ "9ns", "--fd", "3", "--name", bad }; + try std.testing.expectEqual(@as(u8, 125), (try parseArgs(arena, &args)).exit); + } } { // Two transports, no transport, unknown option, missing value: all 125. @@ -439,6 +523,43 @@ test "parseArgs" { } } +test "defaultName" { + var buf: [512]u8 = undefined; + const Case = struct { cfg: Config, want: []const u8 }; + const cases = [_]Case{ + .{ .cfg = .{ .address = .{ .unix = "/tmp/9debug.sock" } }, .want = "9debug" }, + .{ .cfg = .{ .address = .{ .unix = "/run/user/1000/acme" } }, .want = "acme" }, + .{ .cfg = .{ .address = .{ .unix = "ramfs.9p" } }, .want = "ramfs" }, + .{ .cfg = .{ .address = .{ .unix = "/x/y.socket" } }, .want = "y" }, + .{ .cfg = .{ .address = .{ .unix = "/x/.sock" } }, .want = ".sock" }, // the whole name, not empty + .{ .cfg = .{ .address = .{ .unix = "/x/y/" } }, .want = "y" }, + .{ .cfg = .{ .address = .{ .unix = "/" } }, .want = "9p" }, + .{ .cfg = .{ .address = .{ .unix = "/x/.." } }, .want = "9p" }, + .{ .cfg = .{ .address = .{ .tcp = .{ .host = "127.0.0.1", .port = 564 } } }, .want = "tcp-127.0.0.1-564" }, + .{ .cfg = .{ .address = .{ .tcp = .{ .host = "::1", .port = 9999 } } }, .want = "tcp---1-9999" }, + .{ .cfg = .{ .address = .{ .fd = 3 } }, .want = "fd3" }, + .{ .cfg = .{ .spawn_cmd = "/x/9proc-demo --stdio" }, .want = "9proc-demo" }, + .{ .cfg = .{ .spawn_cmd = " ramfs\t-s" }, .want = "ramfs" }, + .{ .cfg = .{ .spawn_cmd = " " }, .want = "9p" }, + .{ .cfg = .{}, .want = "9p" }, + }; + for (cases) |c| try std.testing.expectEqualStrings(c.want, defaultName(&buf, c.cfg)); +} + +test "validName" { + try std.testing.expect(validName("a")); + try std.testing.expect(validName("tcp-127.0.0.1-564")); + try std.testing.expect(validName("...")); + try std.testing.expect(!validName("")); + try std.testing.expect(!validName(".")); + try std.testing.expect(!validName("..")); + try std.testing.expect(!validName("a/b")); + try std.testing.expect(!validName("a\x00b")); +} + test { _ = ns; + _ = @import("nine.zig"); + _ = @import("bridge.zig"); + _ = @import("fuse.zig"); } diff --git a/9ns/src/nine.zig b/9ns/src/nine.zig index 70633e6..c89a343 100644 --- a/9ns/src/nine.zig +++ b/9ns/src/nine.zig @@ -29,8 +29,23 @@ pub const dontcare = cloud9.Stat{ .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, TooLarge, OutOfMemory }; + 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 }; @@ -53,6 +68,9 @@ pub const Session = struct { /// 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, /// Connect to `address`, then negotiate the protocol version. /// `msize` is the maximum message size to ask for (0 = the buffers' size). @@ -117,9 +135,17 @@ pub const Session = struct { } /// 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). 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) { + 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 => { @@ -128,29 +154,52 @@ pub const Session = struct { }, }; try s.flush(); + var flush_tag: ?u16 = null; 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, + 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 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; + switch (try s.wait()) { + .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 (flush_tag == null) { + flush_tag = s.client.submit(.{ .flush = .{ .oldtag = tag } }) catch return error.Protocol; + try s.flush(); + }, + } } } + /// 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 { @@ -245,6 +294,50 @@ pub const Session = struct { return chunkSize(s.client.maxWrite(), s.iounits.get(fid) orelse 0); } + const Ready = enum { socket, cancel }; + + /// Blocks until the socket is readable (`.socket`), the interrupt source + /// wants the request in flight cancelled (`.cancel`), or `stop_fd` fires + /// (`error.Stopped`). Anything the interrupt source consumes without asking + /// for a cancellation simply resumes the wait. + fn wait(s: *Session) 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) return .socket; + const prc = linux.poll(&pfds, @intCast(n), -1); + switch (linux.errno(prc)) { + .SUCCESS => {}, + .INTR, .AGAIN => continue, + else => return error.Io, + } + // 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) { @@ -280,8 +373,13 @@ fn readWith( 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 = try rpcFn(s, .{ .read = .{ .fid = fid, .offset = offset + done, .count = want } }); + 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; @@ -302,29 +400,20 @@ fn writeWith( 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 = try rpcFn(s, .{ .write = .{ .fid = fid, .offset = offset + done, .data = data[done..][0..want] } }); + 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 readSome(fd: i32, stop_fd: i32, buf: []u8) Session.Error!usize { +fn readSocket(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, @@ -339,6 +428,9 @@ fn readSome(fd: i32, stop_fd: i32, buf: []u8) Session.Error!usize { 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 }, @@ -461,6 +553,8 @@ test { } 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")); @@ -524,10 +618,19 @@ const FakeFile = struct { 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); @@ -583,20 +686,60 @@ test "write chunks and stops at a short write" { 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 --------------------------------------------------------- /// 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 { +/// 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, - - 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 { + /// 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, or + /// first the Rread the flush was aimed at and then the Rflush (the race). + on_flush: enum { hangup, rflush, reply_then_rflush } = .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); } @@ -610,6 +753,7 @@ const FakeServer = struct { 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; @@ -649,18 +793,41 @@ const FakeServer = struct { .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] } }); + 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), - // A flush is the test's "hang up now" signal. - .tflush => return, + .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, + .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(); @@ -677,8 +844,199 @@ const FakeServer = struct { 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 "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))); diff --git a/9ns/src/ns.zig b/9ns/src/ns.zig index 6da2c7f..b4342e5 100644 --- a/9ns/src/ns.zig +++ b/9ns/src/ns.zig @@ -198,25 +198,47 @@ fn buildArgv(gpa: Allocator, argv: []const []const u8) ![:null]?[*:0]const u8 { /// Make sure `path` is a directory, inside the *current* mount namespace: /// /// * already a directory → done; -/// * else `mkdir`; on `EACCES`/`EPERM`/`EROFS` shadow the parent directory -/// with a tmpfs that re-exposes every existing entry (bind mounts for -/// directories and files, recreated symlinks) and `mkdir` inside it; +/// * else find the deepest existing ancestor and `mkdir` the missing +/// components under it one by one (`mkdir -p`); the first of them failing +/// with `EACCES`/`EPERM`/`EROFS` (the normal case for `/mnt/9p/` as +/// a plain user) means **shadow that ancestor**: mount a `tmpfs` over it +/// that re-exposes every existing entry (bind mounts for directories and +/// files, recreated symlinks), then create the missing components inside; /// * anything else fails with the errno and a hint. /// +/// So `/mnt/9p/x` on a host without `/mnt/9p` shadows `/mnt` and creates +/// `9p/x`; with a root-owned `/mnt/9p` it shadows `/mnt/9p`; inside a 9ns +/// namespace, where `/mnt/9p` is ours, it just creates `x`. `/` and `/proc` +/// are never shadowed, nor a directory with more than `max_shadow_entries`. +/// /// Every failure prints `9ns: : E` to stderr before /// returning. Meant to be called in the child of `spawn` (or from a /// throwaway namespace: `unshare -Urm`). pub fn ensureMountpoint(gpa: Allocator, path: [:0]const u8) !void { - if (fileType(linux.AT.FDCWD, path, false)) |ft| { + if (try existingKind(path)) |ft| { if (ft == .dir) return; std.debug.print("9ns: mountpoint {s}: exists but is not a directory\n", .{path}); return error.Mountpoint; } - if (fileType(linux.AT.FDCWD, path, true) == .symlink) { - std.debug.print("9ns: mountpoint {s}: dangling symlink\n", .{path}); - return error.Mountpoint; + // Deepest existing ancestor: walk up until something is there. + var base: []const u8 = path; + while (true) { + base = std.fs.path.dirname(base) orelse "/"; + var base_buf: [path_max]u8 = undefined; + const base_z = std.fmt.bufPrintZ(&base_buf, "{s}", .{base}) catch { + std.debug.print("9ns: mountpoint {s}: path too long\n", .{path}); + return error.Mountpoint; + }; + const kind = try existingKind(base_z) orelse continue; + if (kind != .dir) { + std.debug.print("9ns: mountpoint {s}: {s} is not a directory\n", .{ path, base }); + return error.Mountpoint; + } + break; } - const mk = linux.errno(linux.mkdirat(linux.AT.FDCWD, path, 0o755)); + // Missing components, deepest ancestor first. + const missing = path[base.len..]; + const mk = mkdirComponents(path, base.len, missing); switch (mk) { .SUCCESS => return, .ACCES, .PERM, .ROFS => {}, @@ -225,21 +247,20 @@ pub fn ensureMountpoint(gpa: Allocator, path: [:0]const u8) !void { return error.Mountpoint; }, } - const parent = std.fs.path.dirname(path) orelse "/"; - if (std.mem.eql(u8, parent, "/") or isSameDirectory(parent, "/")) { + if (std.mem.eql(u8, base, "/") or isSameDirectory(base, "/")) { std.debug.print("9ns: mkdir {s}: E{t}; refusing to shadow / (pass --mount an existing directory)\n", .{ path, mk }); return error.Mountpoint; } // The shadow rebuilds entries from /proc/self/fd//; a tmpfs // over /proc (or a subtree of it) would take that away from itself. - if (std.mem.eql(u8, parent, "/proc") or std.mem.startsWith(u8, parent, "/proc/")) { - std.debug.print("9ns: mkdir {s}: E{t}; refusing to shadow {s} (pass --mount an existing directory)\n", .{ path, mk, parent }); + if (std.mem.eql(u8, base, "/proc") or std.mem.startsWith(u8, base, "/proc/")) { + std.debug.print("9ns: mkdir {s}: E{t}; refusing to shadow {s} (pass --mount an existing directory)\n", .{ path, mk, base }); return error.Mountpoint; } - const parent_z = try gpa.dupeZ(u8, parent); - defer gpa.free(parent_z); - try shadowDirectory(gpa, parent_z); - switch (linux.errno(linux.mkdirat(linux.AT.FDCWD, path, 0o755))) { + const base_z = try gpa.dupeZ(u8, base); + defer gpa.free(base_z); + try shadowDirectory(gpa, base_z); + switch (mkdirComponents(path, base.len, missing)) { .SUCCESS => {}, else => |e| { std.debug.print("9ns: mkdir {s} (in shadow tmpfs): E{t}\n", .{ path, e }); @@ -248,6 +269,36 @@ pub fn ensureMountpoint(gpa: Allocator, path: [:0]const u8) !void { } } +/// What `path` is, following symlinks: null when nothing is there; an +/// error (reported) for a dangling symlink. +fn existingKind(path: [*:0]const u8) !?FileType { + if (fileType(linux.AT.FDCWD, path, false)) |ft| return ft; + if (fileType(linux.AT.FDCWD, path, true) == .symlink) { + std.debug.print("9ns: mountpoint {s}: dangling symlink\n", .{std.mem.span(path)}); + return error.Mountpoint; + } + return null; +} + +/// `mkdir` each component of `missing` (which is `path[base_len..]`, so +/// it starts with '/') under the existing prefix `path[0..base_len]`, in +/// order. Returns the errno of the first failure (`.SUCCESS` when all were +/// created); `EEXIST` on a component is fine (another process, or a retry). +fn mkdirComponents(path: []const u8, base_len: usize, missing: []const u8) E { + var buf: [path_max]u8 = undefined; + var end: usize = base_len; + var it = std.mem.tokenizeScalar(u8, missing, '/'); + while (it.next()) |comp| { + end += 1 + comp.len; + const prefix = std.fmt.bufPrintZ(&buf, "{s}", .{path[0..end]}) catch return .NAMETOOLONG; + switch (linux.errno(linux.mkdirat(linux.AT.FDCWD, prefix, 0o755))) { + .SUCCESS, .EXIST => {}, + else => |e| return e, + } + } + return .SUCCESS; +} + const FileType = enum { dir, symlink, other }; /// True when both paths resolve (following symlinks, including magic ones -- cgit v1.3