diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-25 16:40:49 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-09-25 17:52:33 -0300 |
| commit | 1c1b192d4ef59199a7196229d56901b1f6512678 (patch) | |
| tree | 9f6e4c1efcf877d1231aa58975ccd923ad4f2f2e /9agents/src/tree.zig | |
| parent | 73602127d15d10a1932b6fe916bd18a608054980 (diff) | |
| download | cloud9-1c1b192d4ef59199a7196229d56901b1f6512678.tar.gz cloud9-1c1b192d4ef59199a7196229d56901b1f6512678.zip | |
Each message of a claude or codex agent's transcript is a file named
<id>-<kind>: the id is its position in the transcript from 0, as 8
digits so the names sort; the kind is user, assistant, thinking, system
or agent, and a tool call is named after its tool (00000002-bash) and
its result after the call (00000003-bash-result). A file is the message
as text (a tool call: its name, then its input); its mtime is when it
was said.
chat.zig indexes a transcript incrementally, keyed by (dev, ino) and
checked by birth time and a hash of its first and last indexed bytes, so
a file rewritten in place is indexed again. It holds offsets, never
bytes: every read reads the file, and a read resumes where the last one
stopped so a long message is not decoded from its start each time. The
indexes and the runner live in lazily backed mappings (NORESERVE,
NOHUGEPAGE), and the fid table is 32768 so a kernel mount can hold every
message of a long chat.
A 9agents started inside a user namespace (from a 9ns-wrapped terminal)
cannot read /proc/<pid>/fd of processes outside it, so the fd route
never resolved codex there; /active now also finds a codex session
through the thread writer locks it holds (/proc/locks), via = lock.
Fixes to /active found on the way: a fifo in place of a transcript or
record no longer blocks the daemon; a transcript is opened from its
pinned root a component at a time, so a symlinked directory or a
sessionId with .. cannot reach outside it; slot generations are 24 bits,
so a stale node id cannot come to name another agent; listing a harness
rescans; an agent that had not resolved yet is asked again.
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Diffstat (limited to '9agents/src/tree.zig')
| -rw-r--r-- | 9agents/src/tree.zig | 566 |
1 files changed, 528 insertions, 38 deletions
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; +} |
