Transcript Incremental Projection
Nucleic-Session: 0FFD007B-0696-4517-9429-129C7B0FD5AC Co-authored-by: Nucleic <[email protected]>
This commit is contained in:
@@ -797,6 +797,61 @@ final class RemoteStore: ObservableObject {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// MARK: - Streaming delta coalescing
|
||||||
|
|
||||||
|
/// Streaming transcript deltas arrive one wire frame at a time — often dozens per second
|
||||||
|
/// while the agent talks — and every `openEvents` mutation fires `objectWillChange`, which
|
||||||
|
/// re-evaluates *every* view observing the store (the Sessions/Home tabs stay mounted behind
|
||||||
|
/// the pushed session detail, so they pay this too). Buffer incoming deltas and publish at
|
||||||
|
/// most one append per `openEventsFlushInterval`: the first delta after a quiet gap applies
|
||||||
|
/// immediately (the leading edge — first-token latency stays imperceptible), followers ride
|
||||||
|
/// the next scheduled flush. ~10 UI updates/sec still reads as live streaming; the view tree
|
||||||
|
/// stops being invalidated per wire frame. Merge/close/flush paths drain the buffer first, so
|
||||||
|
/// nothing downstream ever sees a partial stream.
|
||||||
|
private var pendingOpenEvents: [AgentEvent] = []
|
||||||
|
private var openEventsFlushTask: Task<Void, Never>?
|
||||||
|
private var lastOpenEventsFlushAt = Date.distantPast
|
||||||
|
private static let openEventsFlushInterval: TimeInterval = 0.1
|
||||||
|
|
||||||
|
private func enqueueOpenEvents(_ events: [AgentEvent]) {
|
||||||
|
pendingOpenEvents.append(contentsOf: events)
|
||||||
|
guard openEventsFlushTask == nil else { return } // a trailing flush is already scheduled
|
||||||
|
let elapsed = Date().timeIntervalSince(lastOpenEventsFlushAt)
|
||||||
|
if elapsed >= Self.openEventsFlushInterval {
|
||||||
|
drainPendingOpenEvents()
|
||||||
|
} else {
|
||||||
|
let delay = Self.openEventsFlushInterval - elapsed
|
||||||
|
openEventsFlushTask = Task { [weak self] in
|
||||||
|
try? await Task.sleep(nanoseconds: UInt64(delay * 1_000_000_000))
|
||||||
|
guard let self, !Task.isCancelled else { return }
|
||||||
|
self.openEventsFlushTask = nil
|
||||||
|
self.drainPendingOpenEvents()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Publish the buffered deltas (and schedule persistence). Called on the flush cadence, and
|
||||||
|
/// eagerly by anything that merges, persists, or clears `openEvents`, so those paths always
|
||||||
|
/// operate on the complete stream.
|
||||||
|
private func drainPendingOpenEvents() {
|
||||||
|
openEventsFlushTask?.cancel()
|
||||||
|
openEventsFlushTask = nil
|
||||||
|
guard !pendingOpenEvents.isEmpty else { return }
|
||||||
|
lastOpenEventsFlushAt = Date()
|
||||||
|
openEvents.append(contentsOf: pendingOpenEvents)
|
||||||
|
pendingOpenEvents.removeAll(keepingCapacity: true)
|
||||||
|
persistOpenTranscript()
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Drop buffered deltas without publishing — for session switch/close/unpair, where the
|
||||||
|
/// buffer belongs to a transcript that is being cleared (a stale session's tail must never
|
||||||
|
/// leak into the next session's freshly-opened transcript).
|
||||||
|
private func discardPendingOpenEvents() {
|
||||||
|
openEventsFlushTask?.cancel()
|
||||||
|
openEventsFlushTask = nil
|
||||||
|
pendingOpenEvents.removeAll(keepingCapacity: true)
|
||||||
|
}
|
||||||
|
|
||||||
/// Debounced write of the open transcript, called after each batch of events lands.
|
/// Debounced write of the open transcript, called after each batch of events lands.
|
||||||
private func persistOpenTranscript() {
|
private func persistOpenTranscript() {
|
||||||
guard !demoMode, let id = openSessionID, !openEvents.isEmpty else { return }
|
guard !demoMode, let id = openSessionID, !openEvents.isEmpty else { return }
|
||||||
|
|||||||
Reference in New Issue
Block a user