summaryrefslogtreecommitdiff
path: root/9agents/src
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
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')
-rw-r--r--9agents/src/active.zig207
-rw-r--r--9agents/src/chat.zig1626
-rw-r--r--9agents/src/main.zig43
-rw-r--r--9agents/src/tree.zig566
4 files changed, 2384 insertions, 58 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").?);
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/<harness>/<pid>/chat/`: a harness transcript
+//! read as numbered messages, one file each.
+//!
+//! chat/00000000-user 00000001-assistant 00000002-bash 00000003-bash-result ...
+//!
+//! The name is `<id>-<kind>`. 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 `<tool>-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 "" };
+ }
+
+ /// `<id>-<label>`, the id as `id_digits` digits, written into `buf`.
+ pub fn name(c: *const Chat, i: u32, buf: []u8) ?[]const u8 {
+ const m = c.msg(i) orelse return null;
+ const l = c.label(m);
+ return std.fmt.bufPrint(buf, "{d:0>8}-{s}{s}", .{ i, l[0], l[1] }) catch null;
+ }
+
+ /// The message a file name names: `<id>-<label>` with the id in its
+ /// one spelling (exactly `id_digits` digits) and the label it really
+ /// has, so every message answers to exactly one name.
+ pub fn lookup(c: *const Chat, file_name: []const u8) ?u32 {
+ if (file_name.len <= id_digits or file_name[id_digits] != '-') return null;
+ const digits = file_name[0..id_digits];
+ for (digits) |d| if (!std.ascii.isDigit(d)) return null;
+ const i = std.fmt.parseInt(u32, digits, 10) catch return null;
+ const m = c.msg(i) orelse return null;
+ const rest = file_name[id_digits + 1 ..];
+ const l = c.label(m);
+ if (rest.len != l[0].len + l[1].len) return null;
+ if (!std.mem.startsWith(u8, rest, l[0]) or !std.mem.endsWith(u8, rest, l[1])) return null;
+ return i;
+ }
+
+ /// The index (plus one) a tool name is kept under, adding it if new;
+ /// 0 when it has no usable name. Lowercased, and anything but letters,
+ /// digits, `_`, `.` and `-` becomes `-`, so it is safe in a file name.
+ /// A name that would read as another kind, or as a result, is not used.
+ fn internTool(c: *Chat, raw: []const u8) u8 {
+ var buf: [tool_name_cap]u8 = undefined;
+ const n = @min(raw.len, buf.len);
+ if (n == 0) return 0;
+ for (raw[0..n], buf[0..n]) |ch, *o| o.* = switch (ch) {
+ 'A'...'Z' => ch + ('a' - 'A'),
+ 'a'...'z', '0'...'9', '_', '.', '-' => ch,
+ else => '-',
+ };
+ const t = buf[0..n];
+ if (Kind.fromText(t) != null or std.mem.endsWith(u8, t, "-result")) return 0;
+ for (0..c.ntools) |i| {
+ if (std.mem.eql(u8, c.tool_names[i][0..c.tool_lens[i]], t)) return @intCast(i + 1);
+ }
+ if (c.ntools == max_tools) return 0;
+ @memcpy(c.tool_names[c.ntools][0..n], t);
+ c.tool_lens[c.ntools] = @intCast(n);
+ c.ntools += 1;
+ return c.ntools;
+ }
+
+ /// Remembers which tool a call went to, for its result to find.
+ fn rememberCall(c: *Chat, id: u64, tool: u8) void {
+ if (id == 0 or tool == 0) return;
+ c.calls[c.call_next] = .{ .id = id, .tool = tool };
+ c.call_next = @intCast((c.call_next + 1) % call_ring);
+ }
+
+ /// The tool a result answers, from its call's id; 0 when not known.
+ fn toolOfCall(c: *const Chat, id: u64) u8 {
+ if (id == 0) return 0;
+ for (c.calls) |k| if (k.id == id) return k.tool;
+ return 0;
+ }
+
+ // -- building --
+
+ fn open(c: *Chat, kind: Kind, time: u32) bool {
+ if (c.n == max_msgs) return false;
+ c.msgs[c.n] = .{ .first = c.nspans, .count = 0, .kind = kind, .time = time, .size = 0 };
+ return true;
+ }
+
+ fn add(c: *Chat, s: Span) bool {
+ const m = &c.msgs[c.n];
+ if (c.nspans == max_spans or m.count == std.math.maxInt(u16)) return false;
+ const size = std.math.add(u32, m.size, @intCast(s.size())) catch return false;
+ c.spans[c.nspans] = s;
+ c.nspans += 1;
+ m.count += 1;
+ m.size = size;
+ return true;
+ }
+
+ /// Keeps the message being built if it got a span, drops it if not.
+ fn close(c: *Chat) void {
+ const m = &c.msgs[c.n];
+ if (m.count > 0) c.n += 1 else c.nspans = m.first;
+ }
+};
+
+/// The indexes, and the window they read through. One instance serves the
+/// daemon; the caller serializes access (the harness mutex).
+pub const Pool = struct {
+ chats: [chat_slots]Chat,
+ clock: u64,
+ /// Where the last read of a long string stopped (see `render`).
+ cursor: Cursor,
+ win: [window]u8,
+
+ /// A pool in its own anonymous mapping, `MAP_NORESERVE`: the ~100 MB
+ /// of tables are address space until an index writes them, and no
+ /// part of it is ever written just to initialize it.
+ pub fn create() error{OutOfMemory}!*Pool {
+ const rc = linux.mmap(null, @sizeOf(Pool), .{ .READ = true, .WRITE = true }, .{
+ .TYPE = .PRIVATE,
+ .ANONYMOUS = true,
+ .NORESERVE = true,
+ }, -1, 0);
+ if (linux.errno(rc) != .SUCCESS) return error.OutOfMemory;
+ // Small pages: with transparent huge pages on, the first write to
+ // each chat's header would back a whole 2 MB page of its tables.
+ _ = linux.madvise(@ptrFromInt(rc), @sizeOf(Pool), linux.MADV.NOHUGEPAGE);
+ const p: *Pool = @ptrFromInt(rc);
+ p.init();
+ return p;
+ }
+
+ pub fn destroy(p: *Pool) void {
+ _ = linux.munmap(@ptrCast(p), @sizeOf(Pool));
+ }
+
+ /// Resets the headers only: the tables are written as they fill.
+ pub fn init(p: *Pool) void {
+ p.clock = 0;
+ p.cursor = .{};
+ for (&p.chats) |*c| {
+ c.used = false;
+ c.n = 0;
+ c.nspans = 0;
+ c.epoch = 0;
+ }
+ }
+
+ /// The index of the transcript open as `file`, brought up to date with
+ /// it. Null when the file is not a regular file.
+ pub fn sync(p: *Pool, io: Io, file: Io.File, format: Format) ?*Chat {
+ var sx: linux.Statx = undefined;
+ const want: linux.STATX = .{ .TYPE = true, .INO = true, .SIZE = true, .BTIME = true };
+ const rc = linux.statx(file.handle, "", linux.AT.EMPTY_PATH, want, &sx);
+ if (linux.errno(rc) != .SUCCESS) return null;
+ if (sx.mode & linux.S.IFMT != linux.S.IFREG) return null;
+ const dev = (@as(u64, sx.dev_major) << 32) | sx.dev_minor;
+ const btime: i128 = if (sx.mask.BTIME) @as(i128, sx.btime.sec) * std.time.ns_per_s + sx.btime.nsec else 0;
+
+ p.clock += 1;
+ const c = p.find(format, dev, sx.ino);
+ // A transcript only grows, and keeps the bytes it had. Anything
+ // else — shorter, born again under a reused inode, or rewritten in
+ // place — is another file, and the offsets mean nothing in it.
+ if (sx.size < c.scanned or c.btime != btime or !p.unchanged(io, file, c)) {
+ c.reset(format, dev, sx.ino);
+ }
+ c.btime = btime;
+ c.stamp = p.clock;
+ c.size = sx.size;
+ if (c.scanned < sx.size and !c.full) {
+ p.catchUp(io, file, c);
+ const pr = p.probe(io, file, c.scanned);
+ c.head = pr.head;
+ c.tail = pr.tail;
+ }
+ return c;
+ }
+
+ const probe_len: usize = 64;
+
+ /// Hashes of the first and last `probe_len` bytes before `end`.
+ fn probe(p: *Pool, io: Io, file: Io.File, end: u64) struct { head: u64, tail: u64 } {
+ _ = p;
+ if (end == 0) return .{ .head = 0, .tail = 0 };
+ var buf: [probe_len]u8 = undefined;
+ const k: usize = @intCast(@min(end, probe_len));
+ const hn = file.readPositionalAll(io, buf[0..k], 0) catch 0;
+ const head = std.hash.Wyhash.hash(0, buf[0..hn]);
+ const tn = file.readPositionalAll(io, buf[0..k], end - k) catch 0;
+ const tail = std.hash.Wyhash.hash(1, buf[0..tn]);
+ return .{ .head = head, .tail = tail };
+ }
+
+ fn unchanged(p: *Pool, io: Io, file: Io.File, c: *const Chat) bool {
+ if (c.scanned == 0) return true;
+ const pr = p.probe(io, file, c.scanned);
+ return pr.head == c.head and pr.tail == c.tail;
+ }
+
+ fn find(p: *Pool, format: Format, dev: u64, ino: u64) *Chat {
+ var victim = &p.chats[0];
+ for (&p.chats) |*c| {
+ if (c.used and c.format == format and c.dev == dev and c.ino == ino) return c;
+ if (!c.used) {
+ if (victim.used) victim = c;
+ } else if (victim.used and c.stamp < victim.stamp) victim = c;
+ }
+ victim.reset(format, dev, ino);
+ return victim;
+ }
+
+ fn catchUp(p: *Pool, io: Io, file: Io.File, c: *Chat) void {
+ var src: Src = .{ .io = io, .file = file, .end = c.size, .buf = &p.win };
+ var pos = c.scanned;
+ while (pos < c.size) {
+ // A line without its newline is still being written.
+ // A long unterminated tail is searched once, not per request.
+ const eol = src.find(@max(pos, c.unterminated), '\n') orelse {
+ c.unterminated = c.size;
+ break;
+ };
+ const mark_n = c.n;
+ const mark_s = c.nspans;
+ var ps: P = .{ .src = &src, .lim = eol };
+ const ok = switch (c.format) {
+ .claude => claudeLine(c, &ps, pos),
+ .codex => codexLine(c, &ps, pos),
+ };
+ if (!ok) {
+ // A cap: this line is left whole for nobody, and the
+ // listing says so instead of ending early.
+ c.n = mark_n;
+ c.nspans = mark_s;
+ c.full = true;
+ break;
+ }
+ pos = eol + 1;
+ }
+ c.scanned = pos;
+ }
+
+ /// Renders message `i` from byte `off` into `out`, reading `file`.
+ /// Returns the bytes written; fewer than asked only at the end.
+ ///
+ /// A string's output offset does not map onto its input offset, so a
+ /// read deep into a long one would decode everything before it, and a
+ /// client reading a message front to back would pay for it again on
+ /// every read. So the pool remembers where the last read stopped — the
+ /// input offset of the last escape or run it started and the output
+ /// offset that produced — and a read continuing the same message
+ /// starts from there.
+ pub fn render(p: *Pool, io: Io, file: Io.File, c: *const Chat, i: u32, off: u64, out: []u8) usize {
+ const m = c.msg(i) orelse return 0;
+ var src: Src = .{ .io = io, .file = file, .end = c.size, .buf = &p.win };
+ var at: u64 = 0;
+ var n: usize = 0;
+ for (c.spans[m.first..][0..m.count], m.first..) |s, k| {
+ if (n == out.len) break;
+ const size = s.size();
+ defer at += size;
+ if (at + size <= off) continue;
+ var sink: Sink = .{ .skip = off -| at, .out = out[n..] };
+ switch (s.enc) {
+ .lit => sink.put(literals[@intCast(s.off)]),
+ .raw => {
+ // One byte in, one byte out: start at the offset asked.
+ const from = @min(sink.skip, s.len);
+ sink.skip -= from;
+ copy(&src, @as(u64, s.off) + from, s.len - from, &sink);
+ },
+ .str => {
+ const cur = p.cursor;
+ var from: u64 = s.off;
+ var base: u64 = 0; // output offset `from` produces
+ if (cur.chat == c and cur.epoch == c.epoch and cur.span == k and cur.out <= sink.skip) {
+ from = cur.in;
+ base = cur.out;
+ sink.skip -= cur.out;
+ }
+ sink.mark_in = from;
+ decodeSpan(&src, from, @as(u64, s.off) + s.len - from, &sink);
+ if (sink.done()) p.cursor = .{
+ .chat = c,
+ .epoch = c.epoch,
+ .span = @intCast(k),
+ .in = sink.mark_in,
+ .out = base + sink.mark_out,
+ };
+ },
+ }
+ if (s.nl) sink.put("\n");
+ n += sink.n;
+ }
+ return n;
+ }
+};
+
+// ---- reading the file through a window ----------------------------------------
+
+const Src = struct {
+ io: Io,
+ file: Io.File,
+ /// Nothing at or past this offset is read.
+ end: u64,
+ buf: []u8,
+ base: u64 = 0,
+ len: usize = 0,
+
+ /// The bytes from `pos` on, at least `lookahead` of them unless the
+ /// file ends first. Empty at the end, or when the file was cut short
+ /// under the read.
+ fn at(s: *Src, pos: u64) []const u8 {
+ if (pos >= s.end) return &.{};
+ const have_end = s.base + s.len;
+ const inside = pos >= s.base and pos < have_end;
+ if (!inside or (have_end - pos < lookahead and have_end < s.end)) s.load(pos);
+ if (pos < s.base or pos >= s.base + s.len) return &.{};
+ return s.buf[@intCast(pos - s.base)..s.len];
+ }
+
+ fn load(s: *Src, pos: u64) void {
+ const want: usize = @intCast(@min(s.buf.len, s.end - pos));
+ s.base = pos;
+ s.len = s.file.readPositionalAll(s.io, s.buf[0..want], pos) catch 0;
+ }
+
+ /// The offset of the next `byte` at or after `pos`.
+ fn find(s: *Src, pos: u64, byte: u8) ?u64 {
+ var q = pos;
+ while (true) {
+ const b = s.at(q);
+ if (b.len == 0) return null;
+ if (std.mem.indexOfScalar(u8, b, byte)) |i| return q + i;
+ q += b.len;
+ }
+ }
+};
+
+/// Where the last read of a string stopped: decoding span `span` from
+/// file offset `in` produces its output from offset `out` on.
+const Cursor = struct {
+ chat: ?*const Chat = null,
+ epoch: u32 = 0,
+ span: u32 = 0,
+ in: u64 = 0,
+ out: u64 = 0,
+};
+
+/// Where rendered bytes go: `skip` of them are dropped, then `out` fills.
+/// `total` and `last` see every byte, which is how a span is measured.
+const Sink = struct {
+ skip: u64 = 0,
+ out: []u8 = &.{},
+ n: usize = 0,
+ total: u64 = 0,
+ last: u8 = '\n',
+ /// Measuring: nothing is written, everything is counted.
+ measuring: bool = false,
+ /// The input offset of the last run or escape begun before `out`
+ /// filled, and the `total` it began at: where a later read resumes.
+ mark_in: u64 = 0,
+ mark_out: u64 = 0,
+
+ fn mark(s: *Sink, pos: u64) void {
+ if (s.done()) return;
+ s.mark_in = pos;
+ s.mark_out = s.total;
+ }
+
+ fn put(s: *Sink, bytes: []const u8) void {
+ if (bytes.len == 0) return;
+ s.total += bytes.len;
+ s.last = bytes[bytes.len - 1];
+ var b = bytes;
+ if (s.skip > 0) {
+ const d: usize = @intCast(@min(s.skip, b.len));
+ s.skip -= d;
+ b = b[d..];
+ }
+ const k = @min(s.out.len - s.n, b.len);
+ @memcpy(s.out[s.n..][0..k], b[0..k]);
+ s.n += k;
+ }
+
+ fn done(s: *const Sink) bool {
+ return !s.measuring and s.n == s.out.len;
+ }
+};
+
+fn copy(src: *Src, off: u64, len: u64, sink: *Sink) void {
+ var pos = off;
+ const end = off + len;
+ while (pos < end and !sink.done()) {
+ const b = src.at(pos);
+ if (b.len == 0) return;
+ const chunk = b[0..@intCast(@min(b.len, end - pos))];
+ sink.put(chunk);
+ pos += chunk.len;
+ }
+}
+
+// ---- JSON strings -----------------------------------------------------------
+
+/// Unescapes the body of a JSON string held in `[off, off+len)`.
+fn decodeSpan(src: *Src, off: u64, len: u64, sink: *Sink) void {
+ var pos = off;
+ const end = off + len;
+ while (pos < end and !sink.done()) {
+ const b = src.at(pos);
+ if (b.len == 0) return;
+ const chunk = b[0..@intCast(@min(b.len, end - pos))];
+ const k = decode(chunk, pos, pos + chunk.len == end, sink);
+ if (k == 0) return; // cannot happen with `lookahead`; never spin
+ pos += k;
+ }
+}
+
+/// Unescapes as much of `chunk` (at file offset `at`) as it can,
+/// returning the bytes consumed. An escape cut by the chunk's end is left
+/// for the next call unless `final` says there is no more.
+fn decode(chunk: []const u8, at: u64, final: bool, sink: *Sink) usize {
+ var i: usize = 0;
+ while (i < chunk.len) {
+ sink.mark(at + i);
+ const j = std.mem.indexOfScalarPos(u8, chunk, i, '\\') orelse {
+ sink.put(chunk[i..]);
+ return chunk.len;
+ };
+ sink.put(chunk[i..j]);
+ i = j;
+ if (!final and chunk.len - i < 12) return i;
+ sink.mark(at + i);
+ if (i + 1 == chunk.len) {
+ sink.put("\\"); // a lone backslash at the very end
+ return chunk.len;
+ }
+ const e = chunk[i + 1];
+ const simple: ?u8 = switch (e) {
+ 'n' => '\n',
+ 't' => '\t',
+ 'r' => '\r',
+ 'b' => 0x08,
+ 'f' => 0x0c,
+ '"', '\\', '/' => e,
+ else => null,
+ };
+ if (simple) |c| {
+ sink.put(&.{c});
+ i += 2;
+ continue;
+ }
+ if (e != 'u') {
+ sink.put(chunk[i .. i + 2]); // not an escape JSON has: as written
+ i += 2;
+ continue;
+ }
+ const hi = hex4(chunk[i + 2 ..]) orelse {
+ sink.put(chunk[i .. i + 2]);
+ i += 2;
+ continue;
+ };
+ i += 6;
+ var cp: u21 = hi;
+ if (hi >= 0xd800 and hi < 0xdc00) {
+ // A high surrogate wants its low half next, as `\uDCxx`.
+ const rest = chunk[i..];
+ const lo: ?u16 = if (rest.len >= 6 and rest[0] == '\\' and rest[1] == 'u') hex4(rest[2..]) else null;
+ if (lo) |l| if (l >= 0xdc00 and l < 0xe000) {
+ cp = 0x10000 + ((@as(u21, hi) - 0xd800) << 10) + (l - 0xdc00);
+ i += 6;
+ };
+ if (cp == hi) cp = 0xfffd;
+ } else if (hi >= 0xdc00 and hi < 0xe000) cp = 0xfffd;
+ var utf: [4]u8 = undefined;
+ const n = std.unicode.utf8Encode(cp, &utf) catch 0;
+ sink.put(utf[0..n]);
+ }
+ return i;
+}
+
+/// Four hex digits exactly: `parseInt` would also take a sign or `_`.
+fn hex4(b: []const u8) ?u16 {
+ if (b.len < 4) return null;
+ var v: u16 = 0;
+ for (b[0..4]) |c| v = (v << 4) | (std.fmt.charToDigit(c, 16) catch return null);
+ return v;
+}
+
+// ---- a JSON walker over the window ------------------------------------------
+
+/// Enough of a JSON parser to find values by key and know where each one
+/// starts and ends. Positions are file offsets; `lim` is the line's end.
+/// It checks nesting, not grammar: a line it cannot walk yields nothing.
+const P = struct {
+ src: *Src,
+ lim: u64,
+
+ fn peek(p: *P, pos: u64) ?u8 {
+ if (pos >= p.lim) return null;
+ const b = p.src.at(pos);
+ return if (b.len == 0) null else b[0];
+ }
+
+ /// The window from `pos`, cut at the line's end.
+ fn bytes(p: *P, pos: u64) []const u8 {
+ if (pos >= p.lim) return &.{};
+ const b = p.src.at(pos);
+ return b[0..@intCast(@min(b.len, p.lim - pos))];
+ }
+
+ fn ws(p: *P, pos: u64) u64 {
+ var q = pos;
+ while (p.peek(q)) |c| switch (c) {
+ ' ', '\t', '\r', '\n' => q += 1,
+ else => break,
+ };
+ return q;
+ }
+
+ /// `pos` at an opening quote; the offset just past the closing one.
+ /// Most of a transcript's bytes are inside strings (tool output,
+ /// encrypted reasoning), so this is the hot loop: it jumps from quote
+ /// to quote, and a quote closes the string when the run of
+ /// backslashes before it is even.
+ fn strEnd(p: *P, pos: u64) ?u64 {
+ var q = pos + 1;
+ // Backslashes ending the bytes before `q`, carried across windows.
+ var run: usize = 0;
+ while (true) {
+ const b = p.bytes(q);
+ if (b.len == 0) return null;
+ var from: usize = 0;
+ while (std.mem.indexOfScalarPos(u8, b, from, '"')) |i| {
+ var k = i;
+ while (k > from and b[k - 1] == '\\') k -= 1;
+ const slashes = (i - k) + (if (k == 0) run else 0);
+ if (slashes % 2 == 0) return q + i + 1;
+ from = i + 1;
+ run = 0;
+ }
+ var t = b.len;
+ while (t > from and b[t - 1] == '\\') t -= 1;
+ run = (b.len - t) + (if (t == 0) run else 0);
+ q += b.len;
+ }
+ }
+
+ /// The offset just past the value at `pos`.
+ fn skip(p: *P, pos: u64) ?u64 {
+ const c = p.peek(pos) orelse return null;
+ switch (c) {
+ '"' => return p.strEnd(pos),
+ '{', '[' => {
+ var depth: usize = 0;
+ var q = pos;
+ while (true) {
+ const b = p.bytes(q);
+ if (b.len == 0) return null;
+ const i = std.mem.indexOfAny(u8, b, "\"{}[]") orelse {
+ q += b.len;
+ continue;
+ };
+ q += i;
+ switch (b[i]) {
+ '"' => q = p.strEnd(q) orelse return null,
+ '{', '[' => {
+ depth += 1;
+ q += 1;
+ },
+ else => {
+ if (depth == 0) return null;
+ depth -= 1;
+ q += 1;
+ if (depth == 0) return q;
+ },
+ }
+ }
+ },
+ else => {
+ var q = pos;
+ while (p.peek(q)) |d| switch (d) {
+ ',', '}', ']', ' ', '\t', '\r', '\n' => break,
+ else => q += 1,
+ };
+ return if (q == pos) null else q;
+ },
+ }
+ }
+
+ fn isString(p: *P, pos: u64) bool {
+ return p.peek(pos) == '"';
+ }
+
+ /// A short string value without escapes, copied into `buf`: the
+ /// `type`s, roles and timestamps the walk steers by.
+ fn word(p: *P, pos: ?u64, buf: []u8) []const u8 {
+ const at = pos orelse return "";
+ if (!p.isString(at)) return "";
+ const end = p.strEnd(at) orelse return "";
+ const len = end - at - 2;
+ if (len > buf.len or len > lookahead) return "";
+ const b = p.bytes(at + 1);
+ if (b.len < len) return "";
+ const s = b[0..@intCast(len)];
+ if (std.mem.indexOfScalar(u8, s, '\\') != null) return "";
+ @memcpy(buf[0..s.len], s);
+ return buf[0..s.len];
+ }
+
+ fn isTrue(p: *P, pos: ?u64) bool {
+ const at = pos orelse return false;
+ const b = p.bytes(at);
+ if (!std.mem.startsWith(u8, b, "true")) return false;
+ return b.len == 4 or !std.ascii.isAlphanumeric(b[4]);
+ }
+
+ /// Elements of the array at `pos` (or members of an object), in order.
+ fn items(p: *P, pos: ?u64) Items {
+ const at = pos orelse return .{ .p = p, .pos = 0, .done = true };
+ const open = p.peek(at) orelse return .{ .p = p, .pos = 0, .done = true };
+ if (open != '[' and open != '{') return .{ .p = p, .pos = 0, .done = true };
+ return .{ .p = p, .pos = at + 1, .object = open == '{' };
+ }
+
+ /// The keys this module ever asks a record for.
+ const Key = enum {
+ type,
+ role,
+ text,
+ content,
+ name,
+ input,
+ thinking,
+ arguments,
+ output,
+ action,
+ summary,
+ author,
+ message,
+ payload,
+ timestamp,
+ isMeta,
+ isCompactSummary,
+ isSidechain,
+ id,
+ tool_use_id,
+ call_id,
+ attachment,
+ prompt,
+ origin,
+ kind,
+ };
+
+ const Fields = struct {
+ at: [std.enums.values(Key).len]?u64 = @splat(null),
+ bad: bool = false,
+
+ fn get(f: *const Fields, k: Key) ?u64 {
+ return f.at[@intFromEnum(k)];
+ }
+ };
+
+ /// Where each known key's value starts in the object at `pos`. The
+ /// first occurrence wins. `bad` when the object does not walk.
+ fn fields(p: *P, pos: ?u64) Fields {
+ var f: Fields = .{};
+ const at = pos orelse {
+ f.bad = true;
+ return f;
+ };
+ if (p.peek(at) != '{') {
+ f.bad = true;
+ return f;
+ }
+ var it = p.items(at);
+ while (it.next()) |m| {
+ const key = p.bytes(m.key);
+ if (key.len < m.key_len) continue;
+ inline for (comptime std.enums.values(Key)) |k| {
+ if (f.at[@intFromEnum(k)] == null and std.mem.eql(u8, key[0..m.key_len], @tagName(k))) {
+ f.at[@intFromEnum(k)] = m.val;
+ }
+ }
+ }
+ f.bad = it.bad;
+ return f;
+ }
+};
+
+const Item = struct {
+ /// First byte of the key inside its quotes (objects only).
+ key: u64 = 0,
+ key_len: usize = 0,
+ val: u64,
+};
+
+const Items = struct {
+ p: *P,
+ pos: u64,
+ object: bool = false,
+ first: bool = true,
+ done: bool = false,
+ bad: bool = false,
+
+ fn next(it: *Items) ?Item {
+ if (it.done) return null;
+ const p = it.p;
+ var q = p.ws(it.pos);
+ const c = p.peek(q) orelse return it.fail();
+ if (c == (if (it.object) @as(u8, '}') else ']')) {
+ it.done = true;
+ return null;
+ }
+ if (!it.first) {
+ if (c != ',') return it.fail();
+ q = p.ws(q + 1);
+ }
+ it.first = false;
+ var item: Item = .{ .val = q };
+ if (it.object) {
+ if (!p.isString(q)) return it.fail();
+ const kend = p.strEnd(q) orelse return it.fail();
+ item.key = q + 1;
+ item.key_len = @intCast(@min(kend - q - 2, lookahead));
+ q = p.ws(kend);
+ if (p.peek(q) != ':') return it.fail();
+ q = p.ws(q + 1);
+ item.val = q;
+ }
+ it.pos = p.skip(q) orelse return it.fail();
+ return item;
+ }
+
+ fn fail(it: *Items) ?Item {
+ it.done = true;
+ it.bad = true;
+ return null;
+ }
+};
+
+// ---- spans ------------------------------------------------------------------
+
+/// The body of the string at `pos`, measured.
+fn strSpan(p: *P, pos: u64) ?Span {
+ const end = p.strEnd(pos) orelse return null;
+ const len = end - pos - 2;
+ if (len >= std.math.maxInt(u32)) return null;
+ var sink: Sink = .{ .measuring = true };
+ decodeSpan(p.src, pos + 1, len, &sink);
+ if (sink.total >= std.math.maxInt(u32)) return null;
+ return .{
+ .off = @intCast(pos + 1),
+ .len = @intCast(len),
+ .out = @intCast(sink.total),
+ .enc = .str,
+ .nl = sink.total > 0 and sink.last != '\n',
+ };
+}
+
+/// The value at `pos` as it is written.
+fn rawSpan(p: *P, pos: u64) ?Span {
+ const end = p.skip(pos) orelse return null;
+ const len = end - pos;
+ if (len == 0 or len >= std.math.maxInt(u32)) return null;
+ const last = p.src.at(end - 1);
+ return .{
+ .off = @intCast(pos),
+ .len = @intCast(len),
+ .out = @intCast(len),
+ .enc = .raw,
+ .nl = last.len == 0 or last[0] != '\n',
+ };
+}
+
+/// A string as its text, anything else as its JSON.
+fn valueSpan(p: *P, pos: u64) ?Span {
+ return if (p.isString(pos)) strSpan(p, pos) else rawSpan(p, pos);
+}
+
+fn litSpan(which: u64) Span {
+ const t = literals[@intCast(which)];
+ return .{ .off = @intCast(which), .len = 0, .out = @intCast(t.len), .enc = .lit, .nl = t.len > 0 };
+}
+
+/// Adds a span if there is one; false only when a cap is reached.
+fn addSpan(c: *Chat, s: ?Span) bool {
+ const span = s orelse return true;
+ return c.add(span);
+}
+
+/// One message of `kind` holding the value at `pos` (a string as text).
+fn one(c: *Chat, p: *P, kind: Kind, time: u32, pos: ?u64) bool {
+ const at = pos orelse return true;
+ if (!c.open(kind, time)) return false;
+ if (!addSpan(c, valueSpan(p, at))) return false;
+ c.close();
+ return true;
+}
+
+/// A call's or result's id, hashed; 0 when there is none.
+fn callId(p: *P, pos: ?u64) u64 {
+ var buf: [lookahead]u8 = undefined;
+ const w = p.word(pos, &buf);
+ if (w.len == 0) return 0;
+ return std.hash.Wyhash.hash(0x63616c6c, w) | 1;
+}
+
+/// A tool call, named after its tool: the name on one line, then its
+/// input. `tool` is the name the file takes; `name` the value shown.
+fn call(c: *Chat, p: *P, time: u32, tool: []const u8, id: ?u64, name: ?u64, input: ?u64) bool {
+ if (!c.open(.tool_call, time)) return false;
+ const t = c.internTool(tool);
+ c.msgs[c.n].tool = t;
+ if (name) |n| if (!addSpan(c, valueSpan(p, n))) return false;
+ if (input) |i| if (!addSpan(c, valueSpan(p, i))) return false;
+ const kept = c.msgs[c.n].count > 0;
+ c.close();
+ if (kept) c.rememberCall(callId(p, id), t);
+ return true;
+}
+
+/// A tool's answer, named after the call it answers when that is known.
+fn result(c: *Chat, p: *P, time: u32, id: ?u64, content: ?u64) bool {
+ return blocks(c, p, .tool_result, time, c.toolOfCall(callId(p, id)), null, content);
+}
+
+/// Content that is a string, or an array of blocks whose text is shown,
+/// images named, and anything else served as its JSON — one message.
+fn blocks(c: *Chat, p: *P, kind: Kind, time: u32, tool: u8, head: ?u64, pos: ?u64) bool {
+ if (!c.open(kind, time)) return false;
+ c.msgs[c.n].tool = tool;
+ if (head) |h| if (!addSpan(c, valueSpan(p, h))) return false;
+ if (pos) |at| {
+ if (p.isString(at)) {
+ if (!addSpan(c, strSpan(p, at))) return false;
+ } else if (p.peek(at) == '[') {
+ var it = p.items(at);
+ while (it.next()) |el| {
+ const f = p.fields(el.val);
+ var tb: [32]u8 = undefined;
+ const ty = p.word(f.get(.type), &tb);
+ const s: ?Span = if (f.get(.text)) |t|
+ (if (p.isString(t)) strSpan(p, t) else null)
+ else if (std.mem.endsWith(u8, ty, "image"))
+ litSpan(lit_image)
+ else if (std.mem.eql(u8, ty, "encrypted_content"))
+ null // nothing a reader could use
+ else
+ rawSpan(p, el.val);
+ if (!addSpan(c, s)) return false;
+ }
+ } else if (!addSpan(c, rawSpan(p, at))) return false;
+ }
+ // A result that says nothing is still an answer to its call.
+ if (c.msgs[c.n].count == 0 and kind == .tool_result) {
+ if (!c.add(litSpan(lit_empty))) return false;
+ }
+ c.close();
+ return true;
+}
+
+// ---- Claude Code --------------------------------------------------------------
+
+/// One line of a Claude Code transcript. `user`, `assistant` and `system`
+/// records carry the conversation; the rest (attachments, mode changes,
+/// titles, snapshots) are the harness's bookkeeping and are not messages.
+/// Each content block is its own message, as Claude Code itself writes
+/// them. False only when a cap is reached.
+fn claudeLine(c: *Chat, p: *P, start: u64) bool {
+ const top = p.fields(p.ws(start));
+ if (top.bad) return true;
+ if (p.isTrue(top.get(.isSidechain))) return true; // a subagent's thread
+ var tb: [32]u8 = undefined;
+ var ts: [40]u8 = undefined;
+ const ty = p.word(top.get(.type), &tb);
+ const time = parseTime(p.word(top.get(.timestamp), &ts));
+ // Text the harness put in the user's mouth.
+ const injected = p.isTrue(top.get(.isMeta)) or p.isTrue(top.get(.isCompactSummary));
+ const said: Kind = if (injected) .system else .user;
+
+ if (std.mem.eql(u8, ty, "attachment")) {
+ // What a person typed while the model was working is queued, and
+ // lands here instead of in a `user` record. Queued prompts from
+ // anyone else (a peer, a task finishing) are the harness talking.
+ const a = p.fields(top.get(.attachment));
+ if (a.bad) return true;
+ var ab: [32]u8 = undefined;
+ if (!std.mem.eql(u8, p.word(a.get(.type), &ab), "queued_command")) return true;
+ const prompt = a.get(.prompt) orelse return true;
+ if (!p.isString(prompt)) return true;
+ const o = p.fields(a.get(.origin));
+ var ob: [32]u8 = undefined;
+ const human = !o.bad and std.mem.eql(u8, p.word(o.get(.kind), &ob), "human");
+ return one(c, p, if (human and !injected) .user else .system, time, prompt);
+ }
+ if (std.mem.eql(u8, ty, "system")) {
+ const content = top.get(.content) orelse return true;
+ if (!p.isString(content)) return true;
+ return one(c, p, .system, time, content);
+ }
+ const is_user = std.mem.eql(u8, ty, "user");
+ if (!is_user and !std.mem.eql(u8, ty, "assistant")) return true;
+
+ const m = p.fields(top.get(.message));
+ if (m.bad) return true;
+ const content = m.get(.content) orelse return true;
+ if (p.isString(content)) return one(c, p, if (is_user) said else .assistant, time, content);
+
+ var nm: [lookahead]u8 = undefined;
+ var it = p.items(content);
+ while (it.next()) |el| {
+ const b = p.fields(el.val);
+ if (b.bad) continue;
+ var bt: [32]u8 = undefined;
+ const bty = p.word(b.get(.type), &bt);
+ const ok = if (std.mem.eql(u8, bty, "text"))
+ one(c, p, if (is_user) said else .assistant, time, b.get(.text))
+ else if (std.mem.eql(u8, bty, "thinking")) blk: {
+ // Claude Code keeps only the signature of redacted thinking:
+ // an empty string is nothing to show.
+ const t = b.get(.thinking) orelse break :blk true;
+ if (std.mem.startsWith(u8, p.bytes(t), "\"\"")) break :blk true;
+ break :blk one(c, p, .thinking, time, t);
+ } else if (std.mem.endsWith(u8, bty, "tool_use"))
+ call(c, p, time, p.word(b.get(.name), &nm), b.get(.id), b.get(.name), b.get(.input))
+ else if (std.mem.endsWith(u8, bty, "tool_result"))
+ result(c, p, time, b.get(.tool_use_id), b.get(.content))
+ else if (is_user and std.mem.eql(u8, bty, "image")) blk: {
+ if (!c.open(said, time)) break :blk false;
+ if (!c.add(litSpan(lit_image))) break :blk false;
+ c.close();
+ break :blk true;
+ } else true;
+ if (!ok) return false;
+ }
+ return true;
+}
+
+// ---- Codex --------------------------------------------------------------------
+
+/// One line of a codex rollout. The conversation is the `response_item`
+/// records; `event_msg` repeats parts of it for the UI and is skipped, as
+/// are the session and turn metadata. False only when a cap is reached.
+fn codexLine(c: *Chat, p: *P, start: u64) bool {
+ const top = p.fields(p.ws(start));
+ if (top.bad) return true;
+ var tb: [32]u8 = undefined;
+ var ts: [40]u8 = undefined;
+ if (!std.mem.eql(u8, p.word(top.get(.type), &tb), "response_item")) return true;
+ const time = parseTime(p.word(top.get(.timestamp), &ts));
+ const f = p.fields(top.get(.payload));
+ if (f.bad) return true;
+ var pb: [40]u8 = undefined;
+ const ty = p.word(f.get(.type), &pb);
+
+ if (std.mem.eql(u8, ty, "message")) {
+ var rb: [32]u8 = undefined;
+ const role = p.word(f.get(.role), &rb);
+ const kind: Kind = if (std.mem.eql(u8, role, "user"))
+ (if (codexInjected(p, f.get(.content))) .system else .user)
+ else if (std.mem.eql(u8, role, "assistant"))
+ .assistant
+ else
+ .system; // developer, system
+ if (!c.open(kind, time)) return false;
+ var it = p.items(f.get(.content));
+ while (it.next()) |el| {
+ const b = p.fields(el.val);
+ var bt: [32]u8 = undefined;
+ const bty = p.word(b.get(.type), &bt);
+ const s: ?Span = if (b.get(.text)) |t|
+ (if (p.isString(t)) strSpan(p, t) else null)
+ else if (std.mem.endsWith(u8, bty, "image")) litSpan(lit_image) else null;
+ if (!addSpan(c, s)) return false;
+ }
+ c.close();
+ return true;
+ }
+ if (std.mem.eql(u8, ty, "reasoning")) {
+ // The summary, where codex wrote one; the reasoning itself is
+ // encrypted and never shown. Nothing readable, no message.
+ if (!c.open(.thinking, time)) return false;
+ for ([_]P.Key{ .summary, .content }) |k| {
+ var it = p.items(f.get(k));
+ while (it.next()) |el| {
+ const b = p.fields(el.val);
+ const t = b.get(.text) orelse continue;
+ if (!p.isString(t)) continue;
+ if (!addSpan(c, strSpan(p, t))) return false;
+ }
+ }
+ c.close();
+ return true;
+ }
+ if (std.mem.eql(u8, ty, "agent_message")) {
+ return blocks(c, p, .agent, time, 0, f.get(.author), f.get(.content));
+ }
+ var nm: [lookahead]u8 = undefined;
+ if (std.mem.eql(u8, ty, "function_call")) {
+ return call(c, p, time, p.word(f.get(.name), &nm), f.get(.call_id), f.get(.name), f.get(.arguments));
+ }
+ if (std.mem.eql(u8, ty, "custom_tool_call")) {
+ return call(c, p, time, p.word(f.get(.name), &nm), f.get(.call_id), f.get(.name), f.get(.input));
+ }
+ if (std.mem.endsWith(u8, ty, "_output")) {
+ return result(c, p, time, f.get(.call_id), f.get(.output));
+ }
+ if (std.mem.endsWith(u8, ty, "_call")) {
+ // local_shell_call, web_search_call, ...: no name of their own, so
+ // the record's type names the tool (less its `_call`) and its
+ // action is the input.
+ return call(c, p, time, ty[0 .. ty.len - "_call".len], f.get(.call_id), f.get(.type), f.get(.action));
+ }
+ return true;
+}
+
+/// The openings codex gives the context it writes in the user's role:
+/// the AGENTS.md it loaded, the environment, goals, plugin lists.
+const codex_injected = [_][]const u8{
+ "# AGENTS.md instructions",
+ "<environment_context>",
+ "<user_instructions>",
+ "<codex_internal_context",
+ "<recommended_plugins>",
+};
+
+/// Whether a user-role codex message is codex speaking: its first text
+/// opens with one of codex's own wrappers.
+fn codexInjected(p: *P, content: ?u64) bool {
+ var it = p.items(content);
+ const el = it.next() orelse return false;
+ const b = p.fields(el.val);
+ const t = b.get(.text) orelse return false;
+ if (!p.isString(t)) return false;
+ const head = p.bytes(t + 1);
+ for (codex_injected) |m| if (std.mem.startsWith(u8, head, m)) return true;
+ return false;
+}
+
+// ---- timestamps ---------------------------------------------------------------
+
+/// `2026-09-24T12:51:26.691Z` as Unix seconds; 0 for anything else.
+pub fn parseTime(s: []const u8) u32 {
+ if (s.len < 19 or s[4] != '-' or s[7] != '-' or s[10] != 'T' or s[13] != ':' or s[16] != ':') return 0;
+ const y = std.fmt.parseInt(i64, s[0..4], 10) catch return 0;
+ const mo = std.fmt.parseInt(i64, s[5..7], 10) catch return 0;
+ const d = std.fmt.parseInt(i64, s[8..10], 10) catch return 0;
+ const hh = std.fmt.parseInt(i64, s[11..13], 10) catch return 0;
+ const mi = std.fmt.parseInt(i64, s[14..16], 10) catch return 0;
+ const ss = std.fmt.parseInt(i64, s[17..19], 10) catch return 0;
+ if (mo < 1 or mo > 12 or d < 1 or d > 31 or hh > 23 or mi > 59 or ss > 60) return 0;
+ // Days from the civil date (Howard Hinnant's algorithm).
+ const yy = if (mo <= 2) y - 1 else y;
+ const era = @divFloor(yy, 400);
+ const yoe = yy - era * 400;
+ const mp = @mod(mo + 9, 12);
+ const doy = @divFloor(153 * mp + 2, 5) + d - 1;
+ const doe = yoe * 365 + @divFloor(yoe, 4) - @divFloor(yoe, 100) + doy;
+ const days = era * 146097 + doe - 719468;
+ const t = days * 86400 + hh * 3600 + mi * 60 + ss;
+ return std.math.cast(u32, t) orelse 0;
+}
+
+// ---- unit tests ---------------------------------------------------------------
+
+const testing = std.testing;
+
+const Fixture = struct {
+ dir: testing.TmpDir,
+ pool: *Pool,
+ file: Io.File = undefined,
+
+ fn init() !Fixture {
+ const pool = try Pool.create();
+ return .{ .dir = testing.tmpDir(.{}), .pool = pool };
+ }
+
+ fn deinit(f: *Fixture) void {
+ f.file.close(testing.io);
+ f.dir.cleanup();
+ f.pool.destroy();
+ }
+
+ /// Writes the transcript (replacing it) and opens it for reading.
+ fn write(f: *Fixture, bytes: []const u8) !void {
+ var w = try f.dir.dir.createFile(testing.io, "t.jsonl", .{});
+ defer w.close(testing.io);
+ try w.writeStreamingAll(testing.io, bytes);
+ f.file = try f.dir.dir.openFile(testing.io, "t.jsonl", .{});
+ }
+
+ fn append(f: *Fixture, bytes: []const u8) !void {
+ var w = try f.dir.dir.openFile(testing.io, "t.jsonl", .{ .mode = .read_write });
+ defer w.close(testing.io);
+ const st = try w.stat(testing.io);
+ try w.writePositionalAll(testing.io, bytes, st.size);
+ }
+
+ fn sync(f: *Fixture, format: Format) !*Chat {
+ return f.pool.sync(testing.io, f.file, format) orelse error.NotIndexed;
+ }
+
+ /// The names of every message, space-separated.
+ fn names(f: *Fixture, c: *const Chat, buf: []u8) []const u8 {
+ _ = f;
+ var w = Io.Writer.fixed(buf);
+ var nb: [96]u8 = undefined;
+ for (0..c.n) |i| {
+ if (i > 0) w.writeAll(" ") catch {};
+ w.writeAll(c.name(@intCast(i), &nb).?) catch {};
+ }
+ return w.buffered();
+ }
+
+ fn text(f: *Fixture, c: *const Chat, i: u32, buf: []u8) []const u8 {
+ const n = f.pool.render(testing.io, f.file, c, i, 0, buf);
+ return buf[0..n];
+ }
+};
+
+const claude_fixture =
+ \\{"type":"permission-mode","permissionMode":"default"}
+ \\{"type":"user","message":{"role":"user","content":"fix the \"race\"\nplease"},"timestamp":"2026-09-24T12:51:26.691Z","isMeta":false}
+ \\{"type":"assistant","message":{"role":"assistant","content":[{"type":"thinking","thinking":"","signature":"abc"}]}}
+ \\{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"caf\u00e9 \ud83d\ude00 done"}]}}
+ \\{"type":"assistant","message":{"content":[{"type":"tool_use","id":"toolu_1","name":"Bash","input":{"command":"ls -la"}}],"role":"assistant"}}
+ \\{"type":"user","message":{"role":"user","content":[{"tool_use_id":"toolu_1","type":"tool_result","content":"a\nb\n","is_error":false}]}}
+ \\{"type":"attachment","attachment":{"type":"total_tokens_reminder"}}
+ \\{"type":"user","message":{"role":"user","content":[{"type":"tool_result","tool_use_id":"t2","content":[{"type":"text","text":"shot"},{"type":"image","source":{"data":"AAAA"}}]}]}}
+ \\{"type":"user","isMeta":true,"message":{"role":"user","content":"Base directory for this skill"}}
+ \\{"type":"user","isSidechain":true,"message":{"role":"user","content":"not ours"}}
+ \\{"type":"system","subtype":"local_command","content":"<local-command-stdout>ok</local-command-stdout>"}
+ \\{"type":"user","message":{"content":"broken"
+ \\{"type":"user","message":{"role":"user","content":[{"type":"text","text":"thanks"}]}}
+ \\{"type":"attachment","attachment":{"type":"queued_command","prompt":"also restart it","commandMode":"prompt","origin":{"kind":"human"},"humanTurn":true}}
+ \\{"type":"attachment","attachment":{"type":"queued_command","prompt":"<task-notification>done</task-notification>","commandMode":"task-notification"}}
+ \\
+;
+
+test "chat: a Claude Code transcript reads as numbered messages" {
+ var f = try Fixture.init();
+ defer f.deinit();
+ try f.write(claude_fixture);
+ const c = try f.sync(.claude);
+ var nb: [512]u8 = undefined;
+ try testing.expectEqualStrings(
+ "00000000-user 00000001-assistant 00000002-bash 00000003-bash-result 00000004-tool-result 00000005-system 00000006-system 00000007-user 00000008-user 00000009-system",
+ f.names(c, &nb),
+ );
+ var tb: [256]u8 = undefined;
+ try testing.expectEqualStrings("fix the \"race\"\nplease\n", f.text(c, 0, &tb));
+ try testing.expectEqualStrings("café 😀 done\n", f.text(c, 1, &tb));
+ try testing.expectEqualStrings("Bash\n{\"command\":\"ls -la\"}\n", f.text(c, 2, &tb));
+ try testing.expectEqualStrings("a\nb\n", f.text(c, 3, &tb));
+ try testing.expectEqualStrings("shot\n[image]\n", f.text(c, 4, &tb));
+ try testing.expectEqualStrings("Base directory for this skill\n", f.text(c, 5, &tb));
+ try testing.expectEqualStrings("<local-command-stdout>ok</local-command-stdout>\n", f.text(c, 6, &tb));
+ try testing.expectEqualStrings("thanks\n", f.text(c, 7, &tb));
+ try testing.expectEqualStrings("also restart it\n", f.text(c, 8, &tb));
+ try testing.expectEqualStrings("<task-notification>done</task-notification>\n", f.text(c, 9, &tb));
+ // The size a stat reports is exactly what a read returns.
+ for (0..c.n) |i| {
+ try testing.expectEqual(@as(usize, c.msgs[i].size), f.text(c, @intCast(i), &tb).len);
+ }
+ try testing.expectEqual(@as(u32, 1790254286), c.msgs[0].time);
+ try testing.expectEqual(@as(u32, 0), c.msgs[1].time);
+}
+
+test "chat: a codex rollout reads as numbered messages" {
+ var f = try Fixture.init();
+ defer f.deinit();
+ try f.write(
+ \\{"timestamp":"2026-09-24T12:51:26.691Z","type":"session_meta","payload":{"id":"x","cwd":"/"}}
+ \\{"timestamp":"2026-09-24T12:51:27.000Z","type":"response_item","payload":{"type":"message","role":"developer","content":[{"type":"input_text","text":"be careful"}]}}
+ \\{"timestamp":"2026-09-24T12:51:28.000Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"one"},{"type":"input_text","text":"two\n"}]}}
+ \\{"timestamp":"2026-09-24T12:51:28.000Z","type":"event_msg","payload":{"type":"user_message","message":"one"}}
+ \\{"timestamp":"2026-09-24T12:51:29.000Z","type":"response_item","payload":{"type":"reasoning","summary":[],"content":null,"encrypted_content":"gAAA"}}
+ \\{"timestamp":"2026-09-24T12:51:29.000Z","type":"response_item","payload":{"type":"reasoning","summary":[{"type":"summary_text","text":"**plan**"}]}}
+ \\{"timestamp":"2026-09-24T12:51:30.000Z","type":"response_item","payload":{"type":"function_call","name":"shell","arguments":"{\"command\":[\"ls\"]}","call_id":"c1"}}
+ \\{"timestamp":"2026-09-24T12:51:31.000Z","type":"response_item","payload":{"type":"function_call_output","call_id":"c1","output":"file\n"}}
+ \\{"timestamp":"2026-09-24T12:51:32.000Z","type":"response_item","payload":{"type":"custom_tool_call","name":"apply_patch","input":"*** Begin Patch","call_id":"c2"}}
+ \\{"timestamp":"2026-09-24T12:51:33.000Z","type":"response_item","payload":{"type":"custom_tool_call_output","call_id":"c2","output":[{"type":"input_text","text":"Done"},{"type":"input_text","text":"!"}]}}
+ \\{"timestamp":"2026-09-24T12:51:34.000Z","type":"response_item","payload":{"type":"web_search_call","status":"completed","action":{"type":"search","query":"zig"}}}
+ \\{"timestamp":"2026-09-24T12:51:35.000Z","type":"response_item","payload":{"type":"agent_message","author":"/root/w1","recipient":"/root","content":[{"type":"input_text","text":"hi"},{"type":"encrypted_content","encrypted_content":"gAAA"}]}}
+ \\{"timestamp":"2026-09-24T12:51:36.000Z","type":"compacted","payload":{"message":"","replacement_history":[{"type":"message"}]}}
+ \\{"timestamp":"2026-09-24T12:51:37.000Z","type":"response_item","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"all done"}]}}
+ \\{"timestamp":"2026-09-24T12:51:38.000Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"<environment_context>\\n <cwd>/x</cwd>"}]}}
+ \\{"timestamp":"2026-09-24T12:51:39.000Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"# AGENTS.md instructions for /x"}]}}
+ \\
+ );
+ const c = try f.sync(.codex);
+ var nb: [512]u8 = undefined;
+ try testing.expectEqualStrings(
+ "00000000-system 00000001-user 00000002-thinking 00000003-shell 00000004-shell-result 00000005-apply_patch 00000006-apply_patch-result 00000007-web_search 00000008-agent 00000009-assistant 00000010-system 00000011-system",
+ f.names(c, &nb),
+ );
+ var tb: [256]u8 = undefined;
+ try testing.expectEqualStrings("be careful\n", f.text(c, 0, &tb));
+ try testing.expectEqualStrings("one\ntwo\n", f.text(c, 1, &tb));
+ try testing.expectEqualStrings("**plan**\n", f.text(c, 2, &tb));
+ try testing.expectEqualStrings("shell\n{\"command\":[\"ls\"]}\n", f.text(c, 3, &tb));
+ try testing.expectEqualStrings("file\n", f.text(c, 4, &tb));
+ try testing.expectEqualStrings("apply_patch\n*** Begin Patch\n", f.text(c, 5, &tb));
+ try testing.expectEqualStrings("Done\n!\n", f.text(c, 6, &tb));
+ try testing.expectEqualStrings("web_search_call\n{\"type\":\"search\",\"query\":\"zig\"}\n", f.text(c, 7, &tb));
+ try testing.expectEqualStrings("/root/w1\nhi\n", f.text(c, 8, &tb));
+ try testing.expectEqualStrings("all done\n", f.text(c, 9, &tb));
+ try testing.expectEqual(@as(u32, 1790254287), c.msgs[0].time);
+}
+
+test "chat: a tool call is named after its tool, and its result after the call" {
+ var f = try Fixture.init();
+ defer f.deinit();
+ try f.write(
+ \\{"type":"assistant","message":{"content":[{"type":"tool_use","id":"a","name":"Read","input":{}}]}}
+ \\{"type":"assistant","message":{"content":[{"type":"tool_use","id":"b","name":"mcp__x__Send Mail/now","input":{}}]}}
+ \\{"type":"assistant","message":{"content":[{"type":"tool_use","id":"c","name":"User","input":{}}]}}
+ \\{"type":"assistant","message":{"content":[{"type":"tool_use","id":"d","name":"read","input":{}}]}}
+ \\{"type":"user","message":{"content":[{"type":"tool_result","tool_use_id":"b","content":"sent"}]}}
+ \\{"type":"user","message":{"content":[{"type":"tool_result","tool_use_id":"a","content":"text"}]}}
+ \\{"type":"user","message":{"content":[{"type":"tool_result","tool_use_id":"c","content":"?"}]}}
+ \\{"type":"user","message":{"content":[{"type":"tool_result","tool_use_id":"zz","content":"?"}]}}
+ \\
+ );
+ const c = try f.sync(.claude);
+ var nb: [512]u8 = undefined;
+ // Parallel calls answered out of order still find their tools; a name
+ // is made safe for a file; one that would read as another kind is not
+ // used, and neither is a result whose call is unknown.
+ try testing.expectEqualStrings(
+ "00000000-read 00000001-mcp__x__send-mail-now 00000002-tool-call 00000003-read " ++
+ "00000004-mcp__x__send-mail-now-result 00000005-read-result 00000006-tool-result 00000007-tool-result",
+ f.names(c, &nb),
+ );
+ try testing.expectEqual(@as(?u32, 4), c.lookup("00000004-mcp__x__send-mail-now-result"));
+ try testing.expect(c.lookup("00000004-tool-result") == null);
+ try testing.expect(c.lookup("00000000-tool-call") == null);
+ try testing.expect(c.lookup("00000000-read-result") == null);
+ var tb: [64]u8 = undefined;
+ try testing.expectEqualStrings("Read\n{}\n", f.text(c, 0, &tb));
+}
+
+test "chat: ids hold as the transcript grows; a half-written line waits" {
+ var f = try Fixture.init();
+ defer f.deinit();
+ try f.write(
+ \\{"type":"user","message":{"content":"first"}}
+ \\{"type":"assistant","message":{"content":[{"type":"text","text":"sec
+ );
+ var c = try f.sync(.claude);
+ try testing.expectEqual(@as(u32, 1), c.n); // the second line has no newline yet
+ try f.append("ond\"}]}}\n{\"type\":\"user\",\"message\":{\"content\":\"third\"}}\n");
+ c = try f.sync(.claude);
+ var nb: [128]u8 = undefined;
+ try testing.expectEqualStrings("00000000-user 00000001-assistant 00000002-user", f.names(c, &nb));
+ var tb: [64]u8 = undefined;
+ try testing.expectEqualStrings("second\n", f.text(c, 1, &tb));
+
+ // One name per message: the kind must match, the id has one spelling.
+ try testing.expectEqual(@as(?u32, 1), c.lookup("00000001-assistant"));
+ try testing.expect(c.lookup("00000001-user") == null);
+ try testing.expect(c.lookup("1-assistant") == null); // unpadded
+ try testing.expect(c.lookup("000000001-assistant") == null); // too wide
+ try testing.expect(c.lookup("+0000001-assistant") == null);
+ try testing.expect(c.lookup("0000000x-assistant") == null);
+ try testing.expect(c.lookup("00000003-user") == null);
+ try testing.expect(c.lookup("1") == null);
+ try testing.expect(c.lookup("-assistant") == null);
+
+ // A transcript rewritten shorter is a different transcript.
+ f.file.close(testing.io);
+ try f.write("{\"type\":\"assistant\",\"message\":{\"content\":\"new\"}}\n");
+ c = try f.sync(.claude);
+ try testing.expectEqualStrings("00000000-assistant", f.names(c, &nb));
+}
+
+test "chat: a transcript rewritten in place is indexed again" {
+ var f = try Fixture.init();
+ defer f.deinit();
+ try f.write("{\"type\":\"user\",\"message\":{\"content\":\"first\"}}\n");
+ var c = try f.sync(.claude);
+ var nb: [128]u8 = undefined;
+ try testing.expectEqualStrings("00000000-user", f.names(c, &nb));
+ // The same inode, truncated and rewritten longer: only the bytes can
+ // tell. (`cp backup live.jsonl` does exactly this.)
+ f.file.close(testing.io);
+ try f.write(
+ \\{"type":"assistant","message":{"content":[{"type":"text","text":"a different start"}]}}
+ \\{"type":"user","message":{"content":"then this"}}
+ \\
+ );
+ c = try f.sync(.claude);
+ try testing.expectEqualStrings("00000000-assistant 00000001-user", f.names(c, &nb));
+ var tb: [64]u8 = undefined;
+ try testing.expectEqualStrings("a different start\n", f.text(c, 0, &tb));
+}
+
+test "chat: reads at any offset match the whole, across the window's edge" {
+ var f = try Fixture.init();
+ defer f.deinit();
+ // A string far longer than the window, dense with escapes, so some
+ // escape straddles every window boundary the render crosses.
+ const unit = "ab\\u00e9\\n\\ud83d\\ude00\\\"x";
+ const plain = "abé\n😀\"x";
+ const reps = 12000;
+ var bytes: std.ArrayList(u8) = .empty;
+ defer bytes.deinit(testing.allocator);
+ try bytes.appendSlice(testing.allocator, "{\"type\":\"user\",\"message\":{\"content\":\"");
+ for (0..reps) |_| try bytes.appendSlice(testing.allocator, unit);
+ try bytes.appendSlice(testing.allocator, "\"}}\n");
+ try f.write(bytes.items);
+ const c = try f.sync(.claude);
+ try testing.expectEqual(@as(u32, 1), c.n);
+ const want_len = plain.len * reps + 1;
+ try testing.expectEqual(@as(u64, want_len), c.msgs[0].size);
+
+ const whole = try testing.allocator.alloc(u8, want_len);
+ defer testing.allocator.free(whole);
+ try testing.expectEqual(want_len, f.pool.render(testing.io, f.file, c, 0, 0, whole));
+ for (0..reps) |i| try testing.expectEqualStrings(plain, whole[i * plain.len ..][0..plain.len]);
+ try testing.expectEqual(@as(u8, '\n'), whole[want_len - 1]);
+
+ // Piecewise, at odd sizes, the way a client's reads arrive: each read
+ // resumes where the last stopped.
+ var off: usize = 0;
+ var piece: [4093]u8 = undefined;
+ while (off < want_len) {
+ const n = f.pool.render(testing.io, f.file, c, 0, off, &piece);
+ try testing.expect(n > 0);
+ try testing.expectEqualSlices(u8, whole[off..][0..n], piece[0..n]);
+ off += n;
+ }
+ try testing.expectEqual(@as(usize, 0), f.pool.render(testing.io, f.file, c, 0, want_len, &piece));
+ // Out of order — backwards, and repeating a read — must not trust the
+ // cursor a later read left behind.
+ var back: usize = want_len;
+ while (back > 0) {
+ const at = back -| 3001;
+ const n = f.pool.render(testing.io, f.file, c, 0, at, piece[0..@min(piece.len, back - at)]);
+ try testing.expectEqualSlices(u8, whole[at..][0..n], piece[0..n]);
+ const again = f.pool.render(testing.io, f.file, c, 0, at, piece[0..@min(piece.len, back - at)]);
+ try testing.expectEqual(n, again);
+ try testing.expectEqualSlices(u8, whole[at..][0..n], piece[0..n]);
+ back = at;
+ }
+}
+
+test "chat: a string ends at the quote an even run of backslashes leaves bare" {
+ var f = try Fixture.init();
+ defer f.deinit();
+ // Runs of escaped backslashes long enough to straddle the window's
+ // edge, each ending in an escaped quote that must not close the string.
+ var bytes: std.ArrayList(u8) = .empty;
+ defer bytes.deinit(testing.allocator);
+ const a = testing.allocator;
+ try bytes.appendSlice(a, "{\"type\":\"user\",\"message\":{\"content\":\"x");
+ for (0..3) |_| {
+ for (0..40000) |_| try bytes.appendSlice(a, "\\\\");
+ try bytes.appendSlice(a, "\\\"");
+ }
+ try bytes.appendSlice(a, "end\"}}\n{\"type\":\"assistant\",\"message\":{\"content\":\"after\"}}\n");
+ try f.write(bytes.items);
+ const c = try f.sync(.claude);
+ var nb: [64]u8 = undefined;
+ try testing.expectEqualStrings("00000000-user 00000001-assistant", f.names(c, &nb));
+ const want = 1 + 3 * (40000 + 1) + "end\n".len;
+ try testing.expectEqual(@as(u32, want), c.msgs[0].size);
+ const out = try a.alloc(u8, want);
+ defer a.free(out);
+ try testing.expectEqual(@as(usize, want), f.pool.render(testing.io, f.file, c, 0, 0, out));
+ try testing.expectEqualStrings("\\\"end\n", out[want - 6 ..]);
+ try testing.expectEqual(@as(usize, 3), std.mem.count(u8, out, "\""));
+ var tb: [16]u8 = undefined;
+ try testing.expectEqualStrings("after\n", f.text(c, 1, &tb));
+}
+
+test "chat: a transcript past the cap says so instead of ending early" {
+ var f = try Fixture.init();
+ defer f.deinit();
+ var bytes: std.ArrayList(u8) = .empty;
+ defer bytes.deinit(testing.allocator);
+ for (0..max_msgs + 1) |_| try bytes.appendSlice(testing.allocator, "{\"type\":\"user\",\"message\":{\"content\":\"\"}}\n");
+ try f.write(bytes.items);
+ const c = try f.sync(.claude);
+ try testing.expect(c.full);
+ try testing.expectEqual(@as(u32, max_msgs), c.n);
+}
+
+test "chat: malformed JSON never escapes the line it is on" {
+ var f = try Fixture.init();
+ defer f.deinit();
+ try f.write(
+ \\{"type":"user","message":{"content":"a"}}
+ \\]]]}}}{{{["""\\"
+ \\{"type":"user","message":{"content":[{"type":"text","text":"b"},]}}
+ \\{"type":"assistant","message":{"content":[{"type":"tool_use","name":"X","input":{"u":"\u"}}]}}
+ \\{"type":"user","message":{"content":"c"}}
+ \\
+ );
+ const c = try f.sync(.claude);
+ var nb: [128]u8 = undefined;
+ // The trailing comma ends that walk after "b"; the stray `\u` is not
+ // an escape, so it is served as written.
+ try testing.expectEqualStrings("00000000-user 00000001-user 00000002-x 00000003-user", f.names(c, &nb));
+ var tb: [64]u8 = undefined;
+ try testing.expectEqualStrings("X\n{\"u\":\"\\u\"}\n", f.text(c, 2, &tb));
+ // A sign or an underscore is not a hex digit.
+ var sink: Sink = .{ .out = &tb };
+ _ = decode("\\u+0e9\\u0_e9\\u00e9", 0, true, &sink);
+ try testing.expectEqualStrings("\\u+0e9\\u0_e9é", tb[0..sink.n]);
+ try testing.expectEqualStrings("c\n", f.text(c, 3, &tb));
+}
+
+test "chat: timestamps" {
+ try testing.expectEqual(@as(u32, 0), parseTime("2026-09-24T12:51:26"[0..10]));
+ try testing.expectEqual(@as(u32, 86400), parseTime("1970-01-02T00:00:00Z"));
+ try testing.expectEqual(@as(u32, 1790254286), parseTime("2026-09-24T12:51:26.691Z"));
+ try testing.expectEqual(@as(u32, 0), parseTime("not a time at all, no"));
+}
diff --git a/9agents/src/main.zig b/9agents/src/main.zig
index 2ac9078..3e7c290 100644
--- a/9agents/src/main.zig
+++ b/9agents/src/main.zig
@@ -43,6 +43,7 @@ const usage_text =
\\ /omp /hermes /dsh full mirrors, raw
\\ /skills/{claude,codex,omp} the union skills view
\\ /active/<harness>/<pid>/ what is running right now
+ \\ /active/<harness>/<pid>/chat/<id>-<kind> its conversation (claude, codex)
\\
\\By default the daemon posts itself under the name `agents`, so it is
\\dialable at $XDG_RUNTIME_DIR/9p/agents and mountable by 9ns --mntgen
@@ -71,10 +72,27 @@ const Runner = serve.Runner(tree.Harness, tree.opts, .{
.listeners = 2, // the posted name plus one of --unix/--tcp
});
-// Static memory, the 9proc-demo way: the harness and the runner are large
-// (the path table, the per-connection buffers) and live in .bss.
+// Static memory, the 9proc-demo way: the harness (the path table) lives in
+// .bss. The runner and the chat pool are larger — every connection's fid
+// table, every transcript index — and each gets a mapping of its own
+// (`mapped`), backed only as far as it is used.
var harness_mem: tree.Harness = undefined;
-var runner_mem: Runner = undefined;
+
+/// A `T` in its own anonymous `MAP_NORESERVE` mapping: address space until
+/// written, and never written just to initialize it (a Debug build fills
+/// an `undefined` global with 0xaa, which for these is hundreds of MB).
+fn mapped(comptime T: type) error{OutOfMemory}!*T {
+ const rc = linux.mmap(null, @sizeOf(T), .{ .READ = true, .WRITE = true }, .{
+ .TYPE = .PRIVATE,
+ .ANONYMOUS = true,
+ .NORESERVE = true,
+ }, -1, 0);
+ if (linux.errno(rc) != .SUCCESS) return error.OutOfMemory;
+ // Small pages: with transparent huge pages on, touching each
+ // connection's few header fields would back 2 MB apiece.
+ _ = linux.madvise(@ptrFromInt(rc), @sizeOf(T), linux.MADV.NOHUGEPAGE);
+ return @ptrFromInt(rc);
+}
var stop_requested = std.atomic.Value(bool).init(false);
fn onSignal(sig: linux.SIG) callconv(.c) void {
@@ -243,6 +261,9 @@ fn run(init: std.process.Init) !void {
return error.Usage;
}
+ // The chat indexes live in their own lazily backed mapping: ~100 MB of
+ // address space, none of it memory until a chat is read.
+ const chats = try tree.chat.Pool.create();
harness_mem.init(.{
.io = io,
.pid = @intCast(linux.getpid()),
@@ -253,6 +274,7 @@ fn run(init: std.process.Init) !void {
.runtime = post.getenv(envp, "XDG_RUNTIME_DIR") orelse "",
.allow_move = allow_move,
.envp = envp,
+ .chats = chats,
});
if (allow_move) {
std.debug.print("9agents: moves allowed — a write to /active/<h>/<pid>/zmx re-execs that agent under {s}\n", .{zmx_abs});
@@ -273,12 +295,13 @@ fn run(init: std.process.Init) !void {
.posted, .unix, .tcp => {},
}
- runner_mem.init(.{ .io = io, .root = tree.root, .handler = .{ .ctx = &harness_mem, .serve = serveReq } });
- defer runner_mem.stop();
+ const runner = try mapped(Runner);
+ runner.init(.{ .io = io, .root = tree.root, .handler = .{ .ctx = &harness_mem, .serve = serveReq } });
+ defer runner.stop();
var posted_something = false;
if (!no_post) {
- runner_mem.listenPosted(envp, post_name, 16) catch |err| {
+ runner.listenPosted(envp, post_name, 16) catch |err| {
if (err == error.AlreadyPosted) {
std.debug.print("9agents: the name `{s}` is already posted by a live server\n", .{post_name});
return err;
@@ -293,7 +316,7 @@ fn run(init: std.process.Init) !void {
}
switch (mode) {
.unix => {
- _ = try runner_mem.listen(.{ .unix = try arena.dupeZ(u8, unix_path) }, 16);
+ _ = try runner.listen(.{ .unix = try arena.dupeZ(u8, unix_path) }, 16);
std.debug.print("9agents: listening on {s}\n", .{unix_path});
},
.tcp => {
@@ -301,7 +324,7 @@ fn run(init: std.process.Init) !void {
std.debug.print("9agents: bad --tcp address: {s}\n", .{tcp_addr});
return error.Usage;
};
- const bound = try runner_mem.listen(.{ .tcp = addr }, 16);
+ const bound = try runner.listen(.{ .tcp = addr }, 16);
std.debug.print("9agents: listening on tcp!{f}\n", .{bound});
},
else => {},
@@ -336,7 +359,9 @@ fn serveFd(io: Io, n: i32) !void {
var rbuf: [tree.msize]u8 = undefined;
var wbuf: [2 * tree.msize]u8 = undefined;
var stage: [tree.msize]u8 = undefined;
- var engine = Engine.init(.{ .in = &in, .out = &out, .root = tree.root });
+ // Far too big for the stack at this fid capacity.
+ const engine = try mapped(Engine);
+ engine.initIn(.{ .in = &in, .out = &out, .root = tree.root });
var reader = stream.reader(io, &rbuf);
var writer = stream.writer(io, &wbuf);
while (true) {
diff --git a/9agents/src/tree.zig b/9agents/src/tree.zig
index ed0699c..5e513f1 100644
--- a/9agents/src/tree.zig
+++ b/9agents/src/tree.zig
@@ -43,6 +43,7 @@ const E = fs.E;
const Io = std.Io;
const linux = std.os.linux;
pub const active = @import("active.zig");
+pub const chat = @import("chat.zig");
// ---- comptime bounds ---------------------------------------------------------
@@ -64,7 +65,10 @@ pub const list_name_bytes: usize = 128 * 1024;
pub const base_capacity: usize = 512;
pub const opts: fs.Options = .{
- .fid_capacity = 512,
+ // A kernel mount (9ns) holds a fid for every inode the kernel caches,
+ // and `ls -l` of one chat caches thousands. The engine is set up in
+ // place (`initIn`), so fids cost memory only once a client holds them.
+ .fid_capacity = 32768,
.slot_capacity = 8,
.name_capacity = 255,
.username_capacity = 28,
@@ -214,6 +218,10 @@ pub const AFile = enum(u8) {
agent,
agent_model,
agent_transcript,
+ /// `/active/<harness>/<pid>/chat`
+ chat,
+ /// `/active/<harness>/<pid>/chat/<id>-<kind>`
+ message,
/// The fields of one entry, in listing order. A field with no value
/// is not listed and does not resolve.
@@ -238,14 +246,31 @@ pub const Act = struct {
/// Only for `harness_dir`.
kind: active.Kind = .claude,
slot: u8 = 0,
- gen: u8 = 0,
+ gen: u24 = 0,
agent: u8 = 0,
+ /// Only for `message`: its id in the chat.
+ msg: u32 = 0,
};
+/// The serial of an `/active` node: the live slot (6 bits), its
+/// generation (24 bits: a slot churned through millions of processes
+/// before an old id could name a new one) and the subagent or message
+/// (18 bits).
+const act_slot_bits = 6;
+const act_gen_bits = 24;
+const act_arg_bits = 18;
+comptime {
+ std.debug.assert(act_slot_bits + act_gen_bits + act_arg_bits == 48);
+ std.debug.assert(active.max_live <= 1 << act_slot_bits);
+ std.debug.assert(chat.max_msgs <= 1 << act_arg_bits);
+ std.debug.assert(active.max_agents <= 1 << act_arg_bits);
+}
+
pub fn activeNode(a: Act) u64 {
+ const arg: u48 = if (a.file == .message) a.msg else a.agent;
const serial: u48 = switch (a.file) {
.harness_dir => @intFromEnum(a.kind),
- else => @as(u48, a.slot) | (@as(u48, a.gen) << 8) | (@as(u48, a.agent) << 16),
+ else => @as(u48, a.slot) | (@as(u48, a.gen) << act_slot_bits) | (arg << (act_slot_bits + act_gen_bits)),
};
return @bitCast(Node{
.idx = @intFromEnum(a.file),
@@ -422,8 +447,11 @@ pub const Harness = struct {
self_pid: u32 = 0,
btime: i64 = 0,
allow_move: bool = false,
+ /// The chat indexes behind `/active/<h>/<pid>/chat`. Null leaves the
+ /// chats out of the tree.
+ chats: ?*chat.Pool = null,
act_text: [1024]u8 = undefined,
- act_name: [64]u8 = undefined,
+ act_name: [96]u8 = undefined,
act_path: [active.path_capacity]u8 = undefined,
pub const InitOptions = struct {
@@ -447,6 +475,11 @@ pub const Harness = struct {
envp: ?cloud9.post.Env = null,
/// Whether writing `/active/<h>/<pid>/zmx` may move a session.
allow_move: bool = false,
+ /// The chat indexes (`chat.Pool.create`). Null leaves every agent
+ /// without a `chat/`. Kept outside the harness because it is large
+ /// and mostly untouched: its own mapping costs nothing until a chat
+ /// is read, where assigning the harness would write it all.
+ chats: ?*chat.Pool = null,
};
pub fn init(h: *Harness, o: InitOptions) void {
@@ -471,6 +504,7 @@ pub const Harness = struct {
}
h.allow_move = o.allow_move;
h.envp = o.envp;
+ h.chats = o.chats;
inline for (0..5) |i| {
const path = o.bases[i];
h.base_set[i] = path.len > 0 and path.len <= base_capacity;
@@ -550,11 +584,14 @@ pub const Harness = struct {
const k = std.enums.fromInt(active.Kind, @as(u8, @truncate(n.serial))) orelse return null;
return .{ .act = .{ .file = f, .kind = k } };
}
- const slot: u8 = @truncate(n.serial);
- const gen: u8 = @truncate(n.serial >> 8);
- const agent: u8 = @truncate(n.serial >> 16);
+ const slot: u6 = @truncate(n.serial);
+ const gen: u24 = @truncate(n.serial >> act_slot_bits);
+ const arg: u18 = @truncate(n.serial >> (act_slot_bits + act_gen_bits));
+ // The argument is a message id or a subagent index, by file.
+ const agent: u8 = if (f == .message) 0 else std.math.cast(u8, arg) orelse return null;
+ const msg: u32 = if (f == .message) arg else 0;
const l = h.live.at(slot, gen) orelse return null;
- return .{ .act = .{ .file = f, .kind = l.kind, .slot = slot, .gen = gen, .agent = agent } };
+ return .{ .act = .{ .file = f, .kind = l.kind, .slot = slot, .gen = gen, .agent = agent, .msg = msg } };
},
}
}
@@ -749,14 +786,25 @@ const readme_text =
\\active/<harness>/<pid>/ holds pid ppid started cwd name title
\\session via status transcript zmx.
\\ session its session id, empty when it did not resolve
- \\ via how that id was found: registry, fd, dir or none
+ \\ via how that id was found: registry, fd, dir, lock or none
\\ status only harnesses that publish one have it
\\ transcript the live .jsonl; its mtime is the last activity
\\ zmx the zmx session it runs in, empty outside zmx
+ \\ chat/ the conversation, one file per message (claude, codex)
+ \\
+ \\chat/<id>-<kind>: ids are 8 digits, counting from 0 in the order the
+ \\transcript holds them, and never change; kind is user, assistant,
+ \\thinking, system, agent, or for a tool call the tool (bash, read,
+ \\shell, ...) and for its answer the tool and -result (bash-result).
+ \\A file is the message as text; a call is the tool's name, then its
+ \\input. Not yet: chats of sessions that are not running, of
+ \\subagents, or of omp, hermes and dsh; and a session changed inside
+ \\a running agent (/clear, /new) keeps showing the one it began with.
\\
\\ ls active/*/* what is running
\\ cat active/claude/345104/session what it is
\\ tail -f active/claude/345104/transcript
+ \\ cat active/claude/345104/chat/00000012-assistant
\\
\\Credentials are excluded by name and never listed, at any depth;
\\symlinks are never served. Writing a name into an agent's zmx file
@@ -885,11 +933,11 @@ fn activeText(h: *Harness, a: Act) ?[]const u8 {
.via => line(h, l.via.text()),
.name => line(h, active.registryField(src, l, "name", &scratch) orelse return null),
.status => line(h, active.registryField(src, l, "status", &scratch) orelse return null),
- .title => line(h, active.headField(h.io, l.kind, l.transcript.slice(), .title, &scratch) orelse return null),
- .model => line(h, active.headField(h.io, l.kind, l.transcript.slice(), .model, &scratch) orelse return null),
+ .title => line(h, headOf(h, l, l.transcript.slice(), .title, &scratch) orelse return null),
+ .model => line(h, headOf(h, l, l.transcript.slice(), .model, &scratch) orelse return null),
.agent_model => blk: {
const path = activePath(h, a) orelse break :blk null;
- break :blk line(h, active.headField(h.io, l.kind, path, .model, &scratch) orelse return null);
+ break :blk line(h, headOf(h, l, path, .model, &scratch) orelse return null);
},
// `zmx` is always there, so it can always be written to; empty
// means the agent runs outside zmx.
@@ -912,26 +960,104 @@ fn activePath(h: *Harness, a: Act) ?[]const u8 {
}
}
-/// Opens an absolute path without following a symlink at the last
-/// component. These paths are derived from a pinned root or from the
-/// fd the harness itself holds, never from client bytes, but a name
-/// swapped underneath must still fail rather than redirect.
-fn openAbsNoFollow(path: []const u8) ?i32 {
- if (path.len == 0 or path.len >= active.path_capacity) return null;
- var z: [active.path_capacity]u8 = @splat(0);
- @memcpy(z[0..path.len], path);
- const rc = linux.open(@ptrCast(&z), .{
- .ACCMODE = .RDONLY,
- .NOFOLLOW = true,
- .CLOEXEC = true,
- .NONBLOCK = true,
- }, 0);
- if (linux.errno(rc) != .SUCCESS) return null;
- return @intCast(rc);
+/// Opens a file an agent's record points at — its transcript, a
+/// subagent's — the way the mirror opens everything: from the harness's
+/// pinned root, one component at a time, never following a symlink
+/// (`openIn`). The paths come from a harness's own records and the fds it
+/// holds, never from a client, but a directory swapped for a symlink, or
+/// a record naming `..`, must still not reach outside the root. A path
+/// that is not below the root is not opened at all. Non-blocking, so a
+/// fifo cannot park the daemon; the caller checks what it got.
+fn openAgentFile(h: *Harness, k: active.Kind, path: []const u8) ?Io.File {
+ const r: Root = @enumFromInt(@intFromEnum(k));
+ const b = h.base(r) orelse return null;
+ if (path.len <= b.len + 1 or !std.mem.startsWith(u8, path, b) or path[b.len] != '/') return null;
+ const fd = openIn(h, r, path[b.len + 1 ..], .{ .ACCMODE = .RDONLY, .NONBLOCK = true, .CLOEXEC = true }) orelse return null;
+ return .{ .handle = fd, .flags = .{ .nonblocking = false } };
+}
+
+/// `openAgentFile`, for a regular file only, ready for plain preads.
+fn openAgentRegular(h: *Harness, k: active.Kind, path: []const u8) ?Io.File {
+ const file = openAgentFile(h, k, path) orelse return null;
+ const st = file.stat(h.io) catch {
+ file.close(h.io);
+ return null;
+ };
+ if (st.kind != .file) {
+ file.close(h.io);
+ return null;
+ }
+ _ = linux.fcntl(file.handle, linux.F.SETFL, 0);
+ return file;
+}
+
+/// A title or model from the head of a transcript.
+fn headOf(h: *Harness, l: *const active.Live, path: []const u8, which: active.Head, out: []u8) ?[]const u8 {
+ if (path.len == 0) return null;
+ const file = openAgentRegular(h, l.kind, path) orelse return null;
+ defer file.close(h.io);
+ return active.headField(h.io, l.kind, file, which, out);
+}
+
+/// The transcript format an agent's chat is read in, if its harness has
+/// one this server reads.
+fn chatFormat(k: active.Kind) ?chat.Format {
+ return switch (k) {
+ .claude => .claude,
+ .codex => .codex,
+ else => null,
+ };
+}
+
+/// An agent has a `chat/` when the view is on, its harness's transcripts
+/// are readable here, and its transcript resolved and still opens from
+/// its root: a `chat/` that lists but cannot be read would be a lie.
+fn hasChat(h: *Harness, l: *const active.Live) bool {
+ if (h.chats == null or chatFormat(l.kind) == null or l.transcript.len == 0) return false;
+ const file = openAgentRegular(h, l.kind, l.transcript.slice()) orelse return false;
+ file.close(h.io);
+ return true;
+}
+
+const OpenChat = struct {
+ c: *chat.Chat,
+ file: Io.File,
+};
+
+/// Opens an agent's transcript and brings its chat index up to date with
+/// it. The caller closes `file`; the index is only good while it is open.
+fn openChat(h: *Harness, a: Act) ?OpenChat {
+ const l = h.live.at(a.slot, a.gen) orelse return null;
+ if (h.chats == null or chatFormat(l.kind) == null) return null;
+ const file = openAgentRegular(h, l.kind, l.transcript.slice()) orelse return null;
+ const c = h.chats.?.sync(h.io, file, chatFormat(l.kind).?) orelse {
+ file.close(h.io);
+ return null;
+ };
+ return .{ .c = c, .file = file };
}
fn activeAttr(h: *Harness, a: Act) ?fs.Attr {
switch (a.file) {
+ .chat => {
+ const l = h.live.at(a.slot, a.gen) orelse return null;
+ if (!hasChat(h, l)) return null;
+ return .{ .name = "chat", .node = activeNode(a), .dir = true, .mode = 0o555 };
+ },
+ .message => {
+ const oc = openChat(h, a) orelse return null;
+ defer oc.file.close(h.io);
+ const m = oc.c.msg(a.msg) orelse return null;
+ return .{
+ .name = oc.c.name(a.msg, &h.act_name) orelse return null,
+ .node = activeNode(a),
+ .size = m.size,
+ .mode = 0o444,
+ // When it was said, where the harness wrote that down.
+ .mtime = m.time,
+ .version = qidVers(m.time, m.size),
+ };
+ },
.harness_dir => {
if (!hasLive(h, a.kind)) return null;
return .{ .name = a.kind.text(), .node = activeNode(a), .dir = true, .mode = 0o555 };
@@ -957,8 +1083,10 @@ fn activeAttr(h: *Harness, a: Act) ?fs.Attr {
},
.transcript, .agent_transcript => {
const path = activePath(h, a) orelse return null;
- const st = Io.Dir.statFile(.cwd(), h.io, path, .{ .follow_symlinks = false }) catch return null;
- if (st.kind != .file) return null;
+ const l = h.live.at(a.slot, a.gen) orelse return null;
+ const file = openAgentRegular(h, l.kind, path) orelse return null;
+ defer file.close(h.io);
+ const st = file.stat(h.io) catch return null;
const mtime = st.mtime.toSeconds();
return .{
.name = a.file.fileName(),
@@ -1011,6 +1139,9 @@ fn activeLookup(h: *Harness, req: fs.Req, a: Act, name: []const u8) Answer {
if (std.mem.eql(u8, name, "agents")) {
return activeReply(h, req, .{ .file = .agents, .kind = a.kind, .slot = a.slot, .gen = a.gen }, name);
}
+ if (std.mem.eql(u8, name, "chat")) {
+ return activeReply(h, req, .{ .file = .chat, .kind = a.kind, .slot = a.slot, .gen = a.gen }, name);
+ }
for (AFile.entry_files) |f| {
if (!std.mem.eql(u8, f.fileName(), name)) continue;
return activeReply(h, req, .{ .file = f, .kind = a.kind, .slot = a.slot, .gen = a.gen }, name);
@@ -1043,6 +1174,12 @@ fn activeLookup(h: *Harness, req: fs.Req, a: Act, name: []const u8) Answer {
}
return fail(req.tag, E.NOENT);
},
+ .chat => {
+ const oc = openChat(h, a) orelse return fail(req.tag, E.NOENT);
+ defer oc.file.close(h.io);
+ const i = oc.c.lookup(name) orelse return fail(req.tag, E.NOENT);
+ return activeReply(h, req, .{ .file = .message, .kind = a.kind, .slot = a.slot, .gen = a.gen, .msg = i }, name);
+ },
else => return fail(req.tag, E.NOTDIR),
}
}
@@ -1058,6 +1195,9 @@ fn activeReaddir(h: *Harness, req: fs.Req, a: Act) Answer {
var st: Staging = .{ .buf = &h.stage_buf, .skip = req.off };
switch (a.file) {
.harness_dir => {
+ // A listing is what rescans, here as at /active: an agent that
+ // exited does not list, one that started does.
+ if (req.off == 0) rescan(h);
for (&h.live.slots, 0..) |*sl, i| {
if (!sl.used or sl.live.kind != a.kind) continue;
var num: [24]u8 = undefined;
@@ -1080,6 +1220,24 @@ fn activeReaddir(h: *Harness, req: fs.Req, a: Act) Answer {
if (l.agent_dir.len > 0) {
st.add(activeNode(.{ .file = .agents, .kind = a.kind, .slot = a.slot, .gen = a.gen }), true, "agents");
}
+ if (hasChat(h, l)) {
+ st.add(activeNode(.{ .file = .chat, .kind = a.kind, .slot = a.slot, .gen = a.gen }), true, "chat");
+ }
+ },
+ .chat => {
+ const oc = openChat(h, a) orelse return fail(req.tag, E.NOENT);
+ defer oc.file.close(h.io);
+ // Loud, like every other cap here: a listing that stopped at
+ // the cap would read as a chat that ended there.
+ if (oc.c.full) return failWhy(req.tag, E.NFILE, "chat is longer than this server can index (messages, or parts of one message)");
+ // The records before the offset are only counted, never named.
+ var i: u32 = @intCast(@min(st.skip, oc.c.n));
+ st.skip = 0;
+ var nb: [96]u8 = undefined;
+ while (i < oc.c.n and !st.full) : (i += 1) {
+ const nm = oc.c.name(i, &nb) orelse continue;
+ st.add(activeNode(.{ .file = .message, .kind = a.kind, .slot = a.slot, .gen = a.gen, .msg = i }), false, nm);
+ }
},
.agents => {
const l = h.live.at(a.slot, a.gen) orelse return fail(req.tag, E.NOENT);
@@ -1109,16 +1267,20 @@ fn activeReaddir(h: *Harness, req: fs.Req, a: Act) Answer {
fn activeRead(h: *Harness, req: fs.Req, a: Act) Answer {
switch (a.file) {
- .harness_dir, .entry, .agents, .agent => return fail(req.tag, E.ISDIR),
+ .harness_dir, .entry, .agents, .agent, .chat => return fail(req.tag, E.ISDIR),
+ .message => {
+ const oc = openChat(h, a) orelse return fail(req.tag, E.NOENT);
+ defer oc.file.close(h.io);
+ if (oc.c.msg(a.msg) == null) return fail(req.tag, E.NOENT);
+ const want = @min(req.size, h.data_buf.len);
+ const n = h.chats.?.render(h.io, oc.file, oc.c, a.msg, req.off, h.data_buf[0..want]);
+ return .{ .reply = .{ .tag = req.tag }, .bytes = h.data_buf[0..n] };
+ },
.transcript, .agent_transcript => {
const path = activePath(h, a) orelse return fail(req.tag, E.NOENT);
- const fd = openAbsNoFollow(path) orelse return fail(req.tag, E.NOENT);
- const file: Io.File = .{ .handle = fd, .flags = .{ .nonblocking = false } };
+ const l = h.live.at(a.slot, a.gen) orelse return fail(req.tag, E.NOENT);
+ const file = openAgentRegular(h, l.kind, path) orelse return fail(req.tag, E.NOENT);
defer file.close(h.io);
- const st = file.stat(h.io) catch return fail(req.tag, E.IO);
- if (st.kind == .directory) return fail(req.tag, E.ISDIR);
- if (st.kind != .file) return fail(req.tag, E.PERM);
- _ = linux.fcntl(fd, linux.F.SETFL, 0);
const want = @min(req.size, h.data_buf.len);
const n = file.readPositionalAll(h.io, h.data_buf[0..want], req.off) catch
return fail(req.tag, E.IO);
@@ -1445,7 +1607,8 @@ fn lookupParent(h: *Harness, req: fs.Req, t: Target) Answer {
.act => |a| switch (a.file) {
.harness_dir => .{ .top = .active },
.entry => .{ .act = .{ .file = .harness_dir, .kind = a.kind } },
- .agents => .{ .act = .{ .file = .entry, .kind = a.kind, .slot = a.slot, .gen = a.gen } },
+ .agents, .chat => .{ .act = .{ .file = .entry, .kind = a.kind, .slot = a.slot, .gen = a.gen } },
+ .message => .{ .act = .{ .file = .chat, .kind = a.kind, .slot = a.slot, .gen = a.gen } },
.agent => .{ .act = .{ .file = .agents, .kind = a.kind, .slot = a.slot, .gen = a.gen } },
.agent_model, .agent_transcript => .{ .act = .{
.file = .agent,
@@ -1714,6 +1877,8 @@ const Rig = struct {
self_pid: u32 = 4242,
zmx: []const u8 = "zmx",
allow_move: bool = false,
+ /// The chat indexes, with the fixture process tree only.
+ pool: ?*chat.Pool = null,
path_buf: [std.fs.max_path_bytes]u8 = undefined,
home: []const u8 = undefined,
h: *Harness = undefined,
@@ -1749,6 +1914,7 @@ const Rig = struct {
if (rig.with_proc) {
try Io.Dir.cwd().createDirPath(io, proc_root);
try rig.put("proc/stat", "cpu 1 2 3\nbtime 1000000\nprocesses 7\n");
+ rig.pool = try chat.Pool.create();
}
rig.harness_mem.init(.{
.io = io,
@@ -1758,11 +1924,13 @@ const Rig = struct {
.home = rig.home,
.zmx = rig.zmx,
.allow_move = rig.allow_move,
+ .chats = rig.pool,
});
rig.h = &rig.harness_mem;
}
fn end(rig: *Rig) void {
+ if (rig.pool) |p| p.destroy();
rig.dir.cleanup();
testing.allocator.free(rig.home);
}
@@ -2411,3 +2579,325 @@ test "active: the daemon never lists the process tree it lives in" {
try testing.expect(stageHas(pids.bytes, "1042")); // an unrelated session lists
try testing.expect(!stageHas(pids.bytes, "1040")); // its own parent does not
}
+
+// ---- /active/<h>/<pid>/chat ----------------------------------------------------
+
+/// A Claude Code agent whose transcript holds `lines`; answers its entry.
+fn claudeWithChat(rig: *Rig, pid: u32, lines: []const u8) !u64 {
+ try rig.fakeProc(.{ .pid = pid, .comm = "claude", .starttime = 5000, .argv = "claude\x00" });
+ var rel: [64]u8 = undefined;
+ var rec: [512]u8 = undefined;
+ try rig.put(try std.fmt.bufPrint(&rel, ".claude/sessions/{d}.json", .{pid}), try std.fmt.bufPrint(&rec,
+ \\{{"pid":{d},"sessionId":"sess-chat","cwd":"{s}","procStart":"5000"}}
+ , .{ pid, rig.home }));
+ try rig.put(try claudeTranscript(rig, &rec), lines);
+ var num: [16]u8 = undefined;
+ return activeEntryNode(rig, "claude", try std.fmt.bufPrint(&num, "{d}", .{pid}));
+}
+
+/// The fixture transcript of `claudeWithChat`, relative to the fake home.
+fn claudeTranscript(rig: *Rig, buf: []u8) ![]const u8 {
+ var slug_buf: [512]u8 = undefined;
+ const slug = active.slugOf(.claude, rig.home, rig.home, &slug_buf).?;
+ return std.fmt.bufPrint(buf, ".claude/projects/{s}/sess-chat.jsonl", .{slug});
+}
+
+fn readAll(rig: *Rig, node: u64, buf: []u8) ![]const u8 {
+ const r = handle(rig.h, .{ .tag = 70, .op = .read, .node = node, .off = 0, .size = @intCast(buf.len) });
+ try testing.expect(r.reply.status == .ok);
+ @memcpy(buf[0..r.bytes.len], r.bytes);
+ return buf[0..r.bytes.len];
+}
+
+test "chat: an agent's conversation is a directory of numbered messages" {
+ var rig: Rig = .{ .dir = undefined, .with_proc = true };
+ try rig.start();
+ defer rig.end();
+ const entry = try claudeWithChat(&rig, 1101,
+ \\{"type":"user","message":{"role":"user","content":"hello"},"timestamp":"2026-09-24T12:51:26.691Z"}
+ \\{"type":"assistant","message":{"content":[{"type":"tool_use","id":"t1","name":"Read","input":{"file_path":"/x"}}]}}
+ \\{"type":"user","message":{"content":[{"type":"tool_result","tool_use_id":"t1","content":"contents"}]}}
+ \\{"type":"assistant","message":{"content":[{"type":"text","text":"done"}]}}
+ \\
+ );
+ const listing = handle(rig.h, .{ .tag = 1, .op = .readdir, .node = entry, .off = 0, .size = msize });
+ try testing.expect(stageHas(listing.bytes, "chat"));
+
+ const dir = rig.lookupName(2, entry, "chat");
+ try testing.expect(dir.reply.status == .ok and dir.reply.attr.dir);
+ const msgs = handle(rig.h, .{ .tag = 3, .op = .readdir, .node = dir.reply.attr.node, .off = 0, .size = msize });
+ try testing.expectEqual(@as(usize, 4), stageCount(msgs.bytes));
+ for ([_][]const u8{ "00000000-user", "00000001-read", "00000002-read-result", "00000003-assistant" }) |name| {
+ try testing.expect(stageHas(msgs.bytes, name));
+ }
+
+ const first = rig.lookupName(4, dir.reply.attr.node, "00000000-user");
+ try testing.expect(first.reply.status == .ok and !first.reply.attr.dir);
+ try testing.expectEqual(@as(u32, 1790254286), first.reply.attr.mtime);
+ var buf: [256]u8 = undefined;
+ try testing.expectEqualStrings("hello\n", try readAll(&rig, first.reply.attr.node, &buf));
+ try testing.expectEqual(@as(u64, "hello\n".len), first.reply.attr.size);
+ const call = rig.lookupName(5, dir.reply.attr.node, "00000001-read");
+ try testing.expectEqualStrings("Read\n{\"file_path\":\"/x\"}\n", try readAll(&rig, call.reply.attr.node, &buf));
+
+ // One name per message: the wrong kind or another spelling is no file.
+ try expectNoent(rig.lookupName(6, dir.reply.attr.node, "00000000-assistant"));
+ try expectNoent(rig.lookupName(7, dir.reply.attr.node, "0-user"));
+ try expectNoent(rig.lookupName(8, dir.reply.attr.node, "00000004-user"));
+ try expectNoent(rig.lookupName(9, dir.reply.attr.node, "0"));
+
+ // Read-only like the rest, and ".." climbs back the way it came.
+ try expectEperm(handle(rig.h, .{ .tag = 10, .op = .write, .node = first.reply.attr.node, .data = "x" }));
+ const up = rig.lookupName(11, first.reply.attr.node, "..");
+ try testing.expectEqual(dir.reply.attr.node, up.reply.attr.node);
+ const upup = rig.lookupName(12, dir.reply.attr.node, "..");
+ try testing.expectEqual(entry, upup.reply.attr.node);
+
+ // The harness writes on: the next message takes the next id, and the
+ // ones before it keep theirs.
+ var rel: [640]u8 = undefined;
+ var more: [512]u8 = undefined;
+ const tr = try claudeTranscript(&rig, &rel);
+ var tb: [4096]u8 = undefined;
+ var path_buf: [std.fs.max_path_bytes]u8 = undefined;
+ const path = try std.fmt.bufPrint(&path_buf, "{s}/{s}", .{ rig.home, tr });
+ const old = try Io.Dir.cwd().readFile(testing.io, path, &tb);
+ try rig.put(tr, try std.fmt.bufPrint(&more, "{s}{s}", .{ old, "{\"type\":\"user\",\"message\":{\"content\":\"again\"}}\n" }));
+ const grown = handle(rig.h, .{ .tag = 13, .op = .readdir, .node = dir.reply.attr.node, .off = 0, .size = msize });
+ try testing.expectEqual(@as(usize, 5), stageCount(grown.bytes));
+ try testing.expect(stageHas(grown.bytes, "00000004-user"));
+ const again = rig.lookupName(14, dir.reply.attr.node, "00000004-user");
+ try testing.expectEqualStrings("again\n", try readAll(&rig, again.reply.attr.node, &buf));
+ try testing.expectEqualStrings("hello\n", try readAll(&rig, first.reply.attr.node, &buf));
+}
+
+test "chat: a codex agent's rollout, found through the fd it holds" {
+ var rig: Rig = .{ .dir = undefined, .with_proc = true };
+ try rig.start();
+ defer rig.end();
+ const rollout = ".codex/sessions/2026/09/21/rollout-2026-09-21T10-00-00-01a0bcd2-5653-7cc0-aa36-676f6467a72a.jsonl";
+ try rig.put(rollout,
+ \\{"timestamp":"2026-09-21T10:00:00.000Z","type":"session_meta","payload":{"id":"x"}}
+ \\{"timestamp":"2026-09-21T10:00:01.000Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"go"}]}}
+ \\{"timestamp":"2026-09-21T10:00:02.000Z","type":"response_item","payload":{"type":"function_call","name":"shell","arguments":"{\"cmd\":\"ls\"}"}}
+ \\
+ );
+ try rig.fakeProc(.{ .pid = 1102, .comm = "codex", .argv = "codex\x00" });
+ var link_buf: [std.fs.max_path_bytes]u8 = undefined;
+ var target_buf: [std.fs.max_path_bytes]u8 = undefined;
+ try Io.Dir.cwd().createDirPath(testing.io, try std.fmt.bufPrint(&link_buf, "{s}/proc/1102/fd", .{rig.home}));
+ const link = try std.fmt.bufPrintZ(&link_buf, "{s}/proc/1102/fd/3", .{rig.home});
+ const target = try std.fmt.bufPrint(&target_buf, "{s}/{s}", .{ rig.home, rollout });
+ try Io.Dir.cwd().symLink(testing.io, target, link, .{});
+
+ const entry = try activeEntryNode(&rig, "codex", "1102");
+ const dir = rig.lookupName(1, entry, "chat");
+ try testing.expect(dir.reply.status == .ok);
+ const msgs = handle(rig.h, .{ .tag = 2, .op = .readdir, .node = dir.reply.attr.node, .off = 0, .size = msize });
+ try testing.expectEqual(@as(usize, 2), stageCount(msgs.bytes));
+ const call = rig.lookupName(3, dir.reply.attr.node, "00000001-shell");
+ var buf: [128]u8 = undefined;
+ try testing.expectEqualStrings("shell\n{\"cmd\":\"ls\"}\n", try readAll(&rig, call.reply.attr.node, &buf));
+}
+
+test "chat: a codex whose fds cannot be read is found by the thread locks it holds" {
+ // Seen from inside a user namespace, /proc/<pid>/fd of an agent outside
+ // it is closed, so the fixture has none: only the locks tie the process
+ // to its threads.
+ var rig: Rig = .{ .dir = undefined, .with_proc = true };
+ try rig.start();
+ defer rig.end();
+ const root_id = "01a0d36b-6d58-7143-9ab9-9d9df6d0089a"; // minted 2026-09-24
+ const sub_id = "01a0da0f-d626-7382-ad35-3a1a44af26f2"; // minted 2026-09-25
+ const other_id = "01a0d8ca-8010-7fb0-bd56-568713201370"; // someone else's
+ try rig.put(".codex/sessions/2026/09/24/rollout-2026-09-24T09-37-08-" ++ root_id ++ ".jsonl",
+ \\{"timestamp":"2026-09-24T12:37:08.000Z","type":"session_meta","payload":{"id":"01a0d36b-6d58-7143-9ab9-9d9df6d0089a","thread_source":"user"}}
+ \\{"timestamp":"2026-09-24T12:37:09.000Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"root"}]}}
+ \\
+ );
+ try rig.put(".codex/sessions/2026/09/25/rollout-2026-09-25T16-34-26-" ++ sub_id ++ ".jsonl",
+ \\{"timestamp":"2026-09-25T19:34:26.000Z","type":"session_meta","payload":{"id":"01a0da0f-d626-7382-ad35-3a1a44af26f2","parent_thread_id":"01a0d36b-6d58-7143-9ab9-9d9df6d0089a","thread_source":"subagent"}}
+ \\
+ );
+ var inode: [3]u64 = undefined;
+ for ([_][]const u8{ root_id, sub_id, other_id }, 0..) |id, i| {
+ var rel: [128]u8 = undefined;
+ const lock = try std.fmt.bufPrint(&rel, ".codex/thread-writer-locks/{s}.lock", .{id});
+ try rig.put(lock, "");
+ var abs: [std.fs.max_path_bytes]u8 = undefined;
+ const st = try Io.Dir.cwd().statFile(testing.io, try std.fmt.bufPrint(&abs, "{s}/{s}", .{ rig.home, lock }), .{});
+ inode[i] = st.inode;
+ }
+ var locks: [1024]u8 = undefined;
+ try rig.put("proc/locks", try std.fmt.bufPrint(&locks,
+ \\1: FLOCK ADVISORY WRITE 1201 00:37:{d} 0 EOF
+ \\2: FLOCK ADVISORY WRITE 1201 00:37:{d} 0 EOF
+ \\2: -> FLOCK ADVISORY WRITE 1201 00:37:{d} 0 EOF
+ \\3: FLOCK ADVISORY WRITE 999 00:37:{d} 0 EOF
+ \\4: POSIX ADVISORY READ 1201 00:37:{d} 128 128
+ \\
+ , .{ inode[1], inode[0], inode[2], inode[2], inode[2] }));
+ try rig.fakeProc(.{ .pid = 1201, .comm = "codex", .argv = "codex\x00--dangerously-bypass-approvals-and-sandbox\x00" });
+
+ const entry = try activeEntryNode(&rig, "codex", "1201");
+ try testing.expectEqualStrings("lock\n", try activeField(&rig, entry, "via"));
+ try testing.expectEqualStrings(root_id ++ "\n", try activeField(&rig, entry, "session"));
+ const dir = rig.lookupName(1, entry, "chat");
+ const first = rig.lookupName(2, dir.reply.attr.node, "00000000-user");
+ var buf: [64]u8 = undefined;
+ try testing.expectEqualStrings("root\n", try readAll(&rig, first.reply.attr.node, &buf));
+}
+
+test "chat: a fifo where the transcript should be never blocks the daemon" {
+ var rig: Rig = .{ .dir = undefined, .with_proc = true };
+ try rig.start();
+ defer rig.end();
+ // The transcript must exist to resolve; then it becomes a fifo with no
+ // writer, which a blocking open would wait on for ever.
+ const entry = try claudeWithChat(&rig, 1301, "{\"type\":\"user\",\"message\":{\"content\":\"hi\"}}\n");
+ var rel: [640]u8 = undefined;
+ const tr = try claudeTranscript(&rig, &rel);
+ try rig.del(tr);
+ var abs: [std.fs.max_path_bytes]u8 = undefined;
+ const path = try std.fmt.bufPrintZ(&abs, "{s}/{s}", .{ rig.home, tr });
+ try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.mknodat(linux.AT.FDCWD, path, linux.S.IFIFO | 0o600, 0)));
+ for ([_][]const u8{ "title", "model", "transcript", "chat" }) |name| {
+ const f = rig.lookupName(1, entry, name);
+ if (f.reply.status == .ok) {
+ const dir = handle(rig.h, .{ .tag = 2, .op = .readdir, .node = f.reply.attr.node, .off = 0, .size = msize });
+ try testing.expect(dir.reply.status == .err);
+ }
+ }
+ const listing = handle(rig.h, .{ .tag = 3, .op = .readdir, .node = entry, .off = 0, .size = msize });
+ try testing.expect(listing.reply.status == .ok);
+}
+
+test "chat: a transcript reached through a symlinked directory is not served" {
+ var rig: Rig = .{ .dir = undefined, .with_proc = true };
+ try rig.start();
+ defer rig.end();
+ const entry = try claudeWithChat(&rig, 1302, "{\"type\":\"user\",\"message\":{\"content\":\"inside\"}}\n");
+ // The project directory is swapped for a link to a copy outside every
+ // root, after the agent resolved.
+ try rig.put("outside/sess-chat.jsonl", "{\"type\":\"user\",\"message\":{\"content\":\"OUTSIDE\"}}\n");
+ var slug_buf: [512]u8 = undefined;
+ const slug = active.slugOf(.claude, rig.home, rig.home, &slug_buf).?;
+ var a: [std.fs.max_path_bytes]u8 = undefined;
+ var b: [std.fs.max_path_bytes]u8 = undefined;
+ var o: [std.fs.max_path_bytes]u8 = undefined;
+ const dir = try std.fmt.bufPrint(&a, "{s}/.claude/projects/{s}", .{ rig.home, slug });
+ try Io.Dir.renameAbsolute(dir, try std.fmt.bufPrint(&b, "{s}/moved", .{rig.home}), testing.io);
+ try Io.Dir.cwd().symLink(testing.io, try std.fmt.bufPrint(&o, "{s}/outside", .{rig.home}), dir, .{});
+ try expectNoent(rig.lookupName(1, entry, "chat"));
+ try expectNoent(rig.lookupName(2, entry, "transcript"));
+}
+
+test "chat: a session id that is a path does not resolve" {
+ var rig: Rig = .{ .dir = undefined, .with_proc = true };
+ try rig.start();
+ defer rig.end();
+ try rig.put("outside/x.jsonl", "{\"type\":\"user\",\"message\":{\"content\":\"OUTSIDE\"}}\n");
+ try rig.fakeProc(.{ .pid = 1303, .comm = "claude", .starttime = 5000, .argv = "claude\x00" });
+ var rec: [512]u8 = undefined;
+ try rig.put(".claude/sessions/1303.json", try std.fmt.bufPrint(&rec,
+ \\{{"pid":1303,"sessionId":"../../../outside/x","cwd":"{s}","procStart":"5000"}}
+ , .{rig.home}));
+ const entry = try activeEntryNode(&rig, "claude", "1303");
+ try testing.expectEqualStrings("none\n", try activeField(&rig, entry, "via"));
+ try expectNoent(rig.lookupName(1, entry, "chat"));
+}
+
+test "chat: a slot churned past any u8 generation never serves a stale id" {
+ var rig: Rig = .{ .dir = undefined, .with_proc = true };
+ try rig.start();
+ defer rig.end();
+ const first = try claudeWithChat(&rig, 1304, "{\"type\":\"user\",\"message\":{\"content\":\"SECRET\"}}\n");
+ const dir = rig.lookupName(1, first, "chat");
+ const msg = rig.lookupName(2, dir.reply.attr.node, "00000000-user");
+ try testing.expect(msg.reply.status == .ok);
+ const act = rig.lookupName(3, root, "active");
+ for (0..300) |i| {
+ try rig.reapProc(1304);
+ _ = handle(rig.h, .{ .tag = 4, .op = .readdir, .node = act.reply.attr.node, .off = 0, .size = msize });
+ try rig.fakeProc(.{ .pid = 1304, .comm = "claude", .starttime = 6000 + i, .argv = "claude\x00" });
+ _ = handle(rig.h, .{ .tag = 5, .op = .readdir, .node = act.reply.attr.node, .off = 0, .size = msize });
+ }
+ try expectNoent(handle(rig.h, .{ .tag = 6, .op = .read, .node = msg.reply.attr.node, .off = 0, .size = 64 }));
+ try expectNoent(handle(rig.h, .{ .tag = 7, .op = .getattr, .node = first }));
+}
+
+test "active: listing a harness rescans, so an exited agent is gone" {
+ var rig: Rig = .{ .dir = undefined, .with_proc = true };
+ try rig.start();
+ defer rig.end();
+ try rig.fakeProc(.{ .pid = 1305, .comm = "claude", .argv = "claude\x00" });
+ try rig.fakeProc(.{ .pid = 1306, .comm = "claude", .argv = "claude\x00" });
+ _ = try activeEntryNode(&rig, "claude", "1305");
+ const act = rig.lookupName(1, root, "active");
+ const claude = rig.lookupName(2, act.reply.attr.node, "claude");
+ try rig.reapProc(1305);
+ const pids = handle(rig.h, .{ .tag = 3, .op = .readdir, .node = claude.reply.attr.node, .off = 0, .size = msize });
+ try testing.expect(!stageHas(pids.bytes, "1305"));
+ try testing.expect(stageHas(pids.bytes, "1306"));
+}
+
+test "chat: a harness whose transcripts are not read here has no chat" {
+ var rig: Rig = .{ .dir = undefined, .with_proc = true };
+ try rig.start();
+ defer rig.end();
+ try rig.fakeProc(.{ .pid = 1103, .comm = "claude", .argv = "claude\x00" }); // no session record
+ try rig.fakeProc(.{ .pid = 1104, .comm = "hermes", .argv = "hermes\x00" });
+ for ([_][2][]const u8{ .{ "claude", "1103" }, .{ "hermes", "1104" } }) |hp| {
+ const entry = try activeEntryNode(&rig, hp[0], hp[1]);
+ const listing = handle(rig.h, .{ .tag = 1, .op = .readdir, .node = entry, .off = 0, .size = msize });
+ try testing.expect(!stageHas(listing.bytes, "chat"));
+ try expectNoent(rig.lookupName(2, entry, "chat"));
+ }
+}
+
+test "chat: a long chat lists across several reads, every message once" {
+ var rig: Rig = .{ .dir = undefined, .with_proc = true };
+ try rig.start();
+ defer rig.end();
+ const count = 3000;
+ var lines: std.ArrayList(u8) = .empty;
+ defer lines.deinit(testing.allocator);
+ for (0..count) |_| try lines.appendSlice(testing.allocator, "{\"type\":\"user\",\"message\":{\"content\":\"m\"}}\n");
+ const entry = try claudeWithChat(&rig, 1105, lines.items);
+ const dir = rig.lookupName(1, entry, "chat");
+ var seen: [count]bool = @splat(false);
+ var off: u64 = 0;
+ var reads: usize = 0;
+ while (reads < 64) : (reads += 1) {
+ const page = handle(rig.h, .{ .tag = 2, .op = .readdir, .node = dir.reply.attr.node, .off = off, .size = msize });
+ try testing.expect(page.reply.status == .ok);
+ if (page.bytes.len == 0) break;
+ var i: usize = 0;
+ while (i + 10 <= page.bytes.len) {
+ const len: usize = page.bytes[i + 9];
+ const name = page.bytes[i + 10 ..][0..len];
+ i += 10 + len;
+ const dash = std.mem.indexOfScalar(u8, name, '-').?;
+ const id = try std.fmt.parseInt(usize, name[0..dash], 10);
+ try testing.expect(!seen[id]);
+ seen[id] = true;
+ }
+ off += stageCount(page.bytes);
+ }
+ try testing.expect(reads > 1);
+ for (seen) |s| try testing.expect(s);
+ // Every id resolves after the walk too, not just the ones a byte holds.
+ for ([_][]const u8{ "00000255-user", "00000256-user", "00002999-user" }) |name| {
+ const m = rig.lookupName(3, dir.reply.attr.node, name);
+ try testing.expect(m.reply.status == .ok);
+ const st = handle(rig.h, .{ .tag = 4, .op = .getattr, .node = m.reply.attr.node });
+ try testing.expect(st.reply.status == .ok);
+ try testing.expectEqualStrings(name, st.reply.attr.name);
+ var buf: [16]u8 = undefined;
+ try testing.expectEqualStrings("m\n", try readAll(&rig, m.reply.attr.node, &buf));
+ }
+}
+
+test {
+ _ = chat;
+}