fix(app): apply whole-review fixes (HIGH/MEDIUM/LOW) + dead code & stale comments
CI / Gateway plugin tests (push) Failing after 6m35s
CI / Kotlin tests (android host + desktop) (push) Successful in 6m57s

HIGH:
- ntfy listener: replace blocking exhausted() loop with SSE read + capped
  exponential-backoff reconnect; 60s read timeout
- GatewayClient.stop(): reset HTTP leg (http, httpCursor, sseFailures,
  usingLongPoll, lastAck)
- non-atomic shared state -> synchronized/@Volatile/AtomicLong/
  CopyOnWriteArrayList

MEDIUM:
- mediaId path-traversal guard (isValidMediaId) at network/app/fs boundaries
- loadFromCache: move ts=0 pending bubbles to end, keep stored order
- attachment placeholder tracked by identity, not filename
- optimistic ChannelStore updates on favorite/icon/automation/default
- secure-store caching (desktop map, android store)
- secret passed to keyring via stdin (macOS + Linux)
- SecureStore.clear() clears deviceId/syncCursor/fcmToken/ntfy*
- random ids for system messages; SSE EOF reconnect delay
- PowerShell $ escaping; dispose() cancels job before saving flows
- wire up "Forget pairing" in Settings
- move machine-specific org.gradle.java.home to user-level gradle.properties

LOW + dead code + stale comments:
- .aac->audio/aac; locale-fixed cost/size; 3-digit hex; hour+ latency
- Backdrop.DEFAULT defined once; notification id 24-bit; channel id cap
- remove dead FileSource, unused protocol/theme/media constants, empty
  onDispose, SDK_INT<O guard, hostFromUrl
- fix stale WS/SSE, M1/M5, and milestone KDoc comments

Verified: Kotlin desktop+android host tests, 100/100 Python gateway tests,
LSP clean, installed & running on device.
This commit is contained in:
ARIA committed 2026-08-23 12:57:35 +02:00
1 parent a4e4a4ea63
commit 863ab34915
32 files changed
+720 -390

No files matched your search

