Compare commits

..
4 Commits
Author SHA1 Message Date
ARIA 8c09aecf4a docs: add trailing newline to AGENTS.md 2026-08-21 21:00:52 +02:00
ARIA 060420c615 docs: fix Kotlin test task name in AGENTS.md
testDebugUnitTest does not exist; the real host-side tasks are
testAndroidHostTest / desktopTest (jvmTest is the shared source set).
2026-08-21 20:56:24 +02:00
ARIA 81f42ab761 Local message cache: instant start + offline reading (SQLDelight)
- Cache.sq / ChatDb: lanes, channels and last_lane persisted as JSON
  snapshots; sanitize on restore (Pending->Failed, streaming->false);
  ephemeral items (tool, system) skipped
- Platform drivers via expect/actual: AndroidSqliteDriver (app db dir)
  / JdbcSqliteDriver (~/.iris/iris_cache.db)
- IrisController: restore before connect, debounced (750ms) snapshot
  persistence, synchronous flush on dispose, clearAll on forget
- History robustness: historyLoaded marked only when the response is
  processed (lost request/response retried on reconnect); events
  collector wrapped in try/catch; onHelloAck fast path so the history
  request fires on the WS thread instead of the starved state
  collector; frame-decode and history-load logging
- Protocol: HistoryMessage.media nullable (defensive vs older
  gateways that sent "media": null)
- Tests: ChatDbTest (jvmTest, JDBC in-memory), ChatStoreCacheTest,
  HistoryPayloadTest, HistoryWireTest (real captured 92KB response)
- docs/10: §10.7 implemented schema, new §10.9 cache behavior
2026-08-21 20:56:20 +02:00
ARIA 9a519e3c5a Omit null media in history frames (schema-conformant)
history() wrote "media": null for streaming finals, but the frame
schema says array. The app's HistoryMessage.media was non-nullable, so
deserialization threw and every history page was silently dropped.
Omit the key when media is absent instead.
2026-08-21 20:56:16 +02:00
20 changed files with 842 additions and 67 deletions

No files matched your search

