summaryrefslogtreecommitdiff
path: root/src/fuse.zig
diff options
context:
space:
mode:
Diffstat (limited to 'src/fuse.zig')
-rw-r--r--src/fuse.zig2709
1 files changed, 2709 insertions, 0 deletions
diff --git a/src/fuse.zig b/src/fuse.zig
new file mode 100644
index 00000000..311d887b
--- /dev/null
+++ b/src/fuse.zig
@@ -0,0 +1,2709 @@
+//! The `/dev/fuse` transport for pardes's acme control filesystem: wire codec,
+//! mount and unmount through `fusermount3`, one `poll()` thread, and the park
+//! table that turns acme's blocking `event` read into "ask me again later".
+//!
+//! Raw protocol, no libfuse. libfuse is a thread pool, a request dispatcher and
+//! a session lifetime — three things pardes already has and would have to fight.
+//! What is left once those are removed is a struct layout and a read/write loop,
+//! which is this file. It links nothing; the only external program it runs is
+//! the setuid `fusermount3` helper, because an unprivileged process cannot
+//! `mount(2)` in the initial user namespace and that helper exists precisely to
+//! hand back a `/dev/fuse` descriptor for a mount it made on our behalf.
+//!
+//! The whole file is one side of a strict division of labour:
+//!
+//! - `acmefs.zig` owns the semantics and knows nothing about FUSE. It speaks
+//! `Req`/`Reply` and never blocks.
+//! - this file owns the kernel's opinions and knows nothing about panes. It
+//! answers, in place, every request the core has no business seeing (INIT,
+//! FORGET, INTERRUPT, DESTROY and the whole ENOSYS family), and translates
+//! the eleven that remain.
+//! - the host loop (tty/gui) owns the ordering: `retry()` to null, `next()`
+//! to null, one `update()` per request, effects drained in between.
+//!
+//! THREADING. The main thread owns the descriptor for read and for write. The
+//! poll thread never touches its data, never sees a `Req`, and never calls into
+//! the core; it waits for POLLIN, calls the host's wake callback, and then
+//! blocks until the main thread has drained. That last handshake is not
+//! decoration: `poll()` is level triggered, so a poller that re-polls
+//! immediately would spin a core at 100% for as long as one unanswered request
+//! sits in the kernel queue. A host with no threads at all skips `wakeThread`
+//! and drains from its frame poll; it loses wake latency and nothing else.
+//!
+//! BLOCKING. A FUSE server blocks a reader by simply not answering, and that is
+//! the one and only way (the kernel gives no meaning to an EAGAIN reply). So
+//! `Status.again` means "held": the request moves into the park table with its
+//! bytes copied out of the read buffer, and `retry()` offers it back once per
+//! frame until the core has something to say. Two obligations come with that:
+//!
+//! 1. a SIGKILLed reader whose request is never answered ends in
+//! *uninterruptible* sleep (`fuse_dev`'s final `wait_event` is not
+//! killable), so it survives its own kill until we reply. FUSE_INTERRUPT
+//! is the escape hatch and is honoured below.
+//! 2. teardown must answer everything still parked, and must abort the
+//! connection by closing the descriptor before unmounting, or a reader
+//! that raced the shutdown is stuck in D state with nobody left to wake
+//! it.
+//!
+//! Linux only, guarded the way `file_watch.zig` guards inotify: every entry
+//! point returns the inert answer off Linux, so a macOS or web build compiles
+//! and mounts nothing. Only `mount()` can create an `Fs`, so off Linux no other
+//! function in this file is ever reached.
+//!
+//! Verified against `/usr/include/linux/fuse.h` (7.45) and `fs/fuse/{dev,inode,
+//! file,dir,readdir}.c`; the comptime size assertions below turn a header drift
+//! into a compile error rather than a wedged mount nobody can unmount.
+const std = @import("std");
+const builtin = @import("builtin");
+const libc = std.c;
+const linux = std.os.linux;
+const acmefs = @import("acmefs.zig");
+
+/// Everything below the mount is Linux kernel ABI. Off Linux the module still
+/// compiles (it is imported by the shared native shell) and does nothing.
+const supported = builtin.os.tag == .linux;
+
+// ---------------------------------------------------------------------------
+// wire protocol
+// ---------------------------------------------------------------------------
+
+/// The protocol version this server speaks. A mismatch in the *major* aborts
+/// the connection outright (`fuse_init_finish`: `arg->major !=
+/// FUSE_KERNEL_VERSION` -> `ok = false` -> the mount is dead on arrival), so
+/// there is nothing to negotiate there.
+const kernel_version: u32 = 7;
+
+/// The highest minor these structs were checked against (see the module
+/// header). The INIT reply carries `@min(kernel_minor, what the kernel
+/// offered)`: `fuse_init_finish` stores our number as `fc->minor`, and the
+/// kernel then sizes the replies it reads back from us by it (the
+/// `FUSE_COMPAT_*_SIZE` family in `fs/fuse/`), so echoing a *newer* kernel's
+/// minor promises reply fields these structs do not have. Capping costs
+/// nothing: with `flags = 0` no feature depends on the number.
+const kernel_minor: u32 = 45;
+
+/// `fuse_dev_do_read` refuses to hand over a request when the server's read
+/// buffer is smaller than this, and answers the *client* EIO instead: every
+/// syscall through the mount fails and nothing says why.
+const min_read_buffer: usize = 8192;
+
+/// `FUSE_REC_ALIGN`. A dirent record that is not a multiple of 8 desynchronises
+/// the kernel's parse of the rest of the reply, so one bad name turns the whole
+/// directory into garbage rather than into an error.
+const rec_align: usize = 8;
+
+/// `FUSE_NAME_OFFSET` — the fixed part of a `fuse_dirent`, before the name.
+const dirent_name_offset: usize = @sizeOf(fuse_dirent);
+
+fn recAlign(n: usize) usize {
+ return (n + rec_align - 1) & ~(rec_align - 1);
+}
+
+/// The subset of `enum fuse_opcode` this server can receive. Non-exhaustive on
+/// purpose: a newer kernel adds opcodes, and `@enumFromInt` of an unlisted
+/// value into an exhaustive enum is undefined behaviour — the one bug in a
+/// protocol decoder that cannot be diagnosed from the outside.
+const Opcode = enum(u32) {
+ lookup = 1,
+ forget = 2,
+ getattr = 3,
+ setattr = 4,
+ readlink = 5,
+ symlink = 6,
+ mknod = 8,
+ mkdir = 9,
+ unlink = 10,
+ rmdir = 11,
+ rename = 12,
+ link = 13,
+ open = 14,
+ read = 15,
+ write = 16,
+ statfs = 17,
+ release = 18,
+ fsync = 20,
+ setxattr = 21,
+ getxattr = 22,
+ listxattr = 23,
+ removexattr = 24,
+ flush = 25,
+ init = 26,
+ opendir = 27,
+ readdir = 28,
+ releasedir = 29,
+ fsyncdir = 30,
+ getlk = 31,
+ setlk = 32,
+ setlkw = 33,
+ access = 34,
+ create = 35,
+ interrupt = 36,
+ bmap = 37,
+ destroy = 38,
+ ioctl = 39,
+ poll = 40,
+ notify_reply = 41,
+ batch_forget = 42,
+ fallocate = 43,
+ readdirplus = 44,
+ rename2 = 45,
+ lseek = 46,
+ copy_file_range = 47,
+ setupmapping = 48,
+ removemapping = 49,
+ syncfs = 50,
+ tmpfile = 51,
+ statx = 52,
+ copy_file_range_64 = 53,
+ _,
+};
+
+/// `FATTR_SIZE`. The only setattr bit this filesystem reads: without
+/// `FUSE_ATOMIC_O_TRUNC` (which `flags = 0` deliberately does not negotiate)
+/// the kernel strips `O_TRUNC` from the OPEN and issues a separate
+/// `SETATTR(size = 0)`, so this bit *is* how `> file` reaches the core.
+const FATTR_SIZE: u32 = 1 << 3;
+
+/// `FUSE_GETATTR_FH` — says the `fh` field of `fuse_getattr_in` is meaningful.
+/// Reading `fh` without checking it hands the core a stale handle from an
+/// unrelated open.
+const FUSE_GETATTR_FH: u32 = 1 << 0;
+
+/// `FOPEN_DIRECT_IO`. Without it the kernel serves reads out of the page cache
+/// and coalesces them, which for this filesystem is wrong in both directions:
+/// a second `cat` of `index` would return the first one's bytes, and a blocking
+/// `event` read would never reach us at all.
+const FOPEN_DIRECT_IO: u32 = 1 << 0;
+
+const fuse_in_header = extern struct {
+ len: u32,
+ opcode: u32,
+ unique: u64,
+ nodeid: u64,
+ uid: u32,
+ gid: u32,
+ pid: u32,
+ total_extlen: u16,
+ padding: u16,
+};
+
+const fuse_out_header = extern struct {
+ len: u32,
+ @"error": i32,
+ unique: u64,
+};
+
+const fuse_init_in = extern struct {
+ major: u32,
+ minor: u32,
+ max_readahead: u32,
+ flags: u32,
+ flags2: u32,
+ unused: [11]u32,
+};
+
+const fuse_init_out = extern struct {
+ major: u32,
+ minor: u32,
+ max_readahead: u32,
+ flags: u32,
+ max_background: u16,
+ congestion_threshold: u16,
+ max_write: u32,
+ time_gran: u32,
+ max_pages: u16,
+ map_alignment: u16,
+ flags2: u32,
+ max_stack_depth: u32,
+ request_timeout: u16,
+ unused: [11]u16,
+};
+
+const fuse_attr = extern struct {
+ ino: u64,
+ size: u64,
+ blocks: u64,
+ atime: u64,
+ mtime: u64,
+ ctime: u64,
+ atimensec: u32,
+ mtimensec: u32,
+ ctimensec: u32,
+ mode: u32,
+ nlink: u32,
+ uid: u32,
+ gid: u32,
+ rdev: u32,
+ blksize: u32,
+ flags: u32,
+};
+
+const fuse_entry_out = extern struct {
+ nodeid: u64,
+ generation: u64,
+ entry_valid: u64,
+ attr_valid: u64,
+ entry_valid_nsec: u32,
+ attr_valid_nsec: u32,
+ attr: fuse_attr,
+};
+
+const fuse_attr_out = extern struct {
+ attr_valid: u64,
+ attr_valid_nsec: u32,
+ dummy: u32,
+ attr: fuse_attr,
+};
+
+const fuse_getattr_in = extern struct {
+ getattr_flags: u32,
+ dummy: u32,
+ fh: u64,
+};
+
+const fuse_setattr_in = extern struct {
+ valid: u32,
+ padding: u32,
+ fh: u64,
+ size: u64,
+ lock_owner: u64,
+ atime: u64,
+ mtime: u64,
+ ctime: u64,
+ atimensec: u32,
+ mtimensec: u32,
+ ctimensec: u32,
+ mode: u32,
+ unused4: u32,
+ uid: u32,
+ gid: u32,
+ unused5: u32,
+};
+
+const fuse_open_in = extern struct {
+ flags: u32,
+ open_flags: u32,
+};
+
+const fuse_open_out = extern struct {
+ fh: u64,
+ open_flags: u32,
+ backing_id: i32,
+};
+
+const fuse_read_in = extern struct {
+ fh: u64,
+ offset: u64,
+ size: u32,
+ read_flags: u32,
+ lock_owner: u64,
+ flags: u32,
+ padding: u32,
+};
+
+const fuse_write_in = extern struct {
+ fh: u64,
+ offset: u64,
+ size: u32,
+ write_flags: u32,
+ lock_owner: u64,
+ flags: u32,
+ padding: u32,
+};
+
+const fuse_write_out = extern struct {
+ size: u32,
+ padding: u32,
+};
+
+const fuse_release_in = extern struct {
+ fh: u64,
+ flags: u32,
+ release_flags: u32,
+ lock_owner: u64,
+};
+
+const fuse_flush_in = extern struct {
+ fh: u64,
+ unused: u32,
+ padding: u32,
+ lock_owner: u64,
+};
+
+const fuse_forget_in = extern struct {
+ nlookup: u64,
+};
+
+const fuse_batch_forget_in = extern struct {
+ count: u32,
+ dummy: u32,
+};
+
+const fuse_interrupt_in = extern struct {
+ unique: u64,
+};
+
+const fuse_kstatfs = extern struct {
+ blocks: u64,
+ bfree: u64,
+ bavail: u64,
+ files: u64,
+ ffree: u64,
+ bsize: u32,
+ namelen: u32,
+ frsize: u32,
+ padding: u32,
+ spare: [6]u32,
+};
+
+const fuse_statfs_out = extern struct {
+ st: fuse_kstatfs,
+};
+
+/// The `name` array is flexible in C and therefore absent here; this struct IS
+/// `FUSE_NAME_OFFSET`, and `dirent_name_offset` is taken from its size so the
+/// encoder and the kernel cannot disagree about where a name starts.
+const fuse_dirent = extern struct {
+ ino: u64,
+ off: u64,
+ namelen: u32,
+ type: u32,
+};
+
+/// `DT_*` from `linux/dirent.h`, as `fuse_dirent.type` wants them.
+const DT_DIR: u32 = 4;
+const DT_REG: u32 = 8;
+
+/// `S_IFMT` bits. `Reply.Attr.mode` carries permissions only, so the format
+/// nibble is ours to add; a `fuse_attr.mode` with no format bits is a file of
+/// no type and `stat(2)` through the mount returns something no tool expects.
+const S_IFDIR: u32 = 0o040000;
+const S_IFREG: u32 = 0o100000;
+
+// A drifted header is a mount that hangs with no diagnostic, so every struct
+// on the wire asserts its size here. These numbers are `sizeof` from
+// /usr/include/linux/fuse.h at FUSE_KERNEL_MINOR_VERSION 45; they are frozen
+// ABI and are not allowed to change under us silently.
+comptime {
+ std.debug.assert(@sizeOf(fuse_in_header) == 40);
+ std.debug.assert(@sizeOf(fuse_out_header) == 16);
+ std.debug.assert(@sizeOf(fuse_init_in) == 64);
+ std.debug.assert(@sizeOf(fuse_init_out) == 64);
+ std.debug.assert(@sizeOf(fuse_attr) == 88);
+ std.debug.assert(@sizeOf(fuse_entry_out) == 128);
+ std.debug.assert(@sizeOf(fuse_attr_out) == 104);
+ std.debug.assert(@sizeOf(fuse_getattr_in) == 16);
+ std.debug.assert(@sizeOf(fuse_setattr_in) == 88);
+ std.debug.assert(@sizeOf(fuse_open_in) == 8);
+ std.debug.assert(@sizeOf(fuse_open_out) == 16);
+ std.debug.assert(@sizeOf(fuse_read_in) == 40);
+ std.debug.assert(@sizeOf(fuse_write_in) == 40);
+ std.debug.assert(@sizeOf(fuse_write_out) == 8);
+ std.debug.assert(@sizeOf(fuse_release_in) == 24);
+ std.debug.assert(@sizeOf(fuse_flush_in) == 24);
+ std.debug.assert(@sizeOf(fuse_forget_in) == 8);
+ std.debug.assert(@sizeOf(fuse_batch_forget_in) == 8);
+ std.debug.assert(@sizeOf(fuse_interrupt_in) == 8);
+ std.debug.assert(@sizeOf(fuse_kstatfs) == 80);
+ std.debug.assert(@sizeOf(fuse_statfs_out) == 80);
+ std.debug.assert(@sizeOf(fuse_dirent) == 24);
+ // The one field offset the codec depends on beyond struct sizes: the body
+ // of every request starts here, and 40 is a multiple of 8, which is what
+ // lets the parse point a struct at the read buffer instead of copying.
+ std.debug.assert(@sizeOf(fuse_in_header) % rec_align == 0);
+}
+
+// ---------------------------------------------------------------------------
+// the neutral readdir staging format
+// ---------------------------------------------------------------------------
+
+/// How `acmefs` hands a directory listing to this file. The core is protocol
+/// neutral by design, so it must not stage `fuse_dirent`s: those carry an
+/// alignment rule, a cookie rule and a `DT_*` table that are the kernel's
+/// business, not the editor's. It stages this instead, packed and repeated,
+/// little endian, into `State.out`:
+///
+/// node: u64 the acmefs node id of the entry, never 0 (see below)
+/// kind: u8 0 = regular file, 1 = directory
+/// namelen: u8 1..255, never 0
+/// name: [namelen]u8
+///
+/// `node` travels so that the `d_ino` a `getdents64` sees is the same number a
+/// later `stat` reports. Synthesising one here instead would make `find -inum`
+/// and every hardlink-detecting tool lie about this filesystem.
+///
+/// `node` is never 0. It used to be, for the entries under `new/`: those name
+/// panes that do not exist, because acme creates the pane when the name is
+/// LOOKED UP. `new/` now stages nothing at all — every name in it is a
+/// *creating* lookup, so any tool that stats what a readdir reported (`ls -l`,
+/// `find`, tab completion) would make one pane per entry — which is why there
+/// is no longer a sentinel `d_ino` for an unresolved name on the wire.
+///
+/// The core stages entries starting at index `req.off` (the cookie the kernel
+/// echoed back) in a stable order. This encoder assigns cookie `off = req.off +
+/// n + 1` to the nth entry it emits, and may emit only a *prefix* of what was
+/// staged when the kernel's requested `size` runs out — the remainder comes
+/// back as another readdir at the higher cookie, so staging has to be
+/// idempotent per cookie rather than a stream. Zero staged bytes means EOF; it
+/// is not an error, and the kernel stops asking.
+///
+/// No `.` or `..`: the kernel synthesises neither and needs neither, and a
+/// filesystem that emits them has to answer `LOOKUP("..")` too.
+pub const dirent_stage_prefix = 10;
+
+/// Encode staged entries into kernel `fuse_dirent` records. Returns the bytes
+/// written to `out`. Pure: this is where the alignment and cookie rules live,
+/// and it is tested directly.
+fn encodeDirents(out: []u8, staged: []const u8, cookie: u64) usize {
+ var in: usize = 0;
+ var w: usize = 0;
+ var n: u64 = 0;
+ while (in + dirent_stage_prefix <= staged.len) {
+ const node = std.mem.readInt(u64, staged[in..][0..8], .little);
+ const kind = staged[in + 8];
+ const namelen: usize = staged[in + 9];
+ // A zero name length would make the record self-referential (the
+ // kernel would parse the padding as the next entry), and a truncated
+ // record means the core staged something we cannot read. Stop rather
+ // than guess: a short reply is a legal readdir, a malformed one is not.
+ if (namelen == 0 or in + dirent_stage_prefix + namelen > staged.len) break;
+ const name = staged[in + dirent_stage_prefix ..][0..namelen];
+ const record = recAlign(dirent_name_offset + namelen);
+ if (w + record > out.len) break;
+
+ // Written field by field rather than through a struct pointer: `out`
+ // is a caller's slice of unknown alignment, and one @alignCast that is
+ // wrong here is a misaligned store into a kernel-bound buffer.
+ std.mem.writeInt(u64, out[w..][0..8], node, .little);
+ std.mem.writeInt(u64, out[w + 8 ..][0..8], cookie + n + 1, .little);
+ std.mem.writeInt(u32, out[w + 16 ..][0..4], @intCast(namelen), .little);
+ std.mem.writeInt(u32, out[w + 20 ..][0..4], if (kind == 1) DT_DIR else DT_REG, .little);
+ @memcpy(out[w + dirent_name_offset ..][0..namelen], name);
+ // The kernel never shows the padding to anyone, but zeroing it keeps
+ // the wire deterministic, which is what the encoder test asserts on.
+ @memset(out[w + dirent_name_offset + namelen ..][0 .. record - dirent_name_offset - namelen], 0);
+
+ in += dirent_stage_prefix + namelen;
+ w += record;
+ n += 1;
+ }
+ return w;
+}
+
+// ---------------------------------------------------------------------------
+// fusermount3
+// ---------------------------------------------------------------------------
+
+/// The environment variable `fusermount3` reads to find the socket it must send
+/// the `/dev/fuse` descriptor back over. Spelled with the leading underscore in
+/// libfuse (`FUSE_COMMFD_ENV`); it is a private contract between the two
+/// programs, not a user knob.
+const commfd_env = "_FUSE_COMMFD";
+
+/// Where the helper might be. Arch puts it in /usr/bin with /usr/sbin a symlink
+/// to it, Debian derivatives use /usr/bin, and a machine with only libfuse2
+/// installed spells it without the 3 — that binary speaks the same
+/// socketpair/SCM_RIGHTS protocol, so it is a real fallback and not a guess.
+/// Searched by absolute path rather than through PATH because the thing being
+/// executed is setuid root: PATH is attacker-influenced input.
+const fusermount_paths = [_][:0]const u8{
+ "/usr/bin/fusermount3",
+ "/usr/sbin/fusermount3",
+ "/bin/fusermount3",
+ "/sbin/fusermount3",
+ "/usr/local/bin/fusermount3",
+ "/usr/bin/fusermount",
+ "/usr/sbin/fusermount",
+ "/bin/fusermount",
+};
+
+/// The `-o` string. Every option here is a deliberate refusal:
+///
+/// - `fsname`/`subtype` are cosmetic but load bearing: they are what `mount`,
+/// `df` and `/proc/self/mountinfo` show, and an unnamed fuse mount in a bug
+/// report is indistinguishable from anyone else's.
+/// - `nosuid,nodev` are what fusermount3 forces anyway; naming them keeps the
+/// intent in the source rather than in someone else's default.
+/// - NOT `allow_other`: it needs `user_allow_other` in /etc/fuse.conf, which
+/// is commented out on a stock Arch install, and asking for it makes
+/// fusermount3 fail the whole mount instead of ignoring the option. It
+/// would also be wrong — this filesystem executes text on write.
+/// - NOT `default_permissions`: with it the kernel enforces the mode bits we
+/// report, which sounds like a free wall but moves access control from the
+/// core (which knows that `cons` is write-only) into a mode field, so a
+/// wrong nibble in a table becomes an EACCES nobody can explain. Same
+/// reason INIT negotiates no flags: fewer kernel behaviours to honour.
+fn mountOpts(buf: *[128:0]u8) [:0]const u8 {
+ return std.fmt.bufPrintSentinel(buf, "fsname=pardes,subtype=pardes,nosuid,nodev", .{}, 0) catch unreachable;
+}
+
+/// `_FUSE_COMMFD=<n>`, the child's end of the socketpair by number. libfuse
+/// passes the descriptor this way rather than on the command line because
+/// fusermount3 is setuid: its argv is world readable through /proc, its
+/// environment is not.
+fn commfdEnv(buf: *[32:0]u8, fd: c_int) [:0]const u8 {
+ return std.fmt.bufPrintSentinel(buf, commfd_env ++ "={d}", .{fd}, 0) catch unreachable;
+}
+
+/// `fusermount3 -o <opts> -- <mountpoint>`. The `--` is not optional: a
+/// mountpoint that begins with a dash would otherwise be parsed as a flag by a
+/// setuid program.
+fn mountArgv(
+ argv: *[6:null]?[*:0]const u8,
+ prog: [*:0]const u8,
+ opts: [*:0]const u8,
+ mountpoint: [*:0]const u8,
+) void {
+ argv.* = .{ prog, "-o", opts, "--", mountpoint, null };
+}
+
+/// `fusermount3 -u -q -z -- <mountpoint>`. Lazy (`-z`) because the mount may
+/// still have an open descriptor on it — a pane shell that inherited a cwd
+/// inside the mount, say — and a non-lazy unmount would fail with EBUSY and
+/// leave the mount behind for good. Quiet (`-q`) because the common case at
+/// exit is a mount the kernel already tore down, and its complaint would be the
+/// last thing on the user's terminal.
+fn unmountArgv(argv: *[7:null]?[*:0]const u8, prog: [*:0]const u8, mountpoint: [*:0]const u8) void {
+ argv.* = .{ prog, "-u", "-q", "-z", "--", mountpoint, null };
+}
+
+/// CMSG_ALIGN/CMSG_LEN/CMSG_SPACE. Only ever evaluated on the Linux path,
+/// where the alignment is `sizeof(size_t)`; other platforms align control
+/// messages to 4 and would need their own numbers.
+fn cmsgAlign(n: usize) usize {
+ const a: usize = @alignOf(usize);
+ return (n + a - 1) & ~(a - 1);
+}
+fn cmsgLen(n: usize) usize {
+ return cmsgAlign(@sizeOf(libc.cmsghdr)) + n;
+}
+fn cmsgSpace(n: usize) usize {
+ return cmsgAlign(@sizeOf(libc.cmsghdr)) + cmsgAlign(n);
+}
+
+/// Build the child's environment: ours, plus `_FUSE_COMMFD`, minus any
+/// `_FUSE_COMMFD` we inherited. The subtraction matters — `getenv` returns the
+/// *first* match, so an inherited stale entry (pardes launched from inside
+/// something that mounts) would win over the one we just appended and
+/// fusermount3 would send the descriptor to a closed socket.
+fn buildEnv(gpa: std.mem.Allocator, commfd: [:0]const u8) ![]?[*:0]const u8 {
+ var count: usize = 0;
+ while (libc.environ[count] != null) count += 1;
+ const env = try gpa.alloc(?[*:0]const u8, count + 2);
+ var n: usize = 0;
+ for (0..count) |i| {
+ const entry = libc.environ[i].?;
+ if (std.mem.startsWith(u8, std.mem.span(entry), commfd_env ++ "=")) continue;
+ env[n] = entry;
+ n += 1;
+ }
+ env[n] = commfd.ptr;
+ env[n + 1] = null;
+ return env[0 .. n + 2];
+}
+
+/// Resolve the helper once, by absolute path. Doing it in the parent rather
+/// than by chaining execve attempts in the child keeps `argv[0]` honest (it is
+/// what `ps` and fusermount3's own diagnostics print) and turns "fuse3 is not
+/// installed" into its own error instead of an exit status.
+fn findFusermount() ?[:0]const u8 {
+ for (fusermount_paths) |candidate| {
+ if (libc.access(candidate.ptr, libc.X_OK) == 0) return candidate;
+ }
+ return null;
+}
+
+/// fork + execve the helper and wait for it. Not `std.process.Child`: that has
+/// no way to hand a child an arbitrary descriptor, and the entire protocol here
+/// is "the child writes to descriptor N". Everything the child does before
+/// execve is async-signal-safe (close, execve, _exit) because the parent may
+/// well be multithreaded by the time this runs.
+fn spawnHelper(
+ prog: [*:0]const u8,
+ argv: [*:null]const ?[*:0]const u8,
+ envp: [*:null]const ?[*:0]const u8,
+ close_in_child: c_int,
+) !u8 {
+ const pid = libc.fork();
+ if (pid < 0) return error.ForkFailed;
+ if (pid == 0) {
+ // The parent's end of the socketpair. Left open, the parent's recvmsg
+ // could never see EOF when the helper dies without sending anything,
+ // and a refused mount would hang instead of failing.
+ if (close_in_child >= 0) _ = libc.close(close_in_child);
+ _ = libc.execve(prog, argv, envp);
+ // 127 is the shell's convention for "not found". Reachable only when
+ // the binary vanished between the access(2) above and now.
+ libc._exit(127);
+ }
+ var status: c_int = 0;
+ while (true) {
+ const got = libc.waitpid(pid, &status, 0);
+ if (got == pid) break;
+ if (got < 0 and libc.errno(got) == .INTR) continue;
+ // Reaped by somebody else's SIGCHLD handler: the status is gone, and
+ // the descriptor either arrived or it did not. Claim success and let
+ // the recvmsg be the judge.
+ return 0;
+ }
+ // WIFEXITED/WEXITSTATUS spelled out: std has no portable macro, and a
+ // helper killed by a signal is not a helper that refused the mount.
+ if (status & 0x7f != 0) return error.FusermountKilled;
+ return @intCast((status >> 8) & 0xff);
+}
+
+/// Receive the `/dev/fuse` descriptor. fusermount3 sends it as an SCM_RIGHTS
+/// control message alongside exactly one byte of ordinary data, and the byte is
+/// not padding: a control message with no data attached may be dropped, so both
+/// sides are required to send at least one.
+///
+/// `MSG_CMSG_CLOEXEC` is the important flag. Every pane shell is forked from
+/// this process and inherits open descriptors; a bash holding a copy of this
+/// one keeps the FUSE connection alive after pardes exits, and the mount stays
+/// up, unkillable, answering nothing, until that shell dies.
+fn receiveFd(sock: c_int) !c_int {
+ var byte: [1]u8 = undefined;
+ var iov = [1]std.posix.iovec{.{ .base = &byte, .len = 1 }};
+ var control: [cmsgSpace(@sizeOf(c_int))]u8 align(@alignOf(libc.cmsghdr)) = undefined;
+ while (true) {
+ var msg: libc.msghdr = .{
+ .name = null,
+ .namelen = 0,
+ .iov = &iov,
+ .iovlen = 1,
+ .control = &control,
+ .controllen = @intCast(control.len),
+ .flags = 0,
+ };
+ const n = libc.recvmsg(sock, &msg, linux.MSG.CMSG_CLOEXEC);
+ if (n < 0) {
+ if (libc.errno(n) == .INTR) continue;
+ return error.CommSocketFailed;
+ }
+ // EOF: the helper exited without sending anything, which is what a
+ // refused mount looks like from here.
+ if (n == 0) return error.FusermountRefused;
+ if (@as(usize, @intCast(msg.controllen)) < cmsgLen(@sizeOf(c_int))) return error.NoDescriptor;
+ const cmsg: *const libc.cmsghdr = @ptrCast(&control);
+ if (cmsg.level != libc.SOL.SOCKET or cmsg.type != libc.SCM.RIGHTS) return error.NoDescriptor;
+ if (@as(usize, @intCast(cmsg.len)) < cmsgLen(@sizeOf(c_int))) return error.NoDescriptor;
+ var fd: c_int = -1;
+ @memcpy(
+ std.mem.asBytes(&fd),
+ control[cmsgAlign(@sizeOf(libc.cmsghdr))..][0..@sizeOf(c_int)],
+ );
+ if (fd < 0) return error.NoDescriptor;
+ return fd;
+ }
+}
+
+/// `mkdir -p` for the mount point, 0700. The leaf is this process's own pid
+/// directory and the parent is `.../pardes`, which on a fresh machine does not
+/// exist; without the -p the whole feature would switch itself off in silence
+/// on exactly the machines that never used it before. Same shape as
+/// `nested.zig`'s ensureSocketDir, and 0700 for the same reason: what lives
+/// under here takes commands.
+fn ensureDir(path: [:0]const u8) void {
+ var partial: [4096:0]u8 = undefined;
+ if (path.len >= partial.len) return;
+ @memcpy(partial[0 .. path.len + 1], path[0 .. path.len + 1]);
+ for (1..path.len) |i| {
+ if (path[i] != '/') continue;
+ partial[i] = 0;
+ _ = libc.mkdir(partial[0..i :0], 0o700);
+ partial[i] = '/';
+ }
+ _ = libc.mkdir(path, 0o700);
+}
+
+/// Unmount and remove `<dir>/<pid>` for every pid that is gone. A pardes killed
+/// with SIGKILL runs no defer, so its mount outlives it as an ENOTCONN stump
+/// that `ls` reports as a permission error and that nothing else will ever
+/// clean up — the snapshot suite alone would leave one per aborted run.
+/// Bounded: one readdir of a directory only we write to, one kill(0) each.
+/// Mirrors nested.zig's socket sweep deliberately, including the ESRCH rule:
+/// 0 means alive, EPERM means alive and someone else's, only ESRCH is a corpse.
+pub fn sweepStale(dir: []const u8) void {
+ if (comptime !supported) return;
+ var dir_buf: [4096:0]u8 = undefined;
+ const dir_z = std.fmt.bufPrintSentinel(&dir_buf, "{s}", .{dir}, 0) catch return;
+ const d = libc.opendir(dir_z) orelse return;
+ defer _ = libc.closedir(d);
+ const me = libc.getpid();
+ while (libc.readdir(d)) |ent| {
+ const name = std.mem.sliceTo(&ent.name, 0);
+ // Strictly digits: parseInt would accept `+7` and `-7`, and this
+ // function unmounts and removes whatever it answers about.
+ if (name.len == 0) continue;
+ for (name) |ch| if (!std.ascii.isDigit(ch)) break;
+ if (std.mem.indexOfNone(u8, name, "0123456789") != null) continue;
+ const pid = std.fmt.parseInt(libc.pid_t, name, 10) catch continue;
+ if (pid == me) continue;
+ const rc = libc.kill(pid, @enumFromInt(0));
+ if (rc == 0 or libc.errno(rc) != .SRCH) continue;
+ var path_buf: [4096:0]u8 = undefined;
+ const path = std.fmt.bufPrintSentinel(&path_buf, "{s}/{s}", .{ dir, name }, 0) catch continue;
+ // Always ours to remove: the name is a pid under a directory only
+ // pardes writes to, and taking the stump away is the point of a sweep.
+ unmountPath(path, true);
+ }
+}
+
+/// Run the helper's unmount and, when the directory is ours, take it away.
+/// Best effort in both halves: an already-unmounted point makes fusermount3
+/// complain (which -q swallows) and a non-empty one makes rmdir fail, and
+/// neither is worth a diagnostic at exit.
+///
+/// `remove_dir` is not a convenience. The *unmount* is always right — the mount
+/// is ours whoever made the directory — but the *rmdir* is only right for a
+/// point pardes derived itself (`<parent>/<pid>`, which `ensureDir` created).
+/// A `--fs=<dir>` the user named is theirs, and removing it is the same
+/// overreach `sweepStale` is already refused under an explicit `--fs` for.
+fn unmountPath(path: [:0]const u8, remove_dir: bool) void {
+ if (findFusermount()) |prog| {
+ var argv: [7:null]?[*:0]const u8 = undefined;
+ unmountArgv(&argv, prog.ptr, path.ptr);
+ // A minimal environment: the helper wants nothing of ours, and the one
+ // variable that WOULD change its behaviour is the comm descriptor it
+ // must not find here.
+ const envp = [_:null]?[*:0]const u8{null};
+ _ = spawnHelper(prog.ptr, &argv, &envp, -1) catch {};
+ }
+ if (remove_dir) _ = libc.rmdir(path);
+}
+
+// ---------------------------------------------------------------------------
+// the park table
+// ---------------------------------------------------------------------------
+
+/// How many kernel requests may be outstanding at once. Every slot is either in
+/// flight (handed to the core, not yet answered) or parked (the core said
+/// `.again`). In-flight slots are transient — the host answers each request
+/// inside the same drain step — so in practice this counts BLOCKED READERS: one
+/// slot per process sitting on `event` or `log`. A session with 32 of those has
+/// 32 scripts watching it.
+///
+/// Overflow is a refusal, not a queue: `take` answers EAGAIN and the descriptor
+/// keeps being read. See its comment for why the tempting alternative (stop
+/// reading and let the kernel hold the surplus) is a deadlock.
+const max_slots = 32;
+
+/// Bytes of request payload a slot can own. A parked request's `data` cannot go
+/// on borrowing the read buffer (the next `next()` overwrites it), so it is
+/// copied in at parse time when it fits. This covers every payload that can
+/// realistically block: a LOOKUP name is at most 255 bytes and a ctl verb line
+/// or an event write-back is a few dozen. A WRITE larger than this is left
+/// borrowed and answered EAGAIN if the core ever tries to park it — a write is
+/// a transaction in this design and is not supposed to block, and growing this
+/// table by 64 KiB a slot to make an impossible case zero-copy is the wrong
+/// trade.
+const park_data_max = 512;
+
+const Slot = struct {
+ used: bool = false,
+ /// The core answered `.again`; `retry()` will offer it back.
+ parked: bool = false,
+ /// Already offered in this retry round. Reset when a round finds nothing,
+ /// which is what gives every parked request exactly one attempt per frame
+ /// instead of letting the oldest one starve the rest.
+ retried: bool = false,
+ /// `req.data` points into `data` below rather than into the read buffer.
+ copied: bool = false,
+ /// Arrival order, so retries are FIFO: the reader that blocked first is
+ /// offered first.
+ seq: u64 = 0,
+ op: Opcode = @enumFromInt(0),
+ req: acmefs.Req = undefined,
+ data: [park_data_max]u8 = undefined,
+};
+
+// ---------------------------------------------------------------------------
+// Fs
+// ---------------------------------------------------------------------------
+
+pub const Fs = struct {
+ pub const Options = struct {
+ /// Absolute path of the mount point. Absolute because it is handed to a
+ /// setuid program that resolves it against its own cwd, and because the
+ /// unmount at exit must name the same place after any chdir.
+ mount: []const u8,
+ /// The largest WRITE payload the kernel may send in one request, and
+ /// therefore the size of the read buffer. 64 KiB matches what a `cp`
+ /// into `body` will use; smaller only splits the same bytes into more
+ /// round trips.
+ max_write: u32 = 64 * 1024,
+ /// Whether pardes made this directory and may therefore remove it at
+ /// exit. True for the derived `<parent>/<pid>`, false for a
+ /// `--fs=<dir>` the user named. See `unmountPath`.
+ owns_dir: bool = false,
+ };
+
+ gpa: std.mem.Allocator,
+ /// The `/dev/fuse` descriptor. -1 once torn down; every entry point checks
+ /// it, so a double deinit and a post-unmount drain are both no-ops.
+ fd: c_int = -1,
+ /// Set when the connection is gone (ENODEV/ECONNABORTED, or DESTROY).
+ /// `next()` stops reading; replies are still written because a slot may be
+ /// mid-flight and the write simply fails.
+ dead: bool = false,
+ path: [:0]u8,
+ /// Mirrors `Options.owns_dir`; gates the rmdir in `deinit`.
+ owns_dir: bool = false,
+ /// One request per read(2), so this is sized for the largest request that
+ /// exists: header + fuse_write_in + max_write. Below FUSE_MIN_READ_BUFFER
+ /// the kernel refuses to hand over requests at all and answers the client
+ /// EIO. 8-aligned so the parse can point structs at it.
+ buf: []align(8) u8,
+ /// Encoded `fuse_dirent`s. Separate from `buf` because a readdir reply is
+ /// built while its request is still being read from `buf`.
+ dirents: [8192]u8 align(8) = undefined,
+
+ uid: u32,
+ gid: u32,
+ max_write: u32,
+ /// The minor the kernel offered, echoed back at INIT. Kept for the record:
+ /// it is the one number in this file that a future feature would consult.
+ minor: u32 = 0,
+
+ slots: [max_slots]Slot = @splat(.{}),
+ seq: u64 = 0,
+ thread: ?std.Thread = null,
+ /// main -> poller, an `eventfd(2)`. The main thread adds 1 per completed
+ /// drain and the poller's blocking read takes the whole counter in one go,
+ /// which is the "collapse the acknowledgements that piled up while we were
+ /// not waiting" behaviour a pipe needed three functions and a nonblocking
+ /// toggle to fake. Not a condition variable, because the poller is blocked
+ /// in `poll()` most of the time and an fd is the only thing that both
+ /// `poll()` and a blocking read can wait on — which is what lets shutdown
+ /// break it out of either state.
+ ///
+ /// The counter cannot say "stop": a stop and a drain acknowledgement that
+ /// race are summed into one indistinguishable number. `stopping` is the
+ /// sticky half of the signal, and is re-read after every wake; the eventfd
+ /// only ever means "look again". The store/write and read/load pair is a
+ /// release/acquire edge over the eventfd's own wait-queue lock, so a poller
+ /// that observes the increment observes the flag with it.
+ ctl: c_int = -1,
+ stopping: std.atomic.Value(bool) = .init(false),
+ wake_ctx: ?*anyopaque = null,
+ wake_fn: ?*const fn (?*anyopaque) void = null,
+
+ /// Mount, hand out the descriptor, and complete the INIT handshake. On
+ /// return the filesystem is live: the kernel will start sending lookups the
+ /// moment anything touches the directory.
+ pub fn mount(gpa: std.mem.Allocator, opts: Options) !*Fs {
+ if (comptime !supported) return error.Unsupported;
+ if (opts.mount.len == 0 or opts.mount[0] != '/') return error.MountPathNotAbsolute;
+
+ const path = try gpa.dupeZ(u8, opts.mount);
+ errdefer gpa.free(path);
+ ensureDir(path);
+
+ const buf_len = @max(
+ min_read_buffer,
+ @sizeOf(fuse_in_header) + @sizeOf(fuse_write_in) + @as(usize, opts.max_write),
+ );
+ const buf = try gpa.alignedAlloc(u8, .@"8", buf_len);
+ errdefer gpa.free(buf);
+
+ const fd = try mountFusermount(gpa, path);
+ errdefer _ = libc.close(fd);
+
+ const fs = try gpa.create(Fs);
+ errdefer gpa.destroy(fs);
+ fs.* = .{
+ .gpa = gpa,
+ .fd = fd,
+ .path = path,
+ .buf = buf,
+ .uid = libc.getuid(),
+ .gid = libc.getgid(),
+ .max_write = opts.max_write,
+ .owns_dir = opts.owns_dir,
+ };
+ // Still blocking here on purpose: INIT is already queued (fusermount3
+ // completed mount(2) before it sent us the descriptor), and a
+ // non-blocking read would make the handshake a spin loop.
+ try fs.handshake();
+ try fs.setNonblocking();
+ return fs;
+ }
+
+ /// socketpair, fork the setuid helper, take the descriptor it sends back.
+ fn mountFusermount(gpa: std.mem.Allocator, path: [:0]const u8) !c_int {
+ const prog = findFusermount() orelse return error.FusermountMissing;
+ var sv: [2]c_int = undefined;
+ if (libc.socketpair(libc.AF.UNIX, libc.SOCK.STREAM, 0, &sv) != 0) return error.SocketPairFailed;
+ // Both ends close-on-exec first, then the child's end is un-marked just
+ // before the fork. The window in between is what any *other* thread's
+ // fork would inherit, and pane shells are forked with forkpty and
+ // inherit everything open.
+ setCloexec(sv[0]);
+ setCloexec(sv[1]);
+ errdefer _ = libc.close(sv[0]);
+
+ var opts_buf: [128:0]u8 = undefined;
+ var commfd_buf: [32:0]u8 = undefined;
+ const opts = mountOpts(&opts_buf);
+ const commfd = commfdEnv(&commfd_buf, sv[1]);
+
+ const envp = try buildEnv(gpa, commfd);
+ defer gpa.free(envp);
+ var argv: [6:null]?[*:0]const u8 = undefined;
+ mountArgv(&argv, prog.ptr, opts.ptr, path.ptr);
+
+ clearCloexec(sv[1]);
+ const code = spawnHelper(prog.ptr, &argv, @ptrCast(envp.ptr), sv[0]) catch |err| {
+ _ = libc.close(sv[1]);
+ return err;
+ };
+ // Ours to close either way: the child has its own copy, and while we
+ // hold one the recvmsg below can never see EOF when the helper dies.
+ _ = libc.close(sv[1]);
+ if (code == 127) return error.FusermountMissing;
+
+ const fd = try receiveFd(sv[0]);
+ if (code != 0) {
+ _ = libc.close(fd);
+ return error.FusermountFailed;
+ }
+ _ = libc.close(sv[0]);
+ return fd;
+ }
+
+ /// Read the kernel's INIT and answer it. Negotiating nothing is the design:
+ /// every flag is a kernel behaviour we would then have to honour forever,
+ /// and this filesystem wants none of them — no readdirplus (whose ENOSYS
+ /// has no fallback and would fail every getdents), no atomic O_TRUNC (so
+ /// `> file` arrives as a plain SETATTR the core already handles), no POSIX
+ /// or BSD locks (flags = 0 makes the kernel set `no_lock`/`no_flock` and
+ /// answer them itself).
+ fn handshake(fs: *Fs) !void {
+ const n = readFull(fs.fd, fs.buf);
+ if (n < @sizeOf(fuse_in_header) + @sizeOf(fuse_init_in)) return error.InitFailed;
+ const h: *const fuse_in_header = @ptrCast(fs.buf.ptr);
+ if (@as(Opcode, @enumFromInt(h.opcode)) != .init) return error.InitFailed;
+ const in: *const fuse_init_in = @ptrCast(@as([*]align(8) u8, @alignCast(fs.buf.ptr + @sizeOf(fuse_in_header))));
+ // A major mismatch is fatal and there is nothing to negotiate: the
+ // kernel aborts the connection, and answering anyway just delays the
+ // failure to the first syscall through the mount.
+ if (in.major != kernel_version) return error.InitVersion;
+ fs.minor = in.minor;
+
+ const out: fuse_init_out = .{
+ .major = kernel_version,
+ // Capped, not echoed: see `kernel_minor`.
+ .minor = @min(in.minor, kernel_minor),
+ // Zero, not "some readahead": with FOPEN_DIRECT_IO there is no page
+ // cache to read ahead into, and a nonzero value here only invites
+ // the kernel to ask for bytes nobody wanted.
+ .max_readahead = 0,
+ .flags = 0,
+ // Left at zero so the kernel keeps its own defaults; a nonzero
+ // max_background is the one that silently caps concurrency.
+ .max_background = 0,
+ .congestion_threshold = 0,
+ .max_write = fs.max_write,
+ // 1 ns. Timestamps on this filesystem are all zero anyway, but a
+ // time_gran of 0 is not a legal granularity.
+ .time_gran = 1,
+ .max_pages = 0,
+ .map_alignment = 0,
+ .flags2 = 0,
+ .max_stack_depth = 0,
+ // 0 = no server timeout. A timeout would let the kernel abort the
+ // connection while a legitimately parked `event` read waits.
+ .request_timeout = 0,
+ .unused = @splat(0),
+ };
+ fs.answer(h.unique, std.mem.asBytes(&out), &.{});
+ return;
+ }
+
+ fn setNonblocking(fs: *Fs) !void {
+ const flags = libc.fcntl(fs.fd, libc.F.GETFL, @as(c_int, 0));
+ if (flags < 0) return error.FcntlFailed;
+ var o: libc.O = @bitCast(@as(u32, @bitCast(flags)));
+ o.NONBLOCK = true;
+ if (libc.fcntl(fs.fd, libc.F.SETFL, @as(c_int, @bitCast(@as(u32, @bitCast(o))))) < 0)
+ return error.FcntlFailed;
+ }
+
+ /// Answer everything still held, abort the connection, unmount, remove the
+ /// directory. The order is not interchangeable:
+ ///
+ /// 1. reply -ENODEV to every slot, so a reader blocked on `event` gets an
+ /// error rather than being left in uninterruptible sleep.
+ /// 2. close the descriptor, which aborts the connection — the backstop
+ /// for anything that raced step 1, since the kernel then fails every
+ /// pending request itself.
+ /// 3. only then unmount, because a mount whose server is gone is exactly
+ /// what `fusermount3 -u -z` is for.
+ /// 4. remove the directory, but only when pardes made it: the derived
+ /// `<parent>/<pid>` is ours, a `--fs=<dir>` the user named is not.
+ pub fn deinit(fs: *Fs) void {
+ const gpa = fs.gpa;
+ fs.stopThread();
+ if (fs.fd >= 0) {
+ for (&fs.slots) |*s| {
+ if (!s.used) continue;
+ fs.answerErr(s.req.tag, .NODEV);
+ s.* = .{};
+ }
+ _ = libc.close(fs.fd);
+ fs.fd = -1;
+ }
+ if (comptime supported) unmountPath(fs.path, fs.owns_dir);
+ gpa.free(fs.path);
+ gpa.free(fs.buf);
+ gpa.destroy(fs);
+ }
+
+ // -- request pump -------------------------------------------------------
+
+ /// Parse the next pending kernel request, or null when the descriptor is
+ /// drained. Call in a loop until null; the loop is the batch, and one wake
+ /// serves all of it.
+ ///
+ /// The returned `Req.data` borrows storage owned by this `Fs` and is valid
+ /// until the next `next()` call. The core copies whatever it keeps — the
+ /// same rule as `.pty_read`.
+ ///
+ /// Requests the core has no business seeing are answered here and the loop
+ /// continues, so a caller never observes them.
+ ///
+ /// Running this to null is also what acknowledges the batch to the poll
+ /// thread, so a host that stops early keeps the poller waiting and loses
+ /// wake latency until the next frame. It is not a correctness bug — the
+ /// remaining requests simply wait in the kernel — but the loop is the
+ /// contract.
+ pub fn next(fs: *Fs) ?acmefs.Req {
+ if (comptime !supported) return null;
+ // Only the EAGAIN arm below releases the poller, and deliberately so.
+ // Every other null return from here implies `dead`, which is write-once
+ // and means reads on the descriptor are failing: posting would send the
+ // poller back into `poll()` on a still-open fd that reports POLLIN
+ // forever, wake the host, drain to this same null, and spin two threads
+ // at 100%. Parking the poller in `consume()` is the right resting state
+ // for a connection that can never produce work again; `stopThread`
+ // releases it. `fd < 0` is unreachable here, since only `deinit` sets it
+ // and it joins the poller first.
+ if (fs.fd < 0 or fs.dead) return null;
+ while (true) {
+ const n = libc.read(fs.fd, fs.buf.ptr, fs.buf.len);
+ if (n < 0) switch (libc.errno(n)) {
+ .INTR => continue,
+ .AGAIN => {
+ // Drained: release the poller (see `post`).
+ fs.post();
+ return null;
+ },
+ // The request was interrupted or aborted between being queued
+ // and being read; there is nothing to answer.
+ .NOENT => continue,
+ // ENODEV (connection aborted, or we were unmounted from under
+ // ourselves) and ECONNABORTED are terminal. Anything else here
+ // is not a thing /dev/fuse does, and treating the unknown as
+ // terminal beats a loop that reads -1 forever.
+ else => {
+ fs.dead = true;
+ return null;
+ },
+ };
+ if (n == 0) {
+ fs.dead = true;
+ return null;
+ }
+ const total: usize = @intCast(n);
+ // Cannot happen (the kernel writes whole requests) but the parse
+ // below indexes on it.
+ if (total < @sizeOf(fuse_in_header)) continue;
+ if (fs.dispatch(total)) |req| return req;
+ }
+ }
+
+ /// Offer parked requests back, one per call. Call in a loop until null,
+ /// once per frame, before `next()`: the null both ends the round and resets
+ /// it, so every parked request gets exactly one attempt per frame and a
+ /// permanently blocked reader cannot starve the others.
+ pub fn retry(fs: *Fs) ?acmefs.Req {
+ if (comptime !supported) return null;
+ if (fs.fd < 0) return null;
+ var best: ?usize = null;
+ for (&fs.slots, 0..) |*s, i| {
+ if (!s.used or !s.parked or s.retried) continue;
+ if (best == null or s.seq < fs.slots[best.?].seq) best = i;
+ }
+ const i = best orelse {
+ for (&fs.slots) |*s| s.retried = false;
+ return null;
+ };
+ fs.slots[i].retried = true;
+ // In flight again: `reply()` re-parks it if the core still has nothing.
+ fs.slots[i].parked = false;
+ return fs.slots[i].req;
+ }
+
+ /// Write the core's answer, or park the request when it said `.again`.
+ /// Called from the `.fs_reply` effect; `bytes` is the payload resolved by
+ /// `pardes.fsPayload` and is borrowed only for the duration of this call.
+ pub fn reply(fs: *Fs, r: *const acmefs.Reply, bytes: []const u8) void {
+ if (comptime !supported) return;
+ const i = fs.findSlot(r.tag) orelse return; // interrupted, or torn down
+ const s = &fs.slots[i];
+
+ if (r.status == .again) {
+ // The one case a park is refused: a payload too large to have been
+ // copied at parse time still borrows the read buffer, so parking it
+ // would park a dangling slice. EAGAIN is honest — the writer can
+ // retry — and by construction unreachable, since the core answers
+ // writes as transactions and only reads ever block.
+ if (!s.copied and s.req.data.len != 0) {
+ fs.answerErr(s.req.tag, .AGAIN);
+ fs.release(i);
+ return;
+ }
+ s.parked = true;
+ return;
+ }
+
+ if (r.status == .err) {
+ fs.answerErr(s.req.tag, @enumFromInt(if (r.errno == 0) @intFromEnum(libc.E.IO) else r.errno));
+ fs.release(i);
+ return;
+ }
+
+ switch (s.req.op) {
+ .lookup => {
+ const out: fuse_entry_out = .{
+ .nodeid = r.attr.node,
+ // Node ids are never reused in this filesystem (pane
+ // serials are monotonic), which is exactly the condition
+ // for a constant generation to be safe.
+ .generation = 0,
+ // No caching, at all. Every file here changes under the
+ // reader's feet, and a cached negative lookup would make
+ // `new/<name>` (which CREATES a pane) work exactly once.
+ .entry_valid = 0,
+ .attr_valid = 0,
+ .entry_valid_nsec = 0,
+ .attr_valid_nsec = 0,
+ .attr = fs.attr(r.attr, r.attr.node),
+ };
+ fs.answer(s.req.tag, std.mem.asBytes(&out), &.{});
+ },
+ .getattr, .setattr => {
+ const out: fuse_attr_out = .{
+ .attr_valid = 0,
+ .attr_valid_nsec = 0,
+ .dummy = 0,
+ .attr = fs.attr(r.attr, s.req.node),
+ };
+ fs.answer(s.req.tag, std.mem.asBytes(&out), &.{});
+ },
+ .open => {
+ const out: fuse_open_out = .{
+ .fh = r.handle,
+ // Direct IO for files; nothing for directories, where the
+ // flag has no meaning and FOPEN_CACHE_DIR (which we do not
+ // set) is the caching knob. An uncached directory is the
+ // point: `new/` and the pane list change constantly.
+ .open_flags = if (s.op == .opendir) 0 else FOPEN_DIRECT_IO,
+ .backing_id = 0,
+ };
+ fs.answer(s.req.tag, std.mem.asBytes(&out), &.{});
+ },
+ .read => {
+ // Never more than was asked for: a read reply longer than
+ // `size` is a protocol error the kernel answers with EIO.
+ const len = @min(bytes.len, s.req.size);
+ fs.answer(s.req.tag, &.{}, bytes[0..len]);
+ },
+ .readdir => {
+ const room = @min(@as(usize, s.req.size), fs.dirents.len);
+ const len = encodeDirents(fs.dirents[0..room], bytes, s.req.off);
+ fs.answer(s.req.tag, &.{}, fs.dirents[0..len]);
+ },
+ .write => {
+ // The core's own count, not the request size: `data` refusing a
+ // partial grapheme is a real short write, and claiming the
+ // whole request would tell the writer its trailing bytes
+ // landed when they did not. Clamped anyway, because a count
+ // larger than what was offered makes the kernel advance a file
+ // offset past bytes that never existed.
+ const out: fuse_write_out = .{
+ .size = @min(r.written, s.req.size),
+ .padding = 0,
+ };
+ fs.answer(s.req.tag, std.mem.asBytes(&out), &.{});
+ },
+ .release => fs.answer(s.req.tag, &.{}, &.{}),
+ .statfs => {
+ // Synthetic numbers, but not arbitrary ones: `namelen` is what
+ // pathconf(_PC_NAME_MAX) returns and a zero there makes some
+ // tools refuse to create any name at all, and `bsize` is what
+ // `stat` reports as the IO block size.
+ const out: fuse_statfs_out = .{ .st = .{
+ .blocks = 0,
+ .bfree = 0,
+ .bavail = 0,
+ .files = 0,
+ .ffree = 0,
+ .bsize = 4096,
+ .namelen = 255,
+ .frsize = 4096,
+ .padding = 0,
+ .spare = @splat(0),
+ } };
+ fs.answer(s.req.tag, std.mem.asBytes(&out), &.{});
+ },
+ }
+ fs.release(i);
+ }
+
+ /// Translate one request. Null means it was answered here.
+ fn dispatch(fs: *Fs, total: usize) ?acmefs.Req {
+ const h: *const fuse_in_header = @ptrCast(fs.buf.ptr);
+ // Bounded by the header's own length, not just by what the read
+ // returned. They agree on /dev/fuse, and taking the smaller of the two
+ // is what keeps a WRITE from claiming payload it did not bring even if
+ // some future kernel ever pads a request.
+ const end = @min(total, @max(@as(usize, h.len), @sizeOf(fuse_in_header)));
+ const body: []align(8) const u8 = @alignCast(fs.buf[@sizeOf(fuse_in_header)..end]);
+ const op: Opcode = @enumFromInt(h.opcode);
+ switch (op) {
+ // Already answered in the handshake. A second INIT cannot happen;
+ // answering it again is cheaper than a special case that could.
+ .init => {
+ fs.answerErr(h.unique, .INVAL);
+ return null;
+ },
+ // NEVER replied to. The kernel does not track these as pending
+ // requests, so a reply carries a `unique` it will not recognise —
+ // -ENOENT at best, and at worst a reply matched against a *live*
+ // request that happens to share the number. Ignoring the refcount
+ // itself is fine: this filesystem's node table is bounded by the
+ // pane count, so nothing grows.
+ .forget, .batch_forget => return null,
+ // Answer the ORIGINAL with EINTR and drop it. This is the only
+ // thing standing between a SIGKILLed reader of `event` and
+ // permanent uninterruptible sleep: after the fatal signal the
+ // kernel's last wait is not killable, so the process survives its
+ // own kill until this reply lands. No reply to the interrupt
+ // itself — its unique is `original | 1` and the kernel keeps no
+ // pending entry for it, while answering -ENOSYS would switch
+ // interrupts off for the whole connection and take the escape
+ // hatch away.
+ .interrupt => {
+ if (body.len >= @sizeOf(fuse_interrupt_in)) {
+ const in: *const fuse_interrupt_in = @ptrCast(body.ptr);
+ if (fs.findSlot(in.unique)) |i| {
+ fs.answerErr(fs.slots[i].req.tag, .INTR);
+ fs.release(i);
+ }
+ }
+ return null;
+ },
+ // A missing reply here hangs `umount` outright.
+ .destroy => {
+ fs.answer(h.unique, &.{}, &.{});
+ fs.dead = true;
+ return null;
+ },
+ // -ENOSYS rather than an empty reply: the kernel sets `no_flush`
+ // and stops sending them, so this costs one round trip for the
+ // whole connection instead of one per close(2). Nothing here has
+ // buffered state for a flush to commit.
+ .flush => {
+ fs.answerErr(h.unique, .NOSYS);
+ return null;
+ },
+ .lookup => {
+ // The name is the whole body, NUL terminated. An empty name is
+ // not a lookup of anything.
+ const name = std.mem.sliceTo(body, 0);
+ if (name.len == 0) {
+ fs.answerErr(h.unique, .INVAL);
+ return null;
+ }
+ return fs.take(op, .{
+ .tag = h.unique,
+ .op = .lookup,
+ .node = h.nodeid,
+ .data = name,
+ });
+ },
+ .getattr => {
+ const in = fs.arg(fuse_getattr_in, body) orelse return null;
+ return fs.take(op, .{
+ .tag = h.unique,
+ .op = .getattr,
+ .node = h.nodeid,
+ // `fh` is only meaningful with the flag; reading it blind
+ // hands the core a handle from an unrelated open.
+ .handle = if (in.getattr_flags & FUSE_GETATTR_FH != 0) @truncate(in.fh) else 0,
+ });
+ },
+ .setattr => {
+ const in = fs.arg(fuse_setattr_in, body) orelse return null;
+ return fs.take(op, .{
+ .tag = h.unique,
+ .op = .setattr,
+ .node = h.nodeid,
+ .handle = @truncate(in.fh),
+ // The `> file` path, and the only setattr this filesystem
+ // has an opinion about. A truncate to a nonzero length is
+ // not expressible in the core's ABI and is reported as no
+ // truncate at all: the reply still carries the current
+ // attributes, so ftruncate(fd, n) succeeds and changes
+ // nothing, which is what every synthetic file here wants.
+ .truncate = in.valid & FATTR_SIZE != 0 and in.size == 0,
+ });
+ },
+ .open, .opendir => {
+ // The flags are read only to reject a short body: this
+ // filesystem's permission model is the mode bits each synthetic
+ // file reports from GETATTR, which the kernel enforces itself,
+ // so the access mode has nothing left to say here.
+ if (fs.arg(fuse_open_in, body) == null) return null;
+ return fs.take(op, .{
+ .tag = h.unique,
+ .op = .open,
+ .node = h.nodeid,
+ });
+ },
+ .read, .readdir => {
+ const in = fs.arg(fuse_read_in, body) orelse return null;
+ return fs.take(op, .{
+ .tag = h.unique,
+ .op = if (op == .readdir) .readdir else .read,
+ .node = h.nodeid,
+ .handle = @truncate(in.fh),
+ .off = in.offset,
+ .size = in.size,
+ });
+ },
+ .write => {
+ const in = fs.arg(fuse_write_in, body) orelse return null;
+ const payload = body[@sizeOf(fuse_write_in)..];
+ // Trust the header's length over the struct's: a `size` larger
+ // than what arrived would read past the request.
+ const len = @min(@as(usize, in.size), payload.len);
+ return fs.take(op, .{
+ .tag = h.unique,
+ .op = .write,
+ .node = h.nodeid,
+ .handle = @truncate(in.fh),
+ .off = in.offset,
+ .size = @intCast(len),
+ .data = payload[0..len],
+ });
+ },
+ .release, .releasedir => {
+ const in = fs.arg(fuse_release_in, body) orelse return null;
+ return fs.take(op, .{
+ .tag = h.unique,
+ .op = .release,
+ .node = h.nodeid,
+ .handle = @truncate(in.fh),
+ });
+ },
+ .statfs => return fs.take(op, .{
+ .tag = h.unique,
+ .op = .statfs,
+ .node = h.nodeid,
+ }),
+ // Everything else. -ENOSYS is not a shrug: for most of these the
+ // kernel caches the answer and stops asking (`no_access`,
+ // `no_getxattr`, `no_statx`, `no_poll`, `no_lseek`, `no_create`),
+ // so one refusal switches the whole feature off for the connection.
+ // The mutations (mkdir, unlink, rename, link, symlink) are refused
+ // because this tree is generated: its shape follows the pane list
+ // and there is nothing for a user to create or remove in it.
+ // READDIRPLUS is not in this list by accident — it is unreachable,
+ // because INIT never sets FUSE_DO_READDIRPLUS, and it has to stay
+ // that way: its -ENOSYS has NO fallback in the kernel and would
+ // fail every getdents through the mount.
+ else => {
+ fs.answerErr(h.unique, .NOSYS);
+ return null;
+ },
+ }
+ }
+
+ /// Point a request struct at the read buffer. Null (and an EINVAL reply)
+ /// when the kernel sent less than the struct, which cannot happen but would
+ /// otherwise be a read past the buffer.
+ fn arg(fs: *Fs, comptime T: type, body: []align(8) const u8) ?*const T {
+ if (body.len < @sizeOf(T)) {
+ const h: *const fuse_in_header = @ptrCast(fs.buf.ptr);
+ fs.answerErr(h.unique, .INVAL);
+ return null;
+ }
+ return @ptrCast(body.ptr);
+ }
+
+ /// Move a parsed request into a slot and hand it to the caller. Small
+ /// payloads are copied in here so that a later park has stable bytes; a
+ /// large one stays borrowed (see `park_data_max`).
+ ///
+ /// Null (and an EAGAIN reply) when the table is full. That is the whole
+ /// reason `next` reads unconditionally instead of gating on a free slot:
+ /// gating looks like polite backpressure and is a deadlock. With 32 readers
+ /// blocked on `event`, refusing to read the descriptor means the INTERRUPT
+ /// that would free a slot is never read either, so a SIGKILLed reader stays
+ /// in uninterruptible sleep forever and every unrelated `ls` of the mount
+ /// hangs behind it. Reading and answering EAGAIN keeps FORGET, INTERRUPT,
+ /// DESTROY and the ENOSYS family flowing — none of which need a slot — and
+ /// turns "too many blocked readers" into one failed syscall the caller can
+ /// see and retry.
+ fn take(fs: *Fs, op: Opcode, req: acmefs.Req) ?acmefs.Req {
+ const i = fs.freeSlot() orelse {
+ fs.answerErr(req.tag, .AGAIN);
+ return null;
+ };
+ const s = &fs.slots[i];
+ s.* = .{
+ .used = true,
+ .seq = fs.seq,
+ .op = op,
+ .req = req,
+ };
+ fs.seq += 1;
+ if (req.data.len != 0 and req.data.len <= park_data_max) {
+ @memcpy(s.data[0..req.data.len], req.data);
+ s.copied = true;
+ s.req.data = s.data[0..req.data.len];
+ }
+ return s.req;
+ }
+
+ fn freeSlot(fs: *Fs) ?usize {
+ for (&fs.slots, 0..) |*s, i| if (!s.used) return i;
+ return null;
+ }
+
+ fn findSlot(fs: *Fs, tag: u64) ?usize {
+ for (&fs.slots, 0..) |*s, i| if (s.used and s.req.tag == tag) return i;
+ return null;
+ }
+
+ fn release(fs: *Fs, i: usize) void {
+ fs.slots[i] = .{};
+ }
+
+ /// `Reply.Attr` -> `fuse_attr`. `node` is the fallback inode for replies
+ /// that do not name one (a getattr answers about a node the request already
+ /// identified); a zero `st_ino` is a value no filesystem is allowed to
+ /// report and some tools treat it as a deleted entry.
+ fn attr(fs: *const Fs, a: acmefs.Reply.Attr, node: u64) fuse_attr {
+ const ino = if (a.node != 0) a.node else node;
+ return .{
+ .ino = ino,
+ .size = a.size,
+ // 512-byte units, as `stat` wants them. Rounded up so a nonempty
+ // file never reports zero blocks, which `du` reads as a hole.
+ .blocks = (a.size + 511) / 512,
+ .atime = 0,
+ .mtime = 0,
+ .ctime = 0,
+ .atimensec = 0,
+ .mtimensec = 0,
+ .ctimensec = 0,
+ .mode = (if (a.dir) S_IFDIR else S_IFREG) | @as(u32, a.mode),
+ // 2 for a directory (itself and `.`) is what every tool expects;
+ // `find` in particular uses it to decide whether to recurse.
+ .nlink = if (a.dir) 2 else 1,
+ // The mounting user owns everything: without `allow_other` nobody
+ // else can reach the mount at all, and reporting some other owner
+ // would only make `ls -l` lie.
+ .uid = fs.uid,
+ .gid = fs.gid,
+ .rdev = 0,
+ .blksize = 4096,
+ .flags = 0,
+ };
+ }
+
+ // -- reply framing ------------------------------------------------------
+
+ /// One `writev` per reply: header, then the op's fixed out struct, then the
+ /// payload. Split into iovecs rather than assembled in a buffer so that a
+ /// megabyte read out of a pane's text is written straight from the core's
+ /// bytes — the whole point of `Reply.Payload.region`.
+ fn answer(fs: *Fs, unique: u64, fixed: []const u8, payload: []const u8) void {
+ var header: fuse_out_header = .{
+ .len = @intCast(@sizeOf(fuse_out_header) + fixed.len + payload.len),
+ .@"error" = 0,
+ .unique = unique,
+ };
+ var iov: [3]std.posix.iovec_const = undefined;
+ var n: usize = 1;
+ iov[0] = .{ .base = std.mem.asBytes(&header).ptr, .len = @sizeOf(fuse_out_header) };
+ if (fixed.len != 0) {
+ iov[n] = .{ .base = fixed.ptr, .len = fixed.len };
+ n += 1;
+ }
+ if (payload.len != 0) {
+ iov[n] = .{ .base = payload.ptr, .len = payload.len };
+ n += 1;
+ }
+ fs.writeReply(iov[0..n], header.len);
+ }
+
+ /// An error reply is header-only: the kernel checks `nbytes ==
+ /// sizeof(oh)` when `error != 0` and answers -EINVAL otherwise, which
+ /// leaves the original request pending forever.
+ fn answerErr(fs: *Fs, unique: u64, e: libc.E) void {
+ var header: fuse_out_header = .{
+ .len = @sizeOf(fuse_out_header),
+ .@"error" = -@as(i32, @intFromEnum(e)),
+ .unique = unique,
+ };
+ const iov = [1]std.posix.iovec_const{
+ .{ .base = std.mem.asBytes(&header).ptr, .len = @sizeOf(fuse_out_header) },
+ };
+ fs.writeReply(&iov, header.len);
+ }
+
+ fn writeReply(fs: *Fs, iov: []const std.posix.iovec_const, expect: u32) void {
+ if (fs.fd < 0) return;
+ while (true) {
+ const n = libc.writev(fs.fd, iov.ptr, @intCast(iov.len));
+ if (n < 0) switch (libc.errno(n)) {
+ .INTR => continue,
+ // /dev/fuse writes never block, so this is not the usual
+ // EAGAIN; retrying is the only thing that can make progress and
+ // it cannot loop forever because the kernel is not waiting on
+ // us.
+ .AGAIN => continue,
+ // The request is no longer pending: it was interrupted or the
+ // connection was aborted between the read and this write.
+ // Dropping it is correct — there is nothing left to answer.
+ .NOENT => return,
+ else => {
+ fs.dead = true;
+ return;
+ },
+ };
+ // A short write to /dev/fuse is not a thing (the kernel takes the
+ // whole reply or none of it), so this can only mean the reply was
+ // malformed and the request is still pending. Nothing useful is
+ // left to do about it here, and pretending otherwise would hide it.
+ std.debug.assert(@as(u32, @intCast(n)) == expect);
+ return;
+ }
+ }
+
+ // -- poll thread --------------------------------------------------------
+
+ /// Start the one background thread: it waits for POLLIN and calls `wake`.
+ /// It never touches the descriptor's data, never sees a request and never
+ /// calls the core; the host's `wake` is expected to do nothing but post an
+ /// event on the loop, exactly like the inotify thread's.
+ ///
+ /// Optional by design. A host with no threads simply does not call this and
+ /// drains from its frame poll instead; it loses wake latency and nothing
+ /// else, which is what makes the no-parallelism backend work unchanged.
+ pub fn wakeThread(fs: *Fs, ctx: ?*anyopaque, wake: *const fn (?*anyopaque) void) !void {
+ if (comptime !supported) return;
+ if (fs.thread != null) return;
+ // Blocking on purpose: `consume` is a blocking read on this descriptor.
+ // The write side cannot block anyway — an eventfd write only waits for
+ // a counter one short of `maxInt(u64)` to be drained, which is not
+ // reachable at one increment per drain.
+ const efd = libc.eventfd(0, linux.EFD.CLOEXEC);
+ if (efd < 0) return error.EventFdFailed;
+ fs.ctl = efd;
+ fs.wake_ctx = ctx;
+ fs.wake_fn = wake;
+ fs.thread = std.Thread.spawn(.{}, pollLoop, .{fs}) catch |err| {
+ _ = libc.close(efd);
+ fs.ctl = -1;
+ return err;
+ };
+ }
+
+ fn stopThread(fs: *Fs) void {
+ if (comptime !supported) return;
+ // `ctl` and `thread` are set and cleared together, so there is no
+ // descriptor to close on the path where no poller was ever started.
+ const t = fs.thread orelse return;
+ // The flag before the wake, never after: a poller that reads the
+ // increment must not then find `stopping` false and go back to sleep on
+ // a counter nobody will raise again. With this order every state the
+ // poller can be in ends in an exit — the loop condition, the `poll()`
+ // (the eventfd becomes readable) and the blocking wait for a drain
+ // acknowledgement (the read returns) all re-read the flag.
+ fs.stopping.store(true, .release);
+ fs.post();
+ t.join();
+ fs.thread = null;
+ _ = libc.close(fs.ctl);
+ fs.ctl = -1;
+ }
+
+ /// Raise the counter by one: "the descriptor has been drained, you may poll
+ /// again", or during teardown "look at `stopping`". Without the drain half
+ /// of that handshake the poller re-polls a level-triggered descriptor that
+ /// is still readable and spins a core until the main thread catches up;
+ /// with it, one wake serves one batch.
+ fn post(fs: *Fs) void {
+ if (fs.ctl < 0) return;
+ const one: u64 = 1;
+ _ = libc.write(fs.ctl, std.mem.asBytes(&one), @sizeOf(u64));
+ }
+
+ fn pollLoop(fs: *Fs) void {
+ if (comptime !supported) return;
+ while (!fs.stopping.load(.acquire)) {
+ var fds = [2]libc.pollfd{
+ .{ .fd = fs.fd, .events = libc.POLL.IN, .revents = 0 },
+ .{ .fd = fs.ctl, .events = libc.POLL.IN, .revents = 0 },
+ };
+ const rc = libc.poll(&fds, 2, -1);
+ if (rc < 0) {
+ if (libc.errno(rc) == .INTR) continue;
+ return;
+ }
+ // Shutdown, or an acknowledgement for a drain that happened without
+ // us. Take the whole counter and re-poll either way: a leftover
+ // count would make the wait below return instantly and turn the
+ // next wake into a spin.
+ if (fds[1].revents != 0 and fs.consume()) return;
+ if (fds[0].revents & (libc.POLL.ERR | libc.POLL.HUP | libc.POLL.NVAL) != 0) return;
+ if (fds[0].revents & libc.POLL.IN == 0) continue;
+
+ (fs.wake_fn.?)(fs.wake_ctx);
+ // Wait for the main thread to finish the batch. This is the whole
+ // anti-spin mechanism; see `post`.
+ if (fs.consume()) return;
+ }
+ }
+
+ /// Block until the counter is nonzero, then take all of it. True when the
+ /// poller must exit, which is `stopping` and nothing else: the count itself
+ /// carries no meaning beyond "look again".
+ ///
+ /// The read blocks, including on the branch that reached here from a
+ /// `poll()` that only *said* the descriptor was readable. That is safe
+ /// because `stopThread` closes `ctl` after `join()` and never before: an
+ /// eventfd raises neither POLLERR nor POLLHUP, so the one revents value
+ /// that would be readable-but-not-readable is POLLNVAL, and a closed
+ /// descriptor is the only thing that produces it.
+ fn consume(fs: *Fs) bool {
+ var v: u64 = undefined;
+ while (true) {
+ const n = libc.read(fs.ctl, std.mem.asBytes(&v), @sizeOf(u64));
+ // A short read and an EOF do not exist on an eventfd: the read
+ // returns 8 or -1. So anything but EINTR means this descriptor is
+ // not the one we opened, and exiting beats spinning on it.
+ if (n < 0) {
+ if (libc.errno(n) == .INTR) continue;
+ return true;
+ }
+ return fs.stopping.load(.acquire);
+ }
+ }
+};
+
+// ---------------------------------------------------------------------------
+// descriptor flags
+// ---------------------------------------------------------------------------
+
+fn setCloexec(fd: c_int) void {
+ const FD_CLOEXEC: c_int = 1;
+ _ = libc.fcntl(fd, libc.F.SETFD, FD_CLOEXEC);
+}
+
+/// The child of the mount fork must KEEP this descriptor across execve — it is
+/// the whole channel the setuid helper answers on.
+fn clearCloexec(fd: c_int) void {
+ _ = libc.fcntl(fd, libc.F.SETFD, @as(c_int, 0));
+}
+
+fn setNonblock(fd: c_int) void {
+ const flags = libc.fcntl(fd, libc.F.GETFL, @as(c_int, 0));
+ if (flags < 0) return;
+ var o: libc.O = @bitCast(@as(u32, @bitCast(flags)));
+ o.NONBLOCK = true;
+ _ = libc.fcntl(fd, libc.F.SETFL, @as(c_int, @bitCast(@as(u32, @bitCast(o)))));
+}
+
+/// One blocking read, EINTR-safe. Used only for the INIT handshake, where the
+/// descriptor is still blocking; every later read goes through `next()`.
+fn readFull(fd: c_int, buf: []u8) usize {
+ while (true) {
+ const n = libc.read(fd, buf.ptr, buf.len);
+ if (n < 0) {
+ if (libc.errno(n) == .INTR) continue;
+ return 0;
+ }
+ return @intCast(n);
+ }
+}
+
+// ---------------------------------------------------------------------------
+// tests
+// ---------------------------------------------------------------------------
+//
+// No test here mounts anything: a real mount needs the setuid helper, a
+// writable runtime directory and a kernel that will let go of it again, which
+// is a snapshot test's job and not a unit test's. What is testable without a
+// mount is everything that has ever actually been wrong in a FUSE server —
+// struct sizes, dirent alignment, cookies, the INIT reply, the park table, and
+// the argv handed to a setuid program. Those are what follows, driven through a
+// socketpair standing in for /dev/fuse.
+
+const testing = std.testing;
+
+/// Build an `Fs` with no mount, wired to `fd`. The socketpair replaces
+/// /dev/fuse for the codec tests: the kernel's side of the conversation is
+/// written by hand and the reply is read back and compared byte for byte.
+fn testFs(gpa: std.mem.Allocator, fd: c_int) !*Fs {
+ const fs = try gpa.create(Fs);
+ fs.* = .{
+ .gpa = gpa,
+ .fd = fd,
+ .path = try gpa.dupeZ(u8, "/nonexistent"),
+ .buf = try gpa.alignedAlloc(u8, .@"8", min_read_buffer),
+ .uid = 1000,
+ .gid = 1000,
+ .max_write = 4096,
+ };
+ return fs;
+}
+
+fn testFsFree(fs: *Fs) void {
+ const gpa = fs.gpa;
+ gpa.free(fs.path);
+ gpa.free(fs.buf);
+ gpa.destroy(fs);
+}
+
+/// Frame a request the way the kernel does and push it at the server.
+fn pushRequest(fd: c_int, unique: u64, op: Opcode, nodeid: u64, body: []const u8) !void {
+ var buf: [4096]u8 align(8) = undefined;
+ const h: fuse_in_header = .{
+ .len = @intCast(@sizeOf(fuse_in_header) + body.len),
+ .opcode = @intFromEnum(op),
+ .unique = unique,
+ .nodeid = nodeid,
+ .uid = 1000,
+ .gid = 1000,
+ .pid = 1,
+ .total_extlen = 0,
+ .padding = 0,
+ };
+ @memcpy(buf[0..@sizeOf(fuse_in_header)], std.mem.asBytes(&h));
+ @memcpy(buf[@sizeOf(fuse_in_header)..][0..body.len], body);
+ const total = @sizeOf(fuse_in_header) + body.len;
+ try testing.expectEqual(@as(isize, @intCast(total)), libc.write(fd, &buf, total));
+}
+
+/// Read one reply back off the socketpair.
+fn readReply(fd: c_int, buf: []u8) ![]u8 {
+ const n = libc.read(fd, buf.ptr, buf.len);
+ try testing.expect(n >= @sizeOf(fuse_out_header));
+ return buf[0..@intCast(n)];
+}
+
+fn outHeader(bytes: []const u8) fuse_out_header {
+ var h: fuse_out_header = undefined;
+ @memcpy(std.mem.asBytes(&h), bytes[0..@sizeOf(fuse_out_header)]);
+ return h;
+}
+
+/// A socketpair standing in for /dev/fuse. SEQPACKET, not STREAM, and that is
+/// the whole point: the kernel's character device hands over exactly one
+/// request per read(2) and takes exactly one reply per write(2), and a stream
+/// socket would coalesce three requests into one read and let a codec that
+/// ignores `fuse_in_header.len` pass anyway.
+///
+/// Both ends non-blocking. The server's end so `next()` meets EAGAIN where it
+/// would on the real descriptor; the kernel's end so a test can assert that
+/// NOTHING was written — which is what "a held request has no reply" and "a
+/// FORGET is never answered" mean, and a blocking read would simply hang there
+/// instead of failing.
+fn testPair() ![2]c_int {
+ var sv: [2]c_int = undefined;
+ if (libc.socketpair(libc.AF.UNIX, libc.SOCK.SEQPACKET, 0, &sv) != 0) return error.SocketPairFailed;
+ setNonblock(sv[0]);
+ setNonblock(sv[1]);
+ return sv;
+}
+
+test "lookup round trip: parse borrows the name, reply frames an entry" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ try pushRequest(sv[1], 100, .lookup, 1, "index\x00");
+ const req = fs.next() orelse return error.NoRequest;
+ try testing.expectEqual(acmefs.Op.lookup, req.op);
+ try testing.expectEqual(@as(u64, 100), req.tag);
+ try testing.expectEqual(@as(u64, 1), req.node);
+ try testing.expectEqualStrings("index", req.data);
+ // Drained, and nothing else was invented.
+ try testing.expect(fs.next() == null);
+
+ fs.reply(&.{
+ .tag = 100,
+ .attr = .{ .node = 7, .size = 42, .mode = 0o444 },
+ }, &.{});
+
+ var buf: [512]u8 = undefined;
+ const got = try readReply(sv[1], &buf);
+ const h = outHeader(got);
+ try testing.expectEqual(@as(u32, @sizeOf(fuse_out_header) + @sizeOf(fuse_entry_out)), h.len);
+ try testing.expectEqual(@as(u32, @intCast(got.len)), h.len);
+ try testing.expectEqual(@as(i32, 0), h.@"error");
+ try testing.expectEqual(@as(u64, 100), h.unique);
+
+ var entry: fuse_entry_out = undefined;
+ @memcpy(std.mem.asBytes(&entry), got[@sizeOf(fuse_out_header)..][0..@sizeOf(fuse_entry_out)]);
+ try testing.expectEqual(@as(u64, 7), entry.nodeid);
+ // Caching off in both directions, or `new/<name>` creates a pane once and
+ // then serves the cached negative lookup forever.
+ try testing.expectEqual(@as(u64, 0), entry.entry_valid);
+ try testing.expectEqual(@as(u64, 0), entry.attr_valid);
+ try testing.expectEqual(@as(u64, 7), entry.attr.ino);
+ try testing.expectEqual(@as(u64, 42), entry.attr.size);
+ try testing.expectEqual(S_IFREG | @as(u32, 0o444), entry.attr.mode);
+ try testing.expectEqual(@as(u32, 1), entry.attr.nlink);
+ try testing.expectEqual(@as(u32, 1000), entry.attr.uid);
+ // The slot went back.
+ try testing.expect(fs.freeSlot() != null);
+ try testing.expectEqual(@as(?usize, null), fs.findSlot(100));
+}
+
+test "read reply is capped at the requested size and written as one frame" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ const in: fuse_read_in = .{
+ .fh = 3,
+ .offset = 8,
+ .size = 4,
+ .read_flags = 0,
+ .lock_owner = 0,
+ .flags = 0,
+ .padding = 0,
+ };
+ try pushRequest(sv[1], 200, .read, 5, std.mem.asBytes(&in));
+ const req = fs.next() orelse return error.NoRequest;
+ try testing.expectEqual(acmefs.Op.read, req.op);
+ try testing.expectEqual(@as(u32, 3), req.handle);
+ try testing.expectEqual(@as(u64, 8), req.off);
+ try testing.expectEqual(@as(u32, 4), req.size);
+
+ // The core offers more than was asked for; a reply longer than `size` is
+ // answered EIO by the kernel, so it has to be clamped here.
+ fs.reply(&.{ .tag = 200, .payload = .{ .staged = 9 } }, "abcdefghi");
+ var buf: [512]u8 = undefined;
+ const got = try readReply(sv[1], &buf);
+ try testing.expectEqual(@as(usize, @sizeOf(fuse_out_header) + 4), got.len);
+ try testing.expectEqual(@as(u32, @intCast(got.len)), outHeader(got).len);
+ try testing.expectEqualStrings("abcd", got[@sizeOf(fuse_out_header)..]);
+}
+
+test "an error reply is header only" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ try pushRequest(sv[1], 300, .lookup, 1, "nope\x00");
+ _ = fs.next() orelse return error.NoRequest;
+ fs.reply(&.{ .tag = 300, .status = .err, .errno = @intFromEnum(libc.E.NOENT) }, &.{});
+
+ var buf: [512]u8 = undefined;
+ const got = try readReply(sv[1], &buf);
+ // len MUST be exactly the header when error is set; anything else makes the
+ // kernel answer -EINVAL and leaves the request pending forever.
+ try testing.expectEqual(@as(usize, @sizeOf(fuse_out_header)), got.len);
+ const h = outHeader(got);
+ try testing.expectEqual(@as(u32, @sizeOf(fuse_out_header)), h.len);
+ try testing.expectEqual(-@as(i32, @intFromEnum(libc.E.NOENT)), h.@"error");
+}
+
+test "opcodes the core never sees are answered here" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+ var buf: [512]u8 = undefined;
+
+ // FORGET and BATCH_FORGET get NO reply, ever: the kernel keeps no pending
+ // entry for them, so a reply would carry a unique it does not recognise.
+ const forget: fuse_forget_in = .{ .nlookup = 1 };
+ try pushRequest(sv[1], 400, .forget, 7, std.mem.asBytes(&forget));
+ const batch: fuse_batch_forget_in = .{ .count = 0, .dummy = 0 };
+ try pushRequest(sv[1], 402, .batch_forget, 0, std.mem.asBytes(&batch));
+ // ...and a mutation is refused, which is the first thing that produces a
+ // reply, proving nothing was written for the two above.
+ try pushRequest(sv[1], 404, .mkdir, 1, "x\x00");
+ try testing.expect(fs.next() == null);
+
+ const got = try readReply(sv[1], &buf);
+ try testing.expectEqual(@as(usize, @sizeOf(fuse_out_header)), got.len);
+ const h = outHeader(got);
+ try testing.expectEqual(@as(u64, 404), h.unique);
+ try testing.expectEqual(-@as(i32, @intFromEnum(libc.E.NOSYS)), h.@"error");
+
+ // DESTROY must be answered or umount hangs.
+ try pushRequest(sv[1], 406, .destroy, 0, &.{});
+ try testing.expect(fs.next() == null);
+ const destroyed = try readReply(sv[1], &buf);
+ try testing.expectEqual(@as(usize, @sizeOf(fuse_out_header)), destroyed.len);
+ try testing.expectEqual(@as(i32, 0), outHeader(destroyed).@"error");
+ try testing.expectEqual(@as(u64, 406), outHeader(destroyed).unique);
+
+ // FLUSH is refused so the kernel stops sending one per close(2).
+ fs.dead = false;
+ const flush: fuse_flush_in = .{ .fh = 1, .unused = 0, .padding = 0, .lock_owner = 0 };
+ try pushRequest(sv[1], 408, .flush, 1, std.mem.asBytes(&flush));
+ try testing.expect(fs.next() == null);
+ const flushed = try readReply(sv[1], &buf);
+ try testing.expectEqual(-@as(i32, @intFromEnum(libc.E.NOSYS)), outHeader(flushed).@"error");
+}
+
+test "setattr size=0 is the truncate the kernel sends instead of O_TRUNC" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ var in: fuse_setattr_in = std.mem.zeroes(fuse_setattr_in);
+ in.valid = FATTR_SIZE;
+ in.size = 0;
+ try pushRequest(sv[1], 500, .setattr, 9, std.mem.asBytes(&in));
+ const req = fs.next() orelse return error.NoRequest;
+ try testing.expectEqual(acmefs.Op.setattr, req.op);
+ try testing.expect(req.truncate);
+
+ // A nonzero size is not a truncate this ABI can express, and must not be
+ // reported as one: the core would clear a pane on `ftruncate(fd, 10)`.
+ in.size = 10;
+ try pushRequest(sv[1], 502, .setattr, 9, std.mem.asBytes(&in));
+ fs.reply(&.{ .tag = 500, .attr = .{ .node = 9 } }, &.{});
+ const req2 = fs.next() orelse return error.NoRequest;
+ try testing.expect(!req2.truncate);
+
+ fs.reply(&.{ .tag = 502, .attr = .{ .node = 9, .size = 3, .dir = true } }, &.{});
+ var buf: [512]u8 = undefined;
+ _ = try readReply(sv[1], &buf); // the first reply
+ const got = try readReply(sv[1], &buf);
+ var out: fuse_attr_out = undefined;
+ @memcpy(std.mem.asBytes(&out), got[@sizeOf(fuse_out_header)..][0..@sizeOf(fuse_attr_out)]);
+ try testing.expectEqual(S_IFDIR | @as(u32, 0o600), out.attr.mode);
+ try testing.expectEqual(@as(u32, 2), out.attr.nlink);
+ try testing.expectEqual(@as(u64, 0), out.attr_valid);
+}
+
+test "open reports direct io for files and nothing for directories" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+ var buf: [512]u8 = undefined;
+
+ // O_WRONLY
+ const wr: fuse_open_in = .{ .flags = 1, .open_flags = 0 };
+ try pushRequest(sv[1], 600, .open, 4, std.mem.asBytes(&wr));
+ const req = fs.next() orelse return error.NoRequest;
+ try testing.expectEqual(acmefs.Op.open, req.op);
+ fs.reply(&.{ .tag = 600, .handle = 11 }, &.{});
+ var got = try readReply(sv[1], &buf);
+ var open_out: fuse_open_out = undefined;
+ @memcpy(std.mem.asBytes(&open_out), got[@sizeOf(fuse_out_header)..][0..@sizeOf(fuse_open_out)]);
+ try testing.expectEqual(@as(u64, 11), open_out.fh);
+ try testing.expectEqual(FOPEN_DIRECT_IO, open_out.open_flags);
+
+ // O_RDONLY on a directory
+ const rd: fuse_open_in = .{ .flags = 0, .open_flags = 0 };
+ try pushRequest(sv[1], 602, .opendir, 1, std.mem.asBytes(&rd));
+ const dir_req = fs.next() orelse return error.NoRequest;
+ try testing.expectEqual(acmefs.Op.open, dir_req.op);
+ fs.reply(&.{ .tag = 602, .handle = 12 }, &.{});
+ got = try readReply(sv[1], &buf);
+ @memcpy(std.mem.asBytes(&open_out), got[@sizeOf(fuse_out_header)..][0..@sizeOf(fuse_open_out)]);
+ // No FOPEN_CACHE_DIR either: the pane list changes between two `ls`.
+ try testing.expectEqual(@as(u32, 0), open_out.open_flags);
+}
+
+test "write borrows the payload and reports the core's own count" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ var body: [@sizeOf(fuse_write_in) + 5]u8 = undefined;
+ const in: fuse_write_in = .{
+ .fh = 2,
+ .offset = 0,
+ .size = 5,
+ .write_flags = 0,
+ .lock_owner = 0,
+ .flags = 0,
+ .padding = 0,
+ };
+ @memcpy(body[0..@sizeOf(fuse_write_in)], std.mem.asBytes(&in));
+ @memcpy(body[@sizeOf(fuse_write_in)..], "hello");
+ try pushRequest(sv[1], 700, .write, 6, &body);
+ const req = fs.next() orelse return error.NoRequest;
+ try testing.expectEqual(acmefs.Op.write, req.op);
+ try testing.expectEqualStrings("hello", req.data);
+
+ fs.reply(&.{ .tag = 700, .written = 5 }, &.{});
+ var buf: [512]u8 = undefined;
+ var got = try readReply(sv[1], &buf);
+ var out: fuse_write_out = undefined;
+ @memcpy(std.mem.asBytes(&out), got[@sizeOf(fuse_out_header)..][0..@sizeOf(fuse_write_out)]);
+ try testing.expectEqual(@as(u32, 5), out.size);
+
+ // A short count is a real answer — `data` refusing a partial grapheme —
+ // and must reach write(2) as a short write rather than as a full one.
+ try pushRequest(sv[1], 704, .write, 6, &body);
+ _ = fs.next() orelse return error.NoRequest;
+ fs.reply(&.{ .tag = 704, .written = 3 }, &.{});
+ got = try readReply(sv[1], &buf);
+ @memcpy(std.mem.asBytes(&out), got[@sizeOf(fuse_out_header)..][0..@sizeOf(fuse_write_out)]);
+ try testing.expectEqual(@as(u32, 3), out.size);
+
+ // A count larger than what was offered would advance the file offset past
+ // bytes that never existed.
+ try pushRequest(sv[1], 706, .write, 6, &body);
+ _ = fs.next() orelse return error.NoRequest;
+ fs.reply(&.{ .tag = 706, .written = 99 }, &.{});
+ got = try readReply(sv[1], &buf);
+ @memcpy(std.mem.asBytes(&out), got[@sizeOf(fuse_out_header)..][0..@sizeOf(fuse_write_out)]);
+ try testing.expectEqual(@as(u32, 5), out.size);
+}
+
+test "a write whose size lies about the payload is clamped to what arrived" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ var body: [@sizeOf(fuse_write_in) + 2]u8 = undefined;
+ var in: fuse_write_in = std.mem.zeroes(fuse_write_in);
+ in.size = 4096; // more than the two bytes that follow
+ @memcpy(body[0..@sizeOf(fuse_write_in)], std.mem.asBytes(&in));
+ @memcpy(body[@sizeOf(fuse_write_in)..], "hi");
+ try pushRequest(sv[1], 702, .write, 6, &body);
+ const req = fs.next() orelse return error.NoRequest;
+ try testing.expectEqualStrings("hi", req.data);
+ try testing.expectEqual(@as(u32, 2), req.size);
+}
+
+test "dirent encoding: 8-byte records, cookies from the request offset" {
+ var staged: [64]u8 = undefined;
+ var w: usize = 0;
+ // node=2 kind=file name="addr"
+ std.mem.writeInt(u64, staged[w..][0..8], 2, .little);
+ staged[w + 8] = 0;
+ staged[w + 9] = 4;
+ @memcpy(staged[w + 10 ..][0..4], "addr");
+ w += 14;
+ // node=3 kind=dir name="new"
+ std.mem.writeInt(u64, staged[w..][0..8], 3, .little);
+ staged[w + 8] = 1;
+ staged[w + 9] = 3;
+ @memcpy(staged[w + 10 ..][0..3], "new");
+ w += 13;
+
+ var out: [128]u8 = undefined;
+ const n = encodeDirents(&out, staged[0..w], 5);
+ // 24 + 4 -> 32; 24 + 3 -> 32. A record that is not a multiple of 8
+ // desynchronises the kernel's parse of everything after it.
+ try testing.expectEqual(@as(usize, 64), n);
+ try testing.expectEqual(@as(u64, 0), n % rec_align);
+
+ try testing.expectEqual(@as(u64, 2), std.mem.readInt(u64, out[0..8], .little));
+ // Cookies continue from the request's offset: the kernel sends the last
+ // `off` it saw as the next request's offset, so restarting at 1 would loop
+ // the directory forever.
+ try testing.expectEqual(@as(u64, 6), std.mem.readInt(u64, out[8..16], .little));
+ try testing.expectEqual(@as(u32, 4), std.mem.readInt(u32, out[16..20], .little));
+ try testing.expectEqual(DT_REG, std.mem.readInt(u32, out[20..24], .little));
+ try testing.expectEqualStrings("addr", out[24..28]);
+ // Padding zeroed, so the wire is deterministic.
+ try testing.expectEqualSlices(u8, &.{ 0, 0, 0, 0 }, out[28..32]);
+
+ try testing.expectEqual(@as(u64, 3), std.mem.readInt(u64, out[32..40], .little));
+ try testing.expectEqual(@as(u64, 7), std.mem.readInt(u64, out[40..48], .little));
+ try testing.expectEqual(DT_DIR, std.mem.readInt(u32, out[52..56], .little));
+ try testing.expectEqualStrings("new", out[56..59]);
+}
+
+test "dirent encoding stops cleanly when the reply buffer or the staging runs out" {
+ var staged: [64]u8 = undefined;
+ std.mem.writeInt(u64, staged[0..8], 9, .little);
+ staged[8] = 0;
+ staged[9] = 4;
+ @memcpy(staged[10..14], "body");
+ std.mem.writeInt(u64, staged[14..22], 10, .little);
+ staged[22] = 0;
+ staged[23] = 4;
+ @memcpy(staged[24..28], "ctl!");
+
+ // Room for one record only: the second comes back at the higher cookie.
+ var out: [40]u8 = undefined;
+ try testing.expectEqual(@as(usize, 32), encodeDirents(&out, staged[0..28], 0));
+
+ // A truncated staging record is dropped rather than guessed at.
+ try testing.expectEqual(@as(usize, 32), encodeDirents(&out, staged[0..26], 0));
+ // Zero staged bytes is EOF, not an error.
+ try testing.expectEqual(@as(usize, 0), encodeDirents(&out, &.{}, 4));
+ // A zero name length would make the kernel parse the padding as an entry.
+ var bad: [10]u8 = @splat(0);
+ try testing.expectEqual(@as(usize, 0), encodeDirents(&out, &bad, 0));
+}
+
+test "readdir reply carries encoded dirents built from the staged names" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ const in: fuse_read_in = .{
+ .fh = 1,
+ .offset = 0,
+ .size = 4096,
+ .read_flags = 0,
+ .lock_owner = 0,
+ .flags = 0,
+ .padding = 0,
+ };
+ try pushRequest(sv[1], 800, .readdir, 1, std.mem.asBytes(&in));
+ const req = fs.next() orelse return error.NoRequest;
+ try testing.expectEqual(acmefs.Op.readdir, req.op);
+
+ var staged: [16]u8 = undefined;
+ std.mem.writeInt(u64, staged[0..8], 4, .little);
+ staged[8] = 1;
+ staged[9] = 5;
+ @memcpy(staged[10..15], "panes");
+ fs.reply(&.{ .tag = 800, .payload = .{ .staged = 15 } }, staged[0..15]);
+
+ var buf: [512]u8 = undefined;
+ const got = try readReply(sv[1], &buf);
+ try testing.expectEqual(@as(usize, @sizeOf(fuse_out_header) + 32), got.len);
+ const rec = got[@sizeOf(fuse_out_header)..];
+ try testing.expectEqual(@as(u64, 4), std.mem.readInt(u64, rec[0..8], .little));
+ try testing.expectEqual(@as(u64, 1), std.mem.readInt(u64, rec[8..16], .little));
+ try testing.expectEqual(DT_DIR, std.mem.readInt(u32, rec[20..24], .little));
+ try testing.expectEqualStrings("panes", rec[24..29]);
+}
+
+test "INIT reply negotiates nothing and caps the minor at ours" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+ fs.max_write = 64 * 1024;
+
+ const in: fuse_init_in = .{
+ .major = 7,
+ .minor = 45,
+ .max_readahead = 131072,
+ // Everything the kernel is willing to do. The point of the test is that
+ // none of it comes back.
+ .flags = 0xffff_ffff,
+ .flags2 = 0xffff_ffff,
+ .unused = @splat(0),
+ };
+ try pushRequest(sv[1], 1, .init, 0, std.mem.asBytes(&in));
+ try fs.handshake();
+ try testing.expectEqual(@as(u32, 45), fs.minor);
+
+ var buf: [512]u8 = undefined;
+ const got = try readReply(sv[1], &buf);
+ try testing.expectEqual(@as(usize, @sizeOf(fuse_out_header) + @sizeOf(fuse_init_out)), got.len);
+ try testing.expectEqual(@as(u64, 1), outHeader(got).unique);
+ var out: fuse_init_out = undefined;
+ @memcpy(std.mem.asBytes(&out), got[@sizeOf(fuse_out_header)..][0..@sizeOf(fuse_init_out)]);
+ try testing.expectEqual(@as(u32, 7), out.major);
+ try testing.expectEqual(@as(u32, 45), out.minor);
+ // The one assertion this test exists for. Every bit here is a kernel
+ // behaviour we would owe forever: readdirplus whose ENOSYS has no fallback,
+ // atomic O_TRUNC that would bypass the SETATTR the core handles, locks.
+ try testing.expectEqual(@as(u32, 0), out.flags);
+ try testing.expectEqual(@as(u32, 0), out.flags2);
+ try testing.expectEqual(@as(u32, 0), out.max_readahead);
+ try testing.expectEqual(@as(u32, 64 * 1024), out.max_write);
+ // A time granularity of zero is not a legal value.
+ try testing.expectEqual(@as(u32, 1), out.time_gran);
+ try testing.expectEqual(@as(u16, 0), out.request_timeout);
+
+ // A newer kernel's minor is CAPPED, not echoed. `fc->minor` is our own
+ // declared level and it is what sizes the replies the kernel reads back
+ // from us, so claiming 7.99 on these structs promises fields they do not
+ // have. This assertion is the one the old `@min(in.minor, in.minor)` could
+ // not make.
+ const newer: fuse_init_in = .{
+ .major = 7,
+ .minor = kernel_minor + 54,
+ .max_readahead = 0,
+ .flags = 0,
+ .flags2 = 0,
+ .unused = @splat(0),
+ };
+ try pushRequest(sv[1], 2, .init, 0, std.mem.asBytes(&newer));
+ try fs.handshake();
+ const capped = try readReply(sv[1], &buf);
+ @memcpy(std.mem.asBytes(&out), capped[@sizeOf(fuse_out_header)..][0..@sizeOf(fuse_init_out)]);
+ try testing.expectEqual(kernel_minor, out.minor);
+
+ // A foreign major is fatal, and answering it anyway only moves the failure
+ // to the first syscall through the mount.
+ const bad: fuse_init_in = .{
+ .major = 8,
+ .minor = 0,
+ .max_readahead = 0,
+ .flags = 0,
+ .flags2 = 0,
+ .unused = @splat(0),
+ };
+ try pushRequest(sv[1], 3, .init, 0, std.mem.asBytes(&bad));
+ try testing.expectError(error.InitVersion, fs.handshake());
+}
+
+test "park table: again holds the request, retry offers it back once per round" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ const in: fuse_read_in = .{
+ .fh = 1,
+ .offset = 0,
+ .size = 64,
+ .read_flags = 0,
+ .lock_owner = 0,
+ .flags = 0,
+ .padding = 0,
+ };
+ // Two blocked readers of `event`, in arrival order.
+ try pushRequest(sv[1], 900, .read, 20, std.mem.asBytes(&in));
+ try pushRequest(sv[1], 902, .read, 21, std.mem.asBytes(&in));
+ const a = fs.next() orelse return error.NoRequest;
+ fs.reply(&.{ .tag = a.tag, .status = .again }, &.{});
+ const b = fs.next() orelse return error.NoRequest;
+ fs.reply(&.{ .tag = b.tag, .status = .again }, &.{});
+ try testing.expect(fs.next() == null);
+ // Nothing was written: a held request has no reply, which is the only way
+ // FUSE expresses blocking.
+ var buf: [512]u8 = undefined;
+ try testing.expect(libc.read(sv[1], &buf, buf.len) < 0);
+
+ // One round offers each parked request exactly once, oldest first, and then
+ // ends. Without the per-round flag the oldest would be offered forever and
+ // the second reader would never be looked at again.
+ const r1 = fs.retry() orelse return error.NoRetry;
+ try testing.expectEqual(@as(u64, 900), r1.tag);
+ fs.reply(&.{ .tag = r1.tag, .status = .again }, &.{});
+ const r2 = fs.retry() orelse return error.NoRetry;
+ try testing.expectEqual(@as(u64, 902), r2.tag);
+ fs.reply(&.{ .tag = r2.tag, .status = .again }, &.{});
+ try testing.expect(fs.retry() == null);
+
+ // ...and the next round starts over.
+ const r3 = fs.retry() orelse return error.NoRetry;
+ try testing.expectEqual(@as(u64, 900), r3.tag);
+ fs.reply(&.{ .tag = r3.tag, .payload = .{ .staged = 3 } }, "ev\n");
+ const got = try readReply(sv[1], &buf);
+ try testing.expectEqualStrings("ev\n", got[@sizeOf(fuse_out_header)..]);
+ // The answered one is gone; the other is still held.
+ try testing.expectEqual(@as(?usize, null), fs.findSlot(900));
+ try testing.expect(fs.findSlot(902) != null);
+}
+
+test "park table: interrupt answers the original with EINTR and drops it" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ const in: fuse_read_in = .{
+ .fh = 1,
+ .offset = 0,
+ .size = 64,
+ .read_flags = 0,
+ .lock_owner = 0,
+ .flags = 0,
+ .padding = 0,
+ };
+ try pushRequest(sv[1], 1000, .read, 20, std.mem.asBytes(&in));
+ try pushRequest(sv[1], 1002, .read, 21, std.mem.asBytes(&in));
+ const a = fs.next() orelse return error.NoRequest;
+ fs.reply(&.{ .tag = a.tag, .status = .again }, &.{});
+ const b = fs.next() orelse return error.NoRequest;
+ fs.reply(&.{ .tag = b.tag, .status = .again }, &.{});
+
+ // The kernel's interrupt names the ORIGINAL unique in its body; its own
+ // unique is `original | 1`, which is why it must not be echoed.
+ const intr: fuse_interrupt_in = .{ .unique = 1002 };
+ try pushRequest(sv[1], 1002 | 1, .interrupt, 0, std.mem.asBytes(&intr));
+ try testing.expect(fs.next() == null);
+
+ var buf: [512]u8 = undefined;
+ const got = try readReply(sv[1], &buf);
+ // Exactly one reply, to the interrupted request, not to the interrupt.
+ // Getting this wrong leaves a SIGKILLed reader in uninterruptible sleep.
+ try testing.expectEqual(@as(usize, @sizeOf(fuse_out_header)), got.len);
+ const h = outHeader(got);
+ try testing.expectEqual(@as(u64, 1002), h.unique);
+ try testing.expectEqual(-@as(i32, @intFromEnum(libc.E.INTR)), h.@"error");
+ try testing.expectEqual(@as(?usize, null), fs.findSlot(1002));
+ try testing.expect(fs.findSlot(1000) != null);
+
+ // An interrupt for something we do not hold is ignored, not answered.
+ const stale: fuse_interrupt_in = .{ .unique = 4242 };
+ try pushRequest(sv[1], 4243, .interrupt, 0, std.mem.asBytes(&stale));
+ try testing.expect(fs.next() == null);
+ try testing.expect(libc.read(sv[1], &buf, buf.len) < 0);
+}
+
+test "park table: a full table answers EAGAIN and keeps the descriptor flowing" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ const in: fuse_read_in = .{
+ .fh = 1,
+ .offset = 0,
+ .size = 8,
+ .read_flags = 0,
+ .lock_owner = 0,
+ .flags = 0,
+ .padding = 0,
+ };
+ for (0..max_slots) |i| {
+ try pushRequest(sv[1], 2000 + i * 2, .read, 30, std.mem.asBytes(&in));
+ const req = fs.next() orelse return error.NoRequest;
+ fs.reply(&.{ .tag = req.tag, .status = .again }, &.{});
+ }
+ var buf: [512]u8 = undefined;
+
+ // One more than the table holds. It is READ and refused, not left queued.
+ // Gating the read on a free slot is a deadlock dressed as backpressure:
+ // the INTERRUPT that frees a slot would never be read either, so a
+ // SIGKILLed reader would stay in uninterruptible sleep and every unrelated
+ // `ls` of the mount would hang behind the 32 blocked ones. Measured: that
+ // wedges a real mount.
+ try pushRequest(sv[1], 9998, .read, 30, std.mem.asBytes(&in));
+ try testing.expect(fs.next() == null);
+ const refused = try readReply(sv[1], &buf);
+ try testing.expectEqual(@as(u64, 9998), outHeader(refused).unique);
+ try testing.expectEqual(-@as(i32, @intFromEnum(libc.E.AGAIN)), outHeader(refused).@"error");
+
+ // And the requests that need no slot keep being answered with the table
+ // still full — DESTROY above all, since a missing reply to it hangs umount.
+ try pushRequest(sv[1], 9990, .access, 1, &.{});
+ try testing.expect(fs.next() == null);
+ const nosys = try readReply(sv[1], &buf);
+ try testing.expectEqual(-@as(i32, @intFromEnum(libc.E.NOSYS)), outHeader(nosys).@"error");
+
+ // An interrupt still lands, which is what lets a full table recover at all.
+ const intr: fuse_interrupt_in = .{ .unique = 2000 };
+ try pushRequest(sv[1], 2001, .interrupt, 0, std.mem.asBytes(&intr));
+ try testing.expect(fs.next() == null);
+ const killed = try readReply(sv[1], &buf);
+ try testing.expectEqual(@as(u64, 2000), outHeader(killed).unique);
+ try testing.expectEqual(-@as(i32, @intFromEnum(libc.E.INTR)), outHeader(killed).@"error");
+
+ // ...and the freed slot takes the next request.
+ try pushRequest(sv[1], 9996, .read, 30, std.mem.asBytes(&in));
+ const late = fs.next() orelse return error.NoRequest;
+ try testing.expectEqual(@as(u64, 9996), late.tag);
+}
+
+test "park table: a payload too large to copy is refused rather than dangled" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ const payload_len = park_data_max + 1;
+ var body: [@sizeOf(fuse_write_in) + payload_len]u8 = undefined;
+ var in: fuse_write_in = std.mem.zeroes(fuse_write_in);
+ in.size = payload_len;
+ @memcpy(body[0..@sizeOf(fuse_write_in)], std.mem.asBytes(&in));
+ @memset(body[@sizeOf(fuse_write_in)..], 'z');
+ try pushRequest(sv[1], 3000, .write, 6, &body);
+ const req = fs.next() orelse return error.NoRequest;
+ try testing.expectEqual(@as(usize, payload_len), req.data.len);
+
+ // Parking this would park a slice of the read buffer, which the next
+ // `next()` overwrites. EAGAIN is the honest answer.
+ fs.reply(&.{ .tag = 3000, .status = .again }, &.{});
+ var buf: [512]u8 = undefined;
+ const got = try readReply(sv[1], &buf);
+ try testing.expectEqual(-@as(i32, @intFromEnum(libc.E.AGAIN)), outHeader(got).@"error");
+ try testing.expectEqual(@as(?usize, null), fs.findSlot(3000));
+
+ // A payload that fits IS copied, so parking it is safe even after the read
+ // buffer has been reused.
+ var small: [@sizeOf(fuse_write_in) + 4]u8 = undefined;
+ in.size = 4;
+ @memcpy(small[0..@sizeOf(fuse_write_in)], std.mem.asBytes(&in));
+ @memcpy(small[@sizeOf(fuse_write_in)..], "keep");
+ try pushRequest(sv[1], 3002, .write, 6, &small);
+ const kept = fs.next() orelse return error.NoRequest;
+ fs.reply(&.{ .tag = kept.tag, .status = .again }, &.{});
+ // Something else lands in the read buffer...
+ try pushRequest(sv[1], 3004, .statfs, 1, &.{});
+ _ = fs.next() orelse return error.NoRequest;
+ // ...and the parked bytes survived it.
+ const again = fs.retry() orelse return error.NoRetry;
+ try testing.expectEqualStrings("keep", again.data);
+}
+
+test "a reply for a tag we no longer hold is dropped, not written" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ // The interrupt path already answered and freed this one; a second reply
+ // would carry a unique the kernel does not recognise, and could in
+ // principle be matched against a live request that reused the number.
+ fs.reply(&.{ .tag = 12345 }, &.{});
+ var buf: [512]u8 = undefined;
+ try testing.expect(libc.read(sv[1], &buf, buf.len) < 0);
+}
+
+test "statfs reports a usable namelen" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ try pushRequest(sv[1], 1100, .statfs, 1, &.{});
+ const req = fs.next() orelse return error.NoRequest;
+ try testing.expectEqual(acmefs.Op.statfs, req.op);
+ fs.reply(&.{ .tag = 1100 }, &.{});
+
+ var buf: [512]u8 = undefined;
+ const got = try readReply(sv[1], &buf);
+ var out: fuse_statfs_out = undefined;
+ @memcpy(std.mem.asBytes(&out), got[@sizeOf(fuse_out_header)..][0..@sizeOf(fuse_statfs_out)]);
+ // Zero here makes pathconf(_PC_NAME_MAX) return 0 and some tools then
+ // refuse to create any name at all.
+ try testing.expectEqual(@as(u32, 255), out.st.namelen);
+ try testing.expectEqual(@as(u32, 4096), out.st.bsize);
+}
+
+test "the fusermount command line and environment" {
+ var opts_buf: [128:0]u8 = undefined;
+ const opts = mountOpts(&opts_buf);
+ try testing.expectEqualStrings("fsname=pardes,subtype=pardes,nosuid,nodev", opts);
+ // allow_other needs user_allow_other in /etc/fuse.conf, which is commented
+ // out on a stock install, and asking for it FAILS the whole mount rather
+ // than being ignored. default_permissions would move access control out of
+ // the core and into a mode nibble.
+ try testing.expect(std.mem.indexOf(u8, opts, "allow_other") == null);
+ try testing.expect(std.mem.indexOf(u8, opts, "default_permissions") == null);
+
+ var env_buf: [32:0]u8 = undefined;
+ try testing.expectEqualStrings("_FUSE_COMMFD=7", commfdEnv(&env_buf, 7));
+
+ var argv: [6:null]?[*:0]const u8 = undefined;
+ mountArgv(&argv, "/usr/bin/fusermount3", opts.ptr, "/run/user/1000/pardes/42");
+ try testing.expectEqualStrings("/usr/bin/fusermount3", std.mem.span(argv[0].?));
+ try testing.expectEqualStrings("-o", std.mem.span(argv[1].?));
+ try testing.expectEqualStrings("fsname=pardes,subtype=pardes,nosuid,nodev", std.mem.span(argv[2].?));
+ // Without the `--` a mountpoint beginning with a dash is parsed as a flag
+ // by a setuid program.
+ try testing.expectEqualStrings("--", std.mem.span(argv[3].?));
+ try testing.expectEqualStrings("/run/user/1000/pardes/42", std.mem.span(argv[4].?));
+ try testing.expectEqual(@as(?[*:0]const u8, null), argv[5]);
+
+ var uargv: [7:null]?[*:0]const u8 = undefined;
+ unmountArgv(&uargv, "/usr/bin/fusermount3", "/run/user/1000/pardes/42");
+ try testing.expectEqualStrings("-u", std.mem.span(uargv[1].?));
+ try testing.expectEqualStrings("-q", std.mem.span(uargv[2].?));
+ // Lazy, or a pane shell with a cwd inside the mount makes the unmount fail
+ // with EBUSY and the mount outlives the editor.
+ try testing.expectEqualStrings("-z", std.mem.span(uargv[3].?));
+ try testing.expectEqualStrings("--", std.mem.span(uargv[4].?));
+ try testing.expectEqual(@as(?[*:0]const u8, null), uargv[6]);
+}
+
+test "the child environment drops an inherited comm descriptor" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ var buf: [32:0]u8 = undefined;
+ const commfd = commfdEnv(&buf, 5);
+ const env = try buildEnv(gpa, commfd);
+ defer gpa.free(env);
+
+ // Exactly one _FUSE_COMMFD, and it is ours: getenv returns the FIRST match,
+ // so an inherited stale entry would win and fusermount3 would send the
+ // descriptor to a closed socket.
+ var seen: usize = 0;
+ var i: usize = 0;
+ while (env[i]) |entry| : (i += 1) {
+ if (std.mem.startsWith(u8, std.mem.span(entry), commfd_env ++ "=")) {
+ seen += 1;
+ try testing.expectEqualStrings("_FUSE_COMMFD=5", std.mem.span(entry));
+ }
+ }
+ try testing.expectEqual(@as(usize, 1), seen);
+ try testing.expectEqual(@as(?[*:0]const u8, null), env[env.len - 1]);
+}
+
+test "poll thread: one wake per drained batch, and stop joins from either state" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ var wakes: std.atomic.Value(u32) = .init(0);
+ const Sink = struct {
+ fn wake(ctx: ?*anyopaque) void {
+ const c: *std.atomic.Value(u32) = @ptrCast(@alignCast(ctx.?));
+ _ = c.fetchAdd(1, .release);
+ }
+ };
+ try fs.wakeThread(&wakes, Sink.wake);
+
+ // One pending request, one wake. A FORGET is answered inside `next()` and
+ // never surfaces, so draining to null is the whole batch — and it is that
+ // null which raises the eventfd and lets the poller poll again.
+ const forget: fuse_forget_in = .{ .nlookup = 1 };
+ try pushRequest(sv[1], 7000, .forget, 2, std.mem.asBytes(&forget));
+ while (wakes.load(.acquire) == 0) std.Thread.yield() catch {};
+ try testing.expectEqual(@as(?acmefs.Req, null), fs.next());
+
+ // The poller is now in one of the two states a stop has to break: still in
+ // the blocking wait, or back in `poll()` because the drain above beat the
+ // stop there. Which one is a race, deliberately unresolved — the assertion
+ // is that either joins, and a hang here is this test's only failure mode.
+ fs.stopThread();
+ try testing.expect(fs.thread == null);
+ try testing.expectEqual(@as(c_int, -1), fs.ctl);
+}
+
+test "poll thread: stop breaks a poller that never saw a request" {
+ if (comptime !supported) return;
+ const gpa = testing.allocator;
+ const sv = try testPair();
+ defer {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ }
+ const fs = try testFs(gpa, sv[0]);
+ defer testFsFree(fs);
+
+ const Sink = struct {
+ fn wake(_: ?*anyopaque) void {
+ unreachable; // nothing is ever pending on this descriptor
+ }
+ };
+ try fs.wakeThread(null, Sink.wake);
+ // Covers the two states with no acknowledgement in them at all: blocked in
+ // `poll()` with an idle descriptor, and not yet past the loop condition.
+ fs.stopThread();
+ try testing.expect(fs.thread == null);
+}
+
+test "mount refuses a relative point" {
+ if (comptime !supported) return;
+ try testing.expectError(
+ error.MountPathNotAbsolute,
+ Fs.mount(testing.allocator, .{ .mount = "relative/dir" }),
+ );
+}