diff options
Diffstat (limited to '9agents/src')
| -rw-r--r-- | 9agents/src/active.zig | 207 | ||||
| -rw-r--r-- | 9agents/src/chat.zig | 1626 | ||||
| -rw-r--r-- | 9agents/src/main.zig | 43 | ||||
| -rw-r--r-- | 9agents/src/tree.zig | 566 |
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; +} |
