diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-25 16:40:49 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-09-25 17:52:33 -0300 |
| commit | 1c1b192d4ef59199a7196229d56901b1f6512678 (patch) | |
| tree | 9f6e4c1efcf877d1231aa58975ccd923ad4f2f2e /9agents/src/active.zig | |
| parent | 73602127d15d10a1932b6fe916bd18a608054980 (diff) | |
| download | cloud9-1c1b192d4ef59199a7196229d56901b1f6512678.tar.gz cloud9-1c1b192d4ef59199a7196229d56901b1f6512678.zip | |
Each message of a claude or codex agent's transcript is a file named
<id>-<kind>: the id is its position in the transcript from 0, as 8
digits so the names sort; the kind is user, assistant, thinking, system
or agent, and a tool call is named after its tool (00000002-bash) and
its result after the call (00000003-bash-result). A file is the message
as text (a tool call: its name, then its input); its mtime is when it
was said.
chat.zig indexes a transcript incrementally, keyed by (dev, ino) and
checked by birth time and a hash of its first and last indexed bytes, so
a file rewritten in place is indexed again. It holds offsets, never
bytes: every read reads the file, and a read resumes where the last one
stopped so a long message is not decoded from its start each time. The
indexes and the runner live in lazily backed mappings (NORESERVE,
NOHUGEPAGE), and the fid table is 32768 so a kernel mount can hold every
message of a long chat.
A 9agents started inside a user namespace (from a 9ns-wrapped terminal)
cannot read /proc/<pid>/fd of processes outside it, so the fd route
never resolved codex there; /active now also finds a codex session
through the thread writer locks it holds (/proc/locks), via = lock.
Fixes to /active found on the way: a fifo in place of a transcript or
record no longer blocks the daemon; a transcript is opened from its
pinned root a component at a time, so a symlinked directory or a
sessionId with .. cannot reach outside it; slot generations are 24 bits,
so a stale node id cannot come to name another agent; listing a harness
rescans; an agent that had not resolved yet is asked again.
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Diffstat (limited to '9agents/src/active.zig')
| -rw-r--r-- | 9agents/src/active.zig | 207 |
1 files changed, 196 insertions, 11 deletions
diff --git a/9agents/src/active.zig b/9agents/src/active.zig index daaddc9..ebf7ed6 100644 --- a/9agents/src/active.zig +++ b/9agents/src/active.zig @@ -14,6 +14,7 @@ //! the exclusions are testable against a fixture tree. const std = @import("std"); const Io = std.Io; +const linux = std.os.linux; // ---- comptime bounds --------------------------------------------------------- @@ -58,6 +59,8 @@ pub const Via = enum(u8) { registry, fd, dir, + /// The writer lock the harness holds on its thread (codex). + lock, pub fn text(v: Via) []const u8 { return @tagName(v); @@ -272,7 +275,6 @@ pub fn sessionIdOf(kind: Kind, file_name: []const u8) []const u8 { switch (kind) { .omp => if (std.mem.lastIndexOfScalar(u8, stem, '_')) |i| return stem[i + 1 ..], .codex => { - const uuid_len = "01a0bcd2-5653-7cc0-aa36-676f6467a72a".len; if (stem.len >= uuid_len) return stem[stem.len - uuid_len ..]; }, else => {}, @@ -336,9 +338,27 @@ pub const Sources = struct { /// them in hundredths of a second whatever CONFIG_HZ is. const user_hz: u64 = 100; +/// The head of a small file: `/proc` records, a harness's session +/// record. Opened without blocking and without following a final +/// symlink, and read only when it is a regular file — a fifo left where +/// a record should be would otherwise park the whole daemon in `open`. fn readSmall(io: Io, path: []const u8, buf: []u8) ?[]const u8 { - const out = Io.Dir.readFile(.cwd(), io, path, buf) catch return null; - return out; + var z: [path_capacity + 1]u8 = @splat(0); + if (path.len == 0 or path.len >= z.len) return null; + @memcpy(z[0..path.len], path); + const rc = linux.open(@ptrCast(&z), .{ .ACCMODE = .RDONLY, .NONBLOCK = true, .NOFOLLOW = true, .CLOEXEC = true }, 0); + if (linux.errno(rc) != .SUCCESS) return null; + const file: Io.File = .{ .handle = @intCast(rc), .flags = .{ .nonblocking = false } }; + defer file.close(io); + return readHead(io, file, buf); +} + +/// The first bytes of an open regular file, or null for anything else. +pub fn readHead(io: Io, file: Io.File, buf: []u8) ?[]const u8 { + const st = file.stat(io) catch return null; + if (st.kind != .file) return null; + const n = file.readPositionalAll(io, buf, 0) catch return null; + return buf[0..n]; } fn linkOf(io: Io, path: []const u8, buf: []u8) ?[]const u8 { @@ -401,7 +421,7 @@ pub const Slot = struct { used: bool = false, /// Bumped whenever the slot comes to hold a different process, so a /// node id minted for the old one stops resolving. - gen: u8 = 0, + gen: u24 = 0, seen: bool = false, live: Live = .{}, }; @@ -423,7 +443,7 @@ pub const Table = struct { } /// The live record at `index`, when its generation still matches. - pub fn at(t: *Table, index: usize, gen: u8) ?*Live { + pub fn at(t: *Table, index: usize, gen: u24) ?*Live { if (index >= t.slots.len) return null; const sl = &t.slots[index]; if (!sl.used or sl.gen != gen) return null; @@ -488,9 +508,12 @@ pub const Table = struct { const kind = harnessOf(cmdline) orelse return; // Already known and still the same process: keep the slot and - // everything resolved for it. + // everything resolved for it. One that did not resolve is asked + // again: a harness seen in its first instant may not have written + // its session record or taken its thread lock yet. if (t.slotOfPid(pid, starttime)) |sl| { sl.seen = true; + if (sl.live.via == .none and sl.live.kind != .hermes) resolveSession(&sl.live, s, pid_dir); return; } @@ -543,7 +566,11 @@ fn resolveSession(l: *Live, s: Sources, pid_dir: []const u8) void { switch (l.kind) { .claude => resolveClaude(l, s), .omp => resolveOmp(l, s, pid_dir), - .dsh, .codex => resolveByFd(l, s, pid_dir), + .dsh => resolveByFd(l, s, pid_dir), + .codex => { + resolveByFd(l, s, pid_dir); + if (l.via == .none) resolveCodexByLock(l, s); + }, // hermes has no per-session file to find. .hermes => {}, } @@ -568,6 +595,8 @@ fn resolveClaude(l: *Live, s: Sources) void { const sid = jsonField(text, "sessionId") orelse return; if (sid.len == 0 or sid.len > text_capacity) return; + // It names a file below the root: an id, never a path. + for (sid) |ch| if (!std.ascii.isAlphanumeric(ch) and ch != '-' and ch != '_') return; l.session.set(sid); l.via = .registry; @@ -680,6 +709,151 @@ fn resolveByFd(l: *Live, s: Sources, pid_dir: []const u8) void { } } +/// A daemon started inside a user namespace — anything launched from a +/// terminal that `9ns` wraps, whose namespace maps only the user's uid — +/// cannot read `/proc/<pid>/fd` or `cwd` of a process outside it: the +/// kernel's ptrace access check fails across the namespace, whatever the +/// uid. The fd route then finds nothing. What stays visible is the +/// locks: codex holds a write `flock` on +/// `thread-writer-locks/<thread>.lock` for each thread it is writing, and +/// `/proc/locks` names the holder's pid and the locked file's inode. The +/// threads a process holds are its own and its subagents'; its own is the +/// root, whose rollout says `"thread_source":"user"`. Several roots (a +/// `/new` in the same process) resolve to the newest. +fn resolveCodexByLock(l: *Live, s: Sources) void { + const base = rootOf(s, .codex) orelse return; + // A codex running a team holds a lock per subagent thread too; the + // root's must not fall off the end of this list. + var inos: [512]u64 = undefined; + const held = heldLocks(s, l.pid, &inos); + if (held.len == 0) return; + + var dir_buf: [path_capacity]u8 = undefined; + const lock_dir = join(&dir_buf, &.{ base, "thread-writer-locks" }) orelse return; + const dir = Io.Dir.openDirAbsolute(s.io, lock_dir, .{ .iterate = true }) catch return; + defer Io.Dir.close(dir, s.io); + var rb: [Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined; + var it = Io.Dir.Reader.init(dir, &rb); + + var best: i64 = -1; + var best_path: [path_capacity]u8 = undefined; + var best_len: usize = 0; + var best_id: [uuid_len]u8 = undefined; + while (true) { + const entry = (it.next(s.io) catch break) orelse break; + if (!std.mem.endsWith(u8, entry.name, ".lock")) continue; + const id = entry.name[0 .. entry.name.len - ".lock".len]; + if (id.len != uuid_len) continue; + // By inode alone: on btrfs `/proc/locks` prints the superblock's + // device and stat the subvolume's, so the devices never agree. + // Every candidate is in this one directory, where inodes are + // unique anyway. + const st = Io.Dir.statFile(dir, s.io, entry.name, .{ .follow_symlinks = false }) catch continue; + if (st.kind != .file) continue; + if (std.mem.indexOfScalar(u64, held, st.inode) == null) continue; + var path_buf: [path_capacity]u8 = undefined; + const rollout = rolloutOf(s, base, id, &path_buf) orelse continue; + if (subagentRollout(s.io, rollout)) continue; + const when = mtimeOf(s.io, rollout); + if (when <= best) continue; + @memcpy(best_path[0..rollout.len], rollout); + @memcpy(&best_id, id); + best = when; + best_len = rollout.len; + } + if (best_len == 0) return; + l.transcript.set(best_path[0..best_len]); + l.session.set(&best_id); + l.via = .lock; +} + +const uuid_len = "01a0bcd2-5653-7cc0-aa36-676f6467a72a".len; + +/// The inodes `pid` holds a write `flock` on, from `<proc>/locks`: +/// +/// 1: FLOCK ADVISORY WRITE 657408 00:37:32607155 0 EOF +/// +/// A line whose second word is `->` is a waiter, not a holder. +fn heldLocks(s: Sources, pid: u32, out: []u64) []const u64 { + var path_buf: [path_capacity]u8 = undefined; + const path = join(&path_buf, &.{ s.proc, "locks" }) orelse return out[0..0]; + const file = Io.Dir.openFileAbsolute(s.io, path, .{}) catch return out[0..0]; + defer file.close(s.io); + var buf: [4096]u8 = undefined; + var r = file.readerStreaming(s.io, &buf); + var n: usize = 0; + while (n < out.len) { + const line = (r.interface.takeDelimiter('\n') catch break) orelse break; + var words = std.mem.tokenizeAny(u8, line, " \t"); + _ = words.next(); // "1:" + const kind = words.next() orelse continue; + if (!std.mem.eql(u8, kind, "FLOCK")) continue; + _ = words.next(); // ADVISORY + if (!std.mem.eql(u8, words.next() orelse continue, "WRITE")) continue; + const who = std.fmt.parseInt(u32, words.next() orelse continue, 10) catch continue; + if (who != pid) continue; + const id = words.next() orelse continue; + const colon = std.mem.lastIndexOfScalar(u8, id, ':') orelse continue; + out[n] = std.fmt.parseInt(u64, id[colon + 1 ..], 10) catch continue; + n += 1; + } + return out[0..n]; +} + +/// A codex thread's rollout: `sessions/YYYY/MM/DD/rollout-<when>-<id>.jsonl`. +/// The id is a UUIDv7, whose first 48 bits are the Unix milliseconds it +/// was minted at, so the day directory is known to within the timezone +/// the rollout was named in: the UTC day, or the one either side of it. +pub fn rolloutOf(s: Sources, base: []const u8, id: []const u8, out: []u8) ?[]const u8 { + if (id.len != uuid_len or id[8] != '-') return null; + var hex: [12]u8 = undefined; + @memcpy(hex[0..8], id[0..8]); + @memcpy(hex[8..12], id[9..13]); + const ms = std.fmt.parseInt(u64, &hex, 16) catch return null; + const day: i64 = @intCast(ms / std.time.ms_per_day); + var want_buf: [uuid_len + 8]u8 = undefined; + const want = std.fmt.bufPrint(&want_buf, "-{s}.jsonl", .{id}) catch return null; + for ([_]i64{ day, day - 1, day + 1 }) |d| { + const date = civil(d); + var dir_buf: [path_capacity]u8 = undefined; + const dir_path = std.fmt.bufPrint(&dir_buf, "{s}/sessions/{d:0>4}/{d:0>2}/{d:0>2}", .{ base, date.y, date.m, date.d }) catch continue; + const dir = Io.Dir.openDirAbsolute(s.io, dir_path, .{ .iterate = true }) catch continue; + defer Io.Dir.close(dir, s.io); + var rb: [Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined; + var it = Io.Dir.Reader.init(dir, &rb); + while (true) { + const entry = (it.next(s.io) catch break) orelse break; + if (entry.kind != .file) continue; + if (!std.mem.startsWith(u8, entry.name, "rollout-") or !std.mem.endsWith(u8, entry.name, want)) continue; + return join(out, &.{ dir_path, entry.name }); + } + } + return null; +} + +/// Whether a rollout is a subagent's thread rather than a session a +/// person started: its first record says so. +fn subagentRollout(io: Io, path: []const u8) bool { + var head: [probe_capacity]u8 = undefined; + const text = readSmall(io, path, &head) orelse return true; + return std.mem.indexOf(u8, text, "\"thread_source\":\"subagent\"") != null; +} + +/// The civil date of a day count from the Unix epoch (Howard Hinnant's +/// `civil_from_days`). +fn civil(days: i64) struct { y: u32, m: u32, d: u32 } { + const z = days + 719468; + const era = @divFloor(z, 146097); + const doe = z - era * 146097; + const yoe = @divFloor(doe - @divFloor(doe, 1460) + @divFloor(doe, 36524) - @divFloor(doe, 146096), 365); + const doy = doe - (365 * yoe + @divFloor(yoe, 4) - @divFloor(yoe, 100)); + const mp = @divFloor(5 * doy + 2, 153); + const d = doy - @divFloor(153 * mp + 2, 5) + 1; + const m = if (mp < 10) mp + 3 else mp - 9; + const y = yoe + era * 400 + @intFromBool(m <= 2); + return .{ .y = @intCast(@max(y, 0)), .m = @intCast(m), .d = @intCast(d) }; +} + /// Is this path, relative to the harness's session root, one of its /// session transcripts rather than a subagent's or something else? fn sessionFile(kind: Kind, rel: []const u8) bool { @@ -724,11 +898,14 @@ fn titleKey(k: Kind) []const u8 { /// A field from the first few records of a transcript, where the /// harnesses put the session's title and the model it runs. Only the -/// head is read: these records are written when the session opens. -pub fn headField(io: Io, kind: Kind, transcript: []const u8, which: enum { title, model }, out: []u8) ?[]const u8 { - if (transcript.len == 0) return null; +/// head is read: these records are written when the session opens. The +/// caller opens the transcript, the way every transcript is opened: from +/// its pinned root, a component at a time. +pub const Head = enum { title, model }; + +pub fn headField(io: Io, kind: Kind, transcript: Io.File, which: Head, out: []u8) ?[]const u8 { var head: [probe_capacity]u8 = undefined; - const text = readSmall(io, transcript, &head) orelse return null; + const text = readHead(io, transcript, &head) orelse return null; const key = switch (which) { .title => titleKey(kind), .model => "model", @@ -994,6 +1171,14 @@ test "active: a zmx session name is checked before it is ever used" { try testing.expect(!legalZmxName("x" ** 65)); } +test "active: a codex thread's rollout is found by the day its id was minted" { + try testing.expectEqual(@as(u32, 1970), civil(0).y); + const d = civil(20720); // 2026-09-24 + try testing.expect(d.y == 2026 and d.m == 9 and d.d == 24); + const leap = civil(19782); // 2024-02-29 + try testing.expect(leap.y == 2024 and leap.m == 2 and leap.d == 29); +} + test "active: an environment block answers by key, not by prefix" { const env = "PATH=/usr/bin\x00ZMX_SESSION=harness\x00HOME=/home/goblin\x00"; try testing.expectEqualStrings("harness", envValue(env, "ZMX_SESSION").?); |
