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.zig167
1 files changed, 142 insertions, 25 deletions
diff --git a/9ns/src/nine.zig b/9ns/src/nine.zig
index cf9c49a..6f1fe09 100644
--- a/9ns/src/nine.zig
+++ b/9ns/src/nine.zig
@@ -71,6 +71,20 @@ pub const Session = struct {
/// Optional interrupt source (the bridge's FUSE descriptor) consulted while
/// a reply is outstanding; see `Interrupt`.
interrupt: ?Interrupt = null,
+ /// While set, a cancellation ends the rpc at once with `error.Interrupted`
+ /// and wedges the session, with no Tflush: for the handshake (version,
+ /// attach, the root stat), where there is no session yet to flush a
+ /// request out of, and the honest answer to "stop waiting" is to hang up.
+ abort_on_cancel: bool = false,
+ /// Milliseconds a Tflush may go unanswered before the server is declared
+ /// wedged: the rpc fails with `error.Interrupted` and the session with it.
+ /// The protocol says a client waits for the Rflush; a server that has not
+ /// managed one in this long is not going to, and the process behind the
+ /// interrupt is unkillable until we stop waiting. 0 waits forever.
+ flush_grace_ms: i32 = 3000,
+ /// The server is gone as far as this session is concerned (see
+ /// `abort_on_cancel`, `flush_grace_ms`); every rpc answers `error.Closed`.
+ wedged: bool = false,
/// Connect to `address`, then negotiate the protocol version.
/// `msize` is the maximum message size to ask for (0 = the buffers' size).
@@ -81,23 +95,33 @@ pub const Session = struct {
/// `connect`, with `stop_fd` watched for the whole handshake (the version
/// rpc included). A server that accepts the connection and then never
/// answers the Tversion would otherwise pin the caller in a blocking read
- /// with no way out: 9ns's mntgen dispatcher dials on the strength of the
- /// program's walk, so it must come back when that program is gone. The
- /// field stays set on the returned session, so the attach and stat that
- /// follow a dial keep watching it too; -1 disables the watch.
+ /// with no way out; -1 disables the watch. The field stays set on the
+ /// returned session. An `.fd` address is the caller's to close on failure.
pub fn connectWatched(gpa: std.mem.Allocator, address: Address, msize: u32, stop_fd: i32) !Session {
+ var s = try dial(gpa, address, msize, stop_fd);
+ s.version() catch |e| {
+ s.freeBuffers();
+ if (address != .fd) _ = linux.close(s.fd);
+ return e;
+ };
+ return s;
+ }
+
+ /// The transport and the buffers, no handshake: for a caller that wants
+ /// its interrupt source in place before `version()` (9ns's mntgen
+ /// worker, so a walk interrupted mid-dial can abandon the dial). Owns
+ /// the descriptor from here: `deinit` closes it.
+ pub fn dial(gpa: std.mem.Allocator, address: Address, msize: u32, stop_fd: i32) !Session {
const want: u32 = if (msize == 0) 8192 else @max(msize, 24);
const fd = try openTransport(address);
errdefer if (address != .fd) {
_ = linux.close(fd);
};
-
const in_buf = try gpa.alloc(u8, want);
errdefer gpa.free(in_buf);
const out_buf = try gpa.alloc(u8, want);
errdefer gpa.free(out_buf);
-
- var s: Session = .{
+ return .{
.gpa = gpa,
.fd = fd,
.client = .init(.{ .in = in_buf, .out = out_buf }),
@@ -106,19 +130,26 @@ pub const Session = struct {
.msize = want,
.stop_fd = stop_fd,
};
- const r = try s.rpc(.{ .version = .{ .msize = want } });
+ }
+
+ /// Negotiates the protocol version with the msize `dial` was given.
+ pub fn version(s: *Session) Error!void {
+ const r = try s.rpc(.{ .version = .{ .msize = s.msize } });
if (!std.mem.eql(u8, r.version.version, "9P2000")) return error.Protocol;
s.msize = r.version.msize;
- return s;
}
- /// Closes the descriptor and frees the buffers. Fids are not clunked.
- pub fn deinit(s: *Session) void {
- _ = linux.close(s.fd);
+ fn freeBuffers(s: *Session) void {
s.free_fids.deinit(s.gpa);
s.iounits.deinit(s.gpa);
s.gpa.free(s.in_buf);
s.gpa.free(s.out_buf);
+ }
+
+ /// Closes the descriptor and frees the buffers. Fids are not clunked.
+ pub fn deinit(s: *Session) void {
+ _ = linux.close(s.fd);
+ s.freeBuffers();
s.* = undefined;
}
@@ -154,8 +185,12 @@ pub const Session = struct {
/// 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).
+ /// `error.Interrupted`; the server has dropped the request), or neither
+ /// within `flush_grace_ms` (→ `error.Interrupted`, and the session is
+ /// wedged: the server stopped talking). Under `abort_on_cancel` the
+ /// cancellation itself wedges the session, with nothing sent.
pub fn rpc(s: *Session, req: cloud9.Client.Request) Error!cloud9.Client.Result {
+ if (s.wedged) return error.Closed;
s.ename_len = 0;
const tag = s.client.submit(req) catch |e| switch (e) {
error.NoTags, error.Handshake, error.Dead => return error.Protocol,
@@ -167,6 +202,8 @@ pub const Session = struct {
};
try s.flush();
var flush_tag: ?u16 = null;
+ // Monotonic ms by which the Tflush must have been answered.
+ var flush_deadline: ?i64 = null;
var tmp: [64 * 1024]u8 = undefined;
while (true) {
while (s.client.take()) |done| {
@@ -190,16 +227,28 @@ pub const Session = struct {
// 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;
- switch (try s.wait()) {
+ const grace: i32 = if (flush_deadline) |d| @intCast(@max(d - nowMs(), 0)) else -1;
+ switch (try s.wait(grace)) {
.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();
+ .cancel => {
+ if (s.abort_on_cancel) {
+ s.wedged = true;
+ return error.Interrupted;
+ }
+ if (flush_tag == null) {
+ flush_tag = s.client.submit(.{ .flush = .{ .oldtag = tag } }) catch return error.Protocol;
+ try s.flush();
+ if (s.flush_grace_ms != 0) flush_deadline = nowMs() + s.flush_grace_ms;
+ }
+ },
+ .timeout => {
+ s.wedged = true;
+ return error.Interrupted;
},
}
}
@@ -306,13 +355,14 @@ pub const Session = struct {
return chunkSize(s.client.maxWrite(), s.iounits.get(fid) orelse 0);
}
- const Ready = enum { socket, cancel };
+ const Ready = enum { socket, cancel, timeout };
/// Blocks until the socket is readable (`.socket`), the interrupt source
- /// wants the request in flight cancelled (`.cancel`), or `stop_fd` fires
+ /// wants the request in flight cancelled (`.cancel`), `timeout_ms` passes
+ /// with neither (`.timeout`; -1 waits forever), 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 {
+ fn wait(s: *Session, timeout_ms: i32) Error!Ready {
while (true) {
var pfds: [3]linux.pollfd = undefined;
var n: usize = 0;
@@ -329,13 +379,14 @@ pub const Session = struct {
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);
+ if (n == 1 and timeout_ms < 0) return .socket;
+ const prc = linux.poll(&pfds, @intCast(n), timeout_ms);
switch (linux.errno(prc)) {
.SUCCESS => {},
.INTR, .AGAIN => continue,
else => return error.Io,
}
+ if (prc == 0) return .timeout;
// A reply that is already there wins over everything else.
if (pfds[0].revents != 0) return .socket;
if (stop_at) |i| {
@@ -368,6 +419,13 @@ pub const Session = struct {
}
};
+/// The monotonic clock in milliseconds: deadlines, not timestamps.
+fn nowMs() i64 {
+ var ts: linux.timespec = undefined;
+ _ = linux.clock_gettime(.MONOTONIC, &ts);
+ return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000);
+}
+
fn chunkSize(max: u32, iounit: u32) u32 {
if (iounit != 0 and iounit < max) return iounit;
return max;
@@ -723,6 +781,62 @@ test "interrupted chunk loops: partial count if data moved, Interrupted otherwis
// -- in-process server test ---------------------------------------------------------
+test "flush grace: a Tflush the server never answers wedges the session after the grace" {
+ var pi = try PipeInterrupt.init(7);
+ defer pi.deinit();
+ var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .ignore, .on_hang_inject = &pi };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ s.interrupt = pi.interface();
+ s.flush_grace_ms = 100;
+ // The read at offset 0 hangs; the injected INTERRUPT sends a Tflush; the
+ // server ignores it; the grace runs out.
+ var buf: [100]u8 = undefined;
+ try testing.expectError(error.Interrupted, s.read(1, 0, &buf));
+ try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst));
+ try testing.expect(s.wedged);
+ // From here on the session is closed for business, without another
+ // byte to the server.
+ try testing.expectError(error.Closed, s.read(1, 10, &buf));
+ try testing.expectError(error.Closed, s.clunk(1));
+ try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst));
+ // The socket is still the session's to close: the server's thread ends
+ // when `close` drops it.
+}
+
+test "flush grace: zero waits for the Rflush, however late" {
+ var pi = try PipeInterrupt.init(7);
+ defer pi.deinit();
+ var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush, .on_hang_inject = &pi, .read_delay_ns = 0 };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ s.interrupt = pi.interface();
+ s.flush_grace_ms = 0;
+ var buf: [100]u8 = undefined;
+ try testing.expectError(error.Interrupted, s.read(1, 0, &buf));
+ try testing.expect(!s.wedged);
+ // The session lives: the Rflush released the tag and reads go on.
+ try testing.expectEqual(@as(usize, 40), try s.read(1, 10, buf[0..40]));
+}
+
+test "abort_on_cancel: a cancellation ends the rpc at once, sends no Tflush, wedges the session" {
+ var pi = try PipeInterrupt.init(7);
+ defer pi.deinit();
+ var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush, .on_hang_inject = &pi };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ s.interrupt = pi.interface();
+ s.abort_on_cancel = true;
+ var buf: [100]u8 = undefined;
+ try testing.expectError(error.Interrupted, s.read(1, 0, &buf));
+ try testing.expect(s.wedged);
+ try testing.expectEqual(@as(u32, 0), fs.flushes.load(.seq_cst));
+ try testing.expectError(error.Closed, s.read(1, 10, &buf));
+}
+
/// 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.
/// Test support only (bridge.zig's tests use it too).
@@ -734,9 +848,11 @@ pub const FakeServer = struct {
/// 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,
+ /// What a Tflush gets: the connection dropped (`.hangup`), an Rflush,
+ /// first the Rread the flush was aimed at and then the Rflush (the
+ /// race), or nothing at all (`.ignore`: counted, never answered, the
+ /// read stays hung — a wedged server).
+ on_flush: enum { hangup, rflush, reply_then_rflush, ignore } = .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),
@@ -827,6 +943,7 @@ pub const FakeServer = struct {
switch (fs.on_flush) {
// The test's "hang up now" signal.
.hangup => return,
+ .ignore => {},
.rflush => {
if (hung != null and hung.?.tag == m.oldtag) hung = null;
try srv.reply(tag, .rflush);