diff --git a/NucleicRemote/NucleicRemote/Models/RemoteStore.swift b/NucleicRemote/NucleicRemote/Models/RemoteStore.swift index e02ecd0..5ccc3cb 100644 --- a/NucleicRemote/NucleicRemote/Models/RemoteStore.swift +++ b/NucleicRemote/NucleicRemote/Models/RemoteStore.swift @@ -639,20 +639,23 @@ final class RemoteStore: ObservableObject { guard let self, hostID == self.openSessionHostID, snap.summary.sessionID == self.openSessionID else { return } // Merge, don't replace: a fresh open merges into an empty transcript (the tail window), // while a reconnect's warm-resubscribe delta appends to the history already on screen - // instead of truncating it to the host's 200-event tail. + // instead of truncating it to the host's 200-event tail. Drain the streaming buffer + // first so the seq-dedup merge sees the full stream (a buffered event the snapshot + // also carries would otherwise be re-appended as a duplicate after the merge). + self.drainPendingOpenEvents() self.openEvents = Self.mergedEvents(self.openEvents, snap.recentEvents) self.openApprovals = snap.pendingApprovals self.persistOpenTranscript() } cb.openEvents = { [weak self] batch in guard let self, hostID == self.openSessionHostID, batch.sessionID == self.openSessionID else { return } - self.openEvents.append(contentsOf: batch.events) - self.persistOpenTranscript() + self.enqueueOpenEvents(batch.events) } cb.openBackfill = { [weak self] batch in guard let self, hostID == self.openSessionHostID, batch.sessionID == self.openSessionID else { return } // Full-history backfill precedes what's on screen — merge by seq so it slots in above the - // tail rather than appending out of order. + // tail rather than appending out of order. Drain first (same reason as openSnapshot). + self.drainPendingOpenEvents() self.openEvents = Self.mergedEvents(self.openEvents, batch.events) self.persistOpenTranscript() } @@ -722,41 +725,54 @@ final class RemoteStore: ObservableObject { /// or a representative one. No-op in demo, which seeds the aggregate directly. private func rebuildAggregate() { guard !demoMode else { return } + // Every assignment below goes through `setIfChanged`: `didUpdate` fires on *every* wire + // frame from *any* host (session churn, dashboard refresh, diff-stat ticks), and a plain + // `@Published` assignment fires `objectWillChange` even when the value is identical — + // re-evaluating every view observing the store for nothing. Equality checks over these + // small aggregates are far cheaper than a whole-tree SwiftUI invalidation. let live = aggregatedSessions() if !live.isEmpty { - sessions = live - cachedSummaries = live - persistSummaries(live) + if sessions != live { + sessions = live + cachedSummaries = live + persistSummaries(live) + } } else if connections.values.contains(where: { $0.connectivity.isLive }) { // Connected, but the host genuinely has no sessions — reflect that honestly. - sessions = [] + setIfChanged(\.sessions, []) } else { // Offline: keep showing the saved history rather than blanking the list. - sessions = cachedSummaries + setIfChanged(\.sessions, cachedSummaries) } - dashboard = DashboardSnapshot.merged(connections.values.map(\.dashboard)) - meshPeers = connections.values.flatMap(\.meshPeers) + setIfChanged(\.dashboard, DashboardSnapshot.merged(connections.values.map(\.dashboard))) + setIfChanged(\.meshPeers, connections.values.flatMap(\.meshPeers)) let ctx = contextConnection - hostName = ctx?.hostName ?? "" - capabilities = ctx?.capabilities - ?? WireCapabilities(canModifyToolInput: false, allowAlwaysScopes: []) - modelCatalog = ctx?.modelCatalog ?? .empty + setIfChanged(\.hostName, ctx?.hostName ?? "") + setIfChanged(\.capabilities, ctx?.capabilities + ?? WireCapabilities(canModifyToolInput: false, allowAlwaysScopes: [])) + setIfChanged(\.modelCatalog, ctx?.modelCatalog ?? .empty) if let id = openSessionHostID, let conn = connections[id] { // While a transcript is open, the composer/controls act on *that* Mac — its connectivity // gates send and its scope drives the control affordances. - connectivity = conn.connectivity - grantedScope = conn.grantedScope + setIfChanged(\.connectivity, conn.connectivity) + setIfChanged(\.grantedScope, conn.grantedScope) } else { - connectivity = aggregateConnectivity() + setIfChanged(\.connectivity, aggregateConnectivity()) // Optimistic new-chat gating: enabled if *any* Mac grants control (the owning Mac still // enforces scope when the intent lands there). - grantedScope = connections.values - .filter { $0.connectivity.isLive }.map(\.grantedScope).max() ?? .approve + setIfChanged(\.grantedScope, connections.values + .filter { $0.connectivity.isLive }.map(\.grantedScope).max() ?? .approve) } } + /// Assign a `@Published` property only when the value actually differs, so a no-op rebuild + /// doesn't fire `objectWillChange` (and with it a whole-tree view re-evaluation). + private func setIfChanged(_ keyPath: ReferenceWritableKeyPath, _ value: T) { + if self[keyPath: keyPath] != value { self[keyPath: keyPath] = value } + } + /// Merge transcript events by `seq` (monotonic, globally unique within a session), keeping the /// union sorted. Lets a reconnect's snapshot fold its events into the transcript already on /// screen without duplicating what's shown or dropping history outside the host's tail window. @@ -790,10 +806,67 @@ final class RemoteStore: ObservableObject { Task { [weak self] in let cached = await SessionCache.loadEvents(sessionID) guard let self, self.openSessionID == sessionID, !cached.isEmpty else { return } + // Drain any live deltas first so the seq-dedup merge sees the complete stream. + self.drainPendingOpenEvents() self.openEvents = Self.mergedEvents(cached, self.openEvents) } } + // 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 } @@ -910,6 +983,7 @@ final class RemoteStore: ObservableObject { openSessionID = nil compactDetailPresented = false openSessionHostID = nil + discardPendingOpenEvents() openEvents = []; openApprovals = []; openDiff = nil; diffLoading = false connectivity = .unpaired // No Mac left whose history to hold — drop the offline cache too. @@ -944,6 +1018,7 @@ final class RemoteStore: ObservableObject { openSessionID = nil compactDetailPresented = false openSessionHostID = nil + discardPendingOpenEvents() openEvents = [] openApprovals = [] openDiff = nil @@ -972,6 +1047,7 @@ final class RemoteStore: ObservableObject { // Record which Mac owns this session (mesh P3) so its connection forwards the transcript and // the host-specific projected values follow it. openSessionHostID = connection(owningSession: sessionID)?.hostID + discardPendingOpenEvents() openEvents = [] openApprovals = [] openDiff = nil @@ -1149,6 +1225,8 @@ final class RemoteStore: ObservableObject { compactDetailPresented = false openSessionHostID = nil // Persist the final transcript before clearing it, so it's warm for the next open / offline. + // Buffered streaming deltas are part of that transcript — publish them first. + drainPendingOpenEvents() flushOpenTranscript(target, openEvents) openEvents = [] openApprovals = [] diff --git a/NucleicRemote/NucleicRemote/Views/SessionDetailView.swift b/NucleicRemote/NucleicRemote/Views/SessionDetailView.swift index 343ca8a..fba176f 100644 --- a/NucleicRemote/NucleicRemote/Views/SessionDetailView.swift +++ b/NucleicRemote/NucleicRemote/Views/SessionDetailView.swift @@ -123,17 +123,37 @@ struct SessionDetailView: View { set: { store.setSessionEffort(sessionID, $0) }) } + /// Memo for the context-occupancy scan below. The backward scan usually stops at the last + /// turn's usage event, but a stream with sparse (or no) usage reporting walks the whole + /// transcript — and `body` re-evaluates on every store change and keystroke, so an O(N) scan + /// per evaluation quietly compounds on long sessions. `(count, lastSeq)` pins the stream, as + /// in the projection cache; a class held in `@State` so updating it from `body` doesn't + /// itself invalidate the view. + private final class ContextScanCache { + var count = -1 + var lastSeq: UInt64 = 0 + var used: Int? + } + @State private var contextScanCache = ContextScanCache() + /// Live context-window occupancy (newest turn's input tokens ÷ the model's window), read - /// from the transcript exactly as the Mac header does. + /// from the transcript exactly as the Mac header does. Only the token scan is memoized — + /// the window division stays live, so a model switch reflects immediately. private var contextPercent: Int? { - let used = store.openEvents.reversed().lazy.compactMap { event -> Int? in - switch event.kind { - case .turnCompleted(let turn): return turn.usage?.contextInputTokens - case .usage(let usage): return usage.contextInputTokens - default: return nil - } - }.first { $0 > 0 } - guard let used else { return nil } + let events = store.openEvents + let cache = contextScanCache + if cache.count != events.count || cache.lastSeq != (events.last?.seq ?? 0) { + cache.count = events.count + cache.lastSeq = events.last?.seq ?? 0 + cache.used = events.reversed().lazy.compactMap { event -> Int? in + switch event.kind { + case .turnCompleted(let turn): return turn.usage?.contextInputTokens + case .usage(let usage): return usage.contextInputTokens + default: return nil + } + }.first { $0 > 0 } + } + guard let used = cache.used else { return nil } let window = store.modelCatalog.contextWindow(summary?.model) guard window > 0 else { return nil } return min(100, Int((Double(used) / Double(window)) * 100)) diff --git a/NucleicRemote/NucleicRemote/Views/SessionsView.swift b/NucleicRemote/NucleicRemote/Views/SessionsView.swift index cd4be2a..9b6ef9d 100644 --- a/NucleicRemote/NucleicRemote/Views/SessionsView.swift +++ b/NucleicRemote/NucleicRemote/Views/SessionsView.swift @@ -11,9 +11,15 @@ struct SessionsView: View { private var grouped: [(title: String, rows: [WireSessionSummary])] { let pool = (showArchived ? store.sessions : store.liveSessions) .sorted(by: StatusStyle.attentionThenRecency) - let needs = pool.filter { $0.status.needsYou($0.disposition) } - let running = pool.filter { $0.status == .running || $0.status == .provisioning || $0.status == .idle } - let done = pool.filter { !needs.contains($0) && !running.contains($0) } + // One pass, first bucket wins — the old `done = pool.filter { !needs.contains($0) … }` + // ran O(rows²) full-summary equality scans on every body evaluation. + var needs: [WireSessionSummary] = [], running: [WireSessionSummary] = [], done: [WireSessionSummary] = [] + for summary in pool { + if summary.status.needsYou(summary.disposition) { needs.append(summary) } + else if summary.status == .running || summary.status == .provisioning || summary.status == .idle { + running.append(summary) + } else { done.append(summary) } + } return [("Needs you", needs), ("Running", running), ("Done", done)].filter { !$0.rows.isEmpty } }