From bf6bf7e8bd82774090a95f9c0824e69234dfea7a Mon Sep 17 00:00:00 2001 From: ARIA Date: Thu, 20 Aug 2026 12:00:13 +0200 Subject: [PATCH] M7: polish + E2E + docs (layout pass, theming, states, e2e driver, schema, setup.md, security) --- .pre-commit-config.yaml | 15 + README.md | 17 +- app/shared/build.gradle.kts | 3 + .../iris/platform/AndroidSecureStore.kt | 56 +- .../src/commonMain/kotlin/iris/IrisApp.kt | 22 +- .../commonMain/kotlin/iris/data/ChatStore.kt | 61 +- .../kotlin/iris/net/GatewayClient.kt | 5 +- .../kotlin/iris/protocol/Protocol.kt | 14 +- .../kotlin/iris/state/IrisController.kt | 33 + .../kotlin/iris/ui/screens/ChatScreen.kt | 704 ++++++++++++++---- .../kotlin/iris/ui/screens/ConnectScreen.kt | 122 +-- .../commonMain/kotlin/iris/ui/theme/Theme.kt | 72 ++ .../commonMain/kotlin/iris/util/TimeFormat.kt | 16 + .../jvmMain/kotlin/iris/util/TimeFormatJvm.kt | 27 + docs/04-wire-protocol.md | 17 +- docs/09-pairing-security.md | 24 +- docs/14-milestones.md | 38 +- docs/README.md | 2 +- docs/protocol/frames.schema.json | 58 +- docs/setup.md | 170 +++++ gateway-plugin/adapter.py | 15 + gateway-plugin/protocol.py | 33 +- gateway-plugin/tests/README.md | 71 +- gateway-plugin/tests/e2e.py | 373 ++++++++++ gateway-plugin/tests/ws_probe.py | 495 ++++++++++-- gateway-plugin/ws_server.py | 89 ++- 26 files changed, 2225 insertions(+), 327 deletions(-) create mode 100644 .pre-commit-config.yaml create mode 100644 app/shared/src/commonMain/kotlin/iris/ui/theme/Theme.kt create mode 100644 app/shared/src/commonMain/kotlin/iris/util/TimeFormat.kt create mode 100644 app/shared/src/jvmMain/kotlin/iris/util/TimeFormatJvm.kt create mode 100644 docs/setup.md create mode 100644 gateway-plugin/tests/e2e.py diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml new file mode 100644 index 0000000..a10ae32 --- /dev/null +++ b/.pre-commit-config.yaml @@ -0,0 +1,15 @@ +# Committed pre-commit config so a fresh clone gets the hermes-agent/ guard +# without manual hook installation (scripts/guard_hermes_agent.sh --staged +# fails the commit if any hermes-agent/ path is staged โ€” that directory is a +# read-only research reference, see .gitignore and docs/00-overview.md). +# +# Install: pre-commit install +repos: + - repo: local + hooks: + - id: guard-hermes-agent + name: guard hermes-agent/ (read-only research reference) + entry: scripts/guard_hermes_agent.sh --staged + language: system + pass_filenames: false + always_run: true \ No newline at end of file diff --git a/README.md b/README.md index ee1c553..2b8c883 100644 --- a/README.md +++ b/README.md @@ -15,6 +15,21 @@ A **native Android + Desktop** experience for [hermes-agent](https://github.com/ ๐Ÿ“š **The full implementation reference library is in [`docs/`](docs/README.md).** Read `docs/00-overview.md` first, then follow the numbered docs. +## Quickstart + +Full walkthrough: [`docs/setup.md`](docs/setup.md). + +1. **Gateway:** install the plugin (`ln -s "$PWD/gateway-plugin" ~/.hermes/plugins/android`), run `hermes gateway setup` (generates the pairing token, prints the server URL), then `hermes gateway`. +2. **App:** `cd app && ./gradlew :androidApp:installDebug` (Android; ADB device connected) or `./gradlew :desktopApp:run` (desktop). +3. **Pair:** on the app's Connect screen, enter the server URL (`ws://:8790/ws`) + pairing token, then **Test & Connect**. +4. **Chat.** + +### Pairing notes + +- **Token location:** `~/.hermes/.env` on the gateway host (`ANDROID_TOKEN`), or the `hermes gateway setup` output (printed once, at generation). +- **Manual entry only:** the app has no QR scanner yet โ€” the server prints a QR payload, but you type the URL + token. +- **Default bind is `127.0.0.1`:** for a phone on the LAN, set `ANDROID_WS_HOST` to the gateway's LAN IP. + ## The three deliverables (this monorepo) | Path | What | @@ -28,7 +43,7 @@ project** (`app/`) with a shared KMP module (`app/shared`). ## Status -- **Phase:** Planning complete โ†’ ready to implement (Milestone M0). +- **Phase:** M0โ€“M6 complete; M7 (polish + E2E + docs) in progress. - **Milestones:** see [`docs/14-milestones.md`](docs/14-milestones.md). - **Locked decisions:** see [`docs/16-open-questions.md`](docs/16-open-questions.md). diff --git a/app/shared/build.gradle.kts b/app/shared/build.gradle.kts index 0b25a9e..6f7ebaa 100644 --- a/app/shared/build.gradle.kts +++ b/app/shared/build.gradle.kts @@ -35,6 +35,9 @@ kotlin { } // M4: ExoPlayer (Media3) for inline audio/video playback (Android only). androidMain.dependencies { + // M7: EncryptedSharedPreferences for the pairing token + // (docs/09 ยง9.7; AndroidSecureStore). + implementation("androidx.security:security-crypto:1.1.0") implementation("androidx.media3:media3-exoplayer:1.3.1") implementation("androidx.media3:media3-ui:1.3.1") // SAF picker (rememberLauncherForActivityResult). diff --git a/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt b/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt index b67ad99..233f1e9 100644 --- a/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt +++ b/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt @@ -1,16 +1,64 @@ package iris.platform import android.content.Context +import android.content.SharedPreferences import android.os.Build +import androidx.security.crypto.EncryptedSharedPreferences +import androidx.security.crypto.MasterKey import iris.data.SecureStore import java.util.UUID /** - * Android pairing storage. M1: SharedPreferences (dev). M5 moves the token - * to EncryptedSharedPreferences per docs/10 ยง10.2. + * Android pairing storage. M7: EncryptedSharedPreferences (MasterKey + * AES256_GCM) per docs/09 ยง9.7. The M1/M5 plain "iris" SharedPreferences + * values are migrated on first run (read old key, write encrypted, delete + * old key) so an upgrade never loses the pairing. */ class AndroidSecureStore(context: Context) : SecureStore { - private val prefs = context.applicationContext.getSharedPreferences("iris", Context.MODE_PRIVATE) + private val appContext = context.applicationContext + private val plainPrefs = appContext.getSharedPreferences(PLAIN_PREFS_NAME, Context.MODE_PRIVATE) + private val prefs: SharedPreferences = EncryptedSharedPreferences.create( + appContext, + SECURE_PREFS_NAME, + MasterKey.Builder(appContext) + .setKeyScheme(MasterKey.KeyScheme.AES256_GCM) + .build(), + EncryptedSharedPreferences.PrefKeyEncryptionScheme.AES256_SIV, + EncryptedSharedPreferences.PrefValueEncryptionScheme.AES256_GCM, + ) + + init { + migratePlainValues() + } + + /** One-time migration of the M1/M5 plain values into the encrypted store. */ + private fun migratePlainValues() { + val editor = prefs.edit() + var migrated = false + for (key in listOf(KEY_URL, KEY_TOKEN, KEY_DEVICE_ID, KEY_FCM_TOKEN, KEY_NTFY_TOPIC, KEY_NTFY_SERVER)) { + val old = plainPrefs.getString(key, null) + if (old != null && !prefs.contains(key)) { + editor.putString(key, old) + migrated = true + } + } + val oldCursor = plainPrefs.getLong(KEY_SYNC_CURSOR, 0L) + if (oldCursor != 0L && !prefs.contains(KEY_SYNC_CURSOR)) { + editor.putLong(KEY_SYNC_CURSOR, oldCursor) + migrated = true + } + if (migrated) editor.apply() + // The plain store must not keep a copy of any value. + plainPrefs.edit() + .remove(KEY_URL) + .remove(KEY_TOKEN) + .remove(KEY_DEVICE_ID) + .remove(KEY_SYNC_CURSOR) + .remove(KEY_FCM_TOKEN) + .remove(KEY_NTFY_TOPIC) + .remove(KEY_NTFY_SERVER) + .apply() + } override var serverUrl: String get() = prefs.getString(KEY_URL, "").orEmpty() @@ -59,6 +107,8 @@ class AndroidSecureStore(context: Context) : SecureStore { } private companion object { + const val PLAIN_PREFS_NAME = "iris" + const val SECURE_PREFS_NAME = "iris_secure" const val KEY_URL = "server_url" const val KEY_TOKEN = "token" const val KEY_DEVICE_ID = "device_id" diff --git a/app/shared/src/commonMain/kotlin/iris/IrisApp.kt b/app/shared/src/commonMain/kotlin/iris/IrisApp.kt index a602fd3..0093743 100644 --- a/app/shared/src/commonMain/kotlin/iris/IrisApp.kt +++ b/app/shared/src/commonMain/kotlin/iris/IrisApp.kt @@ -1,9 +1,7 @@ package iris import androidx.compose.foundation.background -import androidx.compose.material3.MaterialTheme import androidx.compose.material3.Surface -import androidx.compose.material3.darkColorScheme import androidx.compose.runtime.Composable import androidx.compose.runtime.DisposableEffect import androidx.compose.runtime.LaunchedEffect @@ -11,21 +9,16 @@ import androidx.compose.runtime.collectAsState import androidx.compose.runtime.getValue import androidx.compose.runtime.remember import androidx.compose.ui.Modifier -import androidx.compose.ui.graphics.Color import iris.data.SecureStore import iris.net.GatewayClient import iris.platform.setActiveController import iris.state.IrisController import iris.ui.screens.ChatScreen import iris.ui.screens.ConnectScreen - -private val IrisDark = darkColorScheme( - background = Color(0xFF1B1E28), - surface = Color(0xFF222634), - onBackground = Color(0xFFE8EAF0), - onSurface = Color(0xFFE8EAF0), - primary = Color(0xFF4F7CFF), -) +import iris.ui.screens.ConnectingScreen +import iris.ui.theme.IrisColors +import iris.ui.theme.IrisTheme +import iris.util.hostFromUrl /** * Root composable shared by the Android and Desktop shells. @@ -57,8 +50,8 @@ fun IrisApp( } val state by controller.client.state.collectAsState() - MaterialTheme(colorScheme = IrisDark) { - Surface(modifier = Modifier.background(IrisDark.background)) { + IrisTheme { + Surface(modifier = Modifier.background(IrisColors.background)) { val s = state when (s) { GatewayClient.State.Disconnected -> @@ -70,6 +63,9 @@ fun IrisApp( prefillToken = store.token, initialError = "Pairing rejected: ${s.message}", ) + // M7: initial connect in flight โ€” dedicated screen, not an empty chat. + GatewayClient.State.Connecting -> + ConnectingScreen(hostFromUrl(store.serverUrl)) else -> ChatScreen(controller) } } diff --git a/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt b/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt index 86a64fd..55f4afd 100644 --- a/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt +++ b/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt @@ -42,6 +42,14 @@ sealed interface ChatItem { val id: String } +/** Delivery status of a user message (M5: read.receipt; M7: failed sends). */ +enum class MsgStatus { + Pending, // optimistic, not yet acknowledged by the gateway + Sent, // gateway accepted it (echo received) + Read, // agent received and started processing it (read.receipt) + Failed, // send failed (error frame); tap the bubble to retry +} + /** A chat message (user / assistant / commentary / streaming bubble). */ data class MessageItem( override val id: String, @@ -49,6 +57,7 @@ data class MessageItem( val text: String, val ts: Long, val pending: Boolean = false, + val status: MsgStatus = MsgStatus.Sent, val reasoning: String? = null, val isCommentary: Boolean = false, val streaming: Boolean = false, @@ -139,7 +148,7 @@ class ChatStore { localSeq++ val id = "local_$localSeq" updateLane(lane) { - it + MessageItem(id = id, role = ROLE_USER, text = text, ts = 0, pending = true, media = media) + it + MessageItem(id = id, role = ROLE_USER, text = text, ts = 0, pending = true, status = MsgStatus.Pending, media = media) } return id } @@ -176,6 +185,7 @@ class ChatStore { reasoning = p.reasoning, pending = false, streaming = false, + status = if (cur.status == MsgStatus.Read) MsgStatus.Read else MsgStatus.Sent, model = p.model, tokens = p.tokens, ts = p.ts ?: cur.ts, @@ -382,6 +392,55 @@ class ChatStore { if (changed) _lanes.value = map } + /** 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 = list.map { item -> + if (item is MessageItem && item.id == messageId && item.role == ROLE_USER && + item.status != MsgStatus.Read + ) { + item.copy(pending = false, status = MsgStatus.Read) + } else 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 = list.map { item -> + if (item is MessageItem && item.role == ROLE_USER && item.status == MsgStatus.Pending) { + item.copy(pending = false, status = MsgStatus.Failed) + } else 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. */ + fun rearmForRetry(lane: String, messageId: String) { + updateLane(lane) { list -> + list.map { item -> + if (item is MessageItem && item.id == messageId && item.status == MsgStatus.Failed) { + item.copy(pending = true, status = MsgStatus.Pending) + } else item + } + } + } + fun clear() { _lanes.value = emptyMap() } diff --git a/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt b/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt index e6db2be..af96cd1 100644 --- a/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt +++ b/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt @@ -217,7 +217,10 @@ class GatewayClient( } TYPE_ERROR -> { val err = frame.payloadAs() - authError.complete(err?.message ?: "auth failed") + 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 -> { diff --git a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt index f8420bf..25d9657 100644 --- a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt +++ b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt @@ -54,9 +54,11 @@ const val TYPE_MEDIA_OFFER = "media.offer" const val TYPE_MEDIA_PULL = "media.pull" const val TYPE_MEDIA_PULL_END = "media.pull.end" -// M5 โ€” push / notifications +// M5 โ€” push / notifications / read receipt / gateway status const val TYPE_NOTIFICATION = "notification" const val TYPE_FCM_REGISTER = "fcm.register" +const val TYPE_READ_RECEIPT = "read.receipt" +const val TYPE_STATUS = "status" // M3 โ€” channels / threads / search / sync const val TYPE_CHANNEL_CREATE = "channel.create" @@ -383,6 +385,16 @@ data class FcmRegisterPayload( @SerialName("ntfy_topic") val ntfyTopic: String? = null, ) +// โ”€โ”€ M5: read receipt / gateway status (server -> app) โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€ + +@Serializable +data class ReadReceiptPayload( + @SerialName("message_id") val messageId: String, +) + +@Serializable +data class StatusPayload(val state: String) + // โ”€โ”€ Frame builders โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€ fun helloFrame( diff --git a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt index 042dc0b..eb11656 100644 --- a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt +++ b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt @@ -3,6 +3,8 @@ package iris.state import iris.data.ChatStore import iris.data.ChannelStore import iris.data.MediaItem +import iris.data.MessageItem +import iris.data.MsgStatus import iris.data.SecureStore import iris.media.MediaCache import iris.media.kindFromMime @@ -16,10 +18,13 @@ import iris.protocol.MediaOfferPayload import iris.protocol.MessagePayload import iris.protocol.MessageStopPayload import iris.protocol.NotificationPayload +import iris.protocol.ReadReceiptPayload import iris.protocol.ROLE_ASSISTANT import iris.protocol.SearchHit import iris.protocol.SearchResultsPayload +import iris.protocol.StatusPayload import iris.protocol.SyncDonePayload +import iris.protocol.TYPE_ERROR import iris.protocol.TYPE_NOTIFICATION import iris.protocol.TYPE_CHANNEL_CREATED import iris.protocol.TYPE_CHANNEL_DELETED @@ -31,7 +36,9 @@ import iris.protocol.TYPE_MESSAGE import iris.protocol.TYPE_MESSAGE_START import iris.protocol.TYPE_MESSAGE_STOP import iris.protocol.TYPE_MESSAGE_UPDATE +import iris.protocol.TYPE_READ_RECEIPT import iris.protocol.TYPE_SEARCH_RESULTS +import iris.protocol.TYPE_STATUS import iris.protocol.TYPE_SYNC_DONE import iris.protocol.TYPE_TOOL_END import iris.protocol.TYPE_TOOL_PROGRESS @@ -89,6 +96,10 @@ class IrisController( private val _homeChannel = MutableStateFlow("android:default") val homeChannel: StateFlow = _homeChannel.asStateFlow() + /** Gateway health state (M5: status frame; null = never received). */ + private val _gatewayStatus = MutableStateFlow(null) + val gatewayStatus: StateFlow = _gatewayStatus.asStateFlow() + // โ”€โ”€ M3: threads toggle (per-app for now; per-channel lands later) โ”€โ”€โ”€โ”€โ”€ private val _threadsEnabled = MutableStateFlow(false) val threadsEnabled: StateFlow = _threadsEnabled.asStateFlow() @@ -259,6 +270,18 @@ class IrisController( TYPE_TYPING -> { frame.payloadAs()?.let { _typing.value = it.on } } + TYPE_READ_RECEIPT -> { + frame.payloadAs()?.let { chat.markRead(it.messageId) } + } + TYPE_STATUS -> { + frame.payloadAs()?.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 } } @@ -375,6 +398,16 @@ class IrisController( _attachments.value = emptyList() } + /** M7: resend a failed user message (tap on the failed bubble). */ + fun retrySend(messageId: String) { + val lane = chat.currentLane.value + val (chatId, threadId) = chat.parseLane(lane) + val item = chat.lanes.value[lane]?.firstOrNull { it.id == messageId } as? MessageItem ?: return + if (item.status != MsgStatus.Failed) return + chat.rearmForRetry(lane, messageId) + client.sendMessage(chatId, item.text, threadId, item.media.map { it.mediaId }) + } + // โ”€โ”€ M4: attachments (pick -> upload -> send) โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€ /** Stage a picked file: upload it, then keep it as a pending attachment. */ diff --git a/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt b/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt index 8590268..af79d1e 100644 --- a/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt +++ b/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt @@ -6,6 +6,7 @@ import androidx.compose.foundation.focusable import androidx.compose.foundation.horizontalScroll import androidx.compose.foundation.layout.Arrangement import androidx.compose.foundation.layout.Box +import androidx.compose.foundation.layout.BoxWithConstraints import androidx.compose.foundation.layout.Column import androidx.compose.foundation.layout.PaddingValues import androidx.compose.foundation.layout.Row @@ -17,13 +18,16 @@ import androidx.compose.foundation.layout.heightIn import androidx.compose.foundation.layout.padding import androidx.compose.foundation.layout.size import androidx.compose.foundation.layout.width +import androidx.compose.foundation.layout.widthIn import androidx.compose.foundation.Image import androidx.compose.foundation.layout.fillMaxHeight import androidx.compose.foundation.lazy.LazyColumn import androidx.compose.foundation.lazy.items import androidx.compose.foundation.lazy.rememberLazyListState import androidx.compose.foundation.rememberScrollState +import androidx.compose.foundation.shape.CircleShape import androidx.compose.foundation.shape.RoundedCornerShape +import androidx.compose.foundation.text.BasicTextField import androidx.compose.foundation.text.KeyboardActions import androidx.compose.foundation.text.KeyboardOptions import androidx.compose.foundation.verticalScroll @@ -31,6 +35,8 @@ import androidx.compose.material3.AlertDialog import androidx.compose.material3.Button import androidx.compose.material3.CircularProgressIndicator import androidx.compose.material3.DrawerValue +import androidx.compose.material3.DropdownMenu +import androidx.compose.material3.DropdownMenuItem import androidx.compose.material3.IconButton import androidx.compose.material3.MaterialTheme import androidx.compose.material3.ModalDrawerSheet @@ -69,12 +75,15 @@ import androidx.compose.ui.text.AnnotatedString import androidx.compose.ui.text.font.FontFamily import androidx.compose.ui.text.font.FontWeight import androidx.compose.ui.text.input.ImeAction +import androidx.compose.ui.text.style.TextOverflow +import androidx.compose.ui.unit.Dp import androidx.compose.ui.unit.dp import androidx.compose.ui.unit.sp import iris.data.ChannelStore import iris.data.ChatItem import iris.data.MediaItem import iris.data.MessageItem +import iris.data.MsgStatus import iris.data.ToolItem import iris.net.GatewayClient import iris.platform.MediaFilePicker @@ -90,6 +99,11 @@ import iris.protocol.ROLE_USER import iris.protocol.SearchHit import iris.state.IrisController import iris.state.ToolDetail +import iris.ui.theme.IrisColors +import iris.ui.theme.avatarColor +import iris.util.formatDayLabel +import iris.util.formatTime +import iris.util.localDayKey import kotlinx.coroutines.launch /** @@ -188,47 +202,114 @@ fun ChatScreen(controller: IrisController) { // M6: the chat column (header + list + composer), shared by both layouts. @Composable fun ChatContent(showHamburger: Boolean, onOpenDrawer: () -> Unit) { + // M7: header overflow actions (rename / forget pairing). + var showOverflow by remember { mutableStateOf(false) } + var showRename by remember { mutableStateOf(false) } + var showForget by remember { mutableStateOf(false) } Column(modifier = Modifier.fillMaxSize()) { - // Header - Row( - modifier = Modifier - .fillMaxWidth() - .padding(horizontal = 8.dp, vertical = 6.dp), - verticalAlignment = Alignment.CenterVertically, - ) { - if (showHamburger) { - IconButton(onClick = onOpenDrawer) { - Text("โ˜ฐ", fontSize = 18.sp) + // Header (M7: avatar + two-line title pill, reference look). Wide + // panes keep everything on one row; narrow phones move the utility + // controls to a second row so the pill stays dominant. + val channelName = currentChannel?.name ?: "Iris" + BoxWithConstraints { + val wide = maxWidth.value >= 600f + if (wide) { + Row( + modifier = Modifier + .fillMaxWidth() + .padding(horizontal = 8.dp, vertical = 6.dp), + verticalAlignment = Alignment.CenterVertically, + ) { + if (showHamburger) { + HeaderIconButton(onClick = onOpenDrawer) { + Text("โ˜ฐ", fontSize = 18.sp) + } + } + TitlePill(channelName, Modifier.weight(1f)) + StatusChip(state) + Spacer(modifier = Modifier.width(4.dp)) + HeaderIconButton(onClick = { showSearch = true }) { + Text("๐Ÿ”", fontSize = 16.sp) + } + HeaderIconButton(onClick = { controller.toggleThreads() }) { + Text(if (threadsEnabled) "๐Ÿงตโœ“" else "๐Ÿงต", fontSize = 16.sp) + } + TextButton(onClick = { controller.cycleToolDetail() }) { + Text("tools: ${toolDetail.label}", fontSize = 11.sp) + } + if (isDesktop) { + IconButton(onClick = { showInspector = !showInspector }) { + Box( + modifier = Modifier + .clip(RoundedCornerShape(8.dp)) + .background( + if (showInspector) IrisColors.chipSelected else Color.Transparent, + ) + .padding(2.dp), + ) { + Text("โ“˜", fontSize = 16.sp) + } + } + } + HeaderOverflow( + expanded = showOverflow, + onToggle = { showOverflow = !showOverflow }, + onRename = { showRename = true }, + onForget = { showForget = true }, + ) } - } - Text( - currentChannel?.name ?: "Iris", - style = MaterialTheme.typography.titleMedium, - fontWeight = FontWeight.SemiBold, - modifier = Modifier.weight(1f), - ) - StatusChip(state) - Spacer(modifier = Modifier.width(4.dp)) - IconButton(onClick = { showSearch = true }) { - Text("๐Ÿ”", fontSize = 16.sp) - } - IconButton(onClick = { controller.toggleThreads() }) { - Text(if (threadsEnabled) "๐Ÿงตโœ“" else "๐Ÿงต", fontSize = 16.sp) - } - TextButton(onClick = { controller.cycleToolDetail() }) { - Text("tools: ${toolDetail.label}", fontSize = 11.sp) - } - if (isDesktop) { - IconButton(onClick = { showInspector = !showInspector }) { - Box( + } else { + Column { + Row( modifier = Modifier - .clip(RoundedCornerShape(8.dp)) - .background( - if (showInspector) Color(0xFF2A3550) else Color.Transparent, - ) - .padding(2.dp), + .fillMaxWidth() + .padding(horizontal = 8.dp, vertical = 6.dp), + verticalAlignment = Alignment.CenterVertically, ) { - Text("โ“˜", fontSize = 16.sp) + if (showHamburger) { + HeaderIconButton(onClick = onOpenDrawer) { + Text("โ˜ฐ", fontSize = 18.sp) + } + } + TitlePill(channelName, Modifier.weight(1f)) + HeaderOverflow( + expanded = showOverflow, + onToggle = { showOverflow = !showOverflow }, + onRename = { showRename = true }, + onForget = { showForget = true }, + ) + } + Row( + modifier = Modifier + .fillMaxWidth() + .padding(start = 8.dp, end = 8.dp, bottom = 4.dp), + verticalAlignment = Alignment.CenterVertically, + horizontalArrangement = Arrangement.spacedBy(4.dp), + ) { + StatusChip(state) + HeaderIconButton(onClick = { showSearch = true }) { + Text("๐Ÿ”", fontSize = 16.sp) + } + HeaderIconButton(onClick = { controller.toggleThreads() }) { + Text(if (threadsEnabled) "๐Ÿงตโœ“" else "๐Ÿงต", fontSize = 16.sp) + } + TextButton(onClick = { controller.cycleToolDetail() }) { + Text("tools: ${toolDetail.label}", fontSize = 11.sp) + } + if (isDesktop) { + IconButton(onClick = { showInspector = !showInspector }) { + Box( + modifier = Modifier + .clip(RoundedCornerShape(8.dp)) + .background( + if (showInspector) IrisColors.chipSelected else Color.Transparent, + ) + .padding(2.dp), + ) { + Text("โ“˜", fontSize = 16.sp) + } + } + } } } } @@ -245,39 +326,48 @@ fun ChatScreen(controller: IrisController) { ) } - // Messages - LazyColumn( - state = listState, - modifier = Modifier - .weight(1f) - .fillMaxWidth(), - contentPadding = PaddingValues(16.dp), - verticalArrangement = Arrangement.spacedBy(8.dp), - ) { - if (items.isEmpty() && !typing) { - item(key = "empty") { - Box(modifier = Modifier.fillMaxWidth(), contentAlignment = Alignment.Center) { - Text( - "Say hello to your agent.", - color = MaterialTheme.colorScheme.onSurfaceVariant, - ) + // Messages (M7: interleaved date separators; bubbles capped at 85% width) + val rows = remember(items) { buildChatRows(items) } + BoxWithConstraints(modifier = Modifier.weight(1f).fillMaxWidth()) { + val bubbleMaxWidth = (maxWidth.value * 0.85f).dp + LazyColumn( + state = listState, + modifier = Modifier.fillMaxSize(), + contentPadding = PaddingValues(16.dp), + verticalArrangement = Arrangement.spacedBy(8.dp), + ) { + if (items.isEmpty() && !typing) { + item(key = "empty") { + Box(modifier = Modifier.fillMaxWidth(), contentAlignment = Alignment.Center) { + Text( + "Say hello to your agent.", + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + } } } - } - items(items, key = { it.id }) { item -> - when (item) { - is MessageItem -> MessageBubble(item) - is ToolItem -> if (toolDetail != ToolDetail.NOTHING) ToolCard(item, toolDetail) + items(rows, key = { it.key }) { row -> + when (row) { + is ChatRow.Day -> DaySeparator(row.label) + is ChatRow.Item -> when (val item = row.item) { + is MessageItem -> MessageBubble( + msg = item, + maxWidth = bubbleMaxWidth, + onRetry = { controller.retrySend(item.id) }, + ) + is ToolItem -> if (toolDetail != ToolDetail.NOTHING) ToolCard(item, toolDetail) + } + } } - } - if (typing) { - item(key = "typing") { - Text( - "typingโ€ฆ", - style = MaterialTheme.typography.bodySmall, - color = MaterialTheme.colorScheme.onSurfaceVariant, - modifier = Modifier.padding(start = 12.dp), - ) + if (typing) { + item(key = "typing") { + Text( + "typingโ€ฆ", + style = MaterialTheme.typography.bodySmall, + color = MaterialTheme.colorScheme.onSurfaceVariant, + modifier = Modifier.padding(start = 12.dp), + ) + } } } } @@ -318,35 +408,116 @@ fun ChatScreen(controller: IrisController) { } } - // Composer + // M7: connection-state banner (persistent, above the composer; no dismiss). + val gatewayStatus by controller.gatewayStatus.collectAsState() + val connBanner = when { + state is GatewayClient.State.Reconnecting -> + "Reconnecting to gatewayโ€ฆ messages will sync automatically when the connection is back." + gatewayStatus == "restarting" -> "Gateway is restartingโ€ฆ" + gatewayStatus == "degraded" -> + "Gateway reports a degraded state โ€” replies may be slow or unavailable." + else -> null + } + if (connBanner != null) { + Box( + modifier = Modifier + .fillMaxWidth() + .padding(horizontal = 12.dp, vertical = 4.dp), + contentAlignment = Alignment.Center, + ) { + Box( + modifier = Modifier + .clip(RoundedCornerShape(16.dp)) + .background(IrisColors.panel) + .padding(horizontal = 14.dp, vertical = 8.dp), + ) { + Text(connBanner, fontSize = 12.sp, color = IrisColors.textTertiary) + } + } + } + + // Composer (M7: rounded pill + accent circular send button) + val canSend = input.isNotBlank() || attachments.any { it.mediaRef != null && it.error == null } Row( modifier = Modifier .fillMaxWidth() - .padding(12.dp), - verticalAlignment = Alignment.Bottom, + .padding(horizontal = 12.dp, vertical = 8.dp) + .clip(RoundedCornerShape(28.dp)) + .background(IrisColors.surface) + .padding(horizontal = 6.dp, vertical = 4.dp), + verticalAlignment = Alignment.CenterVertically, ) { IconButton(onClick = { showPicker = true }) { Text("๐Ÿ“Ž", fontSize = 18.sp) } - OutlinedTextField( - value = input, - onValueChange = { input = it }, + Box( modifier = Modifier .weight(1f) - .heightIn(min = 48.dp, max = 160.dp) - .focusRequester(composerFocusRequester) - .onFocusChanged { composerFocused = it.isFocused }, - placeholder = { Text("Message") }, - keyboardOptions = KeyboardOptions(imeAction = ImeAction.Send), - keyboardActions = KeyboardActions(onSend = { doSend() }), - ) - Spacer(modifier = Modifier.width(8.dp)) - Button( - onClick = { doSend() }, - enabled = input.isNotBlank() || attachments.any { it.mediaRef != null && it.error == null }, + .heightIn(min = 40.dp, max = 160.dp) + .padding(horizontal = 6.dp, vertical = 10.dp), ) { - Text("Send") + if (input.isEmpty()) { + Text( + "Message", + color = IrisColors.textDim, + style = MaterialTheme.typography.bodyLarge, + modifier = Modifier.align(Alignment.TopStart), + ) + } + BasicTextField( + value = input, + onValueChange = { input = it }, + modifier = Modifier + .fillMaxSize() + .focusRequester(composerFocusRequester) + .onFocusChanged { composerFocused = it.isFocused }, + textStyle = MaterialTheme.typography.bodyLarge.copy(color = IrisColors.onSurface), + keyboardOptions = KeyboardOptions(imeAction = ImeAction.Send), + keyboardActions = KeyboardActions(onSend = { doSend() }), + ) } + Box( + modifier = Modifier + .size(40.dp) + .clip(CircleShape) + .background(if (canSend) IrisColors.primary else IrisColors.primary.copy(alpha = 0.35f)) + .clickable(enabled = canSend) { doSend() }, + contentAlignment = Alignment.Center, + ) { + Text("โžค", color = Color.White, fontSize = 16.sp) + } + } + + // M7: overflow actions (rename / forget pairing) + if (showRename) { + NameDialog( + title = "Rename channel", + onConfirm = { name -> + controller.renameChannel(currentChatId, name) + showRename = false + }, + onDismiss = { showRename = false }, + initialName = currentChannel?.name.orEmpty(), + confirmLabel = "Rename", + ) + } + if (showForget) { + AlertDialog( + onDismissRequest = { showForget = false }, + title = { Text("Forget pairing?") }, + text = { + Text("This device will be unpaired from the gateway. You'll need the pairing token to connect again.") + }, + confirmButton = { + TextButton(onClick = { + showForget = false + controller.forget() + }) { Text("Forget") } + }, + dismissButton = { + TextButton(onClick = { showForget = false }) { Text("Cancel") } + }, + ) } } } @@ -367,12 +538,12 @@ fun ChatScreen(controller: IrisController) { onSetDefault = { ch -> controller.setDefaultChannel(ch.chatId) }, onDelete = { ch -> controller.deleteChannel(ch.chatId) }, ) - Box(modifier = Modifier.width(1.dp).fillMaxHeight().background(Color(0xFF2A2E3B))) + Box(modifier = Modifier.width(1.dp).fillMaxHeight().background(IrisColors.divider)) Box(modifier = Modifier.weight(1f).fillMaxHeight()) { ChatContent(showHamburger = false, onOpenDrawer = {}) } if (showInspector) { - Box(modifier = Modifier.width(1.dp).fillMaxHeight().background(Color(0xFF2A2E3B))) + Box(modifier = Modifier.width(1.dp).fillMaxHeight().background(IrisColors.divider)) InspectorPane( controller = controller, state = state, @@ -492,6 +663,147 @@ fun ChatScreen(controller: IrisController) { } } +/** M7: a row in the message list โ€” a chat item or a date separator. */ +private sealed interface ChatRow { + val key: String + + data class Item(val item: ChatItem) : ChatRow { + override val key: String get() = item.id + } + + data class Day(val day: String, val label: String) : ChatRow { + override val key: String get() = "day_$day" + } +} + +/** M7: interleave a date separator wherever the local calendar day changes. */ +private fun buildChatRows(items: List): List { + val rows = mutableListOf() + var lastDay: String? = null + for (item in items) { + val ts = (item as? MessageItem)?.ts + if (ts != null && ts > 0) { + val day = localDayKey(ts) + if (day != lastDay) { + rows += ChatRow.Day(day, formatDayLabel(ts)) + lastDay = day + } + } + rows += ChatRow.Item(item) + } + return rows +} + +/** M7: centered date-separator pill (reference: "7. August"). */ +@Composable +private fun DaySeparator(label: String) { + Box(modifier = Modifier.fillMaxWidth(), contentAlignment = Alignment.Center) { + Box( + modifier = Modifier + .clip(RoundedCornerShape(12.dp)) + .background(IrisColors.panel) + .padding(horizontal = 12.dp, vertical = 4.dp), + ) { + Text(label, fontSize = 12.sp, color = IrisColors.textSecondary) + } + } +} + +/** M7: header title pill โ€” letter avatar + channel name + "Bot" subtitle. */ +@Composable +private fun TitlePill(channelName: String, modifier: Modifier = Modifier) { + Box( + modifier = modifier + .clip(RoundedCornerShape(16.dp)) + .background(IrisColors.surface) + .padding(start = 8.dp, end = 12.dp, top = 4.dp, bottom = 4.dp), + ) { + Row(verticalAlignment = Alignment.CenterVertically) { + LetterAvatar(channelName, 34.dp, IrisColors.primary) + Spacer(modifier = Modifier.width(10.dp)) + Column { + Text( + channelName, + style = MaterialTheme.typography.titleMedium, + fontWeight = FontWeight.SemiBold, + maxLines = 1, + overflow = TextOverflow.Ellipsis, + ) + Text( + "Bot", + style = MaterialTheme.typography.bodySmall, + color = IrisColors.textDim, + ) + } + } + } +} + +/** M7: compact 36dp header icon button (smaller than the 48dp default). */ +@Composable +private fun HeaderIconButton(onClick: () -> Unit, content: @Composable () -> Unit) { + Box( + modifier = Modifier + .size(36.dp) + .clip(RoundedCornerShape(10.dp)) + .clickable(onClick = onClick), + contentAlignment = Alignment.Center, + ) { + content() + } +} + +/** M7: header overflow โ‹ฎ โ€” rename channel / forget pairing. */ +@Composable +private fun HeaderOverflow( + expanded: Boolean, + onToggle: () -> Unit, + onRename: () -> Unit, + onForget: () -> Unit, +) { + HeaderIconButton(onClick = onToggle) { + Text("โ‹ฎ", fontSize = 18.sp) + } + if (expanded) { + DropdownMenu(expanded = expanded, onDismissRequest = onToggle) { + DropdownMenuItem( + text = { Text("Rename channel") }, + onClick = { + onToggle() + onRename() + }, + ) + DropdownMenuItem( + text = { Text("Forget pairing") }, + onClick = { + onToggle() + onForget() + }, + ) + } + } +} + +/** M7: circle with the first letter of [name] (header / rail / drawer). */ +@Composable +private fun LetterAvatar(name: String, size: Dp, color: Color) { + val letter = name.trim().firstOrNull()?.uppercaseChar()?.toString() ?: "?" + Box( + modifier = Modifier + .size(size) + .clip(CircleShape) + .background(color), + contentAlignment = Alignment.Center, + ) { + Text( + letter, + color = Color.White, + fontSize = (size.value * 0.45f).sp, + fontWeight = FontWeight.SemiBold, + ) + } +} + /** Left drawer: channel list + create / set-default / delete. */ @Composable private fun ChannelDrawer( @@ -518,14 +830,18 @@ private fun ChannelDrawer( modifier = Modifier .fillMaxWidth() .clip(RoundedCornerShape(8.dp)) - .background(if (isCurrent) Color(0xFF2A3550) else Color.Transparent) + .background(if (isCurrent) IrisColors.chipSelected else Color.Transparent) .clickable { onOpen(ch) } .padding(horizontal = 10.dp, vertical = 6.dp), verticalAlignment = Alignment.CenterVertically, ) { + LetterAvatar(ch.name, 28.dp, avatarColor(ch.name)) + Spacer(modifier = Modifier.width(8.dp)) Text( ch.name + if (ch.isDefault) " โ˜…" else "", fontSize = 14.sp, + maxLines = 1, + overflow = TextOverflow.Ellipsis, modifier = Modifier.weight(1f), ) if (!ch.isDefault) { @@ -607,8 +923,8 @@ private fun ChannelRail( .clip(RoundedCornerShape(8.dp)) .background( when { - isCurrent -> Color(0xFF2A3550) - isSelected -> Color(0xFF232838) + isCurrent -> IrisColors.chipSelected + isSelected -> IrisColors.rowSelected else -> Color.Transparent }, ) @@ -616,9 +932,23 @@ private fun ChannelRail( .padding(horizontal = 10.dp, vertical = 6.dp), verticalAlignment = Alignment.CenterVertically, ) { + // M7: accent bar marks the active channel. + Box( + modifier = Modifier + .width(3.dp) + .height(22.dp) + .clip(RoundedCornerShape(2.dp)) + .background(if (isCurrent) IrisColors.primary else Color.Transparent), + ) + Spacer(modifier = Modifier.width(8.dp)) + LetterAvatar(ch.name, 28.dp, avatarColor(ch.name)) + Spacer(modifier = Modifier.width(8.dp)) Text( ch.name + if (ch.isDefault) " โ˜…" else "", fontSize = 14.sp, + color = if (isCurrent) IrisColors.primary else IrisColors.onSurface, + maxLines = 1, + overflow = TextOverflow.Ellipsis, modifier = Modifier.weight(1f), ) if (!ch.isDefault) { @@ -693,7 +1023,7 @@ private fun InspectorPane( @Composable private fun InfoRow(label: String, value: String) { Row(modifier = Modifier.fillMaxWidth().padding(vertical = 2.dp)) { - Text(label, fontSize = 12.sp, color = Color(0xFF8A93A6), modifier = Modifier.width(70.dp)) + Text(label, fontSize = 12.sp, color = IrisColors.textDim, modifier = Modifier.width(70.dp)) Text(value, fontSize = 12.sp, modifier = Modifier.weight(1f)) } } @@ -728,7 +1058,7 @@ private fun CommandPalette( ) Spacer(modifier = Modifier.height(8.dp)) if (filtered.isEmpty()) { - Text("No matching commands.", fontSize = 12.sp, color = Color(0xFF8A93A6)) + Text("No matching commands.", fontSize = 12.sp, color = IrisColors.textDim) } else { filtered.forEach { (id, label) -> Row( @@ -778,11 +1108,11 @@ private fun TopicChip(label: String, selected: Boolean, onClick: () -> Unit) { Box( modifier = Modifier .clip(RoundedCornerShape(12.dp)) - .background(if (selected) Color(0xFF4F7CFF) else Color(0xFF2A2E3B)) + .background(if (selected) IrisColors.primary else IrisColors.chip) .clickable(onClick = onClick) .padding(horizontal = 10.dp, vertical = 5.dp), ) { - Text(label, fontSize = 12.sp, color = if (selected) Color.White else Color(0xFFC7CCD8)) + Text(label, fontSize = 12.sp, color = if (selected) Color.White else IrisColors.textTertiary) } } @@ -839,7 +1169,7 @@ private fun SearchOverlay( horizontalArrangement = Arrangement.spacedBy(8.dp), modifier = Modifier.padding(top = 4.dp), ) { - Text("Scope:", fontSize = 12.sp, color = Color(0xFF8A93A6)) + Text("Scope:", fontSize = 12.sp, color = IrisColors.textDim) TextButton(onClick = { scope = "all" }) { Text(if (scope == "all") "โ— all" else "โ—‹ all", fontSize = 12.sp) } @@ -871,7 +1201,7 @@ private fun SearchHitRow(hit: SearchHit, channels: ChannelStore, onClick: () -> modifier = Modifier .fillMaxWidth() .clip(RoundedCornerShape(10.dp)) - .background(Color(0xFF20242E)) + .background(IrisColors.panel) .clickable(onClick = onClick) .padding(10.dp), ) { @@ -880,20 +1210,26 @@ private fun SearchHitRow(hit: SearchHit, channels: ChannelStore, onClick: () -> (channel?.name ?: hit.chatId) + (thread?.let { " / ${it.name}" } ?: ""), fontSize = 12.sp, fontWeight = FontWeight.Medium, - color = Color(0xFFB9C0D0), + color = IrisColors.textSecondary, modifier = Modifier.weight(1f), ) - Text(hit.role, fontSize = 10.sp, color = Color(0xFF8A93A6)) + Text(hit.role, fontSize = 10.sp, color = IrisColors.textDim) } Spacer(modifier = Modifier.height(4.dp)) - Text(hit.snippet, fontSize = 13.sp, color = Color(0xFFD7DBE5), maxLines = 3) + Text(hit.snippet, fontSize = 13.sp, color = IrisColors.textBright, maxLines = 3) } } /** Simple name-entry dialog (new channel / new topic). */ @Composable -private fun NameDialog(title: String, onConfirm: (String) -> Unit, onDismiss: () -> Unit) { - var name by remember { mutableStateOf("") } +private fun NameDialog( + title: String, + onConfirm: (String) -> Unit, + onDismiss: () -> Unit, + initialName: String = "", + confirmLabel: String = "Create", +) { + var name by remember { mutableStateOf(initialName) } AlertDialog( onDismissRequest = onDismiss, title = { Text(title) }, @@ -905,7 +1241,7 @@ private fun NameDialog(title: String, onConfirm: (String) -> Unit, onDismiss: () ) }, confirmButton = { - TextButton(onClick = { onConfirm(name) }, enabled = name.isNotBlank()) { Text("Create") } + TextButton(onClick = { onConfirm(name) }, enabled = name.isNotBlank()) { Text(confirmLabel) } }, dismissButton = { TextButton(onClick = onDismiss) { Text("Cancel") } @@ -916,11 +1252,11 @@ private fun NameDialog(title: String, onConfirm: (String) -> Unit, onDismiss: () @Composable private fun StatusChip(state: GatewayClient.State) { val (label, color) = when (state) { - GatewayClient.State.Disconnected -> "offline" to Color(0xFF9E9E9E) - GatewayClient.State.Connecting -> "connectingโ€ฆ" to Color(0xFFFFC107) - GatewayClient.State.Reconnecting -> "reconnectingโ€ฆ" to Color(0xFFFFC107) - is GatewayClient.State.Connected -> "connected" to Color(0xFF4CAF50) - is GatewayClient.State.AuthFailed -> "auth failed" to Color(0xFFF44336) + GatewayClient.State.Disconnected -> "offline" to IrisColors.statusGrey + GatewayClient.State.Connecting -> "connectingโ€ฆ" to IrisColors.statusAmber + GatewayClient.State.Reconnecting -> "reconnectingโ€ฆ" to IrisColors.statusAmber + is GatewayClient.State.Connected -> "connected" to IrisColors.statusGreen + is GatewayClient.State.AuthFailed -> "auth failed" to IrisColors.statusRed } Box( modifier = Modifier @@ -940,10 +1276,10 @@ private fun NotificationBanner( onDismiss: () -> Unit, ) { val accent = when (banner.kind) { - "approval" -> Color(0xFFFFC107) - "clarify" -> Color(0xFF4F7CFF) - "cron" -> Color(0xFF4CAF50) - else -> Color(0xFF8A93A6) + "approval" -> IrisColors.statusAmber + "clarify" -> IrisColors.primary + "cron" -> IrisColors.statusGreen + else -> IrisColors.textDim } Row( modifier = Modifier @@ -957,7 +1293,7 @@ private fun NotificationBanner( Column(modifier = Modifier.weight(1f)) { Text(banner.title, fontSize = 13.sp, fontWeight = FontWeight.SemiBold, color = accent) if (banner.body.isNotBlank()) { - Text(banner.body, fontSize = 12.sp, color = Color(0xFFD7DBE5), maxLines = 2) + Text(banner.body, fontSize = 12.sp, color = IrisColors.textBright, maxLines = 2) } } TextButton(onClick = onDismiss) { Text("โœ•", fontSize = 12.sp) } @@ -965,27 +1301,34 @@ private fun NotificationBanner( } @Composable -private fun MessageBubble(msg: MessageItem) { +private fun MessageBubble(msg: MessageItem, maxWidth: Dp, onRetry: () -> Unit) { val isUser = msg.role == ROLE_USER val isCommentary = msg.isCommentary val bubbleColor = when { - isUser -> Color(0xFF4F7CFF) - isCommentary -> Color(0xFF23262F) - else -> Color(0xFF2A2E3B) + isUser -> IrisColors.bubbleUser + isCommentary -> IrisColors.bubbleCommentary + else -> IrisColors.bubbleAssistant } val textColor = when { isUser -> Color.White - isCommentary -> Color(0xFFE8EAF0).copy(alpha = 0.55f) - else -> Color(0xFFE8EAF0) + isCommentary -> IrisColors.onSurface.copy(alpha = 0.55f) + else -> IrisColors.onSurface } + val time = formatTime(msg.ts) Row( modifier = Modifier.fillMaxWidth(), horizontalArrangement = if (isUser) Arrangement.End else Arrangement.Start, ) { Column( modifier = Modifier + .widthIn(max = maxWidth) .clip(RoundedCornerShape(14.dp)) .background(bubbleColor) + .then( + if (isUser && msg.status == MsgStatus.Failed) { + Modifier.clickable { onRetry() } + } else Modifier + ) .padding(horizontal = 12.dp, vertical = 8.dp), ) { // Reasoning block above the answer (assistant, non-commentary). @@ -993,12 +1336,49 @@ private fun MessageBubble(msg: MessageItem) { ReasoningBlock(msg.reasoning!!) Spacer(modifier = Modifier.height(6.dp)) } - if (msg.text.isNotBlank() || msg.streaming) { - Text( - msg.text + if (msg.streaming) " โ–‰" else "", - color = textColor, - fontSize = if (isCommentary) 13.sp else 15.sp, - ) + if (isUser) { + // M7: text + inline timestamp / delivery status at the end + // (reference look); the bubble hugs its content. + Row(verticalAlignment = Alignment.Bottom) { + if (msg.text.isNotBlank() || msg.streaming) { + Text( + msg.text + if (msg.streaming) " โ–‰" else "", + color = textColor, + fontSize = 15.sp, + ) + } + Spacer(modifier = Modifier.width(8.dp)) + Row(verticalAlignment = Alignment.CenterVertically) { + if (time.isNotEmpty()) { + Text(time, fontSize = 10.sp, color = textColor.copy(alpha = 0.6f)) + } + when (msg.status) { + MsgStatus.Pending -> + Text(" sendingโ€ฆ", fontSize = 10.sp, color = textColor.copy(alpha = 0.6f)) + MsgStatus.Sent -> + Text(" โœ“", fontSize = 10.sp, color = textColor.copy(alpha = 0.8f)) + MsgStatus.Read -> + Text(" โœ“โœ“", fontSize = 10.sp, color = textColor.copy(alpha = 0.8f)) + MsgStatus.Failed -> Unit + } + } + } + if (msg.status == MsgStatus.Failed) { + Row( + modifier = Modifier.fillMaxWidth(), + horizontalArrangement = Arrangement.End, + ) { + Text("Failed to send โ€” tap to retry", fontSize = 10.sp, color = IrisColors.errorText) + } + } + } else { + if (msg.text.isNotBlank() || msg.streaming) { + Text( + msg.text + if (msg.streaming) " โ–‰" else "", + color = textColor, + fontSize = if (isCommentary) 13.sp else 15.sp, + ) + } } // M4: media attachments (image / player / document chip). if (msg.media.isNotEmpty()) { @@ -1007,23 +1387,31 @@ private fun MessageBubble(msg: MessageItem) { msg.media.forEach { m -> MediaAttachment(m) } } } - if (msg.pending) { - Text("sendingโ€ฆ", color = textColor.copy(alpha = 0.6f), fontSize = 10.sp) - } - // Model / token footer (final assistant answers only). - if (!isUser && !isCommentary && !msg.streaming && - (msg.model != null || msg.tokens != null) - ) { - Spacer(modifier = Modifier.height(4.dp)) - Text( - buildString { - msg.model?.let { append(it) } - if (msg.model != null && msg.tokens != null) append(" ยท ") - msg.tokens?.let { append("${it} tok") } - }, - color = textColor.copy(alpha = 0.4f), - fontSize = 10.sp, - ) + if (!isUser) { + // Model / token footer (final assistant answers only). + if (!isCommentary && !msg.streaming && + (msg.model != null || msg.tokens != null) + ) { + Spacer(modifier = Modifier.height(4.dp)) + Text( + buildString { + msg.model?.let { append(it) } + if (msg.model != null && msg.tokens != null) append(" ยท ") + msg.tokens?.let { append("${it} tok") } + }, + color = textColor.copy(alpha = 0.4f), + fontSize = 10.sp, + ) + } + // M7: timestamp below the footer (reference shows "12:41"). + if (time.isNotEmpty()) { + Row( + modifier = Modifier.fillMaxWidth(), + horizontalArrangement = Arrangement.End, + ) { + Text(time, fontSize = 10.sp, color = textColor.copy(alpha = 0.5f)) + } + } } } } @@ -1063,7 +1451,7 @@ private fun MediaImage(media: MediaItem) { modifier = Modifier .size(120.dp) .clip(RoundedCornerShape(10.dp)) - .background(Color(0xFF20242E)), + .background(IrisColors.panel), contentAlignment = Alignment.Center, ) { if (media.localPath == null) Text("downloadingโ€ฆ", fontSize = 11.sp) @@ -1113,7 +1501,7 @@ private fun MediaDocChip(media: MediaItem) { Row( modifier = Modifier .clip(RoundedCornerShape(10.dp)) - .background(Color(0xFF20242E)) + .background(IrisColors.panel) .clickable { media.localPath?.let { openDocument(it, media.mime) } } @@ -1124,7 +1512,7 @@ private fun MediaDocChip(media: MediaItem) { Spacer(modifier = Modifier.width(8.dp)) Column { Text(media.filename, fontSize = 13.sp, maxLines = 1) - Text(fmtSize(media.size), fontSize = 10.sp, color = Color(0xFF8A93A6)) + Text(fmtSize(media.size), fontSize = 10.sp, color = IrisColors.textDim) } } } @@ -1146,13 +1534,13 @@ private fun AttachmentChip(att: IrisController.PendingAttachment, onRemove: () - Row( modifier = Modifier .clip(RoundedCornerShape(12.dp)) - .background(Color(0xFF2A2E3B)) + .background(IrisColors.chip) .padding(horizontal = 8.dp, vertical = 5.dp), verticalAlignment = Alignment.CenterVertically, ) { Text(icon, fontSize = 12.sp) Spacer(modifier = Modifier.width(4.dp)) - Text(label, fontSize = 12.sp, color = Color(0xFFC7CCD8), maxLines = 1) + Text(label, fontSize = 12.sp, color = IrisColors.textTertiary, maxLines = 1) Spacer(modifier = Modifier.width(4.dp)) TextButton(onClick = onRemove) { Text("โœ•", fontSize = 11.sp) } } @@ -1173,7 +1561,7 @@ private fun ReasoningBlock(reasoning: String) { modifier = Modifier .fillMaxWidth() .clip(RoundedCornerShape(8.dp)) - .background(Color(0xFF1E212B)) + .background(IrisColors.reasoningPanel) .padding(8.dp), ) { Row(verticalAlignment = Alignment.CenterVertically) { @@ -1181,7 +1569,7 @@ private fun ReasoningBlock(reasoning: String) { "๐Ÿ’ญ Reasoning", fontSize = 12.sp, fontWeight = FontWeight.SemiBold, - color = Color(0xFFB9C0D0), + color = IrisColors.textSecondary, modifier = Modifier .weight(1f) .clickable { expanded = !expanded }, @@ -1195,7 +1583,7 @@ private fun ReasoningBlock(reasoning: String) { Text( if (expanded) "โ–พ" else "โ–ธ", fontSize = 12.sp, - color = Color(0xFFB9C0D0), + color = IrisColors.textSecondary, modifier = Modifier.clickable { expanded = !expanded }, ) } @@ -1204,7 +1592,7 @@ private fun ReasoningBlock(reasoning: String) { reasoning, fontFamily = FontFamily.Monospace, fontSize = 12.sp, - color = Color(0xFF9AA3B5), + color = IrisColors.textMuted, ) } } @@ -1223,15 +1611,15 @@ private fun ToolCard(tool: ToolItem, detail: ToolDetail) { else -> "โœ—" } val statusColor = when { - tool.ok -> Color(0xFF4CAF50) - else -> Color(0xFFF44336) + tool.ok -> IrisColors.statusGreen + else -> IrisColors.statusRed } val preview = tool.preview ?: tool.note Column( modifier = Modifier .fillMaxWidth() .clip(RoundedCornerShape(10.dp)) - .background(Color(0xFF20242E)) + .background(IrisColors.panel) .padding(horizontal = 10.dp, vertical = 8.dp), ) { Row(verticalAlignment = Alignment.CenterVertically) { @@ -1249,7 +1637,7 @@ private fun ToolCard(tool: ToolItem, detail: ToolDetail) { "๐Ÿ”ง ${tool.name}", fontSize = 13.sp, fontWeight = FontWeight.Medium, - color = Color(0xFFD7DBE5), + color = IrisColors.textBright, modifier = Modifier .weight(1f) .clickable { expanded = !expanded }, @@ -1258,7 +1646,7 @@ private fun ToolCard(tool: ToolItem, detail: ToolDetail) { Text( "${String.format("%.1f", tool.duration)}s", fontSize = 11.sp, - color = Color(0xFF8A93A6), + color = IrisColors.textDim, ) } } @@ -1267,7 +1655,7 @@ private fun ToolCard(tool: ToolItem, detail: ToolDetail) { Text( preview, fontSize = 12.sp, - color = Color(0xFF9AA3B5), + color = IrisColors.textMuted, maxLines = if (expanded) Int.MAX_VALUE else 1, ) } @@ -1279,7 +1667,7 @@ private fun ToolCard(tool: ToolItem, detail: ToolDetail) { it.toString(), fontFamily = FontFamily.Monospace, fontSize = 11.sp, - color = Color(0xFF8A93A6), + color = IrisColors.textDim, ) } tool.outputPreview?.let { @@ -1288,7 +1676,7 @@ private fun ToolCard(tool: ToolItem, detail: ToolDetail) { it, fontFamily = FontFamily.Monospace, fontSize = 11.sp, - color = Color(0xFF8A93A6), + color = IrisColors.textDim, ) } } diff --git a/app/shared/src/commonMain/kotlin/iris/ui/screens/ConnectScreen.kt b/app/shared/src/commonMain/kotlin/iris/ui/screens/ConnectScreen.kt index 3108cdf..e670036 100644 --- a/app/shared/src/commonMain/kotlin/iris/ui/screens/ConnectScreen.kt +++ b/app/shared/src/commonMain/kotlin/iris/ui/screens/ConnectScreen.kt @@ -1,5 +1,6 @@ package iris.ui.screens +import androidx.compose.foundation.background import androidx.compose.foundation.layout.Arrangement import androidx.compose.foundation.layout.Box import androidx.compose.foundation.layout.Column @@ -9,11 +10,12 @@ import androidx.compose.foundation.layout.fillMaxWidth import androidx.compose.foundation.layout.height import androidx.compose.foundation.layout.padding import androidx.compose.foundation.rememberScrollState +import androidx.compose.foundation.shape.RoundedCornerShape import androidx.compose.foundation.text.KeyboardOptions import androidx.compose.foundation.verticalScroll import androidx.compose.material3.Button +import androidx.compose.material3.CircularProgressIndicator import androidx.compose.material3.MaterialTheme -import androidx.compose.material3.OutlinedButton import androidx.compose.material3.OutlinedTextField import androidx.compose.material3.Text import androidx.compose.runtime.Composable @@ -24,10 +26,12 @@ import androidx.compose.runtime.rememberCoroutineScope import androidx.compose.runtime.setValue import androidx.compose.ui.Alignment import androidx.compose.ui.Modifier +import androidx.compose.ui.draw.clip import androidx.compose.ui.text.input.KeyboardType import androidx.compose.ui.text.input.PasswordVisualTransformation import androidx.compose.ui.unit.dp import iris.state.IrisController +import iris.ui.theme.IrisColors import kotlinx.coroutines.launch /** @@ -64,53 +68,61 @@ fun ConnectScreen( ) Spacer(modifier = Modifier.height(32.dp)) - OutlinedTextField( - value = url, - onValueChange = { url = it }, - label = { Text("Server URL") }, - placeholder = { Text("ws://192.168.1.10:8790/ws") }, - singleLine = true, - keyboardOptions = KeyboardOptions(keyboardType = KeyboardType.Uri), - modifier = Modifier.fillMaxWidth(), - ) - Spacer(modifier = Modifier.height(12.dp)) - OutlinedTextField( - value = token, - onValueChange = { token = it }, - label = { Text("Pairing token") }, - placeholder = { Text("ANDROID_TOKEN (64 hex)") }, - singleLine = true, - visualTransformation = PasswordVisualTransformation(), - modifier = Modifier.fillMaxWidth(), - ) - Spacer(modifier = Modifier.height(24.dp)) - - Button( - onClick = { - if (busy) return@Button - busy = true - error = null - scope.launch { - val result = controller.connect(url.trim(), token.trim()) - busy = false - if (result.isFailure) { - error = result.exceptionOrNull()?.message ?: "connection failed" - } - } - }, - enabled = !busy, - modifier = Modifier.fillMaxWidth(), + Column( + modifier = Modifier + .fillMaxWidth() + .clip(RoundedCornerShape(16.dp)) + .background(IrisColors.surface) + .padding(16.dp), ) { - Text(if (busy) "Testing connectionโ€ฆ" else "Test & Connect") - } - - if (error != null) { - Spacer(modifier = Modifier.height(16.dp)) - Text( - error!!, - color = MaterialTheme.colorScheme.error, - style = MaterialTheme.typography.bodyMedium, + OutlinedTextField( + value = url, + onValueChange = { url = it }, + label = { Text("Server URL") }, + placeholder = { Text("ws://192.168.1.10:8790/ws") }, + singleLine = true, + keyboardOptions = KeyboardOptions(keyboardType = KeyboardType.Uri), + modifier = Modifier.fillMaxWidth(), ) + Spacer(modifier = Modifier.height(12.dp)) + OutlinedTextField( + value = token, + onValueChange = { token = it }, + label = { Text("Pairing token") }, + placeholder = { Text("ANDROID_TOKEN (64 hex)") }, + singleLine = true, + visualTransformation = PasswordVisualTransformation(), + modifier = Modifier.fillMaxWidth(), + ) + Spacer(modifier = Modifier.height(24.dp)) + + Button( + onClick = { + if (busy) return@Button + busy = true + error = null + scope.launch { + val result = controller.connect(url.trim(), token.trim()) + busy = false + if (result.isFailure) { + error = result.exceptionOrNull()?.message ?: "connection failed" + } + } + }, + enabled = !busy, + modifier = Modifier.fillMaxWidth(), + ) { + Text(if (busy) "Testing connectionโ€ฆ" else "Test & Connect") + } + + if (error != null) { + Spacer(modifier = Modifier.height(16.dp)) + Text( + error!!, + color = MaterialTheme.colorScheme.error, + style = MaterialTheme.typography.bodyMedium, + ) + } } Spacer(modifier = Modifier.height(24.dp)) @@ -121,4 +133,22 @@ fun ConnectScreen( color = MaterialTheme.colorScheme.onSurfaceVariant, ) } +} + +/** M7: shown while the initial connection is in flight (state == Connecting). */ +@Composable +fun ConnectingScreen(host: String) { + Box(modifier = Modifier.fillMaxSize(), contentAlignment = Alignment.Center) { + Column( + horizontalAlignment = Alignment.CenterHorizontally, + verticalArrangement = Arrangement.spacedBy(16.dp), + ) { + CircularProgressIndicator() + Text( + "Connecting to $hostโ€ฆ", + style = MaterialTheme.typography.bodyMedium, + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + } + } } \ No newline at end of file diff --git a/app/shared/src/commonMain/kotlin/iris/ui/theme/Theme.kt b/app/shared/src/commonMain/kotlin/iris/ui/theme/Theme.kt new file mode 100644 index 0000000..744ff70 --- /dev/null +++ b/app/shared/src/commonMain/kotlin/iris/ui/theme/Theme.kt @@ -0,0 +1,72 @@ +package iris.ui.theme + +import androidx.compose.material3.MaterialTheme +import androidx.compose.material3.darkColorScheme +import androidx.compose.runtime.Composable +import androidx.compose.ui.graphics.Color + +/** + * Centralized palette (M7). Dark is the only theme; every screen reads colors + * from here instead of hard-coding literals. + */ +object IrisColors { + // Base + val background = Color(0xFF1B1E28) + val surface = Color(0xFF222634) + val onBackground = Color(0xFFE8EAF0) + val onSurface = Color(0xFFE8EAF0) + val primary = Color(0xFF4F7CFF) + + // Text shades (brightest to dimmest) + val textBright = Color(0xFFD7DBE5) + val textTertiary = Color(0xFFC7CCD8) + val textSecondary = Color(0xFFB9C0D0) + val textMuted = Color(0xFF9AA3B5) + val textDim = Color(0xFF8A93A6) + + // Bubbles + val bubbleUser = primary + val bubbleAssistant = Color(0xFF2A2E3B) + val bubbleCommentary = Color(0xFF23262F) + + // Panels, chips, rows + val panel = Color(0xFF20242E) + val chip = Color(0xFF2A2E3B) + val chipSelected = Color(0xFF2A3550) + val rowSelected = Color(0xFF232838) + val reasoningPanel = Color(0xFF1E212B) + val divider = Color(0xFF2A2E3B) + + // Status + val statusGrey = Color(0xFF9E9E9E) + val statusAmber = Color(0xFFFFC107) + val statusGreen = Color(0xFF4CAF50) + val statusRed = Color(0xFFF44336) + + // M7: failed-send hint on the accent user bubble + val errorText = Color(0xFFFF8A80) +} + +private val IrisDark = darkColorScheme( + background = IrisColors.background, + surface = IrisColors.surface, + onBackground = IrisColors.onBackground, + onSurface = IrisColors.onSurface, + onSurfaceVariant = IrisColors.textDim, + primary = IrisColors.primary, + onPrimary = Color.White, + error = IrisColors.statusRed, +) + +@Composable +fun IrisTheme(content: @Composable () -> Unit) { + MaterialTheme(colorScheme = IrisDark) { + content() + } +} + +/** Deterministic pastel avatar color from a channel name (M7 rail/drawer). */ +fun avatarColor(name: String): Color { + val hue = ((name.hashCode() and 0x7FFFFFFF) % 360).toFloat() + return Color.hsv(hue, 0.5f, 0.75f) +} \ No newline at end of file diff --git a/app/shared/src/commonMain/kotlin/iris/util/TimeFormat.kt b/app/shared/src/commonMain/kotlin/iris/util/TimeFormat.kt new file mode 100644 index 0000000..b6f8458 --- /dev/null +++ b/app/shared/src/commonMain/kotlin/iris/util/TimeFormat.kt @@ -0,0 +1,16 @@ +package iris.util + +/** Locale-aware day label for chat date separators (M7), e.g. "7. August". */ +expect fun formatDayLabel(epochMillis: Long): String + +/** Locale-aware "HH:mm" time for bubble timestamps (M7). */ +expect fun formatTime(epochMillis: Long): String + +/** Local calendar-day key used to detect date changes between messages (M7). */ +expect fun localDayKey(epochMillis: Long): String + +/** Host part of a pairing URL ("ws://host:port/ws" -> "host:port"). */ +fun hostFromUrl(url: String): String { + val noScheme = url.trim().substringAfter("://") + return noScheme.substringBefore("/").ifBlank { url.trim() } +} \ No newline at end of file diff --git a/app/shared/src/jvmMain/kotlin/iris/util/TimeFormatJvm.kt b/app/shared/src/jvmMain/kotlin/iris/util/TimeFormatJvm.kt new file mode 100644 index 0000000..fa9b501 --- /dev/null +++ b/app/shared/src/jvmMain/kotlin/iris/util/TimeFormatJvm.kt @@ -0,0 +1,27 @@ +package iris.util + +import java.time.Instant +import java.time.ZoneId +import java.time.format.DateTimeFormatter +import java.util.Locale + +// jvmMain is the intermediate source set for both androidMain (minSdk 29, +// java.time available) and desktopMain, so one actual covers both targets. + +actual fun formatDayLabel(epochMillis: Long): String { + if (epochMillis <= 0) return "" + val zone = ZoneId.systemDefault() + val date = Instant.ofEpochMilli(epochMillis).atZone(zone) + val now = Instant.now().atZone(zone) + val pattern = if (date.year == now.year) "d. MMMM" else "d. MMMM, yyyy" + return DateTimeFormatter.ofPattern(pattern, Locale.getDefault()).format(date) +} + +actual fun formatTime(epochMillis: Long): String { + if (epochMillis <= 0) return "" + return DateTimeFormatter.ofPattern("HH:mm", Locale.getDefault()) + .format(Instant.ofEpochMilli(epochMillis).atZone(ZoneId.systemDefault())) +} + +actual fun localDayKey(epochMillis: Long): String = + Instant.ofEpochMilli(epochMillis).atZone(ZoneId.systemDefault()).toLocalDate().toString() \ No newline at end of file diff --git a/docs/04-wire-protocol.md b/docs/04-wire-protocol.md index 1075e4a..224636b 100644 --- a/docs/04-wire-protocol.md +++ b/docs/04-wire-protocol.md @@ -183,10 +183,21 @@ Agent-sent media is available; app pulls bytes. "media_id":"md_5","kind":"video","mime":"video/mp4","size":123456,"filename":"clip.mp4"}} ``` -### `status` -Gateway lifecycle / session info. +### `read.receipt` +The gateway acknowledges that the agent has received and started processing +the user's message. The app uses it to show โœ“โœ“ on user bubbles. ```json -{"type":"status","payload":{"state":"online","session":{"chat_id":"โ€ฆ","model":"โ€ฆ","tokens":11}}} +{"type":"read.receipt","chat_id":"android:default","payload":{"message_id":"m_9001"}} +``` +Emitted to the originating connection when a `message.send` is accepted for +processing (at the moment it is handed to the agent), for user-originated +messages only. + +### `status` +Gateway health state. Broadcast to all connected clients at startup +(`state: "online"`); `restarting` / `degraded` are reserved for future use. +```json +{"type":"status","payload":{"state":"online"}} ``` `state` โˆˆ `online | restarting | degraded`. diff --git a/docs/09-pairing-security.md b/docs/09-pairing-security.md index 3a2322c..6dc816e 100644 --- a/docs/09-pairing-security.md +++ b/docs/09-pairing-security.md @@ -92,4 +92,26 @@ security principal (the token is). - [ ] Redact all secrets in logs. - [ ] WSS + cert pinning for remote. - [ ] Outbox retention cap + prune. -- [ ] Fail-closed secret reads under multiplexing. \ No newline at end of file +- [ ] Fail-closed secret reads under multiplexing. + +## M7 verification (2026-08-20) + +Status of the ยง9.7 hardening checklist plus the related gaps found in the +M7 research pass. "verified" = implemented and covered by +`hermes-agent/tests/gateway/test_android.py` (35 tests) or the app build; +"gap" = known limitation with the planned mitigation. + +| # | Item | Status | Evidence / mitigation | +|---|------|--------|-----------------------| +| 1 | Constant-time token compare | verified | `gateway-plugin/pairing.py:34` (`hmac.compare_digest`); `test_wrong_token_rejected` | +| 2 | Bounded per-connection send buffer + rate limit on inbound frames | verified | Send: `SEND_TIMEOUT_S` bounds every outbound send (`ws_server.py:47`, `broadcast`/`send_to`). Inbound: per-connection token bucket on JSON frames (20/s, burst 40) โ†’ `error {code:"rate_limited"}` + close on exceed (`ws_server.py:55`, `_TokenBucket`, `_on_frame`); binary upload chunks exempt (see gap 1) | +| 3 | Reject oversized frames / uploads (`max_upload_bytes`) | verified | `serve(max_size=adapter.max_upload_bytes)` (`ws_server.py:139`); per-upload total cap in `media.py` (`create_upload`/`feed`); `test_upload_declared_over_limit_rejected`, `test_upload_midstream_over_limit_rejected` | +| 4 | Verify media sha256 + re-sniff MIME (don't trust client) | verified | `media.py:317` (`complete_upload` digest check), `media.py:147` (`reclassify_kind`); `test_upload_sha256_mismatch_rejected`, `test_reclassify_kind_does_not_trust_client` | +| 5 | Redact all secrets in logs | gap | No mechanical redaction; the token is printed to stdout by design during `hermes gateway setup` (`gateway-plugin/adapter.py:632,650`). Mitigation: stdout is operator-only, not a log file; a redaction pass over gateway logs is planned | +| 6 | WSS + cert pinning for remote | gap (partial) | WSS supported server-side (`ANDROID_WS_CERT`/`ANDROID_WS_KEY`, `ws_server.py:122`); the app builds a default `OkHttpClient` with no `CertificatePinner` (`app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt:87`). Mitigation: remote access requires CA-signed WSS until pinning lands; LAN `ws://` stays the default | +| 7 | Outbox retention cap + prune | verified | `gateway-plugin/outbox.py:48` (`retention_hours` default 72h, `max_rows` cap, `take_overflow_pruned`); `test_outbox_row_cap_prunes_oldest` | +| 8 | Fail-closed secret reads under multiplexing | verified | `_get_scoped_secret` (`gateway-plugin/adapter.py:74`) for `ANDROID_TOKEN`/`ANDROID_WS_CERT`/`ANDROID_WS_KEY`/FCM/ntfy secrets; scoped bind lock in `connect()` (`adapter.py:779`) | +| 9 | Gap: inbound frame rate limiting | implemented | Closes item 2: token bucket in `ws_server.py` (JSON frames only). Binary upload chunks are exempt โ€” a 100 MB upload is 400 ร— 256 KiB frames in a tight loop and would exhaust any sane bucket; uploads are already bounded by per-frame `max_size` + the per-upload total cap | +| 10 | Gap: Android token storage | implemented | `AndroidSecureStore` โ†’ `EncryptedSharedPreferences` (MasterKey AES256_GCM) with one-time migration of the plain `iris` prefs (read old key โ†’ write encrypted โ†’ delete old key); dep in `app/shared/build.gradle.kts` (`app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt`) | +| 11 | Gap: guard not committed | implemented | `.pre-commit-config.yaml` (local hook โ†’ `scripts/guard_hermes_agent.sh --staged`); a fresh clone gets the guard after `pre-commit install` | +| 12 | Gap: in-app QR scanner | gap | Pairing is manual URL+token only; the server prints a QR (`gateway-plugin/adapter.py:648-654`) that any system scanner can read. Plan: in-app camera scan later | \ No newline at end of file diff --git a/docs/14-milestones.md b/docs/14-milestones.md index c6342fa..2fa4ef0 100644 --- a/docs/14-milestones.md +++ b/docs/14-milestones.md @@ -179,18 +179,46 @@ has explicit **acceptance criteria**. Work top-to-bottom; don't skip M0/M1. ## M7 โ€” Polish + E2E + docs **Goal:** ship-quality. -- [ ] Telegram-style layout pass (per reference image): header, bubbles, date +- [x] Telegram-style layout pass (per reference image): header, bubbles, date separators, โœ“โœ“, model/token footer, banner, bottom bar. -- [ ] Theming (dark default, accent), onboarding/pairing UX, empty/loading/ +- [x] Theming (dark default, accent), onboarding/pairing UX, empty/loading/ reconnecting/degraded states with honest copy. -- [ ] Full E2E suite (`13-testing.md` scenarios 1โ€“12) automated where possible. -- [ ] Docs: `docs/protocol/frames.schema.json` finalized; `docs/setup.md` +- [x] Full E2E suite (`13-testing.md` scenarios 1โ€“12) automated where possible. +- [x] Docs: `docs/protocol/frames.schema.json` finalized; `docs/setup.md` (user-facing pairing + FCM/ntfy + remote access); root README. -- [ ] Security hardening checklist (`09-pairing-security.md`) verified. +- [x] Security hardening checklist (`09-pairing-security.md`) verified. - **Demo:** end-to-end on phone + desktop simultaneously; cron into a channel; push; media; search. - **Accept:** all feature-checklist items pass on-device; E2E green; docs complete; `hermes-agent/` still never committed. +- **Status (2026-08-20):** Layout pass verified on-device (MIX 2S): header + with avatar + "Bot" subtitle + overflow menu (rename channel, forget + pairing), centered date-separator pill, bubbles with in-bubble timestamps, + user โœ“/โœ“โœ“ driven by the new `read.receipt` frame (gateway emits it when the + agent takes the message; late-joining clients also get the current `status` + state on hello.ack), model/token footer, letter-avatar channel rail/drawer + with active highlight, restyled bottom bar. Theming centralized in + `ui/theme/Theme.kt` (dark default, single accent; all hard-coded colors + replaced). States: dedicated connecting screen, reconnecting + + degraded/restarting banners (new `status` frame), send-failure rollback + with tap-to-retry. E2E: `tests/e2e.py` driver automates scenarios 1โ€“12 + against the live gateway โ€” 9 PASS / 2 PARTIAL (push device-notification + leg + gateway-kill leg are manual) / 1 SKIP (commentary is + model-dependent) / 0 FAIL; `ws_probe.py` gained `--assert-turn/ + reasoning/tools/commentary/read-receipt/status` plus `--search`, + `--channel-*`, `--watch` modes. Docs: `frames.schema.json` finalized + (mirrors code exactly; 17 unimplemented frames moved to + `x-planned-frames`; 6 deltas + 3 drift fixes), `docs/setup.md` added + (pairing + FCM/ntfy + remote access + troubleshooting), root README + quickstart + status updated. Security: inbound JSON-frame rate limit + (20/s, burst 40, binary upload chunks exempt) with `rate_limited` error + + close; Android token moved to EncryptedSharedPreferences with one-time + plainโ†’encrypted migration; `.pre-commit-config.yaml` commits the + hermes-agent guard; `09-pairing-security.md` M7 verification table (5 + items verified, 3 documented gaps: in-app QR scan, WSS cert pinning, + mechanical log redaction). Known: commentary scenario is model-dependent + (SKIP); FCM path needs a Firebase project to exercise; M6's formal + desktop parity pass + macOS/Windows packaging remain open. --- diff --git a/docs/README.md b/docs/README.md index f90ac1e..9173bec 100644 --- a/docs/README.md +++ b/docs/README.md @@ -62,6 +62,6 @@ project** (`app/`) with a shared KMP module (`app/shared`). ## Status -- **Phase:** Planning complete โ†’ ready to implement (Milestone M0). +- **Phase:** M0โ€“M6 complete; M7 (polish + E2E + docs) in progress. - **Owner decisions locked:** see [`16-open-questions.md`](16-open-questions.md). - **Last updated:** 2026-08-19. \ No newline at end of file diff --git a/docs/protocol/frames.schema.json b/docs/protocol/frames.schema.json index b231ff8..4c1c446 100644 --- a/docs/protocol/frames.schema.json +++ b/docs/protocol/frames.schema.json @@ -20,7 +20,7 @@ "hello.ack": { "description": "Pairing succeeded.", "payload": { - "server_caps": { "type": "object", "properties": { "streaming": {"type":"boolean"}, "reasoning": {"type":"boolean"}, "tools": {"type":"boolean"}, "media": {"type":"boolean"}, "search": {"type":"boolean"}, "push": {"type":"string","enum":["fcm","ntfy","none"]}, "pickers": {"type":"boolean"} } }, + "server_caps": { "type": "object", "properties": { "streaming": {"type":"boolean"}, "reasoning": {"type":"boolean"}, "tools": {"type":"boolean"}, "media": {"type":"boolean"}, "search": {"type":"boolean"}, "push": {"type":"string","enum":["fcm","ntfy","none"]}, "push_ntfy_server": {"type":"string","description":"ntfy server URL for the app's listener; empty string when the backend is not ntfy."}, "pickers": {"type":"boolean"} } }, "sync_cursor": { "type": "integer" }, "channels": { "type": "array", "items": { "$ref": "#/definitions/channel" } } } @@ -41,30 +41,21 @@ }, "message.start": { "payload": { "message_id": { "type": "string" }, "role": { "type": "string" } } }, "message.update": { "payload": { "message_id": { "type": "string" }, "text": { "type": "string", "description": "Full current text (app replaces)." } } }, - "message.stop": { "payload": { "message_id": { "type": "string" }, "final_text": { "type": "string" }, "reasoning": { "type": "string" }, "model": { "type": "string" }, "tokens": { "type": "integer" } } }, + "message.stop": { "payload": { "message_id": { "type": "string" }, "final_text": { "type": "string" }, "reasoning": { "type": "string" }, "model": { "type": "string" }, "tokens": { "type": "integer" }, "ts": { "type": "integer" } } }, "commentary": { "description": "Intermediate assistant beat.", "payload": { "message_id": { "type": "string" }, "text": { "type": "string" } } }, "tool.start": { "payload": { "index": { "type": "integer" }, "name": { "type": "string" }, "preview": { "type": "string" }, "args": { "type": "object" } } }, "tool.progress": { "payload": { "index": { "type": "integer" }, "name": { "type": "string" }, "note": { "type": "string" } } }, "tool.end": { "payload": { "index": { "type": "integer" }, "name": { "type": "string" }, "ok": { "type": "boolean" }, "duration": { "type": "number" }, "output_preview": { "type": "string" } } }, "typing": { "payload": { "on": { "type": "boolean" } } }, - "notification": { "payload": { "kind": { "type": "string", "enum": ["channel_renamed", "channel_created", "cron", "approval", "clarify", "generic"] }, "title": { "type": "string" }, "body": { "type": "string" }, "ts": { "type": "integer" } } }, - "picker.model": { "payload": { "picker_id": { "type": "string" }, "current_model": { "type": "string" }, "current_provider": { "type": "string" }, "providers": { "type": "array", "items": { "type": "object", "properties": { "id": {"type":"string"}, "label": {"type":"string"}, "models": { "type": "array", "items": { "type": "object", "properties": { "id": {"type":"string"}, "label": {"type":"string"} } } } } } } } }, - "picker.choice": { "payload": { "picker_id": { "type": "string" }, "title": { "type": "string" }, "choices": { "type": "array", "items": { "type": "object", "properties": { "value": {"type":"string"}, "label": {"type":"string"}, "is_current": {"type":"boolean"} } } } } }, - "picker.clarify": { "payload": { "picker_id": { "type": "string" }, "question": { "type": "string" }, "choices": { "type": "array", "items": { "type": "object" } } } }, - "picker.approval": { "payload": { "picker_id": { "type": "string" }, "command": { "type": "string" }, "description": { "type": "string" } } }, - "picker.confirm": { "payload": { "picker_id": { "type": "string" }, "title": { "type": "string" }, "message": { "type": "string" } } }, - "channel.list": { "payload": { "channels": { "type": "array", "items": { "$ref": "#/definitions/channel" } } } }, + "notification": { "payload": { "kind": { "type": "string", "enum": ["channel_renamed", "channel_created", "channel_deleted", "cron", "approval", "clarify", "generic"] }, "title": { "type": "string" }, "body": { "type": "string" }, "ts": { "type": "integer" } } }, + "channel.list": { "description": "Full channel directory (response to a channel.list request).", "payload": { "channels": { "type": "array", "items": { "$ref": "#/definitions/channel" } } } }, "channel.created": { "payload": { "$ref": "#/definitions/channel" } }, - "channel.renamed": { "payload": { "chat_id": { "type": "string" }, "name": { "type": "string" } } }, + "channel.renamed": { "description": "Also the response to channel.set_default (carries the full entry incl. is_default).", "payload": { "$ref": "#/definitions/channel" } }, "channel.deleted": { "payload": { "chat_id": { "type": "string" } } }, - "history": { "description": "Response to history request; page of messages oldestโ†’newest.", "payload": { "messages": { "type": "array", "items": { "type": "object", "properties": { "message_id": {"type":"string"}, "role": {"type":"string"}, "text": {"type":"string"}, "reasoning": {"type":"string"}, "model": {"type":"string"}, "tokens": {"type":"integer"}, "ts": {"type":"integer"} } } }, "has_more": { "type": "boolean" }, "oldest_message_id": { "type": "string" } } }, - "commands.catalog": { "description": "Full slash-command catalog.", "payload": { "commands": { "type": "array", "items": { "type": "object", "properties": { "name": {"type":"string"}, "description": {"type":"string"}, "args_hint": {"type":"string"}, "category": {"type":"string"} } } } } }, - "commands.complete": { "description": "Autocomplete matches for a typed prefix.", "payload": { "prefix": { "type": "string" }, "matches": { "type": "array", "items": { "type": "object", "properties": { "name": {"type":"string"}, "description": {"type":"string"}, "args_hint": {"type":"string"} } } } } }, - "agent.busy": { "description": "Agent is processing; app shows thinking indicator.", "payload": { "reason": { "type": "string", "enum": ["processing", "tool", "waiting_input", "cron"] } } }, - "agent.idle": { "description": "Agent turn complete; clear thinking indicator.", "payload": {} }, "search.results": { "payload": { "query": { "type": "string" }, "scope": { "type": "string", "enum": ["all", "chat"] }, "hits": { "type": "array", "items": { "type": "object", "properties": { "message_id": {"type":"string"}, "chat_id": {"type":"string"}, "thread_id": {"type":["string","null"]}, "role": {"type":"string"}, "snippet": {"type":"string"}, "ts": {"type":"integer"} } } } } }, "media.offer": { "description": "Agent-sent media available; app pulls bytes.", "payload": { "$ref": "#/definitions/media_ref" } }, - "status": { "payload": { "state": { "type": "string", "enum": ["online", "restarting", "degraded"] }, "session": { "type": "object" } } }, + "read.receipt": { "description": "Agent received and started processing the user's message; app shows โœ“โœ“ on user bubbles. Emitted to the originating connection when a message.send is accepted for processing.", "payload": { "message_id": { "type": "string" } } }, + "status": { "description": "Gateway health state; broadcast to all connected clients at startup (state=online).", "payload": { "state": { "type": "string", "enum": ["online", "restarting", "degraded"] } } }, "error": { "payload": { "code": { "type": "string", "enum": ["auth", "not_found", "rate_limited", "media_too_large", "unsupported", "internal"] }, "message": { "type": "string" } } }, "pong": { "payload": { "ts": { "type": "integer" } } }, "sync.done": { "payload": { "cursor": { "type": "integer" } } }, @@ -77,18 +68,12 @@ "media.upload.start": { "payload": { "media_ref": { "type": "string" }, "kind": { "$ref": "#/definitions/kind" }, "mime": { "type": "string" }, "size": { "type": "integer" }, "filename": { "type": "string" } } }, "media.upload.end": { "payload": { "media_ref": { "type": "string" }, "sha256": { "type": "string" } } }, "media.pull": { "payload": { "media_id": { "type": "string" } } }, - "picker.select": { "payload": { "picker_id": { "type": "string" }, "value": { "type": "string" } } }, "channel.create": { "payload": { "name": { "type": "string" }, "kind": { "type": "string", "enum": ["channel", "thread"] }, "parent_chat_id": { "type": ["string", "null"] } } }, "channel.rename": { "payload": { "name": { "type": "string" } } }, "channel.set_default": { "payload": {} }, "channel.delete": { "payload": {} }, - "search": { "payload": { "query": { "type": "string" }, "scope": { "type": "string", "enum": ["all", "chat"] }, "chat_id": { "type": "string" }, "thread_id": { "type": "string" } } }, - "history": { "description": "Load a page of messages (initial open / scroll-up).", "payload": { "before_message_id": { "type": "string" }, "limit": { "type": "integer", "minimum": 1, "maximum": 200 } } }, - "commands.catalog": { "description": "Fetch full slash-command catalog.", "payload": {} }, - "commands.complete": { "description": "Autocomplete for typed /prefix.", "payload": { "prefix": { "type": "string" } } }, - "agent.stop": { "description": "Abort current agent turn.", "payload": {} }, - "agent.steer": { "description": "Inject steering message mid-turn.", "payload": { "text": { "type": "string" } } }, - "read.receipt": { "description": "User viewed message; server stores + broadcasts to other devices.", "payload": { "chat_id": { "type": "string" }, "message_id": { "type": "string" } } }, + "channel.list": { "description": "Request the full channel directory; answered by the server_to_app channel.list frame.", "payload": {} }, + "search": { "payload": { "query": { "type": "string" }, "scope": { "type": "string", "enum": ["all", "chat"] }, "chat_id": { "type": "string" }, "thread_id": { "type": "string" }, "limit": { "type": "integer", "description": "Optional; server default 20." } } }, "sync": { "description": "Reconnect catch-up; replays undelivered outbox frames only (not full history).", "payload": { "cursor": { "type": "integer" } } }, "fcm.register": { "payload": { "fcm_token": { "type": "string" }, "ntfy_topic": { "type": "string" } } }, "ping": { "payload": { "ts": { "type": "integer" } } } @@ -96,12 +81,31 @@ }, "definitions": { "kind": { "type": "string", "enum": ["image", "audio", "video", "document", "voice"] }, - "channel": { "type": "object", "properties": { "chat_id": {"type":"string"}, "name": {"type":"string"}, "kind": {"type":"string","enum":["default","channel","thread"]}, "parent_chat_id": {"type":["string","null"]}, "is_default": {"type":"boolean"} } }, - "media_ref": { "type": "object", "properties": { "media_id": {"type":"string"}, "kind": { "$ref": "#/definitions/kind" }, "mime": {"type":"string"}, "size": {"type":"integer"}, "filename": {"type":"string"} } } + "channel": { "type": "object", "properties": { "chat_id": {"type":"string"}, "name": {"type":"string"}, "kind": {"type":"string","enum":["default","channel","thread"]}, "parent_chat_id": {"type":["string","null"]}, "is_default": {"type":"boolean"}, "archived": {"type":"boolean"} } }, + "media_ref": { "type": "object", "properties": { "media_id": {"type":"string"}, "kind": { "$ref": "#/definitions/kind" }, "mime": {"type":"string"}, "size": {"type":"integer"}, "filename": {"type":"string"}, "message_id": {"type":"string","description":"Optional; set on media.offer to associate the offer with the assistant message it belongs to."} } } }, + "x-planned-frames": [ + { "name": "picker.model", "direction": "server_to_app", "note": "Model/provider picker prompt. Planned, not implemented." }, + { "name": "picker.choice", "direction": "server_to_app", "note": "Generic choice picker prompt. Planned, not implemented." }, + { "name": "picker.clarify", "direction": "server_to_app", "note": "Clarify picker prompt. Planned, not implemented (clarifies arrive as notification + message)." }, + { "name": "picker.approval", "direction": "server_to_app", "note": "Approval picker prompt. Planned, not implemented (approvals arrive as notification)." }, + { "name": "picker.confirm", "direction": "server_to_app", "note": "Confirmation picker prompt. Planned, not implemented." }, + { "name": "picker.select", "direction": "app_to_server", "note": "Picker answer. Planned, not implemented." }, + { "name": "history", "direction": "server_to_app", "note": "Paged history response. Planned, not implemented (catch-up is sync/outbox replay)." }, + { "name": "history", "direction": "app_to_server", "note": "Paged history request. Planned, not implemented (catch-up is sync/outbox replay)." }, + { "name": "commands.catalog", "direction": "server_to_app", "note": "Slash-command catalog. Planned, not implemented." }, + { "name": "commands.catalog", "direction": "app_to_server", "note": "Slash-command catalog request. Planned, not implemented." }, + { "name": "commands.complete", "direction": "server_to_app", "note": "Slash-command autocomplete. Planned, not implemented." }, + { "name": "commands.complete", "direction": "app_to_server", "note": "Slash-command autocomplete request. Planned, not implemented." }, + { "name": "agent.busy", "direction": "server_to_app", "note": "Agent-busy indicator. Planned, not implemented (typing frames cover it)." }, + { "name": "agent.idle", "direction": "server_to_app", "note": "Agent-idle indicator. Planned, not implemented." }, + { "name": "agent.stop", "direction": "app_to_server", "note": "Abort current agent turn. Planned, not implemented." }, + { "name": "agent.steer", "direction": "app_to_server", "note": "Steer the agent mid-turn. Planned, not implemented." }, + { "name": "read.receipt", "direction": "app_to_server", "note": "User-viewed receipt (multi-device read state). Planned, not implemented (read.receipt is serverโ†’app only)." } + ], "reliability": { "ordering": "Per-connection (TCP/WS). message.update for a message_id is monotonic; app may coalesce to latest.", - "never_dropped": ["message", "message.stop", "tool.end", "notification", "picker.*", "channel.*", "agent.busy", "agent.idle", "history", "commands.catalog", "commands.complete", "search.results", "error"], + "never_dropped": ["message", "message.stop", "tool.end", "notification", "channel.*", "search.results", "error"], "coalescable_under_backpressure": ["message.update", "tool.progress"], "offline": "Undelivered frames go to the outbox; replayed by sync. Terminal frames always outboxed." } diff --git a/docs/setup.md b/docs/setup.md new file mode 100644 index 0000000..204df37 --- /dev/null +++ b/docs/setup.md @@ -0,0 +1,170 @@ +# Setup โ€” Pairing a Device + +User-facing guide: get a phone or desktop talking to your hermes gateway in +under 10 minutes. Design rationale lives in the numbered docs +([`09-pairing-security.md`](09-pairing-security.md), +[`08-push.md`](08-push.md), [`12-toolchain.md`](12-toolchain.md)); this page is +just the steps. + +## Prerequisites + +| Where | You need | +|---|---| +| Gateway host | hermes installed with its venv (`cd hermes-agent && uv sync`, see [`12-toolchain.md` ยง12.4](12-toolchain.md)) | +| Android build machine | JDK 17, Android SDK with `ANDROID_HOME` set (or `app/local.properties`), ADB with a connected device | +| Desktop build machine | JDK 17 only | + +Gradle needs no system install โ€” both apps use the project wrapper +(`./gradlew`). First-time machine setup: [`12-toolchain.md`](12-toolchain.md). + +## 1. Gateway setup (on the gateway host) + +Install the plugin into the live hermes home (dev: a symlink from the monorepo +root): + +```bash +mkdir -p ~/.hermes/plugins +ln -s "$PWD/gateway-plugin" ~/.hermes/plugins/android +hermes gateway status # should list "android" +``` + +Run the interactive setup: + +```bash +hermes gateway setup +``` + +What it does: + +- Generates `ANDROID_TOKEN` (64 hex chars) if none exists and stores it in + `~/.hermes/.env` (it prints the token once, at generation). +- Prompts for the WS bind host (default `127.0.0.1`), port (default `8790`), + and push backend (`fcm` or `ntfy`, default `fcm`). +- Prints the pairing payload (a QR-encodable `iris://pair?host=โ€ฆ&port=โ€ฆ&token=โ€ฆ` + string) and the server URL (`ws://:8790/ws`). + +Then start the gateway: + +```bash +hermes gateway # or: hermes gateway restart after config changes +``` + +> **Note:** the default bind host `127.0.0.1` only accepts connections from the +> gateway host itself (e.g. a desktop app on the same machine). For a phone on +> the LAN, re-run `hermes gateway setup` (or edit `~/.hermes/.env`) and set +> `ANDROID_WS_HOST` to the host's LAN IP (e.g. `192.168.1.10`). + +## 2. Android app + +Build and install (ADB device connected): + +```bash +cd app +./gradlew :androidApp:installDebug +``` + +First run opens the **Connect** screen: + +1. **Server URL** โ€” `ws://:8790/ws` (the URL printed by + `hermes gateway setup`; use the LAN IP, not `127.0.0.1`, from a phone). +2. **Pairing token** โ€” from the `hermes gateway setup` output, or + `grep ANDROID_TOKEN ~/.hermes/.env` on the gateway host. +3. **Test & Connect** โ€” performs a real `hello` (the auth leg), then saves the + pairing and connects. + +> **Honest limitation:** QR scanning is **not** supported in the app yet. The +> server prints a QR payload, but pairing is manual URL + token entry only. + +## 3. Desktop app + +```bash +cd app +./gradlew :desktopApp:run # dev run +./gradlew :desktopApp:jpackage # native app-image (bundles the JRE) +``` + +Pairing is the same Connect screen (URL + token); the token is stored in the OS +keyring (with an encrypted-file fallback). Desktop push is tray icon + OS +notifications (no FCM). + +> **Known issue:** on Linux with JDK 17 the jpackage launcher prints a +> non-fatal `pure virtual method called` warning (JDK-8348560, a +> jpackage/Linux launcher bug). The app runs and connects regardless. + +## 4. Push notifications + +Push wakes a backgrounded/offline device; on reconnect the app syncs the +outbox, so nothing is lost. Push fires when the device is offline, plus for +high-priority events (approvals, clarifies, cron) even when a device is live. + +### FCM (default; needs a Firebase project) + +1. Create a Firebase project (console.firebase.google.com) and add an Android + app with the app's applicationId; download `google-services.json` into + `app/androidApp/`. +2. Create a service account (Project settings โ†’ Service accounts โ†’ Generate new + private key) and store the JSON path in `~/.hermes/.env`: + `ANDROID_FCM_SERVICE_ACCOUNT=/path/to/service-account.json`. +3. Keep `ANDROID_PUSH_BACKEND=fcm` (the default). + +Without a Firebase project the FCM path is **inert** (the app's FCM service +does nothing) โ€” use ntfy below, or add Firebase later. + +**What you see:** system notifications for new messages when the app is +backgrounded; tapping one deep-links to the chat. + +### ntfy (zero-config fallback) + +``` +ANDROID_PUSH_BACKEND=ntfy +``` + +- The device **generates its own topic** automatically (no `NTFY_TOPIC` needed); + the server publishes to it. +- `NTFY_SERVER_URL` defaults to `https://ntfy.sh`. **Self-hosted ntfy is + recommended** โ€” the public ntfy.sh SSE endpoint is flaky (it has served its + web UI instead of the stream), while a self-hosted instance gives reliable + SSE. For a real trust boundary use a private topic + `NTFY_AUTH_TOKEN`. + +**What you see:** a low-priority foreground "ntfy listener" notification while +the app is off; incoming pushes trigger a silent sync. + +## 5. Remote access + +- **Tailscale / WireGuard (recommended):** the gateway gets a stable tailnet IP; + the app connects to `ws://:8790/ws`. No public exposure. +- **Reverse proxy / tunnel** (Caddy, Cloudflare Tunnel, ngrok): terminate TLS at + the edge, forward the WebSocket to `127.0.0.1:8790`. +- **WSS:** set `ANDROID_WS_CERT` / `ANDROID_WS_KEY` (paths, in + `~/.hermes/.env`) and the server serves `wss://` instead of `ws://`. + +> **Honest limitation:** the app has **no certificate pinning** yet, so +> self-signed certs won't work โ€” remote access requires **CA-signed** WSS for +> now. Plain `ws://` on a trusted LAN (or inside Tailscale) stays the default. + +## 6. Troubleshooting + +| Symptom | Likely cause / fix | +|---|---| +| `auth failed` / `error {code:"auth"}` on connect | Wrong token. Check `ANDROID_TOKEN` in `~/.hermes/.env` on the gateway host (setup prints it only when it generates it). | +| Connection refused | Gateway not running (`hermes gateway status`); wrong URL (port `8790`, path `/ws`, LAN IP instead of `127.0.0.1` from a phone); firewall blocking the port. | +| Push not arriving | Backend not configured (gateway log: `push backend โ€ฆ not configured`); app backgrounded with no working backend; ntfy.sh SSE flakiness โ€” use a self-hosted ntfy. | +| Desktop jpackage launcher warning | Non-fatal (JDK-8348560 on Linux JDK 17); the app runs and connects regardless. | + +Smoke test without the app (from the gateway host): + +```bash +python - <<'PY' +import asyncio, json, websockets +async def main(): + async with websockets.connect("ws://127.0.0.1:8790/ws") as ws: + await ws.send(json.dumps({"v":1,"type":"hello","payload":{ + "token":"","device_id":"test","device_name":"probe", + "caps":{"min_protocol":1}}})) + print("recv:", await ws.recv()) +asyncio.run(main()) +PY +``` + +Expect a `hello.ack`. If you get `error {code:"auth"}`, the token/host/port is +wrong. \ No newline at end of file diff --git a/gateway-plugin/adapter.py b/gateway-plugin/adapter.py index fd29c36..68e077c 100644 --- a/gateway-plugin/adapter.py +++ b/gateway-plugin/adapter.py @@ -691,6 +691,7 @@ class AndroidAdapter(BasePlatformAdapter): self.max_upload_bytes = int( extra.get("max_upload_bytes", DEFAULT_MAX_UPLOAD_BYTES) ) + self._gateway_status = protocol.STATUS_ONLINE # Home channel: the core hook turns the env-seeded ``home_channel`` # dict into a HomeChannel dataclass on the config; config.yaml may @@ -796,6 +797,10 @@ class AndroidAdapter(BasePlatformAdapter): self._connected = False return False + # M5: announce gateway health to connected clients (none yet at + # startup; the frame + plumbing exist for future transitions). + await self._ws_server.broadcast(protocol.status(self._gateway_status)) + # M3: ensure the default (home) channel exists in the directory so the # app's channel list and cron home delivery have a stable anchor. try: @@ -1542,6 +1547,12 @@ class AndroidAdapter(BasePlatformAdapter): media_types=media_types, ) await self.handle_message(event) + # M5: acknowledge the user message to the originating device (the + # app shows โœ“โœ“) at the moment it is handed to the agent. + await self._ws_server.send_to( + device_id, + protocol.read_receipt(chat_id, message_id), + ) # โ”€โ”€ M4: inbound media (app -> agent) โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€ # @@ -2023,6 +2034,10 @@ class AndroidAdapter(BasePlatformAdapter): # โ”€โ”€ hello.ack helpers โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€ + def gateway_status(self) -> str: + """Current gateway health state (sent to each pairing connection).""" + return self._gateway_status + def server_caps(self) -> Dict[str, Any]: """Capability flags advertised in ``hello.ack`` (M4 surface).""" return { diff --git a/gateway-plugin/protocol.py b/gateway-plugin/protocol.py index 868a186..1f94045 100644 --- a/gateway-plugin/protocol.py +++ b/gateway-plugin/protocol.py @@ -12,7 +12,7 @@ Milestone M2: message.start/update/stop, reasoning (on message / message.stop), tool.start/progress/end, commentary. Milestone M3: channel.*, search, sync. Milestone M4: media.*. -Milestone M5: notification, fcm.register, read.receipt. +Milestone M5: notification, fcm.register, read.receipt, status. """ import json @@ -79,6 +79,10 @@ TYPE_MEDIA_PULL_END = "media.pull.end" # Push / notifications (M5) TYPE_NOTIFICATION = "notification" TYPE_FCM_REGISTER = "fcm.register" +TYPE_READ_RECEIPT = "read.receipt" + +# Gateway health (M5) +TYPE_STATUS = "status" # --------------------------------------------------------------------------- # Error codes (``error`` frame payload.code) @@ -116,6 +120,14 @@ NOTIF_GENERIC = "generic" # it decides whether to also show an in-app banner). HIGH_PRIORITY_NOTIF_KINDS = frozenset({NOTIF_APPROVAL, NOTIF_CLARIFY, NOTIF_CRON}) +# --------------------------------------------------------------------------- +# Gateway health states (``status`` frame payload.state) +# --------------------------------------------------------------------------- + +STATUS_ONLINE = "online" +STATUS_RESTARTING = "restarting" +STATUS_DEGRADED = "degraded" + # --------------------------------------------------------------------------- # Envelope @@ -507,6 +519,25 @@ def fcm_register( return Frame(type=TYPE_FCM_REGISTER, payload=payload) +def read_receipt(chat_id: str, message_id: str) -> Frame: + """Ack to the originating device: the agent received and started + processing the user's message (the app shows โœ“โœ“ on the user bubble).""" + return Frame( + type=TYPE_READ_RECEIPT, + chat_id=chat_id, + payload={"message_id": message_id}, + ) + + +# --------------------------------------------------------------------------- +# Gateway health frame (M5) +# --------------------------------------------------------------------------- + +def status(state: str) -> Frame: + """Gateway health state (``state`` is one of the ``STATUS_*`` constants).""" + return Frame(type=TYPE_STATUS, payload={"state": state}) + + # --------------------------------------------------------------------------- # Media frames (M4) # --------------------------------------------------------------------------- diff --git a/gateway-plugin/tests/README.md b/gateway-plugin/tests/README.md index 6a26aed..ed98366 100644 --- a/gateway-plugin/tests/README.md +++ b/gateway-plugin/tests/README.md @@ -4,4 +4,73 @@ Run via hermes's hermetic runner (never bare pytest):: scripts/run_tests.sh tests/gateway/test_android.py -See ``docs/13-testing.md`` for the scenario list. \ No newline at end of file +See ``docs/13-testing.md`` for the scenario list. + +## WS probe (`ws_probe.py`) + +Manual test-client harness: connects to the **real running gateway** and +drives a turn, printing every frame. Run with the hermes venv python +(needs `websockets`); the gateway must already be up:: + + hermes-agent/.venv/bin/python gateway-plugin/tests/ws_probe.py \ + --token --send "hello" + +Beyond the base modes (`--send`, `--upload`, `--pull-offer`, `--sync`, +`--fcm-token`/`--fcm-reg`, `--authfail`, `--url`, `--token`, `--device`, +`--timeout`), the probe has assertion and request modes: + +- `--assert-turn` โ€” assert the turn produced `message.start` โ†’ โ‰ฅ1 + `message.update` โ†’ `message.stop` (scenario 2). +- `--assert-reasoning` โ€” assert the final `message.stop` carries a + non-empty `reasoning` field (scenario 3). +- `--assert-tools` โ€” assert โ‰ฅ1 `tool.start` with a matching `tool.end` + (matched by `index`; scenario 4). +- `--assert-commentary` โ€” assert โ‰ฅ1 `commentary` frame (scenario 5). +- `--assert-read-receipt` โ€” assert a `read.receipt` frame arrives after + the sent message (new M7 frame; requires `--send`). **SKIPs** (exit 0, + prints `== SKIP: โ€ฆ`) when the frame never arrives, e.g. against a + gateway that predates the M7 frames. +- `--assert-status` โ€” assert a `status` frame is received (new M7 frame; + **SKIPs** when absent). +- `--search QUERY [--scope all|chat] [--chat-id C]` โ€” send a `search` + frame (`{query, scope, limit}`) and assert โ‰ฅ1 hit in `search.results` + (scenario 8). With `--send`, the turn is driven first, then the search. +- `--channel-create NAME` / `--channel-delete CHAT_ID` / + `--channel-list` โ€” M3 channel directory management; create prints + `== channel created: ` for scripting. +- `--watch CHAT_ID` โ€” wait up to `--timeout` for a message to land in + `CHAT_ID` (cron delivery E2E, scenario 7). +- `--offer-grace S` โ€” with `--pull-offer`, keep listening S seconds after + the final message for a `media.offer` (offers are emitted post-turn, + right after the final; default 15). + +Exit codes: `0` ok (incl. SKIP for absent M7 frames), `2` connect fail, +`3` no hello.ack, `4` expected hello.ack, `5` authfail expected but +acked, `6` timeout, `7` no final message, `8` upload/sync fail, +`9` pull fail, `10` assert-turn fail, `11` assert-reasoning fail, +`12` assert-tools fail, `13` assert-commentary fail, `14` search fail +(error or zero hits), `15` channel.create/list fail, `16` channel.delete +fail, `17` watch timeout, `18` read.receipt arrived before the sent +message, `19` status frame with empty payload. + +## E2E driver (`e2e.py`) + +Runs the `docs/13-testing.md` ยง13.4 scenarios 1โ€“12 automated-where- +possible against the live gateway, invoking `ws_probe.py` (and the +`hermes` CLI for cron) as subprocesses. Prints PASS / PARTIAL / SKIP / +FAIL per scenario plus a summary table; exits 0 if no FAIL, 1 otherwise:: + + hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py + hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py --skip 3,5,7 + hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py --url ws://host:8790/ws + +The token is read from `$ANDROID_TOKEN`, else `hermes-agent/.env`, else +`~/.hermes/.env`. The gateway must already be running (the driver never +starts or stops it). It is idempotent: channels/jobs it creates are +cleaned up even on failure, and leftover `e2e-*` channels/jobs from +earlier runs are removed at start. + +Scenario notes: 3 (reasoning) and 5 (commentary) are model-dependent and +SKIP rather than FAIL when the current model does not emit them; 11 +(push) and 12 (reconnect/sync) are PARTIAL by design โ€” the WS leg is +automated, the device-notification / gateway-kill leg is manual. \ No newline at end of file diff --git a/gateway-plugin/tests/e2e.py b/gateway-plugin/tests/e2e.py new file mode 100644 index 0000000..7401d22 --- /dev/null +++ b/gateway-plugin/tests/e2e.py @@ -0,0 +1,373 @@ +#!/usr/bin/env python3 +"""E2E driver: docs/13-testing.md ยง13.4 scenarios 1-12 against the live gateway. + +Drives ws_probe.py (and the hermes CLI for cron) as subprocesses. For each +scenario prints PASS / PARTIAL / SKIP / FAIL with a one-line reason, then a +summary table. Exit 0 if no FAIL, 1 otherwise. + +Usage:: + + hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py + hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py --skip 3,5,7 + hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py --url ws://host:8790/ws + +The token is read from $ANDROID_TOKEN, else hermes-agent/.env, else +~/.hermes/.env. The gateway must already be running (this driver never +starts or stops it). Idempotent: channels/jobs it creates are cleaned up +even on failure, and leftover "e2e-*" channels/jobs from earlier runs are +removed at start. + +Scenario notes: + 3 (reasoning) and 5 (commentary) are model-dependent: they SKIP (not + FAIL) when the current model does not emit reasoning / commentary. + 11 (push) and 12 (reconnect) are PARTIAL by design: the WS leg is + automated, the device-notification / gateway-kill leg is manual. +""" + +import argparse +import os +import re +import struct +import subprocess +import sys +import uuid +import zlib +from pathlib import Path + +HERE = Path(__file__).resolve().parent +REPO = HERE.parent.parent +PY = REPO / "hermes-agent" / ".venv" / "bin" / "python" +PROBE = HERE / "ws_probe.py" +HERMES = REPO / "hermes-agent" / ".venv" / "bin" / "hermes" +DEFAULT_URL = "ws://127.0.0.1:8790/ws" + +PASS, PARTIAL, SKIP, FAIL = "PASS", "PARTIAL", "SKIP", "FAIL" + + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + +def find_token(cli_token: str) -> str: + if cli_token: + return cli_token + env = os.getenv("ANDROID_TOKEN") + if env: + return env + for p in (REPO / "hermes-agent" / ".env", Path.home() / ".hermes" / ".env"): + try: + for line in p.read_text().splitlines(): + line = line.strip() + if line.startswith("ANDROID_TOKEN="): + return line.split("=", 1)[1].strip().strip('"').strip("'") + except OSError: + pass + return "" + + +def run_probe(env, url, token, *args, timeout=300): + cmd = [str(PY), str(PROBE), "--url", url, "--token", token, *args] + p = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout, env=env) + return p.returncode, p.stdout, p.stderr + + +def run_hermes(env, *args, timeout=120): + cmd = [str(HERMES), *args] + p = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout, env=env) + return p.returncode, p.stdout, p.stderr + + +def write_png(path: Path, color, size: int = 200) -> None: + """Write a solid-color RGB PNG using only the stdlib (no PIL needed).""" + raw = b"".join(b"\x00" + bytes(color) * size for _ in range(size)) + + def chunk(tag: bytes, data: bytes) -> bytes: + return (struct.pack(">I", len(data)) + tag + data + + struct.pack(">I", zlib.crc32(tag + data) & 0xFFFFFFFF)) + + ihdr = struct.pack(">IIBBBBB", size, size, 8, 2, 0, 0, 0) + path.write_bytes( + b"\x89PNG\r\n\x1a\n" + + chunk(b"IHDR", ihdr) + + chunk(b"IDAT", zlib.compress(raw)) + + chunk(b"IEND", b"") + ) + + +def parse_created_chat_id(out: str) -> str | None: + m = re.search(r"== channel created: (\S+)", out) + return m.group(1) if m else None + + +def sweep_leftovers(env, url, token) -> None: + """Remove e2e-* channels / cron jobs left behind by earlier runs.""" + rc, out, _ = run_probe(env, url, token, "--channel-list") + if rc == 0: + for m in re.finditer(r"== channel: (\S+) name='(e2e-[^']*)'", out): + chat_id, name = m.group(1), m.group(2) + print(f" cleanup: removing leftover channel {chat_id} ({name})") + run_probe(env, url, token, "--channel-delete", chat_id) + rc, out, _ = run_hermes(env, "cron", "list") + if rc == 0: + for m in re.finditer( + r"(\S+) \[(?:active|paused)\]\s*\n\s*Name:\s+(e2e-cron-[^ \n]*)", out + ): + job_id, name = m.group(1), m.group(2) + print(f" cleanup: removing leftover cron job {job_id} ({name})") + run_hermes(env, "cron", "remove", job_id) + + +# --------------------------------------------------------------------------- +# Scenarios (docs/13-testing.md ยง13.4) +# --------------------------------------------------------------------------- + +def s1_pair(env, url, token): + rc, _, _ = run_probe(env, url, "definitely-wrong-token", "--authfail", "--send", "") + if rc != 0: + return FAIL, f"wrong token was not rejected (rc={rc})" + rc, _, _ = run_probe(env, url, token, "--send", "") + if rc != 0: + return FAIL, f"valid token did not pair (rc={rc})" + return PASS, "wrong token rejected; hello.ack on valid token" + + +def s2_text(env, url, token): + prompt = "Write a short poem about the ocean, at least 8 lines" + rc, _, _ = run_probe(env, url, token, "--send", prompt, + "--assert-turn", "--timeout", "120") + if rc == 0: + return PASS, "message.start -> >=1 message.update -> message.stop" + if rc == 10: + return FAIL, "no ordered start/update/stop segment" + return FAIL, f"probe rc={rc}" + + +def s3_reasoning(env, url, token): + prompt = "Work out step by step: what is 17 * 23? Show your reasoning." + rc, _, _ = run_probe(env, url, token, "--send", prompt, + "--assert-reasoning", "--timeout", "120") + if rc == 0: + return PASS, "final message.stop carries non-empty reasoning" + if rc == 11: + return SKIP, "model returned no reasoning (model-dependent)" + return FAIL, f"probe rc={rc}" + + +def s4_tools(env, url, token): + prompt = ("List the files in your current working directory using your " + "shell tool, then tell me how many there are") + rc, _, _ = run_probe(env, url, token, "--send", prompt, + "--assert-tools", "--timeout", "150") + if rc == 0: + return PASS, "tool.start with a matching tool.end" + if rc == 12: + return FAIL, "no tool.start/tool.end pair" + return FAIL, f"probe rc={rc}" + + +def s5_commentary(env, url, token): + prompt = ("Research task: (1) use your shell tool to list the top-level " + "directories in /tmp, (2) report your findings so far, " + "(3) use your shell tool to count files in /tmp, " + "(4) report those findings too, (5) give a final summary of both") + rc, _, _ = run_probe(env, url, token, "--send", prompt, + "--assert-commentary", "--timeout", "150") + if rc == 0: + return PASS, "commentary frame observed" + if rc == 13: + return SKIP, "no commentary (model/agent-dependent per M2)" + return FAIL, f"probe rc={rc}" + + +def s6_channels(env, url, token): + name = f"e2e-chan-{uuid.uuid4().hex[:6]}" + rc, out, _ = run_probe(env, url, token, "--channel-create", name) + if rc != 0: + return FAIL, f"channel.create failed (rc={rc})" + chat_id = parse_created_chat_id(out) + if not chat_id: + return FAIL, "channel.created received but chat_id not parseable" + rc, _, _ = run_probe(env, url, token, "--channel-delete", chat_id) + if rc != 0: + run_probe(env, url, token, "--channel-delete", chat_id) # best-effort + return FAIL, f"channel.delete failed (rc={rc})" + return PASS, f"created {chat_id} + deleted (cleanup)" + + +def s7_cron(env, url, token): + chan_name = f"e2e-cron-chan-{uuid.uuid4().hex[:6]}" + rc, out, _ = run_probe(env, url, token, "--channel-create", chan_name) + if rc != 0: + return SKIP, f"could not create cron target channel (rc={rc})" + chat_id = parse_created_chat_id(out) + if not chat_id: + return FAIL, "channel.created received but chat_id not parseable" + job_name = f"e2e-cron-{uuid.uuid4().hex[:6]}" + deliver = f"android:{chat_id}" + rc, out, err = run_hermes( + env, "cron", "create", "1m", + "Reply with exactly: e2e cron delivery OK", + "--deliver", deliver, "--name", job_name, + ) + job_id = None + if rc == 0: + m = re.search(r"Created job: (\S+)", out) + job_id = m.group(1) if m else None + try: + if rc != 0: + return SKIP, f"hermes cron create failed: {(err or out).strip()[:120]}" + rc, out, _ = run_probe(env, url, token, "--watch", chat_id, + "--timeout", "330", timeout=400) + if rc == 0: + return PASS, f"one-shot cron job fired; message landed in {chat_id}" + return FAIL, f"no message in {chat_id} within 330s (probe rc={rc})" + finally: + if job_id: + run_hermes(env, "cron", "remove", job_id) + else: + # create succeeded but the id was not parseable: find by name. + _, list_out, _ = run_hermes(env, "cron", "list") + m = re.search(r"(\S+) \[active\]\s*\n\s*Name:\s+" + re.escape(job_name), + list_out) + if m: + run_hermes(env, "cron", "remove", m.group(1)) + run_probe(env, url, token, "--channel-delete", chat_id) + + +def s8_search(env, url, token): + marker = f"e2emarker{uuid.uuid4().hex[:8]}" + rc, _, _ = run_probe(env, url, token, "--send", + f"Remember this marker phrase: {marker}. " + "Just acknowledge it briefly.", + "--timeout", "120") + if rc != 0: + return FAIL, f"setup message failed (rc={rc})" + rc, _, _ = run_probe(env, url, token, "--send", "", "--search", marker) + if rc == 0: + return PASS, f"search for {marker!r} returned >=1 hit" + if rc == 14: + return FAIL, f"search for {marker!r} returned 0 hits" + return FAIL, f"probe rc={rc}" + + +def s9_media_in(env, url, token): + png = Path(f"/tmp/e2e_in_{uuid.uuid4().hex[:6]}.png") + write_png(png, (30, 120, 220)) + try: + rc, _, _ = run_probe(env, url, token, "--upload", str(png), + "--send", "describe this image briefly", + "--timeout", "120") + if rc == 0: + return PASS, "upload + vision reply (final message)" + if rc == 8: + return FAIL, "media upload failed" + if rc == 7: + return FAIL, "no final message after upload" + return FAIL, f"probe rc={rc}" + finally: + png.unlink(missing_ok=True) + + +def s10_media_out(env, url, token): + prompt = ("Create a 100x100 orange square PNG in /tmp with your tools. " + "In your final reply, include the MEDIA:/absolute/path tag for " + "that file so it is delivered to me.") + rc, out, _ = run_probe(env, url, token, "--send", prompt, + "--pull-offer", "--timeout", "150") + m = re.search(r"== pulled (\d+) bytes", out) + if rc == 0 and m and int(m.group(1)) > 0: + return PASS, f"media.offer pulled ({m.group(1)} bytes)" + if rc == 9: + return FAIL, "media pull failed" + return FAIL, "no media.offer pulled (agent did not deliver an image)" + + +def s11_push(env, url, token): + rc, out, _ = run_probe(env, url, token, "--fcm-token", "test-token-123", + "--fcm-reg", "--send", "") + if rc != 0: + return FAIL, f"probe rc={rc}" + if "<- error" in out: + return FAIL, "error frame after fcm.register" + return PARTIAL, ("fcm.register accepted (no error frame); " + "device-notification leg is manual") + + +def s12_sync(env, url, token): + rc, out, _ = run_probe(env, url, token, "--sync", "0") + if rc == 0 and "sync done" in out: + return PARTIAL, ("sync replay + sync.done verified; " + "gateway-kill/restart leg is manual") + if rc == 8: + return FAIL, "sync failed" + return FAIL, f"probe rc={rc}" + + +SCENARIOS = [ + (1, "pair", s1_pair), + (2, "text round-trip", s2_text), + (3, "reasoning", s3_reasoning), + (4, "tools", s4_tools), + (5, "commentary", s5_commentary), + (6, "channels", s6_channels), + (7, "cron delivery", s7_cron), + (8, "search", s8_search), + (9, "media in", s9_media_in), + (10, "media out", s10_media_out), + (11, "push", s11_push), + (12, "reconnect/sync", s12_sync), +] + + +def main() -> int: + p = argparse.ArgumentParser( + description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter + ) + p.add_argument("--url", default=os.getenv("ANDROID_WS_URL", DEFAULT_URL)) + p.add_argument("--token", default="") + p.add_argument("--skip", default="", + help="comma-separated scenario numbers to skip (e.g. 3,5,7)") + args = p.parse_args() + + token = find_token(args.token) + if not token: + print("!! ANDROID_TOKEN not found (env, hermes-agent/.env, or ~/.hermes/.env)") + return 1 + skip = {int(x) for x in args.skip.split(",") if x.strip()} + + env = dict(os.environ) + env["ANDROID_TOKEN"] = token + + print(f"== e2e: url={args.url} token={token[:6]}โ€ฆ") + sweep_leftovers(env, args.url, token) + + results = [] + for num, name, fn in SCENARIOS: + if num in skip: + results.append((num, name, SKIP, "skipped by --skip")) + print(f"[{num:2d}] {name:<18} {SKIP:<7} skipped by --skip") + continue + print(f"[{num:2d}] {name:<18} runningโ€ฆ", flush=True) + try: + status, reason = fn(env, args.url, token) + except Exception as e: + status, reason = FAIL, f"driver error: {e}" + results.append((num, name, status, reason)) + print(f"[{num:2d}] {name:<18} {status:<7} {reason}") + + print() + print("=" * 78) + print(f"{'#':<3} {'scenario':<18} {'status':<8} reason") + print("-" * 78) + for num, name, status, reason in results: + print(f"{num:<3} {name:<18} {status:<8} {reason}") + print("-" * 78) + counts = {s: sum(1 for r in results if r[2] == s) + for s in (PASS, PARTIAL, SKIP, FAIL)} + print(f"total: {len(results)} PASS={counts[PASS]} PARTIAL={counts[PARTIAL]} " + f"SKIP={counts[SKIP]} FAIL={counts[FAIL]}") + return 1 if counts[FAIL] else 0 + + +if __name__ == "__main__": + sys.exit(main()) \ No newline at end of file diff --git a/gateway-plugin/tests/ws_probe.py b/gateway-plugin/tests/ws_probe.py index 645e872..b912846 100644 --- a/gateway-plugin/tests/ws_probe.py +++ b/gateway-plugin/tests/ws_probe.py @@ -18,14 +18,64 @@ Options: --send TEXT send this message after pairing (default: "hello") --upload F M4: upload F (chunked media.upload) and attach it to the message.send via media_refs ---pull-offer M4: when a media.offer arrives during the turn, pull the - media (chunked) and verify the byte count + --pull-offer M4: when a media.offer arrives during the turn, pull the + media (chunked) and verify the byte count. Offers are + emitted right AFTER the final message (MEDIA: tag + extraction runs post-turn), so after the final the probe + keeps listening for --offer-grace seconds for one. --sync C M5: after pairing, send sync {cursor: C} and print the - replay + sync.done (no turn is driven) + replay + sync.done (no turn is driven) --fcm-token M5: attach this FCM token to the hello payload --fcm-reg M5: after pairing, send fcm.register with --fcm-token --timeout S seconds to wait for the final reply (default 120) --authfail expect an auth rejection (wrong token) and exit 0 on it + +Assertion modes (checked after the turn; see exit codes below): + --assert-turn M2: the turn produced message.start -> >=1 + message.update -> message.stop (scenario 2) + --assert-reasoning M2: the final message.stop carries a non-empty + reasoning field (scenario 3) + --assert-tools M2: >=1 tool.start with a matching tool.end + (matched by index; scenario 4) + --assert-commentary M2: >=1 commentary frame (scenario 5) + --assert-read-receipt M7: a read.receipt frame arrives after the sent + message. SKIPs (exit 0) when the frame never + arrives (old gateway without the M7 frame). + --assert-status M7: a status frame is received. SKIPs (exit 0) + when the frame never arrives. + +Request modes (no turn driven unless --send/--upload also given): + --search Q [--scope all|chat] [--chat-id C] + M3: send search {query, scope, limit} and assert >=1 hit + in search.results (scenario 8). With --send, the turn is + driven first, then the search runs. + --channel-create NAME M3: send channel.create, print the new chat_id + ("== channel created: "), exit + --channel-delete CHAT M3: send channel.delete, assert channel.deleted + --channel-list M3: send channel.list, print the directory + --watch CHAT_ID wait up to --timeout for a message to land in + CHAT_ID (cron delivery E2E, scenario 7) + +Exit codes: + 0 ok (incl. SKIP for absent M7 frames) + 2 connect failed + 3 no hello.ack + 4 expected hello.ack, got something else + 5 --authfail but the token was accepted + 6 timeout waiting for the final message + 7 no final assistant message + 8 upload/sync failed + 9 media pull failed + 10 --assert-turn failed (no ordered start/update/stop segment) + 11 --assert-reasoning failed (final has no non-empty reasoning) + 12 --assert-tools failed (no tool.start with a matching tool.end) + 13 --assert-commentary failed (no commentary frame) + 14 --search failed (error or zero hits) + 15 --channel-create / --channel-list failed + 16 --channel-delete failed + 17 --watch timed out (no message landed in the channel) + 18 --assert-read-receipt failed (frame arrived before the sent message) + 19 --assert-status failed (status frame arrived with an empty payload) """ import argparse @@ -105,6 +155,19 @@ def _print_frame(raw): extra = f" cursor={payload.get('cursor')}" elif ftype == "sync.done": extra = f" cursor={payload.get('cursor')}" + elif ftype == "search.results": + hits = payload.get("hits") or [] + extra = f" query={payload.get('query')!r} scope={payload.get('scope')} hits={len(hits)}" + elif ftype == "channel.created": + extra = f" chat_id={payload.get('chat_id')} name={payload.get('name')!r}" + elif ftype == "channel.deleted": + extra = f" chat_id={payload.get('chat_id')}" + elif ftype == "channel.list": + extra = f" channels={len(payload.get('channels') or [])}" + elif ftype == "read.receipt": + extra = f" payload={ {k: payload[k] for k in list(payload)[:4]} }" + elif ftype == "status": + extra = f" payload={ {k: payload[k] for k in list(payload)[:4]} }" scope = f" chat={chat}" if chat else "" idpart = f" id={fid}" if fid is not None else "" print(f" <- {ftype}{idpart}{scope}{extra}") @@ -189,6 +252,228 @@ async def pull_media(ws, media_id: str, request_id: int, expected_size: int | No raise RuntimeError(f"pull failed: {data['payload']}") +class _TurnState: + """Assertion-relevant facts collected while driving a turn.""" + + def __init__(self): + self.seq: list[tuple[str, str | None]] = [] # (type, message_id) + self.tool_starts: set[int] = set() + self.tool_ends: set[int] = set() + self.commentary = 0 + self.final_stop_reasoning: str | None = None + self.final_message_reasoning: str | None = None + self.user_echo_seen = False + self.read_receipt: bool | None = None # None = never arrived + self.status_seen = False + self.status_empty = False + self.pulled = False + + def track(self, ftype: str, payload: dict) -> None: + if ftype in ("message.start", "message.update", "message.stop"): + self.seq.append((ftype, payload.get("message_id"))) + if ftype == "message.stop": + r = payload.get("reasoning") + if isinstance(r, str) and r.strip(): + self.final_stop_reasoning = r + if ftype == "message" and payload.get("role") == "assistant": + r = payload.get("reasoning") + if isinstance(r, str) and r.strip(): + self.final_message_reasoning = r + if ftype == "message" and payload.get("role") == "user": + self.user_echo_seen = True + if ftype == "tool.start" and isinstance(payload.get("index"), int): + self.tool_starts.add(payload["index"]) + if ftype == "tool.end" and isinstance(payload.get("index"), int): + self.tool_ends.add(payload["index"]) + if ftype == "commentary": + self.commentary += 1 + if ftype == "read.receipt": + self.read_receipt = self.user_echo_seen + if ftype == "status": + self.status_seen = True + if not payload: + self.status_empty = True + + +def _evaluate_assertions(args, st: _TurnState) -> list[tuple[int, bool, str]]: + """Evaluate the enabled assertion modes. Returns (exit_code, ok, message) + per failed-or-passed assertion; SKIPs are printed here and not returned.""" + results: list[tuple[int, bool, str]] = [] + if args.assert_turn: + ok = False + for mid in {m for _, m in st.seq if m is not None}: + events = [t for t, m in st.seq if m == mid] + if "message.start" in events and "message.stop" in events: + i_start = events.index("message.start") + i_stop = events.index("message.stop") + if any(i_start < i < i_stop + for i, e in enumerate(events) if e == "message.update"): + ok = True + break + results.append((10, ok, + "assert-turn: no message.start -> >=1 message.update -> message.stop")) + if args.assert_reasoning: + reasoning = st.final_stop_reasoning or st.final_message_reasoning + results.append((11, bool(reasoning), + "assert-reasoning: final message has no non-empty reasoning")) + if args.assert_tools: + ok = bool(st.tool_starts) and bool(st.tool_starts & st.tool_ends) + results.append((12, ok, + "assert-tools: no tool.start with a matching tool.end")) + if args.assert_commentary: + results.append((13, st.commentary >= 1, + "assert-commentary: no commentary frame")) + if args.assert_read_receipt: + if st.read_receipt is None: + print("== SKIP: no read.receipt frame (M7 frame not live on this gateway)") + elif not st.read_receipt: + results.append((18, False, + "assert-read-receipt: read.receipt arrived before the sent message")) + if args.assert_status: + if not st.status_seen: + print("== SKIP: no status frame (M7 frame not live on this gateway)") + elif st.status_empty: + results.append((19, False, + "assert-status: status frame arrived with an empty payload")) + return results + + +async def _recv_frames(ws, timeout: float): + """Yield parsed frames (dicts) until *timeout* seconds elapse.""" + deadline = time.time() + timeout + while time.time() < deadline: + try: + raw = await asyncio.wait_for(ws.recv(), timeout=deadline - time.time()) + except asyncio.TimeoutError: + return + if isinstance(raw, (bytes, bytearray)): + continue + data = _print_frame(raw) + if data is not None: + yield data + + +async def _search_mode(ws, args, next_id: int) -> int: + """M3: send a search frame, wait for search.results, assert >=1 hit.""" + req_id = next_id + payload = {"query": args.search, "scope": args.scope, "limit": 20} + if args.scope == "chat": + payload["chat_id"] = args.chat_id + await ws.send(json.dumps({"v": 1, "id": req_id, "type": "search", "payload": payload})) + print(f" -> search id={req_id} query={args.search!r} scope={args.scope}") + async for data in _recv_frames(ws, timeout=30): + if data.get("type") == "search.results" and data.get("id") == req_id: + hits = (data.get("payload") or {}).get("hits") or [] + print(f"== search: {len(hits)} hit(s)") + for h in hits[:10]: + print(f" hit chat={h.get('chat_id')} role={h.get('role')} " + f"snippet={str(h.get('snippet'))[:100]!r}") + if hits: + return 0 + print("!! search: no hits") + return 14 + if data.get("type") == "error": + print(f"!! search failed: {data.get('payload')}") + return 14 + print("!! search: no search.results within 30s") + return 14 + + +async def _channel_create_mode(ws, args) -> int: + """M3: channel.create -> channel.created; print the new chat_id.""" + req_id = 1 + await ws.send(json.dumps({ + "v": 1, "id": req_id, "type": "channel.create", + "payload": {"name": args.channel_create}, + })) + print(f" -> channel.create id={req_id} name={args.channel_create!r}") + async for data in _recv_frames(ws, timeout=30): + if data.get("type") == "channel.created" and data.get("id") == req_id: + chat_id = (data.get("payload") or {}).get("chat_id") + print(f"== channel created: {chat_id}") + await ws.close() + return 0 + if data.get("type") == "error": + print(f"!! channel.create failed: {data.get('payload')}") + await ws.close() + return 15 + print("!! channel.create: no channel.created within 30s") + await ws.close() + return 15 + + +async def _channel_delete_mode(ws, args) -> int: + """M3: channel.delete -> channel.deleted.""" + req_id = 1 + await ws.send(json.dumps({ + "v": 1, "id": req_id, "type": "channel.delete", + "payload": {"chat_id": args.channel_delete}, + })) + print(f" -> channel.delete id={req_id} chat_id={args.channel_delete!r}") + async for data in _recv_frames(ws, timeout=30): + if data.get("type") == "channel.deleted" and data.get("id") == req_id: + print(f"== channel deleted: {args.channel_delete}") + await ws.close() + return 0 + if data.get("type") == "error": + print(f"!! channel.delete failed: {data.get('payload')}") + await ws.close() + return 16 + print("!! channel.delete: no channel.deleted within 30s") + await ws.close() + return 16 + + +async def _channel_list_mode(ws, args) -> int: + """M3: channel.list -> print the directory.""" + req_id = 1 + await ws.send(json.dumps({"v": 1, "id": req_id, "type": "channel.list", "payload": {}})) + print(" -> channel.list") + async for data in _recv_frames(ws, timeout=30): + if data.get("type") == "channel.list" and data.get("id") == req_id: + for c in (data.get("payload") or {}).get("channels") or []: + print(f"== channel: {c.get('chat_id')} name={c.get('name')!r} " + f"default={bool(c.get('is_default'))}") + await ws.close() + return 0 + if data.get("type") == "error": + print(f"!! channel.list failed: {data.get('payload')}") + await ws.close() + return 15 + print("!! channel.list: no response within 30s") + await ws.close() + return 15 + + +async def _watch_mode(ws, args) -> int: + """Wait up to --timeout for a message to land in args.watch (cron E2E).""" + print(f"== watching {args.watch} for a message (timeout {args.timeout:.0f}s)") + deadline = time.time() + args.timeout + while time.time() < deadline: + try: + raw = await asyncio.wait_for(ws.recv(), timeout=deadline - time.time()) + except asyncio.TimeoutError: + print(f"!! timeout after {args.timeout:.0f}s watching {args.watch}") + await ws.close() + return 17 + if isinstance(raw, (bytes, bytearray)): + continue + data = _print_frame(raw) + if data is None: + continue + if data.get("chat_id") != args.watch: + continue + ftype = data.get("type") + payload = data.get("payload") or {} + if ftype == "message" and payload.get("role") in ("assistant", "cron"): + print(f"== message landed in {args.watch}: {str(payload.get('text'))[:120]!r}") + await ws.close() + return 0 + print(f"!! no message landed in {args.watch}") + await ws.close() + return 17 + + async def run(args) -> int: url = args.url token = args.token @@ -266,8 +551,18 @@ async def run(args) -> int: await ws.close() return 8 - if not args.send and not args.upload: - print("== paired OK (no --send/--upload; exiting)") + # M3: request modes (no turn driven). + if args.channel_create: + return await _channel_create_mode(ws, args) + if args.channel_delete: + return await _channel_delete_mode(ws, args) + if args.channel_list: + return await _channel_list_mode(ws, args) + if args.watch: + return await _watch_mode(ws, args) + + if not args.send and not args.upload and not args.search: + print("== paired OK (no --send/--upload/--search; exiting)") await ws.close() return 0 @@ -284,63 +579,114 @@ async def run(args) -> int: return 8 media_refs.append(media_ref) - # Drive a turn. - msg_id = next_id - send_payload: dict = {"text": args.send or ""} - if media_refs: - send_payload["media_refs"] = media_refs - send_frame = { - "v": 1, - "id": msg_id, - "type": "message.send", - "chat_id": "android:default", - "payload": send_payload, - } - await ws.send(json.dumps(send_frame)) - print(f" -> message.send id={msg_id} text={args.send!r} media_refs={media_refs}") - - deadline = time.time() + args.timeout + # Drive a turn (if --send or --upload). + st = _TurnState() got_final = False - seen_final_frame = False - while time.time() < deadline: - try: - raw = await asyncio.wait_for(ws.recv(), timeout=deadline - time.time()) - except asyncio.TimeoutError: - print(f"!! timeout after {args.timeout}s waiting for final message") - await ws.close() - return 6 - data = _print_frame(raw) - if data is None: - continue - ftype = data.get("type") - payload = data.get("payload") or {} - # M4: fetch offered media live (outbound direction). - if ftype == "media.offer" and args.pull_offer and payload.get("media_id"): + if args.send or args.upload: + msg_id = next_id + send_payload: dict = {"text": args.send or ""} + if media_refs: + send_payload["media_refs"] = media_refs + send_frame = { + "v": 1, + "id": msg_id, + "type": "message.send", + "chat_id": "android:default", + "payload": send_payload, + } + await ws.send(json.dumps(send_frame)) + print(f" -> message.send id={msg_id} text={args.send!r} media_refs={media_refs}") + + deadline = time.time() + args.timeout + seen_final_frame = False + while time.time() < deadline: try: - next_id = await pull_media( - ws, payload.get("media_id"), next_id, payload.get("size") - ) - except Exception as e: - print(f"!! pull failed: {e}") + raw = await asyncio.wait_for(ws.recv(), timeout=deadline - time.time()) + except asyncio.TimeoutError: + print(f"!! timeout after {args.timeout}s waiting for final message") await ws.close() - return 9 - # A standalone assistant `message` (non-streaming) is immediately final. - if ftype == "message" and payload.get("role") == "assistant": - got_final = True - break - # A `message.stop` finalizes a streaming segment; the turn is done once - # typing stops afterwards (multi-segment turns have several stops). - if ftype == "message.stop": - seen_final_frame = True - if ftype == "typing" and payload.get("on") is False and seen_final_frame: - got_final = True - break - await ws.close() - if got_final: + return 6 + data = _print_frame(raw) + if data is None: + continue + ftype = data.get("type") + payload = data.get("payload") or {} + st.track(ftype, payload) + # M4: fetch offered media live (outbound direction). + if ftype == "media.offer" and args.pull_offer and payload.get("media_id"): + try: + await pull_media( + ws, payload["media_id"], next_id, payload.get("size") + ) + next_id += 1 + st.pulled = True + except Exception as e: + print(f"!! pull failed: {e}") + await ws.close() + return 9 + # A standalone assistant `message` (non-streaming) is immediately final. + if ftype == "message" and payload.get("role") == "assistant": + got_final = True + break + # A `message.stop` finalizes a streaming segment; the turn is done + # once typing stops afterwards (multi-segment turns have several + # stops). + if ftype == "message.stop": + seen_final_frame = True + if ftype == "typing" and payload.get("on") is False and seen_final_frame: + got_final = True + break + + # M4: media offers are emitted right AFTER the final message (the + # MEDIA: tag is extracted post-turn); give them a grace window. + if got_final and args.pull_offer and not st.pulled: + grace_deadline = time.time() + args.offer_grace + while time.time() < grace_deadline: + try: + raw = await asyncio.wait_for( + ws.recv(), timeout=grace_deadline - time.time() + ) + except asyncio.TimeoutError: + break + if isinstance(raw, (bytes, bytearray)): + continue + data = _print_frame(raw) + if data is None: + continue + if data.get("type") == "media.offer" and (data.get("payload") or {}).get("media_id"): + try: + await pull_media( + ws, data["payload"]["media_id"], next_id, + data["payload"].get("size"), + ) + next_id += 1 + st.pulled = True + except Exception as e: + print(f"!! pull failed: {e}") + await ws.close() + return 9 + break + if not st.pulled: + print(f"== no media.offer within {args.offer_grace:.0f}s grace") + + if not got_final: + print("!! no final assistant message") + await ws.close() + return 7 print("== final assistant message received") - return 0 - print("!! no final assistant message") - return 7 + + # M3: optional search (standalone, or after the turn). + if args.search: + rc = await _search_mode(ws, args, next_id) + await ws.close() + return rc + + await ws.close() + for code, ok, msg in _evaluate_assertions(args, st): + if not ok: + print(f"!! {msg}") + return code + return 0 def main() -> int: @@ -362,9 +708,42 @@ def main() -> int: p.add_argument("--timeout", type=float, default=120.0) p.add_argument("--authfail", action="store_true", help="expect an auth rejection (wrong token)") + p.add_argument("--assert-turn", action="store_true", + help="assert message.start -> >=1 message.update -> message.stop") + p.add_argument("--assert-reasoning", action="store_true", + help="assert the final message.stop carries non-empty reasoning") + p.add_argument("--assert-tools", action="store_true", + help="assert >=1 tool.start with a matching tool.end") + p.add_argument("--assert-commentary", action="store_true", + help="assert >=1 commentary frame") + p.add_argument("--assert-read-receipt", action="store_true", + help="assert a read.receipt arrives after the sent message " + "(SKIP if absent; M7)") + p.add_argument("--assert-status", action="store_true", + help="assert a status frame is received (SKIP if absent; M7)") + p.add_argument("--search", default="", + help="M3: send search {query, scope, limit}, assert >=1 hit") + p.add_argument("--scope", choices=("all", "chat"), default="all", + help="search scope (default all)") + p.add_argument("--chat-id", default="android:default", + help="chat_id for --scope chat (default android:default)") + p.add_argument("--channel-create", default="", + help="M3: create a channel, print its chat_id, exit") + p.add_argument("--channel-delete", default="", + help="M3: delete (archive) a channel, exit") + p.add_argument("--channel-list", action="store_true", + help="M3: list channels, exit") + p.add_argument("--watch", default="", + help="wait up to --timeout for a message to land in this chat_id") + p.add_argument("--offer-grace", type=float, default=15.0, + help="seconds to wait for a media.offer after the final " + "message when --pull-offer (default 15)") args = p.parse_args() if not args.token and not args.authfail: p.error("--token (or $ANDROID_TOKEN) is required") + if args.assert_read_receipt and not args.send: + p.error("--assert-read-receipt requires --send (the receipt must follow " + "the sent message)") return asyncio.run(run(args)) diff --git a/gateway-plugin/ws_server.py b/gateway-plugin/ws_server.py index 24e9c17..e065649 100644 --- a/gateway-plugin/ws_server.py +++ b/gateway-plugin/ws_server.py @@ -11,7 +11,9 @@ Per-connection handler: 2. On success: register in the device registry (SQLite) + connection registry (``device_id -> {ws, caps, fcm_token}``), send ``hello.ack {server_caps, sync_cursor, channels[]}``. - 3. Loop: decode frames, dispatch to adapter inbound handlers. + 3. Loop: decode frames, dispatch to adapter inbound handlers. Inbound JSON + frames are rate-limited per connection (token bucket, ``INBOUND_RATE_PER_S`` + / ``INBOUND_BURST``); binary media-upload chunks are exempt. 4. On close: deregister. Routing: ``broadcast(frame)`` sends to ALL connected devices (single-user @@ -44,12 +46,46 @@ HELLO_TIMEOUT_S = 10.0 # rest of the broadcast). The peer's own ping timeout reaps it afterwards. SEND_TIMEOUT_S = 10.0 +# Inbound JSON control-frame rate limit (per connection, token bucket). +# A legitimate app sends pings + occasional user-initiated requests โ€” far +# below 20/s sustained. Binary media-upload chunks are EXEMPT (see +# ``_on_frame``): a 100 MB upload is 400 x 256 KiB frames in a tight loop +# and would exhaust any sane bucket; uploads are bounded instead by the +# per-frame ``max_size`` and the per-upload total cap (``media.py``). +INBOUND_RATE_PER_S = 20.0 +INBOUND_BURST = 40 + # Close codes (4000-4999 are reserved for applications). CLOSE_AUTH_FAILED = 4401 CLOSE_REPLACED = 4402 +CLOSE_RATE_LIMITED = 4403 CLOSE_SHUTDOWN = 1001 +class _TokenBucket: + """Minimal token bucket (stdlib only). One instance per connection.""" + + __slots__ = ("rate", "burst", "tokens", "updated_at") + + def __init__(self, rate: float, burst: int): + self.rate = rate + self.burst = burst + self.tokens = float(burst) + self.updated_at = time.monotonic() + + def consume(self) -> bool: + """Try to take one token. Refills at ``rate``/s up to ``burst``.""" + now = time.monotonic() + elapsed = now - self.updated_at + if elapsed > 0: + self.tokens = min(self.burst, self.tokens + elapsed * self.rate) + self.updated_at = now + if self.tokens >= 1.0: + self.tokens -= 1.0 + return True + return False + + @dataclass class DeviceConnection: """One live, authenticated device socket.""" @@ -61,6 +97,9 @@ class DeviceConnection: fcm_token: Optional[str] = None ntfy_topic: Optional[str] = None connected_at: float = field(default_factory=time.time) + rate_bucket: _TokenBucket = field( + default_factory=lambda: _TokenBucket(INBOUND_RATE_PER_S, INBOUND_BURST) + ) class WsServer: @@ -254,6 +293,9 @@ class WsServer: ) try: await ws.send(ack.to_json()) + # M7: tell late-joining clients the current gateway health state + # (the startup broadcast only reaches clients already connected). + await ws.send(protocol.status(self._adapter.gateway_status()).to_json()) except Exception: return logger.info("android: device paired: %s (%s)", device_name, device_id) @@ -261,7 +303,11 @@ class WsServer: # 3. frame loop ------------------------------------------------------ try: async for raw in ws: - await self._on_frame(ws, device_id, raw) + # ``_on_frame`` returns False once it has closed the socket + # (rate limit); stop draining the buffered frames so a + # flood doesn't re-trigger the error+close per frame. + if not await self._on_frame(ws, device_id, raw): + break except ConnectionClosed: pass except Exception: @@ -280,16 +326,38 @@ class WsServer: # โ”€โ”€ Inbound dispatch โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€ - async def _on_frame(self, ws: ServerConnection, device_id: str, raw: Any) -> None: + async def _on_frame(self, ws: ServerConnection, device_id: str, raw: Any) -> bool: + """Dispatch one inbound frame. Returns False once the socket has been + closed (rate limit) so the caller stops draining buffered frames.""" # M4: binary frames are media upload chunks (raw bytes, no JSON - # envelope). Route them to the active upload session. + # envelope). Route them to the active upload session. They are + # EXEMPT from the inbound rate limit: a 100 MB upload is 400 x + # 256 KiB frames in a tight loop, which would exhaust any sane + # frame bucket. Uploads are bounded instead by the per-frame + # ``max_size`` and the per-upload total cap (``media.py``). if isinstance(raw, (bytes, bytearray, memoryview)): await self._adapter.on_media_chunk(device_id, bytes(raw)) - return + return True + + # Inbound rate limit (JSON control frames only). On exceed: error + + # close, same pattern as auth rejection. + conn = self._connection_for(ws) + if conn is not None and not conn.rate_bucket.consume(): + logger.warning( + "android: inbound rate limit exceeded for %s; closing", device_id + ) + await self._send_quiet( + ws, + protocol.error( + protocol.ERR_RATE_LIMITED, "inbound frame rate limit exceeded" + ), + ) + await self._close_quiet(ws, CLOSE_RATE_LIMITED, "rate limited") + return False frame = protocol.Frame.from_json(raw) if frame is None: - return # malformed JSON: ignore (forward-compat) + return True # malformed JSON: ignore (forward-compat) if frame.type == protocol.TYPE_PING: ts = frame.payload.get("ts") @@ -319,9 +387,18 @@ class WsServer: elif frame.type == protocol.TYPE_FCM_REGISTER: await self._adapter.on_fcm_register(frame, device_id) # Unknown types are ignored (forward-compat). + return True # โ”€โ”€ Helpers โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€ + def _connection_for(self, ws: ServerConnection) -> Optional[DeviceConnection]: + """The live registry entry for this exact socket (identity match, so + a replaced socket never consumes the new connection's bucket).""" + for conn in self._connections.values(): + if conn.ws is ws: + return conn + return None + async def _send_quiet(self, ws: ServerConnection, frame: protocol.Frame) -> None: try: await ws.send(frame.to_json())