@@ -32,9 +32,9 @@ import iris.util.PairLink
/**
* Root composable shared by the Android and Desktop shells.
*
* M1: routes between the Connect screen (unpaired / auth failed) and the
* Chat screen (paired). Later milestones add the channel list, search,
* settings, and media (docs/10-android-app.md).
* Routes between the Connect screen (unpaired / auth failed) and the main
* app (paired), which hosts the channel list, chat, search, settings, and
* media (docs/10-android-app.md).
*/
@Composable
fun IrisApp(
@@ -71,6 +71,54 @@ class ChannelStore {
_channels.value.filter { it.chatId != p.chatId && it.parentChatId != p.chatId }
}
// ── Optimistic local updates (M-7) ────────────────────────────────────
//
// The gateway only broadcasts channel.created/renamed/deleted — there are
// no favorite/icon/automation/default events. So a toggle sent by THIS
// device would not update its own UI until a full channel.list re-fetch.
// These apply the change locally (optimistically); a later channel.list /
// hello.ack re-seed reconciles any divergence (e.g. a server rejection).
/** Toggle the cosmetic favorite flag locally. */
fun setFavorite(
chatId: String,
on: Boolean,
) = update(chatId) { it.copy(favorite = on) }
/** Toggle the automation flag locally. */
fun setAutomation(
chatId: String,
on: Boolean,
) = update(chatId) { it.copy(automation = on) }
/** Set the icon (base64) and/or avatar color locally; null clears a field. */
fun setIcon(
chatId: String,
icon: String?,
color: String?,
) = update(chatId) { it.copy(icon = icon, color = color) }
/** Make [chatId] the default channel locally (clearing the previous one). */
fun setDefault(chatId: String) {
_channels.value =
sorted(
_channels.value.map {
when {
it.chatId == chatId -> it.copy(isDefault = true)
it.isDefault -> it.copy(isDefault = false)
else -> it
}
},
)
}
private fun update(
chatId: String,
transform: (ChannelInfo) -> ChannelInfo,
) {
_channels.value = sorted(_channels.value.map { if (it.chatId == chatId) transform(it) else it })
}
private fun sorted(list: List<ChannelInfo>): List<ChannelInfo> =
list.sortedWith(
compareByDescending<ChannelInfo> { it.isDefault }
@@ -156,7 +156,12 @@ class ChatStore {
private val _unread = MutableStateFlow<Map<String, Int>>(emptyMap())
val unread: StateFlow<Map<String, Int>> = _unread.asStateFlow()
private var localSeq = 0
/** Single lock for the lane/todo/unread maps: they are mutated from the
* UI thread (addPending via send) and the frame-collector thread
* (Dispatchers.Default). A non-atomic read-modify-write loses a frame
* that lands between the read and the write (e.g. a streaming delta
* dropped while the user sends). */
private val lock = Any()
/** When false, `message.start`/`message.update` frames are ignored and each
* reply materializes as a single final message on `message.stop`
@@ -199,9 +204,26 @@ class ChatStore {
lane: String,
transform: (List<ChatItem>) -> List<ChatItem>,
) {
synchronized(lock) {
val map = _lanes.value.toMutableMap()
map[lane] = transform(map[lane].orEmpty())
_lanes.value = map
}
}
/** Apply [transform] to every lane, writing back only when something
* changed. Callers must hold [lock]. */
private fun mapLanes(transform: (List<ChatItem>) -> List<ChatItem>) {
val map = _lanes.value.toMutableMap()
map[lane] = transform(map[lane].orEmpty())
_lanes.value = map
var changed = false
for ((lane, list) in map) {
val updated = transform(list)
if (updated != list) {
map[lane] = updated
changed = true
}
}
if (changed) _lanes.value = map
}
// ── Optimistic send ───────────────────────────────────────────────────
@@ -230,8 +252,10 @@ class ChatStore {
lane: String,
text: String,
) {
localSeq++
val id = "sys_$localSeq"
// randomId (not a process-local seq): system messages are persisted,
// and a seq that resets on restart would re-mint sys_0 and collide
// with the restored one (upsert overwrite).
val id = randomId("sys")
updateLane(lane) {
it + MessageItem(id = id, role = "system", text = text, ts = nowMillis(), isSystem = true)
}
@@ -354,16 +378,18 @@ class ChatStore {
val (chatId, threadId) = parseLane(lane)
if (threadId == null) return
val flatLane = chatId
val map = _lanes.value.toMutableMap()
val flatList = map[flatLane].orEmpty()
val idx =
flatList.indexOfLast {
it is MessageItem && it.role == ROLE_USER && it.text == p.text &&
(it.pending || it.status == MsgStatus.Failed)
}
if (idx < 0) return
map[flatLane] = flatList.toMutableList().also { it.removeAt(idx) }
_lanes.value = map
synchronized(lock) {
val map = _lanes.value.toMutableMap()
val flatList = map[flatLane].orEmpty()
val idx =
flatList.indexOfLast {
it is MessageItem && it.role == ROLE_USER && it.text == p.text &&
(it.pending || it.status == MsgStatus.Failed)
}
if (idx < 0) return@synchronized
map[flatLane] = flatList.toMutableList().also { it.removeAt(idx) }
_lanes.value = map
}
}
/** Merge server media refs into existing items, keeping local paths. */
@@ -540,9 +566,11 @@ class ChatStore {
frame: Frame,
) {
val p = frame.payloadAs<TodoUpdatePayload>() ?: return
val map = _todos.value.toMutableMap()
if (p.todos.isEmpty()) map.remove(lane) else map[lane] = p.todos
_todos.value = map
synchronized(lock) {
val map = _todos.value.toMutableMap()
if (p.todos.isEmpty()) map.remove(lane) else map[lane] = p.todos
_todos.value = map
}
}
// ── commentary (dimmed interim beat) ──────────────────────────────────
@@ -623,10 +651,8 @@ class ChatStore {
mediaId: String,
localPath: String,
) {
val map = _lanes.value.toMutableMap()
var changed = false
for ((lane, list) in map) {
val updated =
synchronized(lock) {
mapLanes { list ->
list.map { item ->
if (item is MessageItem) {
item.copy(
@@ -639,12 +665,8 @@ class ChatStore {
item
}
}
if (updated != list) {
map[lane] = updated
changed = true
}
}
if (changed) _lanes.value = map
}
// ── picker.choice (interactive slash-command menu) ──────────────────────
@@ -682,10 +704,8 @@ class ChatStore {
pickerId: String,
value: String,
) {
val map = _lanes.value.toMutableMap()
var changed = false
for ((lane, list) in map) {
val updated =
synchronized(lock) {
mapLanes { list ->
list.map { item ->
if (item is PickerItem && item.id == pickerId && item.selected == null) {
item.copy(selected = value)
@@ -693,27 +713,27 @@ class ChatStore {
item
}
}
if (updated != list) {
map[lane] = updated
changed = true
}
}
if (changed) _lanes.value = map
}
/** M8: a new message arrived in [lane] that the user hasn't seen —
* increment its unread count. */
fun markUnread(lane: String) {
val map = _unread.value.toMutableMap()
map[lane] = (map[lane] ?: 0) + 1
_unread.value = map
synchronized(lock) {
val map = _unread.value.toMutableMap()
map[lane] = (map[lane] ?: 0) + 1
_unread.value = map
}
}
/** M8: the user is now viewing [lane]'s newest content — clear its unread
* count. Idempotent (a lane with no unread is a no-op). */
fun markLaneRead(lane: String) {
val map = _unread.value.toMutableMap()
if (map.remove(lane) != null) _unread.value = map
synchronized(lock) {
val map = _unread.value.toMutableMap()
if (map.remove(lane) != null) _unread.value = map
}
}
/** M8: unread count for a single lane (0 when none). */
@@ -721,10 +741,8 @@ class ChatStore {
/** M5: mark the user message [messageId] as read (read.receipt). */
fun markRead(messageId: String) {
val map = _lanes.value.toMutableMap()
var changed = false
for ((lane, list) in map) {
val updated =
synchronized(lock) {
mapLanes { list ->
list.map { item ->
if (item is MessageItem && item.id == messageId && item.role == ROLE_USER &&
item.status != MsgStatus.Read
@@ -734,22 +752,16 @@ class ChatStore {
item
}
}
if (updated != list) {
map[lane] = updated
changed = true
}
}
if (changed) _lanes.value = map
}
/** M7: mark a single user message as failed (the send never reached the
* gateway — network drop, or the gateway rejected it); tap the bubble
* to retry. */
fun failMessage(messageId: String) {
val map = _lanes.value.toMutableMap()
var changed = false
for ((lane, list) in map) {
val updated =
synchronized(lock) {
mapLanes { list ->
list.map { item ->
if (item is MessageItem && item.id == messageId && item.role == ROLE_USER &&
item.status != MsgStatus.Failed
@@ -759,12 +771,8 @@ class ChatStore {
item
}
}
if (updated != list) {
map[lane] = updated
changed = true
}
}
if (changed) _lanes.value = map
}
/**
@@ -775,16 +783,9 @@ class ChatStore {
*/
fun removeMessages(messageIds: Set<String>) {
if (messageIds.isEmpty()) return
val map = _lanes.value.toMutableMap()
var changed = false
for ((lane, list) in map) {
val updated = list.filterNot { it.id in messageIds }
if (updated != list) {
map[lane] = updated
changed = true
}
synchronized(lock) {
mapLanes { list -> list.filterNot { it.id in messageIds } }
}
if (changed) _lanes.value = map
}
/**
@@ -795,10 +796,8 @@ class ChatStore {
* message.stop) will never arrive to close them.
*/
fun finalizeInterrupted() {
val map = _lanes.value.toMutableMap()
var changed = false
for ((lane, list) in map) {
val updated =
synchronized(lock) {
mapLanes { list ->
list.map { item ->
when (item) {
is ToolItem -> if (!item.done) item.copy(done = true, ok = false) else item
@@ -806,20 +805,14 @@ class ChatStore {
is PickerItem -> item
}
}
if (updated != list) {
map[lane] = updated
changed = true
}
}
if (changed) _lanes.value = map
}
/** M7: mark all pending user messages as failed (gateway error frame). */
fun failPending() {
val map = _lanes.value.toMutableMap()
var changed = false
for ((lane, list) in map) {
val updated =
synchronized(lock) {
mapLanes { list ->
list.map { item ->
if (item is MessageItem && item.role == ROLE_USER && item.status == MsgStatus.Pending) {
item.copy(pending = false, status = MsgStatus.Failed)
@@ -827,12 +820,8 @@ class ChatStore {
item
}
}
if (updated != list) {
map[lane] = updated
changed = true
}
}
if (changed) _lanes.value = map
}
/** M7: re-arm a failed user message for a retry send. */
@@ -924,11 +913,27 @@ class ChatStore {
*/
fun loadFromCache(lanes: Map<String, List<ChatItem>>) {
if (lanes.isEmpty()) return
_lanes.value = lanes
synchronized(lock) {
// The cache is ordered by (ts, id); pending/failed sends are
// persisted with ts = 0, so the DB returns them FIRST — but in
// memory addPending appends them to the END of the lane. Move the
// ts=0 message bubbles to the end (preserving their relative
// order); everything else keeps its stored order, so a restored
// lane looks like the live one until loadHistory re-sorts it
// after a connect.
_lanes.value =
lanes.mapValues { (_, items) ->
val (zeroTs, rest) = items.partition { (it as? MessageItem)?.ts == 0L }
rest + zeroTs
}
}
}
fun clear() {
_lanes.value = emptyMap()
_unread.value = emptyMap()
synchronized(lock) {
_lanes.value = emptyMap()
_unread.value = emptyMap()
_todos.value = emptyMap()
}
}
}
@@ -2,8 +2,8 @@ package iris.data
/**
* Pairing settings storage. The token is a secret: platform actuals keep it
* in secure storage (EncryptedSharedPreferences on Android — M5; plain
* SharedPreferences for M1 dev, file on desktop).
* in secure storage (EncryptedSharedPreferences on Android, OS keyring or an
* encrypted file on desktop).
*/
interface SecureStore {
/** http(s)://host:port (legacy ws(s):// URLs are still accepted) */
@@ -1,10 +0,0 @@
package iris.media
/** A readable local file (expect/actual; JVM impl in jvmMain). */
expect class FileSource(path: String) : AutoCloseable {
/** Total size in bytes. */
fun size(): Long
/** Read up to [buf.size] bytes into [buf]; returns bytes read or -1 at EOF. */
fun read(buf: ByteArray): Int
}
@@ -6,32 +6,43 @@ import iris.protocol.KIND_IMAGE
import iris.protocol.KIND_VIDEO
/** Map a MIME type to a media kind (docs/07 §7.1). */
fun kindFromMime(mime: String): String = when {
mime.startsWith("image/") -> KIND_IMAGE
mime.startsWith("video/") -> KIND_VIDEO
mime.startsWith("audio/") -> KIND_AUDIO
else -> KIND_DOCUMENT
}
fun kindFromMime(mime: String): String =
when {
mime.startsWith("image/") -> KIND_IMAGE
mime.startsWith("video/") -> KIND_VIDEO
mime.startsWith("audio/") -> KIND_AUDIO
else -> KIND_DOCUMENT
}
/** Best-effort file extension for a MIME type (cache file naming). */
fun extForMime(mime: String): String = when {
mime == "image/jpeg" -> ".jpg"
mime == "image/png" -> ".png"
mime == "image/webp" -> ".webp"
mime == "image/gif" -> ".gif"
mime == "image/heic" -> ".heic"
mime == "image/heif" -> ".heif"
mime == "video/mp4" -> ".mp4"
mime == "video/webm" -> ".webm"
mime == "video/quicktime" -> ".mov"
mime == "audio/mpeg" -> ".mp3"
mime == "audio/mp4" || mime == "audio/x-m4a" -> ".m4a"
mime == "audio/ogg" -> ".ogg"
mime == "audio/wav" -> ".wav"
mime == "audio/flac" -> ".flac"
mime == "audio/aac" -> ".aac"
mime == "application/pdf" -> ".pdf"
mime == "application/zip" -> ".zip"
mime == "text/plain" -> ".txt"
else -> ".bin"
}
fun extForMime(mime: String): String =
when {
mime == "image/jpeg" -> ".jpg"
mime == "image/png" -> ".png"
mime == "image/webp" -> ".webp"
mime == "image/gif" -> ".gif"
mime == "image/heic" -> ".heic"
mime == "image/heif" -> ".heif"
mime == "video/mp4" -> ".mp4"
mime == "video/webm" -> ".webm"
mime == "video/quicktime" -> ".mov"
mime == "audio/mpeg" -> ".mp3"
mime == "audio/mp4" || mime == "audio/x-m4a" -> ".m4a"
mime == "audio/ogg" -> ".ogg"
mime == "audio/wav" -> ".wav"
mime == "audio/flac" -> ".flac"
mime == "audio/aac" -> ".aac"
mime == "application/pdf" -> ".pdf"
mime == "application/zip" -> ".zip"
mime == "text/plain" -> ".txt"
else -> ".bin"
}
/**
* Validate a server-provided media id before it is used in a file path or a
* `GET /v1/media/{id}` URL (M-2 / S-1). The id comes from the gateway's
* `media.offer` frame; a value like `../../x` would write outside the media
* directory and break the pull URL. Only a conservative token charset is
* accepted.
*/
fun isValidMediaId(id: String): Boolean = id.length in 1..128 && id.all { it.isLetterOrDigit() || it == '_' || it == '-' }
@@ -13,7 +13,6 @@ import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Job
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.delay
@@ -27,7 +26,6 @@ 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 okhttp3.OkHttpClient
import java.util.concurrent.TimeUnit
@@ -42,7 +40,9 @@ import kotlin.random.Random
* - connect: health probe + SSE hello (the HTTP hello.ack)
* - reconnect: exponential backoff + jitter; re-hello on every (re)connect
* - events: server frames on [events]
* - request/response correlation by id
* - correlation: the POST body carries the synchronous reply (e.g. read
* receipt, errors); everything else arrives on the event stream, and the
* app reconciles by frame id (docs/19 §19.7)
*/
class GatewayClient(
private val scope: CoroutineScope,
@@ -87,6 +87,14 @@ class GatewayClient(
private var connectJob: Job? = null
private var nextRequestId = 1
// Incremented from the SSE callback thread (onHttpHello) and the UI
// thread (sendFrame/sendMessage) — keep it atomic or ids collide and
// request/response correlation breaks.
private fun nextId(): Int = synchronized(this) { nextRequestId++ }
// Written by poke() (UI thread) and the connect loop (Default).
@Volatile
private var attempt = 0
// Set by poke() (app returned to the foreground): the connect loop's
@@ -98,17 +106,21 @@ class GatewayClient(
// start(). Drives Connecting (first dial) vs Reconnecting (redial after a
// drop) so the UI can show the right status without a blocking screen.
private var hasConnected = false
private val pending = mutableMapOf<Int, CompletableDeferred<Frame>>()
// HTTP leg: [http] is created lazily from the stored URL; [httpCursor] is
// the resume cursor (SSE id / outbox high-water mark).
// the resume cursor (SSE id / outbox high-water mark), updated from the
// SSE/poll callback threads.
private var http: HttpGateway? = null
@Volatile
private var httpCursor: Long = 0
private var sseFailures = 0
private var usingLongPoll = false
// Last hello.ack payload — used to restore State.Connected after a
// reconnect state race in the connect loop.
// reconnect state race in the connect loop. Written by the SSE callback
// thread, read by the long-poll loop.
@Volatile
private var lastAck: HelloAckPayload? = null
// Epoch ms of the last frame delivered by the receive stream (SSE or
@@ -144,11 +156,22 @@ class GatewayClient(
connectJob = scope.launch { connectLoop() }
}
/** Stop the connect loop. */
/**
* Stop the connect loop and reset the HTTP leg. Without the reset, a
* re-pair to a *different* gateway would keep using the cached
* [HttpGateway] (old URL, old token, old resume cursor) until the process
* is killed.
*/
fun stop() {
connectJob?.cancel()
connectJob = null
_state.value = State.Disconnected
http?.close()
http = null
httpCursor = 0
sseFailures = 0
usingLongPoll = false
lastAck = null
}
/**
@@ -249,22 +272,25 @@ class GatewayClient(
// ── HTTP receive leg ──────────────────────────────────────────────────
/** Lazily build the HTTP client from the stored URL. */
private fun httpGateway(): HttpGateway? {
val url = store.serverUrl.trim()
val token = store.token
if (url.isBlank() || token.isBlank()) return null
return http
?: HttpGateway(
client,
HttpGateway.deriveHttpUrl(url),
token,
store.deviceId,
deviceName = store.deviceName,
fcmToken = { store.fcmToken.ifBlank { null } },
ntfyTopic = { store.ntfyTopic.ifBlank { null } },
).also { http = it }
}
/** Lazily build the HTTP client from the stored URL. Synchronized: two
* threads racing the check-then-create would leak a gateway (and its
* OkHttp clients). */
private fun httpGateway(): HttpGateway? =
synchronized(this) {
val url = store.serverUrl.trim()
val token = store.token
if (url.isBlank() || token.isBlank()) return@synchronized null
http
?: HttpGateway(
client,
HttpGateway.deriveHttpUrl(url),
token,
store.deviceId,
deviceName = store.deviceName,
fcmToken = { store.fcmToken.ifBlank { null } },
ntfyTopic = { store.ntfyTopic.ifBlank { null } },
).also { http = it }
}
/**
* The HTTP receive loop: SSE by default; after two consecutive SSE open
@@ -298,9 +324,16 @@ class GatewayClient(
onHello = { onHttpHello(it) },
onFrame = { emitHttpFrame(it) },
onCursor = { if (it > httpCursor) httpCursor = it },
// A keep-alive comment proves the stream is alive —
// count it toward liveness so an idle-but-healthy
// stream isn't force-reconnected by the stale
// watchdog every 3 minutes.
onKeepAlive = { lastFrameMs = nowMs() },
)
// Clean EOF: reconnect immediately.
// Clean EOF: back off briefly so a server that keeps
// closing cleanly can't tight-loop the reconnect.
backoff = 1_000L
delay(backoff)
} catch (e: HttpGateway.HttpAuthException) {
_state.value = State.AuthFailed("gateway rejected the pairing token (HTTP 401)")
return
@@ -328,7 +361,7 @@ class GatewayClient(
// M5: reconnect catch-up — replay frames parked while offline.
val local = store.syncCursor
if (local < ack.syncCursor) {
val id = nextRequestId++
val id = nextId()
scope.launch { httpGateway()?.postFrame(syncFrame(id, local)) }
}
// Prompt fast path (before the possibly-starved state collector).
@@ -361,9 +394,20 @@ class GatewayClient(
/** Deliver an HTTP-leg frame to the same sinks as any other frame. */
private fun emitHttpFrame(frame: Frame) {
lastFrameMs = nowMs()
_events.tryEmit(frame)
frame.id?.let { id ->
pending[id]?.complete(frame)
if (!_events.tryEmit(frame)) {
// The collector is starved beyond the 128-frame buffer: the frame
// is dropped. Log it — a silent drop here loses user-visible
// content (echoes, deltas, notifications).
IrisLog.e("events buffer full — dropped frame ${frame.type} (id=${frame.id})")
}
}
/** Deliver the POST body's synchronous reply (or null) and handle a 401.
* Shared by [sendMessage] and [sendFrame]. */
private fun handlePostResult(res: HttpGateway.PostResult?) {
res?.frame?.let { emitHttpFrame(it) }
if (res?.status == 401) {
_state.value = State.AuthFailed("gateway rejected the pairing token (HTTP 401)")
}
}
@@ -390,7 +434,7 @@ class GatewayClient(
onResult?.invoke(0)
return
}
val id = nextRequestId++
val id = nextId()
scope.launch {
val res =
httpGateway()?.postFrame(
@@ -399,10 +443,7 @@ class GatewayClient(
// The synchronous reply (e.g. the read receipt, or an error frame
// on 4xx) comes back in the POST body, not on the event stream —
// deliver it or it is lost (docs/19 §19.7).
res?.frame?.let { emitHttpFrame(it) }
if (res?.status == 401) {
_state.value = State.AuthFailed("gateway rejected the pairing token (HTTP 401)")
}
handlePostResult(res)
onResult?.invoke(res?.status ?: 0)
}
}
@@ -438,10 +479,8 @@ class GatewayClient(
}
companion object {
const val UPLOAD_TIMEOUT_MS = 120_000L
const val PULL_TIMEOUT_MS = 300_000L
/** No frames from the receive stream for this long (gateway still
/** No frames (or keep-alive comments) from the receive stream for
* this long (gateway still
* healthy) = stale stream: force a fresh (re)connect. 3 min keeps an
* idle app from churning while still recovering a stuck stream well
* before a user notices (issue #6). */
@@ -455,16 +494,13 @@ class GatewayClient(
*/
fun sendFrame(frame: Frame): Int {
if (_state.value !is State.Connected) return -1
val id = nextRequestId++
val id = nextId()
scope.launch {
val res = httpGateway()?.postFrame(frame.copy(id = id))
// Single-frame responses (commands.catalog, channel.list, search,
// history, sync, errors) come back in the POST body, not on the
// event stream — deliver it or it is lost (docs/19 §19.7).
res?.frame?.let { emitHttpFrame(it) }
if (res?.status == 401) {
_state.value = State.AuthFailed("gateway rejected the pairing token (HTTP 401)")
}
handlePostResult(res)
}
return id
}
@@ -1,6 +1,7 @@
package iris.net
import iris.media.Sha256
import iris.media.isValidMediaId
import iris.protocol.ErrorPayload
import iris.protocol.Frame
import iris.protocol.HelloAckPayload
@@ -60,6 +61,19 @@ class HttpGateway(
private val pollClient: OkHttpClient = client.pollClient()
private val mediaClient: OkHttpClient = client.mediaClient()
/** Close the per-purpose clients (and their idle connections). The
* shared base [client] is owned by the caller. */
fun close() {
healthClient.dispatcher.executorService.shutdown()
streamClient.dispatcher.executorService.shutdown()
pollClient.dispatcher.executorService.shutdown()
mediaClient.dispatcher.executorService.shutdown()
healthClient.connectionPool.evictAll()
streamClient.connectionPool.evictAll()
pollClient.connectionPool.evictAll()
mediaClient.connectionPool.evictAll()
}
/** POST /v1/frame result. [frame] is the handler's synchronous reply
* (error frame on 4xx, e.g. read.receipt on 200) or null for a plain
* 202 accept-and-ack. */
@@ -195,8 +209,10 @@ class HttpGateway(
* leg is proven; the server may still replay a large outbox before the
* hello); [onHello] fires for `event: hello` (the HTTP hello.ack);
* [onFrame] for `event: frame`; [onCursor] with the SSE `id` (outbox
* cursor) when present. Returns on clean EOF; throws [IOException] on
* open/read failure. Callbacks run on the IO thread.
* cursor) when present; [onKeepAlive] for SSE comment lines (the
* heartbeat) so callers can count keep-alives toward liveness. Returns
* on clean EOF; throws [IOException] on open/read failure. Callbacks run
* on the IO thread.
*/
suspend fun events(
cursor: Long,
@@ -204,6 +220,7 @@ class HttpGateway(
onHello: (HelloAckPayload) -> Unit,
onFrame: (Frame) -> Unit,
onCursor: (Long) -> Unit,
onKeepAlive: (() -> Unit)? = null,
) {
withContext(Dispatchers.IO) {
val request =
@@ -257,10 +274,10 @@ class HttpGateway(
}
line.startsWith(":") -> {
Unit
// heartbeat comment
onKeepAlive?.invoke()
}
// heartbeat comment
line.startsWith("id:") -> {
line
.removePrefix("id:")
@@ -399,6 +416,11 @@ class HttpGateway(
onChunk: suspend (ByteArray) -> Unit,
): Result<Unit> =
withContext(Dispatchers.IO) {
// Server-controlled id: reject path-traversal / URL-breaking
// values before they reach the request line (M-2 / S-1).
if (!isValidMediaId(mediaId)) {
return@withContext Result.failure(IllegalStateException("invalid media id"))
}
val request =
Request
.Builder()
@@ -6,7 +6,7 @@ package iris.platform
* The controller (commonMain) needs to know whether the app is in the
* foreground (to decide between an in-app banner and a system notification)
* and to post a system notification when a `notification` frame arrives while
* the app is backgrounded but the WS is still live.
* the app is backgrounded but the gateway connection is still live.
*/
/** True when the app's UI is visible (Android: activity resumed). */
@@ -12,9 +12,9 @@ import kotlinx.serialization.json.put
/**
* Wire protocol frames (mirror of gateway-plugin/protocol.py).
* See docs/04-wire-protocol.md. M1: hello/hello.ack, message, message.send,
* error, ping/pong, typing. M2: message.start/update/stop, tool.start/
* progress/end, commentary, reasoning (on message / message.stop).
* See docs/04-wire-protocol.md. Defines all frame types and payloads:
* hello/hello.ack, message (send + streaming), tool cards, commentary,
* reasoning, channels, media, notifications, sync, and error frames.
*/
const val PROTOCOL_VERSION = 1
@@ -89,20 +89,11 @@ const val TYPE_SYNC = "sync"
const val TYPE_SYNC_DONE = "sync.done"
const val TYPE_HISTORY = "history"
// ── Error codes ─────────────────────────────────────────────────────────
const val ERR_AUTH = "auth"
const val ERR_NOT_FOUND = "not_found"
const val ERR_UNSUPPORTED = "unsupported"
const val ERR_INTERNAL = "internal"
const val ERR_MEDIA_TOO_LARGE = "media_too_large"
// M4 — media kinds (docs/07 §7.1)
const val KIND_IMAGE = "image"
const val KIND_AUDIO = "audio"
const val KIND_VIDEO = "video"
const val KIND_DOCUMENT = "document"
const val KIND_VOICE = "voice"
// ── Roles ───────────────────────────────────────────────────────────────
@@ -516,12 +507,9 @@ data class MessageDeletedPayload(
// ── M5: push / notifications ────────────────────────────────────────────
/** Notification kinds (mirror of protocol.NOTIF_*). */
const val NOTIF_MESSAGE = "message"
const val NOTIF_APPROVAL = "approval"
const val NOTIF_CLARIFY = "clarify"
const val NOTIF_CRON = "cron"
const val NOTIF_CHANNEL = "channel"
const val NOTIF_OUTBOX_PRUNED = "outbox_pruned"
/** Kinds that stay on screen until dismissed (docs/08 §8.3). */
val HIGH_PRIORITY_NOTIF_KINDS = setOf(NOTIF_APPROVAL, NOTIF_CLARIFY, NOTIF_CRON)
@@ -8,6 +8,7 @@ import iris.data.MessageItem
import iris.data.MsgStatus
import iris.data.SecureStore
import iris.media.MediaCache
import iris.media.isValidMediaId
import iris.media.kindFromMime
import iris.net.GatewayClient
import iris.platform.PickedFile
@@ -45,7 +46,6 @@ import iris.protocol.TYPE_ERROR
import iris.protocol.TYPE_HISTORY
import iris.protocol.TYPE_MEDIA_OFFER
import iris.protocol.TYPE_MESSAGE
import iris.protocol.TYPE_MESSAGE_DELETE
import iris.protocol.TYPE_MESSAGE_DELETED
import iris.protocol.TYPE_MESSAGE_START
import iris.protocol.TYPE_MESSAGE_STOP
@@ -90,6 +90,8 @@ import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.debounce
import kotlinx.coroutines.launch
import java.util.Collections
import java.util.concurrent.atomic.AtomicLong
import kotlin.random.Random
/**
@@ -144,8 +146,9 @@ class IrisController(
/** Lanes whose full history has been loaded this session (in-memory; reset
* on a process death, which is exactly when a reload is needed). The
* `sync` delta can seed a lane with recent frames without it being opened,
* so "lane is empty" is not a reliable first-open signal. */
private val historyLoaded = mutableSetOf<String>()
* so "lane is empty" is not a reliable first-open signal. Synchronized:
* written by the frame collector, read by the UI thread (loadHistory). */
private val historyLoaded = Collections.synchronizedSet(mutableSetOf<String>())
/** Gateway health state (M5: status frame; null = never received).
* Note: "restarting" is deliberately NOT stored here — it posts the
@@ -300,6 +303,9 @@ class IrisController(
}
companion object {
/** Max simultaneous banners; persistent ones are exempt from the cap. */
private const val MAX_BANNERS = 5
const val FONT_SCALE_MIN = 0.8f
const val FONT_SCALE_MAX = 1.5f
@@ -411,8 +417,12 @@ class IrisController(
// ── M4: media ─────────────────────────────────────────────────────────
private val mediaCache = MediaCache(mediaCacheBaseDir())
/** A composer attachment: picked file being uploaded (or uploaded). */
/** A composer attachment: picked file being uploaded (or uploaded).
* [id] is a per-attachment identity used to apply the upload result to
* the right placeholder — matching by filename would cross-update two
* files picked with the same name (M-6). */
data class PendingAttachment(
val id: String,
val filename: String,
val mime: String,
val size: Long,
@@ -441,7 +451,10 @@ class IrisController(
private val _banners = MutableStateFlow<List<Banner>>(emptyList())
val banners: StateFlow<List<Banner>> = _banners.asStateFlow()
private var bannerSeq = 0L
// Atomic: pushBanner can be called from the frame collector and the UI
// thread; a lost update would mint a duplicate banner id.
private val bannerSeq = AtomicLong(0)
fun dismissBanner(id: Long) {
_banners.value = _banners.value.filterNot { it.id == id }
@@ -455,8 +468,13 @@ class IrisController(
threadId: String?,
) {
val persistent = kind in HIGH_PRIORITY_NOTIF_KINDS
val banner = Banner(bannerSeq++, kind, title, body, chatId, threadId, persistent)
_banners.value = (_banners.value + banner).takeLast(5)
val banner = Banner(bannerSeq.getAndIncrement(), kind, title, body, chatId, threadId, persistent)
// Persistent banners (approval / clarify / cron) are never evicted by
// the cap; only the transient ones compete for the remaining slots.
val merged = _banners.value + banner
val keptPersistent = merged.filter { it.persistent }
val room = (MAX_BANNERS - keptPersistent.size).coerceAtLeast(0)
_banners.value = keptPersistent + merged.filterNot { it.persistent }.takeLast(room)
if (!persistent) {
scope.launch {
delay(5_000)
@@ -723,7 +741,7 @@ class IrisController(
)
// Mark the lane loaded only when the response
// is actually processed: if the request or
// response is lost in a WS drop, the lane
// response is lost in a connection drop, the lane
// stays unmarked and the next (re)connect
// retries it.
historyLoaded.add(lane)
@@ -743,7 +761,7 @@ class IrisController(
// sync replay) must not re-show the banner.
if (!alreadyConsumed) {
pushBanner(p.kind, p.title, p.body, p.chatId, p.threadId)
// M5: WS is live but the app is backgrounded — the
// M5: the connection 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
@@ -777,7 +795,7 @@ class IrisController(
if (st.state == "restarting") {
// Gateway is going down (restart/stop):
// post the notice IMMEDIATELY — the
// socket can take up to the ping timeout
// connection can take up to the ping timeout
// (~20 s) to actually drop, and waiting
// for that transition would delay the
// message. No banner for this state: the
@@ -883,7 +901,7 @@ class IrisController(
}
/**
* Fast-path connect handler (runs on the WS thread via [GatewayClient.onHelloAck],
* Fast-path connect handler (runs on the gateway callback 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.
@@ -970,6 +988,7 @@ class IrisController(
}
fun setDefaultChannel(chatId: String) {
channels.setDefault(chatId)
client.sendFrame(channelSetDefaultFrame(0, chatId))
}
@@ -978,6 +997,7 @@ class IrisController(
chatId: String,
on: Boolean,
) {
channels.setFavorite(chatId, on)
client.sendFrame(channelFavoriteFrame(0, chatId, on))
}
@@ -987,6 +1007,7 @@ class IrisController(
chatId: String,
on: Boolean,
) {
channels.setAutomation(chatId, on)
client.sendFrame(channelSetAutomationFrame(0, chatId, on))
}
@@ -997,6 +1018,7 @@ class IrisController(
icon: String?,
color: String?,
) {
channels.setIcon(chatId, icon, color)
client.sendFrame(channelIconFrame(0, chatId, icon, color))
}
@@ -1055,9 +1077,10 @@ class IrisController(
val lane = chat.laneKey(chatId, threadId)
if (lane in historyLoaded) return
// Newest page, sized to restore a full working view on restart / first
// open (older pages are reachable via scroll-up pagination). The lane
// latest page only (older pages are not paginated in the current UI).
// 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.
// a request lost in a connection drop must be retryable on reconnect.
client.sendFrame(historyFrame(0, chatId, threadId, limit = 200))
}
@@ -1150,8 +1173,9 @@ class IrisController(
* gateway error frame): queued (Pending) or failed bubbles that go out
* automatically on the next (re)connect. In-memory only — a process
* death leaves them as tap-to-retry (the local cache restore already
* marks pending sends failed). */
private val networkFailed = mutableSetOf<String>()
* marks pending sends failed). Synchronized: written by the send-result
* callback (UI thread) and the frame collector. */
private val networkFailed = Collections.synchronizedSet(mutableSetOf<String>())
/** Resend queued/failed user messages of [lane] now that the link is
* back: a message the server already has (the POST response was lost in
@@ -1205,8 +1229,13 @@ class IrisController(
/** Stage a picked file: upload it, then keep it as a pending attachment. */
fun attachFile(picked: PickedFile) {
val kind = kindFromMime(picked.mime)
// Per-attachment identity: the upload result is applied to THIS
// placeholder by id, so two files with the same name don't
// cross-update (M-6).
val id = "att_${Random.nextLong(1_000_000_000L, 9_999_999_999L)}"
val placeholder =
PendingAttachment(
id = id,
filename = picked.name,
mime = picked.mime,
size = picked.size,
@@ -1226,7 +1255,7 @@ class IrisController(
)
_attachments.value =
_attachments.value.map {
if (it.filename == picked.name && it.uploading) {
if (it.id == id && it.uploading) {
result.fold(
{ ref -> it.copy(uploading = false, mediaRef = ref) },
{ e -> it.copy(uploading = false, error = e.message) },
@@ -1245,6 +1274,12 @@ class IrisController(
/** Pull offered media into the local cache and record the path. */
private fun pullMedia(offer: MediaOfferPayload) {
scope.launch {
// Server-controlled id: reject path-traversal values before they
// reach the cache or the pull URL (M-2 / S-1).
if (!isValidMediaId(offer.mediaId)) {
IrisLog.w("rejecting media offer with invalid id: ${offer.mediaId.take(40)}")
return@launch
}
// Already cached? Skip the pull.
mediaCache.path(offer.mediaId, offer.mime)?.let {
chat.setMediaLocalPath(offer.mediaId, it)
@@ -1283,11 +1318,13 @@ 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).
// M-17: cancel the debounce job FIRST so a pending save can't fire
// after our final flush and clobber it; then do the final synchronous
// flush so the newest frames survive the process death (the debounce
// window may still hold unsaved changes).
job.cancel()
chatDb.saveLanes(chat.lanes.value)
chatDb.saveChannels(channels.channels.value)
job.cancel()
}
}
@@ -118,7 +118,6 @@ import iris.net.GatewayClient
import iris.platform.ImageFilePicker
import iris.platform.MediaFilePicker
import iris.platform.MediaPlayerView
import iris.platform.PickedFile
import iris.platform.decodeBase64Image
import iris.platform.encodeImageAsBase64
import iris.platform.isDesktop
@@ -152,6 +151,7 @@ import iris.util.prepareForMarkdown
import iris.util.preserveNewlinesAsHardBreaks
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import java.util.Locale
import kotlin.math.roundToInt
/**
@@ -1514,7 +1514,9 @@ private fun LetterAvatar(
/** Parse a "#RRGGBB" (or "#AARRGGBB") hex string into a [Color]. */
private fun parseHexColor(hex: String): Color {
val cleaned = hex.removePrefix("#")
var cleaned = hex.removePrefix("#")
// Expand 3-digit shorthand (#F00 -> #FF0000) so short hex is accepted.
if (cleaned.length == 3) cleaned = cleaned.map { "$it$it" }.joinToString("")
val withAlpha = if (cleaned.length == 6) "FF$cleaned" else cleaned
return Color(withAlpha.toLongOrNull(16) ?: 0xFF000000L)
}
@@ -2658,13 +2660,19 @@ private fun formatLatency(seconds: Double): String {
if (seconds < 1) return "<1s"
val total = seconds.roundToInt()
if (total < 60) return "${total}s"
val m = total / 60
val h = total / 3600
val m = (total % 3600) / 60
val s = total % 60
return "${m}m${s.toString().padStart(2, '0')}s"
return if (h > 0) {
"${h}h${m.toString().padStart(2, '0')}m"
} else {
"${m}m${s.toString().padStart(2, '0')}s"
}
}
/** Format a cost in USD: sub-cent at 4 decimals, else 2. */
private fun formatCost(cost: Double): String = if (cost < 0.01) "$%.4f".format(cost) else "$%.2f".format(cost)
/** Format a cost in USD: sub-cent at 4 decimals, else 2. Fixed locale so the
* decimal separator is always a dot (L-12). */
private fun formatCost(cost: Double): String = if (cost < 0.01) "$%.4f".format(Locale.US, cost) else "$%.2f".format(Locale.US, cost)
/** Circular selection indicator shown beside a bubble in selection mode. */
@Composable
@@ -2839,8 +2847,8 @@ private fun MediaDocChip(media: MediaItem) {
private fun fmtSize(bytes: Long): String =
when {
bytes >= 1_048_576 -> "%.1f MB".format(bytes / 1_048_576.0)
bytes >= 1024 -> "%.0f KB".format(bytes / 1024.0)
bytes >= 1_048_576 -> "%.1f MB".format(Locale.US, bytes / 1_048_576.0)
bytes >= 1024 -> "%.0f KB".format(Locale.US, bytes / 1024.0)
else -> "$bytes B"
}
@@ -78,6 +78,7 @@ fun SettingsScreen(
val fontSizeScale by controller.fontSizeScale.collectAsState()
val theme = LocalUserTheme.current
var pickerTarget by remember { mutableStateOf<PickerTarget?>(null) }
var showForgetConfirm by remember { mutableStateOf(false) }
Box(modifier = Modifier.fillMaxSize().background(theme.background)) {
Column(
@@ -391,6 +392,24 @@ fun SettingsScreen(
}
}
}
Text(
"Connection",
style = MaterialTheme.typography.titleSmall,
modifier = Modifier.padding(top = 12.dp, bottom = 4.dp),
)
SettingsCard {
Text("🔌 Forget pairing", fontSize = 14.sp)
Text(
"Clears the gateway token and wipes local chat history",
fontSize = 12.sp,
color = IrisColors.textDim,
)
Spacer(modifier = Modifier.height(8.dp))
TextButton(onClick = { showForgetConfirm = true }) {
Text("Forget pairing", fontSize = 12.sp)
}
}
}
when (pickerTarget) {
@@ -437,6 +456,32 @@ fun SettingsScreen(
Unit
}
}
if (showForgetConfirm) {
AlertDialog(
onDismissRequest = { showForgetConfirm = false },
title = { Text("Forget pairing?") },
text = {
Text(
"This clears the gateway token and wipes all local chat history on this device. You'll need to pair again to reconnect.",
fontSize = 13.sp,
)
},
confirmButton = {
TextButton(onClick = {
showForgetConfirm = false
controller.forget()
}) {
Text("Forget")
}
},
dismissButton = {
TextButton(onClick = { showForgetConfirm = false }) {
Text("Cancel")
}
},
)
}
}
}
@@ -1,8 +1,9 @@
package iris.ui.theme
/**
* A bundled backdrop image shipped in the app (Android `assets/backdrops/`,
* desktop classpath `backdrops/`), loaded via [iris.platform.readBackdropBytes].
* A bundled backdrop image shipped in the app (Android: `androidApp` module
* `assets/backdrops/`; desktop: classpath `backdrops/`), loaded via
* [iris.platform.readBackdropBytes].
*
* Bundled backdrops are referenced in [UserTheme.backgroundImagePath] by the
* sentinel path `backdrop://<id>` (see [path]); user-picked images use a real
@@ -18,9 +19,12 @@ data class Backdrop(
companion object {
const val PATH_PREFIX = "backdrop://"
/** Default background for new chats / fresh installs. */
val DEFAULT = Backdrop("pexels-yunszyveli-12368637", "Yun Syzveli")
val ALL: List<Backdrop> =
listOf(
Backdrop("pexels-yunszyveli-12368637", "Yun Syzveli"),
DEFAULT,
Backdrop("pexels-bogdankrupin-12049700", "Bogdan Krupin"),
Backdrop("pexels-bosichong-27940302", "Bosi Chong"),
Backdrop("pexels-farhan-najeer-644774196-32490483", "Farhan Najeer"),
@@ -28,9 +32,6 @@ data class Backdrop(
Backdrop("pexels-steve-29390703", "Steve"),
)
/** Default background for new chats / fresh installs. */
val DEFAULT: Backdrop = ALL.first { it.id == "pexels-yunszyveli-12368637" }
/** True when [path] references a bundled backdrop. */
fun isBackdropPath(path: String?): Boolean = !path.isNullOrBlank() && path.startsWith(PATH_PREFIX)
@@ -35,8 +35,6 @@ object IrisColors {
val textDim = Color(0xFF8A93A6)
// Bubbles
val bubbleUser = primary
val bubbleAssistant = Color(0xFF2A2E3B)
val bubbleCommentary = Color(0xFF23262F)
// Panels, chips, rows
@@ -47,7 +45,6 @@ object IrisColors {
val divider = Color(0xFF2A2E3B)
// Status
val statusGrey = Color(0xFF9E9E9E)
val statusAmber = Color(0xFFFFC107)
val statusGreen = Color(0xFF4CAF50)
val statusRed = Color(0xFFF44336)
@@ -11,9 +11,3 @@ expect fun localDayKey(epochMillis: Long): String
/** Current wall-clock time in epoch milliseconds (for locally generated items). */
expect fun nowMillis(): Long
/** Host part of a pairing URL (the "host:port" of the full gateway ws URL). */
fun hostFromUrl(url: String): String {
val noScheme = url.trim().substringAfter("://")
return noScheme.substringBefore("/").ifBlank { url.trim() }
}
@@ -5,7 +5,7 @@
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)
id TEXT NOT NULL, -- message id (server id, or local<rand>/sys<rand> for optimistic)
ts INTEGER NOT NULL, -- epoch millis (0 for optimistic, not yet echoed)
payload TEXT NOT NULL, -- serialized MessageItem
PRIMARY KEY (lane, id)
@@ -13,7 +13,7 @@ CREATE TABLE message (
CREATE TABLE tool (
lane TEXT NOT NULL, -- lane key: chatId or chatId::threadId
id TEXT NOT NULL, -- local tool card id (tool_N)
id TEXT NOT NULL, -- local tool card id (tool<rand>)
seq INTEGER NOT NULL, -- card order within the lane (lane position)
payload TEXT NOT NULL, -- serialized ToolItem (carries its anchor_id)
PRIMARY KEY (lane, id)