Transcript Incremental Projection
Nucleic-Session: 0FFD007B-0696-4517-9429-129C7B0FD5AC Co-authored-by: Nucleic <[email protected]>
This commit is contained in:
@@ -812,6 +812,69 @@ final class RemoteStore: ObservableObject {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// MARK: - Background transcript prefetch
|
||||||
|
|
||||||
|
/// Warm the offline transcript cache for recent sessions *before* they're ever opened. A
|
||||||
|
/// fresh open otherwise starts from an empty transcript and waits on the host's snapshot, so
|
||||||
|
/// the view renders at the top and visibly drops to the end as history lands; a session
|
||||||
|
/// opened before seeds instantly from `SessionCache` and opens at its end. Prefetching runs
|
||||||
|
/// the same `fetchTranscript` flow the open path uses — one session at a time, newest first,
|
||||||
|
/// pulling only what the cache is missing — so every recent session opens like a warm one.
|
||||||
|
private var prefetchQueue: [SessionID] = []
|
||||||
|
private var prefetchInFlight: SessionID?
|
||||||
|
/// The summary `lastSeq` each session was last prefetched (or attempted) at, so a session is
|
||||||
|
/// re-queued only after it has moved on — not on every aggregate rebuild.
|
||||||
|
private var prefetchedSeq: [SessionID: UInt64] = [:]
|
||||||
|
|
||||||
|
/// Rebuild the prefetch queue from the freshest summaries. Called when the session list
|
||||||
|
/// actually changes (connect, session churn); cheap when nothing needs pulling.
|
||||||
|
private func scheduleTranscriptPrefetch() {
|
||||||
|
guard !demoMode else { return }
|
||||||
|
prefetchQueue = sessions
|
||||||
|
.filter { !$0.archived && $0.sessionID != openSessionID && $0.sessionID != prefetchInFlight }
|
||||||
|
.sorted { $0.updatedAt > $1.updatedAt }
|
||||||
|
.prefix(8)
|
||||||
|
.filter { $0.lastSeq > (prefetchedSeq[$0.sessionID] ?? 0) }
|
||||||
|
.map(\.sessionID)
|
||||||
|
pumpTranscriptPrefetch()
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Start the next queued prefetch if none is in flight. Serial on purpose — background work
|
||||||
|
/// must trickle behind the live stream, not contend with it.
|
||||||
|
private func pumpTranscriptPrefetch() {
|
||||||
|
guard prefetchInFlight == nil, !prefetchQueue.isEmpty else { return }
|
||||||
|
let id = prefetchQueue.removeFirst()
|
||||||
|
guard let summary = sessions.first(where: { $0.sessionID == id }),
|
||||||
|
let conn = connection(owningSession: id), conn.canPrefetchTranscript
|
||||||
|
else {
|
||||||
|
// Gone / connection busy or incapable — skip; a later list update re-queues it.
|
||||||
|
pumpTranscriptPrefetch()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
prefetchInFlight = id
|
||||||
|
// Mark the attempt at this watermark now, so an empty/unavailable result doesn't
|
||||||
|
// re-queue in a loop; the session re-qualifies once its lastSeq advances.
|
||||||
|
prefetchedSeq[id] = summary.lastSeq
|
||||||
|
Task { [weak self] in
|
||||||
|
// Pull only the gap beyond what's cached; a cold session is bounded to the same
|
||||||
|
// tail window the cache would keep anyway.
|
||||||
|
let cachedLast = await SessionCache.loadEvents(id).last?.seq ?? 0
|
||||||
|
guard let self else { return }
|
||||||
|
if cachedLast >= summary.lastSeq {
|
||||||
|
self.prefetchInFlight = nil // cache already current
|
||||||
|
self.pumpTranscriptPrefetch()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
let coldFloor = summary.lastSeq > UInt64(SessionCache.eventLimit)
|
||||||
|
? summary.lastSeq - UInt64(SessionCache.eventLimit) : 0
|
||||||
|
let afterSeq = cachedLast > 0 ? cachedLast : coldFloor
|
||||||
|
if !conn.prefetchTranscript(id, afterSeq: afterSeq) {
|
||||||
|
self.prefetchInFlight = nil // couldn't start (connection changed) — move on
|
||||||
|
self.pumpTranscriptPrefetch()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// MARK: - Streaming delta coalescing
|
// MARK: - Streaming delta coalescing
|
||||||
|
|
||||||
/// Streaming transcript deltas arrive one wire frame at a time — often dozens per second
|
/// Streaming transcript deltas arrive one wire frame at a time — often dozens per second
|
||||||
|
|||||||
Reference in New Issue
Block a user