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
This commit is contained in:
1 parent
9a519e3c5a
commit
81f42ab761
18 files changed
+1040
-268
No files matched your search
@@ -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,14 +98,17 @@ class GatewayClient(
|
||||
private val _events = MutableSharedFlow<Frame>(extraBufferCapacity = 128)
|
||||
val events: SharedFlow<Frame> = _events.asSharedFlow()
|
||||
|
||||
private val client: OkHttpClient = OkHttpClient.Builder()
|
||||
.pingInterval(20, TimeUnit.SECONDS)
|
||||
.build()
|
||||
private val client: OkHttpClient =
|
||||
OkHttpClient
|
||||
.Builder()
|
||||
.pingInterval(20, TimeUnit.SECONDS)
|
||||
.build()
|
||||
|
||||
private var connectJob: Job? = null
|
||||
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,100 +221,133 @@ 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(
|
||||
request,
|
||||
object : WebSocketListener() {
|
||||
override fun onOpen(webSocket: WebSocket, response: Response) {
|
||||
IrisLog.d("ws open (${response.code})")
|
||||
webSocket.send(
|
||||
helloFrame(
|
||||
token = token,
|
||||
deviceId = store.deviceId,
|
||||
deviceName = store.deviceName,
|
||||
fcmToken = store.fcmToken.ifBlank { null },
|
||||
ntfyTopic = store.ntfyTopic.ifBlank { null },
|
||||
).toWire(),
|
||||
)
|
||||
}
|
||||
|
||||
override fun onMessage(webSocket: WebSocket, text: String) {
|
||||
lastLiveness = TimeSource.Monotonic.markNow()
|
||||
val frame = try {
|
||||
IrisJson.instance.decodeFromString(Frame.serializer(), text)
|
||||
} catch (_: Exception) {
|
||||
return
|
||||
val ws =
|
||||
client.newWebSocket(
|
||||
request,
|
||||
object : WebSocketListener() {
|
||||
override fun onOpen(
|
||||
webSocket: WebSocket,
|
||||
response: Response,
|
||||
) {
|
||||
IrisLog.d("ws open (${response.code})")
|
||||
webSocket.send(
|
||||
helloFrame(
|
||||
token = token,
|
||||
deviceId = store.deviceId,
|
||||
deviceName = store.deviceName,
|
||||
fcmToken = store.fcmToken.ifBlank { null },
|
||||
ntfyTopic = store.ntfyTopic.ifBlank { null },
|
||||
).toWire(),
|
||||
)
|
||||
}
|
||||
when (frame.type) {
|
||||
TYPE_HELLO_ACK -> {
|
||||
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")
|
||||
// M7: post-connect error frames are app events, not
|
||||
// auth failures — let the controller react.
|
||||
_events.tryEmit(frame)
|
||||
}
|
||||
TYPE_PONG -> Unit
|
||||
else -> {
|
||||
_events.tryEmit(frame)
|
||||
frame.id?.let { id ->
|
||||
pending[id]?.complete(frame)
|
||||
// M4: terminal frame of a pull stream — close
|
||||
// the chunk channel so the pull loop exits.
|
||||
if (frame.type == TYPE_MEDIA_PULL_END) {
|
||||
(binarySession as? BinarySession.Pulling)?.let {
|
||||
it.chunks.close()
|
||||
binarySession = null
|
||||
|
||||
override fun onMessage(
|
||||
webSocket: WebSocket,
|
||||
text: String,
|
||||
) {
|
||||
lastLiveness = TimeSource.Monotonic.markNow()
|
||||
val frame =
|
||||
try {
|
||||
IrisJson.instance.decodeFromString(Frame.serializer(), text)
|
||||
} 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) {
|
||||
TYPE_HELLO_ACK -> {
|
||||
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")
|
||||
// M7: post-connect error frames are app events, not
|
||||
// auth failures — let the controller react.
|
||||
_events.tryEmit(frame)
|
||||
}
|
||||
|
||||
TYPE_PONG -> {
|
||||
Unit
|
||||
}
|
||||
|
||||
else -> {
|
||||
_events.tryEmit(frame)
|
||||
frame.id?.let { id ->
|
||||
pending[id]?.complete(frame)
|
||||
// M4: terminal frame of a pull stream — close
|
||||
// the chunk channel so the pull loop exits.
|
||||
if (frame.type == TYPE_MEDIA_PULL_END) {
|
||||
(binarySession as? BinarySession.Pulling)?.let {
|
||||
it.chunks.close()
|
||||
binarySession = null
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
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).
|
||||
(binarySession as? BinarySession.Pulling)
|
||||
?.chunks
|
||||
?.trySend(bytes.toByteArray())
|
||||
}
|
||||
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).
|
||||
(binarySession as? BinarySession.Pulling)
|
||||
?.chunks
|
||||
?.trySend(bytes.toByteArray())
|
||||
}
|
||||
|
||||
override fun onClosed(webSocket: WebSocket, code: Int, reason: String) {
|
||||
IrisLog.w("ws closed code=$code reason=\"$reason\"")
|
||||
closed.complete(Unit)
|
||||
}
|
||||
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?) {
|
||||
IrisLog.e("ws failure: ${t.javaClass.simpleName}: ${t.message} (http=${response?.code})")
|
||||
fail.complete(t.message ?: "connection failed")
|
||||
closed.complete(Unit)
|
||||
}
|
||||
},
|
||||
)
|
||||
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)
|
||||
}
|
||||
},
|
||||
)
|
||||
socket = ws
|
||||
|
||||
val winner = CompletableDeferred<DialResult>()
|
||||
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,17 +357,18 @@ class GatewayClient(
|
||||
fail.invokeOnCompletion { e ->
|
||||
if (e == null) winner.complete(DialResult.Failed(fail.getCompleted()))
|
||||
}
|
||||
val result = withTimeoutOrNull(15_000) { winner.await() }
|
||||
?: DialResult.Failed("timeout waiting for hello.ack")
|
||||
val result =
|
||||
withTimeoutOrNull(15_000) { winner.await() }
|
||||
?: DialResult.Failed("timeout waiting for hello.ack")
|
||||
return Dial(result, ws, closed)
|
||||
}
|
||||
|
||||
// ── Outbound ──────────────────────────────────────────────────────────
|
||||
|
||||
/** Send a text message (fire-and-forget; the server echoes it back).
|
||||
* M4: [mediaRefs] reference completed uploads (media.upload.ack refs).
|
||||
* [autoThread] asks the gateway to mint a fresh thread for the message
|
||||
* (auto-threading, docs/06 §6.3). */
|
||||
* M4: [mediaRefs] reference completed uploads (media.upload.ack refs).
|
||||
* [autoThread] asks the gateway to mint a fresh thread for the message
|
||||
* (auto-threading, docs/06 §6.3). */
|
||||
fun sendMessage(
|
||||
chatId: String,
|
||||
text: String,
|
||||
@@ -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) {
|
||||
for (chunk in chunks) onChunk(chunk)
|
||||
end.await()
|
||||
}
|
||||
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))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -494,4 +581,4 @@ class GatewayClient(
|
||||
}
|
||||
}
|
||||
|
||||
private fun Frame.toWire(): String = IrisJson.instance.encodeToString(Frame.serializer(), this)
|
||||
private fun Frame.toWire(): String = IrisJson.instance.encodeToString(Frame.serializer(), this)
|
||||
@@ -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,175 +431,213 @@ 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 ->
|
||||
when (frame.type) {
|
||||
TYPE_MESSAGE,
|
||||
TYPE_MESSAGE_START,
|
||||
TYPE_MESSAGE_UPDATE,
|
||||
TYPE_MESSAGE_STOP,
|
||||
TYPE_TOOL_START,
|
||||
TYPE_TOOL_PROGRESS,
|
||||
TYPE_TOOL_END,
|
||||
TYPE_COMMENTARY,
|
||||
TYPE_MEDIA_OFFER,
|
||||
-> {
|
||||
chat.onFrame(frame)
|
||||
// M4: pull offered media into the local cache.
|
||||
if (frame.type == TYPE_MEDIA_OFFER) {
|
||||
frame.payloadAs<MediaOfferPayload>()?.let { pullMedia(it) }
|
||||
try {
|
||||
when (frame.type) {
|
||||
TYPE_MESSAGE,
|
||||
TYPE_MESSAGE_START,
|
||||
TYPE_MESSAGE_UPDATE,
|
||||
TYPE_MESSAGE_STOP,
|
||||
TYPE_TOOL_START,
|
||||
TYPE_TOOL_PROGRESS,
|
||||
TYPE_TOOL_END,
|
||||
TYPE_COMMENTARY,
|
||||
TYPE_MEDIA_OFFER,
|
||||
-> {
|
||||
chat.onFrame(frame)
|
||||
// M4: pull offered media into the local cache.
|
||||
if (frame.type == TYPE_MEDIA_OFFER) {
|
||||
frame.payloadAs<MediaOfferPayload>()?.let { pullMedia(it) }
|
||||
}
|
||||
// M5: Telegram/WhatsApp-style — when the app is
|
||||
// backgrounded, a finalized assistant reply posts a
|
||||
// system notification (the message still lands in the
|
||||
// chat via onFrame above). Streaming turns finalize on
|
||||
// message.stop; non-streaming replies are a single
|
||||
// assistant `message`.
|
||||
when (frame.type) {
|
||||
TYPE_MESSAGE_STOP -> {
|
||||
frame.payloadAs<MessageStopPayload>()?.let {
|
||||
if (!isPushedReplay(frame)) {
|
||||
notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.finalText)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
TYPE_MESSAGE -> {
|
||||
frame.payloadAs<MessagePayload>()?.let {
|
||||
if (it.role == ROLE_ASSISTANT && !isPushedReplay(frame)) {
|
||||
notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.text)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
else -> {
|
||||
Unit
|
||||
}
|
||||
}
|
||||
}
|
||||
// M5: Telegram/WhatsApp-style — when the app is
|
||||
// backgrounded, a finalized assistant reply posts a
|
||||
// system notification (the message still lands in the
|
||||
// chat via onFrame above). Streaming turns finalize on
|
||||
// message.stop; non-streaming replies are a single
|
||||
// assistant `message`.
|
||||
when (frame.type) {
|
||||
TYPE_MESSAGE_STOP -> {
|
||||
frame.payloadAs<MessageStopPayload>()?.let {
|
||||
if (!isPushedReplay(frame)) {
|
||||
notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.finalText)
|
||||
|
||||
TYPE_CHANNEL_CREATED,
|
||||
TYPE_CHANNEL_RENAMED,
|
||||
TYPE_CHANNEL_DELETED,
|
||||
TYPE_CHANNEL_LIST,
|
||||
-> {
|
||||
channels.onFrame(frame)
|
||||
// Auto-created thread (Settings → "Threads"): the
|
||||
// gateway minted it for a message in the flat lane we
|
||||
// are viewing — jump into it; the user's message and
|
||||
// the reply land there.
|
||||
if (frame.type == TYPE_CHANNEL_CREATED) {
|
||||
frame.payloadAs<ChannelInfo>()?.let { info ->
|
||||
if (info.kind == "thread" && info.auto && info.parentChatId != null) {
|
||||
val (curChat, curThread) = chat.parseLane(chat.currentLane.value)
|
||||
if (curThread == null && curChat == info.parentChatId) {
|
||||
val lane = chat.laneKey(info.parentChatId, info.chatId)
|
||||
chat.setLane(lane)
|
||||
historyLoaded.add(lane)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
TYPE_MESSAGE -> {
|
||||
frame.payloadAs<MessagePayload>()?.let {
|
||||
if (it.role == ROLE_ASSISTANT && !isPushedReplay(frame)) {
|
||||
notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.text)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
else -> {
|
||||
Unit
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
TYPE_CHANNEL_CREATED,
|
||||
TYPE_CHANNEL_RENAMED,
|
||||
TYPE_CHANNEL_DELETED,
|
||||
TYPE_CHANNEL_LIST,
|
||||
-> {
|
||||
channels.onFrame(frame)
|
||||
// Auto-created thread (Settings → "Threads"): the
|
||||
// gateway minted it for a message in the flat lane we
|
||||
// are viewing — jump into it; the user's message and
|
||||
// the reply land there.
|
||||
if (frame.type == TYPE_CHANNEL_CREATED) {
|
||||
frame.payloadAs<ChannelInfo>()?.let { info ->
|
||||
if (info.kind == "thread" && info.auto && info.parentChatId != null) {
|
||||
// A thread/channel we were viewing got deleted — fall back
|
||||
// to a valid lane (the thread's parent, or the home channel).
|
||||
if (frame.type == TYPE_CHANNEL_DELETED) {
|
||||
frame.payloadAs<ChannelDeletedPayload>()?.let { p ->
|
||||
val (curChat, curThread) = chat.parseLane(chat.currentLane.value)
|
||||
if (curThread == null && curChat == info.parentChatId) {
|
||||
val lane = chat.laneKey(info.parentChatId, info.chatId)
|
||||
chat.setLane(lane)
|
||||
historyLoaded.add(lane)
|
||||
when {
|
||||
curThread == p.chatId -> {
|
||||
openChannel(curChat)
|
||||
}
|
||||
|
||||
curChat == p.chatId -> {
|
||||
channels.defaultChannel()?.let { openChannel(it.chatId) }
|
||||
}
|
||||
|
||||
else -> {
|
||||
Unit
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
// A thread/channel we were viewing got deleted — fall back
|
||||
// to a valid lane (the thread's parent, or the home channel).
|
||||
if (frame.type == TYPE_CHANNEL_DELETED) {
|
||||
frame.payloadAs<ChannelDeletedPayload>()?.let { p ->
|
||||
val (curChat, curThread) = chat.parseLane(chat.currentLane.value)
|
||||
when {
|
||||
curThread == p.chatId -> {
|
||||
openChannel(curChat)
|
||||
}
|
||||
|
||||
curChat == p.chatId -> {
|
||||
channels.defaultChannel()?.let { openChannel(it.chatId) }
|
||||
}
|
||||
TYPE_SEARCH_RESULTS -> {
|
||||
frame.payloadAs<SearchResultsPayload>()?.let { _searchResults.value = it.hits }
|
||||
_searching.value = false
|
||||
}
|
||||
|
||||
else -> {
|
||||
Unit
|
||||
}
|
||||
TYPE_COMMANDS_CATALOG -> {
|
||||
frame.payloadAs<CommandsCatalogPayload>()?.let { _slashCommands.value = it.commands }
|
||||
}
|
||||
|
||||
TYPE_SYNC_DONE -> {
|
||||
// Replayed frames already flowed through [events]; the
|
||||
// cursor is authoritative server-side (outbox).
|
||||
frame.payloadAs<SyncDonePayload>()?.let { store.syncCursor = it.cursor }
|
||||
}
|
||||
|
||||
TYPE_HISTORY -> {
|
||||
// Full message history for a chat/thread (initial open /
|
||||
// restart restore). `sync` only replays the outbox delta
|
||||
// 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.
|
||||
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)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
TYPE_SEARCH_RESULTS -> {
|
||||
frame.payloadAs<SearchResultsPayload>()?.let { _searchResults.value = it.hits }
|
||||
_searching.value = false
|
||||
}
|
||||
|
||||
TYPE_COMMANDS_CATALOG -> {
|
||||
frame.payloadAs<CommandsCatalogPayload>()?.let { _slashCommands.value = it.commands }
|
||||
}
|
||||
|
||||
TYPE_SYNC_DONE -> {
|
||||
// Replayed frames already flowed through [events]; the
|
||||
// cursor is authoritative server-side (outbox).
|
||||
frame.payloadAs<SyncDonePayload>()?.let { store.syncCursor = it.cursor }
|
||||
}
|
||||
|
||||
TYPE_HISTORY -> {
|
||||
// Full message history for a chat/thread (initial open /
|
||||
// restart restore). `sync` only replays the outbox delta
|
||||
// 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 lane = chat.laneKey(chatId, threadId)
|
||||
chat.loadHistory(
|
||||
lane,
|
||||
p.messages.map { it.toMessageItem() },
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
TYPE_NOTIFICATION -> {
|
||||
frame.payloadAs<NotificationPayload>()?.let { p ->
|
||||
pushBanner(p.kind, p.title, p.body, p.chatId, p.threadId)
|
||||
// M5: WS is live but the app is backgrounded — the
|
||||
// in-app banner is invisible, so mirror to a system
|
||||
// notification (the push backend only fires when
|
||||
// there is no live subscriber). Suppressed for
|
||||
// sync replays that already woke the device via
|
||||
// push (docs/08 §8.7).
|
||||
if (!isAppForeground() && !isPushedReplay(frame)) {
|
||||
postSystemNotification(p.chatId, null, p.title, p.body, p.threadId)
|
||||
TYPE_NOTIFICATION -> {
|
||||
frame.payloadAs<NotificationPayload>()?.let { p ->
|
||||
pushBanner(p.kind, p.title, p.body, p.chatId, p.threadId)
|
||||
// M5: WS is live but the app is backgrounded — the
|
||||
// in-app banner is invisible, so mirror to a system
|
||||
// notification (the push backend only fires when
|
||||
// there is no live subscriber). Suppressed for
|
||||
// sync replays that already woke the device via
|
||||
// push (docs/08 §8.7).
|
||||
if (!isAppForeground() && !isPushedReplay(frame)) {
|
||||
postSystemNotification(p.chatId, null, p.title, p.body, p.threadId)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
TYPE_TYPING -> {
|
||||
frame.payloadAs<TypingPayload>()?.let { _typing.value = it.on }
|
||||
}
|
||||
TYPE_TYPING -> {
|
||||
frame.payloadAs<TypingPayload>()?.let { _typing.value = it.on }
|
||||
}
|
||||
|
||||
TYPE_READ_RECEIPT -> {
|
||||
frame.payloadAs<ReadReceiptPayload>()?.let { chat.markRead(it.messageId) }
|
||||
}
|
||||
TYPE_READ_RECEIPT -> {
|
||||
frame.payloadAs<ReadReceiptPayload>()?.let { chat.markRead(it.messageId) }
|
||||
}
|
||||
|
||||
TYPE_MESSAGE_DELETED -> {
|
||||
// A message was deleted (by this or another device).
|
||||
// Idempotent: dropping an unknown id is a no-op.
|
||||
frame.payloadAs<MessageDeletedPayload>()?.let {
|
||||
chat.removeMessages(it.messageIds.toSet())
|
||||
TYPE_MESSAGE_DELETED -> {
|
||||
// A message was deleted (by this or another device).
|
||||
// Idempotent: dropping an unknown id is a no-op.
|
||||
frame.payloadAs<MessageDeletedPayload>()?.let {
|
||||
chat.removeMessages(it.messageIds.toSet())
|
||||
}
|
||||
}
|
||||
|
||||
TYPE_STATUS -> {
|
||||
frame.payloadAs<StatusPayload>()?.let { _gatewayStatus.value = it.state }
|
||||
}
|
||||
|
||||
TYPE_ERROR -> {
|
||||
// M7: error frames are global, not per-message — fail
|
||||
// any optimistic sends still in flight so they don't
|
||||
// sit at "sending…" forever.
|
||||
chat.failPending()
|
||||
}
|
||||
|
||||
else -> {
|
||||
Unit
|
||||
}
|
||||
}
|
||||
|
||||
TYPE_STATUS -> {
|
||||
frame.payloadAs<StatusPayload>()?.let { _gatewayStatus.value = it.state }
|
||||
}
|
||||
|
||||
TYPE_ERROR -> {
|
||||
// M7: error frames are global, not per-message — fail
|
||||
// any optimistic sends still in flight so they don't
|
||||
// sit at "sending…" forever.
|
||||
chat.failPending()
|
||||
}
|
||||
|
||||
else -> {
|
||||
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")
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
Reference in new issue
Block a user