Auto-threading, history pagination, streaming toggle + tool/reasoning display settings

- Auto-threading (Telegram topic-mode workflow): message.send {auto_thread}
  mints a fresh AI-named thread (instant derived title, LLM upgrade via
  channel.renamed); channel.created {auto:true}; the app jumps into the new
  thread and relocates the optimistic pending bubble.
- history frame: paged full message history for initial channel open /
  scroll-up pagination (reconstructed from the outbox log).
- Streaming on/off: gateway side (display.platforms.android.streaming) plus a
  per-device app toggle (Settings → Streaming); reasoning/model/tokens carried
  on message frames.
- Context menu: long-press (touch) / right-click (desktop) thread affordances
  via a KMP rightClick expect/actual.
- Settings → Reasoning: auto-collapse long reasoning blocks (default on).
- Tool detail: the gateway now always supplies full tool data — it forces
  verbose tool progress (full args → tool.start.args) and captures each
  completed call via the post_tool_call hook (output/duration/ok → tool.end).
  The app reveals the full call + output on expand (Truncated) and
  auto-expands cards in Everything mode.
This commit is contained in:
ARIA committed 2026-08-20 16:26:52 +02:00
1 parent e678caafdc
commit efecf2732e
24 files changed
+1109 -75

No files matched your search

@@ -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"
}
}
@@ -0,0 +1,5 @@
package iris.ui
import androidx.compose.ui.Modifier
actual fun Modifier.rightClick(onRightClick: () -> Unit): Modifier = this
@@ -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)
}
}
@@ -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<MessagePayload>() ?: 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<MediaItem>, incoming: List<MediaRef>): List<MediaItem> {
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<MessageStartPayload>() ?: 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<MessageItem>) {
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()
}
@@ -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()
}
@@ -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<Int, CompletableDeferred<Frame>>()
@@ -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<String> = 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 ───────────────────────────────────────────
@@ -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<String> = 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<MediaRef> = emptyList(),
)
@Serializable
data class HistoryPayload(
val messages: List<HistoryMessage> = 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<String> = 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(
@@ -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<String> = _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<String>()
/** Gateway health state (M5: status frame; null = never received). */
private val _gatewayStatus = MutableStateFlow<String?>(null)
val gatewayStatus: StateFlow<String?> = _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<Boolean> = _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<Boolean> = _reasoningAutoCollapse.asStateFlow()
fun toggleReasoningAutoCollapse() {
_reasoningAutoCollapse.value = !_reasoningAutoCollapse.value
store.reasoningAutoCollapse = _reasoningAutoCollapse.value
}
// ── M3: search state ──────────────────────────────────────────────────
private val _searchResults = MutableStateFlow<List<SearchHit>>(emptyList())
val searchResults: StateFlow<List<SearchHit>> = _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<ChannelInfo>()?.let { info ->
if (info.kind == "thread" && info.auto && info.parentChatId != null) {
val (curChat, curThread) = chat.parseLane(chat.currentLane.value)
if (curThread == null && curChat == info.parentChatId) {
val lane = chat.laneKey(info.parentChatId, info.chatId)
chat.setLane(lane)
historyLoaded.add(lane)
}
}
}
}
// A thread/channel we were viewing got deleted — fall back
// to a valid lane (the thread's parent, or the home channel).
if (frame.type == TYPE_CHANNEL_DELETED) {
frame.payloadAs<ChannelDeletedPayload>()?.let { p ->
val (curChat, curThread) = chat.parseLane(chat.currentLane.value)
when {
curThread == p.chatId -> openChannel(curChat)
curChat == p.chatId ->
channels.defaultChannel()?.let { openChannel(it.chatId) }
else -> Unit
}
}
}
}
TYPE_SEARCH_RESULTS -> {
frame.payloadAs<SearchResultsPayload>()?.let { _searchResults.value = it.hits }
_searching.value = false
@@ -265,6 +333,22 @@ class IrisController(
// cursor is authoritative server-side (outbox).
frame.payloadAs<SyncDonePayload>()?.let { store.syncCursor = it.cursor }
}
TYPE_HISTORY -> {
// Full message history for a chat/thread (initial open /
// restart restore). `sync` only replays the outbox delta
// since the saved cursor, so after a process death the
// in-memory store is empty and older messages are not in
// the delta — history loads the full list.
frame.payloadAs<HistoryPayload>()?.let { p ->
val (chatId, threadId) =
(frame.chatId ?: return@let) to frame.threadId
val lane = chat.laneKey(chatId, threadId)
chat.loadHistory(
lane,
p.messages.map { it.toMessageItem() },
)
}
}
TYPE_NOTIFICATION -> {
frame.payloadAs<NotificationPayload>()?.let { p ->
pushBanner(p.kind, p.title, p.body, p.chatId, p.threadId)
@@ -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<PendingAttachment> = 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
@@ -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
@@ -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<ChannelInfo?>(null) }
var deleteThread by remember { mutableStateOf<ChannelInfo?>(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(
@@ -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,
)
}
}
}
@@ -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",
@@ -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()))
@@ -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()
}
}
}
}
+12 -1
View File
@@ -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`.
+9 -2
View File
@@ -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)
+31
View File
@@ -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:
+32 -3
View File
@@ -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" /
+9
View File
@@ -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
+5 -5
View File
@@ -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." },
+360 -12
View File
@@ -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 ``<emoji> <name>(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
(``<emoji> <name>(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:
if parsed is None:
return None
name, preview = parsed
# Terminal code block: the command lives in the fenced lines that
# follow the "<emoji> 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
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",
+108
View File
@@ -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:
+47 -3
View File
@@ -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)
# ---------------------------------------------------------------------------
+2
View File
@@ -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: