package com.example.aiapp /** * 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" /** * 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. * * The connection and its framing belong to [Sse]; what stays here is what this stream's frames * mean. [close] from any thread ends it, and the caller owns reconnecting -- with the last seq it * saw as the new cursor. See SessionScreen. */ class EventStream(settings: ServerSettings, private val sessionId: String) { private val stream = Sse(settings) fun close() = stream.close() /** * Streams events after [after] into [onEvent] until the stream drops. * * [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) { stream.run("/sessions/$sessionId/events?after=$after", onOpen) { name, data -> // 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)) } } }