package com.example.aiapp import java.io.IOException import java.net.HttpURLConnection import java.net.URL /** * How long to wait before opening a dropped stream again. * * Shared by every screen that follows one, so a reconnect is not paced differently depending on * which stream dropped. Short enough that a tunnel coming back is not noticed, long enough that a * server which is genuinely down is not being asked several times a second. */ const val RECONNECT_DELAY_MS = 1500L /** * One server-sent-events connection, framed. * * The framing is the part worth having once: `data:` and `event:` lines accumulate until a blank * line ends the frame, comments start with `:`, and a frame is either named with no payload or a * payload with no name. Two screens follow two different streams and neither should re-derive that. * * 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, and [run] then returns * rather than throwing. Reconnecting belongs to the caller, which is the only one that knows where * to resume from. */ class Sse(private val settings: ServerSettings) { @Volatile private var connection: HttpURLConnection? = null @Volatile private var closed = false fun close() { closed = true connection?.disconnect() } /** * Follows the stream at [path], handing each frame to [onFrame] as its name (null for an * ordinary data frame) and its payload. The path is given here rather than at construction * because a caller that reconnects usually resumes from somewhere new. * * [onOpen] fires once the server has accepted the connection. That is the measured moment the * stream is live, and the only honest thing to clear a previous failure on: clearing on the * first *event* instead left an idle stream displaying an error it had already recovered from. */ fun run(path: String, onOpen: () -> Unit, onFrame: (name: String?, data: String) -> Unit) { // Opening is inside the try, not before it. Everything this method can fail at owes the // caller the same kind of failure, and a connection that could not even be constructed used // to escape as a raw `IOException` from a line no `catch` covered. var connection: HttpURLConnection? = null try { connection = (URL("${settings.baseUrl}$path").openConnection() as HttpURLConnection).also { this.connection = it } connection.applyPinnedTls() connection.connectTimeout = CONNECT_TIMEOUT_MS // No read timeout: between events there is nothing to read for as long as the thing // being followed 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() val data = StringBuilder() var name: String? = null while (true) { val line = reader.readLine() ?: break when { line.isEmpty() -> { if (name != null || data.isNotEmpty()) onFrame(name, 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})", cause = e, ) } } finally { connection?.disconnect() this.connection = null } } }