summaryrefslogtreecommitdiff
path: root/9agents/src/active.zig
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-25 16:40:49 -0300
committerGabriel Schneider <[email protected]>2026-09-25 17:52:33 -0300
commit1c1b192d4ef59199a7196229d56901b1f6512678 (patch)
tree9f6e4c1efcf877d1231aa58975ccd923ad4f2f2e /9agents/src/active.zig
parent73602127d15d10a1932b6fe916bd18a608054980 (diff)
downloadcloud9-main.tar.gz
cloud9-main.zip
9agents: /active/<h>/<pid>/chat, the conversation as one file per messageHEADmain
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.zig207
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").?);