diff options
Diffstat (limited to 'src')
| -rw-r--r-- | src/9p_io.zig | 144 |
1 files changed, 143 insertions, 1 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig index 090eb838..4f7c2e97 100644 --- a/src/9p_io.zig +++ b/src/9p_io.zig @@ -241,6 +241,10 @@ pub const Listener = struct { paused_ms: i64 = 0, path_buf: [sun_path_len]u8 = undefined, path_len: usize = 0, + /// The registry entry this editor posted + /// (`$XDG_RUNTIME_DIR/9p/pardes/<name>`), zero when it did not. + posted_buf: [sun_path_len]u8 = undefined, + posted_len: usize = 0, conns: [quic_slots]Conn = @splat(.{}), control: [2]c_int = .{ -1, -1 }, watcher: ?std.Thread = null, @@ -635,8 +639,91 @@ pub const Listener = struct { } } + /// Advertises this editor in the posted-9P registry so a + /// `9ns --mntgen` mount lists it and dials it on a walk: + /// `$XDG_RUNTIME_DIR/9p/pardes/<name>` is a symlink to the socket + /// pardes already binds. One directory for the program, one entry + /// per running editor — the layout zmx uses for its sessions, so + /// several editors group instead of crowding the registry root. + /// + /// Only the runtime-directory socket posts: an instance that fell + /// back to `~/.local/state/pardes` stays out of the user's + /// registry, the way a private `ZMX_DIR` does for zmx. + /// + /// The socket itself is still bound by `listen` above, not by + /// `cloud9.post`: `post` claims flat names only, and its claim + /// protocol derives its lock directory from the registry path, so + /// a name inside a subdirectory cannot go through it yet. See + /// docs/cloud9.md. + fn postToRegistry(l: *Listener, name: []const u8) void { + if (comptime !supported) return; + if (name.len == 0 or std.mem.indexOfAny(u8, name, "/\x00") != null) return; + const xdg_c = libc.getenv("XDG_RUNTIME_DIR") orelse return; + const xdg = std.mem.span(xdg_c); + if (xdg.len == 0) return; + + var reg_buf: [sun_path_len:0]u8 = undefined; + const reg = std.fmt.bufPrintSentinel(®_buf, "{s}/9p", .{xdg}, 0) catch return; + _ = libc.mkdir(reg, 0o750); + var svc_buf: [sun_path_len:0]u8 = undefined; + const svc = std.fmt.bufPrintSentinel(&svc_buf, "{s}/9p/pardes", .{xdg}, 0) catch return; + if (libc.mkdir(svc, 0o750) != 0 and statNoFollow(svc) == null) { + log.warn("registry post skipped: cannot create {s}", .{svc}); + return; + } + 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}); + return; + }; + 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. + 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)) { + log.warn("registry entry pardes/{s} is live; not re-posted", .{name}); + return; + } + _ = libc.unlink(entry); + } + var target_z: [sun_path_len:0]u8 = undefined; + const target_zs = std.fmt.bufPrintSentinel(&target_z, "{s}", .{target}, 0) catch return; + if (libc.symlink(target_zs, entry) != 0) { + log.warn("registry post skipped: pardes/{s} is occupied", .{name}); + return; + } + @memcpy(l.posted_buf[0..entry.len], entry); + l.posted_len = entry.len; + log.info("posted pardes/{s} -> {s}", .{ name, target }); + } + + /// Removes the registry entry, but only while it is still ours: a + /// name another editor has since claimed is never unlinked. + fn unpostFromRegistry(l: *Listener) void { + if (comptime !supported) return; + if (l.posted_len == 0) return; + var z: [sun_path_len:0]u8 = undefined; + @memcpy(z[0..l.posted_len], l.posted_buf[0..l.posted_len]); + z[l.posted_len] = 0; + const entry = z[0..l.posted_len :0]; + var link_buf: [sun_path_len]u8 = undefined; + const n = libc.readlink(entry, &link_buf, link_buf.len); + if (n >= 0 and std.mem.eql(u8, link_buf[0..@intCast(n)], l.path_buf[0..l.path_len])) { + _ = libc.unlink(entry); + } + l.posted_len = 0; + } + pub fn deinit(l: *Listener, gpa: std.mem.Allocator) void { l.stopping.store(true, .release); + l.unpostFromRegistry(); if (l.watcher) |thread| { l.watch_stop.store(true, .release); _ = libc.write(l.control[1], "q", 1); @@ -680,7 +767,8 @@ pub fn listen(io: std.Io, gpa: std.mem.Allocator, named: []const u8, fallback: [ .handler = .{ .ctx = l, .serve = Listener.onServe, .opened = Listener.onOpened, .closed = Listener.onClosed }, .greet_timeout_ms = Listener.greet_deadline_ms, }); - const p = socketPath(&l.path_buf, dir, if (named.len != 0) named else fallback) orelse { + const entry_name = if (named.len != 0) named else fallback; + const p = socketPath(&l.path_buf, dir, entry_name) orelse { gpa.destroy(l); return null; }; @@ -705,6 +793,7 @@ pub fn listen(io: std.Io, gpa: std.mem.Allocator, named: []const u8, fallback: [ l.deinit(gpa); return null; } + l.postToRegistry(entry_name); if (tcp_dial) |dial| { if (!std.mem.startsWith(u8, dial, "tcp!")) { l.deinit(gpa); @@ -840,6 +929,59 @@ test "same-session TCP mounts compare canonical endpoints and local wildcard des extern "c" fn mkdtemp(template: [*:0]u8) ?[*:0]u8; extern "c" fn rmdir(path: [*:0]const u8) c_int; +test "a listening editor posts itself into the 9P registry and unposts on stop" { + if (comptime !supported) return error.SkipZigTest; + const gpa = testing.allocator; + var directory: [64:0]u8 = undefined; + _ = try std.fmt.bufPrintSentinel(&directory, "/tmp/pardes-post-XXXXXX", .{}, 0); + if (mkdtemp(&directory) == null) return error.TempDirectoryFailed; + const dir = std.mem.span(@as([*:0]const u8, &directory)); + const old_runtime = if (libc.getenv("XDG_RUNTIME_DIR")) |v| try gpa.dupeZ(u8, std.mem.span(v)) else null; + defer { + if (old_runtime) |v| { + _ = setenv("XDG_RUNTIME_DIR", v, 1); + gpa.free(v); + } else _ = unsetenv("XDG_RUNTIME_DIR"); + } + try testing.expectEqual(@as(c_int, 0), setenv("XDG_RUNTIME_DIR", &directory, 1)); + + var entry_buf: [sun_path_len:0]u8 = undefined; + const entry = try std.fmt.bufPrintSentinel(&entry_buf, "{s}/9p/pardes/unit", .{dir}, 0); + var svc_buf: [sun_path_len:0]u8 = undefined; + const svc = try std.fmt.bufPrintSentinel(&svc_buf, "{s}/9p/pardes", .{dir}, 0); + var sock_buf: [sun_path_len:0]u8 = undefined; + const sock = try std.fmt.bufPrintSentinel(&sock_buf, "{s}/" ++ prefix ++ "unit.sock", .{dir}, 0); + defer { + _ = libc.unlink(entry); + _ = rmdir(svc); + var reg_buf: [sun_path_len:0]u8 = undefined; + if (std.fmt.bufPrintSentinel(®_buf, "{s}/9p", .{dir}, 0)) |reg| { + _ = rmdir(reg); + } else |_| {} + _ = libc.unlink(sock); + _ = rmdir(&directory); + } + + var link: [sun_path_len]u8 = undefined; + { + const l = listen(testing.io, gpa, "unit", "", null, null) orelse return error.ListenFailed; + defer l.deinit(gpa); + // The socket stays exactly where pardes has always bound it: + // adopting the registry moves nothing, it only advertises. + try testing.expect(statNoFollow(sock) != null); + // And the registry holds a symlink to it one directory down, so + // several editors group under /mnt/9p/pardes/ instead of + // crowding the registry root — the layout zmx posts its + // sessions in. + const n = libc.readlink(entry, &link, link.len); + try testing.expect(n > 0); + try testing.expectEqualStrings(sock, link[0..@intCast(n)]); + } + // Stopping takes the name back out, so the next editor of that name + // is not refused by its own corpse. + try testing.expect(libc.readlink(entry, &link, link.len) < 0); +} + test "Unix TCP and QUIC share one listener through reads writes reconnects and reset" { if (comptime !supported) return error.SkipZigTest; const gpa = testing.allocator; |
