summaryrefslogtreecommitdiff
path: root/src/9p_io.zig
diff options
context:
space:
mode:
Diffstat (limited to 'src/9p_io.zig')
-rw-r--r--src/9p_io.zig290
1 files changed, 271 insertions, 19 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig
index 4f7c2e97..3d019e0d 100644
--- a/src/9p_io.zig
+++ b/src/9p_io.zig
@@ -596,7 +596,7 @@ pub const Listener = struct {
}
for (l.control) |fd| {
setCloexec(fd);
- setNonblock(fd);
+ _ = setNonblock(fd);
}
l.arm();
l.watcher = try std.Thread.spawn(.{}, watchQuic, .{l});
@@ -671,6 +671,7 @@ pub const Listener = struct {
log.warn("registry post skipped: cannot create {s}", .{svc});
return;
}
+ sweepRegistry(l.io, svc);
var entry_buf: [sun_path_len:0]u8 = undefined;
const entry = std.fmt.bufPrintSentinel(&entry_buf, "{s}/{s}", .{ svc, name }, 0) catch {
log.warn("registry post skipped: name too long: {s}", .{name});
@@ -679,15 +680,15 @@ pub const Listener = struct {
const target = l.path_buf[0..l.path_len];
// Replace only what is provably not live: our own entry, or a
- // dead predecessor's symlink. `isListening` treats uncertainty
- // as live, so it is only asked about an entry that *is* a
- // symlink; anything else is left strictly alone and the
- // symlink below simply fails.
+ // dead predecessor's symlink. `probe` treats uncertainty as
+ // live, so it is only asked about an entry that *is* a symlink;
+ // anything else is left strictly alone and the symlink below
+ // simply fails.
var link_buf: [sun_path_len]u8 = undefined;
const n = libc.readlink(entry, &link_buf, link_buf.len);
if (n >= 0) {
const had = link_buf[0..@intCast(n)];
- if (!std.mem.eql(u8, had, target) and alive(entry)) {
+ if (!std.mem.eql(u8, had, target) and probe(entry) == .live) {
log.warn("registry entry pardes/{s} is live; not re-posted", .{name});
return;
}
@@ -842,12 +843,129 @@ pub fn listen(io: std.Io, gpa: std.mem.Allocator, named: []const u8, fallback: [
const alive = transport.isListening;
-fn setNonblock(fd: c_int) void {
+/// What is at a socket path, in cloud9's three answers (`post.Probe`).
+/// Asked here with libc rather than through `cloud9.post.probe`, whose
+/// raw Linux syscalls darwin cannot compile, but the classification is
+/// that file's and must not become a second opinion: the one definite
+/// refusal is `.stale`, nothing at the path at all is `.none`, and
+/// everything else — connected, busy, refused permission, a surprise —
+/// is `.live`. Uncertainty belongs to the server that owns the socket,
+/// never to a sweeper deciding what to delete.
+const Probe = enum { none, stale, live };
+
+/// `SOCK.STREAM`, and the kernel's own non-blocking bit where there is
+/// one. Darwin's `SOCK.NONBLOCK` is a Zig shim for `std.posix.socket`
+/// to unpack, not an ABI value, so handing it to the raw libc call
+/// would ask a kernel that has never heard of it; there it is an
+/// `fcntl` instead.
+const probe_socket_kind: c_uint = libc.SOCK.STREAM | (if (darwin) 0 else libc.SOCK.NONBLOCK);
+
+fn probe(path: [:0]const u8) Probe {
+ if (path.len + 1 > sun_path_len) return .live; // cannot ask; assume occupied
+ var addr: libc.sockaddr.un = .{ .path = @splat(0) };
+ @memcpy(addr.path[0 .. path.len + 1], path[0 .. path.len + 1]);
+ const fd = libc.socket(libc.AF.UNIX, probe_socket_kind, 0);
+ if (fd < 0) return .live;
+ defer _ = libc.close(fd);
+
+ // Non-blocking is the whole safety of this function, so it is read
+ // back rather than assumed: a BLOCKING connect to a live server
+ // whose backlog is full parks in the kernel with no timeout to end
+ // it — measured, it simply never returns — and a sweep that parks
+ // takes the editor's startup with it. An fd that cannot be proven
+ // non-blocking is never connected at all, which lands on `.live`,
+ // the answer that deletes nothing.
+ if (comptime darwin) _ = setNonblock(fd);
+ if (!isNonblocking(fd)) return .live;
+
+ if (libc.connect(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) == 0) return .live;
+ // cloud9's `post.probe` mapping, answer for answer. The connection
+ // is never wanted: EAGAIN (a full Unix backlog) and EINPROGRESS say
+ // somebody is listening, which is all that was asked, and the fd
+ // closes without ever waiting for it to become writable.
+ return switch (libc.errno(-1)) {
+ .AGAIN, .INPROGRESS, .PERM, .ACCES => .live,
+ .CONNREFUSED => .stale,
+ .NOENT, .NOTDIR => .none,
+ else => .live,
+ };
+}
+
+/// Reaps the registry entries whose editor is gone, once per post.
+///
+/// A clean exit unposts itself (`unpostFromRegistry`), so everything
+/// left behind comes from an exit that could not: an aborted test, a
+/// kill, a crash. No code in the dead process can ever run, so the cure
+/// has to be somebody else's readdir, and the next editor to start is
+/// the somebody. It is only ever a symlink that is followed or
+/// unlinked, because that is the only thing `postToRegistry` makes and
+/// anything else under a name belongs to whoever put it there; and only
+/// a definite refusal counts as gone. The socket a stale entry points
+/// at goes too, but not before a stat agrees it is a socket of ours: a
+/// plain file answers a connect with the same refusal.
+fn sweepRegistry(io: std.Io, svc: [:0]const u8) void {
+ if (comptime !supported) return;
+
+ // The names are staged before anything is unlinked, so the sweep
+ // never asks a directory to keep reading while it is being edited.
+ var names: [4096]u8 = undefined;
+ var staged: usize = 0;
+ {
+ const dir = std.Io.Dir.openDirAbsolute(io, svc, .{ .iterate = true }) catch return;
+ defer dir.close(io);
+ var read_buf: [std.Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined;
+ var reader: std.Io.Dir.Reader = .init(dir, &read_buf);
+ while (true) {
+ const listed = (reader.next(io) catch break) orelse break;
+ if (listed.name.len == 0 or listed.name.len > 255) continue;
+ if (staged + 1 + listed.name.len > names.len) break;
+ names[staged] = @intCast(listed.name.len);
+ @memcpy(names[staged + 1 ..][0..listed.name.len], listed.name);
+ staged += 1 + listed.name.len;
+ }
+ }
+
+ var reaped: usize = 0;
+ var i: usize = 0;
+ while (i < staged) {
+ const name = names[i + 1 ..][0..names[i]];
+ i += 1 + name.len;
+ var entry_buf: [sun_path_len:0]u8 = undefined;
+ const entry = std.fmt.bufPrintSentinel(&entry_buf, "{s}/{s}", .{ svc, name }, 0) catch continue;
+ const facts = statNoFollow(entry) orelse continue;
+ if (facts.mode & 0o170000 != 0o120000) continue; // not a symlink: not ours to judge
+ const state = probe(entry);
+ if (state == .live) continue;
+ var link_buf: [sun_path_len]u8 = undefined;
+ const n = libc.readlink(entry, &link_buf, link_buf.len);
+ if (libc.unlink(entry) != 0) continue;
+ reaped += 1;
+ // A relative target would resolve against this editor's working
+ // directory, which says nothing about what the entry named.
+ if (state != .stale or n <= 0 or link_buf[0] != '/') continue;
+ var target_buf: [sun_path_len:0]u8 = undefined;
+ const target = std.fmt.bufPrintSentinel(&target_buf, "{s}", .{link_buf[0..@intCast(n)]}, 0) catch continue;
+ const t = statNoFollow(target) orelse continue;
+ if (t.mode & 0o170000 == 0o140000 and t.uid == libc.getuid()) _ = libc.unlink(target);
+ }
+ if (reaped != 0) log.info("reaped {d} stale registry entries under {s}", .{ reaped, svc });
+}
+
+/// Whether it worked, because `probe` is not allowed to find out the
+/// hard way: a caller that needs the guarantee has to be able to check.
+fn setNonblock(fd: c_int) bool {
const flags = libc.fcntl(fd, libc.F.GETFL, @as(c_int, 0));
- if (flags < 0) return;
+ if (flags < 0) return false;
var o: libc.O = @bitCast(@as(u32, @bitCast(flags)));
o.NONBLOCK = true;
- _ = libc.fcntl(fd, libc.F.SETFL, @as(c_int, @bitCast(@as(u32, @bitCast(o)))));
+ return libc.fcntl(fd, libc.F.SETFL, @as(c_int, @bitCast(@as(u32, @bitCast(o))))) >= 0;
+}
+
+fn isNonblocking(fd: c_int) bool {
+ const flags = libc.fcntl(fd, libc.F.GETFL, @as(c_int, 0));
+ if (flags < 0) return false;
+ const o: libc.O = @bitCast(@as(u32, @bitCast(flags)));
+ return o.NONBLOCK;
}
const testing = std.testing;
@@ -866,6 +984,108 @@ test "a name that is not one path component is no address at all" {
try testing.expect(socketPath(&buf, "/run", "a\x00b") == null);
}
+test "the registry sweep takes the dead entries and leaves everything else" {
+ if (comptime !supported) return error.SkipZigTest;
+ // Under the real runtime directory, because a sun_path is 108 bytes
+ // and the test cache's temporary directories are longer than that.
+ var dir_buf: [sun_path_len:0]u8 = undefined;
+ const base = socketDir(&dir_buf) orelse return error.SkipZigTest;
+ if (!ensureSocketDir(base)) return error.SkipZigTest;
+
+ var svc_buf: [sun_path_len:0]u8 = undefined;
+ const svc = std.fmt.bufPrintSentinel(&svc_buf, "{s}/sweep-{d}", .{ base, @as(u32, @intCast(libc.getpid())) }, 0) catch
+ return error.SkipZigTest;
+ if (libc.mkdir(svc, 0o700) != 0) return error.SkipZigTest;
+
+ var paths: [5][sun_path_len:0]u8 = undefined;
+ const live_sock = try std.fmt.bufPrintSentinel(&paths[0], "{s}/live.sock", .{svc}, 0);
+ const dead_sock = try std.fmt.bufPrintSentinel(&paths[1], "{s}/dead.sock", .{svc}, 0);
+ const live = try std.fmt.bufPrintSentinel(&paths[2], "{s}/live", .{svc}, 0);
+ const dead = try std.fmt.bufPrintSentinel(&paths[3], "{s}/dead", .{svc}, 0);
+ const stranger = try std.fmt.bufPrintSentinel(&paths[4], "{s}/stranger", .{svc}, 0);
+ defer {
+ for ([_][:0]const u8{ live_sock, dead_sock, live, dead, stranger }) |p| _ = libc.unlink(p);
+ _ = libc.rmdir(svc);
+ }
+
+ const listening = bindSocket(live_sock);
+ try testing.expect(listening >= 0);
+ defer _ = libc.close(listening);
+ try testing.expectEqual(@as(c_int, 0), libc.listen(listening, 1));
+
+ // Bound and then dropped: the file stays, and nobody answers it —
+ // exactly what an aborted test leaves behind.
+ const abandoned = bindSocket(dead_sock);
+ try testing.expect(abandoned >= 0);
+ _ = libc.close(abandoned);
+
+ try testing.expectEqual(@as(c_int, 0), libc.symlink(live_sock, live));
+ try testing.expectEqual(@as(c_int, 0), libc.symlink(dead_sock, dead));
+ // Not a symlink, so not this program's to reason about, even though
+ // connecting to it is refused exactly like the dead socket.
+ try std.Io.Dir.cwd().writeFile(testing.io, .{ .sub_path = stranger, .data = "" });
+
+ // The guarantee itself, asserted rather than inferred from the fact
+ // that the test finished: if a probe socket ever stops being
+ // non-blocking the connect below parks instead of failing, and a
+ // parked test costs whoever is building far more than a red one.
+ const checking = libc.socket(libc.AF.UNIX, probe_socket_kind, 0);
+ try testing.expect(checking >= 0);
+ if (comptime darwin) try testing.expect(setNonblock(checking));
+ try testing.expect(isNonblocking(checking));
+ _ = libc.close(checking);
+
+ // The listener's backlog is filled before anything is asked, because
+ // a full backlog is the one place a blocking connect parks forever
+ // (`unix_wait_for_peer`, no timeout) and a sweep that only answers
+ // while nobody is queued is the sweep that hangs a build. Nothing
+ // below accepts any of these, so the queue stays full throughout.
+ var queued: [8]c_int = @splat(-1);
+ defer for (queued) |fd| {
+ if (fd >= 0) _ = libc.close(fd);
+ };
+ for (&queued) |*slot| {
+ const fd = libc.socket(libc.AF.UNIX, probe_socket_kind, 0);
+ if (fd < 0) break;
+ if (comptime darwin) _ = setNonblock(fd);
+ var addr: libc.sockaddr.un = .{ .path = @splat(0) };
+ @memcpy(addr.path[0 .. live_sock.len + 1], live_sock[0 .. live_sock.len + 1]);
+ _ = libc.connect(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr)));
+ slot.* = fd;
+ }
+
+ const started = transport.nowMs();
+ try testing.expectEqual(Probe.live, probe(live));
+ try testing.expectEqual(Probe.stale, probe(dead));
+
+ sweepRegistry(testing.io, svc);
+
+ // Promptly, and not "eventually": a blocking probe never comes back
+ // at all, so any wall-clock bound at all is the assertion that
+ // matters. A second is several thousand times what three connects
+ // and a readdir cost.
+ try testing.expect(transport.nowMs() - started < 1000);
+
+ try testing.expect(statNoFollow(dead) == null);
+ try testing.expect(statNoFollow(dead_sock) == null);
+ try testing.expect(statNoFollow(live) != null);
+ try testing.expect(statNoFollow(live_sock) != null);
+ try testing.expect(statNoFollow(stranger) != null);
+ try testing.expectEqual(Probe.live, probe(live));
+}
+
+fn bindSocket(path: [:0]const u8) c_int {
+ const fd = libc.socket(libc.AF.UNIX, libc.SOCK.STREAM, 0);
+ if (fd < 0) return fd;
+ var addr: libc.sockaddr.un = .{ .path = @splat(0) };
+ @memcpy(addr.path[0 .. path.len + 1], path[0 .. path.len + 1]);
+ if (libc.bind(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) != 0) {
+ _ = libc.close(fd);
+ return -1;
+ }
+ return fd;
+}
+
test "one connection's buffers are sized from the one msize constant" {
try testing.expect(msize >= ninep.min_msize);
try testing.expectEqual(@as(u32, msize), Runner.msize);
@@ -1102,24 +1322,41 @@ pub fn start(io: std.Io, gpa: std.mem.Allocator, core: *pardes.Pardes) ?*Listene
return listener;
}
-pub fn exportPaneEnv(listener: ?*const Listener, serial: u32, forward_look: bool) void {
- _ = unsetenv("PARDES_FORWARD_LOOK");
+/// What a pane shell is told about the editor above it, set into this
+/// process's environment just before `forkpty` so the child inherits it.
+///
+/// Two independent facts, and they are separate variables because they answer
+/// separate questions. `PARDES_PID` says "you are inside this editor, and it
+/// will take a Look from you": a session whose 9P listener never came up still
+/// owns its children, so the child says so rather than looking, from its own
+/// side, exactly like no pardes at all. `PARDES_9P` and `PARDES_PANE` say how
+/// to reach it, and a `--nested` session exports them too — its socket stays
+/// open to scripts and to mounted shells — while withholding `PARDES_PID`,
+/// which is the whole of what `--nested` means.
+pub fn exportPaneEnv(listener: ?*const Listener, serial: u32, adopts: bool) void {
+ var announced = false;
+ if (adopts) announcing: {
+ var buf: [16]u8 = undefined;
+ const text = std.fmt.bufPrintSentinel(&buf, "{d}", .{@as(u32, @intCast(libc.getpid()))}, 0) catch
+ break :announcing;
+ announced = setenv("PARDES_PID", text, 1) == 0;
+ }
+ if (!announced) _ = unsetenv("PARDES_PID");
+
if (listener) |l| exporting: {
var sock: [sun_path_len]u8 = undefined;
const path = std.fmt.bufPrintSentinel(&sock, "{s}", .{l.path()}, 0) catch break :exporting;
var buf: [16]u8 = undefined;
const id = std.fmt.bufPrintSentinel(&buf, "{d}", .{serial}, 0) catch break :exporting;
if (setenv("PARDES_9P", path, 1) != 0) break :exporting;
- if (setenv("PARDES_PANE", id, 1) != 0) break :exporting;
- if (setenv("PARDES_FORWARD_LOOK", if (forward_look) "1" else "0", 1) == 0) return;
+ if (setenv("PARDES_PANE", id, 1) == 0) return;
}
_ = unsetenv("PARDES_9P");
_ = unsetenv("PARDES_PANE");
- _ = setenv("PARDES_FORWARD_LOOK", "0", 1);
}
-test "9P shell environment preserves identity when nested Look forwarding is disabled" {
- const names = [_][*:0]const u8{ "XDG_RUNTIME_DIR", "HOME", "PARDES_9P", "PARDES_PANE", "PARDES_FORWARD_LOOK" };
+test "9P shell environment states being inside pardes apart from how to reach it" {
+ const names = [_][*:0]const u8{ "XDG_RUNTIME_DIR", "HOME", "PARDES_PID", "PARDES_9P", "PARDES_PANE" };
var saved: [names.len]?[:0]u8 = @splat(null);
for (names, &saved) |name, *value| {
if (libc.getenv(name)) |old| value.* = try testing.allocator.dupeZ(u8, std.mem.span(old));
@@ -1146,17 +1383,32 @@ test "9P shell environment preserves identity when nested Look forwarding is dis
const path = "/tmp/pardes-example.sock";
@memcpy(listener.path_buf[0..path.len], path);
listener.path_len = path.len;
+ var own: [16]u8 = undefined;
+ const own_pid = try std.fmt.bufPrint(&own, "{d}", .{@as(u32, @intCast(libc.getpid()))});
+
+ // A `--nested` session: reachable for scripts and mounts, but nobody's
+ // parent, so a pardes started in one of its shells runs a session of
+ // its own instead of handing its argument over.
exportPaneEnv(listener, 7, false);
try testing.expectEqualStrings(path, std.mem.span(libc.getenv("PARDES_9P").?));
try testing.expectEqualStrings("7", std.mem.span(libc.getenv("PARDES_PANE").?));
- try testing.expectEqualStrings("0", std.mem.span(libc.getenv("PARDES_FORWARD_LOOK").?));
+ try testing.expect(libc.getenv("PARDES_PID") == null);
+
exportPaneEnv(listener, 8, true);
try testing.expectEqualStrings("8", std.mem.span(libc.getenv("PARDES_PANE").?));
- try testing.expectEqualStrings("1", std.mem.span(libc.getenv("PARDES_FORWARD_LOOK").?));
+ try testing.expectEqualStrings(own_pid, std.mem.span(libc.getenv("PARDES_PID").?));
+
+ // A session whose listener never came up is still the session this shell
+ // is inside: the pid stands on its own, and only the address is missing.
+ exportPaneEnv(null, 0, true);
+ try testing.expect(libc.getenv("PARDES_9P") == null);
+ try testing.expect(libc.getenv("PARDES_PANE") == null);
+ try testing.expectEqualStrings(own_pid, std.mem.span(libc.getenv("PARDES_PID").?));
+
exportPaneEnv(null, 0, false);
try testing.expect(libc.getenv("PARDES_9P") == null);
try testing.expect(libc.getenv("PARDES_PANE") == null);
- try testing.expectEqualStrings("0", std.mem.span(libc.getenv("PARDES_FORWARD_LOOK").?));
+ try testing.expect(libc.getenv("PARDES_PID") == null);
}
pub const Client = struct {