Files
nucleic-remote-ios/NucleicRemote/NucleicRemote/Views/Transcript/IncrementalTranscriptProjection.swift
T

427 lines
21 KiB
Swift

import Foundation
import NucleicProtocol
/// Stable-prefix incremental projector (docs/TRANSCRIPT_INCREMENTAL_PROJECTION.md).
///
/// `TranscriptProjection.build` folds the whole stream on every read, so a streaming turn of K
/// deltas over an N-event transcript costs O(N·K) ≈ O(N²) per session. But the stream is
/// immutable except at the tail: everything before the live turn is frozen — its `messageID`s and
/// `toolCallID`s never recur. So this projector **seals** the longest provably-stable prefix of
/// the folded item list once, and re-folds only the unstable suffix per delta: O(live-tail)
/// instead of O(N).
///
/// The seal point (an index into the raw event stream) is chosen so that
/// `fold(prefix) ++ fold(tail)` is byte-identical to `fold(whole)`:
///
/// 1. **No coalescing key crosses the seam.** Text/thinking/tool items coalesce on
/// `messageID`/`toolCallID` (and a subagent's events name their parent Task). An exact
/// interval check over the current stream forbids any seam inside an id's first→last
/// reference span, so no item can straddle it.
/// 2. **No id can recur after the seam.** Future arrivals are fenced by closure rules read off
/// the stream itself: a tool is closed once its result arrived; a message/thinking block once
/// a different message has started in its scope; everything, once its turn completed. These
/// are the premises of the design doc ("the stream is immutable except at the tail"); if one
/// is ever violated — a tail event referencing a sealed id — it is *detected* and the
/// projector resets and re-folds from scratch, so correctness never rests on them.
/// 3. **No tool run is split.** `coalesceToolRuns` merges adjacent `.tool` items, so the seam
/// only falls where the last folded item is a hard separator — a visible non-tool row that a
/// future tool call can't merge across. (Empty redacted-thinking rows are transparent to runs
/// and therefore to this rule too.)
/// 4. **Lock notes fold exactly, even across the seam.** A lock's `released` note lands when the
/// file lands in the parent — potentially many turns after the edit it brackets — so sealed
/// edit cards stay reachable through a registry (`TranscriptProjection.PriorEdit`): a tail
/// note that path-matches a sealed edit patches that card, exactly where a whole-stream fold
/// would put it. No lag heuristic, no divergence.
/// 5. **`worktreeRoot` is a carried constant.** Lock-path normalization needs the seq-0
/// `sessionStarted` cwd; nothing seals until it is known, and it never changes once found.
///
/// The equivalence test (`IncrementalProjectionEquivalenceTests`) replays streams and asserts
/// `incremental(prefix) == build(prefix)` for **every** prefix — the whole correctness argument,
/// checked mechanically.
///
/// A class held in `@State` so reads/updates from `body` don't invalidate the view; not
/// thread-safe (main-actor use only, like the `ProjectionCache` it replaces).
final class IncrementalTranscriptProjection {
// MARK: - Sealed state
/// Folded output of the stable prefix — appended to at each seal, never re-walked.
private var sealedItems: [TranscriptItem] = []
/// Watermark into the raw stream: `events[0..<sealedEventCount]` produced `sealedItems`.
private var sealedEventCount = 0
/// `events[sealedEventCount - 1].seq` at seal time — detects a rewritten/merged prefix
/// (a reconnect backfill slotting events in by seq) that shifts history under the watermark.
private var sealedLastSeq: UInt64 = 0
/// The stream's first seq — detects a session switch / stream reset.
private var streamFirstSeq: UInt64?
/// Carried constant from the first non-empty `sessionStarted.cwd`. Nothing seals while nil
/// (a root arriving later would retroactively change sealed lock folds).
private var worktreeRoot: String?
/// Every `messageID`/`toolCallID` referenced by a sealed event. A later event referencing one
/// would mutate sealed output — detected here, answered with a full reset (self-healing).
private var sealedIDs: Set<String> = []
/// Sealed edit-class calls in item order, for cross-seam lock-note folding (rule 4).
private var sealedEdits: [TranscriptProjection.PriorEdit] = []
/// Where each sealed edit's card sits in `sealedItems` (it may live inside a `.toolBlock`).
private var sealedEditIndex: [String: Int] = [:]
/// Display toggles are fold inputs; flipping either resets.
private var showRaw = false
private var showLockEvents = true
// MARK: - Read memo
/// `body` re-evaluates far more often than the stream changes (settling layout, scroll
/// geometry); this collapses those redundant reads to a cached return, as the previous
/// `ProjectionCache` did. The stream is append-only with strictly increasing seq, so
/// `(count, firstSeq, lastSeq)` pins it.
private struct MemoKey: Equatable {
var count: Int
var firstSeq: UInt64
var lastSeq: UInt64
var showRaw: Bool
var showLockEvents: Bool
}
private var memoKey: MemoKey?
private var memoValue: [TranscriptItem] = []
/// Hold-back from the stream edge: never seal into the newest few events. Decoders emit
/// tightly-coupled events in one batch (a `toolResult` and its inferred `fileChange`; a
/// whole-message's final chunks) that sync may deliver one at a time — holding the edge back
/// keeps a mid-batch read from sealing an entity whose trailing batch-mates are still in
/// flight. Cheap insurance on top of the closure rules; violations would only cost a reset.
private let edgeLag = 4
/// Test hook: how far the watermark has advanced (the equivalence suite also asserts sealing
/// actually happens, so a regression to "never seal" can't pass silently).
var sealedEventCountForTesting: Int { sealedEventCount }
// MARK: - Read
func items(for events: [AgentEvent], showRaw: Bool, showLockEvents: Bool) -> [TranscriptItem] {
let key = MemoKey(count: events.count, firstSeq: events.first?.seq ?? 0,
lastSeq: events.last?.seq ?? 0, showRaw: showRaw, showLockEvents: showLockEvents)
if memoKey == key { return memoValue }
if needsReset(events, showRaw: showRaw, showLockEvents: showLockEvents) { reset() }
self.showRaw = showRaw
self.showLockEvents = showLockEvents
streamFirstSeq = events.first?.seq
if worktreeRoot == nil {
// Only the unsealed region needs scanning: sealing requires the root, so a sealed
// region can only exist after it was found.
worktreeRoot = TranscriptProjection.worktreeRoot(in: events[sealedEventCount...])
}
var watermark = chooseWatermark(events)
if watermark == nil {
// A tail event referenced a sealed id — a closure premise was violated (late file
// change, resumed message, post-result subagent child). Refold from scratch; with no
// sealed ids the second pass cannot be violated.
reset()
streamFirstSeq = events.first?.seq
worktreeRoot = TranscriptProjection.worktreeRoot(in: events[...])
watermark = chooseWatermark(events)
}
seal(events, upTo: watermark ?? sealedEventCount)
let result = render(events)
memoKey = key
memoValue = result
return result
}
// MARK: - Reset / identity
private func needsReset(_ events: [AgentEvent], showRaw: Bool, showLockEvents: Bool) -> Bool {
if showRaw != self.showRaw || showLockEvents != self.showLockEvents { return true }
if events.count < sealedEventCount { return true }
if sealedEventCount > 0 {
if events.first?.seq != streamFirstSeq { return true }
if events[sealedEventCount - 1].seq != sealedLastSeq { return true }
} else if streamFirstSeq != nil, events.first?.seq != streamFirstSeq {
return true // nothing sealed, but the carried worktreeRoot belongs to the old stream
}
return false
}
private func reset() {
sealedItems = []
sealedEventCount = 0
sealedLastSeq = 0
streamFirstSeq = nil
worktreeRoot = nil
sealedIDs = []
sealedEdits = []
sealedEditIndex = [:]
memoKey = nil
memoValue = []
}
// MARK: - Watermark selection
/// The coalescing ids an event mentions: its `messageID` or `toolCallID`, plus the parent
/// Task id for subagent-owned events. Two events sharing an id must land on the same side of
/// the seam; the parent link chains a subagent's whole scope (and, transitively, deeper
/// descendants) to its spawn.
private static func refs(of event: AgentEvent) -> [String] {
switch event.kind {
case .userText(let c), .assistantText(let c), .thinking(let c):
if let parent = c.parentToolCallID { return [c.messageID, parent] }
return [c.messageID]
case .toolCallStarted(let c), .toolCallCompleted(let c):
if let parent = c.parentToolCallID { return [c.toolCallID, parent] }
return [c.toolCallID]
case .toolCallInputDelta(let d): return [d.toolCallID]
case .toolResult(let r): return [r.toolCallID]
case .fileChange(let f): return f.toolCallID.map { [$0] } ?? []
case .note(let note):
// A note tagged with the edit it brackets (a subagent's lock-acquire / nvrsion-land)
// is now part of that edit's scope: chain it to the edit's id so a seal seam can't
// separate them. If the edit was already sealed (a rare late note after its scope
// closed), the tail reference to a sealed id trips the self-healing reset — a whole
// re-fold that necessarily matches `build`. Untagged notes chain to nothing, as before.
return note.toolCallID.map { [$0] } ?? []
default: return []
}
}
/// One id's life within the unsealed region.
private struct IDSpan {
var firstRef: Int
var lastRef: Int
/// Result arrived → the tool (or Task, with its children) is done.
var resultSeen = false
/// A later chunk with a different messageID in the same scope → this message is done
/// (its authoritative non-partial text can only arrive before the next message starts).
var closedByChunk = false
/// For message/thinking ids: the owning subagent scope ("" = top level). The owner
/// Task's result closes everything inside it.
var chunkScope: String?
}
/// What an event *creates* in the folded item list, for the run-split rule (3).
private enum Creation {
/// A visible, never-dropped, non-tool row — a safe last-item for a seam.
case separator
/// A `.tool` row a future adjacent call could merge with.
case tool
/// A thinking row: a separator iff its final text is non-empty (an empty redacted block
/// is transparent to run coalescing, so it must be transparent to the seam rule too).
case thinking(String)
/// A lock note that may fold away (dropping it can fuse the runs around it), so it
/// counts as nothing — the seam just waits for the next hard separator.
case transparent
}
/// The furthest event index the stream can be sealed to right now, or nil when a region event
/// references an already-sealed id (premise violation → caller resets).
private func chooseWatermark(_ events: [AgentEvent]) -> Int? {
let start = sealedEventCount
let n = events.count - start
// Everything below needs the worktree root (rule 5); without it, just verify no sealed-id
// violation … but nothing is sealed if no root was ever found, so there is nothing to do.
guard worktreeRoot != nil else { return start }
guard n > edgeLag else {
// Too little unsealed to advance, but tail refs must still be validated against
// sealed ids so a violation triggers the reset path.
for r in 0..<n where Self.refs(of: events[start + r]).contains(where: sealedIDs.contains) {
return nil
}
return start
}
var spans: [String: IDSpan] = [:]
var creations: [Creation?] = Array(repeating: nil, count: n)
var thinkingText: [String: String] = [:]
var seenMessageItem = Set<String>()
var seenThinkingItem = Set<String>()
var seenToolItem = Set<String>()
var lastChunkInScope: [String: String] = [:]
var lastBoundary = -1 // region index of the latest turnCompleted/runFinished
for r in 0..<n {
let event = events[start + r]
for id in Self.refs(of: event) {
if sealedIDs.contains(id) { return nil }
if var span = spans[id] {
span.lastRef = r
spans[id] = span
} else {
spans[id] = IDSpan(firstRef: r, lastRef: r)
}
}
switch event.kind {
case .userText(let c), .assistantText(let c):
trackChunk(c, in: &spans, lastChunkInScope: &lastChunkInScope)
if seenMessageItem.insert(c.messageID).inserted { creations[r] = .separator }
case .thinking(let c):
trackChunk(c, in: &spans, lastChunkInScope: &lastChunkInScope)
let existing = thinkingText[c.messageID] ?? ""
thinkingText[c.messageID] = c.isPartial ? existing + c.text : c.text
if seenThinkingItem.insert(c.messageID).inserted { creations[r] = .thinking(c.messageID) }
case .toolCallStarted(let c), .toolCallCompleted(let c):
if seenToolItem.insert(c.toolCallID).inserted { creations[r] = .tool }
case .toolResult(let result):
spans[result.toolCallID]?.resultSeen = true
case .toolCallInputDelta, .fileChange, .approvalResolved, .codexUsage:
break // refs (if any) tracked above; creates nothing
case .turnCompleted, .runFinished:
creations[r] = .separator
lastBoundary = r
case .sessionStarted, .usage, .rateLimit, .approvalRequested, .error:
creations[r] = .separator
case .note(let note):
if note.lockEvent && !showLockEvents { break }
if let lock = note.lock, !lock.paths.isEmpty { creations[r] = .transparent }
else { creations[r] = .separator }
case .raw:
if showRaw { creations[r] = .separator }
}
}
// Future-proofing (rule 2): an id wholly before the seam must be *closed* — provably done
// taking new events. An open id caps the seam at its first reference (it stays whole in
// the tail).
var cap = n - edgeLag
for (_, span) in spans where !isClosed(span, spans: spans, lastBoundary: lastBoundary) {
cap = min(cap, span.firstRef)
}
// Interval isolation (rule 1): no seam inside any id's [firstRef, lastRef] span.
// maxLastFromBefore[w] = the furthest lastRef among ids first referenced before w; a seam
// at w is isolation-safe iff that never reaches w.
var maxLastAtFirst = [Int](repeating: -1, count: n)
for span in spans.values {
maxLastAtFirst[span.firstRef] = max(maxLastAtFirst[span.firstRef], span.lastRef)
}
// Run-split rule (3): replay creations in item order; a seam is placeable after event r
// only while the last solid (visible, surviving) item is a hard separator. The replay
// resolves each thinking row against its *final* region text, which is exactly what the
// sealed fold will contain (open ids were already excluded by `cap`).
var separatorOK = [Bool](repeating: false, count: n)
var lastSolidIsSeparator = true // sealed prefix is empty or ends with a separator (invariant)
for r in 0..<n {
switch creations[r] {
case .separator: lastSolidIsSeparator = true
case .tool: lastSolidIsSeparator = false
case .thinking(let id):
let text = thinkingText[id] ?? ""
if !text.trimmingCharacters(in: .whitespacesAndNewlines).isEmpty {
lastSolidIsSeparator = true
}
case .transparent, nil: break
}
separatorOK[r] = lastSolidIsSeparator
}
var runningMaxLast = -1
var best = 0
if cap >= 1 {
for w in 1...cap {
runningMaxLast = max(runningMaxLast, maxLastAtFirst[w - 1])
if separatorOK[w - 1] && runningMaxLast < w { best = w }
}
}
return start + best
}
/// Track a text/thinking chunk for message-closure: a new messageID in a scope closes the
/// previous one (chunks of one message never resume after the next begins — decoder order).
private func trackChunk(
_ chunk: TextChunk, in spans: inout [String: IDSpan], lastChunkInScope: inout [String: String]
) {
let scope = chunk.parentToolCallID ?? ""
spans[chunk.messageID]?.chunkScope = scope
if let previous = lastChunkInScope[scope], previous != chunk.messageID {
spans[previous]?.closedByChunk = true
}
lastChunkInScope[scope] = chunk.messageID
}
private func isClosed(_ span: IDSpan, spans: [String: IDSpan], lastBoundary: Int) -> Bool {
if span.lastRef < lastBoundary { return true } // its turn completed; ids don't cross turns
if span.resultSeen { return true }
if span.closedByChunk { return true }
if let scope = span.chunkScope, !scope.isEmpty, spans[scope]?.resultSeen == true {
return true // the owning subagent returned; its inner stream is done
}
return false
}
// MARK: - Sealing
private func seal(_ events: [AgentEvent], upTo watermark: Int) {
guard watermark > sealedEventCount else { return }
let slice = events[sealedEventCount..<watermark]
let segment = TranscriptProjection.buildSegment(
slice, worktreeRoot: worktreeRoot, priorEdits: sealedEdits,
showRaw: showRaw, showLockEvents: showLockEvents)
// Lock notes in this segment that folded onto edits sealed earlier: bake them in — the
// note is now sealed too, so the fold is final.
for patch in segment.priorLockPatches {
Self.applyLock(patch.lock, to: patch.toolCallID, at: sealedEditIndex, in: &sealedItems)
}
let editIDs = Set(segment.edits.map(\.toolCallID))
for item in segment.items {
let index = sealedItems.count
switch item.kind {
case .tool(let group) where editIDs.contains(group.toolCallID):
sealedEditIndex[group.toolCallID] = index
case .toolBlock(let groups):
for group in groups where editIDs.contains(group.toolCallID) {
sealedEditIndex[group.toolCallID] = index
}
default:
break
}
sealedItems.append(item)
}
sealedEdits.append(contentsOf: segment.edits)
for event in slice {
for id in Self.refs(of: event) { sealedIDs.insert(id) }
}
sealedEventCount = watermark
sealedLastSeq = events[watermark - 1].seq
}
// MARK: - Rendering
private func render(_ events: [AgentEvent]) -> [TranscriptItem] {
let tailSlice = events[sealedEventCount...]
guard !tailSlice.isEmpty else { return sealedItems }
let tail = TranscriptProjection.buildSegment(
tailSlice, worktreeRoot: worktreeRoot, priorEdits: sealedEdits,
showRaw: showRaw, showLockEvents: showLockEvents)
if tail.priorLockPatches.isEmpty { return sealedItems + tail.items }
// A live (unsealed) lock note folded onto a sealed edit card: patch a copy per read —
// the note may still be re-evaluated until it seals, so the base stays unpatched.
var patched = sealedItems
for patch in tail.priorLockPatches {
Self.applyLock(patch.lock, to: patch.toolCallID, at: sealedEditIndex, in: &patched)
}
return patched + tail.items
}
/// Append a folded lock line to a sealed edit's card, whether it renders alone or inside a
/// coalesced `.toolBlock`.
private static func applyLock(
_ lock: NoteLock, to toolCallID: String, at index: [String: Int],
in items: inout [TranscriptItem]
) {
guard let i = index[toolCallID] else { return }
switch items[i].kind {
case .tool(var group) where group.toolCallID == toolCallID:
group.lockLines.append(lock)
items[i].kind = .tool(group)
case .toolBlock(var groups):
guard let k = groups.firstIndex(where: { $0.toolCallID == toolCallID }) else { return }
groups[k].lockLines.append(lock)
items[i].kind = .toolBlock(groups)
default:
break
}
}
}