Transcript Incremental Projection
Nucleic-Session: 0FFD007B-0696-4517-9429-129C7B0FD5AC Co-authored-by: Nucleic <[email protected]>
This commit is contained in:
@@ -0,0 +1,417 @@
|
||||
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] } ?? []
|
||||
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:
|
||||
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
|
||||
for w in 1...min(cap, n - edgeLag) where w >= 1 {
|
||||
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
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user