summaryrefslogtreecommitdiff
path: root/9ns/src/nine.zig
diff options
context:
space:
mode:
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)));