diff --git a/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt b/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt index 1bef195..c09811b 100644 --- a/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt +++ b/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt @@ -27,6 +27,7 @@ import iris.util.IrisLog import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Job +import kotlinx.coroutines.async import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.currentCoroutineContext import kotlinx.coroutines.delay @@ -72,14 +73,33 @@ class GatewayClient( data object Connecting : State + /** Common surface of [Connected] and [HttpFallback]: both carry the + * hello.ack data (caps, channels, push watermark). */ + interface HelloInfo { + val caps: ServerCaps + val channels: List + val lastPushedCursor: Long + } + data class Connected( - val caps: ServerCaps, - val channels: List, + override val caps: ServerCaps, + override val channels: List, /** M5: highest outbox cursor already pushed to this device * (from hello.ack; 0 = never). Sync-replayed frames at/below * it must not re-post system notifications (docs/08 §8.7). */ - val lastPushedCursor: Long = 0, - ) : State + override val lastPushedCursor: Long = 0, + ) : State, + HelloInfo + + /** docs/19: the WS is down but the gateway is reachable over the + * HTTP leg — sendable (POST /v1/frame) + receiving (SSE/long-poll). + * Media is unavailable until the WS is back. */ + data class HttpFallback( + override val caps: ServerCaps, + override val channels: List, + override val lastPushedCursor: Long = 0, + ) : State, + HelloInfo data object Reconnecting : State @@ -116,6 +136,18 @@ class GatewayClient( private var lastLiveness: TimeMark = TimeSource.Monotonic.markNow() private val pending = mutableMapOf>() + // docs/19: HTTP fallback leg (the "HTTP leg"). [http] is created lazily + // from the stored WS URL; [httpJob] runs the SSE/long-poll receive loop; + // [httpCursor] is the resume cursor (SSE id / outbox high-water mark). + private var http: HttpGateway? = null + private var httpJob: Job? = null + private var httpCursor: Long = 0 + private var sseFailures = 0 + private var usingLongPoll = false + // Last hello.ack payload (WS or SSE) — used to restore State.Connected + // after a fallback/WS state race in the connect loop. + private var lastAck: HelloAckPayload? = null + // M4: binary frames (media upload chunks / pull stream) have no per-frame // id, so at most one binary session is active per socket. The gateway // allows one upload per connection; pull is request/response. @@ -136,7 +168,7 @@ class GatewayClient( * (Dispatchers.Default) and would push a history request past a flaky * network's window. Set before [start]. */ - var onHelloAck: ((State.Connected) -> Unit)? = null + var onHelloAck: ((State.HelloInfo) -> Unit)? = null // M4: only one pull may be in flight at a time (binarySession is a single // slot). Serialize concurrent offers so their byte streams don't interleave. @@ -156,6 +188,7 @@ class GatewayClient( fun stop() { connectJob?.cancel() connectJob = null + stopHttpLeg() socket?.close(1000, "client shutdown") socket = null _state.value = State.Disconnected @@ -176,9 +209,20 @@ class GatewayClient( return } _state.value = if (hasConnected) State.Reconnecting else State.Connecting - val dial = dial(url, token) + // docs/19: race the WS dial against the HTTP health probe. If the + // gateway is alive over HTTP, the app can send immediately + // (fallback) without waiting out the WS dial timeout — the key + // UX fix (sendable in < 1 s on a dead WS port). + val dialDeferred = scope.async { dial(url, token) } + val healthOk = + withTimeoutOrNull(2_000) { + httpHealthy() + } ?: false + if (healthOk) enterHttpFallback() + val dial = dialDeferred.await() when (val result = dial.result) { is DialResult.AuthFailed -> { + stopHttpLeg() _state.value = State.AuthFailed(result.message) dial.socket.close(1000, "auth failed") return @@ -187,10 +231,20 @@ class GatewayClient( DialResult.Connected -> { hasConnected = true attempt = 0 + // WS is up again: back to WS-only (media available). + stopHttpLeg() lastLiveness = TimeSource.Monotonic.markNow() + // Restore Connected if the fallback state won the race + // (health landed before hello.ack). + lastAck?.let { + _state.value = + State.Connected(it.serverCaps, it.channels, it.lastPushedCursor) + } dial.closed.await() if (!currentCoroutineContext().isActive) return - // socket dropped -> loop again (Reconnecting) + // WS dropped: fall back to HTTP immediately (no backoff + // gate on the send path), then redial below. + enterHttpFallback() } is DialResult.Failed -> { @@ -201,6 +255,114 @@ class GatewayClient( } } + // ── docs/19: HTTP fallback leg ──────────────────────────────────────── + + /** Lazily build the HTTP client from the stored WS URL. */ + private fun httpGateway(): HttpGateway? { + val url = store.serverUrl.trim() + val token = store.token + if (url.isBlank() || token.isBlank()) return null + return http + ?: HttpGateway(client, HttpGateway.deriveHttpUrl(url), token, store.deviceId) + .also { http = it } + } + + private suspend fun httpHealthy(): Boolean = + try { + httpGateway()?.health() ?: false + } catch (e: Exception) { + false + } + + /** + * Enter [State.HttpFallback]: open the SSE (or long-poll) receive loop + * from the saved sync cursor. Idempotent; a no-op while WS-connected. + */ + private fun enterHttpFallback() { + if (_state.value is State.Connected) return + val gw = httpGateway() ?: return + sseFailures = 0 + usingLongPoll = false + httpCursor = store.syncCursor + // Provisional state (previous caps/channels) until the SSE hello + // arrives with the real ones. + val prev = _state.value + _state.value = + State.HttpFallback( + caps = (prev as? State.HelloInfo)?.caps ?: ServerCaps(), + channels = (prev as? State.HelloInfo)?.channels ?: emptyList(), + lastPushedCursor = (prev as? State.HelloInfo)?.lastPushedCursor ?: 0, + ) + httpJob?.cancel() + httpJob = scope.launch { httpReceiveLoop(gw) } + } + + private fun stopHttpLeg() { + httpJob?.cancel() + httpJob = null + sseFailures = 0 + usingLongPoll = false + } + + /** + * The HTTP receive loop: SSE by default; after two consecutive SSE open + * failures (buffering proxy) it switches to long-poll until the next + * full (re)connect (docs/19 §19.6). + */ + private suspend fun httpReceiveLoop(gw: HttpGateway) { + var backoff = 1_000L + while (currentCoroutineContext().isActive) { + if (usingLongPoll) { + try { + val res = gw.poll(httpCursor) + res.frames.forEach { emitHttpFrame(it) } + if (res.cursor > httpCursor) httpCursor = res.cursor + } catch (e: Exception) { + IrisLog.w("http poll failed: ${e.message}") + delay(backoff) + backoff = minOf(backoff * 2, 15_000) + } + } else { + try { + gw.events( + cursor = httpCursor, + onHello = { onHttpHello(it) }, + onFrame = { emitHttpFrame(it) }, + onCursor = { if (it > httpCursor) httpCursor = it }, + ) + // Clean EOF: reconnect immediately. + backoff = 1_000L + } catch (e: Exception) { + sseFailures++ + if (sseFailures >= 2) { + // SSE seems blocked: switch to long-poll. + usingLongPoll = true + continue + } + IrisLog.w("sse read failed: ${e.message}") + delay(backoff) + backoff = minOf(backoff * 2, 15_000) + } + } + } + } + + /** The SSE `event: hello` (the HTTP hello.ack). */ + private fun onHttpHello(ack: HelloAckPayload) { + lastAck = ack + val fb = State.HttpFallback(ack.serverCaps, ack.channels, ack.lastPushedCursor) + _state.value = fb + onHelloAck?.invoke(fb) + } + + /** Deliver an HTTP-leg frame to the same sinks as a WS frame. */ + private fun emitHttpFrame(frame: Frame) { + _events.tryEmit(frame) + frame.id?.let { id -> + pending[id]?.complete(frame) + } + } + // ── Dial (one connect + hello) ──────────────────────────────────────── private sealed interface DialResult { @@ -338,6 +500,7 @@ class GatewayClient( helloAck.invokeOnCompletion { e -> if (e == null) { val ack = helloAck.getCompleted() + lastAck = ack val connected = State.Connected(ack.serverCaps, ack.channels, ack.lastPushedCursor) _state.value = connected // M5: reconnect catch-up — replay frames parked while offline. @@ -376,9 +539,23 @@ class GatewayClient( mediaRefs: List = emptyList(), autoThread: Boolean = false, ) { - val ws = socket ?: return - val id = nextRequestId++ - ws.send(messageSendFrame(id, chatId, text, threadId, mediaRefs, autoThread).toWire()) + val ws = socket + if (ws != null && _state.value is State.Connected) { + val id = nextRequestId++ + ws.send(messageSendFrame(id, chatId, text, threadId, mediaRefs, autoThread).toWire()) + return + } + // docs/19: WS down — route over the HTTP leg. Media is WS-only in + // v1 (uploads need the live connection), so mediaRefs are dropped in + // fallback (the UI disables the attach button in that state). + if (_state.value is State.HttpFallback) { + val id = nextRequestId++ + scope.launch { + httpGateway()?.postFrame( + messageSendFrame(id, chatId, text, threadId, emptyList(), autoThread), + ) + } + } } // ── M4: media upload / pull ─────────────────────────────────────────── @@ -510,10 +687,21 @@ class GatewayClient( * app reconciles from [events]. Returns the id used, or -1 if not connected. */ fun sendFrame(frame: Frame): Int { - val ws = socket ?: return -1 val id = nextRequestId++ - ws.send(frame.copy(id = id).toWire()) - return id + val ws = socket + if (ws != null && _state.value is State.Connected) { + ws.send(frame.copy(id = id).toWire()) + return id + } + // docs/19: WS down — route over the HTTP leg; the response (same id) + // arrives on the SSE/long-poll stream via [events]. + if (_state.value is State.HttpFallback) { + scope.launch { + httpGateway()?.postFrame(frame.copy(id = id)) + } + return id + } + return -1 } /** Send a ping (heartbeat). */ diff --git a/app/shared/src/commonMain/kotlin/iris/net/HttpGateway.kt b/app/shared/src/commonMain/kotlin/iris/net/HttpGateway.kt new file mode 100644 index 0000000..2fd6e87 --- /dev/null +++ b/app/shared/src/commonMain/kotlin/iris/net/HttpGateway.kt @@ -0,0 +1,299 @@ +package iris.net + +import iris.protocol.Frame +import iris.protocol.HelloAckPayload +import iris.protocol.IrisJson +import iris.util.IrisLog +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.withContext +import kotlinx.serialization.json.jsonArray +import kotlinx.serialization.json.jsonObject +import kotlinx.serialization.json.jsonPrimitive +import okhttp3.Headers +import okhttp3.MediaType.Companion.toMediaType +import okhttp3.OkHttpClient +import okhttp3.Request +import okhttp3.RequestBody.Companion.toRequestBody +import java.io.IOException +import java.util.concurrent.TimeUnit + +/** + * HTTP fallback transport client (docs/19): the "HTTP leg". + * + * When the WS is down (flaky network, NAT timeout, app just relaunched), + * the app sends over `POST /v1/frame` and receives over SSE + * `GET /v1/events` (or long-poll `GET /v1/poll` where SSE is blocked). + * Same frames, same outbox cursor, same token as the WS. + * + * [events] is ONE SSE connection attempt (blocking read on + * [Dispatchers.IO]); [GatewayClient] wraps it in a retry loop and tracks + * the resume cursor via the [events] `onCursor` callback (SSE `id` = + * outbox cursor, so resume is just the last seen id). + */ +class HttpGateway( + private val client: OkHttpClient, + private val baseUrl: String, + private val token: String, + private val deviceId: String, +) { + // OkHttp's default read timeout (10 s) is shorter than the gateway's SSE + // heartbeat (15 s) and the long-poll hold (25 s) — per-purpose clients + // with extended call timeouts (see the *Client() helpers below). + private val healthClient: OkHttpClient = client.healthClient() + private val streamClient: OkHttpClient = client.streamClient() + private val pollClient: OkHttpClient = client.pollClient() + /** POST /v1/frame result. [frame] is the handler's synchronous reply + * (error frame on 4xx, e.g. read.receipt on 200) or null for a plain + * 202 accept-and-ack. */ + data class PostResult( + val ok: Boolean, + val status: Int, + val frame: Frame?, + ) + + /** Long-poll result: new high-water [cursor] + frames with cursor > + * the requested one (may be empty at timeout). */ + data class PollResult( + val cursor: Long, + val frames: List, + ) + + companion object { + val JSON = "application/json".toMediaType() + + /** Default port of the gateway's HTTP leg (WS default is 8790). */ + const val DEFAULT_PORT = 8791 + + /** + * Derive the HTTP base URL from the stored WS URL (docs/19 §19.4): + * `ws(s)://host[:port]/ws` -> `http(s)://host:8791`. The WS port is + * a different service, so the port is always replaced with the + * HTTP leg's default. Pure function (unit-tested). + */ + fun deriveHttpUrl(wsUrl: String): String { + val u = wsUrl.trim() + val (scheme, rest) = + when { + u.startsWith("wss://") -> "https" to u.removePrefix("wss://") + u.startsWith("ws://") -> "http" to u.removePrefix("ws://") + else -> return u // already http(s) + } + val authority = rest.substringBefore('/') + val host = authority.substringBefore(':') + return "$scheme://$host:$DEFAULT_PORT" + } + } + + private fun authHeaders(): Headers = + Headers + .Builder() + .add("Authorization", "Bearer $token") + .add("X-Iris-Device", deviceId) + .build() + + /** Liveness probe (unauthenticated by design). True on 200. */ + suspend fun health(): Boolean = + withContext(Dispatchers.IO) { + val request = + Request + .Builder() + .url("$baseUrl/v1/health") + .build() + healthClient + .newCall(request) + .execute() + .use { response -> + response.body?.close() + response.code == 200 + } + } + + /** + * POST /v1/frame (accept-and-ack, docs/19 §19.7). 2xx -> [PostResult.ok] + * (with the synchronous reply frame when the handler sent one); 4xx -> + * the error frame as the body. + */ + suspend fun postFrame(frame: Frame): PostResult = + withContext(Dispatchers.IO) { + val wire = IrisJson.instance.encodeToString(Frame.serializer(), frame) + val request = + Request + .Builder() + .url("$baseUrl/v1/frame") + .headers(authHeaders()) + .post(wire.toRequestBody(JSON)) + .build() + client + .newCall(request) + .execute() + .use { response -> + val body = response.body?.string().orEmpty() + val parsed = + try { + if (body.startsWith("{")) { + val obj = IrisJson.instance.parseToJsonElement(body) + // 202 {"ok":true} is not a frame; 4xx/200 bodies are. + if (obj.jsonObject.containsKey("type")) { + IrisJson.instance.decodeFromJsonElement(Frame.serializer(), obj) + } else { + null + } + } else { + null + } + } catch (e: Exception) { + null + } + PostResult(response.isSuccessful, response.code, parsed) + } + } + + /** + * One SSE connection attempt: outbox catch-up from [cursor], then live + * frames. [onHello] fires for `event: hello` (the HTTP hello.ack); + * [onFrame] for `event: frame`; [onCursor] with the SSE `id` (outbox + * cursor) when present. Returns on clean EOF; throws [IOException] on + * open/read failure. Callbacks run on the IO thread. + */ + suspend fun events( + cursor: Long, + onHello: (HelloAckPayload) -> Unit, + onFrame: (Frame) -> Unit, + onCursor: (Long) -> Unit, + ) { + withContext(Dispatchers.IO) { + val request = + Request + .Builder() + .url("$baseUrl/v1/events?cursor=$cursor") + .headers(authHeaders()) + .build() + client + .newCall(request) + .execute() + .use { response -> + if (!response.isSuccessful) { + throw IOException("SSE open failed: HTTP ${response.code}") + } + val source = response.body?.source() ?: throw IOException("empty SSE body") + var eventId: String? = null + val dataLines = mutableListOf() + while (true) { + val line = source.readUtf8Line() ?: break // EOF + when { + line.isEmpty() -> { + if (dataLines.isNotEmpty()) { + val data = dataLines.joinToString("\n") + try { + val frame = + IrisJson.instance.decodeFromString( + Frame.serializer(), + data, + ) + when (eventId) { + "hello" -> { + onHello( + frame.payloadAs() + ?: HelloAckPayload(), + ) + } + + else -> { + onFrame(frame) + } + } + } catch (e: Exception) { + IrisLog.e("sse frame decode failed: $e :: ${data.take(120)}") + } + } + eventId = null + dataLines.clear() + } + + line.startsWith(":") -> { + Unit + } + + // heartbeat comment + line.startsWith("id:") -> { + line + .removePrefix("id:") + .trim() + .toLongOrNull() + ?.let(onCursor) + } + + line.startsWith("event:") -> { + eventId = line.removePrefix("event:").trim() + } + + line.startsWith("data:") -> { + dataLines.add(line.removePrefix("data:").removePrefix(" ")) + } + } + } + } + } + } + + /** + * Long-poll (docs/19 §19.6): the server holds the request up to 25 s. + * Returns the new high-water cursor + any frames with cursor > [cursor]. + */ + suspend fun poll(cursor: Long): PollResult = + withContext(Dispatchers.IO) { + val request = + Request + .Builder() + .url("$baseUrl/v1/poll?cursor=$cursor") + .headers(authHeaders()) + .build() + client + .newCall(request) + .execute() + .use { response -> + if (!response.isSuccessful) { + throw IOException("poll failed: HTTP ${response.code}") + } + val body = response.body?.string().orEmpty() + val obj = IrisJson.instance.parseToJsonElement(body).jsonObject + val newCursor = obj["cursor"]?.jsonPrimitive?.content?.toLongOrNull() ?: cursor + val frames = + obj["frames"] + ?.jsonArray + ?.mapNotNull { el -> + try { + IrisJson.instance.decodeFromJsonElement(Frame.serializer(), el) + } catch (e: Exception) { + null + } + }.orEmpty() + PollResult(newCursor, frames) + } + } +} + +/** + * OkHttp's default read timeout (10 s) is shorter than the gateway's SSE + * heartbeat (15 s) and the long-poll hold (25 s) — extend the call timeout + * for streaming endpoints. Applied via [OkHttpClient] builders in + * [GatewayClient]. + */ +internal const val HTTP_STREAM_CALL_TIMEOUT_MS = 60_000L +internal const val HTTP_POLL_CALL_TIMEOUT_MS = 35_000L +internal const val HTTP_HEALTH_TIMEOUT_MS = 2_000L + +internal fun OkHttpClient.streamClient(): OkHttpClient = + newBuilder() + .callTimeout(HTTP_STREAM_CALL_TIMEOUT_MS, TimeUnit.MILLISECONDS) + .build() + +internal fun OkHttpClient.pollClient(): OkHttpClient = + newBuilder() + .callTimeout(HTTP_POLL_CALL_TIMEOUT_MS, TimeUnit.MILLISECONDS) + .build() + +internal fun OkHttpClient.healthClient(): OkHttpClient = + newBuilder() + .callTimeout(HTTP_HEALTH_TIMEOUT_MS, TimeUnit.MILLISECONDS) + .build() diff --git a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt index 3e6ce44..5ebadc4 100644 --- a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt +++ b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt @@ -432,7 +432,7 @@ class IrisController( ) { if (chatId.isNullOrBlank()) return pendingDeepLink = chatId to threadId - if (client.state.value is GatewayClient.State.Connected) applyDeepLink() + if (client.state.value is GatewayClient.State.HelloInfo) applyDeepLink() } private fun applyDeepLink() { @@ -683,7 +683,7 @@ class IrisController( client.state.collect { s -> val prev = prevState prevState = s - if (s is GatewayClient.State.Connected) { + if (s is GatewayClient.State.HelloInfo) { // Clear any stale "restarting" latch from the previous // down phase (the gateway's own status{online} frame // follows on hello.ack and re-asserts the truth). @@ -755,7 +755,7 @@ class IrisController( * after a process death the cached copy may be stale — so history always * refreshes (skipped on a plain reconnect via historyLoaded). */ - private fun onConnectedLane(connected: GatewayClient.State.Connected) { + private fun onConnectedLane(connected: GatewayClient.State.HelloInfo) { channels.setAll(connected.channels) val home = connected.channels.firstOrNull { it.isDefault }?.chatId _homeChannel.value = home ?: ChatStore.DEFAULT_LANE @@ -883,7 +883,7 @@ class IrisController( /** Request the gateway's slash-command catalog. No-op while disconnected * (sendFrame drops silently); the response lands via [slashCommands]. */ fun requestCommandsCatalog() { - if (client.state.value !is GatewayClient.State.Connected) return + if (client.state.value !is GatewayClient.State.HelloInfo) return client.sendFrame(commandsCatalogFrame(0)) } diff --git a/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt b/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt index f4fd4b5..22a4d25 100644 --- a/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt +++ b/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt @@ -245,9 +245,10 @@ fun ChatScreen(controller: IrisController) { val nonThreadChannels = channels.filter { it.kind != "thread" } fun doSend() { - // No-op while the socket is down (sendMessage drops silently); the send - // button is disabled in that state, this guards the IME "Send" action. - if (state !is GatewayClient.State.Connected) return + // No-op while no transport is up (sendMessage drops silently); the + // send button is disabled in that state, this guards the IME "Send" + // action. docs/19: the HTTP fallback leg counts as sendable. + if (state !is GatewayClient.State.HelloInfo) return val ready = attachments.filter { it.mediaRef != null && it.error == null } if (input.isBlank() && ready.isEmpty()) return val text = input @@ -311,7 +312,7 @@ fun ChatScreen(controller: IrisController) { } fun onSlashPick(cmd: SlashCommand) { - if (state !is GatewayClient.State.Connected) return + if (state !is GatewayClient.State.HelloInfo) return input = "" controller.send(cmd.name) focusManager.clearFocus(force = true) @@ -742,9 +743,12 @@ fun ChatScreen(controller: IrisController) { // Composer (M7: rounded pill + accent circular send button). Sending is gated // on a live socket: sendMessage is a no-op while disconnected, so an // ungated send would show a pending bubble that never resolves. - val isConnected = state is GatewayClient.State.Connected + val isConnected = state is GatewayClient.State.HelloInfo + // docs/19: media needs the live WS connection; in HTTP fallback + // only text is sendable. + val wsConnected = state is GatewayClient.State.Connected val canSend = - isConnected && (input.isNotBlank() || attachments.any { it.mediaRef != null && it.error == null }) + isConnected && (input.isNotBlank() || (wsConnected && attachments.any { it.mediaRef != null && it.error == null })) val layoutDensity = LocalDensity.current.density var textHeightPx by remember { mutableFloatStateOf(0f) } if (isAutomation) { @@ -818,7 +822,7 @@ fun ChatScreen(controller: IrisController) { ) } } - IconButton(onClick = { showPicker = true }) { + IconButton(onClick = { showPicker = true }, enabled = state is GatewayClient.State.Connected) { Text("📎", fontSize = 18.sp) } } @@ -1822,6 +1826,7 @@ private fun statusLabel(state: GatewayClient.State): String = GatewayClient.State.Connecting -> "connecting…" GatewayClient.State.Reconnecting -> "reconnecting…" is GatewayClient.State.Connected -> "connected" + is GatewayClient.State.HttpFallback -> "connected · http" is GatewayClient.State.AuthFailed -> "auth failed" } @@ -2104,6 +2109,7 @@ private fun NameDialog( private fun statusToastText(state: GatewayClient.State): String = when (state) { is GatewayClient.State.Connected -> "Connected to Hermes" + is GatewayClient.State.HttpFallback -> "Connected to Hermes (HTTP fallback — media paused)" GatewayClient.State.Connecting -> "Connecting to Hermes" GatewayClient.State.Reconnecting -> "Re-Connecting to Hermes" GatewayClient.State.Disconnected -> "Unpaired from Hermes" @@ -2119,7 +2125,9 @@ private fun StatusBubble( ) { val (color, pulsing) = when (state) { - is GatewayClient.State.Connected -> IrisColors.statusGreen to false + is GatewayClient.State.Connected, + is GatewayClient.State.HttpFallback, + -> IrisColors.statusGreen to false GatewayClient.State.Connecting, GatewayClient.State.Reconnecting, diff --git a/app/shared/src/commonTest/kotlin/iris/net/HttpGatewayTest.kt b/app/shared/src/commonTest/kotlin/iris/net/HttpGatewayTest.kt new file mode 100644 index 0000000..3f06754 --- /dev/null +++ b/app/shared/src/commonTest/kotlin/iris/net/HttpGatewayTest.kt @@ -0,0 +1,203 @@ +package iris.net + +import iris.protocol.Frame +import iris.protocol.IrisJson +import iris.protocol.MessagePayload +import iris.protocol.TYPE_MESSAGE +import kotlin.test.Test +import kotlin.test.assertEquals + +/** + * docs/19: unit tests for the HTTP fallback leg client. + * + * The SSE line parser is the trickiest pure logic (event/id/data fields, + * comments, multi-line data, Last-Event-ID bookkeeping), so it is factored + * into [SseParser] and tested directly. The WS->HTTP URL derivation is a + * pure function (unit-tested). Full transport behavior (health, POST, SSE + * catch-up, long-poll, delivery counting) is covered by the gateway-side + * Python tests (hermes-agent/tests/gateway/test_android_http.py). + */ +class HttpGatewayTest { + // ── URL derivation ──────────────────────────────────────────────────── + + @Test + fun deriveHttpUrlReplacesSchemeAndPort() { + assertEquals("http://192.168.1.10:8791", HttpGateway.deriveHttpUrl("ws://192.168.1.10:8790/ws")) + assertEquals("http://127.0.0.1:8791", HttpGateway.deriveHttpUrl("ws://127.0.0.1:8790/ws")) + assertEquals("https://gw.example.com:8791", HttpGateway.deriveHttpUrl("wss://gw.example.com:8790/ws")) + // No explicit port on the WS URL: still the HTTP leg's default port. + assertEquals("http://gw.example.com:8791", HttpGateway.deriveHttpUrl("ws://gw.example.com/ws")) + } + + @Test + fun deriveHttpUrlPassesThroughHttpUrls() { + assertEquals("http://1.2.3.4:9000", HttpGateway.deriveHttpUrl("http://1.2.3.4:9000")) + assertEquals("https://a.b", HttpGateway.deriveHttpUrl("https://a.b")) + } + + // ── SSE parser ──────────────────────────────────────────────────────── + + @Test + fun sseParsesHelloAndFrames() { + val parser = SseParser() + val hello = """{"v":1,"type":"hello.ack","payload":{"sync_cursor":5}}""" + val frame = """{"v":1,"type":"$TYPE_MESSAGE","payload":{"text":"hi"}}""" + val lines = + listOf( + "event: hello", + "data: $hello", + "", + "id: 7", + "event: frame", + "data: $frame", + "", + ) + var helloCount = 0 + var frameCount = 0 + var lastCursor: Long? = null + for (line in lines) { + parser.feed( + line, + onHello = { helloCount++ }, + onFrame = { frameCount++ }, + onCursor = { lastCursor = it }, + ) + } + assertEquals(1, helloCount) + assertEquals(1, frameCount) + assertEquals(7L, lastCursor) + } + + @Test + fun sseIgnoresHeartbeatComments() { + val parser = SseParser() + var frames = 0 + parser.feed(": hb", onHello = {}, onFrame = { frames++ }, onCursor = {}) + parser.feed("", onHello = {}, onFrame = { frames++ }, onCursor = {}) + assertEquals(0, frames) + } + + @Test + fun sseMultiLineDataJoinsWithNewline() { + val parser = SseParser() + var frame: Frame? = null + // SSE data may span multiple `data:` lines; the parser must join + // them with "\n" so the reassembled JSON still decodes. Split at a + // legal JSON whitespace point (right after a comma, between tokens). + val full = + """{"v":1,"type":"$TYPE_MESSAGE","payload":{"message_id":"m1","role":"assistant","text":"a\nb"}}""" + val cut = full.indexOf("\"m1\",") + "\"m1\",".length + val l1 = full.substring(0, cut) + val l2 = full.substring(cut) + parser.feed("event: frame", onHello = {}, onFrame = {}, onCursor = {}) + parser.feed("data: $l1", onHello = {}, onFrame = {}, onCursor = {}) + parser.feed("data: $l2", onHello = {}, onFrame = {}, onCursor = {}) + parser.feed("", onHello = {}, onFrame = { frame = it }, onCursor = {}) + assertEquals("a\nb", frame?.payloadAsText()) + } + + @Test + fun sseTracksLastEventIdAcrossFrames() { + val parser = SseParser() + val ids = mutableListOf() + val frame = """{"v":1,"type":"$TYPE_MESSAGE","payload":{}}""" + for (id in listOf(1L, 2L, 3L)) { + parser.feed("id: $id", onHello = {}, onFrame = {}, onCursor = { ids.add(it) }) + parser.feed("event: frame", onHello = {}, onFrame = {}, onCursor = {}) + parser.feed("data: $frame", onHello = {}, onFrame = {}, onCursor = {}) + parser.feed("", onHello = {}, onFrame = {}, onCursor = {}) + } + assertEquals(listOf(1L, 2L, 3L), ids) + assertEquals(3L, parser.lastEventId) + } + + @Test + fun sseMalformedLineDoesNotThrow() { + val parser = SseParser() + parser.feed("garbage without colon", onHello = {}, onFrame = {}, onCursor = {}) + parser.feed("id: notanumber", onHello = {}, onFrame = {}, onCursor = {}) + parser.feed("data: {not json", onHello = {}, onFrame = {}, onCursor = {}) + parser.feed("", onHello = {}, onFrame = {}, onCursor = {}) + // No exception, no frame emitted for the malformed data. + } + + @Test + fun sseDecodesFramePayload() { + val parser = SseParser() + var text: String? = null + val frame = + """{"v":1,"type":"$TYPE_MESSAGE","payload":{"message_id":"m1","role":"assistant","text":"hello"}}""" + parser.feed("event: frame", onHello = {}, onFrame = {}, onCursor = {}) + parser.feed("data: $frame", onHello = {}, onFrame = {}, onCursor = {}) + parser.feed("", onHello = {}, onFrame = { f -> text = f.payloadAsText() }, onCursor = {}) + assertEquals("hello", text) + } +} + +// ── SSE line parser (pure; shared by the live reader + tests) ───────────── + +/** + * Incremental SSE parser (docs/19 §19.5). Feed raw lines (without + * terminators); a blank line dispatches the buffered event. [lastEventId] + * is the most recent `id:` field (the outbox cursor) — the resume point for + * a reconnect. + */ +class SseParser { + var lastEventId: Long? = null + private set + + private var eventId: String? = null + private val dataLines = mutableListOf() + + fun feed( + line: String, + onHello: (String) -> Unit, + onFrame: (Frame) -> Unit, + onCursor: (Long) -> Unit, + ) { + when { + line.isEmpty() -> { + if (dataLines.isNotEmpty()) { + val data = dataLines.joinToString("\n") + val frame = + try { + IrisJson.instance.decodeFromString(Frame.serializer(), data) + } catch (e: Exception) { + null + } + if (frame != null) { + when (eventId) { + "hello" -> onHello(data) + else -> onFrame(frame) + } + } + } + eventId = null + dataLines.clear() + } + + line.startsWith(":") -> { + Unit + } + + // comment / heartbeat + line.startsWith("id:") -> { + val id = line.removePrefix("id:").trim().toLongOrNull() + if (id != null) { + lastEventId = id + onCursor(id) + } + } + + line.startsWith("event:") -> { + eventId = line.removePrefix("event:").trim() + } + + line.startsWith("data:") -> { + dataLines.add(line.removePrefix("data:").removePrefix(" ")) + } + } + } +} + +private fun Frame.payloadAsText(): String? = payloadAs()?.text