diff options
| author | Gabriel Schneider <[email protected]> | 2026-08-25 02:07:23 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-08-25 09:42:07 -0300 |
| commit | 6f48508aa08396bcf9dd4da2cab1d221bcc53f78 (patch) | |
| tree | 83daec3db5ac27ea3172651df2b3e9cb62eddb4a /src/fuse.zig | |
| parent | 28c70cabb6ceb7e5fecfd74f6984f5a995269f01 (diff) | |
| download | pardes-6f48508aa08396bcf9dd4da2cab1d221bcc53f78.tar.gz pardes-6f48508aa08396bcf9dd4da2cab1d221bcc53f78.zip | |
acmefs: pardes --fs serves acme's control filesystem over raw Linux FUSE
Diffstat (limited to 'src/fuse.zig')
| -rw-r--r-- | src/fuse.zig | 2709 |
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" }), + ); +} |
