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).
- 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).
- 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`.
- 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("org.jetbrains.compose") version "1.11.1" 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") {
+28
View File
@@ -6,6 +6,8 @@ plugins {
id("com.android.kotlin.multiplatform.library")
id("org.jetbrains.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"
@@ -25,6 +27,8 @@ val kcefVersion = "2025.03.23"
// release). Provides GFM tables, bold/italic/underscore, and the -code module
// for language-aware syntax highlighting (Highlights).
val markdownVersion = "0.44.0"
// Local cache DB (messages/channels/meta; docs/10 §10.7, docs/16).
val sqldelightVersion = "2.3.2"
kotlin {
android {
@@ -51,6 +55,11 @@ kotlin {
val jvmMain by creating { dependsOn(commonMain.get()) }
val androidMain 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 {
implementation("org.jetbrains.compose.runtime:runtime:$composeVersion")
@@ -74,15 +83,24 @@ kotlin {
// M9: HTML artifact previews — WebView composable (platform WebView
// on Android, KCEF/JCEF on desktop).
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 {
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
// webview library exposes it transitively, but we reference KCEF
// directly (init + progress) so declare it explicitly.
desktopMain.dependencies {
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).
androidMain.dependencies {
@@ -97,6 +115,16 @@ kotlin {
// the ntfy listener is the fallback). The google-services plugin is
// applied conditionally in the app module.
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)
}
/** 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. */
fun onFrame(frame: Frame) {
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.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.serialization.Serializable
import kotlinx.serialization.json.JsonElement
import kotlin.random.Random
@@ -54,7 +55,9 @@ enum class MsgStatus {
/** A chat message (user / assistant / commentary / streaming bubble).
* [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(
override val id: String,
val role: String,
@@ -77,6 +80,7 @@ data class MessageItem(
* has been pulled to the local cache (outbound) — inbound attachments the
* app itself uploaded carry no local path (the agent reads the server copy).
*/
@Serializable
data class MediaItem(
val mediaId: 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() {
_lanes.value = emptyMap()
}
@@ -30,29 +30,29 @@ import kotlinx.coroutines.Job
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.delay
import kotlinx.coroutines.isActive
import kotlinx.coroutines.launch
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.SharedFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.isActive
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withTimeout
import kotlinx.coroutines.withTimeoutOrNull
import kotlin.time.TimeMark
import kotlin.time.TimeSource
import okhttp3.OkHttpClient
import okio.ByteString
import okio.ByteString.Companion.toByteString
import okhttp3.Request
import okhttp3.Response
import okhttp3.WebSocket
import okhttp3.WebSocketListener
import okio.ByteString
import okio.ByteString.Companion.toByteString
import java.util.concurrent.TimeUnit
import kotlin.random.Random
import kotlin.time.TimeMark
import kotlin.time.TimeSource
/**
* OkHttp WebSocket client for the hermes android gateway (docs/10 §10.3).
@@ -69,7 +69,9 @@ class GatewayClient(
) {
sealed interface State {
data object Disconnected : State
data object Connecting : State
data class Connected(
val caps: ServerCaps,
val channels: List<ChannelInfo>,
@@ -78,8 +80,12 @@ class GatewayClient(
* it must not re-post system notifications (docs/08 §8.7). */
val lastPushedCursor: Long = 0,
) : 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)
@@ -92,7 +98,9 @@ class GatewayClient(
private val _events = MutableSharedFlow<Frame>(extraBufferCapacity = 128)
val events: SharedFlow<Frame> = _events.asSharedFlow()
private val client: OkHttpClient = OkHttpClient.Builder()
private val client: OkHttpClient =
OkHttpClient
.Builder()
.pingInterval(20, TimeUnit.SECONDS)
.build()
@@ -100,6 +108,7 @@ class GatewayClient(
private var socket: WebSocket? = null
private var nextRequestId = 1
private var attempt = 0
// True once a connection has been established this session; reset by
// start(). Drives Connecting (first dial) vs Reconnecting (redial after a
// drop) so the UI can show the right status without a blocking screen.
@@ -120,6 +129,15 @@ class GatewayClient(
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
// slot). Serialize concurrent offers so their byte streams don't interleave.
private val pullMutex = Mutex()
@@ -165,6 +183,7 @@ class GatewayClient(
dial.socket.close(1000, "auth failed")
return
}
DialResult.Connected -> {
hasConnected = true
attempt = 0
@@ -173,6 +192,7 @@ class GatewayClient(
if (!currentCoroutineContext().isActive) return
// socket dropped -> loop again (Reconnecting)
}
is DialResult.Failed -> {
attempt++
delay(backoffMs(attempt))
@@ -185,8 +205,14 @@ class GatewayClient(
private sealed interface 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(
@@ -195,17 +221,24 @@ class GatewayClient(
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 helloAck = CompletableDeferred<HelloAckPayload>()
val authError = CompletableDeferred<String>()
val fail = CompletableDeferred<String>()
val request = Request.Builder().url(url).build()
val ws = client.newWebSocket(
val ws =
client.newWebSocket(
request,
object : WebSocketListener() {
override fun onOpen(webSocket: WebSocket, response: Response) {
override fun onOpen(
webSocket: WebSocket,
response: Response,
) {
IrisLog.d("ws open (${response.code})")
webSocket.send(
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()
val frame = try {
val frame =
try {
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
}
when (frame.type) {
@@ -230,6 +270,7 @@ class GatewayClient(
val ack = frame.payloadAs<HelloAckPayload>()
if (ack != null) helloAck.complete(ack)
}
TYPE_ERROR -> {
val err = frame.payloadAs<ErrorPayload>()
if (!authError.isCompleted) authError.complete(err?.message ?: "auth failed")
@@ -237,7 +278,11 @@ class GatewayClient(
// auth failures — let the controller react.
_events.tryEmit(frame)
}
TYPE_PONG -> Unit
TYPE_PONG -> {
Unit
}
else -> {
_events.tryEmit(frame)
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()
// M4: binary frames belong to the active pull stream
// (uploads are outbound; stray inbound chunks are dropped).
@@ -264,12 +312,20 @@ class GatewayClient(
?.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\"")
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})")
fail.complete(t.message ?: "connection failed")
closed.complete(Unit)
@@ -282,13 +338,16 @@ class GatewayClient(
helloAck.invokeOnCompletion { e ->
if (e == null) {
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.
val local = store.syncCursor
if (local < ack.syncCursor) {
val id = nextRequestId++
ws.send(syncFrame(id, local).toWire())
}
// Prompt fast path (before the possibly-starved state collector).
onHelloAck?.invoke(connected)
winner.complete(DialResult.Connected)
}
}
@@ -298,7 +357,8 @@ class GatewayClient(
fail.invokeOnCompletion { e ->
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")
return Dial(result, ws, closed)
}
@@ -363,14 +423,21 @@ class GatewayClient(
return when (frame.type) {
TYPE_MEDIA_UPLOAD_ACK -> {
val p = frame.payloadAs<MediaUploadAckPayload>()
if (p != null && p.ok) Result.success(p.mediaRef)
else Result.failure(IllegalStateException("upload rejected by server"))
if (p != null && p.ok) {
Result.success(p.mediaRef)
} else {
Result.failure(IllegalStateException("upload rejected by server"))
}
}
TYPE_ERROR -> {
val e = frame.payloadAs<ErrorPayload>()
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) {
return Result.failure(e)
@@ -383,7 +450,10 @@ class GatewayClient(
* 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).
*/
suspend fun pullMedia(mediaId: String, onChunk: suspend (ByteArray) -> Unit): Result<Unit> =
suspend fun pullMedia(
mediaId: String,
onChunk: suspend (ByteArray) -> Unit,
): Result<Unit> =
pullMutex.withLock {
val ws = socket ?: return@withLock Result.failure(IllegalStateException("not connected"))
val id = nextRequestId++
@@ -393,21 +463,29 @@ class GatewayClient(
binarySession = BinarySession.Pulling(id, chunks, end)
try {
ws.send(mediaPullFrame(id, mediaId).toWire())
val frame = withTimeout(PULL_TIMEOUT_MS) {
val frame =
withTimeout(PULL_TIMEOUT_MS) {
for (chunk in chunks) onChunk(chunk)
end.await()
}
when (frame.type) {
TYPE_MEDIA_PULL_END -> {
val p = frame.payloadAs<MediaPullEndPayload>()
if (p != null && p.ok) Result.success(Unit)
else Result.failure(IllegalStateException("pull failed"))
if (p != null && p.ok) {
Result.success(Unit)
} else {
Result.failure(IllegalStateException("pull failed"))
}
}
TYPE_ERROR -> {
val e = frame.payloadAs<ErrorPayload>()
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) {
Result.failure(e)
@@ -459,15 +537,24 @@ class GatewayClient(
* Real `hello` test: dial, wait for hello.ack (or auth error), close.
* 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)
return when (val result = dial.result) {
DialResult.Connected -> {
dial.socket.close(1000, "test complete")
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 runtime: RuntimeMeta? = 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
@@ -1,6 +1,7 @@
package iris.state
import iris.data.ChannelStore
import iris.data.ChatDb
import iris.data.ChatStore
import iris.data.MediaItem
import iris.data.MessageItem
@@ -10,6 +11,7 @@ import iris.media.MediaCache
import iris.media.kindFromMime
import iris.net.GatewayClient
import iris.platform.PickedFile
import iris.platform.createCacheDriver
import iris.platform.isAppForeground
import iris.platform.mediaCacheBaseDir
import iris.platform.postSystemNotification
@@ -73,6 +75,7 @@ import iris.protocol.syncFrame
import iris.ui.theme.Backdrop
import iris.ui.theme.BackgroundMode
import iris.ui.theme.UserTheme
import iris.util.IrisLog
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
@@ -80,6 +83,7 @@ import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.debounce
import kotlinx.coroutines.launch
import kotlin.random.Random
@@ -98,6 +102,12 @@ class IrisController(
val chat = ChatStore()
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)
val typing: StateFlow<Boolean> = _typing.asStateFlow()
@@ -210,6 +220,13 @@ class IrisController(
const val FONT_SCALE_MIN = 0.8f
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). */
val RUNTIME_FIELD_KEYS = listOf("model", "context_pct", "cwd", "latency", "cost")
@@ -414,9 +431,28 @@ class IrisController(
}
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
// 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 {
client.events.collect { frame ->
try {
when (frame.type) {
TYPE_MESSAGE,
TYPE_MESSAGE_START,
@@ -527,14 +563,27 @@ class IrisController(
// since the saved cursor, so after a process death the
// in-memory store is empty and older messages are not in
// the delta — history loads the full list.
frame.payloadAs<HistoryPayload>()?.let { p ->
val (chatId, threadId) =
(frame.chatId ?: return@let) to frame.threadId
val p = frame.payloadAs<HistoryPayload>()
if (p == null) {
// 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)
IrisLog.d("history loaded: lane=$lane messages=${p.messages.size}")
chat.loadHistory(
lane,
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
}
}
} 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 {
@@ -592,20 +647,11 @@ class IrisController(
val prev = prevState
prevState = s
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).
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.
if (s.caps.pushNtfyServer.isNotBlank()) {
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)}"
}
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()
}
/**
* 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) ───────────────────────────────────
/** 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)
if (lane in historyLoaded) return
historyLoaded.add(lane)
// 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))
}
@@ -930,6 +1016,8 @@ class IrisController(
fun forget() {
store.clear()
// A different gateway means a different chat universe — wipe the cache.
chatDb.clearAll()
chat.clear()
channels.clear()
client.restart()
@@ -937,6 +1025,10 @@ class IrisController(
fun dispose() {
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()
}
}
@@ -953,7 +1045,7 @@ private fun HistoryMessage.toMessageItem(): MessageItem =
tokens = tokens,
runtime = runtime,
media =
media.map {
(media ?: emptyList()).map {
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 |
| WS client | OkHttp (`WebSocketListener`) |
| 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) |
| Push | Firebase Messaging (FCM) [primary] / ntfy listener [fallback] |
| DI | Hilt |
@@ -250,13 +250,58 @@ color. Accent = user's chosen brand color (default indigo, like the reference).
## 10.7 SQLDelight schema (cache)
- `channels(chat_id PK, name, kind, parent_chat_id, is_default, last_preview,
last_ts, unread)`.
- `messages(id PK, chat_id, thread_id, role, text, reasoning, model, tokens,
ts, status[pending|sent|read], media_json)`.
- `media(media_id PK, local_path, kind, mime, size, ts)`.
- `meta(key PK, value)` — sync cursor, settings, device_id, server url, pinned
cert fingerprint.
Implemented in `app/shared/src/commonMain/sqldelight/iris/db/Cache.sq`
(database `IrisDatabase`, package `iris.db`). Rows are stored as **JSON
payloads** so the schema does not drift with the Kotlin model fields
(`MessageItem` / `ChannelInfo` are `@Serializable`):
- `message(lane PK, id PK, ts, payload)` — one row per persisted
`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
+8 -5
View File
@@ -207,8 +207,7 @@ class Outbox:
role = payload.get("role")
if role not in ("user", "assistant"):
continue
final.append(
{
msg = {
"cursor": int(r["cursor"]),
"message_id": payload.get("message_id"),
"role": role,
@@ -218,10 +217,15 @@ class Outbox:
"tokens": payload.get("tokens"),
"runtime": payload.get("runtime"),
"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":
# Streaming finals carry no media (offers are separate
# frames); omit the key (schema: array, not null).
final.append(
{
"cursor": int(r["cursor"]),
@@ -233,7 +237,6 @@ class Outbox:
"tokens": payload.get("tokens"),
"runtime": payload.get("runtime"),
"ts": payload.get("ts"),
"media": None,
}
)
# Deduplicate by message_id (keep the latest occurrence), keep order.