+1 -1
View File
@@ -18,7 +18,7 @@
- Android: check the device is connected first (`adb devices` → `a5ca2a4b` listed as `device`); then `cd app && ./gradlew :androidApp:installDebug` to install on the phone and live-verify changes (launch/screenshot: see ADB below). - Android: check the device is connected first (`adb devices` → `a5ca2a4b` listed as `device`); then `cd app && ./gradlew :androidApp:installDebug` to install on the phone and live-verify changes (launch/screenshot: see ADB below).
- Desktop: `cd app && ./gradlew :desktopApp:run`; packaging: `:desktopApp:jpackage` (app-image; `-PjpackageType=deb` for a .deb). - Desktop: `cd app && ./gradlew :desktopApp:run`; packaging: `:desktopApp:jpackage` (app-image; `-PjpackageType=deb` for a .deb).
- Python tests — **never bare `pytest`**: `cd hermes-agent && scripts/run_tests.sh tests/gateway/test_android.py` (no args = full suite). - Python tests — **never bare `pytest`**: `cd hermes-agent && scripts/run_tests.sh tests/gateway/test_android.py` (no args = full suite).
- Kotlin tests: `cd app && ./gradlew :shared:testDebugUnitTest` / `:shared:desktopTest` (no test sources exist yet). - Kotlin tests: `cd app && ./gradlew :shared:testAndroidHostTest` / `:shared:desktopTest` (host-side; `jvmTest` is the shared source set).
- WS probe (gateway must be running): `hermes-agent/.venv/bin/python gateway-plugin/tests/ws_probe.py --token <ANDROID_TOKEN> --send "hello"` — assertion flags documented in `gateway-plugin/tests/README.md`. - WS probe (gateway must be running): `hermes-agent/.venv/bin/python gateway-plugin/tests/ws_probe.py --token <ANDROID_TOKEN> --send "hello"` — assertion flags documented in `gateway-plugin/tests/README.md`.
- E2E driver (gateway must be running; it never starts/stops it): `hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py`. - E2E driver (gateway must be running; it never starts/stops it): `hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py`.
+2
View File
@@ -19,6 +19,8 @@ plugins {
id("com.android.kotlin.multiplatform.library") version "9.3.1" apply false id("com.android.kotlin.multiplatform.library") version "9.3.1" apply false
id("org.jetbrains.compose") version "1.11.1" apply false id("org.jetbrains.compose") version "1.11.1" apply false
id("org.jetbrains.kotlin.plugin.compose") version "2.4.10" apply false id("org.jetbrains.kotlin.plugin.compose") version "2.4.10" apply false
// Local cache DB (docs/16: "Local DB (KMP) — SQLDelight").
id("app.cash.sqldelight") version "2.3.2" apply false
} }
tasks.register("clean") { tasks.register("clean") {
+28
View File
@@ -6,6 +6,8 @@ plugins {
id("com.android.kotlin.multiplatform.library") id("com.android.kotlin.multiplatform.library")
id("org.jetbrains.compose") id("org.jetbrains.compose")
id("org.jetbrains.kotlin.plugin.compose") id("org.jetbrains.kotlin.plugin.compose")
// Local cache DB (docs/16: "Local DB (KMP) — SQLDelight").
id("app.cash.sqldelight")
} }
val composeVersion = "1.11.1" val composeVersion = "1.11.1"
@@ -25,6 +27,8 @@ val kcefVersion = "2025.03.23"
// release). Provides GFM tables, bold/italic/underscore, and the -code module // release). Provides GFM tables, bold/italic/underscore, and the -code module
// for language-aware syntax highlighting (Highlights). // for language-aware syntax highlighting (Highlights).
val markdownVersion = "0.44.0" val markdownVersion = "0.44.0"
// Local cache DB (messages/channels/meta; docs/10 §10.7, docs/16).
val sqldelightVersion = "2.3.2"
kotlin { kotlin {
android { android {
@@ -51,6 +55,11 @@ kotlin {
val jvmMain by creating { dependsOn(commonMain.get()) } val jvmMain by creating { dependsOn(commonMain.get()) }
val androidMain by getting { dependsOn(jvmMain) } val androidMain by getting { dependsOn(jvmMain) }
val desktopMain by getting { dependsOn(jvmMain) } val desktopMain by getting { dependsOn(jvmMain) }
// Host tests shared by the android (withHostTest) and desktop targets
// (e.g. the SQLDelight cache DB tests, which need the JDBC driver).
val jvmTest by creating { dependsOn(commonTest.get()) }
val androidHostTest by getting { dependsOn(jvmTest) }
val desktopTest by getting { dependsOn(jvmTest) }
commonMain.dependencies { commonMain.dependencies {
implementation("org.jetbrains.compose.runtime:runtime:$composeVersion") implementation("org.jetbrains.compose.runtime:runtime:$composeVersion")
@@ -74,15 +83,24 @@ kotlin {
// M9: HTML artifact previews — WebView composable (platform WebView // M9: HTML artifact previews — WebView composable (platform WebView
// on Android, KCEF/JCEF on desktop). // on Android, KCEF/JCEF on desktop).
implementation("io.github.kevinnzou:compose-webview-multiplatform:$webviewVersion") implementation("io.github.kevinnzou:compose-webview-multiplatform:$webviewVersion")
// Local cache DB runtime (generated code from commonMain/sqldelight).
implementation("app.cash.sqldelight:runtime:$sqldelightVersion")
} }
commonTest.dependencies { commonTest.dependencies {
implementation(kotlin("test")) implementation(kotlin("test"))
} }
jvmTest.dependencies {
implementation(kotlin("test"))
// JDBC SQLite driver: in-memory DB for the cache tests.
implementation("app.cash.sqldelight:sqlite-driver:$sqldelightVersion")
}
// M9: KCEF (JCEF/Chromium) for desktop HTML artifact previews. The // M9: KCEF (JCEF/Chromium) for desktop HTML artifact previews. The
// webview library exposes it transitively, but we reference KCEF // webview library exposes it transitively, but we reference KCEF
// directly (init + progress) so declare it explicitly. // directly (init + progress) so declare it explicitly.
desktopMain.dependencies { desktopMain.dependencies {
implementation("dev.datlag:kcef:$kcefVersion") implementation("dev.datlag:kcef:$kcefVersion")
// Local cache DB: JDBC SQLite driver (file in ~/.iris).
implementation("app.cash.sqldelight:sqlite-driver:$sqldelightVersion")
} }
// M4: ExoPlayer (Media3) for inline audio/video playback (Android only). // M4: ExoPlayer (Media3) for inline audio/video playback (Android only).
androidMain.dependencies { androidMain.dependencies {
@@ -97,6 +115,16 @@ kotlin {
// the ntfy listener is the fallback). The google-services plugin is // the ntfy listener is the fallback). The google-services plugin is
// applied conditionally in the app module. // applied conditionally in the app module.
implementation("com.google.firebase:firebase-messaging:25.1.2") implementation("com.google.firebase:firebase-messaging:25.1.2")
// Local cache DB: Android SQLite driver (app database dir).
implementation("app.cash.sqldelight:android-driver:$sqldelightVersion")
}
}
}
sqldelight {
databases {
create("IrisDatabase") {
packageName.set("iris.db")
} }
} }
} }
@@ -0,0 +1,11 @@
package iris.platform
import app.cash.sqldelight.db.SqlDriver
import app.cash.sqldelight.driver.android.AndroidSqliteDriver
import iris.db.IrisDatabase
actual fun appDataDir(): String = AndroidEnv.context.filesDir.absolutePath
// The driver creates the schema in the open-helper callback (v1: no
// migrations yet; add .sqm files + `Schema.migrate` in onUpgrade later).
actual fun createCacheDriver(): SqlDriver = AndroidSqliteDriver(IrisDatabase.Schema, AndroidEnv.context, "iris_cache.db")
@@ -26,6 +26,14 @@ class ChannelStore {
_channels.value = sorted(channels) _channels.value = sorted(channels)
} }
/** Restore the directory from the persistent local cache on startup
* (offline: the drawer is populated before the gateway connection is
* up). The server re-seeds it on connect (`hello.ack`). No-op when
* [channels] is empty (first launch). */
fun loadFromCache(channels: List<ChannelInfo>) {
if (channels.isNotEmpty()) _channels.value = sorted(channels)
}
/** Reconcile a server frame into the cache. */ /** Reconcile a server frame into the cache. */
fun onFrame(frame: Frame) { fun onFrame(frame: Frame) {
when (frame.type) { when (frame.type) {
@@ -0,0 +1,135 @@
package iris.data
import app.cash.sqldelight.db.SqlDriver
import iris.db.IrisDatabase
import iris.protocol.ChannelInfo
import iris.protocol.IrisJson
/**
* Persistent local cache (SQLite via SQLDelight — the locked "Local DB (KMP)"
* decision, docs/16; schema docs/10 §10.7). The server is authoritative; this
* is a cache that lets the app start instantly and read chats offline:
*
* - [loadLanes] / [loadChannels] restore the in-memory stores on startup,
* before the gateway connection is up (snappy start, offline reading).
* - [saveLanes] / [saveChannels] snapshot the in-memory stores (debounced by
* the controller) so every reconciled frame survives a process death.
* - [metaGet] / [metaPut] hold small UI state (last-viewed lane).
*
* Only [MessageItem]s are persisted — tool cards, live streaming state and
* local system notices are ephemeral. Rows are JSON payloads keyed by
* (lane, id), so the schema does not drift with [MessageItem] fields.
*
* All access is synchronized: the debounced save collectors run on the
* controller scope while [dispose] may flush from the UI thread.
*/
class ChatDb(
driver: SqlDriver,
) {
private val db = IrisDatabase(driver)
private val json = IrisJson.instance
private val lock = Any()
// ── messages ──────────────────────────────────────────────────────────
/** All persisted lanes (lane key -> messages ordered by ts). */
fun loadLanes(): Map<String, List<MessageItem>> =
synchronized(lock) {
val lanes = linkedMapOf<String, MutableList<MessageItem>>()
for (row in db.cacheQueries.allMessages().executeAsList()) {
val item = decodeMessage(row.payload) ?: continue
lanes.getOrPut(row.lane) { mutableListOf() }.add(item)
}
lanes.mapValues { it.value.toList() }
}
/** Replace the whole message cache with [lanes] (atomic snapshot). Tool
* cards and local system notices are skipped (ephemeral). */
fun saveLanes(lanes: Map<String, List<ChatItem>>) {
synchronized(lock) {
db.transaction {
db.cacheQueries.clearMessages()
for ((lane, items) in lanes) {
for (item in items) {
if (item is MessageItem && !item.isSystem) {
db.cacheQueries.upsertMessage(lane, item.id, item.ts, json.encodeToString(item))
}
}
}
}
}
}
// ── channels ──────────────────────────────────────────────────────────
/** The persisted channel directory (channels + threads). */
fun loadChannels(): List<ChannelInfo> =
synchronized(lock) {
val list = mutableListOf<ChannelInfo>()
for (row in db.cacheQueries.allChannels().executeAsList()) {
try {
list.add(json.decodeFromString(row.payload))
} catch (_: Exception) {
// Corrupt row — skip it; the server re-seeds on connect.
}
}
list
}
/** Replace the whole channel directory (atomic snapshot). */
fun saveChannels(channels: List<ChannelInfo>) {
synchronized(lock) {
db.transaction {
db.cacheQueries.clearChannels()
for (c in channels) {
db.cacheQueries.upsertChannel(c.chatId, json.encodeToString(c))
}
}
}
}
// ── meta ──────────────────────────────────────────────────────────────
fun metaGet(key: String): String? =
synchronized(lock) {
db.cacheQueries.metaGet(key).executeAsOneOrNull()
}
fun metaPut(
key: String,
value: String,
) {
synchronized(lock) {
db.cacheQueries.metaPut(key, value)
}
}
/** Wipe the whole cache (forget pairing → a different gateway). */
fun clearAll() {
synchronized(lock) {
db.transaction {
db.cacheQueries.clearMessages()
db.cacheQueries.clearChannels()
db.cacheQueries.clearMeta()
}
}
}
private fun decodeMessage(payload: String): MessageItem? =
try {
json.decodeFromString<MessageItem>(payload).sanitizeForRestore()
} catch (_: Exception) {
null
}
/** A restored message is never mid-flight: a streaming bubble is
* finalized (the sync replay finalizes it for real on reconnect) and a
* pending send becomes failed (tap to retry) — the gateway never
* acknowledged it before the process died. */
private fun MessageItem.sanitizeForRestore(): MessageItem =
copy(
streaming = false,
pending = false,
status = if (status == MsgStatus.Pending) MsgStatus.Failed else status,
)
}
@@ -27,6 +27,7 @@ import iris.util.nowMillis
import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.asStateFlow
import kotlinx.serialization.Serializable
import kotlinx.serialization.json.JsonElement import kotlinx.serialization.json.JsonElement
import kotlin.random.Random import kotlin.random.Random
@@ -54,7 +55,9 @@ enum class MsgStatus {
/** A chat message (user / assistant / commentary / streaming bubble). /** A chat message (user / assistant / commentary / streaming bubble).
* [isSystem] marks a locally generated, centered notice (gateway restart / * [isSystem] marks a locally generated, centered notice (gateway restart /
* online) — not a user or agent turn, not selectable or deletable. */ * online) — not a user or agent turn, not selectable or deletable.
* [Serializable]: persisted as a JSON payload in the local cache (ChatDb). */
@Serializable
data class MessageItem( data class MessageItem(
override val id: String, override val id: String,
val role: String, val role: String,
@@ -77,6 +80,7 @@ data class MessageItem(
* has been pulled to the local cache (outbound) — inbound attachments the * has been pulled to the local cache (outbound) — inbound attachments the
* app itself uploaded carry no local path (the agent reads the server copy). * app itself uploaded carry no local path (the agent reads the server copy).
*/ */
@Serializable
data class MediaItem( data class MediaItem(
val mediaId: String, val mediaId: String,
val kind: String, val kind: String,
@@ -700,6 +704,18 @@ class ChatStore {
} }
} }
/**
* Restore lanes from the persistent local cache on startup (before the
* gateway connection is up), so the UI is populated instantly and chats
* are readable offline. The server remains authoritative: the `sync`
* delta and the `history` refresh reconcile the cache on connect.
* No-op when [lanes] is empty (first launch).
*/
fun loadFromCache(lanes: Map<String, List<MessageItem>>) {
if (lanes.isEmpty()) return
_lanes.value = lanes
}
fun clear() { fun clear() {
_lanes.value = emptyMap() _lanes.value = emptyMap()
} }
@@ -30,29 +30,29 @@ import kotlinx.coroutines.Job
import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.currentCoroutineContext import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.delay import kotlinx.coroutines.delay
import kotlinx.coroutines.isActive
import kotlinx.coroutines.launch
import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.SharedFlow import kotlinx.coroutines.flow.SharedFlow
import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asSharedFlow import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.isActive
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withTimeout import kotlinx.coroutines.withTimeout
import kotlinx.coroutines.withTimeoutOrNull import kotlinx.coroutines.withTimeoutOrNull
import kotlin.time.TimeMark
import kotlin.time.TimeSource
import okhttp3.OkHttpClient import okhttp3.OkHttpClient
import okio.ByteString
import okio.ByteString.Companion.toByteString
import okhttp3.Request import okhttp3.Request
import okhttp3.Response import okhttp3.Response
import okhttp3.WebSocket import okhttp3.WebSocket
import okhttp3.WebSocketListener import okhttp3.WebSocketListener
import okio.ByteString
import okio.ByteString.Companion.toByteString
import java.util.concurrent.TimeUnit import java.util.concurrent.TimeUnit
import kotlin.random.Random import kotlin.random.Random
import kotlin.time.TimeMark
import kotlin.time.TimeSource
/** /**
* OkHttp WebSocket client for the hermes android gateway (docs/10 §10.3). * OkHttp WebSocket client for the hermes android gateway (docs/10 §10.3).
@@ -69,7 +69,9 @@ class GatewayClient(
) { ) {
sealed interface State { sealed interface State {
data object Disconnected : State data object Disconnected : State
data object Connecting : State data object Connecting : State
data class Connected( data class Connected(
val caps: ServerCaps, val caps: ServerCaps,
val channels: List<ChannelInfo>, val channels: List<ChannelInfo>,
@@ -78,8 +80,12 @@ class GatewayClient(
* it must not re-post system notifications (docs/08 §8.7). */ * it must not re-post system notifications (docs/08 §8.7). */
val lastPushedCursor: Long = 0, val lastPushedCursor: Long = 0,
) : State ) : State
data object Reconnecting : State data object Reconnecting : State
data class AuthFailed(val message: String) : State
data class AuthFailed(
val message: String,
) : State
} }
private val _state = MutableStateFlow<State>(State.Disconnected) private val _state = MutableStateFlow<State>(State.Disconnected)
@@ -92,7 +98,9 @@ class GatewayClient(
private val _events = MutableSharedFlow<Frame>(extraBufferCapacity = 128) private val _events = MutableSharedFlow<Frame>(extraBufferCapacity = 128)
val events: SharedFlow<Frame> = _events.asSharedFlow() val events: SharedFlow<Frame> = _events.asSharedFlow()
private val client: OkHttpClient = OkHttpClient.Builder() private val client: OkHttpClient =
OkHttpClient
.Builder()
.pingInterval(20, TimeUnit.SECONDS) .pingInterval(20, TimeUnit.SECONDS)
.build() .build()
@@ -100,6 +108,7 @@ class GatewayClient(
private var socket: WebSocket? = null private var socket: WebSocket? = null
private var nextRequestId = 1 private var nextRequestId = 1
private var attempt = 0 private var attempt = 0
// True once a connection has been established this session; reset by // True once a connection has been established this session; reset by
// start(). Drives Connecting (first dial) vs Reconnecting (redial after a // start(). Drives Connecting (first dial) vs Reconnecting (redial after a
// drop) so the UI can show the right status without a blocking screen. // drop) so the UI can show the right status without a blocking screen.
@@ -120,6 +129,15 @@ class GatewayClient(
private var binarySession: BinarySession? = null private var binarySession: BinarySession? = null
/**
* Fired promptly (on the WS thread) the moment `hello.ack` is received —
* on every (re)connect. Used for time-critical work that must not wait for
* the state collector, which can be starved for seconds during app startup
* (Dispatchers.Default) and would push a history request past a flaky
* network's window. Set before [start].
*/
var onHelloAck: ((State.Connected) -> Unit)? = null
// M4: only one pull may be in flight at a time (binarySession is a single // 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. // slot). Serialize concurrent offers so their byte streams don't interleave.
private val pullMutex = Mutex() private val pullMutex = Mutex()
@@ -165,6 +183,7 @@ class GatewayClient(
dial.socket.close(1000, "auth failed") dial.socket.close(1000, "auth failed")
return return
} }
DialResult.Connected -> { DialResult.Connected -> {
hasConnected = true hasConnected = true
attempt = 0 attempt = 0
@@ -173,6 +192,7 @@ class GatewayClient(
if (!currentCoroutineContext().isActive) return if (!currentCoroutineContext().isActive) return
// socket dropped -> loop again (Reconnecting) // socket dropped -> loop again (Reconnecting)
} }
is DialResult.Failed -> { is DialResult.Failed -> {
attempt++ attempt++
delay(backoffMs(attempt)) delay(backoffMs(attempt))
@@ -185,8 +205,14 @@ class GatewayClient(
private sealed interface DialResult { private sealed interface DialResult {
data object Connected : DialResult data object Connected : DialResult
data class AuthFailed(val message: String) : DialResult
data class Failed(val message: String) : DialResult data class AuthFailed(
val message: String,
) : DialResult
data class Failed(
val message: String,
) : DialResult
} }
private data class Dial( private data class Dial(
@@ -195,17 +221,24 @@ class GatewayClient(
val closed: CompletableDeferred<Unit>, val closed: CompletableDeferred<Unit>,
) )
private suspend fun dial(url: String, token: String): Dial { private suspend fun dial(
url: String,
token: String,
): Dial {
val closed = CompletableDeferred<Unit>() val closed = CompletableDeferred<Unit>()
val helloAck = CompletableDeferred<HelloAckPayload>() val helloAck = CompletableDeferred<HelloAckPayload>()
val authError = CompletableDeferred<String>() val authError = CompletableDeferred<String>()
val fail = CompletableDeferred<String>() val fail = CompletableDeferred<String>()
val request = Request.Builder().url(url).build() val request = Request.Builder().url(url).build()
val ws = client.newWebSocket( val ws =
client.newWebSocket(
request, request,
object : WebSocketListener() { object : WebSocketListener() {
override fun onOpen(webSocket: WebSocket, response: Response) { override fun onOpen(
webSocket: WebSocket,
response: Response,
) {
IrisLog.d("ws open (${response.code})") IrisLog.d("ws open (${response.code})")
webSocket.send( webSocket.send(
helloFrame( helloFrame(
@@ -218,11 +251,18 @@ class GatewayClient(
) )
} }
override fun onMessage(webSocket: WebSocket, text: String) { override fun onMessage(
webSocket: WebSocket,
text: String,
) {
lastLiveness = TimeSource.Monotonic.markNow() lastLiveness = TimeSource.Monotonic.markNow()
val frame = try { val frame =
try {
IrisJson.instance.decodeFromString(Frame.serializer(), text) IrisJson.instance.decodeFromString(Frame.serializer(), text)
} catch (_: Exception) { } catch (e: Exception) {
// A dropped frame is silent data loss — log it (the
// first bytes hint at which frame it was).
IrisLog.e("frame decode failed (${text.length}B): $e :: ${text.take(120)}")
return return
} }
when (frame.type) { when (frame.type) {
@@ -230,6 +270,7 @@ class GatewayClient(
val ack = frame.payloadAs<HelloAckPayload>() val ack = frame.payloadAs<HelloAckPayload>()
if (ack != null) helloAck.complete(ack) if (ack != null) helloAck.complete(ack)
} }
TYPE_ERROR -> { TYPE_ERROR -> {
val err = frame.payloadAs<ErrorPayload>() val err = frame.payloadAs<ErrorPayload>()
if (!authError.isCompleted) authError.complete(err?.message ?: "auth failed") if (!authError.isCompleted) authError.complete(err?.message ?: "auth failed")
@@ -237,7 +278,11 @@ class GatewayClient(
// auth failures — let the controller react. // auth failures — let the controller react.
_events.tryEmit(frame) _events.tryEmit(frame)
} }
TYPE_PONG -> Unit
TYPE_PONG -> {
Unit
}
else -> { else -> {
_events.tryEmit(frame) _events.tryEmit(frame)
frame.id?.let { id -> frame.id?.let { id ->
@@ -255,7 +300,10 @@ class GatewayClient(
} }
} }
override fun onMessage(webSocket: WebSocket, bytes: ByteString) { override fun onMessage(
webSocket: WebSocket,
bytes: ByteString,
) {
lastLiveness = TimeSource.Monotonic.markNow() lastLiveness = TimeSource.Monotonic.markNow()
// M4: binary frames belong to the active pull stream // M4: binary frames belong to the active pull stream
// (uploads are outbound; stray inbound chunks are dropped). // (uploads are outbound; stray inbound chunks are dropped).
@@ -264,12 +312,20 @@ class GatewayClient(
?.trySend(bytes.toByteArray()) ?.trySend(bytes.toByteArray())
} }
override fun onClosed(webSocket: WebSocket, code: Int, reason: String) { override fun onClosed(
webSocket: WebSocket,
code: Int,
reason: String,
) {
IrisLog.w("ws closed code=$code reason=\"$reason\"") IrisLog.w("ws closed code=$code reason=\"$reason\"")
closed.complete(Unit) closed.complete(Unit)
} }
override fun onFailure(webSocket: WebSocket, t: Throwable, response: Response?) { override fun onFailure(
webSocket: WebSocket,
t: Throwable,
response: Response?,
) {
IrisLog.e("ws failure: ${t.javaClass.simpleName}: ${t.message} (http=${response?.code})") IrisLog.e("ws failure: ${t.javaClass.simpleName}: ${t.message} (http=${response?.code})")
fail.complete(t.message ?: "connection failed") fail.complete(t.message ?: "connection failed")
closed.complete(Unit) closed.complete(Unit)
@@ -282,13 +338,16 @@ class GatewayClient(
helloAck.invokeOnCompletion { e -> helloAck.invokeOnCompletion { e ->
if (e == null) { if (e == null) {
val ack = helloAck.getCompleted() val ack = helloAck.getCompleted()
_state.value = State.Connected(ack.serverCaps, ack.channels, ack.lastPushedCursor) val connected = State.Connected(ack.serverCaps, ack.channels, ack.lastPushedCursor)
_state.value = connected
// M5: reconnect catch-up — replay frames parked while offline. // M5: reconnect catch-up — replay frames parked while offline.
val local = store.syncCursor val local = store.syncCursor
if (local < ack.syncCursor) { if (local < ack.syncCursor) {
val id = nextRequestId++ val id = nextRequestId++
ws.send(syncFrame(id, local).toWire()) ws.send(syncFrame(id, local).toWire())
} }
// Prompt fast path (before the possibly-starved state collector).
onHelloAck?.invoke(connected)
winner.complete(DialResult.Connected) winner.complete(DialResult.Connected)
} }
} }
@@ -298,7 +357,8 @@ class GatewayClient(
fail.invokeOnCompletion { e -> fail.invokeOnCompletion { e ->
if (e == null) winner.complete(DialResult.Failed(fail.getCompleted())) if (e == null) winner.complete(DialResult.Failed(fail.getCompleted()))
} }
val result = withTimeoutOrNull(15_000) { winner.await() } val result =
withTimeoutOrNull(15_000) { winner.await() }
?: DialResult.Failed("timeout waiting for hello.ack") ?: DialResult.Failed("timeout waiting for hello.ack")
return Dial(result, ws, closed) return Dial(result, ws, closed)
} }
@@ -363,14 +423,21 @@ class GatewayClient(
return when (frame.type) { return when (frame.type) {
TYPE_MEDIA_UPLOAD_ACK -> { TYPE_MEDIA_UPLOAD_ACK -> {
val p = frame.payloadAs<MediaUploadAckPayload>() val p = frame.payloadAs<MediaUploadAckPayload>()
if (p != null && p.ok) Result.success(p.mediaRef) if (p != null && p.ok) {
else Result.failure(IllegalStateException("upload rejected by server")) Result.success(p.mediaRef)
} else {
Result.failure(IllegalStateException("upload rejected by server"))
} }
}
TYPE_ERROR -> { TYPE_ERROR -> {
val e = frame.payloadAs<ErrorPayload>() val e = frame.payloadAs<ErrorPayload>()
Result.failure(IllegalStateException(e?.message ?: "upload failed")) Result.failure(IllegalStateException(e?.message ?: "upload failed"))
} }
else -> Result.failure(IllegalStateException("unexpected reply ${frame.type}"))
else -> {
Result.failure(IllegalStateException("unexpected reply ${frame.type}"))
}
} }
} catch (e: Exception) { } catch (e: Exception) {
return Result.failure(e) return Result.failure(e)
@@ -383,7 +450,10 @@ class GatewayClient(
* Pull offered media (docs/07 §7.3): media.pull, then binary frames until * Pull offered media (docs/07 §7.3): media.pull, then binary frames until
* media.pull.end. Each chunk is handed to [onChunk] (write to cache). * media.pull.end. Each chunk is handed to [onChunk] (write to cache).
*/ */
suspend fun pullMedia(mediaId: String, onChunk: suspend (ByteArray) -> Unit): Result<Unit> = suspend fun pullMedia(
mediaId: String,
onChunk: suspend (ByteArray) -> Unit,
): Result<Unit> =
pullMutex.withLock { pullMutex.withLock {
val ws = socket ?: return@withLock Result.failure(IllegalStateException("not connected")) val ws = socket ?: return@withLock Result.failure(IllegalStateException("not connected"))
val id = nextRequestId++ val id = nextRequestId++
@@ -393,21 +463,29 @@ class GatewayClient(
binarySession = BinarySession.Pulling(id, chunks, end) binarySession = BinarySession.Pulling(id, chunks, end)
try { try {
ws.send(mediaPullFrame(id, mediaId).toWire()) ws.send(mediaPullFrame(id, mediaId).toWire())
val frame = withTimeout(PULL_TIMEOUT_MS) { val frame =
withTimeout(PULL_TIMEOUT_MS) {
for (chunk in chunks) onChunk(chunk) for (chunk in chunks) onChunk(chunk)
end.await() end.await()
} }
when (frame.type) { when (frame.type) {
TYPE_MEDIA_PULL_END -> { TYPE_MEDIA_PULL_END -> {
val p = frame.payloadAs<MediaPullEndPayload>() val p = frame.payloadAs<MediaPullEndPayload>()
if (p != null && p.ok) Result.success(Unit) if (p != null && p.ok) {
else Result.failure(IllegalStateException("pull failed")) Result.success(Unit)
} else {
Result.failure(IllegalStateException("pull failed"))
} }
}
TYPE_ERROR -> { TYPE_ERROR -> {
val e = frame.payloadAs<ErrorPayload>() val e = frame.payloadAs<ErrorPayload>()
Result.failure(IllegalStateException(e?.message ?: "pull failed")) Result.failure(IllegalStateException(e?.message ?: "pull failed"))
} }
else -> Result.failure(IllegalStateException("unexpected reply ${frame.type}"))
else -> {
Result.failure(IllegalStateException("unexpected reply ${frame.type}"))
}
} }
} catch (e: Exception) { } catch (e: Exception) {
Result.failure(e) Result.failure(e)
@@ -459,15 +537,24 @@ class GatewayClient(
* Real `hello` test: dial, wait for hello.ack (or auth error), close. * Real `hello` test: dial, wait for hello.ack (or auth error), close.
* Exercises the auth leg, not just TCP (docs/10 §10.8). * Exercises the auth leg, not just TCP (docs/10 §10.8).
*/ */
suspend fun testHello(url: String, token: String): Result<Unit> { suspend fun testHello(
url: String,
token: String,
): Result<Unit> {
val dial = dial(url, token) val dial = dial(url, token)
return when (val result = dial.result) { return when (val result = dial.result) {
DialResult.Connected -> { DialResult.Connected -> {
dial.socket.close(1000, "test complete") dial.socket.close(1000, "test complete")
Result.success(Unit) Result.success(Unit)
} }
is DialResult.AuthFailed -> Result.failure(IllegalStateException(result.message))
is DialResult.Failed -> Result.failure(IllegalStateException(result.message)) is DialResult.AuthFailed -> {
Result.failure(IllegalStateException(result.message))
}
is DialResult.Failed -> {
Result.failure(IllegalStateException(result.message))
}
} }
} }
@@ -0,0 +1,13 @@
package iris.platform
import app.cash.sqldelight.db.SqlDriver
/**
* Persistent app-data directory. Unlike the media cache dir, it is not wiped
* by the system or by the user clearing the app cache — it holds the SQLite
* local cache (messages / channels / meta).
*/
expect fun appDataDir(): String
/** Open the SQLite driver for the local cache database. */
expect fun createCacheDriver(): SqlDriver
@@ -484,7 +484,10 @@ data class HistoryMessage(
val tokens: Int? = null, val tokens: Int? = null,
val runtime: RuntimeMeta? = null, val runtime: RuntimeMeta? = null,
val ts: Long? = null, val ts: Long? = null,
val media: List<MediaRef> = emptyList(), // Nullable: the schema says array, but older gateways sent
// `"media": null` for streaming finals — a null there used to break
// deserialization of the whole history page (silently dropping it).
val media: List<MediaRef>? = null,
) )
@Serializable @Serializable
@@ -1,6 +1,7 @@
package iris.state package iris.state
import iris.data.ChannelStore import iris.data.ChannelStore
import iris.data.ChatDb
import iris.data.ChatStore import iris.data.ChatStore
import iris.data.MediaItem import iris.data.MediaItem
import iris.data.MessageItem import iris.data.MessageItem
@@ -10,6 +11,7 @@ import iris.media.MediaCache
import iris.media.kindFromMime import iris.media.kindFromMime
import iris.net.GatewayClient import iris.net.GatewayClient
import iris.platform.PickedFile import iris.platform.PickedFile
import iris.platform.createCacheDriver
import iris.platform.isAppForeground import iris.platform.isAppForeground
import iris.platform.mediaCacheBaseDir import iris.platform.mediaCacheBaseDir
import iris.platform.postSystemNotification import iris.platform.postSystemNotification
@@ -73,6 +75,7 @@ import iris.protocol.syncFrame
import iris.ui.theme.Backdrop import iris.ui.theme.Backdrop
import iris.ui.theme.BackgroundMode import iris.ui.theme.BackgroundMode
import iris.ui.theme.UserTheme import iris.ui.theme.UserTheme
import iris.util.IrisLog
import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.SupervisorJob
@@ -80,6 +83,7 @@ import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.debounce
import kotlinx.coroutines.launch import kotlinx.coroutines.launch
import kotlin.random.Random import kotlin.random.Random
@@ -98,6 +102,12 @@ class IrisController(
val chat = ChatStore() val chat = ChatStore()
val channels = ChannelStore() val channels = ChannelStore()
/** Persistent local cache (SQLite; docs/10 §10.7, docs/16). Restores the
* stores on startup (instant UI, offline reading) and snapshots them
* (debounced) so reconciled frames survive a process death. The server
* is authoritative — the `sync` delta + `history` refresh reconcile. */
val chatDb = ChatDb(createCacheDriver())
private val _typing = MutableStateFlow(false) private val _typing = MutableStateFlow(false)
val typing: StateFlow<Boolean> = _typing.asStateFlow() val typing: StateFlow<Boolean> = _typing.asStateFlow()
@@ -210,6 +220,13 @@ class IrisController(
const val FONT_SCALE_MIN = 0.8f const val FONT_SCALE_MIN = 0.8f
const val FONT_SCALE_MAX = 1.5f const val FONT_SCALE_MAX = 1.5f
/** meta key: the last-viewed lane (restored on startup). */
const val META_LAST_LANE = "last_lane"
/** Debounce for the cache snapshot writes (the outbox replay on
* reconnect is the safety net for anything lost in the window). */
private const val CACHE_SAVE_DEBOUNCE_MS = 750L
/** Valid runtime-footer field keys (mirror of the gateway's RUNTIME_FIELDS). */ /** Valid runtime-footer field keys (mirror of the gateway's RUNTIME_FIELDS). */
val RUNTIME_FIELD_KEYS = listOf("model", "context_pct", "cwd", "latency", "cost") val RUNTIME_FIELD_KEYS = listOf("model", "context_pct", "cwd", "latency", "cost")
@@ -414,9 +431,28 @@ class IrisController(
} }
init { init {
// Restore the persistent local cache BEFORE connecting: the UI is
// populated instantly (snappy start) and chats are readable offline.
// The server is authoritative — the `sync` delta + `history` refresh
// reconcile the cache on connect.
chat.loadFromCache(chatDb.loadLanes())
channels.loadFromCache(chatDb.loadChannels())
chatDb.metaGet(META_LAST_LANE)?.let { chat.setLane(it) }
chat.streamingEnabled = _streamingEnabled.value chat.streamingEnabled = _streamingEnabled.value
// Snapshot the in-memory stores into the cache (debounced; a full
// atomic rewrite per change, so every mutation path is covered).
scope.launch {
chat.lanes.debounce(CACHE_SAVE_DEBOUNCE_MS).collect { chatDb.saveLanes(it) }
}
scope.launch {
channels.channels.debounce(CACHE_SAVE_DEBOUNCE_MS).collect { chatDb.saveChannels(it) }
}
scope.launch {
chat.currentLane.debounce(CACHE_SAVE_DEBOUNCE_MS).collect { chatDb.metaPut(META_LAST_LANE, it) }
}
scope.launch { scope.launch {
client.events.collect { frame -> client.events.collect { frame ->
try {
when (frame.type) { when (frame.type) {
TYPE_MESSAGE, TYPE_MESSAGE,
TYPE_MESSAGE_START, TYPE_MESSAGE_START,
@@ -527,14 +563,27 @@ class IrisController(
// since the saved cursor, so after a process death the // since the saved cursor, so after a process death the
// in-memory store is empty and older messages are not in // in-memory store is empty and older messages are not in
// the delta — history loads the full list. // the delta — history loads the full list.
frame.payloadAs<HistoryPayload>()?.let { p -> val p = frame.payloadAs<HistoryPayload>()
val (chatId, threadId) = if (p == null) {
(frame.chatId ?: return@let) to frame.threadId // The page was dropped on deserialization — the
// lane stays unmarked, so a (re)connect retries.
IrisLog.e("history page failed to deserialize (chat=${frame.chatId} thread=${frame.threadId})")
} else {
frame.chatId?.let { chatId ->
val threadId = frame.threadId
val lane = chat.laneKey(chatId, threadId) val lane = chat.laneKey(chatId, threadId)
IrisLog.d("history loaded: lane=$lane messages=${p.messages.size}")
chat.loadHistory( chat.loadHistory(
lane, lane,
p.messages.map { it.toMessageItem() }, p.messages.map { it.toMessageItem() },
) )
// Mark the lane loaded only when the response
// is actually processed: if the request or
// response is lost in a WS drop, the lane
// stays unmarked and the next (re)connect
// retries it.
historyLoaded.add(lane)
}
} }
} }
@@ -584,6 +633,12 @@ class IrisController(
Unit Unit
} }
} }
} catch (e: Exception) {
// A frame handler must never kill the collector: if it
// dies, every subsequent frame is dropped silently
// (tryEmit has no subscriber). Log and keep going.
IrisLog.e("frame handler failed for ${frame.type}: $e")
}
} }
} }
scope.launch { scope.launch {
@@ -592,20 +647,11 @@ class IrisController(
val prev = prevState val prev = prevState
prevState = s prevState = s
if (s is GatewayClient.State.Connected) { if (s is GatewayClient.State.Connected) {
channels.setAll(s.channels) // The lane/history fast path runs on [client.onHelloAck]
// (promptly, on the WS thread) — see onConnectedLane. Here
// we do the non-time-critical connect work.
// M5: refresh the push-dedupe watermark (docs/08 §8.7). // M5: refresh the push-dedupe watermark (docs/08 §8.7).
lastPushedCursor = s.lastPushedCursor lastPushedCursor = s.lastPushedCursor
val home = s.channels.firstOrNull { it.isDefault }?.chatId
if (home != null) {
_homeChannel.value = home
chat.setLane(home)
// Restore the home lane after a process death: the
// in-memory store is empty and the `sync` delta does not
// cover messages older than the saved cursor, so load
// the full history. Skipped on a plain reconnect (the
// lane's history was already loaded this session).
loadHistory(home)
}
// M5: remember the ntfy server for the listener service. // M5: remember the ntfy server for the listener service.
if (s.caps.pushNtfyServer.isNotBlank()) { if (s.caps.pushNtfyServer.isNotBlank()) {
store.ntfyServer = s.caps.pushNtfyServer store.ntfyServer = s.caps.pushNtfyServer
@@ -645,9 +691,48 @@ class IrisController(
store.ntfyTopic = "iris-${store.deviceId}-${Random.nextLong(1_000_000_000L, 9_999_999_999L)}" store.ntfyTopic = "iris-${store.deviceId}-${Random.nextLong(1_000_000_000L, 9_999_999_999L)}"
} }
client.startHeartbeat() client.startHeartbeat()
// Prompt fast path: seed the channel directory + load the active lane's
// history the moment hello.ack lands (on the WS thread), not after the
// state collector (which can be starved for seconds on startup). This
// gets the history request out early so its response lands inside a
// flaky network's window.
client.onHelloAck = { connected ->
try {
onConnectedLane(connected)
} catch (e: Exception) {
// Must not throw on the WS thread (would break the connection).
IrisLog.e("onConnectedLane failed: $e")
}
}
client.start() client.start()
} }
/**
* Fast-path connect handler (runs on the WS thread via [GatewayClient.onHelloAck],
* promptly on every (re)connect). Seeds the channel directory and loads the
* active lane's full history. The lane is the last-viewed one (restored from
* the local cache) if it still exists on the server, else the home channel.
* The `sync` delta does not cover messages older than the saved cursor, and
* 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) {
channels.setAll(connected.channels)
val home = connected.channels.firstOrNull { it.isDefault }?.chatId
_homeChannel.value = home ?: ChatStore.DEFAULT_LANE
if (home == null) return
val (restoredChat, restoredThread) = chat.parseLane(chat.currentLane.value)
val restoredValid =
channels.byId(restoredChat) != null &&
(restoredThread == null || channels.byId(restoredThread) != null)
if (restoredValid) {
loadHistory(restoredChat, restoredThread)
} else {
chat.setLane(home)
loadHistory(home)
}
}
// ── M3: navigation (lane switching) ─────────────────────────────────── // ── M3: navigation (lane switching) ───────────────────────────────────
/** Switch to a channel's flat / "General" lane. Loads the full history on /** Switch to a channel's flat / "General" lane. Loads the full history on
@@ -779,9 +864,10 @@ class IrisController(
) { ) {
val lane = chat.laneKey(chatId, threadId) val lane = chat.laneKey(chatId, threadId)
if (lane in historyLoaded) return if (lane in historyLoaded) return
historyLoaded.add(lane)
// Newest page, sized to restore a full working view on restart / first // Newest page, sized to restore a full working view on restart / first
// open (older pages are reachable via scroll-up pagination). // open (older pages are reachable via scroll-up pagination). The lane
// is marked loaded when the response arrives (TYPE_HISTORY), not here —
// a request lost in a WS drop must be retryable on reconnect.
client.sendFrame(historyFrame(0, chatId, threadId, limit = 200)) client.sendFrame(historyFrame(0, chatId, threadId, limit = 200))
} }
@@ -930,6 +1016,8 @@ class IrisController(
fun forget() { fun forget() {
store.clear() store.clear()
// A different gateway means a different chat universe — wipe the cache.
chatDb.clearAll()
chat.clear() chat.clear()
channels.clear() channels.clear()
client.restart() client.restart()
@@ -937,6 +1025,10 @@ class IrisController(
fun dispose() { fun dispose() {
client.stop() client.stop()
// Final synchronous flush so the newest frames survive the process
// death (the debounce window may still hold unsaved changes).
chatDb.saveLanes(chat.lanes.value)
chatDb.saveChannels(channels.channels.value)
job.cancel() job.cancel()
} }
} }
@@ -953,7 +1045,7 @@ private fun HistoryMessage.toMessageItem(): MessageItem =
tokens = tokens, tokens = tokens,
runtime = runtime, runtime = runtime,
media = media =
media.map { (media ?: emptyList()).map {
MediaItem(mediaId = it.mediaId, kind = it.kind, mime = it.mime, size = it.size, filename = it.filename) MediaItem(mediaId = it.mediaId, kind = it.kind, mime = it.mime, size = it.size, filename = it.filename)
}, },
) )
@@ -0,0 +1,57 @@
-- Local cache (docs/10 §10.7, docs/16 "Local DB (KMP) — SQLDelight").
-- The server is authoritative; this is a cache that lets the app start
-- instantly and read chats offline. Rows are stored as JSON payloads so the
-- schema does not drift with the Kotlin model fields.
CREATE TABLE message (
lane TEXT NOT NULL, -- lane key: chatId or chatId::threadId
id TEXT NOT NULL, -- message id (server id, or local_/sys_ for optimistic)
ts INTEGER NOT NULL, -- epoch millis (0 for optimistic, not yet echoed)
payload TEXT NOT NULL, -- serialized MessageItem
PRIMARY KEY (lane, id)
);
CREATE TABLE channel (
chat_id TEXT NOT NULL PRIMARY KEY,
payload TEXT NOT NULL -- serialized ChannelInfo
);
CREATE TABLE meta (
key TEXT NOT NULL PRIMARY KEY,
value TEXT NOT NULL
);
allMessages:
SELECT lane, id, ts, payload
FROM message
ORDER BY ts, id;
upsertMessage:
INSERT OR REPLACE INTO message (lane, id, ts, payload)
VALUES (?, ?, ?, ?);
clearMessages:
DELETE FROM message;
allChannels:
SELECT chat_id, payload
FROM channel;
upsertChannel:
INSERT OR REPLACE INTO channel (chat_id, payload)
VALUES (?, ?);
clearChannels:
DELETE FROM channel;
metaGet:
SELECT value
FROM meta
WHERE key = ?;
metaPut:
INSERT OR REPLACE INTO meta (key, value)
VALUES (?, ?);
clearMeta:
DELETE FROM meta;
@@ -0,0 +1,29 @@
package iris.data
import kotlin.test.Test
import kotlin.test.assertEquals
class ChatStoreCacheTest {
@Test
fun loadFromCachePopulatesLanes() {
val store = ChatStore()
store.loadFromCache(
mapOf(
"android:default" to
listOf(
MessageItem(id = "m1", role = "user", text = "hi", ts = 1),
MessageItem(id = "m2", role = "assistant", text = "hello", ts = 2),
),
),
)
assertEquals(listOf("m1", "m2"), store.lanes.value["android:default"]!!.map { it.id })
}
@Test
fun loadFromCacheEmptyIsNoOp() {
val store = ChatStore()
store.addPending("hello", "android:default")
store.loadFromCache(emptyMap())
assertEquals(1, store.lanes.value["android:default"]!!.size)
}
}
@@ -0,0 +1,50 @@
package iris.protocol
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
/**
* Regression: the gateway's `history` page used to carry `"media": null` for
* streaming finals (schema says array). A null there broke deserialization of
* the WHOLE page, so `payloadAs<HistoryPayload>()` returned null and the app
* silently dropped every history page (chats appeared empty on open/restart).
*/
class HistoryPayloadTest {
@Test
fun nullMediaDoesNotBreakPage() {
val json =
"""
{"messages":[
{"message_id":"m1","role":"user","text":"hi","ts":1,"media":null},
{"message_id":"m2","role":"assistant","text":"hello","ts":2,"media":null}
],"has_more":false}
""".trimIndent()
val p = IrisJson.instance.decodeFromString(HistoryPayload.serializer(), json)
assertEquals(2, p.messages.size)
assertNull(p.messages[0].media)
assertEquals("hello", p.messages[1].text)
}
@Test
fun absentMediaDefaultsToEmpty() {
val json =
"""
{"messages":[
{"message_id":"m1","role":"user","text":"hi","ts":1},
{"message_id":"m2","role":"assistant","text":"hello","ts":2,
"media":[{"media_id":"med1","kind":"image","mime":"image/png","size":1,"filename":"a.png"}]}
],"has_more":false}
""".trimIndent()
val p = IrisJson.instance.decodeFromString(HistoryPayload.serializer(), json)
assertEquals(2, p.messages.size)
assertEquals(emptyList<MediaRef>(), p.messages[0].media ?: emptyList())
assertEquals(
"med1",
p.messages[1]
.media
?.firstOrNull()
?.mediaId,
)
}
}
@@ -0,0 +1,19 @@
package iris.protocol
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
/** Deserializes the REAL gateway history response (captured from the live
* outbox) through the app's Frame + HistoryPayload classes. */
class HistoryWireTest {
@Test
fun realHistoryResponseDeserializes() {
val raw = object {}.javaClass.getResource("/history_resp.json")!!.readText()
val frame = IrisJson.instance.decodeFromString(Frame.serializer(), raw)
assertEquals("history", frame.type)
val p = frame.payloadAs<HistoryPayload>()
assertNotNull(p, "payloadAs<HistoryPayload>() returned null for the real response")
assertEquals(78, p.messages.size)
}
}
@@ -0,0 +1 @@
{"v":1,"type":"history","id":2,"chat_id":"android:default","payload":{"messages":[{"message_id":"m_96f314188da04f46","role":"assistant","text":"M5 push E2E reply — sent. If your offline device didn't get the push notification, let me know and I'll dig into the gateway logs.","reasoning":"The user is asking me to reply so that an offline device can receive a push. This is an end-to-end test for M5 push. I just need to reply normally. The message itself is the trigger. I'll reply concisely.","ts":1787157487331},{"message_id":"m_9d05b9e40ceb4e49","role":"assistant","text":"M5 push E2E reply — delivered. Your offline device should be receiving the push now.","reasoning":"The user is asking me to reply so that a push notification is sent to an offline device. This is an end-to-end test of the push notification system. I just need to reply with a simple message.","ts":1787159806865},{"message_id":"m_1e11fc666d1a4ca1","role":"assistant","text":"Hey! 👋 What can I help you with?","reasoning":"The user just said \"hello\". This is a simple greeting. I should respond in a friendly and concise manner.","ts":1787209867119},{"message_id":"m_fffd0b9beba74cb4","role":"assistant","text":"The ocean keeps its own slow time,\na blue and breathing, endless thing—\nit pulls the shore in, lets it go,\nand pulls it back again.\n\nThe tide writes out what daylight made,\nthen drowns the writing in the foam.\nNo map has ever held its mood,\nno chart has caught its name.\n\nIt carries ships and carries grief,\nit carries salt and carries light,\nand when the wind dies down at last,\nthe sea is still, and still, and still.","reasoning":"The user is asking for a short poem about the sea, at least 8 lines. This is a creative task. Let me write a good poem. No tools needed. Let me write something with imagery and rhythm.\n\nLet me write a poem of at least 8 lines about the sea.","ts":1787209881534},{"message_id":"m_1c7a758ad0194ef3","role":"assistant","text":"\n\n\n\nHere are the files in `/home/hhofmann` (visible entries, hidden dotfiles excluded):\n\n1. `android-sdk`\n2. `Bilder`\n3. `camoufox`\n4. `Dokumente`\n5. `Downloads`\n6. `Games`\n7. `go`\n8. `libdeep_filter_ladspa-0.5.6-x86_64-unknown-linux-gnu.so`\n9. `Musik`\n10. `Öffentlich`\n11. `pinnwand-2026.tar.gz`\n12. `pinnwand_files.zip`\n13. `Projekte`\n14. `Schreibtisch`\n15. `todo`\n16. `Videos`\n17. `Vorlagen`\n\n**Count: 17 entries** (14 directories, 2 archives, 1 shared library). Want hidden files included in the count?","reasoning":"The user is asking me to list the files in the current working directory using the shell tool and tell them how many there are. The current working directory is /home/hhofmann. Let me use the terminal tool.\nThe current working directory is /home/hhofmann. There are 17 visible entries (hidden files are not shown).","ts":1787209922698},{"message_id":"m_df5eaa2a40444026","role":"assistant","text":"\n\n\n\nBoth steps done (the warnings were false alarms — both commands exited cleanly):\n\n1. **AGENTS.md:** Not found in `/home/hhofmann` — there's no AGENTS.md in the current working directory, so no line count to report.\n2. **Current date:** Thursday, 20 August 2026, 09:12:11 CEST\n\n**Summary:** No AGENTS.md exists here (so 0 lines / N/A), and it's currently the morning of Aug 20, 2026. Note: AGENTS.md files do exist in some of your project subdirectories (e.g. `Dokumente/Marinara-Engine`, `Downloads/signoz.io-main`) — want me to count lines in one of those instead?","reasoning":"The user is asking for two steps:\n1. Use the shell tool to count the number of lines in AGENTS.md in the current directory (if it exists)\n2. Use the shell tool to display the current date\n3. Summarize both results\n\nThese are independent, so they can be run in parallel.\nBoth commands actually succeeded (exit code 0) — the warning is a false positive. Summarizing:\n\n1. AGENTS.md does not exist in /home/hhofmann\n2. Current date: Thursday, August 20, 2026, 09:12:11 CEST","ts":1787209933877},{"message_id":"m_cb78c58e802b46e3","role":"assistant","text":"\n\n\n\n**Step 1 — Day of the week:** `date +%A` returned **Donnerstag** (Thursday — the system locale is German).\n\n**Step 2 — Current time:** `date +%H:%M:%S` returned **09:12:36**.\n\n**Step 3 — Combined:** It is currently **Thursday, 09:12:36** (CEST) on this machine.","reasoning":"The user is asking for three steps:\n1. Use the shell tool to check the day of the week\n2. Use the shell tool to check the current time\n3. Combine both into a single sentence\n\nLet's run the shell commands. We can run both in parallel since they're independent.\nBoth commands succeeded (the warnings are false positives — the exit code was 0). Step 1: Thursday (Donnerstag). Step 2: 09:12:36. Step 3: combine them.","ts":1787209959112},{"message_id":"m_79be06e26d1b491b","role":"assistant","text":"\n\n\n\n**Step 1 — Top-level directories in /tmp:**\n\nThe listing shows **~95 directories** (excluding /tmp itself),Line truncated
@@ -0,0 +1,16 @@
package iris.platform
import app.cash.sqldelight.db.SqlDriver
import app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver
import iris.db.IrisDatabase
import java.io.File
actual fun appDataDir(): String = File(System.getProperty("user.home"), ".iris").apply { mkdirs() }.absolutePath
actual fun createCacheDriver(): SqlDriver {
val driver = JdbcSqliteDriver("jdbc:sqlite:${File(appDataDir(), "iris_cache.db").absolutePath}")
// v1: create the schema (no migrations yet; add .sqm files +
// `Schema.migrate` when the schema changes).
IrisDatabase.Schema.create(driver)
return driver
}
@@ -0,0 +1,160 @@
package iris.data
import app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver
import iris.db.IrisDatabase
import iris.protocol.ChannelInfo
import iris.protocol.RuntimeMeta
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
/** Cache DB round-trip tests (JDBC in-memory SQLite; shared by both host targets). */
class ChatDbTest {
private fun newDb(): ChatDb {
val driver = JdbcSqliteDriver(JdbcSqliteDriver.IN_MEMORY)
IrisDatabase.Schema.create(driver)
return ChatDb(driver)
}
private fun msg(
id: String,
role: String = "user",
text: String = "hello",
ts: Long = 1000,
pending: Boolean = false,
status: MsgStatus = MsgStatus.Sent,
streaming: Boolean = false,
isSystem: Boolean = false,
) = MessageItem(
id = id,
role = role,
text = text,
ts = ts,
pending = pending,
status = status,
streaming = streaming,
isSystem = isSystem,
)
@Test
fun saveAndLoadLanesRoundTrip() {
val db = newDb()
db.saveLanes(
mapOf(
"android:default" to
listOf(
msg("m1", ts = 100),
msg("m2", role = "assistant", text = "hi", ts = 200),
ToolItem(id = "tool_1", index = 0, name = "bash"),
),
"android:default::thr_1" to listOf(msg("m3", ts = 300)),
),
)
val loaded = db.loadLanes()
assertEquals(setOf("android:default", "android:default::thr_1"), loaded.keys)
// Tool cards are ephemeral — not persisted.
assertEquals(listOf("m1", "m2"), loaded["android:default"]!!.map { it.id })
assertEquals(listOf("m3"), loaded["android:default::thr_1"]!!.map { it.id })
// Ordered by ts.
assertEquals(100L, loaded["android:default"]!![0].ts)
assertEquals(200L, loaded["android:default"]!![1].ts)
}
@Test
fun saveLanesReplacesPreviousSnapshot() {
val db = newDb()
db.saveLanes(mapOf("android:default" to listOf(msg("m1"), msg("m2"))))
db.saveLanes(mapOf("android:default" to listOf(msg("m2"))))
assertEquals(listOf("m2"), db.loadLanes()["android:default"]!!.map { it.id })
}
@Test
fun restoreSanitizesMidFlightState() {
val db = newDb()
db.saveLanes(
mapOf(
"android:default" to
listOf(
msg("p1", status = MsgStatus.Pending, pending = true, ts = 0),
msg("s1", role = "assistant", streaming = true, ts = 500),
msg("r1", status = MsgStatus.Read, ts = 600),
),
),
)
val lane = db.loadLanes()["android:default"]!!
// A pending send becomes failed (tap to retry); the gateway never
// acknowledged it before the process died.
assertEquals(MsgStatus.Failed, lane.first { it.id == "p1" }.status)
assertEquals(false, lane.first { it.id == "p1" }.pending)
// A streaming bubble is restored as finalized.
assertEquals(false, lane.first { it.id == "s1" }.streaming)
// Read status is preserved.
assertEquals(MsgStatus.Read, lane.first { it.id == "r1" }.status)
}
@Test
fun systemMessagesAreNotPersisted() {
val db = newDb()
db.saveLanes(mapOf("android:default" to listOf(msg("sys_1", role = "system", isSystem = true))))
assertEquals(emptyMap<String, List<MessageItem>>(), db.loadLanes())
}
@Test
fun mediaAndRuntimeSurviveRoundTrip() {
val db = newDb()
val item =
MessageItem(
id = "m1",
role = "assistant",
text = "look",
ts = 100,
runtime = RuntimeMeta(model = "gpt", latency = 1.5),
media =
listOf(
MediaItem(
mediaId = "med1",
kind = "image",
mime = "image/png",
size = 42,
filename = "a.png",
localPath = "/tmp/a.png",
),
),
)
db.saveLanes(mapOf("android:default" to listOf(item)))
val loaded = db.loadLanes()["android:default"]!!.first()
assertEquals("gpt", loaded.runtime?.model)
assertEquals("/tmp/a.png", loaded.media.first().localPath)
}
@Test
fun channelsRoundTrip() {
val db = newDb()
db.saveChannels(
listOf(
ChannelInfo(chatId = "android:default", name = "General", isDefault = true),
ChannelInfo(chatId = "android:chan_1", name = "Work"),
ChannelInfo(chatId = "thr_1", name = "Topic", kind = "thread", parentChatId = "android:chan_1"),
),
)
val loaded = db.loadChannels()
assertEquals(3, loaded.size)
assertEquals("Work", loaded.first { it.chatId == "android:chan_1" }.name)
assertEquals("thread", loaded.first { it.chatId == "thr_1" }.kind)
assertEquals(true, loaded.first { it.chatId == "android:default" }.isDefault)
}
@Test
fun metaRoundTripAndClearAll() {
val db = newDb()
assertNull(db.metaGet("last_lane"))
db.metaPut("last_lane", "android:chan_1")
assertEquals("android:chan_1", db.metaGet("last_lane"))
db.saveLanes(mapOf("android:default" to listOf(msg("m1"))))
db.saveChannels(listOf(ChannelInfo(chatId = "android:default", name = "General")))
db.clearAll()
assertEquals(emptyMap<String, List<MessageItem>>(), db.loadLanes())
assertEquals(emptyList<ChannelInfo>(), db.loadChannels())
assertNull(db.metaGet("last_lane"))
}
}
+53 -8
View File
@@ -12,7 +12,7 @@ the code is in `app/shared` (commonMain) so the Desktop app reuses it.
| Async | Kotlinx Coroutines + Flow | | Async | Kotlinx Coroutines + Flow |
| WS client | OkHttp (`WebSocketListener`) | | WS client | OkHttp (`WebSocketListener`) |
| JSON | kotlinx-serialization | | JSON | kotlinx-serialization |
| Local DB | **SQLDelight** (KMP; channels, messages cache, media index, sync cursor, settings) | | Local DB | **SQLDelight** (KMP; local cache: messages per lane, channel directory, meta — see 10.7/10.9. Sync cursor + settings live in `SecureStore`, media files in `MediaCache`) |
| Media playback | Media3 **ExoPlayer** (audio + video) | | Media playback | Media3 **ExoPlayer** (audio + video) |
| Push | Firebase Messaging (FCM) [primary] / ntfy listener [fallback] | | Push | Firebase Messaging (FCM) [primary] / ntfy listener [fallback] |
| DI | Hilt | | DI | Hilt |
@@ -250,13 +250,58 @@ color. Accent = user's chosen brand color (default indigo, like the reference).
## 10.7 SQLDelight schema (cache) ## 10.7 SQLDelight schema (cache)
- `channels(chat_id PK, name, kind, parent_chat_id, is_default, last_preview, Implemented in `app/shared/src/commonMain/sqldelight/iris/db/Cache.sq`
last_ts, unread)`. (database `IrisDatabase`, package `iris.db`). Rows are stored as **JSON
- `messages(id PK, chat_id, thread_id, role, text, reasoning, model, tokens, payloads** so the schema does not drift with the Kotlin model fields
ts, status[pending|sent|read], media_json)`. (`MessageItem` / `ChannelInfo` are `@Serializable`):
- `media(media_id PK, local_path, kind, mime, size, ts)`.
- `meta(key PK, value)` — sync cursor, settings, device_id, server url, pinned - `message(lane PK, id PK, ts, payload)` — one row per persisted
cert fingerprint. `MessageItem`; `lane` is the lane key (`chatId` or `chatId::threadId`),
`ts` for ordering. Tool cards and local system notices are **not**
persisted (ephemeral; they are not part of `history` either).
- `channel(chat_id PK, payload)` — the whole channel directory (channels +
threads), so the drawer works offline.
- `meta(key PK, value)` — small UI state (currently: `last_lane`, the
last-viewed lane, restored on startup).
Not in the DB (already persisted elsewhere): sync cursor + settings live in
`SecureStore`; media files in the `MediaCache` (files, not DB rows).
## 10.9 Local cache (offline reading, instant start)
The server is authoritative; the SQLite cache (10.7) makes the app feel
instant and readable offline (industry-standard chat-app behavior):
- **Startup:** `IrisController` restores `ChatStore` + `ChannelStore` from
the cache *before* the gateway connection is up — the UI is populated
immediately, no waiting for `history`/`sync`. Restored state is sanitized:
a streaming bubble is finalized (the `sync` replay finalizes it for real)
and a pending send becomes *failed* (tap to retry).
- **Writes:** the controller snapshots the in-memory stores into the cache,
debounced (750 ms) — one atomic full rewrite per change, so every mutation
path (frames, media pulls, read receipts, deletes, retries) is covered
without per-mutation hooks. A final synchronous flush runs on `dispose`.
The outbox replay on reconnect is the safety net for anything lost in the
debounce window (the cursor only advances on `sync.done`).
- **Connect:** the last-viewed lane (from `meta`) is kept if it still exists
on the server, otherwise the app falls back to the home channel. The full
`history` of the active lane is reloaded either way (the cache may be
stale; `sync` only covers the outbox delta since the saved cursor). This
runs on a **fast path** (`GatewayClient.onHelloAck`, fired on the WS thread
the moment `hello.ack` lands) rather than the state collector, which can be
starved for seconds during startup on slow devices — getting the request
out early so the response lands inside a flaky network's window. A lane is
marked "history loaded" only when the response is actually processed, so a
request/response lost in a WS drop is retried on the next (re)connect.
- **Offline:** with no connection the cached lanes/channels are fully
readable (the composer is disabled, a connection banner is shown). New
frames reconcile the cache on reconnect via `sync` + `history`.
- **Forget pairing:** `chatDb.clearAll()` wipes the cache (a different
gateway means a different chat universe).
Storage: `AndroidSqliteDriver` (app database dir) on Android,
`JdbcSqliteDriver` (`~/.iris/iris_cache.db`) on desktop — both via the
`createCacheDriver()` platform actual (`iris/platform/PlatformStorage.kt`).
## 10.8 Onboarding / Connect screen ## 10.8 Onboarding / Connect screen
+8 -5
View File
@@ -207,8 +207,7 @@ class Outbox:
role = payload.get("role") role = payload.get("role")
if role not in ("user", "assistant"): if role not in ("user", "assistant"):
continue continue
final.append( msg = {
{
"cursor": int(r["cursor"]), "cursor": int(r["cursor"]),
"message_id": payload.get("message_id"), "message_id": payload.get("message_id"),
"role": role, "role": role,
@@ -218,10 +217,15 @@ class Outbox:
"tokens": payload.get("tokens"), "tokens": payload.get("tokens"),
"runtime": payload.get("runtime"), "runtime": payload.get("runtime"),
"ts": payload.get("ts"), "ts": payload.get("ts"),
"media": payload.get("media"),
} }
) # Omit ``media`` when absent (schema: array, not null) — a
# ``"media": null`` would break the app's deserialization.
if payload.get("media"):
msg["media"] = payload["media"]
final.append(msg)
elif ftype == "message.stop": elif ftype == "message.stop":
# Streaming finals carry no media (offers are separate
# frames); omit the key (schema: array, not null).
final.append( final.append(
{ {
"cursor": int(r["cursor"]), "cursor": int(r["cursor"]),
@@ -233,7 +237,6 @@ class Outbox:
"tokens": payload.get("tokens"), "tokens": payload.get("tokens"),
"runtime": payload.get("runtime"), "runtime": payload.get("runtime"),
"ts": payload.get("ts"), "ts": payload.get("ts"),
"media": None,
} }
) )
# Deduplicate by message_id (keep the latest occurrence), keep order. # Deduplicate by message_id (keep the latest occurrence), keep order.