diff options
Diffstat (limited to '9ns/src/nine.zig')
| -rw-r--r-- | 9ns/src/nine.zig | 55 |
1 files changed, 55 insertions, 0 deletions
diff --git a/9ns/src/nine.zig b/9ns/src/nine.zig index c89a343..cf9c49a 100644 --- a/9ns/src/nine.zig +++ b/9ns/src/nine.zig @@ -75,6 +75,17 @@ pub const Session = struct { /// Connect to `address`, then negotiate the protocol version. /// `msize` is the maximum message size to ask for (0 = the buffers' size). pub fn connect(gpa: std.mem.Allocator, address: Address, msize: u32) !Session { + return connectWatched(gpa, address, msize, -1); + } + + /// `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. + pub fn connectWatched(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) { @@ -93,6 +104,7 @@ pub const Session = struct { .in_buf = in_buf, .out_buf = out_buf, .msize = want, + .stop_fd = stop_fd, }; const r = try s.rpc(.{ .version = .{ .msize = want } }); if (!std.mem.eql(u8, r.version.version, "9P2000")) return error.Protocol; @@ -1037,6 +1049,49 @@ test "rpc wait loop: a read interrupted after some data is a short read" { try testing.expectEqual(@as(usize, 0), s.client.pending()); } +test "connectWatched: a silent server cannot pin the handshake past stop_fd" { + // A server that accepts and then never answers: the version handshake has + // nothing to read. With a readable stop_fd the connect must come back with + // error.Stopped instead of blocking in readSocket (the fd is blocking), and + // the caller's descriptor must survive: `Address.fd` is not ours to close. + var sv: [2]i32 = undefined; + try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.socketpair(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0, &sv))); + defer _ = linux.close(sv[0]); + defer _ = linux.close(sv[1]); + var p: [2]i32 = undefined; + try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.pipe2(&p, .{ .CLOEXEC = true, .NONBLOCK = true }))); + defer _ = linux.close(p[0]); + defer _ = linux.close(p[1]); + try testing.expectEqual(@as(usize, 1), linux.write(p[1], "x", 1)); + + const Probe = struct { + const Self = @This(); + done: std.atomic.Value(bool) = .init(false), + stopped: std.atomic.Value(bool) = .init(false), + fd_open: std.atomic.Value(bool) = .init(false), + + fn run(w: *Self, client: i32, stop: i32) void { + if (Session.connectWatched(testing.allocator, .{ .fd = client }, 8192, stop)) |session| { + var s = session; + s.deinit(); + } else |e| w.stopped.store(e == error.Stopped, .release); + w.fd_open.store(linux.errno(linux.fcntl(client, linux.F.GETFD, 0)) == .SUCCESS, .release); + w.done.store(true, .release); + } + }; + var w: Probe = .{}; + const th = try std.Thread.spawn(.{}, Probe.run, .{ &w, sv[0], p[0] }); + var waited_ms: usize = 0; + while (!w.done.load(.acquire) and waited_ms < 3000) : (waited_ms += 10) { + const ts: linux.timespec = .{ .sec = 0, .nsec = 10 * std.time.ns_per_ms }; + _ = linux.nanosleep(&ts, null); + } + try testing.expect(w.done.load(.acquire)); + try testing.expect(w.stopped.load(.acquire)); + try testing.expect(w.fd_open.load(.acquire)); + th.join(); +} + 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))); |
