diff options
Diffstat (limited to '9ns/src/nine.zig')
| -rw-r--r-- | 9ns/src/nine.zig | 167 |
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); |
