From 1c1b192d4ef59199a7196229d56901b1f6512678 Mon Sep 17 00:00:00 2001 From: Gabriel Schneider Date: Fri, 25 Sep 2026 16:40:49 -0300 Subject: 9agents: /active///chat, the conversation as one file per message Each message of a claude or codex agent's transcript is a file named -: 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//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) --- 9agents/src/active.zig | 207 +++++- 9agents/src/chat.zig | 1626 ++++++++++++++++++++++++++++++++++++++++++++++++ 9agents/src/main.zig | 43 +- 9agents/src/tree.zig | 566 +++++++++++++++-- 4 files changed, 2384 insertions(+), 58 deletions(-) create mode 100644 9agents/src/chat.zig (limited to '9agents/src') 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//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/.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 `/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--.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").?); diff --git a/9agents/src/chat.zig b/9agents/src/chat.zig new file mode 100644 index 0000000..d6b4773 --- /dev/null +++ b/9agents/src/chat.zig @@ -0,0 +1,1626 @@ +//! The chat behind `/active///chat/`: a harness transcript +//! read as numbered messages, one file each. +//! +//! chat/00000000-user 00000001-assistant 00000002-bash 00000003-bash-result ... +//! +//! The name is `-`. The id is the message's position in the +//! transcript, counted from 0 and written as 8 digits so the names sort +//! in order; transcripts are append-only, so an id never moves and a new +//! message takes the next one. The kind says who spoke, normalized across +//! harnesses (`Kind`); a tool call is named after its tool and its result +//! after the call (`Chat.label`). The file holds the message rendered as +//! text: JSON strings unescaped, and each part on its own line — a tool +//! call is its name, then its input. +//! +//! A transcript is JSONL, and finding message N means knowing where every +//! message before it starts, so each transcript gets an index: per message, +//! the spans of the file that render it (`Span`). The index is not a copy +//! of anything. It holds offsets, never bytes, and every read still reads +//! the file. It is keyed by the file's (dev, ino) and extended from where it +//! stopped when the file grows; a file that shrinks or is replaced starts +//! over. Only complete lines are indexed, so a record the harness is +//! half-way through writing waits for its newline. +//! +//! The caps are loud, as everywhere else in 9agents: a transcript with more +//! messages than an index holds is marked `full`, and the listing refuses +//! rather than end early. +const std = @import("std"); +const Io = std.Io; +const linux = std.os.linux; + +// ---- comptime bounds --------------------------------------------------------- + +/// Messages one index holds. The longest codex rollout on this machine +/// (49k records) is 32k messages; Claude Code sessions stay near 3k. The +/// pool is mapped `MAP_NORESERVE` (`Pool.create`), so a generous cap costs +/// address space, not memory: a page is only backed once an index uses it. +pub const max_msgs: usize = 1 << 18; +/// Spans one index holds: most messages are one span, a tool call two. +pub const max_spans: usize = 2 * max_msgs; +/// Digits in a message id: every id has the same width, so a plain sort +/// of the names is the order the messages were said in. +pub const id_digits: usize = 8; +comptime { + std.debug.assert(max_msgs <= 100_000_000); // every id fits the width +} +/// Transcripts indexed at once, reused least-recently-used first. +pub const chat_slots: usize = 8; +/// The read window over a transcript. +pub const window: usize = 64 * 1024; +/// The window always shows at least this much past the position asked for, +/// unless the file ends first: enough for any key compared and for the +/// longest JSON escape (a surrogate pair, 12 bytes). +const lookahead: usize = 64; + +/// The transcript formats this module reads. +pub const Format = enum(u8) { claude, codex }; + +/// Who spoke. The name of a message file ends in `text()`. +pub const Kind = enum(u8) { + /// A person's prompt. + user, + /// The model's reply text. + assistant, + /// The model's reasoning, where the harness writes it in the clear. + thinking, + /// A tool the model called: its name, then its input. Named after + /// the tool when it has a usable name; `tool-call` when not. + tool_call, + /// What the tool answered. Named `-result` when its call is + /// known; `tool-result` when not. + tool_result, + /// The harness speaking: developer instructions, slash-command + /// output, injected context, compaction summaries. + system, + /// Another agent of a team (codex): its name, then what it said. + agent, + + pub fn text(k: Kind) []const u8 { + return switch (k) { + .tool_call => "tool-call", + .tool_result => "tool-result", + else => @tagName(k), + }; + } + + pub fn fromText(s: []const u8) ?Kind { + for (std.enums.values(Kind)) |k| if (std.mem.eql(u8, k.text(), s)) return k; + return null; + } +}; + +/// How a span's bytes become text. +pub const Enc = enum(u2) { + /// The bytes as they are (a JSON value served verbatim). + raw, + /// The body of a JSON string, unescaped. + str, + /// A fixed text from `literals`; `off` is its index. + lit, +}; + +/// Texts for what a transcript holds that is not text. +const literals = [_][]const u8{ "[image]", "" }; +const lit_image: u64 = 0; +const lit_empty: u64 = 1; + +/// One stretch of the transcript that renders part of a message. +/// Packed to 16 bytes: an index is mostly these. +pub const Span = packed struct(u128) { + /// File offset of the first byte; for `lit`, an index into `literals`. + off: u48, + /// Bytes in the file. + len: u32, + /// Bytes it renders to, not counting `nl`. + out: u32, + enc: Enc, + /// A newline follows, because the text does not end in one. + nl: bool, + _: u13 = 0, + + fn size(s: Span) u64 { + return @as(u64, s.out) + @intFromBool(s.nl); + } +}; + +pub const Msg = struct { + /// Index of the first span in `Chat.spans`. + first: u32, + count: u16, + kind: Kind, + /// Unix seconds of the record, 0 when it carries no timestamp. + time: u32, + /// Rendered size: the sum of the spans'. + size: u32, + /// For a tool call or result, the tool, as an index into the chat's + /// tool names plus one; 0 when unknown. + tool: u8 = 0, +}; + +/// Distinct tool names one chat names its calls by. Past this, a call is +/// named `tool-call`, as one whose tool has no usable name is. +pub const max_tools: usize = 255; +/// Longest tool name kept for a file name. +pub const tool_name_cap: usize = 48; +/// Calls remembered for matching their results: results follow their +/// calls closely, a few dozen at most when calls run in parallel. +const call_ring: usize = 256; + +const Call = struct { + /// Hash of the call's id (never 0 when set). + id: u64 = 0, + tool: u8 = 0, +}; + +/// The index of one transcript. +pub const Chat = struct { + used: bool = false, + format: Format = .claude, + dev: u64 = 0, + ino: u64 = 0, + /// The file's size when last synced; reads never go past it. + size: u64 = 0, + /// Offset of the first line not yet indexed. Always a line start. + scanned: u64 = 0, + /// Least-recently-used ordering. + stamp: u64 = 0, + /// A cap was reached. Messages past `n` exist but are not indexed. + full: bool = false, + /// When the file was created, where the filesystem says: a new file + /// that reuses a deleted one's inode is born later. + btime: i128 = 0, + /// Hashes of the first and last bytes indexed (up to `probe_len` each, + /// the last ending at `scanned`). A transcript rewritten in place to + /// the same size or longer keeps its inode and grows past `scanned`; + /// only its bytes can say it is no longer the file that was indexed. + head: u64 = 0, + tail: u64 = 0, + /// Bumped by every reset, so a read cursor cannot outlive its index. + epoch: u32 = 0, + /// Offset up to which the unindexed tail is known to hold no newline. + unterminated: u64 = 0, + n: u32 = 0, + nspans: u32 = 0, + ntools: u8 = 0, + call_next: u16 = 0, + tool_lens: [max_tools]u8, + tool_names: [max_tools][tool_name_cap]u8, + calls: [call_ring]Call, + msgs: [max_msgs]Msg, + spans: [max_spans]Span, + + fn reset(c: *Chat, format: Format, dev: u64, ino: u64) void { + c.used = true; + c.epoch +%= 1; + c.btime = 0; + c.head = 0; + c.tail = 0; + c.format = format; + c.dev = dev; + c.ino = ino; + c.size = 0; + c.scanned = 0; + c.full = false; + c.unterminated = 0; + c.n = 0; + c.nspans = 0; + c.ntools = 0; + c.call_next = 0; + c.calls = @splat(.{}); + } + + pub fn msg(c: *const Chat, i: u32) ?Msg { + return if (i < c.n) c.msgs[i] else null; + } + + /// What a message's name says after its id: its kind — or, for a + /// tool call, the tool, and for its result the tool and `-result`. + fn label(c: *const Chat, m: Msg) struct { []const u8, []const u8 } { + if (m.tool == 0) return .{ m.kind.text(), "" }; + const t = c.tool_names[m.tool - 1][0..c.tool_lens[m.tool - 1]]; + return .{ t, if (m.kind == .tool_result) "-result" else "" }; + } + + /// `-