package com.example.aiapp import android.content.Context import java.io.File import java.util.concurrent.atomic.AtomicReference /** * Where the session screen gets a transcript from: this phone's copy first, the server for the * rest. * * One seam rather than a cache the screen has to remember to consult. Everything it fetched before * is asked of this, and everything the server sends is written into the cache on the way past, so * the screen never learns which side answered. What it does learn, through [DebugStats], is how * often each one did. * * See TRANSCRIPT_CACHE.md. The one rule worth keeping in mind: the cache is never load-bearing. * Every read has a network path beside it producing the same result. */ class TranscriptSource( private val settings: ServerSettings, private val address: TranscriptAddress, val cache: SessionCache, ) { private val stream = AtomicReference(null) /** * The cached opening window, or null when there is nothing usable to draw. * * Drawn *before* [probe] returns, which is the whole point of the feature: the rows are on * screen while the check that they are still the server's rows is in flight, and a failed check * replaces them exactly as a `reset` does. */ fun cachedOpening(limit: Int = OPENING_WINDOW): List? { if (cache.tail() == null) return null val lines = cache.newest(limit) if (lines.isEmpty()) return null return try { lines.map { parseSeqEvent(it) } } catch (e: org.json.JSONException) { // Lines this build cannot read at all, which the cache's own checks cannot see: it // reads a seq off a line, not an event. Nothing to serve, so a cold open. cache.purge() null } } /** * Whether the server's event at the cached cursor is still the cached one. * * The screen must not resume a stream from a cached seq unless it is the same conversation. A * transcript is append-only in ordinary use, but the file can be replaced or truncated -- a * sandbox re-seeded with the same ids, a backup restored, a session re-imported -- and the * server's catch-up on such a file would hand this phone a continuation of a *different* * conversation, spliced onto the cached one with no seam. Caught with one request of a few * hundred bytes, in the slot the opening page's request used to be in. * * False purges the cache and means "open cold". A throw is the server not being askable, which * is neither: the cached rows stay on screen and the caller tries again on the reconnect * schedule. * * What this cannot see is a line changed in the middle of the file with the tail intact. That * is what the Reload button in session settings is for. */ suspend fun probe(): Boolean { val tail = cache.tail() ?: return false // `before = seq + 1` is the newest event with seq <= the cursor, which is the event *at* // the cursor when the server still has one there. val answer = fetchTranscript(settings, address, before = tail.seq + 1, limit = 1) val matches = answer.size == 1 && try { answer[0].second == parseSeqEvent(tail.line) } catch (e: org.json.JSONException) { false } if (!matches) cache.purge() return matches } /** * Today's opening fetch, kept as the start of the live run. Only called when the cache has * nothing to open with, or when [probe] said what it had was not the server's. */ suspend fun fetchOpening(): List { DebugStats.count("transcript page from server") val page = fetchTranscript(settings, address, limit = OPENING_WINDOW) page.forEach { (line, entry) -> cache.append(line, entry.seq) } cache.flush() return page.map { it.second } } /** * The page before [before]: from the cache when it holds it, otherwise from the server bounded * by what the cache already has. * * The bound is what keeps the cache worth having. A coalesced page reaches back as far as its * row count takes it -- a single reply is hundreds of lines -- so a page fetched after the * reader has been away would run straight past the cached run and overlap it, and an * overlapping page cannot be stored. Told where this phone's copy starts, the server stops * there instead. */ suspend fun page(before: Long, limit: Int, coalesce: Boolean): List { cache.page(before, limit, rows = coalesce)?.let { lines -> DebugStats.count("transcript page from cache") return lines.map { parseSeqEvent(it) } } DebugStats.count("transcript page from server") val page = fetchTranscript( settings, address, before = before, limit = limit, coalesce = coalesce, after = cache.coveredUpTo(before)?.minus(1), ) if (page.isNotEmpty()) { // `before` rather than the newest line's seq: a coalesced page covers everything up to // the cursor it was asked with, and nothing in its lines says so. cache.storePage(page.map { it.first }, page.first().second.seq, before, rows = coalesce) } return page.map { it.second } } /** * [EventStream.run], with every frame written to the cache before [onEvent] sees it. * * Before, so that an event held back for a reader who is scrolled away is already on disk -- * what the cache holds is what the server sent, not what the screen has got round to drawing. * Flushed on each status change, which is a turn's boundary and the granularity a crash may as * well lose. */ fun follow(after: Long, onOpen: () -> Unit, onReset: () -> Unit, onEvent: (SeqEvent) -> Unit) { val opened = EventStream(settings, address) stream.getAndSet(opened)?.close() try { opened.run(after, onOpen, onReset) { raw, entry -> cache.append(raw, entry.seq) if (entry.event is SessionEvent.Status) cache.flush() onEvent(entry) } } finally { cache.flush() } } /** Ends the stream, from any thread, and leaves the cache with everything it was given. */ fun close() { stream.getAndSet(null)?.close() cache.flush() } } /** * How many events the screen opens with, cached or fetched. * * The server's own default for a page, named here because the cached opening has to be the same * size as the fetched one -- a reader must not get a shorter first screen for having been here * before. */ private const val OPENING_WINDOW = 80 /** * Where this server's cached transcripts live. * * Under `cacheDir` because that is exactly what it is for: bytes the phone can regenerate from the * server, which Android may delete under storage pressure without asking. Keyed by host and port * because two servers can hold a session with the same id, and a line from one shown against the * other is the whole invariant broken. `v1` is the layout's version. */ fun cacheRoot(context: Context, settings: ServerSettings): File { val transcripts = File(context.cacheDir, "transcripts") transcripts.listFiles()?.forEach { if (it.name != CACHE_VERSION) it.deleteRecursively() } return File(transcripts, "$CACHE_VERSION/${settings.host}_${settings.port}") } private const val CACHE_VERSION = "v1"