diff --git a/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt b/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt index a390cb2..14027de 100644 --- a/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt +++ b/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt @@ -105,6 +105,14 @@ class AndroidSecureStore(context: Context) : SecureStore { get() = prefs.getString(KEY_TOOL_DETAIL, "truncated").orEmpty().ifBlank { "truncated" } set(value) = prefs.edit().putString(KEY_TOOL_DETAIL, value).apply() + override var streamingEnabled: Boolean + get() = prefs.getBoolean(KEY_STREAMING_ENABLED, true) + set(value) = prefs.edit().putBoolean(KEY_STREAMING_ENABLED, value).apply() + + override var reasoningAutoCollapse: Boolean + get() = prefs.getBoolean(KEY_REASONING_AUTO_COLLAPSE, true) + set(value) = prefs.edit().putBoolean(KEY_REASONING_AUTO_COLLAPSE, value).apply() + override fun savePairing(url: String, token: String) { serverUrl = url this.token = token @@ -126,5 +134,7 @@ class AndroidSecureStore(context: Context) : SecureStore { const val KEY_NTFY_SERVER = "ntfy_server" const val KEY_THREADS_ENABLED = "threads_enabled" const val KEY_TOOL_DETAIL = "tool_detail" + const val KEY_STREAMING_ENABLED = "streaming_enabled" + const val KEY_REASONING_AUTO_COLLAPSE = "reasoning_auto_collapse" } } \ No newline at end of file diff --git a/app/shared/src/androidMain/kotlin/iris/ui/ContextMenu.kt b/app/shared/src/androidMain/kotlin/iris/ui/ContextMenu.kt new file mode 100644 index 0000000..1ee7129 --- /dev/null +++ b/app/shared/src/androidMain/kotlin/iris/ui/ContextMenu.kt @@ -0,0 +1,5 @@ +package iris.ui + +import androidx.compose.ui.Modifier + +actual fun Modifier.rightClick(onRightClick: () -> Unit): Modifier = this \ No newline at end of file diff --git a/app/shared/src/commonMain/kotlin/iris/IrisApp.kt b/app/shared/src/commonMain/kotlin/iris/IrisApp.kt index 0093743..df4f5b9 100644 --- a/app/shared/src/commonMain/kotlin/iris/IrisApp.kt +++ b/app/shared/src/commonMain/kotlin/iris/IrisApp.kt @@ -15,10 +15,8 @@ import iris.platform.setActiveController import iris.state.IrisController import iris.ui.screens.ChatScreen import iris.ui.screens.ConnectScreen -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. @@ -63,9 +61,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)) + // Connecting / Reconnecting / Connected all render the chat; the header + // StatusChip + connection banner show the link state without + // blocking the view (M7's full-screen spinner is gone). 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 55f4afd..eb0c051 100644 --- a/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt +++ b/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt @@ -106,6 +106,12 @@ class ChatStore { private var localSeq = 0 private var toolSeq = 0 + /** When false, `message.start`/`message.update` frames are ignored and each + * reply materializes as a single final message on `message.stop` + * (Settings → "Streaming"). The gateway keeps streaming (other devices + * may want it); this is a per-device display preference. */ + var streamingEnabled: Boolean = true + companion object { const val DEFAULT_LANE = "android:default" @@ -176,6 +182,10 @@ class ChatStore { private fun onMessage(lane: String, frame: Frame) { val p = frame.payloadAs() ?: return + // Auto-threaded send: the gateway created a fresh thread for this + // message; the optimistic bubble is still in the parent flat lane. + // Relocate it so the echo (below) lands in the thread lane. + if (p.role == ROLE_USER) relocatePendingToThreadLane(lane, p) updateLane(lane) { list -> val byId = list.indexOfFirst { it.id == p.messageId } if (byId >= 0) { @@ -222,6 +232,27 @@ class ChatStore { } } + /** + * Remove the matching optimistic user bubble from the parent flat lane of + * [lane] (an auto-created thread). The gateway moved the message into the + * fresh thread, so the pending bubble must follow it; the echo then replaces + * it in the thread lane. No-op when [lane] is not a thread lane or no + * pending bubble with the same text exists in the flat lane. + */ + private fun relocatePendingToThreadLane(lane: String, p: MessagePayload) { + val (chatId, threadId) = parseLane(lane) + if (threadId == null) return + val flatLane = chatId + val map = _lanes.value.toMutableMap() + val flatList = map[flatLane].orEmpty() + val idx = flatList.indexOfLast { + it is MessageItem && it.pending && it.role == ROLE_USER && it.text == p.text + } + if (idx < 0) return + map[flatLane] = flatList.toMutableList().also { it.removeAt(idx) } + _lanes.value = map + } + /** Merge server media refs into existing items, keeping local paths. */ private fun mergeMedia(existing: List, incoming: List): List { if (incoming.isEmpty()) return existing @@ -239,6 +270,7 @@ class ChatStore { // ── message.start (open a live streaming bubble) ────────────────────── private fun onMessageStart(lane: String, frame: Frame) { + if (!streamingEnabled) return val p = frame.payloadAs() ?: return updateLane(lane) { list -> if (list.any { it.id == p.messageId }) list @@ -441,6 +473,27 @@ class ChatStore { } } + /** + * Load a history page into [lane] (oldest → newest). The history is the + * authoritative full list of final messages for the lane; it replaces the + * lane's final messages and preserves non-final items (tool cards, live + * streaming bubbles, commentary) that are not part of the history. Used to + * restore the view on first open of a chat / after a process death, where + * the in-memory store is empty and the `sync` delta does not cover older + * messages. + */ + fun loadHistory(lane: String, messages: List) { + updateLane(lane) { list -> + val historyIds = messages.map { it.id }.toSet() + val preserved = list.filter { item -> + item !is MessageItem || item.id !in historyIds + } + (messages + preserved).sortedBy { item -> + (item as? MessageItem)?.ts?.takeIf { it > 0 } ?: Long.MAX_VALUE + } + } + } + fun clear() { _lanes.value = emptyMap() } diff --git a/app/shared/src/commonMain/kotlin/iris/data/SecureStore.kt b/app/shared/src/commonMain/kotlin/iris/data/SecureStore.kt index e85d363..9ced608 100644 --- a/app/shared/src/commonMain/kotlin/iris/data/SecureStore.kt +++ b/app/shared/src/commonMain/kotlin/iris/data/SecureStore.kt @@ -36,6 +36,12 @@ interface SecureStore { /** UI setting: tool verbosity — "everything" / "truncated" / "nothing". */ var toolDetail: String + /** UI setting: stream assistant replies live (token by token). */ + var streamingEnabled: Boolean + + /** UI setting: auto-collapse long reasoning blocks in the chat view. */ + var reasoningAutoCollapse: Boolean + fun savePairing(url: String, token: String) fun clear() } \ No newline at end of file diff --git a/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt b/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt index af96cd1..f05506f 100644 --- a/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt +++ b/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt @@ -92,6 +92,10 @@ class GatewayClient( private var socket: WebSocket? = null private var nextRequestId = 1 private var attempt = 0 + // True once a connection has been established this session; reset by + // start(). Drives Connecting (first dial) vs Reconnecting (redial after a + // drop) so the UI can show the right status without a blocking screen. + private var hasConnected = false private var lastLiveness: TimeMark = TimeSource.Monotonic.markNow() private val pending = mutableMapOf>() @@ -118,6 +122,7 @@ class GatewayClient( fun start() { if (connectJob?.isActive == true) return attempt = 0 + hasConnected = false connectJob = scope.launch { connectLoop() } } @@ -144,7 +149,7 @@ class GatewayClient( _state.value = State.Disconnected return } - _state.value = if (attempt == 0) State.Connecting else State.Reconnecting + _state.value = if (hasConnected) State.Reconnecting else State.Connecting val dial = dial(url, token) when (val result = dial.result) { is DialResult.AuthFailed -> { @@ -153,6 +158,7 @@ class GatewayClient( return } DialResult.Connected -> { + hasConnected = true attempt = 0 lastLiveness = TimeSource.Monotonic.markNow() dial.closed.await() @@ -289,16 +295,19 @@ class GatewayClient( // ── Outbound ────────────────────────────────────────────────────────── /** Send a text message (fire-and-forget; the server echoes it back). - * M4: [mediaRefs] reference completed uploads (media.upload.ack refs). */ + * M4: [mediaRefs] reference completed uploads (media.upload.ack refs). + * [autoThread] asks the gateway to mint a fresh thread for the message + * (auto-threading, docs/06 §6.3). */ fun sendMessage( chatId: String, text: String, threadId: String? = null, mediaRefs: List = emptyList(), + autoThread: Boolean = false, ) { val ws = socket ?: return val id = nextRequestId++ - ws.send(messageSendFrame(id, chatId, text, threadId, mediaRefs).toWire()) + ws.send(messageSendFrame(id, chatId, text, threadId, mediaRefs, autoThread).toWire()) } // ── M4: media upload / pull ─────────────────────────────────────────── diff --git a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt index 25d9657..05f6f0e 100644 --- a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt +++ b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt @@ -73,6 +73,7 @@ const val TYPE_SEARCH = "search" const val TYPE_SEARCH_RESULTS = "search.results" const val TYPE_SYNC = "sync" const val TYPE_SYNC_DONE = "sync.done" +const val TYPE_HISTORY = "history" // ── Error codes ───────────────────────────────────────────────────────── @@ -149,6 +150,9 @@ data class ChannelInfo( @SerialName("is_default") val isDefault: Boolean = false, @SerialName("parent_chat_id") val parentChatId: String? = null, val archived: Boolean = false, + /** Gateway minted this thread for an incoming message (auto-threading); + * the name is a derived title, upgraded by the LLM via channel.renamed. */ + val auto: Boolean = false, ) @Serializable @@ -289,6 +293,9 @@ data class MessageSendPayload( val text: String, @SerialName("reply_to") val replyTo: String? = null, @SerialName("media_refs") val mediaRefs: List = emptyList(), + /** Ask the gateway to mint a fresh thread for this message (auto- + * threading; only honored in a channel's flat lane, docs/06 §6.3). */ + @SerialName("auto_thread") val autoThread: Boolean = false, ) // ── typing / error / ping ─────────────────────────────────────────────── @@ -356,6 +363,28 @@ data class SyncPayload(val cursor: Long) @Serializable data class SyncDonePayload(val cursor: Long) +// ── history (full message history for a chat/thread) ──────────────────── + +/** A final message in a history page (user echo / assistant reply). */ +@Serializable +data class HistoryMessage( + @SerialName("message_id") val messageId: String, + val role: String, + val text: String, + val reasoning: String? = null, + val model: String? = null, + val tokens: Int? = null, + val ts: Long? = null, + val media: List = emptyList(), +) + +@Serializable +data class HistoryPayload( + val messages: List = emptyList(), + @SerialName("has_more") val hasMore: Boolean = false, + @SerialName("oldest_message_id") val oldestMessageId: String? = null, +) + // ── M5: push / notifications ──────────────────────────────────────────── /** Notification kinds (mirror of protocol.NOTIF_*). */ @@ -424,6 +453,7 @@ fun messageSendFrame( text: String, threadId: String? = null, mediaRefs: List = emptyList(), + autoThread: Boolean = false, ): Frame = Frame( id = id, @@ -432,7 +462,7 @@ fun messageSendFrame( threadId = threadId, payload = IrisJson.instance.encodeToJsonElement( MessageSendPayload.serializer(), - MessageSendPayload(text = text, mediaRefs = mediaRefs), + MessageSendPayload(text = text, mediaRefs = mediaRefs, autoThread = autoThread), ), ) @@ -499,6 +529,27 @@ fun syncFrame(id: Int, cursor: Long): Frame = ), ) +/** Request a page of full message history for a chat/thread (initial open / + * restart restore). [beforeMessageId] returns the page older than that id + * (null = newest page). */ +fun historyFrame( + id: Int, + chatId: String, + threadId: String? = null, + beforeMessageId: String? = null, + limit: Int = 50, +): Frame = + Frame( + id = id, + type = TYPE_HISTORY, + chatId = chatId, + threadId = threadId, + payload = buildJsonObject { + if (beforeMessageId != null) put("before_message_id", beforeMessageId) + put("limit", limit) + }, + ) + // ── M4 frame builders ──────────────────────────────────────────────────── fun mediaUploadStartFrame( diff --git a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt index ad88a3a..32a0610 100644 --- a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt +++ b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt @@ -13,7 +13,11 @@ import iris.platform.PickedFile import iris.platform.isAppForeground import iris.platform.mediaCacheBaseDir import iris.platform.postSystemNotification +import iris.protocol.ChannelDeletedPayload +import iris.protocol.ChannelInfo import iris.protocol.HIGH_PRIORITY_NOTIF_KINDS +import iris.protocol.HistoryMessage +import iris.protocol.HistoryPayload import iris.protocol.MediaOfferPayload import iris.protocol.MessagePayload import iris.protocol.MessageStopPayload @@ -25,6 +29,7 @@ import iris.protocol.SearchResultsPayload import iris.protocol.StatusPayload import iris.protocol.SyncDonePayload import iris.protocol.TYPE_ERROR +import iris.protocol.TYPE_HISTORY import iris.protocol.TYPE_NOTIFICATION import iris.protocol.TYPE_CHANNEL_CREATED import iris.protocol.TYPE_CHANNEL_DELETED @@ -50,6 +55,7 @@ import iris.protocol.channelDeleteFrame import iris.protocol.channelListFrame import iris.protocol.channelRenameFrame import iris.protocol.channelSetDefaultFrame +import iris.protocol.historyFrame import iris.protocol.searchFrame import iris.protocol.syncFrame import kotlinx.coroutines.CoroutineScope @@ -104,6 +110,12 @@ class IrisController( private val _homeChannel = MutableStateFlow("android:default") val homeChannel: StateFlow = _homeChannel.asStateFlow() + /** Lanes whose full history has been loaded this session (in-memory; reset + * on a process death, which is exactly when a reload is needed). The + * `sync` delta can seed a lane with recent frames without it being opened, + * so "lane is empty" is not a reliable first-open signal. */ + private val historyLoaded = mutableSetOf() + /** Gateway health state (M5: status frame; null = never received). */ private val _gatewayStatus = MutableStateFlow(null) val gatewayStatus: StateFlow = _gatewayStatus.asStateFlow() @@ -118,6 +130,30 @@ class IrisController( store.threadsEnabled = _threadsEnabled.value } + // ── Streaming toggle (Settings → "Streaming") ────────────────────────── + // Persisted. When off, the app ignores message.start/update frames and + // shows each reply as a single final message (per-device display + // preference, like tool verbosity — the gateway keeps streaming). + private val _streamingEnabled = MutableStateFlow(store.streamingEnabled) + val streamingEnabled: StateFlow = _streamingEnabled.asStateFlow() + + fun toggleStreaming() { + _streamingEnabled.value = !_streamingEnabled.value + store.streamingEnabled = _streamingEnabled.value + chat.streamingEnabled = _streamingEnabled.value + } + + // ── Reasoning auto-collapse (Settings → "Reasoning") ─────────────────── + // Persisted. When on, long reasoning blocks start collapsed (short ones + // stay expanded); when off, all reasoning blocks start expanded. + private val _reasoningAutoCollapse = MutableStateFlow(store.reasoningAutoCollapse) + val reasoningAutoCollapse: StateFlow = _reasoningAutoCollapse.asStateFlow() + + fun toggleReasoningAutoCollapse() { + _reasoningAutoCollapse.value = !_reasoningAutoCollapse.value + store.reasoningAutoCollapse = _reasoningAutoCollapse.value + } + // ── M3: search state ────────────────────────────────────────────────── private val _searchResults = MutableStateFlow>(emptyList()) val searchResults: StateFlow> = _searchResults.asStateFlow() @@ -215,6 +251,7 @@ class IrisController( } init { + chat.streamingEnabled = _streamingEnabled.value scope.launch { client.events.collect { frame -> when (frame.type) { @@ -255,7 +292,38 @@ class IrisController( TYPE_CHANNEL_CREATED, TYPE_CHANNEL_RENAMED, TYPE_CHANNEL_DELETED, - TYPE_CHANNEL_LIST -> channels.onFrame(frame) + TYPE_CHANNEL_LIST -> { + channels.onFrame(frame) + // Auto-created thread (Settings → "Threads"): the + // gateway minted it for a message in the flat lane we + // are viewing — jump into it; the user's message and + // the reply land there. + if (frame.type == TYPE_CHANNEL_CREATED) { + frame.payloadAs()?.let { info -> + if (info.kind == "thread" && info.auto && info.parentChatId != null) { + val (curChat, curThread) = chat.parseLane(chat.currentLane.value) + if (curThread == null && curChat == info.parentChatId) { + val lane = chat.laneKey(info.parentChatId, info.chatId) + chat.setLane(lane) + historyLoaded.add(lane) + } + } + } + } + // A thread/channel we were viewing got deleted — fall back + // to a valid lane (the thread's parent, or the home channel). + if (frame.type == TYPE_CHANNEL_DELETED) { + frame.payloadAs()?.let { p -> + val (curChat, curThread) = chat.parseLane(chat.currentLane.value) + when { + curThread == p.chatId -> openChannel(curChat) + curChat == p.chatId -> + channels.defaultChannel()?.let { openChannel(it.chatId) } + else -> Unit + } + } + } + } TYPE_SEARCH_RESULTS -> { frame.payloadAs()?.let { _searchResults.value = it.hits } _searching.value = false @@ -265,6 +333,22 @@ class IrisController( // cursor is authoritative server-side (outbox). frame.payloadAs()?.let { store.syncCursor = it.cursor } } + TYPE_HISTORY -> { + // Full message history for a chat/thread (initial open / + // restart restore). `sync` only replays the outbox delta + // since the saved cursor, so after a process death the + // in-memory store is empty and older messages are not in + // the delta — history loads the full list. + frame.payloadAs()?.let { p -> + val (chatId, threadId) = + (frame.chatId ?: return@let) to frame.threadId + val lane = chat.laneKey(chatId, threadId) + chat.loadHistory( + lane, + p.messages.map { it.toMessageItem() }, + ) + } + } TYPE_NOTIFICATION -> { frame.payloadAs()?.let { p -> pushBanner(p.kind, p.title, p.body, p.chatId, p.threadId) @@ -304,6 +388,12 @@ class IrisController( if (home != null) { _homeChannel.value = home chat.setLane(home) + // Restore the home lane after a process death: the + // in-memory store is empty and the `sync` delta does not + // cover messages older than the saved cursor, so load + // the full history. Skipped on a plain reconnect (the + // lane's history was already loaded this session). + loadHistory(home) } // M5: remember the ntfy server for the listener service. if (s.caps.pushNtfyServer.isNotBlank()) { @@ -324,14 +414,19 @@ class IrisController( // ── M3: navigation (lane switching) ─────────────────────────────────── - /** Switch to a channel's flat / "General" lane. */ + /** Switch to a channel's flat / "General" lane. Loads the full history on + * first open this session so the view is populated, not just the `sync` + * delta. */ fun openChannel(chatId: String) { chat.setLane(chatId) + loadHistory(chatId) } - /** Switch to a specific thread lane under [chatId]. */ + /** Switch to a specific thread lane under [chatId]. Loads the full history + * on first open this session. */ fun openThread(chatId: String, threadId: String) { chat.setLane(chat.laneKey(chatId, threadId)) + loadHistory(chatId, threadId) } // ── M3: channel directory ops (server is authoritative) ─────────────── @@ -388,6 +483,19 @@ class IrisController( client.sendFrame(syncFrame(0, cursor)) } + /** Request the full message history for a chat/thread (newest page). + * Populates the lane via the `history` response; used on first open of a + * chat and after a process death (restart) to restore the view, since the + * `sync` delta does not cover messages older than the saved cursor. */ + fun loadHistory(chatId: String, threadId: String? = null) { + val lane = chat.laneKey(chatId, threadId) + if (lane in historyLoaded) return + historyLoaded.add(lane) + // Newest page, sized to restore a full working view on restart / first + // open (older pages are reachable via scroll-up pagination). + client.sendFrame(historyFrame(0, chatId, threadId, limit = 200)) + } + /** Optimistic send: show immediately in the current lane, then hand to the gateway. * M4: [attachments] (uploaded) are attached via media_refs. */ fun send(text: String, attachments: List = emptyList()) { @@ -404,10 +512,17 @@ class IrisController( ) } chat.addPending(trimmed, lane, media) - client.sendMessage(chatId, trimmed, threadId, refs) + client.sendMessage(chatId, trimmed, threadId, refs, autoThread = wantsAutoThread(trimmed, threadId)) _attachments.value = emptyList() } + /** Auto-threading (Settings → "Threads", docs/06 §6.3): a message in a + * channel's flat lane gets its own fresh thread, AI-named by the gateway + * (Telegram topic-mode workflow). Slash commands and media-only sends + * stay in the flat lane. */ + private fun wantsAutoThread(text: String, threadId: String?): Boolean = + _threadsEnabled.value && threadId == null && text.isNotEmpty() && !text.startsWith("/") + /** M7: resend a failed user message (tap on the failed bubble). */ fun retrySend(messageId: String) { val lane = chat.currentLane.value @@ -415,7 +530,13 @@ class IrisController( 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 }) + client.sendMessage( + chatId, + item.text, + threadId, + item.media.map { it.mediaId }, + autoThread = wantsAutoThread(item.text, threadId), + ) } // ── M4: attachments (pick -> upload -> send) ────────────────────────── @@ -492,6 +613,21 @@ class IrisController( } } +/** Convert a history message (server) into a chat [MessageItem]. */ +private fun HistoryMessage.toMessageItem(): MessageItem = + MessageItem( + id = messageId, + role = role, + text = text, + ts = ts ?: 0, + reasoning = reasoning, + model = model, + tokens = tokens, + media = media.map { + MediaItem(mediaId = it.mediaId, kind = it.kind, mime = it.mime, size = it.size, filename = it.filename) + }, + ) + /** How much tool detail to show (Settings → "Tool detail"). */ enum class ToolDetail { EVERYTHING, // name + full args (collapsible) + output preview diff --git a/app/shared/src/commonMain/kotlin/iris/ui/ContextMenu.kt b/app/shared/src/commonMain/kotlin/iris/ui/ContextMenu.kt new file mode 100644 index 0000000..f844db8 --- /dev/null +++ b/app/shared/src/commonMain/kotlin/iris/ui/ContextMenu.kt @@ -0,0 +1,10 @@ +package iris.ui + +import androidx.compose.ui.Modifier + +/** + * Opens [onRightClick] on a right (secondary) mouse-button press. Desktop-only + * behavior; on touch platforms this is a no-op (the context menu is opened via + * long-press instead). + */ +expect fun Modifier.rightClick(onRightClick: () -> Unit): Modifier \ No newline at end of file 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 c1eb0bc..543731d 100644 --- a/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt +++ b/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt @@ -1,7 +1,9 @@ package iris.ui.screens +import androidx.compose.foundation.ExperimentalFoundationApi import androidx.compose.foundation.background import androidx.compose.foundation.clickable +import androidx.compose.foundation.combinedClickable import androidx.compose.foundation.focusable import androidx.compose.foundation.horizontalScroll import androidx.compose.foundation.layout.Arrangement @@ -102,6 +104,7 @@ import iris.protocol.ROLE_USER import iris.protocol.SearchHit import iris.state.IrisController import iris.state.ToolDetail +import iris.ui.rightClick import iris.ui.theme.IrisColors import iris.ui.theme.avatarColor import iris.util.formatDayLabel @@ -123,6 +126,7 @@ fun ChatScreen(controller: IrisController) { val typing by controller.typing.collectAsState() val toolDetail by controller.toolDetail.collectAsState() val threadsEnabled by controller.threadsEnabled.collectAsState() + val reasoningAutoCollapse by controller.reasoningAutoCollapse.collectAsState() val channels by controller.channels.channels.collectAsState() val (currentChatId, currentThreadId) = controller.chat.parseLane(currentLane) @@ -152,6 +156,9 @@ fun ChatScreen(controller: IrisController) { val nonThreadChannels = channels.filter { it.kind != "thread" } fun doSend() { + // No-op while the socket is down (sendMessage drops silently); the send + // button is disabled in that state, this guards the IME "Send" action. + if (state !is GatewayClient.State.Connected) return val ready = attachments.filter { it.mediaRef != null && it.error == null } if (input.isBlank() && ready.isEmpty()) return val text = input @@ -211,6 +218,9 @@ fun ChatScreen(controller: IrisController) { var showOverflow by remember { mutableStateOf(false) } var showRename by remember { mutableStateOf(false) } var showForget by remember { mutableStateOf(false) } + // Topic context menu (long-press / right-click): rename / delete a thread. + var renameThread by remember { mutableStateOf(null) } + var deleteThread by remember { mutableStateOf(null) } Column(modifier = Modifier.fillMaxSize()) { // Header (M7: avatar + two-line title pill, reference look). Wide // panes keep everything on one row; narrow phones move the utility @@ -316,6 +326,8 @@ fun ChatScreen(controller: IrisController) { onGeneral = { controller.openChannel(currentChatId) }, onThread = { t -> controller.openThread(currentChatId, t.chatId) }, onNewThread = { showNewThread = true }, + onRenameThread = { renameThread = it }, + onDeleteThread = { deleteThread = it }, ) } @@ -346,6 +358,7 @@ fun ChatScreen(controller: IrisController) { is MessageItem -> MessageBubble( msg = item, maxWidth = bubbleMaxWidth, + reasoningAutoCollapse = reasoningAutoCollapse, onRetry = { controller.retrySend(item.id) }, ) is ToolItem -> if (toolDetail != ToolDetail.NOTHING) ToolCard(item, toolDetail) @@ -404,6 +417,8 @@ fun ChatScreen(controller: IrisController) { // M7: connection-state banner (persistent, above the composer; no dismiss). val gatewayStatus by controller.gatewayStatus.collectAsState() val connBanner = when { + state is GatewayClient.State.Connecting -> + "Connecting to gateway…" state is GatewayClient.State.Reconnecting -> "Reconnecting to gateway… messages will sync automatically when the connection is back." gatewayStatus == "restarting" -> "Gateway is restarting…" @@ -429,8 +444,12 @@ fun ChatScreen(controller: IrisController) { } } - // Composer (M7: rounded pill + accent circular send button) - val canSend = input.isNotBlank() || attachments.any { it.mediaRef != null && it.error == null } + // Composer (M7: rounded pill + accent circular send button). Sending is gated + // on a live socket: sendMessage is a no-op while disconnected, so an + // ungated send would show a pending bubble that never resolves. + val isConnected = state is GatewayClient.State.Connected + val canSend = + isConnected && (input.isNotBlank() || attachments.any { it.mediaRef != null && it.error == null }) val layoutDensity = LocalDensity.current.density var textHeightPx by remember { mutableFloatStateOf(0f) } Row( @@ -524,6 +543,38 @@ fun ChatScreen(controller: IrisController) { }, ) } + // Topic context menu: rename a thread. + renameThread?.let { t -> + NameDialog( + title = "Rename topic", + onConfirm = { name -> + controller.renameChannel(t.chatId, name) + renameThread = null + }, + onDismiss = { renameThread = null }, + initialName = t.name, + confirmLabel = "Rename", + ) + } + // Topic context menu: delete a thread (and its messages from the chat). + deleteThread?.let { t -> + AlertDialog( + onDismissRequest = { deleteThread = null }, + title = { Text("Delete topic?") }, + text = { + Text("Delete \"${t.name}\"? It will be removed from this channel.") + }, + confirmButton = { + TextButton(onClick = { + deleteThread = null + controller.deleteChannel(t.chatId) + }) { Text("Delete") } + }, + dismissButton = { + TextButton(onClick = { deleteThread = null }) { Text("Cancel") } + }, + ) + } } } @@ -657,6 +708,7 @@ fun ChatScreen(controller: IrisController) { "new-thread" to "New topic (Ctrl+T)", "search" to "Search (Ctrl+F)", "toggle-threads" to "Toggle threads", + "toggle-reasoning-collapse" to "Toggle reasoning auto-collapse", "cycle-tool-detail" to "Cycle tool verbosity", "toggle-inspector" to "Toggle inspector", ), @@ -666,6 +718,7 @@ fun ChatScreen(controller: IrisController) { "new-thread" -> showNewThread = true "search" -> showSearch = true "toggle-threads" -> controller.toggleThreads() + "toggle-reasoning-collapse" -> controller.toggleReasoningAutoCollapse() "cycle-tool-detail" -> controller.cycleToolDetail() "toggle-inspector" -> showInspector = !showInspector } @@ -1098,6 +1151,8 @@ private fun TopicSwitcher( onGeneral: () -> Unit, onThread: (ChannelInfo) -> Unit, onNewThread: () -> Unit, + onRenameThread: (ChannelInfo) -> Unit, + onDeleteThread: (ChannelInfo) -> Unit, ) { Row( modifier = Modifier @@ -1108,22 +1163,74 @@ private fun TopicSwitcher( ) { TopicChip("General", currentThreadId == null, onClick = onGeneral) threads.forEach { t -> - TopicChip(t.name, currentThreadId == t.chatId, onClick = { onThread(t) }) + TopicChip( + label = t.name, + selected = currentThreadId == t.chatId, + onClick = { onThread(t) }, + onRename = { onRenameThread(t) }, + onDelete = { onDeleteThread(t) }, + ) } TopicChip("+ topic", false, onClick = onNewThread) } } +/** + * A topic chip. When [onRename]/[onDelete] are provided the chip also opens a + * context menu — long-press on touch, right-click on a desktop mouse — with + * "Rename" and "Delete" actions. + */ +@OptIn(ExperimentalFoundationApi::class) @Composable -internal fun TopicChip(label: String, selected: Boolean, onClick: () -> Unit) { +internal fun TopicChip( + label: String, + selected: Boolean, + onClick: () -> Unit, + onRename: (() -> Unit)? = null, + onDelete: (() -> Unit)? = null, +) { + var menuOpen by remember { mutableStateOf(false) } + val hasMenu = onRename != null || onDelete != null Box( modifier = Modifier .clip(RoundedCornerShape(12.dp)) .background(if (selected) IrisColors.primary else IrisColors.chip) - .clickable(onClick = onClick) + .then( + if (hasMenu) { + // Tap (both platforms) + long-press (touch / held mouse) + + // right-click (desktop mouse, no-op on touch). + Modifier + .combinedClickable(onClick = onClick, onLongClick = { menuOpen = true }) + .rightClick { menuOpen = true } + } else { + Modifier.clickable(onClick = onClick) + } + ) .padding(horizontal = 10.dp, vertical = 5.dp), ) { Text(label, fontSize = 12.sp, color = if (selected) Color.White else IrisColors.textTertiary) + if (hasMenu) { + DropdownMenu(expanded = menuOpen, onDismissRequest = { menuOpen = false }) { + onRename?.let { + DropdownMenuItem( + text = { Text("Rename") }, + onClick = { + menuOpen = false + it() + }, + ) + } + onDelete?.let { + DropdownMenuItem( + text = { Text("Delete") }, + onClick = { + menuOpen = false + it() + }, + ) + } + } + } } } @@ -1312,7 +1419,7 @@ private fun NotificationBanner( } @Composable -private fun MessageBubble(msg: MessageItem, maxWidth: Dp, onRetry: () -> Unit) { +private fun MessageBubble(msg: MessageItem, maxWidth: Dp, reasoningAutoCollapse: Boolean, onRetry: () -> Unit) { val isUser = msg.role == ROLE_USER val isCommentary = msg.isCommentary val bubbleColor = when { @@ -1344,7 +1451,7 @@ private fun MessageBubble(msg: MessageItem, maxWidth: Dp, onRetry: () -> Unit) { ) { // Reasoning block above the answer (assistant, non-commentary). if (!isUser && !isCommentary && !msg.reasoning.isNullOrBlank()) { - ReasoningBlock(msg.reasoning!!) + ReasoningBlock(msg.reasoning!!, reasoningAutoCollapse) Spacer(modifier = Modifier.height(6.dp)) } if (isUser) { @@ -1559,13 +1666,15 @@ private fun AttachmentChip(att: IrisController.PendingAttachment, onRemove: () - /** * Collapsible reasoning panel (M2): header "💭 Reasoning", monospace body, - * copy button. Collapsed by default when long; tap to toggle. + * copy button. With [autoCollapse] on (Settings → "Reasoning"), long blocks + * start collapsed and short ones expanded; with it off, all blocks start + * expanded. Tap to toggle. */ @Composable -private fun ReasoningBlock(reasoning: String) { +private fun ReasoningBlock(reasoning: String, autoCollapse: Boolean) { val clipboard = LocalClipboardManager.current - var expanded by remember { - mutableStateOf(reasoning.length <= 240) + var expanded by remember(autoCollapse) { + mutableStateOf(!autoCollapse || reasoning.length <= 240) } var copied by remember { mutableStateOf(false) } Column( @@ -1615,7 +1724,9 @@ private fun ReasoningBlock(reasoning: String) { */ @Composable private fun ToolCard(tool: ToolItem, detail: ToolDetail) { - var expanded by remember { mutableStateOf(false) } + // Everything starts expanded (show the full call + output); Truncated + // starts collapsed and reveals them on tap. + var expanded by remember(detail) { mutableStateOf(detail == ToolDetail.EVERYTHING) } val status = when { !tool.done -> null tool.ok -> "✓" @@ -1670,8 +1781,8 @@ private fun ToolCard(tool: ToolItem, detail: ToolDetail) { maxLines = if (expanded) Int.MAX_VALUE else 1, ) } - // Full args + output preview (Everything, expanded). - if (detail == ToolDetail.EVERYTHING && expanded) { + // Full args + output (revealed on expand in both Truncated and Everything). + if (expanded) { tool.args?.let { Spacer(modifier = Modifier.height(4.dp)) Text( 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 e670036..a2378b2 100644 --- a/app/shared/src/commonMain/kotlin/iris/ui/screens/ConnectScreen.kt +++ b/app/shared/src/commonMain/kotlin/iris/ui/screens/ConnectScreen.kt @@ -1,8 +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 import androidx.compose.foundation.layout.Spacer import androidx.compose.foundation.layout.fillMaxSize @@ -14,7 +12,6 @@ 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.OutlinedTextField import androidx.compose.material3.Text @@ -135,20 +132,3 @@ fun ConnectScreen( } } -/** 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/screens/SettingsScreen.kt b/app/shared/src/commonMain/kotlin/iris/ui/screens/SettingsScreen.kt index abdfdfb..a8ba875 100644 --- a/app/shared/src/commonMain/kotlin/iris/ui/screens/SettingsScreen.kt +++ b/app/shared/src/commonMain/kotlin/iris/ui/screens/SettingsScreen.kt @@ -41,6 +41,8 @@ import iris.ui.theme.IrisColors @Composable fun SettingsScreen(controller: IrisController, onBack: () -> Unit) { val threadsEnabled by controller.threadsEnabled.collectAsState() + val streamingEnabled by controller.streamingEnabled.collectAsState() + val reasoningAutoCollapse by controller.reasoningAutoCollapse.collectAsState() val toolDetail by controller.toolDetail.collectAsState() Box(modifier = Modifier.fillMaxSize().background(IrisColors.background)) { @@ -74,7 +76,7 @@ fun SettingsScreen(controller: IrisController, onBack: () -> Unit) { Column(modifier = Modifier.weight(1f)) { Text("🧵 Threads", fontSize = 14.sp) Text( - "Show topics in the chat view", + "Show topics; new messages in General start their own topic (AI-named)", fontSize = 12.sp, color = IrisColors.textDim, ) @@ -85,6 +87,44 @@ fun SettingsScreen(controller: IrisController, onBack: () -> Unit) { ) } } + SettingsCard { + Row( + verticalAlignment = Alignment.CenterVertically, + modifier = Modifier.fillMaxWidth(), + ) { + Column(modifier = Modifier.weight(1f)) { + Text("⚡ Streaming", fontSize = 14.sp) + Text( + "Show replies as they are generated", + fontSize = 12.sp, + color = IrisColors.textDim, + ) + } + Switch( + checked = streamingEnabled, + onCheckedChange = { controller.toggleStreaming() }, + ) + } + } + SettingsCard { + Row( + verticalAlignment = Alignment.CenterVertically, + modifier = Modifier.fillMaxWidth(), + ) { + Column(modifier = Modifier.weight(1f)) { + Text("🧠 Reasoning", fontSize = 14.sp) + Text( + "Auto-collapse long reasoning blocks", + fontSize = 12.sp, + color = IrisColors.textDim, + ) + } + Switch( + checked = reasoningAutoCollapse, + onCheckedChange = { controller.toggleReasoningAutoCollapse() }, + ) + } + } Text( "Tools", diff --git a/app/shared/src/desktopMain/kotlin/iris/platform/DesktopSecureStore.kt b/app/shared/src/desktopMain/kotlin/iris/platform/DesktopSecureStore.kt index 11eb9c5..09b34cf 100644 --- a/app/shared/src/desktopMain/kotlin/iris/platform/DesktopSecureStore.kt +++ b/app/shared/src/desktopMain/kotlin/iris/platform/DesktopSecureStore.kt @@ -36,6 +36,8 @@ class DesktopSecureStore : SecureStore { val ntfyServer: String = "", val threadsEnabled: Boolean = false, val toolDetail: String = "truncated", + val streamingEnabled: Boolean = true, + val reasoningAutoCollapse: Boolean = true, ) init { @@ -149,6 +151,20 @@ class DesktopSecureStore : SecureStore { save(d.copy(toolDetail = value)) } + override var streamingEnabled: Boolean + get() = load().streamingEnabled + set(value) { + val d = load() + save(d.copy(streamingEnabled = value)) + } + + override var reasoningAutoCollapse: Boolean + get() = load().reasoningAutoCollapse + set(value) { + val d = load() + save(d.copy(reasoningAutoCollapse = value)) + } + override fun savePairing(url: String, token: String) { val d = load() save(d.copy(serverUrl = url.trim())) diff --git a/app/shared/src/desktopMain/kotlin/iris/ui/ContextMenu.kt b/app/shared/src/desktopMain/kotlin/iris/ui/ContextMenu.kt new file mode 100644 index 0000000..5d79882 --- /dev/null +++ b/app/shared/src/desktopMain/kotlin/iris/ui/ContextMenu.kt @@ -0,0 +1,20 @@ +package iris.ui + +import androidx.compose.ui.ExperimentalComposeUiApi +import androidx.compose.ui.Modifier +import androidx.compose.ui.input.pointer.PointerButton +import androidx.compose.ui.input.pointer.PointerEventType +import androidx.compose.ui.input.pointer.pointerInput + +@OptIn(ExperimentalComposeUiApi::class) +actual fun Modifier.rightClick(onRightClick: () -> Unit): Modifier = + pointerInput(Unit) { + awaitPointerEventScope { + while (true) { + val event = awaitPointerEvent() + if (event.type == PointerEventType.Press && event.button == PointerButton.Secondary) { + onRightClick() + } + } + } + } \ No newline at end of file diff --git a/docs/04-wire-protocol.md b/docs/04-wire-protocol.md index 224636b..2767aed 100644 --- a/docs/04-wire-protocol.md +++ b/docs/04-wire-protocol.md @@ -116,6 +116,10 @@ subscribe; the server pushes to every open WS). {"type":"channel.created","payload":{"chat_id":"android:chan_7","name":"Cron Reports", "kind":"channel","parent_chat_id":null}} ``` +`channel.created` may carry `"auto":true` for a thread the gateway minted +itself for an incoming message (auto-threading): the name is an instant +derived title, and a follow-up `channel.renamed` upgrades it to the model's +title. ### `history` Response to a `history` request. Returns a page of messages for a chat/thread. @@ -225,9 +229,16 @@ First frame; auth + caps. Send text (or a `/slash-command`). ```json {"type":"message.send","id":10,"chat_id":"android:default","thread_id":null, - "payload":{"text":"/model qwen3-27b","reply_to":"m_9001","media_refs":["mu_1"]}} + "payload":{"text":"/model qwen3-27b","reply_to":"m_9001","media_refs":["mu_1"], + "auto_thread":false}} ``` `media_refs` reference completed `media.upload`s to attach. +`auto_thread` (optional, default false) asks the gateway to mint a fresh +thread for the message (auto-threading, `06-channels-cron-search.md` §6.3): +honored only in a channel's flat lane (`thread_id` null) with non-empty text +that is not a slash command. The gateway then broadcasts +`channel.created {auto:true}` and the echo / agent turn carry the new +`thread_id`. ### `media.upload.start` / (binary) / `media.upload.end` See `07-media.md`. diff --git a/docs/05-streaming.md b/docs/05-streaming.md index 5d85a64..cd34a84 100644 --- a/docs/05-streaming.md +++ b/docs/05-streaming.md @@ -27,8 +27,15 @@ replace the bubble text (cheap: it's a full snapshot). On `message.stop`, finalize (attach reasoning/model/tokens footer, stop the cursor). Auto-scroll while the user is at the bottom. -**Streaming on/off.** Controlled by hermes `display.platforms.android.streaming` -(default follows global). When off, the app just gets one final `message` frame. +**Streaming on/off.** Two levels: +- **Gateway side:** hermes `display.platforms.android.streaming` (default + follows global). When off, the app just gets one final `message` frame. +- **App side (per device):** Settings → "Streaming" toggle (default on). When + off, the app ignores `message.start`/`message.update` frames and + materializes each reply as a single final message on `message.stop` + (`ChatStore.onMessageStop` already handles a stop without a live bubble). + The gateway keeps streaming — frames are broadcast to all devices, so this + is a display preference like tool verbosity, not a wire flag. ## 5.2 Reasoning (shown *before* the message) diff --git a/docs/06-channels-cron-search.md b/docs/06-channels-cron-search.md index 05e74de..76890e5 100644 --- a/docs/06-channels-cron-search.md +++ b/docs/06-channels-cron-search.md @@ -45,6 +45,34 @@ gateway identity concepts**. - The gateway's `create_handoff_thread` is used where hermes wants to open a named thread (e.g. continuable cron). +### 6.3.1 Auto-threading (Telegram topic-mode workflow) + +With **Threads ON**, a message sent in a channel's flat lane (no active +thread) gets its **own fresh thread**, the way Telegram topic mode mints a +topic per new conversation — and the **AI names it** instead of the user: + +1. The app sends `message.send {…, auto_thread:true}` (only from a flat lane, + with non-empty text, never for slash commands). +2. The gateway mints a thread (`channel.create`-equivalent, `kind:thread`) + under the channel, **named instantly** from the user's opening message via + hermes' session-title derivation (`agent/title_generator.derive_title` — + a deterministic slice of the user's own words, no model call), and + broadcasts `channel.created {auto:true}`. +3. The user echo, the agent turn, and all streaming frames carry the new + `thread_id` — the whole conversation lives in the thread. +4. In the background, the gateway upgrades the name with the model's title + (`agent/title_generator.generate_title`, the `title_generation` auxiliary + task) and broadcasts `channel.renamed`. This is hermes' two-stage session + titling (derived < llm < user) applied to the thread name; failures leave + the derived name in place. + +App side: on `channel.created {auto:true}` under the flat lane it is viewing, +the app **jumps into the new thread** and relocates the optimistic pending +bubble from the flat lane into it (the echo arrives in the thread lane). +Follow-ups sent inside the thread stay there; the next flat-lane message +starts another thread. Media-only sends and slash commands stay in the flat +lane (nothing to title / session-scoped, not conversation starters). + ## 6.4 User-created channels (for cron delegation) - **Requirement:** the user creates new channels so **cron job outputs can be @@ -55,6 +83,9 @@ gateway identity concepts**. - **`channel.rename` / `channel.set_default` / `channel.delete`** manage the directory (rename broadcasts `channel.renamed`; delete is soft — marks archived, keeps history for search). +- **App affordance for threads:** long-press (touch) / right-click (desktop) a + topic chip → "Rename" (`channel.rename`) or "Delete" (`channel.delete`). + Deleting the open thread falls the app back to the channel's flat lane. - **Cron targeting** (the key payoff): because the plugin registers `parse_target_ref_fn` and `cron_deliver_env_var`, cron jobs and the `send_message` tool can target any channel/thread: diff --git a/docs/10-android-app.md b/docs/10-android-app.md index b0c4c97..a3ba7d2 100644 --- a/docs/10-android-app.md +++ b/docs/10-android-app.md @@ -93,18 +93,35 @@ app/shared/src/ clarifies) render as **native pickers** from `picker.*` frames (a dialog / sheet with the options; answer via `picker.select`). +### Streaming (app-controlled) +- **Settings → "Streaming"** toggle (default on). When off, the app ignores + `message.start`/`message.update` frames and shows each reply as a single + final message on `message.stop` (the typing indicator covers the wait). + Per-device display preference — the gateway keeps streaming for other + devices (see `05-streaming.md` §5.1). + ### Tool output (app-controlled verbosity) - `ToolCard` renders `tool.start/progress/end` frames. +- The gateway always supplies the **full** tool data: it forces + `display.platforms.android.tool_progress: verbose` (so the progress line + carries the full args JSON → `tool.start.args`) and captures each completed + call via the `post_tool_call` hook (→ `tool.end` `output_preview` / + `duration` / `ok`). The app decides how much to show. - **Settings → "Tool detail": Everything / Truncated / Nothing.** - - Everything: name + full args (collapsible) + output preview. - - Truncated (default): `emoji name: "short preview"` one-liner, expandable. - - Nothing: suppress tool frames. + - Everything: cards start expanded — name + full args + output. + - Truncated (default): compact one-liner; tap a card to reveal the full + args + output (what was actually called). + - Nothing: suppress tool cards. - Spinner while running; ✓/✗ + duration on `tool.end`. ### Reasoning before message - `ReasoningBlock` (collapsible, "💭 Reasoning" header, monospace body, **copy** button) rendered **above** the message body from the `reasoning` field. Matches the reference screenshot. +- **Settings → "Reasoning"** toggle (default on): auto-collapse — long + reasoning blocks start collapsed, short ones expanded. Off: all reasoning + blocks start expanded. Per-device display preference; the per-block + tap-to-toggle always works. ### Intermediate messages - `commentary` frames → dimmed/smaller bubble, distinct from final answers. @@ -122,8 +139,20 @@ app/shared/src/ highlight bar) — matches the reference left sidebar. - **Thread toggle** in the default chat header: "Threads on/off". On → topic switcher above the message list (each topic = a `thread_id`). +- **Auto-threading** (Threads on, Telegram topic-mode workflow, + `06-channels-cron-search.md` §6.3.1): a message sent in a channel's flat + lane goes out as `message.send {auto_thread:true}`; the gateway mints a + fresh thread (instant derived name, AI-named a moment later via + `channel.renamed`) and the app jumps into it, relocating the optimistic + bubble from the flat lane. - **New channel** (FAB / channel-list menu) → `channel.create` → appears in list; menu offers "Set as cron target". +- **Topic context menu** — long-press a topic chip (touch) or right-click it + (desktop mouse) → "Rename" / "Delete". Rename → `channel.rename` (prefilled + dialog); Delete → confirm → `channel.delete` (soft-delete; the thread leaves + the switcher and, if it was the open lane, the app falls back to the channel's + flat lane). The right-click handler is a skiko `expect`/`actual` + (`iris/ui/ContextMenu.kt`); on touch it is a no-op (long-press covers it). ### Search - Search bar (chat header or top) with a **scope toggle**: "Search everywhere" / diff --git a/docs/13-testing.md b/docs/13-testing.md index fd85be4..0dad897 100644 --- a/docs/13-testing.md +++ b/docs/13-testing.md @@ -32,6 +32,11 @@ without the app (critical for verifying frame shapes early). - **Push backend selection:** fcm vs ntfy chosen by config; `configured()` reflects missing creds. - **Channel directory:** create/rename/set_default; cron target resolution. + - **Auto-threading:** `message.send {auto_thread:true}` in a flat lane mints a + thread (derived name, `channel.created {auto:true}`) and the echo/event + carry the new `thread_id`; no-op with an existing `thread_id`, for slash + commands, or for media-only sends; the LLM upgrade renames the thread + (`channel.renamed`). - **No `~/.hermes` writes in tests** — use the `_isolate_hermes_home` fixture pattern (temp `HERMES_HOME`). Profile tests also mock `Path.home()`. @@ -113,6 +118,10 @@ adb logcat -d > /tmp/logcat.txt trigger a message → FCM notification appears → tap → syncs + deep-links. 12. **Reconnect/sync:** kill the WS (stop gateway briefly) → restart → app reconnects → `sync` catches up (no lost/dup messages). +13. **Auto-threading:** Threads on → send a message in the Default channel's + flat lane → a new topic appears (derived name), the app jumps into it, the + reply streams there, and the topic is renamed to the AI's title a moment + later. Slash commands / media-only sends stay in the flat lane. ## 13.5 Debugging tips diff --git a/docs/protocol/frames.schema.json b/docs/protocol/frames.schema.json index 4c1c446..5214d72 100644 --- a/docs/protocol/frames.schema.json +++ b/docs/protocol/frames.schema.json @@ -49,7 +49,7 @@ "typing": { "payload": { "on": { "type": "boolean" } } }, "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.created": { "description": "May carry auto:true for a thread the gateway minted itself for an incoming message (auto-threading); the name is a derived title, upgraded by a follow-up channel.renamed.", "payload": { "$ref": "#/definitions/channel" } }, "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" } } }, "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"} } } } } }, @@ -59,12 +59,13 @@ "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" } } }, + "history": { "description": "Paged full message history for a chat/thread (response to a history request). Reconstructed from the outbox log; used to populate the view on first open / after a process death, since sync only replays the outbox delta.", "payload": { "messages": { "type": "array", "items": { "type": "object", "properties": { "message_id": {"type":"string"}, "role": {"type":"string","enum":["user","assistant"]}, "text": {"type":"string"}, "reasoning": {"type":"string"}, "model": {"type":"string"}, "tokens": {"type":"integer"}, "ts": {"type":"integer"}, "media": {"type":"array","items":{"$ref":"#/definitions/media_ref"}} } } }, "has_more": { "type": "boolean", "description": "True when older pages exist." }, "oldest_message_id": { "type": "string", "description": "before_message_id for the next (older) page." } } } }, "media.pull.end": { "payload": { "ok": { "type": "boolean" } } }, "media.upload.ack": { "description": "Response to media.upload.end; ref is cached and usable in message.send media_refs.", "payload": { "ok": { "type": "boolean" }, "media_ref": { "type": "string" } } } }, "app_to_server": { "hello": { "description": "First frame; auth + caps.", "payload": { "token": { "type": "string" }, "device_id": { "type": "string" }, "device_name": { "type": "string" }, "caps": { "type": "object", "properties": { "min_protocol": {"type":"integer"}, "media": {"type":"boolean"}, "push": {"type":"string"} } }, "fcm_token": { "type": "string" }, "ntfy_topic": { "type": "string" } } }, - "message.send": { "payload": { "text": { "type": "string" }, "reply_to": { "type": "string" }, "media_refs": { "type": "array", "items": { "type": "string" } } } }, + "message.send": { "payload": { "text": { "type": "string" }, "reply_to": { "type": "string" }, "media_refs": { "type": "array", "items": { "type": "string" } }, "auto_thread": { "type": "boolean", "description": "Optional, default false. Ask the gateway to mint a fresh thread for this message (auto-threading). Honored only in a channel's flat lane (thread_id null) with non-empty non-slash text; the gateway broadcasts channel.created {auto:true} and the echo/turn carry the new thread_id." } } }, "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" } } }, @@ -75,13 +76,14 @@ "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" } } }, + "history": { "description": "Load a page of full message history for a chat/thread (initial open, scroll-up pagination).", "payload": { "before_message_id": { "type": "string", "description": "Return messages older than this (omit for newest page)." }, "limit": { "type": "integer", "description": "Max messages (default 50, max 200)." } } }, "fcm.register": { "payload": { "fcm_token": { "type": "string" }, "ntfy_topic": { "type": "string" } } }, "ping": { "payload": { "ts": { "type": "integer" } } } } }, "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"}, "archived": {"type":"boolean"} } }, + "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"}, "auto": {"type":"boolean","description":"Optional; true on channel.created for a gateway-minted auto-thread."} } }, "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": [ @@ -91,8 +93,6 @@ { "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." }, diff --git a/gateway-plugin/adapter.py b/gateway-plugin/adapter.py index 68e077c..7cc0232 100644 --- a/gateway-plugin/adapter.py +++ b/gateway-plugin/adapter.py @@ -58,12 +58,14 @@ Or via environment variables (overrides config.yaml; secrets live in .env): """ import asyncio +import json import logging import os import re import threading import time import uuid +from collections import deque from dataclasses import dataclass, field from typing import Any, Dict, List, Optional, Tuple @@ -185,6 +187,75 @@ def _reset_reasoning() -> None: _reasoning_flushed.clear() +# --------------------------------------------------------------------------- +# M2 — tool-result capture (post_tool_call hook) +# +# The gateway renders tool *progress* lines to the platform but never streams +# the tool *output* (it is the agent's concern, persisted to history, not +# presentation). To let the app show the full call + result on demand +# (Settings → Tool detail), we capture each completed tool call via the +# ``post_tool_call`` hook and attach it to the ``tool.end`` frame. +# +# Global FIFO (like the reasoning buffer): a personal android gateway serves +# one active turn at a time, and records are matched to the open tool by name +# in completion order. Bounded so a runaway turn can't grow it without limit. +# --------------------------------------------------------------------------- + +_tool_results: "deque[Dict[str, Any]]" = deque() +_tool_results_lock = threading.Lock() +_MAX_TOOL_RESULTS = 200 +_MAX_OUTPUT_PREVIEW = 8000 + + +def _on_post_tool_call(**kwargs: Any) -> None: + """Plugin hook: capture a completed tool call's result + timing.""" + result = kwargs.get("result") + record = { + "tool_name": kwargs.get("tool_name") or "", + "result": (str(result) if result is not None else "")[:_MAX_OUTPUT_PREVIEW], + "duration_ms": kwargs.get("duration_ms") or 0, + "status": kwargs.get("status") or "ok", + } + with _tool_results_lock: + _tool_results.append(record) + while len(_tool_results) > _MAX_TOOL_RESULTS: + _tool_results.popleft() + + +def _take_tool_result(tool_name: str) -> Optional[Dict[str, Any]]: + """Pop the first completed record matching *tool_name* (FIFO), else None.""" + if not tool_name: + return None + with _tool_results_lock: + for i, rec in enumerate(_tool_results): + if rec["tool_name"] == tool_name: + del _tool_results[i] + return rec + return None + + +def _reset_tool_results() -> None: + """Clear captured records (turn boundary — drop anything unconsumed).""" + with _tool_results_lock: + _tool_results.clear() + + +def _tool_end_fields(tool_name: str) -> Dict[str, Any]: + """Build the ``tool.end`` enrichment (ok/duration/output_preview) from the + captured hook record for *tool_name*; empty dict when none is available + (e.g. tool_progress off, or the call came from another session).""" + rec = _take_tool_result(tool_name) + if rec is None: + return {} + fields: Dict[str, Any] = { + "ok": rec["status"] == "ok", + "output_preview": rec["result"] or None, + } + if rec["duration_ms"]: + fields["duration"] = round(rec["duration_ms"] / 1000.0, 3) + return fields + + # --------------------------------------------------------------------------- # Defaults # --------------------------------------------------------------------------- @@ -283,6 +354,25 @@ def _thread_id_from_metadata(metadata: Optional[Dict[str, Any]]) -> Optional[str return None +def _derive_thread_name(text: str) -> str: + """Instant auto-thread name from the user's opening message (no model). + + Reuses hermes' session-title derivation (``agent/title_generator.py``): + a deterministic slice of the user's own words, so the thread is named the + moment it is created. The LLM upgrade (``_schedule_thread_title_upgrade``) + replaces it moments later — the same two-stage titling hermes uses for + sessions (derived < llm < user). + """ + try: + from agent.title_generator import derive_title + + title = derive_title(text) + except Exception: + logger.debug("Thread name derivation failed", exc_info=True) + title = None + return (title or "").strip() or "New thread" + + def _strip_streaming_cursor(text: str) -> str: if text and text.endswith(_STREAMING_CURSOR): return text[: -len(_STREAMING_CURSOR)] @@ -406,6 +496,50 @@ def _extract_code_block(content: str) -> Optional[str]: return None +def _extract_verbose_args(line: str, content: str) -> Optional[Dict[str, Any]]: + """Recover the full args dict from a verbose tool line, else ``None``. + + In verbose mode the gateway renders `` (keys)`` on one line + and the full args JSON on the line that follows. When *line* is such a + header, return the parsed JSON object from the following line. + """ + parts = line.strip().split(None, 1) + if len(parts) < 2 or not _TOOL_NAME_ARGS_RE.match(parts[1]): + return None + lines = content.splitlines() + for i, ln in enumerate(lines): + if ln.strip() != line.strip(): + continue + for follow in lines[i + 1:]: + follow = follow.strip() + if not follow: + continue + if follow.startswith("{"): + try: + obj = json.loads(follow) + return obj if isinstance(obj, dict) else None + except Exception: + return None + return None # next non-empty line is not the args JSON + return None + + +def _short_preview_from_args(args: Dict[str, Any], cap: int = 60) -> Optional[str]: + """Derive a short one-line preview from a verbose args dict. + + The verbose line carries no explicit preview, so the Truncated display + would otherwise lose its one-liner. Use the first non-empty string value + (whitespace-collapsed, capped) as a stand-in. + """ + if not isinstance(args, dict): + return None + for value in args.values(): + if isinstance(value, str) and value.strip(): + s = " ".join(value.split()) + return s[: cap - 1] + "…" if len(s) > cap else s + return None + + def _is_tool_progress(content: str) -> bool: """Heuristic: does *content* look like gateway tool-progress line(s)? @@ -439,6 +573,9 @@ class _TurnState: tool_index: int = 0 # Index of the most recently started tool (awaiting tool.end). open_tool_index: Optional[int] = None + # Name of the most recently started tool (matches the post_tool_call + # record when the tool completes, so tool.end can carry its output). + open_tool_name: Optional[str] = None # Tool lines already emitted as tool.start (dedup across edits). seen_tool_lines: set = field(default_factory=set) @@ -601,6 +738,46 @@ async def _standalone_send( } +# --------------------------------------------------------------------------- +# Verbose tool progress (full args on the progress line) +# --------------------------------------------------------------------------- + +def _ensure_verbose_tool_progress() -> None: + """Ensure the android platform renders tool progress in ``verbose`` mode. + + Verbose mode makes the gateway's tool-progress line carry the FULL + argument JSON (not just a ~40-char preview), which the adapter parses + into the ``tool.start`` frame's ``args`` field; the app then decides how + much to show (Settings → Tool detail). The tool *output* is captured + separately via the ``post_tool_call`` hook (verbose mode does not stream + it). + + Best-effort and idempotent: writes + ``display.platforms.android.tool_progress: verbose`` to config.yaml only + when it isn't already set. The gateway's config cache is mtime-keyed, so + the write takes effect on the next turn without a restart. Never raises. + """ + try: + from hermes_cli.config import load_config_readonly + + cfg = load_config_readonly() or {} + display = cfg.get("display") or {} + platforms = display.get("platforms") or {} + android = platforms.get("android") or {} + if android.get("tool_progress") == "verbose": + return # already set + from utils import atomic_roundtrip_yaml_update + + atomic_roundtrip_yaml_update( + get_hermes_home() / "config.yaml", + "display.platforms.android.tool_progress", + "verbose", + ) + logger.info("android: set display.platforms.android.tool_progress=verbose") + except Exception: + logger.debug("android: could not ensure verbose tool_progress", exc_info=True) + + # --------------------------------------------------------------------------- # Interactive setup (hermes gateway setup flow) # --------------------------------------------------------------------------- @@ -654,6 +831,10 @@ def interactive_setup() -> None: print_info(f"Pairing URL: {qr_payload(host or DEFAULT_HOST, _parse_port(port), token)}") print_info(f"Server URL: {url}") + # Always render tool progress verbosely so the app receives the full tool + # call args (it decides how much to show via Settings → Tool detail). + _ensure_verbose_tool_progress() + print_success("Android configuration saved to ~/.hermes/.env") print_info("Restart the gateway for changes to take effect: hermes gateway restart") @@ -675,6 +856,10 @@ class AndroidAdapter(BasePlatformAdapter): platform = Platform("android") super().__init__(config=config, platform=platform) + # Ensure verbose tool progress (full args on the progress line) so the + # app can show the full tool call on demand. Best-effort; idempotent. + _ensure_verbose_tool_progress() + extra = getattr(config, "extra", {}) or {} # Connection settings (env vars override config.yaml) @@ -1094,38 +1279,59 @@ class AndroidAdapter(BasePlatformAdapter): parsed = self._parse_tool_line_or_block(line, content) if parsed is None: continue - name, preview = parsed - # A new tool begins: close the previously-open one. + name, preview, args = parsed + # A new tool begins: close the previously-open one, attaching the + # output/duration/ok captured by the post_tool_call hook. if state.open_tool_index is not None: + extra = _tool_end_fields(state.open_tool_name or "") await self._broadcast_or_log( chat_id, - protocol.tool_end(chat_id, state.open_tool_index, "", ok=True, thread_id=thread_id), + protocol.tool_end( + chat_id, state.open_tool_index, state.open_tool_name or "", + ok=extra.get("ok", True), + duration=extra.get("duration"), + output_preview=extra.get("output_preview"), + thread_id=thread_id, + ), ) state.tool_index += 1 state.open_tool_index = state.tool_index + state.open_tool_name = name await self._broadcast_or_log( chat_id, protocol.tool_start( chat_id, state.tool_index, name, - preview=preview, thread_id=thread_id, + preview=preview, args=args, thread_id=thread_id, ), ) return SendResult(success=True, message_id=message_id) @staticmethod - def _parse_tool_line_or_block(line: str, content: str) -> Optional[Tuple[str, Optional[str]]]: - """Parse a tool line, expanding a terminal code block to its command.""" + def _parse_tool_line_or_block( + line: str, content: str + ) -> Optional[Tuple[str, Optional[str], Optional[Dict[str, Any]]]]: + """Parse a tool line into ``(name, preview, args)``. + + Expands a terminal code block to its command, and a verbose header + (`` (keys)``) to its full args JSON (the JSON sits on the + following line). ``args`` is ``None`` unless the line is a verbose + header with a parseable JSON body. + """ parsed = _parse_tool_line(line) - if parsed is not None: - name, preview = parsed - # Terminal code block: the command lives in the fenced lines that - # follow the " terminal" head line. - if name == "terminal" and preview is None and "```" in content: - cmd = _extract_code_block(content) - if cmd: - return name, cmd - return parsed - return None + if parsed is None: + return None + name, preview = parsed + # Terminal code block: the command lives in the fenced lines that + # follow the " terminal" head line. + if name == "terminal" and preview is None and "```" in content: + cmd = _extract_code_block(content) + if cmd: + return name, cmd, None + # Verbose mode: recover the full args from the following JSON line. + args = _extract_verbose_args(line, content) + if args is not None and preview is None: + preview = _short_preview_from_args(args) + return name, preview, args async def _close_open_tool( self, chat_id: str, state: _TurnState, thread_id: Optional[str] @@ -1137,11 +1343,19 @@ class AndroidAdapter(BasePlatformAdapter): tool it was waiting on has returned). """ if state.open_tool_index is not None: + extra = _tool_end_fields(state.open_tool_name or "") await self._broadcast_or_log( chat_id, - protocol.tool_end(chat_id, state.open_tool_index, "", ok=True, thread_id=thread_id), + protocol.tool_end( + chat_id, state.open_tool_index, state.open_tool_name or "", + ok=extra.get("ok", True), + duration=extra.get("duration"), + output_preview=extra.get("output_preview"), + thread_id=thread_id, + ), ) state.open_tool_index = None + state.open_tool_name = None def _reset_tool_state(self, state: _TurnState) -> None: """Clear per-turn tool bookkeeping (called at turn finalization).""" @@ -1149,6 +1363,10 @@ class AndroidAdapter(BasePlatformAdapter): state.seen_tool_lines = set() state.tool_index = 0 state.open_tool_index = None + state.open_tool_name = None + # Drop any captured tool results not consumed by a tool.end this turn + # (e.g. tool_progress off) so they can't leak into the next turn. + _reset_tool_results() async def _broadcast_or_log(self, chat_id: str, frame: "protocol.Frame") -> None: delivered = await self._ws_server.broadcast(frame) @@ -1465,6 +1683,33 @@ class AndroidAdapter(BasePlatformAdapter): if not isinstance(reply_to, str) or not reply_to.strip(): reply_to = None + # Auto-threading (the app's Threads setting, docs/06 §6.3): a message + # in a channel's flat lane gets its own fresh thread, the way Telegram + # topic mode mints a topic per new conversation. The thread is named + # instantly from the user's opening message (derived title) and the + # LLM upgrades the name in the background. The user echo, the agent + # turn, and all streaming frames then carry the new thread_id. + # Skipped for slash commands (session-scoped, not conversation + # starters) and replies (they continue where the user is). + auto_thread = payload.get("auto_thread") is True + if ( + auto_thread + and thread_id is None + and text.strip() + and not text.lstrip().startswith("/") + and reply_to is None + ): + entry = self._channels.create( + name=_derive_thread_name(text), + kind="thread", + parent_chat_id=chat_id, + ) + thread_id = entry["chat_id"] + # Bare broadcast (like channel.create): the directory is + # re-served on hello.ack, so no outbox entry is needed. + await self._ws_server.broadcast(protocol.channel_created(entry, auto=True)) + self._schedule_thread_title_upgrade(entry["chat_id"], text) + # M4: resolve media refs (single-use; unknown ref -> error). media_urls: List[str] = [] media_types: List[str] = [] @@ -1494,6 +1739,12 @@ class AndroidAdapter(BasePlatformAdapter): # Echo to all devices: the sender confirms (server-assigned id), # other devices see the message too (single-user, multi-device). + # Routed through _broadcast_or_log (not a bare broadcast) so the echo + # is appended to the outbox: the app's ChatStore is in-memory only, so + # after a process death / activity recreation the only way the user's + # own message is restored is via the sync replay. Without this, user + # messages vanish on reconnect while bot messages (already parked) + # survive. message_id = f"m_{uuid.uuid4().hex[:16]}" echo = protocol.message( chat_id=chat_id, @@ -1505,7 +1756,7 @@ class AndroidAdapter(BasePlatformAdapter): reply_to=reply_to, ts=int(time.time() * 1000), ) - await self._ws_server.broadcast(echo) + await self._broadcast_or_log(chat_id, echo) # Refs are consumed by this message (no replay). for ref in media_refs: self._media.pop_inbound(ref) @@ -1554,6 +1805,48 @@ class AndroidAdapter(BasePlatformAdapter): protocol.read_receipt(chat_id, message_id), ) + def _schedule_thread_title_upgrade(self, thread_id: str, text: str) -> None: + """Upgrade an auto-created thread's name with the model's title. + + Stage 2 of hermes' two-stage session titling (``agent/title_generator + .py``): the thread was created with an instant derived name; this + background call on the ``title_generation`` auxiliary task replaces it + with the model's title and broadcasts ``channel.renamed``. Best-effort + — any failure (config, model, network) leaves the derived name in + place, and a thread the user already renamed or archived is untouched. + """ + loop = asyncio.get_running_loop() + + def _work() -> None: + try: + from agent.title_generator import generate_title + + title = generate_title(text) + except Exception: + logger.debug("Thread title upgrade failed", exc_info=True) + return + if not title: + return + entry = self._channels.get(thread_id) + if entry is None or entry.get("archived"): + return + if (entry.get("name") or "") == title: + return + renamed = self._channels.rename(thread_id, title) + if renamed is None: + return + try: + asyncio.run_coroutine_threadsafe( + self._ws_server.broadcast(protocol.channel_renamed(renamed)), + loop, + ) + except Exception: + logger.debug("Thread title rename broadcast failed", exc_info=True) + + threading.Thread( + target=_work, daemon=True, name="android-thread-title" + ).start() + # ── M4: inbound media (app -> agent) ────────────────────────────────── # # ``media.upload.start`` -> raw binary frames (one at a time per @@ -1875,6 +2168,54 @@ class AndroidAdapter(BasePlatformAdapter): done = protocol.sync_done(self._outbox.latest_cursor(), id=frame.id) await self._ws_server.send_to(device_id, done) + # ── Full message history (initial channel open / scroll-up) ─────────── + + async def on_history(self, frame: protocol.Frame, device_id: str) -> None: + """Handle an inbound ``history`` request. + + ``sync`` only replays the outbox delta since the device's cursor, so + after a process death the app's in-memory ChatStore is empty and the + delta does not cover older messages. ``history`` loads the full + message list for a chat/thread (reconstructed from the outbox log) so + the app can populate the view on first open / restart. + """ + payload = frame.payload + chat_id = frame.chat_id or payload.get("chat_id") + logger.info("android: history request from %s chat_id=%r", device_id, chat_id) + if not isinstance(chat_id, str) or not chat_id.strip(): + await self._ws_server.send_to( + device_id, + protocol.error(protocol.ERR_UNSUPPORTED, "history requires a chat_id", id=frame.id), + ) + return + chat_id = chat_id.strip() + thread_id = frame.thread_id or payload.get("thread_id") + if not isinstance(thread_id, str) or not thread_id.strip(): + thread_id = None + before = payload.get("before_message_id") + if not isinstance(before, str) or not before.strip(): + before = None + limit_raw = payload.get("limit") + try: + limit = int(limit_raw) if limit_raw is not None else 50 + except (TypeError, ValueError): + limit = 50 + page = self._outbox.history( + chat_id, + thread_id=thread_id, + before_message_id=before, + limit=limit, + ) + resp = protocol.history( + chat_id, + page["messages"], + page["has_more"], + thread_id=thread_id, + oldest_message_id=page["oldest_message_id"], + id=frame.id, + ) + await self._ws_server.send_to(device_id, resp) + # ── M5: push token registration ─────────────────────────────────────── async def on_fcm_register(self, frame: protocol.Frame, device_id: str) -> None: @@ -2123,6 +2464,13 @@ def register(ctx): ctx.register_hook("on_stream_delta", _on_stream_delta) except Exception: logger.debug("android: on_stream_delta hook registration failed", exc_info=True) + # M2: capture each completed tool call's result + timing so the tool.end + # frame can carry the output (the gateway never streams tool output to + # platforms). The app shows it on demand (Settings → Tool detail). + try: + ctx.register_hook("post_tool_call", _on_post_tool_call) + except Exception: + logger.debug("android: post_tool_call hook registration failed", exc_info=True) ctx.register_platform( name="android", label="Android", diff --git a/gateway-plugin/outbox.py b/gateway-plugin/outbox.py index b3c41be..0770b29 100644 --- a/gateway-plugin/outbox.py +++ b/gateway-plugin/outbox.py @@ -169,6 +169,114 @@ class Outbox: ) return out + # ── history (full message history for a chat/thread) ────────────────── + + def history( + self, + chat_id: str, + thread_id: Optional[str] = None, + before_message_id: Optional[str] = None, + limit: int = 50, + ) -> Dict[str, Any]: + """Final messages for a chat/thread, for the ``history`` frame. + + Reconstructs the message list from the outbox log: a final message is + either a standalone ``message`` frame (user echo / non-streaming + assistant) or a ``message.stop`` frame (streaming assistant final). + Intermediate frames (``message.start``/``message.update``, tool, + commentary, notification, …) are skipped. Messages are deduplicated by + ``message_id`` and returned oldest → newest. + + Pagination: ``before_message_id`` returns the page of messages older + than that id (the newest page when omitted). Returns + ``{messages, has_more, oldest_message_id}``. + """ + limit = max(1, min(int(limit or 50), 200)) + with self._lock: + rows = self._conn.execute( + "SELECT cursor, frame FROM outbox WHERE chat_id = ? " + "ORDER BY cursor ASC", + (chat_id,), + ).fetchall() + final: List[Dict[str, Any]] = [] + for r in rows: + try: + frame = json.loads(r["frame"]) + except (json.JSONDecodeError, TypeError): + continue + if not isinstance(frame, dict): + continue + if thread_id is not None and frame.get("thread_id") != thread_id: + continue + ftype = frame.get("type") + payload = frame.get("payload") + if not isinstance(payload, dict): + continue + if ftype == "message": + role = payload.get("role") + if role not in ("user", "assistant"): + continue + final.append( + { + "cursor": int(r["cursor"]), + "message_id": payload.get("message_id"), + "role": role, + "text": payload.get("text", ""), + "reasoning": payload.get("reasoning"), + "model": payload.get("model"), + "tokens": payload.get("tokens"), + "ts": payload.get("ts"), + "media": payload.get("media"), + } + ) + elif ftype == "message.stop": + final.append( + { + "cursor": int(r["cursor"]), + "message_id": payload.get("message_id"), + "role": "assistant", + "text": payload.get("final_text", ""), + "reasoning": payload.get("reasoning"), + "model": payload.get("model"), + "tokens": payload.get("tokens"), + "ts": payload.get("ts"), + "media": None, + } + ) + # Deduplicate by message_id (keep the latest occurrence), keep order. + by_id: Dict[str, Dict[str, Any]] = {} + for m in final: + mid = m.get("message_id") + if mid: + by_id[mid] = m + ordered = sorted(by_id.values(), key=lambda m: m["cursor"]) + # Paginate: messages older than ``before_message_id`` (newest page when + # omitted / not found — a pruned anchor falls back to the newest page). + if before_message_id: + idx = next( + (i for i, m in enumerate(ordered) if m["message_id"] == before_message_id), + None, + ) + pool = ordered[:idx] if idx is not None else ordered + else: + pool = ordered + has_more = len(pool) > limit + page = pool[-limit:] if has_more else list(pool) + oldest_message_id = page[0]["message_id"] if page else None + for m in page: + m.pop("cursor", None) + # Omit absent optional fields (the app's serializer treats a + # missing key as its default, but a JSON ``null`` for a + # non-nullable field like ``media`` would fail to parse). + for key in ("reasoning", "model", "tokens", "ts", "media"): + if m.get(key) is None: + m.pop(key, None) + return { + "messages": page, + "has_more": has_more, + "oldest_message_id": oldest_message_id, + } + # ── retention ───────────────────────────────────────────────────────── def _maybe_prune(self) -> None: diff --git a/gateway-plugin/protocol.py b/gateway-plugin/protocol.py index 1f94045..5093fac 100644 --- a/gateway-plugin/protocol.py +++ b/gateway-plugin/protocol.py @@ -68,6 +68,9 @@ TYPE_SEARCH_RESULTS = "search.results" TYPE_SYNC = "sync" TYPE_SYNC_DONE = "sync.done" +# Full message history (initial channel open / scroll-up pagination) +TYPE_HISTORY = "history" + # Media (M4) TYPE_MEDIA_UPLOAD_START = "media.upload.start" TYPE_MEDIA_UPLOAD_END = "media.upload.end" @@ -430,9 +433,18 @@ def _channel_payload(entry: Dict[str, Any]) -> Dict[str, Any]: return payload -def channel_created(entry: Dict[str, Any]) -> Frame: - """Broadcast: a channel/thread was created.""" - return Frame(type=TYPE_CHANNEL_CREATED, payload=_channel_payload(entry)) +def channel_created(entry: Dict[str, Any], auto: bool = False) -> Frame: + """Broadcast: a channel/thread was created. + + ``auto=True`` marks a thread the gateway minted itself for an incoming + message (auto-threading, docs/06 §6.3): the app jumps into it and the + name is an instant derived title, upgraded by the LLM via a follow-up + ``channel.renamed``. + """ + payload = _channel_payload(entry) + if auto: + payload["auto"] = True + return Frame(type=TYPE_CHANNEL_CREATED, payload=payload) def channel_renamed(entry: Dict[str, Any]) -> Frame: @@ -484,6 +496,38 @@ def sync_done(cursor: int, *, id: Optional[int] = None) -> Frame: return Frame(type=TYPE_SYNC_DONE, id=id, payload={"cursor": cursor}) +# --------------------------------------------------------------------------- +# History frame (full message history for a chat/thread) +# --------------------------------------------------------------------------- + +def history( + chat_id: str, + messages: List[Dict[str, Any]], + has_more: bool, + *, + thread_id: Optional[str] = None, + oldest_message_id: Optional[str] = None, + id: Optional[int] = None, +) -> Frame: + """Response to a ``history`` request: a page of final messages for a + chat/thread, ordered oldest → newest. ``has_more`` signals older pages + exist; ``oldest_message_id`` is the ``before_message_id`` for the next + (older) page.""" + payload: Dict[str, Any] = { + "messages": messages, + "has_more": has_more, + } + if oldest_message_id: + payload["oldest_message_id"] = oldest_message_id + return Frame( + type=TYPE_HISTORY, + id=id, + chat_id=chat_id, + thread_id=thread_id, + payload=payload, + ) + + # --------------------------------------------------------------------------- # Push / notification frames (M5) # --------------------------------------------------------------------------- diff --git a/gateway-plugin/ws_server.py b/gateway-plugin/ws_server.py index e065649..88485a2 100644 --- a/gateway-plugin/ws_server.py +++ b/gateway-plugin/ws_server.py @@ -378,6 +378,8 @@ class WsServer: await self._adapter.on_search(frame, device_id) elif frame.type == protocol.TYPE_SYNC: await self._adapter.on_sync(frame, device_id) + elif frame.type == protocol.TYPE_HISTORY: + await self._adapter.on_history(frame, device_id) elif frame.type == protocol.TYPE_MEDIA_UPLOAD_START: await self._adapter.on_media_upload_start(frame, device_id) elif frame.type == protocol.TYPE_MEDIA_UPLOAD_END: