diff options
Diffstat (limited to 'src/9p_io.zig')
| -rw-r--r-- | src/9p_io.zig | 290 |
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 { |
