package com.example.aiapp import java.io.IOException import java.net.HttpURLConnection import java.net.URL /** * The SSE half of the API: one long-lived GET per open session screen, replaying the transcript * after a cursor and then following it live. * * Blocking -- run() occupies its thread until the stream ends. [close] (from any thread) is the * cancellation path: it disconnects the socket, which unblocks the read; run() then returns instead * of throwing, so a deliberate close doesn't surface as a connection error. The caller owns * reconnecting (with the last seq it saw as the new cursor) -- see SessionScreen. */ /** * The frame name the server uses to say a cursor was too far behind to continue from. Must match * `send_backlog` in the backend's routes.rs. */ private const val RESET_EVENT = "reset" class EventStream(private val settings: ServerSettings, private val sessionId: String) { @Volatile private var connection: HttpURLConnection? = null @Volatile private var closed = false fun close() { closed = true connection?.disconnect() } /** * Streams events after [after] into [onEvent] until the stream drops. * * [onOpen] fires once the server has accepted the connection. That is the measured moment the * stream is live again, and the only honest thing to clear a previous failure on: an earlier * version cleared on the first event instead, so an idle session went on displaying a * connection error that had already been recovered from, indefinitely. * * [onReset] fires when the server answers that the cursor is too far behind to continue from: * everything already displayed is stale and the events that follow are a fresh window, so the * caller drops what it holds and rebuilds -- the same thing it does when the screen opens. It * arrives before those events, so a caller that clears on it stays in order. */ fun run(after: Long, onOpen: () -> Unit, onReset: () -> Unit, onEvent: (SeqEvent) -> Unit) { val connection = URL("${settings.baseUrl}/sessions/$sessionId/events?after=$after").openConnection() as HttpURLConnection this.connection = connection try { connection.applyPinnedTls() connection.connectTimeout = CONNECT_TIMEOUT_MS // No read timeout: between events there is nothing to read for // as long as the session is idle; the server's keep-alives and // a dead socket erroring out are the liveness story. connection.readTimeout = 0 connection.setRequestProperty("Authorization", "Bearer ${settings.token}") connection.setRequestProperty("Accept", "text/event-stream") if (connection.responseCode != 200) { val detail = connection.errorStream?.bufferedReader()?.readText()?.trim() throw ApiException(detail ?: "HTTP ${connection.responseCode} for the event stream") } onOpen() val reader = connection.inputStream.bufferedReader() // SSE framing: `data:` and `event:` lines accumulate until a // blank line ends the frame. `id:` (the seq) is also inside the // JSON payload, so it needs no separate handling; comment lines // (keep-alives) start with ':' and are skipped. val data = StringBuilder() var name: String? = null while (true) { val line = reader.readLine() ?: break when { line.isEmpty() -> { // A named frame carries no payload and a data frame // has no name, so this is one or the other. if (name == RESET_EVENT) onReset() else if (data.isNotEmpty()) onEvent(parseSeqEvent(data.toString())) data.clear() name = null } line.startsWith("data:") -> data.append(line.removePrefix("data:").trim()) line.startsWith("event:") -> name = line.removePrefix("event:").trim() else -> {} // id:, comments -- nothing to do } } } catch (e: ApiException) { throw e } catch (e: IOException) { if (!closed) { throw ApiException( "Can't reach the server -- retrying. (${e.message ?: e::class.simpleName})", e, ) } } finally { connection.disconnect() this.connection = null } } }