summaryrefslogtreecommitdiff
path: root/9ns/src/nine.zig
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-19 23:55:47 -0300
committerGabriel Schneider <[email protected]>2026-09-19 23:55:47 -0300
commit3e9f8805f293f622bb885cf849b5ce47dc062ad1 (patch)
tree0f31107fd20e7a9aa065826c891d2d326efc2309 /9ns/src/nine.zig
parentba996acfcad1698adbf4a1834fe50e73b1c6cab9 (diff)
downloadcloud9-3e9f8805f293f622bb885cf849b5ce47dc062ad1.tar.gz
cloud9-3e9f8805f293f622bb885cf849b5ce47dc062ad1.zip
9ns: --name and /mnt/9p/<name> mounts, qid.path as inode number, interrupts as Tflush
- --name NAME (default derived from the transport: socket basename, tcp-IP-PORT, spawned command, fdN) mounts at /mnt/9p/<name>; --mount still overrides. ensureMountpoint walks down and creates missing components, shadowing the deepest unwritable ancestor. NINE_MOUNT is the only exported variable. - The inode number reported to the kernel is the 9P qid.path for every node, root included; a server handing qid.path 1 to a file (Pardes /self) no longer collides with the root. - FUSE_INTERRUPT for the request in flight becomes Tflush; a blocked read returns EINTR when the server answers the flush, chunked transfers return short counts, other requests arriving meanwhile are stashed and served next. Servers ignoring Tflush still block until they answer. - 9ns-test now covers nine/bridge/fuse; new adv_bridge_interrupt suite (28); 9ns-itest grows to 88 checks. Co-Authored-By: Claude Fable 5.1 <[email protected]>
Diffstat (limited to '9ns/src/nine.zig')
-rw-r--r--9ns/src/nine.zig438
1 files changed, 398 insertions, 40 deletions
diff --git a/9ns/src/nine.zig b/9ns/src/nine.zig
index 70633e6..c89a343 100644
--- a/9ns/src/nine.zig
+++ b/9ns/src/nine.zig
@@ -29,8 +29,23 @@ pub const dontcare = cloud9.Stat{
.muid = "",
};
+/// How the owner of a session (the FUSE bridge) gets a say while an rpc waits
+/// for its reply. `watch` names a descriptor to poll alongside the socket, or
+/// -1 to poll nothing extra right now; when it becomes readable `onReadable`
+/// consumes whatever is there and returns true if the request in flight should
+/// be cancelled with a Tflush. `armed` reports whether such a cancellation was
+/// requested earlier for the operation in progress; the chunked read/write
+/// loops stop between chunks when it is set (a reply that raced the flush still
+/// leaves the caller wanting out).
+pub const Interrupt = struct {
+ ctx: *anyopaque,
+ watch: *const fn (ctx: *anyopaque) i32,
+ onReadable: *const fn (ctx: *anyopaque) Session.Error!bool,
+ armed: *const fn (ctx: *anyopaque) bool,
+};
+
pub const Session = struct {
- pub const Error = error{ Nine, Protocol, Io, Closed, Stopped, TooLarge, OutOfMemory };
+ pub const Error = error{ Nine, Protocol, Io, Closed, Stopped, Interrupted, TooLarge, OutOfMemory };
pub const Walk = struct { nwqid: u16, wqid: [cloud9.max_welem]cloud9.Qid };
pub const Open = struct { qid: cloud9.Qid, iounit: u32 };
@@ -53,6 +68,9 @@ pub const Session = struct {
/// readable (the bridge's "child exited" pipe) the pending rpc fails with
/// `error.Stopped` instead of blocking on a server that never answers.
stop_fd: i32 = -1,
+ /// Optional interrupt source (the bridge's FUSE descriptor) consulted while
+ /// a reply is outstanding; see `Interrupt`.
+ interrupt: ?Interrupt = null,
/// Connect to `address`, then negotiate the protocol version.
/// `msize` is the maximum message size to ask for (0 = the buffers' size).
@@ -117,9 +135,17 @@ pub const Session = struct {
}
/// Generic RPC. Result slices borrow the input buffer until the next call.
+ ///
+ /// While the reply is outstanding the socket is polled together with
+ /// `stop_fd` (→ `error.Stopped`) and the interrupt source's descriptor. When
+ /// the latter asks for a cancellation a Tflush for the request's tag goes out
+ /// and the wait continues until either the original reply arrives (the flush
+ /// lost the race; the result is returned as if nothing happened and the
+ /// Rflush is swallowed by a later call) or the Rflush does (→
+ /// `error.Interrupted`; the server has dropped the request).
pub fn rpc(s: *Session, req: cloud9.Client.Request) Error!cloud9.Client.Result {
s.ename_len = 0;
- _ = s.client.submit(req) catch |e| switch (e) {
+ const tag = s.client.submit(req) catch |e| switch (e) {
error.NoTags, error.Handshake, error.Dead => return error.Protocol,
error.NoSpace, error.TooLarge => return error.TooLarge,
error.BadRequest => {
@@ -128,29 +154,52 @@ pub const Session = struct {
},
};
try s.flush();
+ var flush_tag: ?u16 = null;
var tmp: [64 * 1024]u8 = undefined;
while (true) {
- if (s.client.take()) |done| {
- switch (done.result) {
- .fail => |ename| {
- s.setEname(ename);
- return error.Nine;
- },
- else => return done.result,
+ while (s.client.take()) |done| {
+ if (done.tag == tag) {
+ switch (done.result) {
+ .fail => |ename| {
+ s.setEname(ename);
+ return error.Nine;
+ },
+ else => return done.result,
+ }
}
+ if (flush_tag != null and done.tag == flush_tag.?) return error.Interrupted;
+ // An Rflush for a flush whose original reply won the race in an
+ // earlier call: the client has released both tags; nothing to do.
+ if (done.op == .flush) continue;
+ return error.Protocol;
}
if (s.client.dead) return error.Protocol;
// After take() returned null the previous frame is gone, so the free
// space is at least what the pending frame still needs.
const room = s.client.in.len - s.client.in_len;
if (room == 0) return error.Protocol;
- const n = try readSome(s.fd, s.stop_fd, tmp[0..@min(room, tmp.len)]);
- if (n == 0) return error.Closed;
- const pushed = s.client.push(tmp[0..n]);
- if (pushed != n) return error.Protocol;
+ switch (try s.wait()) {
+ .socket => {
+ const n = try readSocket(s.fd, tmp[0..@min(room, tmp.len)]);
+ if (n == 0) return error.Closed;
+ const pushed = s.client.push(tmp[0..n]);
+ if (pushed != n) return error.Protocol;
+ },
+ .cancel => if (flush_tag == null) {
+ flush_tag = s.client.submit(.{ .flush = .{ .oldtag = tag } }) catch return error.Protocol;
+ try s.flush();
+ },
+ }
}
}
+ /// True when the interrupt source has asked for the operation in progress to
+ /// stop; consulted between the chunks of a read or write.
+ pub fn interruptArmed(s: *const Session) bool {
+ const i = s.interrupt orelse return false;
+ return i.armed(i.ctx);
+ }
+
/// Walk `names` from `fid` to `newfid`. A partial walk leaves `newfid` unbound
/// (9P semantics) and reports `error.Nine` with ename "file does not exist".
pub fn walk(s: *Session, fid: u32, newfid: u32, names: []const []const u8) Error!Walk {
@@ -245,6 +294,50 @@ pub const Session = struct {
return chunkSize(s.client.maxWrite(), s.iounits.get(fid) orelse 0);
}
+ const Ready = enum { socket, cancel };
+
+ /// Blocks until the socket is readable (`.socket`), the interrupt source
+ /// wants the request in flight cancelled (`.cancel`), or `stop_fd` fires
+ /// (`error.Stopped`). Anything the interrupt source consumes without asking
+ /// for a cancellation simply resumes the wait.
+ fn wait(s: *Session) Error!Ready {
+ while (true) {
+ var pfds: [3]linux.pollfd = undefined;
+ var n: usize = 0;
+ pfds[n] = .{ .fd = s.fd, .events = linux.POLL.IN, .revents = 0 };
+ n += 1;
+ const stop_at: ?usize = if (s.stop_fd >= 0) n else null;
+ if (stop_at != null) {
+ pfds[n] = .{ .fd = s.stop_fd, .events = linux.POLL.IN, .revents = 0 };
+ n += 1;
+ }
+ const ifd: i32 = if (s.interrupt) |i| i.watch(i.ctx) else -1;
+ const int_at: ?usize = if (ifd >= 0) n else null;
+ if (int_at != null) {
+ pfds[n] = .{ .fd = ifd, .events = linux.POLL.IN, .revents = 0 };
+ n += 1;
+ }
+ if (n == 1) return .socket;
+ const prc = linux.poll(&pfds, @intCast(n), -1);
+ switch (linux.errno(prc)) {
+ .SUCCESS => {},
+ .INTR, .AGAIN => continue,
+ else => return error.Io,
+ }
+ // A reply that is already there wins over everything else.
+ if (pfds[0].revents != 0) return .socket;
+ if (stop_at) |i| {
+ if (pfds[i].revents != 0) return error.Stopped;
+ }
+ if (int_at) |i| {
+ if (pfds[i].revents != 0) {
+ const src = s.interrupt.?;
+ if (try src.onReadable(src.ctx)) return .cancel;
+ }
+ }
+ }
+ }
+
/// Writes everything in the client's output buffer to the socket.
fn flush(s: *Session) Error!void {
while (s.client.output().len != 0) {
@@ -280,8 +373,13 @@ fn readWith(
if (max_chunk == 0) return error.Protocol;
var done: usize = 0;
while (done < buf.len) {
+ // Like read(2): an interruption after some data arrived is a short read.
+ if (done != 0 and s.interruptArmed()) break;
const want: u32 = @intCast(@min(buf.len - done, max_chunk));
- const r = try rpcFn(s, .{ .read = .{ .fid = fid, .offset = offset + done, .count = want } });
+ const r = rpcFn(s, .{ .read = .{ .fid = fid, .offset = offset + done, .count = want } }) catch |e| {
+ if (e == error.Interrupted and done != 0) break;
+ return e;
+ };
const data = r.read;
@memcpy(buf[done..][0..data.len], data);
done += data.len;
@@ -302,29 +400,20 @@ fn writeWith(
if (max_chunk == 0) return error.Protocol;
var done: usize = 0;
while (done < data.len) {
+ if (done != 0 and s.interruptArmed()) break;
const want: usize = @min(data.len - done, max_chunk);
- const r = try rpcFn(s, .{ .write = .{ .fid = fid, .offset = offset + done, .data = data[done..][0..want] } });
+ const r = rpcFn(s, .{ .write = .{ .fid = fid, .offset = offset + done, .data = data[done..][0..want] } }) catch |e| {
+ if (e == error.Interrupted and done != 0) break;
+ return e;
+ };
done += r.write;
if (r.write < want) break;
}
return done;
}
-fn readSome(fd: i32, stop_fd: i32, buf: []u8) Session.Error!usize {
+fn readSocket(fd: i32, buf: []u8) Session.Error!usize {
while (true) {
- if (stop_fd >= 0) {
- var pfds = [_]linux.pollfd{
- .{ .fd = fd, .events = linux.POLL.IN, .revents = 0 },
- .{ .fd = stop_fd, .events = linux.POLL.IN, .revents = 0 },
- };
- const prc = linux.poll(&pfds, pfds.len, -1);
- switch (linux.errno(prc)) {
- .SUCCESS => {},
- .INTR, .AGAIN => continue,
- else => return error.Io,
- }
- if (pfds[1].revents != 0 and pfds[0].revents == 0) return error.Stopped;
- }
const rc = linux.read(fd, buf.ptr, buf.len);
switch (linux.errno(rc)) {
.SUCCESS => return rc,
@@ -339,6 +428,9 @@ fn readSome(fd: i32, stop_fd: i32, buf: []u8) Session.Error!usize {
pub fn enameToErrno(ename: []const u8) linux.E {
const Rule = struct { needle: []const u8, err: linux.E };
const rules = [_]Rule{
+ // A server that answers a flushed request with an error (Pardes says
+ // "Interrupted system call") should look like a flush to the caller.
+ .{ .needle = "interrupt", .err = .INTR },
.{ .needle = "not exist", .err = .NOENT },
.{ .needle = "not found", .err = .NOENT },
.{ .needle = "no such", .err = .NOENT },
@@ -461,6 +553,8 @@ test {
}
test "ename → errno mapping" {
+ try testing.expectEqual(linux.E.INTR, enameToErrno("Interrupted system call"));
+ try testing.expectEqual(linux.E.INTR, enameToErrno("read interrupted"));
try testing.expectEqual(linux.E.NOENT, enameToErrno("file does not exist"));
try testing.expectEqual(linux.E.NOENT, enameToErrno("No Such File"));
try testing.expectEqual(linux.E.NOENT, enameToErrno("directory entry not found"));
@@ -524,10 +618,19 @@ const FakeFile = struct {
calls: usize = 0,
max_count: u32 = 0,
short_write_at: ?usize = null,
+ /// Fail the call with `error.Interrupted` once this many calls were made.
+ interrupt_at: ?usize = null,
+ /// Report an armed interrupt once this many calls were made.
+ armed_at: ?usize = null,
scratch: [4096]u8 = undefined,
+ fn interruptArmed(f: *const FakeFile) bool {
+ return if (f.armed_at) |at| f.calls >= at else false;
+ }
+
fn rpc(f: *FakeFile, req: cloud9.Client.Request) Session.Error!cloud9.Client.Result {
f.calls += 1;
+ if (f.interrupt_at) |at| if (f.calls > at) return error.Interrupted;
switch (req) {
.read => |r| {
f.max_count = @max(f.max_count, r.count);
@@ -583,20 +686,60 @@ test "write chunks and stops at a short write" {
try testing.expectEqual(@as(usize, 2), f.calls);
}
+test "interrupted chunk loops: partial count if data moved, Interrupted otherwise" {
+ var buf: [4000]u8 = undefined;
+ // The second chunk's rpc is interrupted: the first chunk is returned.
+ var f: FakeFile = .{ .len = 2500, .interrupt_at = 1 };
+ try testing.expectEqual(@as(usize, 1000), try readWith(&f, FakeFile.rpc, 7, 0, &buf, 1000));
+ // The first chunk's rpc is interrupted: nothing was transferred.
+ f = .{ .len = 2500, .interrupt_at = 0 };
+ try testing.expectError(error.Interrupted, readWith(&f, FakeFile.rpc, 7, 0, &buf, 1000));
+ // A reply that raced the flush arms the interrupt: stop before the next chunk.
+ f = .{ .len = 2500, .armed_at = 2 };
+ try testing.expectEqual(@as(usize, 2000), try readWith(&f, FakeFile.rpc, 7, 0, &buf, 1000));
+ try testing.expectEqual(@as(usize, 2), f.calls);
+ // Same for writes.
+ var data: [2500]u8 = undefined;
+ for (&data, 0..) |*b, i| b.* = @truncate(i);
+ f = .{ .len = 0, .interrupt_at = 2 };
+ try testing.expectEqual(@as(usize, 2000), try writeWith(&f, FakeFile.rpc, 7, 0, &data, 1000));
+ f = .{ .len = 0, .interrupt_at = 0 };
+ try testing.expectError(error.Interrupted, writeWith(&f, FakeFile.rpc, 7, 0, &data, 1000));
+ f = .{ .len = 0, .armed_at = 1 };
+ try testing.expectEqual(@as(usize, 1000), try writeWith(&f, FakeFile.rpc, 7, 0, &data, 1000));
+}
+
// -- in-process server test ---------------------------------------------------------
/// A tiny 9P2000 backend on a cloud9.Server: answers version/attach/walk/stat/open/
/// read/clunk/remove with canned data. Runs in its own thread over a socketpair.
-const FakeServer = struct {
+/// Test support only (bridge.zig's tests use it too).
+pub const FakeServer = struct {
fd: i32,
msize: u32,
max_read_count: u32 = 0,
file_len: usize,
+ /// A Tread at this offset is never answered (a blocked stream read); the
+ /// server keeps serving whatever else arrives, notably a Tflush.
+ hang_offset: ?u64 = null,
+ /// What a Tflush gets: the connection dropped (`.hangup`), an Rflush, or
+ /// first the Rread the flush was aimed at and then the Rflush (the race).
+ on_flush: enum { hangup, rflush, reply_then_rflush } = .hangup,
+ /// Observed by the test thread: number of Tflush seen and the last oldtag.
+ flushes: std.atomic.Value(u32) = .init(0),
+ flush_oldtag: std.atomic.Value(u32) = .init(0xFFFF),
+ /// The tag of the hung Tread, for the test to compare with `flush_oldtag`.
+ hung_tag: std.atomic.Value(u32) = .init(0xFFFF),
+ /// When set, the server injects that source's INTERRUPT the moment a read
+ /// hangs, so the cancellation provably arrives while the wait is on.
+ on_hang_inject: ?*PipeInterrupt = null,
+ /// Delay before every Rread, so a test can be sure the client is waiting.
+ read_delay_ns: u64 = 0,
- const file_qid: cloud9.Qid = .{ .type = 0, .version = 3, .path = 0x1234 };
- const dir_qid: cloud9.Qid = .{ .type = cloud9.qtdir, .version = 1, .path = 0x1 };
+ pub const file_qid: cloud9.Qid = .{ .type = 0, .version = 3, .path = 0x1234 };
+ pub const dir_qid: cloud9.Qid = .{ .type = cloud9.qtdir, .version = 1, .path = 0x1 };
- fn run(fs: *FakeServer) void {
+ pub fn run(fs: *FakeServer) void {
fs.loop() catch |e| std.debug.print("fake server: {s}\n", .{@errorName(e)});
_ = linux.close(fs.fd);
}
@@ -610,6 +753,7 @@ const FakeServer = struct {
var srv: cloud9.Server = .init(.{ .in = in, .out = out });
var tmp: [4096]u8 = undefined;
var data: [8192]u8 = undefined;
+ var hung: ?struct { tag: u16, offset: u64, count: u32 } = null;
while (true) {
while (try srv.receive()) |req| {
const tag = req.tag;
@@ -649,18 +793,41 @@ const FakeServer = struct {
.topen => |m| try srv.reply(tag, .{ .ropen = .{ .qid = file_qid, .iounit = if (m.mode == cloud9.owrite) 700 else 0 } }),
.tread => |m| {
fs.max_read_count = @max(fs.max_read_count, m.count);
- var n: usize = 0;
- if (m.offset < fs.file_len) n = @min(@as(usize, m.count), fs.file_len - @as(usize, @intCast(m.offset)));
- n = @min(n, data.len);
- for (data[0..n], 0..) |*b, i| b.* = @truncate(m.offset + i);
- try srv.reply(tag, .{ .rread = .{ .data = data[0..n] } });
+ if (fs.hang_offset != null and fs.hang_offset.? == m.offset) {
+ hung = .{ .tag = tag, .offset = m.offset, .count = m.count };
+ fs.hung_tag.store(tag, .seq_cst);
+ if (fs.on_hang_inject) |p| try p.inject(p.unique);
+ } else {
+ if (fs.read_delay_ns != 0) {
+ const ts: linux.timespec = .{ .sec = @intCast(fs.read_delay_ns / std.time.ns_per_s), .nsec = @intCast(fs.read_delay_ns % std.time.ns_per_s) };
+ _ = linux.nanosleep(&ts, null);
+ }
+ try srv.reply(tag, .{ .rread = .{ .data = fs.fill(&data, m.offset, m.count) } });
+ }
},
.twrite => |m| try srv.reply(tag, .{ .rwrite = .{ .count = @intCast(m.data.len) } }),
.tclunk => try srv.reply(tag, .rclunk),
.tremove => try srv.reply(tag, .{ .rerror = .{ .ename = "permission denied" } }),
.twstat => try srv.reply(tag, .rwstat),
- // A flush is the test's "hang up now" signal.
- .tflush => return,
+ .tflush => |m| {
+ fs.flush_oldtag.store(m.oldtag, .seq_cst);
+ _ = fs.flushes.fetchAdd(1, .seq_cst);
+ switch (fs.on_flush) {
+ // The test's "hang up now" signal.
+ .hangup => return,
+ .rflush => {
+ if (hung != null and hung.?.tag == m.oldtag) hung = null;
+ try srv.reply(tag, .rflush);
+ },
+ .reply_then_rflush => {
+ if (hung) |h| if (h.tag == m.oldtag) {
+ try srv.reply(h.tag, .{ .rread = .{ .data = fs.fill(&data, h.offset, h.count) } });
+ hung = null;
+ };
+ try srv.reply(tag, .rflush);
+ },
+ }
+ },
else => try srv.reply(tag, .{ .rerror = .{ .ename = "not supported" } }),
}
srv.release();
@@ -677,8 +844,199 @@ const FakeServer = struct {
if (srv.push(tmp[0..rc]) != rc) return error.Overflow;
}
}
+
+ /// File contents: byte i == i & 0xff, `file_len` bytes long.
+ fn fill(fs: *const FakeServer, data: []u8, offset: u64, count: u32) []const u8 {
+ var n: usize = 0;
+ if (offset < fs.file_len) n = @min(@as(usize, count), fs.file_len - @as(usize, @intCast(offset)));
+ n = @min(n, data.len);
+ for (data[0..n], 0..) |*b, i| b.* = @truncate(offset + i);
+ return data[0..n];
+ }
+
+ /// A connected session (fid 0 attached, fid 1 walked to "file" and opened
+ /// for reading) plus the server thread; `close` when done.
+ pub const Pair = struct {
+ server: *FakeServer,
+ session: Session,
+ thread: std.Thread,
+
+ pub fn close(p: *Pair) void {
+ p.session.deinit();
+ p.thread.join();
+ }
+ };
+
+ pub fn start(fs: *FakeServer) !Pair {
+ var fds: [2]i32 = undefined;
+ if (linux.errno(linux.socketpair(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0, &fds)) != .SUCCESS) return error.Io;
+ fs.fd = fds[1];
+ const th = try std.Thread.spawn(.{}, FakeServer.run, .{fs});
+ var s = try Session.connect(testing.allocator, .{ .fd = fds[0] }, fs.msize);
+ errdefer s.deinit();
+ _ = try s.attach(0, "me", "");
+ const fid = s.allocFid();
+ _ = try s.walk(0, fid, &.{"file"});
+ _ = try s.open(fid, cloud9.oread);
+ return .{ .server = fs, .session = s, .thread = th };
+ }
+};
+
+/// A fake interrupt source for the rpc wait loop: a SOCK_SEQPACKET pair stands
+/// in for the FUSE descriptor (one datagram per request, like /dev/fuse
+/// delivers one request per read). Mirrors the bridge's rules: an INTERRUPT for
+/// `unique` arms and cancels, anything else is consumed and ignored.
+pub const PipeInterrupt = struct {
+ read_end: i32,
+ write_end: i32,
+ unique: u64,
+ armed_flag: bool = false,
+ /// Requests consumed that were not the matching INTERRUPT.
+ ignored: usize = 0,
+
+ const fuse_interrupt_opcode: u32 = 36;
+
+ pub fn init(unique: u64) !PipeInterrupt {
+ var fds: [2]i32 = undefined;
+ const flags = linux.SOCK.SEQPACKET | linux.SOCK.CLOEXEC | linux.SOCK.NONBLOCK;
+ if (linux.errno(linux.socketpair(linux.AF.UNIX, flags, 0, &fds)) != .SUCCESS) return error.Io;
+ return .{ .read_end = fds[0], .write_end = fds[1], .unique = unique };
+ }
+
+ pub fn deinit(p: *PipeInterrupt) void {
+ _ = linux.close(p.read_end);
+ _ = linux.close(p.write_end);
+ }
+
+ pub fn interface(p: *PipeInterrupt) Interrupt {
+ return .{ .ctx = p, .watch = watch, .onReadable = onReadable, .armed = armed };
+ }
+
+ /// Writes a FUSE_INTERRUPT request (InHeader + InterruptIn) naming `target`.
+ pub fn inject(p: *PipeInterrupt, target: u64) !void {
+ var wire: [48]u8 = undefined;
+ std.mem.writeInt(u32, wire[0..4], 48, .little); // len
+ std.mem.writeInt(u32, wire[4..8], fuse_interrupt_opcode, .little); // opcode
+ std.mem.writeInt(u64, wire[8..16], 0x8000_0000_0000_0001, .little); // the interrupt's own unique
+ @memset(wire[16..40], 0); // nodeid, uid, gid, pid, extlen, padding
+ std.mem.writeInt(u64, wire[40..48], target, .little); // InterruptIn.unique
+ if (linux.write(p.write_end, &wire, wire.len) != wire.len) return error.Io;
+ }
+
+ fn watch(ctx: *anyopaque) i32 {
+ const p: *PipeInterrupt = @ptrCast(@alignCast(ctx));
+ return p.read_end;
+ }
+
+ fn onReadable(ctx: *anyopaque) Session.Error!bool {
+ const p: *PipeInterrupt = @ptrCast(@alignCast(ctx));
+ var buf: [4096]u8 = undefined;
+ const rc = linux.read(p.read_end, &buf, buf.len);
+ if (linux.errno(rc) != .SUCCESS or rc != 48) return error.Io;
+ const opcode = std.mem.readInt(u32, buf[4..8], .little);
+ const target = std.mem.readInt(u64, buf[40..48], .little);
+ if (opcode == fuse_interrupt_opcode and target == p.unique) {
+ p.armed_flag = true;
+ return true;
+ }
+ p.ignored += 1;
+ return false;
+ }
+
+ fn armed(ctx: *anyopaque) bool {
+ const p: *PipeInterrupt = @ptrCast(@alignCast(ctx));
+ return p.armed_flag;
+ }
};
+test "rpc wait loop: INTERRUPT → Tflush → Rflush → error.Interrupted; session still usable" {
+ var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ var pi = try PipeInterrupt.init(77);
+ defer pi.deinit();
+ s.interrupt = pi.interface();
+ defer s.interrupt = null;
+
+ // An INTERRUPT for some other request is consumed and ignored: the wait
+ // goes on, and the reply (a read past the hang offset) arrives normally.
+ var buf: [100]u8 = undefined;
+ try pi.inject(78);
+ try testing.expectEqual(@as(usize, 40), try s.read(1, 10, &buf));
+ try testing.expect(!pi.armed_flag);
+
+ // The read at offset 0 hangs; the INTERRUPT for our request cancels it
+ // (the stray one above is consumed along the way if the reply beat it).
+ try pi.inject(77);
+ try testing.expectError(error.Interrupted, s.read(1, 0, &buf));
+ try testing.expect(pi.armed_flag);
+ try testing.expectEqual(@as(usize, 1), pi.ignored);
+ try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst));
+ try testing.expectEqual(fs.hung_tag.load(.seq_cst), fs.flush_oldtag.load(.seq_cst));
+
+ // Both tags are free again: further rpcs work.
+ const st = try s.stat(1);
+ try testing.expectEqual(@as(u64, 50), st.length);
+ try testing.expectEqual(@as(usize, 40), try s.read(1, 10, &buf));
+ try testing.expectEqual(@as(usize, 0), s.client.pending());
+}
+
+test "rpc wait loop: the reply beats the Rflush → data returned, stray Rflush swallowed" {
+ var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .reply_then_rflush };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ var pi = try PipeInterrupt.init(5);
+ defer pi.deinit();
+ s.interrupt = pi.interface();
+ defer s.interrupt = null;
+
+ var buf: [40]u8 = undefined;
+ try pi.inject(5);
+ try testing.expectEqual(@as(usize, 30), try s.read(1, 0, buf[0..30]));
+ for (buf[0..30], 0..) |b, i| try testing.expectEqual(@as(u8, @truncate(i)), b);
+ try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst));
+ try testing.expectEqual(fs.hung_tag.load(.seq_cst), fs.flush_oldtag.load(.seq_cst));
+
+ // The Rflush is still in flight (or already buffered): the next rpcs must
+ // step over it, and afterwards nothing is pending in the client.
+ pi.armed_flag = false;
+ const st = try s.stat(1);
+ try testing.expectEqual(@as(u64, 50), st.length);
+ try testing.expectEqual(@as(usize, 40), try s.read(1, 10, &buf));
+ try testing.expectEqual(@as(usize, 0), s.client.pending());
+}
+
+test "rpc wait loop: a read interrupted after some data is a short read" {
+ // Chunks of maxRead = 1013 (msize 1024); the third chunk (offset 2026) hangs
+ // and the server fires the INTERRUPT at that moment.
+ var pi = try PipeInterrupt.init(9);
+ defer pi.deinit();
+ var fs: FakeServer = .{ .fd = -1, .msize = 1024, .file_len = 5000, .hang_offset = 2026, .on_flush = .rflush, .on_hang_inject = &pi };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ s.interrupt = pi.interface();
+ defer s.interrupt = null;
+
+ var buf: [4000]u8 = undefined;
+ try testing.expectEqual(@as(usize, 2026), try s.read(1, 0, &buf));
+ for (buf[0..2026], 0..) |b, i| try testing.expectEqual(@as(u8, @truncate(i)), b);
+ try testing.expect(pi.armed_flag);
+ try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst));
+ try testing.expectEqual(@as(usize, 0), s.client.pending());
+
+ // An INTERRUPT already waiting when the read starts: the first chunk's
+ // reply races the flush and wins, the armed flag then stops the loop.
+ pi.armed_flag = false;
+ try pi.inject(9);
+ try testing.expectEqual(@as(usize, 1013), try s.read(1, 0, &buf));
+ // The stray Rflush is consumed by the next call.
+ _ = try s.stat(1);
+ try testing.expectEqual(@as(usize, 0), s.client.pending());
+}
+
test "session against an in-process cloud9.Server" {
var fds: [2]i32 = undefined;
try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.socketpair(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0, &fds)));