diff --git a/NucleicRemote/NucleicRemote/Models/RemoteStore.swift b/NucleicRemote/NucleicRemote/Models/RemoteStore.swift index d041181..f62f4e5 100644 --- a/NucleicRemote/NucleicRemote/Models/RemoteStore.swift +++ b/NucleicRemote/NucleicRemote/Models/RemoteStore.swift @@ -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? + 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. private func persistOpenTranscript() { guard !demoMode, let id = openSessionID, !openEvents.isEmpty else { return }