native macOS codings agent orchestrator prowl.onev.cat
Something went wrong. Try again.
13 kB · 315 lines
Swift
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316import Foundation
/// What a Codex rollout says about ownership: its thread, its parent thread, and the/// client message ids of the user messages it holds, in file order.nonisolated struct CodexDaemonRollout: Sendable, Equatable { let id: String let parentID: String? var clientIDs: [String] var path = "" /// A fork copies history before its own `thread_settings_applied`; that part is not live. var isForked = false /// Byte offset of the live `task_started` line that precedes each client message id. var turnStartOffsets: [String: UInt64] = [:]}
/// The rollouts a pane's Codex TUI currently drives through the daemon (docs-ai 073.002).nonisolated struct CodexDaemonBinding: Sendable, Equatable { let rootID: String /// The root thread and every descendant the daemon holds open. let paths: [String] /// Where reading the bound rollout starts: the turn that holds the pane's newest submit. let liveOffsets: [String: UInt64]}
/// The submits of the TUI that last started a pane's session log.nonisolated struct CodexTUISessionRecord: Sendable, Equatable { struct Submit: Sendable, Equatable { let clientID: String let submittedAt: Date }
let surfaceID: UUID let startedAt: Date let submits: [Submit]}
/// Maps a thread that Codex's shared app-server daemon runs to the pane whose TUI drives it/// (docs-ai 073).////// Each pane runs Codex with a `CodexTUISessionLog`. Every submit writes a `UserTurn` op with/// a fresh `client_user_message_id`, and the daemon stores the same value as the `client_id`/// of the `UserMessage` item in the target thread's rollout. A pane is bound to the thread of/// its newest submit that a rollout already holds: a newer unmatched id is a turn whose item/// is not written yet, or a steer that waits for the next tool boundary. A child thread/// belongs to the root its rollout header names.actor CodexDaemonThreadMapper { static let shared = CodexDaemonThreadMapper()
typealias Rollout = CodexDaemonRollout typealias SessionLog = CodexTUISessionRecord
private struct Cursor { let inode: UInt64 var offset: UInt64 var pending = Data() var headerRead = false var live = true var lastTurnStart: UInt64? var rollout: Rollout? }
private let openFilePaths: @Sendable (pid_t) -> [String]? private let sessionLogDirectory: URL private let fileManager: FileManager private var cursors: [String: Cursor] = [:] /// A first lookup reads each open rollout once; later lookups read only appended bytes. private let readBudget = 256 * 1_024 * 1_024 private let sessionLogLimit = 64 * 1_024 * 1_024
init( sessionLogDirectory: URL = ProwlPaths.codexTUISessionLogDirectory, fileManager: FileManager = .default, openFilePaths: @escaping @Sendable (pid_t) -> [String]? = { pid in var complete = false let paths = ProcessDetection.openFilePaths(pid: pid, complete: &complete) return complete ? paths : nil } ) { self.sessionLogDirectory = sessionLogDirectory self.fileManager = fileManager self.openFilePaths = openFilePaths }
/// The pane driving `threadID`, or nil when the evidence is missing, incomplete, or ambiguous. func pane(threadID: String, daemonPID: pid_t) -> CodexThreadPane? { guard let paths = openFilePaths(daemonPID), let rollouts = refreshRollouts(paths.filter(Self.isRolloutPath)) else { return nil } return Self.resolve(threadID: threadID, rollouts: rollouts, logs: sessionLogs()) }
/// What the pane's current Codex TUI drives, or nil before its first indexed submit. The /// session log must belong to the TUI process that started at `tuiStartedAt`. func binding(surfaceID: UUID, daemonPID: pid_t, tuiStartedAt: Date) -> CodexDaemonBinding? { let url = CodexTUISessionLog.url(for: surfaceID, in: sessionLogDirectory) guard let size = (try? url.resourceValues(forKeys: [.fileSizeKey]))?.fileSize, size <= sessionLogLimit, let data = try? Data(contentsOf: url), let log = Self.parseSessionLog(data, surfaceID: surfaceID), CodexTUISessionLog.belongs(sessionStartedAt: log.startedAt, toProcessStartedAt: tuiStartedAt), !log.submits.isEmpty, let paths = openFilePaths(daemonPID), let rollouts = refreshRollouts(paths.filter(Self.isRolloutPath)) else { return nil } return Self.binding(for: log, rollouts: rollouts) }
// MARK: - Resolution
static func binding(for log: SessionLog, rollouts: [Rollout]) -> CodexDaemonBinding? { let byID = Dictionary(rollouts.map { ($0.id, $0) }, uniquingKeysWith: { first, _ in first }) let owner = owners(rollouts) guard let submit = log.submits.last(where: { owner[$0.clientID] != nil }), let bound = owner[submit.clientID].flatMap({ byID[$0] }) else { return nil } let rootID = root(of: bound.id, in: byID) let family = rollouts.filter { root(of: $0.id, in: byID) == rootID } return CodexDaemonBinding( rootID: rootID, paths: family.map(\.path).sorted(), liveOffsets: bound.turnStartOffsets[submit.clientID].map { [bound.path: $0] } ?? [:] ) }
private static func owners(_ rollouts: [Rollout]) -> [String: String] { var owner: [String: String] = [:] for rollout in rollouts { for clientID in rollout.clientIDs { owner[clientID] = rollout.id } } return owner }
private static func root(of id: String, in byID: [String: Rollout]) -> String { var current = id var visited: Set<String> = [] while let parent = byID[current]?.parentID, visited.insert(current).inserted { current = parent } return current }
static func resolve(threadID: String, rollouts: [Rollout], logs: [SessionLog]) -> CodexThreadPane? { let byID = Dictionary(rollouts.map { ($0.id, $0) }, uniquingKeysWith: { first, _ in first }) let owner = owners(rollouts) func root(_ id: String) -> String { Self.root(of: id, in: byID) } let target = root(threadID) let bound = logs.compactMap { log -> (log: SessionLog, submittedAt: Date)? in guard let submit = log.submits.last(where: { owner[$0.clientID] != nil }), let thread = owner[submit.clientID], root(thread) == target else { return nil } return (log, submit.submittedAt) } guard let newest = bound.max(by: { $0.submittedAt < $1.submittedAt }), bound.filter({ $0.submittedAt == newest.submittedAt }).count == 1 else { return nil } return CodexThreadPane(surfaceID: newest.log.surfaceID, sessionStartedAt: newest.log.startedAt) }
// MARK: - Session logs
private func sessionLogs() -> [SessionLog] { guard let entries = try? fileManager.contentsOfDirectory( at: sessionLogDirectory, includingPropertiesForKeys: [.fileSizeKey], options: [.skipsHiddenFiles]) else { return [] } return entries.compactMap { entry in guard let surfaceID = CodexTUISessionLog.surfaceID(forFileName: entry.lastPathComponent), let size = (try? entry.resourceValues(forKeys: [.fileSizeKey]))?.fileSize, size <= sessionLogLimit, let data = try? Data(contentsOf: entry) else { return nil } return Self.parseSessionLog(data, surfaceID: surfaceID) } }
/// The submits of the TUI that last started this log. Each TUI launch truncates the file. static func parseSessionLog(_ data: Data, surfaceID: UUID) -> SessionLog? { var startedAt: Date? var submits: [SessionLog.Submit] = [] for line in data.split(separator: 10) { let isStart = line.range(of: Data("\"session_start\"".utf8)) != nil guard isStart || line.range(of: Data("\"UserTurn\"".utf8)) != nil, let record = try? JSONSerialization.jsonObject(with: line) as? [String: Any], let timestamp = (record["ts"] as? String).flatMap(parseTimestamp) else { continue } if record["kind"] as? String == "session_start" { startedAt = timestamp submits = [] } else if record["kind"] as? String == "op", let turn = (record["payload"] as? [String: Any])?["UserTurn"] as? [String: Any], let clientID = turn["client_user_message_id"] as? String, !clientID.isEmpty { submits.append(.init(clientID: clientID, submittedAt: timestamp)) } } return startedAt.map { SessionLog(surfaceID: surfaceID, startedAt: $0, submits: submits) } }
private static func parseTimestamp(_ value: String) -> Date? { try? Date.ISO8601FormatStyle(includingFractionalSeconds: true).parse(value) }
// MARK: - Rollouts
static func isRolloutPath(_ path: String) -> Bool { let name = (path as NSString).lastPathComponent return name.hasPrefix("rollout-") && name.hasSuffix(".jsonl") }
/// Reads the bytes appended since the last lookup. Returns nil when a file changed identity /// or the budget ran out, because a partial index could bind a pane to a stale thread. private func refreshRollouts(_ paths: [String]) -> [Rollout]? { var next = cursors.filter { paths.contains($0.key) } var budget = readBudget for path in paths { guard let attributes = try? fileManager.attributesOfItem(atPath: path), let inode = attributes[.systemFileNumber] as? UInt64, let size = attributes[.size] as? UInt64 else { return nil } var cursor = next[path].flatMap { $0.inode == inode && size >= $0.offset ? $0 : nil } ?? Cursor(inode: inode, offset: 0) guard size > cursor.offset else { next[path] = cursor continue } guard size - cursor.offset <= UInt64(budget), let handle = FileHandle(forReadingAtPath: path) else { cursors = next return nil } defer { try? handle.close() } do { try handle.seek(toOffset: cursor.offset) let bytes = try handle.read(upToCount: Int(size - cursor.offset)) ?? Data() budget -= bytes.count cursor.offset += UInt64(bytes.count) cursor.pending.append(bytes) } catch { return nil } Self.consume(&cursor, path: path) next[path] = cursor } cursors = next return next.values.compactMap(\.rollout) }
private static func consume(_ cursor: inout Cursor, path: String) { let base = cursor.offset - UInt64(cursor.pending.count) var start = cursor.pending.startIndex while let end = cursor.pending[start...].firstIndex(of: 10) { let line = cursor.pending[start..<end] let lineOffset = base + UInt64(start - cursor.pending.startIndex) start = end + 1 if !cursor.headerRead { cursor.headerRead = true cursor.rollout = header(line) cursor.rollout?.path = path cursor.live = cursor.rollout?.isForked != true continue } guard let id = cursor.rollout?.id else { continue } if !cursor.live { cursor.live = isLiveBoundary(line, threadID: id) } else if isTaskStarted(line) { cursor.lastTurnStart = lineOffset } else if let clientID = userMessageClientID(line) { cursor.rollout?.clientIDs.append(clientID) cursor.rollout?.turnStartOffsets[clientID] = cursor.lastTurnStart } } cursor.pending.removeSubrange(cursor.pending.startIndex..<start) }
static func header(_ line: Data) -> Rollout? { guard let record = try? JSONSerialization.jsonObject(with: line) as? [String: Any], record["type"] as? String == "session_meta", let payload = record["payload"] as? [String: Any], let id = payload["id"] as? String, !id.isEmpty else { return nil } let source = payload["source"] as? [String: Any] let spawn = (source?["subagent"] as? [String: Any])?["thread_spawn"] as? [String: Any] return Rollout( id: id, parentID: spawn?["parent_thread_id"] as? String, clientIDs: [], isForked: payload["forked_from_id"] is String) }
private static func eventPayload(_ line: Data, type: String) -> [String: Any]? { guard line.range(of: Data("\"\(type)\"".utf8)) != nil, let record = try? JSONSerialization.jsonObject(with: line) as? [String: Any], record["type"] as? String == "event_msg", let payload = record["payload"] as? [String: Any], payload["type"] as? String == type else { return nil } return payload }
static func isTaskStarted(_ line: Data) -> Bool { eventPayload(line, type: "task_started") != nil }
static func isLiveBoundary(_ line: Data, threadID: String) -> Bool { eventPayload(line, type: "thread_settings_applied")?["thread_id"] as? String == threadID }
static func userMessageClientID(_ line: Data) -> String? { guard line.range(of: Data("\"client_id\"".utf8)) != nil, let record = try? JSONSerialization.jsonObject(with: line) as? [String: Any], record["type"] as? String == "event_msg", let payload = record["payload"] as? [String: Any], payload["type"] as? String == "item_completed", let item = payload["item"] as? [String: Any], item["type"] as? String == "UserMessage", let clientID = item["client_id"] as? String, !clientID.isEmpty else { return nil } return clientID }}