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.zig55
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)));