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. */ 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. */ fun run(after: Long, 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") } val reader = connection.inputStream.bufferedReader() // SSE framing: `data:` lines accumulate until a blank line ends // the event. `id:` (the seq) is also inside the JSON payload, // so only data lines matter; comment lines (keep-alives) start // with ':' and are skipped. val data = StringBuilder() while (true) { val line = reader.readLine() ?: break when { line.isEmpty() -> { if (data.isNotEmpty()) { onEvent(parseSeqEvent(data.toString())) data.clear() } } line.startsWith("data:") -> data.append(line.removePrefix("data:").trim()) else -> {} // id:, event:, comments -- nothing to do } } } catch (e: ApiException) { throw e } catch (e: IOException) { if (!closed) { throw ApiException( "Lost the event stream (${e::class.simpleName}: ${e.message})", e, ) } } finally { connection.disconnect() this.connection = null } } }