Compare commits

..
Author SHA1 Message Date
Bailey Dixon 76ead50c60 fix(android): remove provisional threads safely 2026-08-28 22:08:47 -04:00
34 changed files with 585 additions and 707 deletions
+1 -1
View File
@@ -17,7 +17,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/), and this
### Fixed
- **Opening Android no longer claims or interrupts a turn already running in Hermes Desktop/TUI.** Passive foreground and session browsing now use read-only Gateway status plus profile-scoped history; live-session resume remains reserved for explicit Android actions and exact Android-owned recovery.
- **Android provisional Threads can be removed without touching server history.** The drawer now offers a local-only removal action, reconciles promoted phone sessions without duplicate rows, and keeps Thread routing isolated to the active saved connection.
- **The visible Android Sphere keeps its smooth procedural motion across startup and chat.** Backgrounded and motion-disabled surfaces remain still without reducing foreground animation to a stepped ambient pulse.
### Removed
-3
View File
@@ -26,9 +26,6 @@ model device-certified:
renders as Working.
- Run a background process that outlives its parent turn and verify Background
work remains separate from the conversation's Idle state.
- On a physical phone, open and repeatedly foreground Android while the same
session is working in official Desktop/TUI; verify Android sends no live
attach/interrupt RPC, the producer completes, and final history appears.
- Pursue an upstream `session.active_list` profile field/filter or an aggregate
activity route with explicit profile ownership so multi-profile clients do
not need to resolve process-wide rows from durable keys.
@@ -239,60 +239,6 @@ class GatewayForegroundRecoveryInstrumentedTest {
assertEquals(0, fixture.requestsTo("/v1/chat/completions"))
}
@Test
fun desktopOwnedTurn_remainsReadOnlyAcrossAndroidForegroundLifecycle() {
viewModel.setChatVisible(false)
viewModel.updateGatewayClient(null)
gatewayClient.shutdown()
gatewayScope.cancel()
val controlMethods = setOf(
"session.resume",
"session.activate",
"session.interrupt",
"prompt.submit",
)
val baseline = controlMethods.associateWith(fixture::rpcCount)
val baselineActiveList = fixture.rpcCount("session.active_list")
fixture.activeSessionStatus = "working"
gatewayScope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
val okHttp = OkHttpClient()
gatewayClient = GatewayChatClient(
initialDashboardClient = DashboardApiClient(
baseUrl = fixture.server.url("/").toString().trimEnd('/'),
okHttpClient = okHttp,
),
okHttpClient = okHttp,
callbackDispatcher = { block -> Handler(Looper.getMainLooper()).post(block) },
scope = gatewayScope,
reconnectJitterUnit = { 0.0 },
)
viewModel.setChatTurnCheckpointStore(null)
viewModel.updateGatewayClient(gatewayClient)
viewModel.setChatVisible(true)
compose.activityRule.scenario.moveToState(Lifecycle.State.STARTED)
compose.activityRule.scenario.moveToState(Lifecycle.State.RESUMED)
viewModel.setChatVisible(false)
viewModel.setChatVisible(true)
fixture.awaitRpcCount("session.active_list", baselineActiveList + 1)
controlMethods.forEach { method ->
assertEquals(
"passive lifecycle sent $method",
baseline.getValue(method),
fixture.rpcCount(method),
)
}
viewModel.updateGatewayClient(null)
gatewayClient.shutdown()
assertEquals(
"observer teardown interrupted the Desktop turn",
baseline.getValue("session.interrupt"),
fixture.rpcCount("session.interrupt"),
)
}
private companion object {
const val STORED_SESSION_ID = "20260821_120000_fixture"
const val LIVE_SESSION_ID = "fixture-live-1"
@@ -317,9 +263,6 @@ internal class AndroidGatewayContractFixture {
@Volatile
var recoveryRunning = false
@Volatile
var activeSessionStatus: String? = null
private val listener = object : WebSocketListener() {
override fun onOpen(webSocket: WebSocket, response: Response) {
sockets.add(webSocket)
@@ -339,18 +282,6 @@ internal class AndroidGatewayContractFixture {
"session.activate" -> sessionSnapshot(
(params["session_id"] as? JsonPrimitive)?.contentOrNull ?: "fixture-live-1",
)
"session.active_list" -> buildJsonObject {
put("sessions", kotlinx.serialization.json.buildJsonArray {
activeSessionStatus?.let { status ->
add(buildJsonObject {
put("id", LIVE_SESSION_ID)
put("session_key", STORED_SESSION_ID)
put("status", status)
put("last_active", 1.0)
})
}
})
}
"prompt.submit", "session.interrupt" -> buildJsonObject { put("ok", true) }
else -> JsonObject(emptyMap())
}
@@ -414,15 +345,6 @@ internal class AndroidGatewayContractFixture {
error("Gateway RPC $method not observed; saw ${rpcLog.map { it.first }}")
}
fun awaitRpcCount(method: String, count: Int) {
val deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5)
while (System.nanoTime() < deadline) {
if (rpcCount(method) >= count) return
Thread.sleep(20)
}
error("Gateway RPC $method count $count not observed; saw ${rpcLog.map { it.first }}")
}
fun requestsTo(path: String): Int = requestPaths.count { it.startsWith(path) }
fun rpcCount(method: String): Int = rpcLog.count { it.first == method }
@@ -431,9 +353,4 @@ internal class AndroidGatewayContractFixture {
allSockets.forEach { socket -> runCatching { socket.close(1001, "teardown") } }
runCatching { server.shutdown() }
}
private companion object {
const val STORED_SESSION_ID = "20260821_120000_fixture"
const val LIVE_SESSION_ID = "fixture-live-1"
}
}
@@ -38,6 +38,8 @@ data class ProactiveInboxEntry(
val connectionId: String? = null,
/** Relay proved this row came from its bounded offline queue. */
val arrivedWhileAway: Boolean = false,
/** Exact Android notification slot, when recorded by the receiving build. */
val notificationId: Int? = null,
)
private val Context.proactiveInboxStore: DataStore<Preferences> by
@@ -58,15 +60,19 @@ private const val MAX_ENTRIES = 100
* bounded store also backs the provisional Thread until the user's first reply
* promotes it to a real `source=phone` session.
*/
class ProactiveInboxRepository(private val context: Context) {
class ProactiveInboxRepository internal constructor(
private val store: DataStore<Preferences>,
) {
constructor(context: Context) : this(context.proactiveInboxStore)
private val json = Json { ignoreUnknownKeys = true }
val entries: Flow<List<ProactiveInboxEntry>> =
context.proactiveInboxStore.data.map { prefs -> decode(prefs[INBOX_JSON]) }
store.data.map { prefs -> decode(prefs[INBOX_JSON]) }
suspend fun add(entry: ProactiveInboxEntry) {
context.proactiveInboxStore.edit { prefs ->
store.edit { prefs ->
val current = decode(prefs[INBOX_JSON]).toMutableList()
current.removeAll { it.id == entry.id }
current.add(0, entry)
@@ -76,7 +82,40 @@ class ProactiveInboxRepository(private val context: Context) {
}
suspend fun clear() {
context.proactiveInboxStore.edit { it.remove(INBOX_JSON) }
store.edit { it.remove(INBOX_JSON) }
}
/**
* Remove one provisional Thread owned by one saved connection.
*
* This only edits the bounded local inbox. A promoted Thread is server
* history and is deliberately outside this repository, so this operation
* can never delete it. Legacy entries without a connection owner are
* removed with the active row because they are rendered in that row; rows
* explicitly owned by another connection remain isolated.
*/
suspend fun removeThread(
chatId: String,
connectionId: String,
): List<ProactiveInboxEntry> {
val normalizedChatId = chatId.ifBlank { "phone" }
var removed = emptyList<ProactiveInboxEntry>()
store.edit { prefs ->
val current = decode(prefs[INBOX_JSON])
removed = current.filter {
(it.connectionId == null || it.connectionId == connectionId) &&
(it.chatId ?: "phone") == normalizedChatId
}
if (removed.isNotEmpty()) {
val retained = current.filterNot { it in removed }
if (retained.isEmpty()) {
prefs.remove(INBOX_JSON)
} else {
prefs[INBOX_JSON] = json.encodeToString(retained)
}
}
}
return removed
}
private fun decode(raw: String?): List<ProactiveInboxEntry> {
@@ -102,22 +102,18 @@ class ProactiveMessageHandler(
/** Route a parsed message: into the open Thread if it belongs there, else
* the durable inbox log + the surface its hint selects. */
private fun dispatch(msg: ProactiveMessage) {
// Persist first even when the currently open Thread consumes the live
// message. Agent-initiated outbound sends do not create a gateway
// session until the phone replies, so this cache is the provisional
// Thread transcript during that gap.
toInbox?.invoke(msg)
// The surfacing hint selects the additional surface. Thread injection
// is best-effort presentation of the persisted row, not itself a reason
// to suppress an explicitly requested notification.
when (msg.surfacing?.lowercase()) {
val notificationId = when (msg.surfacing?.lowercase()) {
"inbox" -> {
injectIntoThread?.invoke(msg)
null
}
"session" -> {
val delivered = injectIntoThread?.invoke(msg) == true ||
toSession?.invoke(msg) == true
if (!delivered) notify(msg)
if (delivered) null else notify(msg)
}
// null / "default" / "notification" / anything unrecognized.
else -> {
@@ -125,9 +121,13 @@ class ProactiveMessageHandler(
notify(msg)
}
}
// Every message remains in the bounded local cache. Persist the exact
// posted notification slot as part of that row so a later local Thread
// removal can cancel only its own notification.
toInbox?.invoke(msg.copy(notificationId = notificationId))
}
private fun notify(msg: ProactiveMessage) {
private fun notify(msg: ProactiveMessage): Int? =
ProactiveMessageNotifier.notify(
context = context,
title = msg.title,
@@ -135,7 +135,6 @@ class ProactiveMessageHandler(
messageId = msg.messageId,
chatId = msg.chatId,
)
}
private fun parse(payload: JsonObject): ProactiveMessage? {
val text = payload["text"]?.jsonPrimitive?.contentOrNull
@@ -172,4 +171,6 @@ data class ProactiveMessage(
val replyTo: String? = null,
/** True only when Relay explicitly marked this as a reconnect queue flush. */
val arrivedWhileAway: Boolean = false,
/** Exact Android notification slot when this delivery posted one. */
val notificationId: Int? = null,
)
@@ -906,30 +906,6 @@ class GatewayChatClient(
scope.launch { prewarmAwait(storedSessionId) }
}
/**
* Establish only the shared Gateway socket for read-only observation.
*
* `session.resume` and `session.activate` attach a live runtime to this
* transport. Opening Chat, foreground restoration, and selecting a saved
* transcript must not claim a turn that another Desktop/TUI client owns,
* so those paths use this socket-only warmup and observe through REST
* history plus `session.active_list` instead.
*/
fun observe(onReady: (() -> Unit)? = null) {
scope.launch {
if (observeAwait() && onReady != null) callbackDispatcher(onReady)
}
}
/** Suspending [observe]; returns true once the read-only socket is ready. */
suspend fun observeAwait(): Boolean = try {
connectMutex.withLock { ensureConnected() }
true
} catch (e: Exception) {
Log.d(TAG, "Gateway observation warmup skipped: ${e.message}")
false
}
/**
* Suspending [prewarm]: establishes the socket and (when [storedSessionId]
* is non-null) resumes the existing session, returning only once that work
@@ -74,13 +74,13 @@ object ProactiveMessageNotifier {
text: String,
messageId: String?,
chatId: String?,
) {
): Int? {
ensureChannel(context)
if (!hasPostNotificationsPermission(context)) {
Log.i(TAG, "POST_NOTIFICATIONS not granted — skipping proactive notification")
return
return null
}
if (text.isBlank()) return
if (text.isBlank()) return null
val tapIntent = Intent(context, MainActivity::class.java).apply {
flags = Intent.FLAG_ACTIVITY_NEW_TASK or Intent.FLAG_ACTIVITY_CLEAR_TOP
@@ -89,7 +89,7 @@ object ProactiveMessageNotifier {
val pendingFlags = PendingIntent.FLAG_UPDATE_CURRENT or PendingIntent.FLAG_IMMUTABLE
// Distinct requestCode per slot so each notification gets its own
// PendingIntent rather than all sharing slot 0's intent.
val notificationId = slotFor(messageId)
val notificationId = slotFor(messageId, chatId)
val tapPending =
PendingIntent.getActivity(context, notificationId, tapIntent, pendingFlags)
@@ -108,9 +108,15 @@ object ProactiveMessageNotifier {
.setCategory(NotificationCompat.CATEGORY_MESSAGE)
.setPriority(NotificationCompat.PRIORITY_HIGH)
runCatching {
return runCatching {
NotificationManagerCompat.from(context).notify(notificationId, builder.build())
}.onFailure { Log.w(TAG, "notify failed", it) }
notificationId
}.onFailure { Log.w(TAG, "notify failed", it) }.getOrNull()
}
/** Cancel one exact slot previously returned by [notificationIdFor]. */
fun cancel(context: Context, notificationId: Int) {
NotificationManagerCompat.from(context).cancel(notificationId)
}
/**
@@ -206,8 +212,12 @@ object ProactiveMessageNotifier {
}
/** Derive a stable notification slot from the message id. */
private fun slotFor(messageId: String?): Int {
val key = messageId?.takeIf { it.isNotBlank() } ?: return ID_BASE
internal fun notificationIdFor(messageId: String?, chatId: String?): Int =
slotFor(messageId, chatId)
private fun slotFor(messageId: String?, chatId: String?): Int {
val key = messageId?.takeIf { it.isNotBlank() }
?: "chat:${chatId?.takeIf { it.isNotBlank() } ?: "phone"}"
// Keep within a small positive window above the base so re-delivery of
// the same id collapses to one slot and distinct ids spread out.
return ID_BASE + (key.hashCode() and 0xFFFF)
@@ -2217,9 +2217,12 @@ fun RelayApp() {
(it.connectionId == null || it.connectionId == activeConnectionId) &&
(it.chatId ?: "phone") == chatId
}
if (entries.isEmpty()) return@LaunchedEffect
chatViewModel.openProactiveThread(chatId, entries)
if (entries.isNotEmpty()) {
chatViewModel.openProactiveThread(chatId, entries)
}
}
// Consume the request even when deletion removed its
// local row before a stale notification tap arrived.
backStackEntry.arguments?.putString(
Screen.Chat.ARG_PROACTIVE_CHAT_ID,
null,
@@ -237,6 +237,8 @@ fun SessionDrawerContent(
onNewThread: ((String) -> Unit)? = null,
provisionalThreads: List<ProvisionalThreadRow> = emptyList(),
onSelectProvisionalThread: ((String) -> Unit)? = null,
/** Deletes only the local provisional inbox row; never a server session. */
onDeleteProvisionalThread: ((String) -> Unit)? = null,
/** Gateway sources currently hidden from the drawer (default: cron+webhook). */
hiddenSources: Set<String> = emptySet(),
/** Toggle a source's visibility (persisted). Null hides the source filter. */
@@ -789,13 +791,17 @@ fun SessionDrawerContent(
showTokens = viewOptions.showTokens,
showCost = viewOptions.showCost,
nowMillis = drawerNowMillis,
actionsEnabled = !provisional && (
actionsEnabled = if (provisional) {
onDeleteProvisionalThread != null &&
supervisedSessionActions?.delete != false
} else {
supervisedSessionActions == null ||
supervisedSessionActions.pin ||
supervisedSessionActions.rename ||
supervisedSessionActions.delete ||
(supervisedSessionActions.archive && archiveSupported)
),
},
provisional = provisional,
isActive = !showAllProfiles && session.sessionId == currentSessionId,
activityState = activityState,
animationEnabled = animationEnabled && isOpen,
@@ -934,15 +940,38 @@ fun SessionDrawerContent(
// Delete confirmation dialog
deleteDialogTarget?.let { (row, allProfiles) ->
val session = row.session
val provisional = session.sessionId.startsWith(PROVISIONAL_THREAD_PREFIX)
AlertDialog(
onDismissRequest = { deleteDialogTarget = null },
title = { Text(stringResource(R.string.drawer_delete_session_title)) },
title = {
Text(
stringResource(
if (provisional) {
R.string.drawer_remove_provisional_thread_title
} else {
R.string.drawer_delete_session_title
},
),
)
},
text = {
Text(stringResource(R.string.drawer_delete_session_prefix) + (session.title ?: stringResource(R.string.drawer_untitled)) + stringResource(R.string.drawer_delete_session_suffix))
val title = session.title ?: stringResource(R.string.drawer_untitled)
Text(
if (provisional) {
stringResource(R.string.drawer_remove_provisional_thread_message, title)
} else {
stringResource(R.string.drawer_delete_session_prefix) + title +
stringResource(R.string.drawer_delete_session_suffix)
},
)
},
confirmButton = {
TextButton(onClick = {
if (allProfiles) {
if (session.sessionId.startsWith(PROVISIONAL_THREAD_PREFIX)) {
onDeleteProvisionalThread?.invoke(
session.sessionId.removePrefix(PROVISIONAL_THREAD_PREFIX),
)
} else if (allProfiles) {
onDeleteProfileSession?.invoke(row.profile, session.sessionId)
} else {
onDeleteSession(session.sessionId)
@@ -1358,6 +1387,7 @@ private fun SessionItem(
showCost: Boolean,
nowMillis: Long,
actionsEnabled: Boolean,
provisional: Boolean,
isActive: Boolean,
activityState: SessionActivityState?,
animationEnabled: Boolean,
@@ -1528,7 +1558,7 @@ private fun SessionItem(
expanded = menuOpen,
onDismissRequest = { menuOpen = false },
) {
if (supervisedSessionActions?.pin != false) DropdownMenuItem(
if (!provisional && supervisedSessionActions?.pin != false) DropdownMenuItem(
text = {
Text(
if (pinned) {
@@ -1554,7 +1584,7 @@ private fun SessionItem(
onTogglePinned()
},
)
if (supervisedSessionActions == null) DropdownMenuItem(
if (!provisional && supervisedSessionActions == null) DropdownMenuItem(
text = { Text(stringResource(R.string.chat_copy_session_id)) },
leadingIcon = {
Icon(Icons.Filled.ContentCopy, contentDescription = null)
@@ -1564,7 +1594,7 @@ private fun SessionItem(
onCopySessionId()
},
)
if (supervisedSessionActions?.rename != false) DropdownMenuItem(
if (!provisional && supervisedSessionActions?.rename != false) DropdownMenuItem(
text = { Text(stringResource(R.string.drawer_rename)) },
leadingIcon = {
Icon(Icons.Filled.Edit, contentDescription = null)
@@ -1574,7 +1604,7 @@ private fun SessionItem(
onRename()
},
)
if (archiveSupported && supervisedSessionActions?.archive != false) {
if (!provisional && archiveSupported && supervisedSessionActions?.archive != false) {
DropdownMenuItem(
text = { Text(if (archived) stringResource(R.string.drawer_restore) else stringResource(R.string.drawer_archive)) },
leadingIcon = {
@@ -1119,16 +1119,12 @@ fun ChatScreen(
}
// Recover any durable in-flight chat checkpoint whenever Chat returns to
// the foreground. setChatVisible owns that edge; an ordinary Gateway open
// warms only the observation socket and never attaches a saved session.
// the foreground. On Gateway this also pre-warms/re-attaches the socket;
// sessions-SSE falls back to bounded persisted-history reconciliation.
val appForeground by com.hermesandroid.relay.util.AppForegroundTracker.isForeground.collectAsState()
LaunchedEffect(isGatewayTransport, appForeground, chatReady) {
val chatVisible = appForeground && chatReady
val visibilityChanged = chatViewModel.setChatVisible(chatVisible)
if (isGatewayTransport && chatVisible && !visibilityChanged) {
// Gateway availability can settle after Chat was already visible.
// Repeat the socket-only warmup for that edge; ordinary observation
// still cannot resume or activate a session.
chatViewModel.setChatVisible(appForeground && chatReady)
if (appForeground && chatReady) {
chatViewModel.prewarmGateway()
}
if (isGatewayTransport && appForeground && chatReady) {
@@ -2435,6 +2431,26 @@ fun ChatScreen(
activeConnectionId = activeConnection?.id,
realThreadChatIds = phoneThreadChatIds.values,
)
val realPhoneSessionIds = remember(sessions) {
sessions.asSequence()
.filter { it.source.equals("phone", ignoreCase = true) }
.map { it.sessionId }
.toSet()
}
val provisionalThreadChatIds = provisionalThreadEntries.keys
LaunchedEffect(
activeConnection?.id,
realPhoneSessionIds,
provisionalThreadChatIds,
) {
// A reply promotes the local provisional row to a real Gateway
// source=phone session. Refresh the relay-owned chat_id index at
// that boundary so the local duplicate disappears immediately,
// without guessing a chat_id from the opaque session id.
if (realPhoneSessionIds.isNotEmpty() && provisionalThreadChatIds.isNotEmpty()) {
connectionViewModel.refreshPhoneThreadChatIds()
}
}
val provisionalThreads = provisionalThreadEntries.map { (chatId, entries) ->
val latest = entries.maxBy { it.receivedAt }
ProvisionalThreadRow(
@@ -2553,6 +2569,11 @@ fun ChatScreen(
)
scope.launch { drawerState.close() }
},
onDeleteProvisionalThread = { chatId ->
activeConnection?.id?.let { connectionId ->
connectionViewModel.removeProvisionalThread(chatId, connectionId)
}
},
hiddenSources = hiddenSources,
onToggleSourceHidden = { source, hidden ->
connectionViewModel.setSourceHidden(source, hidden)
@@ -408,9 +408,6 @@ class ChatViewModel : ViewModel() {
private val sessionActivityGeneration = AtomicLong(0L)
private val sessionActivityPollMutex = Mutex()
private var sessionActivityPollJob: Job? = null
private var passiveGatewayHistoryRefreshJob: Job? = null
private var passivelyObservedGatewaySessionId: String? = null
private var passiveObservationCatchupPendingSessionId: String? = null
private var sessionActivityDirectory: Set<SessionActivityOwner> = emptySet()
private var lastProjectedProcessIds: Set<String> = emptySet()
private var lastProjectedProcessOwner: SessionActivityOwner? = null
@@ -732,6 +729,7 @@ class ChatViewModel : ViewModel() {
* including one this app didn't create, or any Thread after a restart.
*/
fun seedThreadChatIds(map: Map<String, String>) {
threadChatIds.clear()
threadChatIds.putAll(map)
}
@@ -2255,8 +2253,6 @@ class ChatViewModel : ViewModel() {
}
private suspend fun pollSessionActivity(client: GatewayChatClient) {
var hasPassiveCurrentLiveWork = false
var hasPassiveCatchupPending = false
sessionActivityPollMutex.withLock {
if (gatewayClient !== client || !chatVisible || streamingEndpoint != "gateway") return
val generation = sessionActivityGeneration.get()
@@ -2277,43 +2273,6 @@ class ChatViewModel : ViewModel() {
when (val result = client.listActiveSessions()) {
is GatewayActiveSessionsResult.Success -> {
if (gatewayClient !== client || generation != sessionActivityGeneration.get()) return
val currentStoredId = currentOwner?.storedSessionId
val passiveCurrentRows = if (currentStoredId == null) {
emptyList()
} else {
result.sessions.filter { row ->
row.storedSessionId == currentStoredId &&
client.knownSessionOwner(row.runtimeSessionId) == null
}
}
hasPassiveCurrentLiveWork = passiveCurrentRows.any { row ->
row.status != GatewayActiveSessionStatus.Idle
}
val initialCatchupPending =
passiveObservationCatchupPendingSessionId == currentStoredId
val needsFinalPassiveRefresh =
passivelyObservedGatewaySessionId == currentStoredId &&
!hasPassiveCurrentLiveWork
if (hasPassiveCurrentLiveWork) {
currentStoredId?.let(::refreshPassivelyObservedGatewayHistory)
if (initialCatchupPending) {
passiveObservationCatchupPendingSessionId = null
}
} else if (needsFinalPassiveRefresh || initialCatchupPending) {
val scheduled = currentStoredId?.let { storedId ->
refreshPassivelyObservedGatewayHistory(
storedSessionId = storedId,
retryUntilChanged = true,
)
} == true
if (scheduled && initialCatchupPending) {
passiveObservationCatchupPendingSessionId = null
}
}
passivelyObservedGatewaySessionId =
currentStoredId?.takeIf { hasPassiveCurrentLiveWork }
hasPassiveCatchupPending =
passiveObservationCatchupPendingSessionId == currentStoredId
val resolved = resolveGatewayActiveSessions(
sessions = result.sessions,
directory = directory,
@@ -2375,18 +2334,6 @@ class ChatViewModel : ViewModel() {
GatewayActiveSessionsResult.Unsupported,
is GatewayActiveSessionsResult.TransientFailure -> {
if (gatewayClient !== client || generation != sessionActivityGeneration.get()) return
val currentStoredId = currentOwner?.storedSessionId
if (passiveObservationCatchupPendingSessionId == currentStoredId) {
val scheduled = currentStoredId?.let { storedId ->
refreshPassivelyObservedGatewayHistory(
storedSessionId = storedId,
retryUntilChanged = true,
)
} == true
if (scheduled) passiveObservationCatchupPendingSessionId = null
}
hasPassiveCatchupPending =
passiveObservationCatchupPendingSessionId == currentStoredId
val scopes = directory.mapTo(mutableSetOf()) {
SessionActivityScope.of(it.connectionId, it.profile)
}.apply { add(currentScope) }
@@ -2405,9 +2352,7 @@ class ChatViewModel : ViewModel() {
record.freshness == SessionActivityFreshness.Confirmed &&
record.phase(System.currentTimeMillis()) != SessionActivityPhase.Idle
}
val delayMs = if (
hasConfirmedLiveWork || hasPassiveCurrentLiveWork || hasPassiveCatchupPending
) 1_500L else 30_000L
val delayMs = if (hasConfirmedLiveWork) 1_500L else 30_000L
sessionActivityPollJob = viewModelScope.launch {
delay(delayMs)
if (gatewayClient === client && chatVisible) pollSessionActivity(client)
@@ -2480,10 +2425,6 @@ class ChatViewModel : ViewModel() {
clearProjectedBackgroundProcesses()
sessionActivityPollJob?.cancel()
sessionActivityPollJob = null
passiveGatewayHistoryRefreshJob?.cancel()
passiveGatewayHistoryRefreshJob = null
passivelyObservedGatewaySessionId = null
passiveObservationCatchupPendingSessionId = null
sessionActivityGeneration.incrementAndGet()
sessionActivityDirectory = emptySet()
lastLocalActivityOwner = null
@@ -2603,9 +2544,8 @@ class ChatViewModel : ViewModel() {
// Foreground can race OkHttp's delayed close callback:
// the first prewarm sees the old socket as Ready, then
// the callback moves it to Idle. Re-run from this exact
// client transition so the observation socket is
// restored; only an exact Android checkpoint may
// resume/activate a live runtime.
// client transition so the visible durable session is
// resumed and its authoritative history reconciled.
prewarmGateway()
}
}
@@ -3060,103 +3000,6 @@ class ChatViewModel : ViewModel() {
}
}
/**
* Refresh a Desktop/TUI-owned turn through the profile-scoped history
* surface without attaching its live runtime. `session.active_list` drives
* the bounded cadence; one final read follows Working/Waiting -> Idle.
*/
private fun refreshPassivelyObservedGatewayHistory(
storedSessionId: String,
retryUntilChanged: Boolean = false,
): Boolean {
if (passiveGatewayHistoryRefreshJob?.isActive == true) return false
if (_isLoadingHistory.value) return false
val handler = chatHandler ?: return false
val contextKey = activeProfileContextKey
val profileName = currentSessionProfileName()
val refreshJob = viewModelScope.launch(start = CoroutineStart.LAZY) {
try {
repeat(if (retryUntilChanged) 8 else 1) { attempt ->
val serverMessages = runCatching {
loadGatewaySessionHistory(
sessionId = storedSessionId,
requireProfileScope = true,
profileName = profileName,
)
}.getOrNull() ?: return@launch
if (
chatHandler !== handler ||
activeProfileContextKey != contextKey ||
currentSessionProfileName() != profileName ||
handler.currentSessionId.value != storedSessionId ||
_isLoadingHistory.value ||
activeStream != null ||
handler.isStreaming.value
) return@launch
val visibleSignature = handler.messages.value
.filterNot { it.clientOnly }
.map { message ->
Triple(
message.role.name.lowercase(),
message.content,
message.thinkingContent,
)
}
val serverSignature = serverMessages.map { message ->
Triple(
message.role.lowercase(),
message.contentText.orEmpty(),
message.resolvedReasoning.orEmpty(),
)
}
if (visibleSignature != serverSignature) {
handler.loadMessageHistory(serverMessages)
refreshSessions()
scheduleTitleReconcile(storedSessionId)
return@launch
}
if (attempt < 7 && retryUntilChanged) delay(250L)
}
} finally {
if (passiveGatewayHistoryRefreshJob === coroutineContext[Job]) {
passiveGatewayHistoryRefreshJob = null
}
}
}
passiveGatewayHistoryRefreshJob = refreshJob
refreshJob.start()
return true
}
/** Open the read-only socket off Main, then publish observation ownership on Main. */
private fun observeGatewaySession(
client: GatewayChatClient?,
handler: ChatHandler,
storedSessionId: String,
) {
val observer = client ?: return
val contextKey = activeProfileContextKey
val profileName = currentSessionProfileName()
observer.observe {
if (
chatVisible &&
gatewayClient === observer &&
chatHandler === handler &&
activeProfileContextKey == contextKey &&
currentSessionProfileName() == profileName &&
handler.currentSessionId.value == storedSessionId
) {
passiveObservationCatchupPendingSessionId = storedSessionId
refreshPassivelyObservedGatewayHistory(
storedSessionId = storedSessionId,
retryUntilChanged = true,
)
requestSessionActivityRefresh()
}
}
}
/** One-shot `config.get personality` over a ready socket → drives the collector. */
private fun seedServerPersonality(client: GatewayChatClient) {
viewModelScope.launch {
@@ -3167,11 +3010,11 @@ class ChatViewModel : ViewModel() {
}
/**
* Warm the Gateway socket when Chat is visible without claiming a runtime
* that may belong to Desktop/TUI. Exact Android-owned checkpoints recover
* through `session.activate`/`session.resume`; an ordinary open observes
* through REST history and `session.active_list` until the user performs
* an explicit action that needs session ownership.
* Warm the gateway socket (and resume the current session) when the chat
* surface is visible and the gateway is the resolved transport, so the
* first send is warm instead of paying the cold connect + `session.resume`
* on the send path. No-op without a gateway client; idempotent when warm.
* Driven by a foreground/visibility effect in ChatScreen.
*/
fun prewarmGateway() {
val client = gatewayClient
@@ -3179,16 +3022,17 @@ class ChatViewModel : ViewModel() {
val sessionId = handler.currentSessionId.value
selectBackgroundProcessSession(sessionId)
if (sessionId == null) {
client?.observe()
client?.prewarm(null)
} else {
// Preserve the original warm-up path before persistence wiring is
// available (early composition and JVM tests). Production installs
// the store from initializeMedia before Chat becomes ready.
if (chatTurnCheckpointStore == null) {
val gateway = client ?: return
// GatewayChatClient owns the socket IO scope, so the dial can
// progress while a paused UI dispatcher is being recreated.
observeGatewaySession(gateway, handler, sessionId)
// GatewayChatClient owns an IO scope, so this can progress even
// while a paused/blocked UI dispatcher is being recreated.
// Its cold-ready listener performs history/process refresh.
gateway.prewarm(sessionId)
return
}
if (activeStream == null && (streamRecovery == null || client != null)) {
@@ -3200,29 +3044,26 @@ class ChatViewModel : ViewModel() {
chatHandler === handler &&
handler.currentSessionId.value == sessionId
) {
observeGatewaySession(client, handler, sessionId)
if (client?.prewarmAwait(sessionId) == true) {
gatewayProcessController.sessionReady(sessionId)
}
}
checkpointRecoveryJob = null
}
return
}
// A locally-owned live mapper may revalidate its existing binding.
// A passive transcript must remain socket-only: resuming it here
// can replace another client's transport and turn Android teardown
// into a later session.interrupt.
if (client?.hasActiveTurnForSession(sessionId) == true) {
viewModelScope.launch {
if (client.prewarmAwait(sessionId) &&
gatewayClient === client &&
chatHandler === handler &&
handler.currentSessionId.value == sessionId
) {
gatewayProcessController.sessionReady(sessionId)
requestSessionActivityRefresh()
}
// prewarm() only emits the existing "cold ready" callback when it
// had to resume. An already-live session still needs its initial
// process snapshot when Chat opens, so confirm it explicitly.
viewModelScope.launch {
if (
client?.prewarmAwait(sessionId) == true &&
gatewayClient === client &&
chatHandler === handler &&
handler.currentSessionId.value == sessionId
) {
gatewayProcessController.sessionReady(sessionId)
}
} else {
observeGatewaySession(client, handler, sessionId)
}
}
}
@@ -3232,7 +3073,7 @@ class ChatViewModel : ViewModel() {
* Gateway chat owns automatic idle-socket reattachment; other tabs and a
* backgrounded app retain the normal no-reconnect behavior.
*/
fun setChatVisible(visible: Boolean): Boolean {
fun setChatVisible(visible: Boolean) {
val changed = chatVisible != visible
chatVisible = visible
if (visible && changed) {
@@ -3241,12 +3082,7 @@ class ChatViewModel : ViewModel() {
} else if (!visible) {
sessionActivityPollJob?.cancel()
sessionActivityPollJob = null
passiveGatewayHistoryRefreshJob?.cancel()
passiveGatewayHistoryRefreshJob = null
passivelyObservedGatewaySessionId = null
passiveObservationCatchupPendingSessionId = null
}
return changed
}
// === Gateway desktop-parity state ===
@@ -4630,9 +4466,7 @@ class ChatViewModel : ViewModel() {
)
if (stillCurrent()) {
handler.loadMessageHistory(messages)
if (streamingEndpoint == "gateway") {
observeGatewaySession(gatewayClient, handler, sessionId)
}
if (streamingEndpoint == "gateway") gatewayClient?.prewarm(sessionId)
}
}
} catch (e: kotlinx.coroutines.CancellationException) {
@@ -5143,9 +4977,7 @@ class ChatViewModel : ViewModel() {
handler.currentSessionId.value == sessionId
) {
handler.loadMessageHistory(messages)
if (streamingEndpoint == "gateway") {
observeGatewaySession(gatewayClient, handler, sessionId)
}
if (streamingEndpoint == "gateway") gatewayClient?.prewarm(sessionId)
}
} catch (e: kotlinx.coroutines.CancellationException) {
throw e
@@ -6735,10 +6567,6 @@ class ChatViewModel : ViewModel() {
* SSE cannot, so it retains the existing interrupt/cancel behavior.
*/
private fun releaseTurnForNavigation(handler: ChatHandler) {
passiveGatewayHistoryRefreshJob?.cancel()
passiveGatewayHistoryRefreshJob = null
passivelyObservedGatewaySessionId = null
passiveObservationCatchupPendingSessionId = null
val gateway = gatewayClient
val canBackground = streamingEndpoint == "gateway" &&
activeStreamIsGateway && activeStream != null && gateway != null
@@ -116,6 +116,7 @@ import com.hermesandroid.relay.accessibility.BridgeStatusReporter
import com.hermesandroid.relay.accessibility.ScreenCapture
import com.hermesandroid.relay.network.relay.BridgeCommandHandler
import com.hermesandroid.relay.network.relay.ProactiveMessageHandler
import com.hermesandroid.relay.notifications.ProactiveMessageNotifier
import com.hermesandroid.relay.network.relay.models.Envelope
// === END PHASE3-accessibility ===
import com.hermesandroid.relay.util.AppForegroundTracker
@@ -164,6 +165,33 @@ internal data class RelayUiInputs(
val configured: Boolean,
)
internal data class PhoneThreadChatIdIndex(
val connectionId: String? = null,
val values: Map<String, String> = emptyMap(),
)
internal fun visiblePhoneThreadChatIds(
activeConnectionId: String?,
index: PhoneThreadChatIdIndex,
): Map<String, String> =
index.values.takeIf { index.connectionId == activeConnectionId }.orEmpty()
internal fun reconcilePhoneThreadChatIdIndex(
current: PhoneThreadChatIdIndex,
requestedConnectionId: String,
activeConnectionId: String?,
fetched: Result<Map<String, String>>,
): PhoneThreadChatIdIndex = fetched.fold(
onSuccess = { values ->
if (requestedConnectionId == activeConnectionId) {
PhoneThreadChatIdIndex(requestedConnectionId, values)
} else {
current
}
},
onFailure = { current },
)
data class HostResourcePressureStatus(
val memoryPressure: String? = null,
val memoryAvailableMb: Int? = null,
@@ -2568,6 +2596,18 @@ class ConnectionViewModel(application: Application) : AndroidViewModel(applicati
val inboxMessages: StateFlow<List<ProactiveInboxEntry>> =
proactiveInbox.entries.stateIn(viewModelScope, SharingStarted.Eagerly, emptyList())
/** Delete one connection-scoped provisional Thread from local storage only. */
fun removeProvisionalThread(chatId: String, connectionId: String) {
viewModelScope.launch {
proactiveInbox.removeThread(chatId = chatId, connectionId = connectionId)
.mapNotNull(ProactiveInboxEntry::notificationId)
.distinct()
.forEach { notificationId ->
ProactiveMessageNotifier.cancel(getApplication(), notificationId)
}
}
}
// The handler centralizes surfacing (notification / inbox / session). The
// inbox sink persists messages here; the session sink lands in Phase 2b.
val proactiveMessageHandler = ProactiveMessageHandler(
@@ -2583,6 +2623,7 @@ class ConnectionViewModel(application: Application) : AndroidViewModel(applicati
chatId = msg.chatId,
connectionId = connectionStore.activeConnectionId.value,
arrivedWhileAway = msg.arrivedWhileAway,
notificationId = msg.notificationId,
),
)
}
@@ -2620,16 +2661,30 @@ class ConnectionViewModel(application: Application) : AndroidViewModel(applicati
// composer's reply routing so a Thread the app didn't create — or any Thread
// after restart — routes to the right conversation. Fail-soft: empty on an
// older relay / fetch error, and the client's learned map still applies.
private val _phoneThreadChatIds = MutableStateFlow<Map<String, String>>(emptyMap())
val phoneThreadChatIds: StateFlow<Map<String, String>> = _phoneThreadChatIds.asStateFlow()
private val _phoneThreadChatIdIndex = MutableStateFlow(PhoneThreadChatIdIndex())
private val phoneThreadChatIdRefreshMutex = Mutex()
val phoneThreadChatIds: StateFlow<Map<String, String>> = combine(
connectionStore.activeConnectionId,
_phoneThreadChatIdIndex,
) { activeConnectionId, index ->
visiblePhoneThreadChatIds(activeConnectionId, index)
}.stateIn(viewModelScope, SharingStarted.Eagerly, emptyMap())
fun refreshPhoneThreadChatIds() {
val connectionId = connectionStore.activeConnectionId.value ?: return
viewModelScope.launch {
relayHttpClient.fetchPhoneThreads().onSuccess { threads ->
val map = threads
.filter { it.sessionId.isNotBlank() && it.chatId.isNotBlank() }
.associate { it.sessionId to it.chatId }
if (map.isNotEmpty()) _phoneThreadChatIds.value = map
phoneThreadChatIdRefreshMutex.withLock {
val fetched = relayHttpClient.fetchPhoneThreads().map { threads ->
threads
.filter { it.sessionId.isNotBlank() && it.chatId.isNotBlank() }
.associate { it.sessionId to it.chatId }
}
_phoneThreadChatIdIndex.value = reconcilePhoneThreadChatIdIndex(
current = _phoneThreadChatIdIndex.value,
requestedConnectionId = connectionId,
activeConnectionId = connectionStore.activeConnectionId.value,
fetched = fetched,
)
}
}
}
@@ -916,6 +916,8 @@
<string name="drawer_delete_session_title">Excluir sessão?</string>
<string name="drawer_delete_session_prefix">Isso excluirá permanentemente \"</string>
<string name="drawer_delete_session_suffix">\" e o histórico de mensagens.</string>
<string name="drawer_remove_provisional_thread_title">Remover Thread?</string>
<string name="drawer_remove_provisional_thread_message">Isso remove \"%1$s\" deste dispositivo. O histórico promovido ou armazenado no servidor não será excluído.</string>
<string name="drawer_untitled">Sem título</string>
<string name="drawer_delete">Excluir</string>
<string name="drawer_thread">Thread</string>
@@ -960,6 +960,8 @@
<string name="drawer_delete_session_title">删除会话?</string>
<string name="drawer_delete_session_prefix">这将永久删除\"</string>
<string name="drawer_delete_session_suffix">\"及其消息历史。</string>
<string name="drawer_remove_provisional_thread_title">移除话题?</string>
<string name="drawer_remove_provisional_thread_message">这会从此设备移除“%1$s”,不会删除已提升或服务器端的历史记录。</string>
<string name="drawer_untitled">未命名</string>
<string name="drawer_delete">删除</string>
<string name="drawer_thread">话题</string>
+2
View File
@@ -965,6 +965,8 @@
<string name="drawer_delete_session_title">Sitzung löschen?</string>
<string name="drawer_delete_session_prefix">Dadurch werden \"</string>
<string name="drawer_delete_session_suffix">\" und der Nachrichtenverlauf dauerhaft gelöscht.</string>
<string name="drawer_remove_provisional_thread_title">Thread entfernen?</string>
<string name="drawer_remove_provisional_thread_message">Dadurch wird „%1$s“ von diesem Gerät entfernt. Hochgestufte oder serverseitige Verläufe werden nicht gelöscht.</string>
<string name="drawer_untitled">Ohne Titel</string>
<string name="drawer_delete">Löschen</string>
<string name="drawer_thread">Thread</string>
+2
View File
@@ -880,6 +880,8 @@
<string name="drawer_delete_session_title">¿Eliminar sesión?</string>
<string name="drawer_delete_session_prefix">Esto eliminará permanentemente \"</string>
<string name="drawer_delete_session_suffix">\" y su historial de mensajes.</string>
<string name="drawer_remove_provisional_thread_title">¿Quitar hilo?</string>
<string name="drawer_remove_provisional_thread_message">Esto quita «%1$s» de este dispositivo. No elimina el historial promocionado ni el del servidor.</string>
<string name="drawer_untitled">Intitulado</string>
<string name="drawer_delete">Borrar</string>
<string name="drawer_thread">Hilo</string>
+2
View File
@@ -976,6 +976,8 @@
<string name="drawer_delete_session_title">セッションを削除しますか?</string>
<string name="drawer_delete_session_prefix">「</string>
<string name="drawer_delete_session_suffix">」とそのメッセージ履歴を完全に削除します。</string>
<string name="drawer_remove_provisional_thread_title">スレッドを削除しますか?</string>
<string name="drawer_remove_provisional_thread_message">「%1$s」をこのデバイスから削除します。昇格済みまたはサーバー上の履歴は削除されません。</string>
<string name="drawer_untitled">無題</string>
<string name="drawer_delete">消去</string>
<string name="drawer_thread">糸</string>
+2
View File
@@ -992,6 +992,8 @@
<string name="drawer_delete_session_title">Удалить сессию?</string>
<string name="drawer_delete_session_prefix">Это навсегда удалит &quot;</string>
<string name="drawer_delete_session_suffix">&quot; и историю сообщений.</string>
<string name="drawer_remove_provisional_thread_title">Удалить поток?</string>
<string name="drawer_remove_provisional_thread_message">Это удалит «%1$s» с этого устройства. Повышенная или серверная история не будет удалена.</string>
<string name="drawer_untitled">Без названия</string>
<string name="drawer_delete">Удалить</string>
<string name="drawer_thread">Ветка</string>
+2
View File
@@ -1090,6 +1090,8 @@
<string name="drawer_delete_session_title">Delete Session?</string>
<string name="drawer_delete_session_prefix">This will permanently delete \"</string>
<string name="drawer_delete_session_suffix">\" and its message history.</string>
<string name="drawer_remove_provisional_thread_title">Remove Thread?</string>
<string name="drawer_remove_provisional_thread_message">This removes \"%1$s\" from this device. It does not delete promoted or server history.</string>
<string name="drawer_untitled">Untitled</string>
<string name="drawer_delete">Delete</string>
<string name="drawer_thread">Thread</string>
@@ -0,0 +1,73 @@
package com.hermesandroid.relay.data
import androidx.datastore.core.DataStore
import androidx.datastore.preferences.core.Preferences
import androidx.datastore.preferences.core.emptyPreferences
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.runBlocking
import org.junit.Assert.assertEquals
import org.junit.Test
class ProactiveInboxStoreTest {
@Test
fun `remove thread matches the rendered connection row and preserves other rows`() = runBlocking {
val repository = ProactiveInboxRepository(InMemoryPreferencesDataStore())
repository.add(entry("owned", "reminders", "connection-a", notificationId = 42))
repository.add(entry("legacy", "reminders", null))
repository.add(entry("other-connection", "reminders", "connection-b"))
repository.add(entry("other-thread", "updates", "connection-a"))
val removed = repository.removeThread("reminders", "connection-a")
assertEquals(setOf("owned", "legacy"), removed.map { it.id }.toSet())
assertEquals(listOf(42), removed.mapNotNull { it.notificationId })
assertEquals(
setOf("other-connection", "other-thread"),
repository.entries.first().map { it.id }.toSet(),
)
}
@Test
fun `blank chat id removes only the local phone fallback row`() = runBlocking {
val repository = ProactiveInboxRepository(InMemoryPreferencesDataStore())
repository.add(entry("default", null, "connection-a"))
repository.add(entry("named", "reminders", "connection-a"))
repository.removeThread("phone", "connection-a")
assertEquals(listOf("named"), repository.entries.first().map { it.id })
}
private fun entry(
id: String,
chatId: String?,
connectionId: String?,
notificationId: Int? = null,
) =
ProactiveInboxEntry(
id = id,
title = "Hermes",
text = id,
receivedAt = 1L,
chatId = chatId,
connectionId = connectionId,
notificationId = notificationId,
)
private class InMemoryPreferencesDataStore : DataStore<Preferences> {
private val state = MutableStateFlow<Preferences>(emptyPreferences())
override val data: Flow<Preferences> = state
override suspend fun updateData(
transform: suspend (t: Preferences) -> Preferences,
): Preferences {
val next = transform(state.value)
state.value = next
return next
}
}
}
@@ -3,9 +3,7 @@ package com.hermesandroid.relay.network.relay
import android.content.Context
import com.hermesandroid.relay.network.relay.models.Envelope
import com.hermesandroid.relay.notifications.ProactiveMessageNotifier
import io.mockk.Runs
import io.mockk.every
import io.mockk.just
import io.mockk.mockk
import io.mockk.mockkObject
import io.mockk.unmockkObject
@@ -26,7 +24,7 @@ class ProactiveMessageHandlerTest {
mockkObject(ProactiveMessageNotifier)
every {
ProactiveMessageNotifier.notify(any(), any(), any(), any(), any())
} just Runs
} returns 42
}
@After
@@ -44,6 +42,7 @@ class ProactiveMessageHandlerTest {
handler.onMessage(messageEnvelope(surfacing = "notification"))
assertEquals(1, persisted.size)
assertEquals(42, persisted.single().notificationId)
verify(exactly = 1) {
ProactiveMessageNotifier.notify(context, "Hermes", "ready", "m-1", "phone")
}
@@ -51,7 +50,8 @@ class ProactiveMessageHandlerTest {
@Test
fun `inbox surfacing persists silently`() {
val handler = ProactiveMessageHandler(context, toInbox = {}).apply {
val persisted = mutableListOf<ProactiveMessage>()
val handler = ProactiveMessageHandler(context, toInbox = persisted::add).apply {
injectIntoThread = { true }
}
@@ -60,6 +60,7 @@ class ProactiveMessageHandlerTest {
verify(exactly = 0) {
ProactiveMessageNotifier.notify(any(), any(), any(), any(), any())
}
assertEquals(null, persisted.single().notificationId)
}
@Test
@@ -231,14 +231,11 @@ class GatewayClientHarness(
val suppressAckMethods: MutableSet<String> = ConcurrentHashMap.newKeySet()
val pendingAcks = LinkedBlockingQueue<PendingAck>()
@Volatile
var suppressGatewayReady: Boolean = false
private val wsListener = object : WebSocketListener() {
override fun onOpen(webSocket: WebSocket, response: okhttp3.Response) {
serverSockets.add(webSocket)
allServerSockets.add(webSocket)
if (!suppressGatewayReady) sendGatewayReady(webSocket)
webSocket.send(eventFrame("gateway.ready", null, null))
}
override fun onMessage(webSocket: WebSocket, text: String) {
@@ -630,10 +627,6 @@ class GatewayClientHarness(
fun awaitServerSocket(): WebSocket =
serverSockets.poll(5, TimeUnit.SECONDS) ?: error("server socket never opened")
fun sendGatewayReady(webSocket: WebSocket) {
webSocket.send(eventFrame("gateway.ready", null, null))
}
fun awaitRpc(method: String): JsonObject {
val deadline = System.currentTimeMillis() + 5_000
while (System.currentTimeMillis() < deadline) {
@@ -1414,27 +1407,6 @@ class GatewayChatClientTest {
assertEquals(listOf("stored-session"), resumedSessions.toList())
}
@Test
fun `observation warmup never claims or interrupts a foreign runtime`() = runBlocking {
val registrations = AtomicInteger(0)
client.setUnsolicitedTurnProvider {
registrations.incrementAndGet()
GatewayInboundTurnRegistration(Recorder().callbacks) { true }
}
assertTrue(client.observeAwait())
val serverWs = harness.awaitServerSocket()
serverWs.send(harness.eventFrame("message.start", null, "foreign-runtime"))
delay(100)
client.shutdown()
assertEquals(0, registrations.get())
assertFalse(harness.rpcLog.any { it.first == "session.resume" })
assertFalse(harness.rpcLog.any { it.first == "session.activate" })
assertFalse(harness.rpcLog.any { it.first == "session.interrupt" })
assertFalse(harness.rpcLog.any { it.first == "prompt.submit" })
}
@Test
fun `newer prewarm selection wins when an older resume completes late`() = runBlocking {
harness.suppressAckMethods += "session.resume"
@@ -0,0 +1,74 @@
package com.hermesandroid.relay.notifications
import android.Manifest
import android.app.NotificationManager
import android.content.Context
import android.os.Build
import org.junit.After
import org.junit.Assert.assertEquals
import org.junit.Assert.assertNotEquals
import org.junit.Before
import org.junit.Test
import org.junit.runner.RunWith
import org.robolectric.RobolectricTestRunner
import org.robolectric.RuntimeEnvironment
import org.robolectric.Shadows.shadowOf
import org.robolectric.annotation.Config
@RunWith(RobolectricTestRunner::class)
@Config(sdk = [Build.VERSION_CODES.UPSIDE_DOWN_CAKE])
class ProactiveMessageNotifierTest {
private lateinit var context: Context
private lateinit var manager: NotificationManager
@Before
fun setUp() {
context = RuntimeEnvironment.getApplication()
manager = context.getSystemService(NotificationManager::class.java)
manager.cancelAll()
shadowOf(RuntimeEnvironment.getApplication()).grantPermissions(
Manifest.permission.POST_NOTIFICATIONS,
)
}
@After
fun tearDown() {
manager.cancelAll()
}
@Test
fun `notification identity is stable per wire message id`() {
assertEquals(
ProactiveMessageNotifier.notificationIdFor("message-1", "reminders"),
ProactiveMessageNotifier.notificationIdFor("message-1", "updates"),
)
assertNotEquals(
ProactiveMessageNotifier.notificationIdFor("message-1", "reminders"),
ProactiveMessageNotifier.notificationIdFor("message-2", "reminders"),
)
}
@Test
fun `blank message ids keep independent thread slots`() {
assertNotEquals(
ProactiveMessageNotifier.notificationIdFor(null, "reminders"),
ProactiveMessageNotifier.notificationIdFor(null, "updates"),
)
}
@Test
fun `cancel removes only the persisted notification slot`() {
ProactiveMessageNotifier.notify(context, "Hermes", "first", "message-1", "reminders")
ProactiveMessageNotifier.notify(context, "Hermes", "second", "message-2", "updates")
ProactiveMessageNotifier.cancel(
context,
ProactiveMessageNotifier.notificationIdFor("message-1", "reminders"),
)
assertEquals(
listOf(ProactiveMessageNotifier.notificationIdFor("message-2", "updates")),
manager.activeNotifications.map { it.id },
)
}
}
@@ -18,6 +18,7 @@ import androidx.compose.ui.test.performScrollToNode
import androidx.test.ext.junit.runners.AndroidJUnit4
import com.hermesandroid.relay.data.ChatSession
import com.hermesandroid.relay.data.SessionActivityState
import com.hermesandroid.relay.data.SupervisedSessionActions
import com.hermesandroid.relay.ui.theme.ProfileAccentSwatches
import org.junit.Rule
import org.junit.Test
@@ -44,6 +45,72 @@ class SessionDrawerTest {
assertEquals(0.45f, UNPINNED_STAR_ALPHA, 0.0f)
}
@Test
fun `provisional thread exposes local delete only`() {
var deletedProvisional: String? = null
var deletedServerSession: String? = null
compose.setContent {
MaterialTheme {
SessionDrawerContent(
sessions = emptyList(),
currentSessionId = null,
threadsCapabilityActive = true,
provisionalThreads = listOf(
ProvisionalThreadRow(
chatId = "reminders",
title = "Reminder",
messageCount = 1,
lastActivityAt = 1L,
),
),
onDeleteProvisionalThread = { deletedProvisional = it },
onNewChat = {},
onSelectSession = {},
onDeleteSession = { deletedServerSession = it },
onRenameSession = { _, _ -> },
)
}
}
compose.onNodeWithContentDescription("Session actions").performClick()
compose.onNodeWithText("Pin session").assertDoesNotExist()
compose.onNodeWithText("Rename").assertDoesNotExist()
compose.onNodeWithText("Delete").performClick()
compose.onNodeWithText("Remove Thread?").assertIsDisplayed()
compose.onNodeWithText(
"This removes \"Reminder\" from this device. It does not delete promoted or server history.",
).assertIsDisplayed()
compose.onNodeWithText("Delete").performClick()
compose.runOnIdle {
assertEquals("reminders", deletedProvisional)
assertEquals(null, deletedServerSession)
}
}
@Test
fun `provisional thread hides actions when supervised deletion is disabled`() {
compose.setContent {
MaterialTheme {
SessionDrawerContent(
sessions = emptyList(),
currentSessionId = null,
supervisedSessionActions = SupervisedSessionActions(delete = false),
provisionalThreads = listOf(
ProvisionalThreadRow("reminders", "Reminder", 1, 1L),
),
onDeleteProvisionalThread = {},
onNewChat = {},
onSelectSession = {},
onDeleteSession = {},
onRenameSession = { _, _ -> },
)
}
}
compose.onNodeWithContentDescription("Session actions").assertDoesNotExist()
}
@Test
fun `archive filter resets when connection cannot restore archived sessions`() {
assertEquals(
@@ -36,6 +36,21 @@ class ProvisionalThreadRowsTest {
assertTrue("phone" in rows)
}
@Test
fun promotedChatIdSuppressesOnlyItsProvisionalRow() {
val rows = buildProvisionalThreadRows(
entries = listOf(
entry("promoted", connectionId = "connection-a", chatId = "reminders"),
entry("still-local", connectionId = "connection-a", chatId = "updates"),
),
activeConnectionId = "connection-a",
realThreadChatIds = listOf("reminders"),
)
assertFalse("reminders" in rows)
assertEquals(listOf("still-local"), rows.getValue("updates").map { it.id })
}
private fun entry(id: String, connectionId: String?, chatId: String?) =
ProactiveInboxEntry(
id = id,
@@ -992,12 +992,11 @@ class ChatViewModelGatewayInboundTurnTest {
awaitCondition { handler.messages.value.any { it.content == "Partial A" } }
viewModel.switchSession(secondSession)
gatewayHarness.awaitRpcCount("session.resume", 2)
awaitCondition { handler.currentSessionId.value == secondSession && !handler.isStreaming.value }
assertEquals(1, gatewayHarness.rpcLog.count { it.first == "session.resume" })
assertTrue(gatewayHarness.rpcLog.none { it.first == "session.interrupt" })
viewModel.sendMessage("Run task B")
gatewayHarness.awaitRpcCount("session.resume", 2)
gatewayHarness.awaitRpcCount("prompt.submit", 2)
serverWs.send(
gatewayHarness.eventFrame(
@@ -1825,12 +1824,11 @@ class ChatViewModelGatewayInboundTurnTest {
awaitCondition { viewModel.queuedMessages.value == listOf("Follow up A") }
viewModel.switchSession(secondSession)
gatewayHarness.awaitRpcCount("session.resume", 2)
awaitCondition { handler.currentSessionId.value == secondSession && !handler.isStreaming.value }
assertEquals(1, gatewayHarness.rpcLog.count { it.first == "session.resume" })
assertTrue("session B must not show A's queue", viewModel.queuedMessages.value.isEmpty())
viewModel.sendMessage("Run task B")
gatewayHarness.awaitRpcCount("session.resume", 2)
gatewayHarness.awaitRpcCount("prompt.submit", 2)
serverWs.send(
gatewayHarness.eventFrame(
@@ -2033,16 +2031,13 @@ class ChatViewModelGatewayInboundTurnTest {
serverWs.close(1012, "test disconnect")
awaitCondition { gatewayClient.connectionState.value == GatewayConnectionState.Idle }
viewModel.setChatVisible(true)
viewModel.prewarmGateway()
gatewayHarness.awaitServerSocket()
gatewayHarness.awaitRpcCount("session.resume", 2)
awaitCondition {
handler.messages.value.singleOrNull()?.content == BACKGROUND_ANSWER
}
assertEquals(1, gatewayHarness.rpcLog.count { it.first == "session.resume" })
assertEquals(0, gatewayHarness.rpcLog.count { it.first == "session.activate" })
assertEquals(0, gatewayHarness.rpcLog.count { it.first == "session.interrupt" })
assertFalse(handler.isStreaming.value)
}
@@ -2053,19 +2048,16 @@ class ChatViewModelGatewayInboundTurnTest {
// Foreground arrives while OkHttp still reports the old socket ready,
// so the one-shot prewarm is an intentional no-op. The delayed close
// callback must itself restore the observation socket and catch up
// history without attaching the live session.
// callback must itself trigger an exact-session reattach.
viewModel.prewarmGateway()
serverWs.close(1012, "late background close")
awaitCondition { gatewayHarness.ticketMints.get() >= 2 }
serverWs = gatewayHarness.awaitServerSocket()
gatewayHarness.awaitRpcCount("session.resume", 2)
awaitCondition {
handler.messages.value.singleOrNull()?.content == BACKGROUND_ANSWER
}
assertEquals(1, gatewayHarness.rpcLog.count { it.first == "session.resume" })
assertEquals(0, gatewayHarness.rpcLog.count { it.first == "session.activate" })
assertEquals(0, gatewayHarness.rpcLog.count { it.first == "session.interrupt" })
assertFalse(handler.isStreaming.value)
}
@@ -2374,189 +2366,36 @@ class ChatViewModelGatewayInboundTurnTest {
}
@Test
fun reconnectCatchupClosesCompletionBetweenFirstReadAndIdleSnapshot() {
viewModel.switchProfileContext(PROFILE_CONTEXT, STORED_SESSION_ID)
awaitCondition { !viewModel.isLoadingHistory.value }
val firstReadStarted = CompletableDeferred<Unit>()
val releaseFirstRead = CompletableDeferred<Unit>()
val readCount = AtomicInteger(0)
viewModel.setProfileMessageLoader {
if (readCount.incrementAndGet() == 1) {
firstReadStarted.complete(Unit)
releaseFirstRead.await()
Result.success(emptyList())
} else {
Result.success(persistedHistory)
}
}
fun reconnectAfterMissedStartRecoversOnExactSessionCompletion() {
serverWs.close(1012, "missed start")
awaitCondition { gatewayClient.connectionState.value == GatewayConnectionState.Idle }
viewModel.setChatVisible(true)
viewModel.prewarmGateway()
serverWs = gatewayHarness.awaitServerSocket()
awaitCondition { firstReadStarted.isCompleted }
// Completion persists after the reconnect's first catch-up read began,
// while the first active-list snapshot is already empty/idle. The
// pending final-read marker must close this exact ordering window.
gatewayHarness.awaitRpcCount("session.resume", 2)
// Reconnected midway through the synthetic turn: no message.start is
// replayed, so the delta is intentionally ignored and completion drives
// authoritative history recovery.
serverWs.send(
gatewayHarness.eventFrame(
"message.delta",
buildJsonObject { put("text", BACKGROUND_ANSWER) },
"live-resumed",
),
)
persistedHistory = persistedAnswerHistory()
releaseFirstRead.complete(Unit)
serverWs.send(
gatewayHarness.eventFrame(
"message.complete",
buildJsonObject { put("text", BACKGROUND_ANSWER) },
"live-resumed",
),
)
awaitCondition { handler.messages.value.any { it.content == BACKGROUND_ANSWER } }
assertEquals(1, gatewayHarness.rpcLog.count { it.first == "session.resume" })
assertEquals(0, gatewayHarness.rpcLog.count { it.first == "session.activate" })
assertEquals(0, gatewayHarness.rpcLog.count { it.first == "session.interrupt" })
assertFalse(handler.isStreaming.value)
}
@Test
fun passiveForegroundObservationNeverClaimsOrInterruptsDesktopTurn() {
val observerProfile = Profile(
name = "observer",
model = "model-a",
description = "Observer",
)
viewModel.setSelectedProfileProvider { observerProfile }
viewModel.setSessionProfileNameProvider { observerProfile.name }
viewModel.setProfileMessageLoaderWithMode { profileName, sessionId, _ ->
assertEquals(STORED_SESSION_ID, sessionId)
Result.success(
if (profileName == observerProfile.name) {
persistedHistory
} else {
listOf(
MessageItem(
id = "wrong-profile",
sessionId = STORED_SESSION_ID,
role = "assistant",
content = JsonPrimitive("Wrong profile history"),
),
)
},
)
}
viewModel.switchProfileContext(
AgentDisplay.profileContextKey("connection-a", observerProfile.name),
STORED_SESSION_ID,
)
awaitCondition { !viewModel.isLoadingHistory.value }
viewModel.setChatVisible(false)
viewModel.updateGatewayClient(null)
gatewayClient.shutdown()
gatewayScope.cancel()
val ownershipMethods = setOf(
"session.resume",
"session.activate",
"session.interrupt",
"prompt.submit",
)
val baseline = ownershipMethods.associateWith { method ->
gatewayHarness.rpcLog.count { it.first == method }
}
val baselineActiveList = gatewayHarness.rpcLog.count { it.first == "session.active_list" }
gatewayHarness.activeSessionListPayload = activeSessionPayload("working")
gatewayScope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
gatewayClient = GatewayChatClient(
initialDashboardClient = DashboardApiClient(
baseUrl = gatewayHarness.server.url("/").toString().trimEnd('/'),
okHttpClient = OkHttpClient(),
),
okHttpClient = OkHttpClient(),
callbackDispatcher = { block -> Handler(Looper.getMainLooper()).post(block) },
scope = gatewayScope,
)
viewModel.setChatTurnCheckpointStore(MemoryCheckpointStore())
viewModel.updateGatewayClient(gatewayClient)
viewModel.setChatVisible(true)
viewModel.prewarmGateway()
awaitCondition {
gatewayHarness.rpcLog.count { it.first == "session.active_list" } > baselineActiveList
}
persistedHistory = listOf(
MessageItem(
id = "desktop-answer",
sessionId = STORED_SESSION_ID,
role = "assistant",
content = JsonPrimitive("Desktop completed without Android attachment."),
),
)
awaitCondition {
handler.messages.value.singleOrNull()?.content ==
"Desktop completed without Android attachment."
}
ownershipMethods.forEach { method ->
assertEquals(
"passive foreground sent $method",
baseline.getValue(method),
gatewayHarness.rpcLog.count { it.first == method },
)
}
viewModel.updateGatewayClient(null)
gatewayClient.shutdown()
assertEquals(
"observer teardown interrupted the Desktop turn",
baseline.getValue("session.interrupt"),
gatewayHarness.rpcLog.count { it.first == "session.interrupt" },
)
}
@Test
fun observerReadyAfterChatHidesCannotRestartPassiveWork() {
viewModel.switchProfileContext(PROFILE_CONTEXT, STORED_SESSION_ID)
awaitCondition { !viewModel.isLoadingHistory.value }
viewModel.setChatVisible(false)
viewModel.updateGatewayClient(null)
gatewayClient.shutdown()
gatewayScope.cancel()
val historyReads = AtomicInteger(0)
viewModel.setProfileMessageLoader {
historyReads.incrementAndGet()
Result.success(persistedHistory)
}
val controlMethods = setOf(
"session.resume",
"session.activate",
"session.interrupt",
"prompt.submit",
)
val baseline = controlMethods.associateWith { method ->
gatewayHarness.rpcLog.count { it.first == method }
}
val baselineActiveList = gatewayHarness.rpcLog.count { it.first == "session.active_list" }
gatewayHarness.suppressGatewayReady = true
gatewayScope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
gatewayClient = GatewayChatClient(
initialDashboardClient = DashboardApiClient(
baseUrl = gatewayHarness.server.url("/").toString().trimEnd('/'),
okHttpClient = OkHttpClient(),
),
okHttpClient = OkHttpClient(),
callbackDispatcher = { block -> Handler(Looper.getMainLooper()).post(block) },
scope = gatewayScope,
)
viewModel.setChatTurnCheckpointStore(MemoryCheckpointStore())
viewModel.updateGatewayClient(gatewayClient)
viewModel.setChatVisible(true)
val delayedSocket = gatewayHarness.awaitServerSocket()
viewModel.setChatVisible(false)
gatewayHarness.sendGatewayReady(delayedSocket)
shadowOf(Looper.getMainLooper()).idleFor(500, TimeUnit.MILLISECONDS)
Thread.sleep(100)
assertEquals(0, historyReads.get())
assertEquals(
baselineActiveList,
gatewayHarness.rpcLog.count { it.first == "session.active_list" },
)
controlMethods.forEach { method ->
assertEquals(baseline.getValue(method), gatewayHarness.rpcLog.count { it.first == method })
}
}
@Test
fun staleHistoryReadCannotEraseATurnCompletedDuringTheFetch() {
val loadCount = AtomicInteger(0)
@@ -0,0 +1,37 @@
package com.hermesandroid.relay.viewmodel
import org.junit.Assert.assertEquals
import org.junit.Test
class PhoneThreadChatIdIndexTest {
@Test
fun `index is visible only to its owning connection`() {
val index = PhoneThreadChatIdIndex(
connectionId = "connection-a",
values = mapOf("session-a" to "reminders"),
)
assertEquals(index.values, visiblePhoneThreadChatIds("connection-a", index))
assertEquals(emptyMap<String, String>(), visiblePhoneThreadChatIds("connection-b", index))
assertEquals(emptyMap<String, String>(), visiblePhoneThreadChatIds(null, index))
}
@Test
fun `a later failed refresh preserves the last successful index`() {
val successful = reconcilePhoneThreadChatIdIndex(
current = PhoneThreadChatIdIndex(),
requestedConnectionId = "connection-a",
activeConnectionId = "connection-a",
fetched = Result.success(mapOf("session-a" to "reminders")),
)
val afterFailure = reconcilePhoneThreadChatIdIndex(
current = successful,
requestedConnectionId = "connection-a",
activeConnectionId = "connection-a",
fetched = Result.failure(IllegalStateException("offline")),
)
assertEquals(successful, afterFailure)
}
}
-33
View File
@@ -3977,36 +3977,3 @@ disappearance, client-side profile isolation, and method-not-found; physical
and current-host certification remains tracked in `TODO.md`. An upstream
profile field/filter or explicitly owned aggregate activity route would remove
the remaining ambiguity for multi-profile clients.
---
## ADR 69 — Passive Android observation never attaches another client's Gateway turn
**Status:** Accepted (2026-08-28).
**Context.** `session.resume` and `session.activate` are live-runtime attachment
operations, not read-only subscriptions. Android previously called
`session.resume` while opening or foregrounding Chat and after loading a saved
session's history. When Desktop/TUI already owned a running turn, that passive
prewarm could rebind the runtime transport to Android. A later Android socket,
route, or client teardown could then strand the producer or promote the foreign
turn into an Android `GatewayTurn` whose cancellation sends `session.interrupt`.
The issue was distinct from the earlier stale-view and missing-terminal recovery
paths, which concern exact Android-owned checkpoints.
**Decision.** Ordinary visibility, foreground restoration, Idle-socket recovery,
and saved-session selection establish only the shared Gateway socket. They use
profile-scoped REST history plus process-wide `session.active_list`; while an
unowned row with the selected durable id is live, Android performs bounded
history refreshes and one final read after settlement. These observer paths send
no `session.resume`, `session.activate`, `prompt.submit`, or `session.interrupt`.
Exact Android-owned checkpoints retain `session.activate` with durable-resume
fallback, and explicit send or session-config actions may resume because the user
is intentionally taking control of that destination.
**Consequences.** Opening Android cannot replace, stop, or later cancel a turn
already running in Desktop/TUI. Live token frames remain with the producing
client; Android observes durable progress and final history without inventing a
multi-subscriber Gateway contract. The first explicit Android mutation may pay
the resume latency that passive prewarm previously hid. Cross-client fixtures
and Android lifecycle coverage enforce the no-control-RPC observation boundary.
-1
View File
@@ -41,7 +41,6 @@ the upstream contract identifiers it depends on.
| `active_status_lifecycle` | `session.active_list` reports starting, working, waiting, and idle, then a complete empty process-wide snapshot permits removal of unambiguously owned prior rows |
| `active_status_profile_scope` | A row has no profile metadata and a caller profile hint has no effect; the client must use exact client-held ownership and reject invented attribution |
| `active_status_unsupported` | An older Gateway returns JSON-RPC method-not-found; the client retains Unknown rather than inventing Idle or Working |
| `cross_client_observation` | A second client observes a Desktop-owned working session through active status and history without resume, activate, submit, or interrupt; the producing client receives the terminal event |
Fixture evidence is a bounded metadata-only ring. It records sequence,
connection number, RPC method, event type, scope classification, and outcome.
+6 -6
View File
@@ -13,7 +13,7 @@
"verification": "ai-translated",
"review_refs": [],
"source_sha256": {
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"main": "3f7ec5aea36744d36ed7e585bf99cce04d8db3a6f9b51399695389b306f40a8e",
"sideload": "4abff4f1069091ec2de735c3037a7ec7d77699cb4321e8511a622437bceaf7c2"
},
"surfaces": {
@@ -48,7 +48,7 @@
"verification": "ai-translated",
"review_refs": [],
"source_sha256": {
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"main": "3f7ec5aea36744d36ed7e585bf99cce04d8db3a6f9b51399695389b306f40a8e",
"sideload": "4abff4f1069091ec2de735c3037a7ec7d77699cb4321e8511a622437bceaf7c2"
},
"surfaces": {
@@ -72,7 +72,7 @@
"verification": "ai-translated",
"review_refs": [],
"source_sha256": {
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"main": "3f7ec5aea36744d36ed7e585bf99cce04d8db3a6f9b51399695389b306f40a8e",
"sideload": "4abff4f1069091ec2de735c3037a7ec7d77699cb4321e8511a622437bceaf7c2"
},
"surfaces": {
@@ -96,7 +96,7 @@
"verification": "ai-translated",
"review_refs": [],
"source_sha256": {
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"main": "3f7ec5aea36744d36ed7e585bf99cce04d8db3a6f9b51399695389b306f40a8e",
"sideload": "4abff4f1069091ec2de735c3037a7ec7d77699cb4321e8511a622437bceaf7c2"
},
"surfaces": {
@@ -120,7 +120,7 @@
"verification": "ai-translated",
"review_refs": [],
"source_sha256": {
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"main": "3f7ec5aea36744d36ed7e585bf99cce04d8db3a6f9b51399695389b306f40a8e",
"sideload": "4abff4f1069091ec2de735c3037a7ec7d77699cb4321e8511a622437bceaf7c2"
},
"surfaces": {
@@ -135,7 +135,7 @@
"verification": "ai-translated",
"review_refs": [],
"source_sha256": {
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"main": "3f7ec5aea36744d36ed7e585bf99cce04d8db3a6f9b51399695389b306f40a8e",
"sideload": "4abff4f1069091ec2de735c3037a7ec7d77699cb4321e8511a622437bceaf7c2"
},
"surfaces": {
+1 -1
View File
@@ -481,7 +481,7 @@ Bottom navigation bar with 4 tabs:
- **Bot group projection** — Android merges the bounded `ui_meta["hermes-bots-groups"]` v3 projection across gateways by durable room identity and newest revision. Rooms and recent messages are visibly read-only; Android does not create, rename, disband, join, send, coordinate member turns, or become a second room-log authority. Binary room images are ignored at this metadata boundary.
- **Session drawer** (swipe from left or hamburger icon) — session list with title, timestamp, message count. Create, switch, rename, delete, pin/unpin, and archive/restore. The process-owned conversation binding is the single connection/profile/session identity for Chat; selecting an All Profiles row atomically makes its owner the selected agent and persists that profile/session, while merely browsing All Profiles changes no agent state. Lifecycle or locale-driven Activity recreation cannot replace an explicit binding with stale persisted state, and asynchronous list/history/mutation work is accepted only for the binding's exact namespace. A profile lock hides All Profiles and rejects stale/deep-linked cross-profile opens. The All Profiles browser mode otherwise survives Activity state restoration and refetches its rows after recreation. Pin and archive are durable upstream session fields loaded and patched through the owning connection/profile's Dashboard session API; Android does not keep a second local flag registry. Archived rows are requested explicitly so they remain restorable after recreation. Failed mutations roll back the optimistic row, while refresh and deletion reconcile from server truth. When a persisted title is absent, use upstream's first-user-message `preview`, matching the Hermes Desktop session picker; show "Untitled" only when neither value exists.
- **Authoritative session activity** — one composite registry keyed by connection, normalized profile, and durable session id drives the drawer, filters, grouping, animation, accessibility, and the visible composer. Exact pending approval/clarify/sudo/secret/MCP requests produce **Needs input**; the Gateway's process-wide `session.active_list` supplies **Starting**, **Working**, and **Idle**; exact terminal or `session.info {running:false}` can settle the matching generation. Because active-list rows normally have no profile metadata, Android assigns a row only through exact foreground/detached ownership already held by that client, or explicit profile metadata if a future upstream sends it. A bounded REST directory never proves global uniqueness. Unresolved rows create no status. Resolved rows from a partial snapshot may update their exact owners, but disappearance settles a scope only when the successful process-wide snapshot was completely and unambiguously resolved for it. Restart/checkpoint recovery is **Checking**; a failed or unsupported live refresh is **Unavailable**, never inferred Idle. REST `is_active` remains recency metadata only. `process.list` may add a separate **Background work** indicator and never keeps the parent conversation Working. Old socket generations, bare session ids from another profile, and delayed snapshots cannot revive newer settled state.
- **Concurrent Gateway chats** — switching sessions, profiles, drafts, or Threads detaches the visible Android-owned turn without sending `session.interrupt`; each Android-owned running chat keeps a connection/profile/session-scoped checkpoint and reattaches to its live Gateway session when reopened. Opening, foregrounding, or selecting a saved session without that exact checkpoint is read-only observation: Android warms only the socket, reads profile-scoped history, and polls `session.active_list` without `session.resume`, `session.activate`, `prompt.submit`, or `session.interrupt`. A Desktop/TUI-owned turn therefore remains owned by its producing client; Android refreshes persisted progress and performs one final history read when the runtime settles. Explicit send/config actions may resume the destination session, explicit Stop still interrupts, and SSE fallback stays single-stream and cancels on navigation.
- **Concurrent Gateway chats** — switching sessions, profiles, drafts, or Threads detaches the visible turn without sending `session.interrupt`; each running chat keeps a connection/profile/session-scoped checkpoint and reattaches to its live Gateway session when reopened. Explicit Stop still interrupts. SSE fallback stays single-stream and cancels on navigation.
- **Queued Gateway follow-ups** — every local queued item is immutably scoped to its originating connection, profile, stored session, transport, and run generation; only that run's completion can make it eligible, and switching sessions shows only that session's queue. Restored text queues retain the same scope, while unavailable/deleted destinations and non-restorable attachment queues fail visibly instead of following the current composer. Drained messages add `queued: true` to `prompt.submit`; ordinary sends omit the field. Authoritative submit rejections (`4004`, `4018`, `4028`, `4029`, `4030`, `4090`, `5008`, `5070`, and `5071`) preserve the server message and never fall through to API-server SSE.
- **Durable composer drafts** — each connection/profile/session owns one app-private draft containing text, quote/edit context, and pending attachment bytes. Metadata and content-addressed blobs live under Android's no-backup directory, are capped at 64 drafts and 128 MB of retained blobs outside the active draft, flush when Chat backgrounds, and are removed after a successful send. Session/profile/connection navigation saves the previous owner before restoring the destination; an opened cross-profile session uses its actual owning profile rather than the global picker.
- **Large paste review** — a default-on Chat setting converts any single insertion of at least 5,000 characters into a visible `pasted-text.txt` attachment before the normal message-length limit rejects it. Gateway uses upstream `file.attach`; API-server SSE and proactive Thread paths materialize the same UTF-8 text into the outgoing prompt and remove only the synthetic attachment from that transport, so the behavior never requires Relay or silently drops content.
+29
View File
@@ -28,6 +28,7 @@ import json
import unittest
from typing import Any
from plugin.phone_platform import _normalize_reply
from plugin.relay.channels.proactive import ProactiveChannel, ProactiveError
@@ -336,6 +337,34 @@ class ProactiveChannelTests(unittest.TestCase):
_run(run())
def test_custom_chat_id_survives_relay_drain_and_adapter_normalization(self) -> None:
async def run() -> None:
ch = ProactiveChannel()
ws = _FakeWs()
await ch.handle(
ws,
{
"type": "proactive.reply",
"payload": {
"text": "continue this thread",
"chat_id": "thread-project-461",
"reply_to": "prompt-1",
"message_id": "reply-1",
},
},
)
replies = await ch.take_replies(timeout=0.1)
self.assertEqual(len(replies), 1)
normalized = _normalize_reply(replies[0], "configured-home")
self.assertIsNotNone(normalized)
assert normalized is not None
self.assertEqual(normalized["chat_id"], "thread-project-461")
self.assertEqual(normalized["reply_to"], "prompt-1")
self.assertEqual(normalized["message_id"], "reply-1")
_run(run())
def test_reply_empty_text_dropped(self) -> None:
async def run() -> None:
ch = ProactiveChannel()
@@ -100,47 +100,6 @@ class FixtureTestCase(unittest.IsolatedAsyncioTestCase):
self.assertEqual(["user", "assistant"], [row["role"] for row in history["messages"]])
self.assertEqual(2, history["pagination"]["returned"])
async def test_cross_client_observer_never_claims_or_interrupts_producer(self) -> None:
fixture, base_url = await self.start("cross_client_observation")
producer, _ = await self.connect(base_url)
await self.rpc(producer, 1, "session.resume", {"session_id": fixture.scenario.stored_session_id})
await producer.receive_json()
await self.rpc(producer, 2, "prompt.submit", {"text": "producer-only content"})
producer_frames = await self.frames_until(
producer,
lambda frame: frame.get("params", {}).get("type") == "message.delta",
)
observer, _ = await self.connect(base_url)
await self.rpc(observer, 3, "session.active_list")
active = (await observer.receive_json())["result"]["sessions"]
self.assertEqual("working", active[0]["status"])
async with self.session.get(
f"{base_url}/api/sessions/{fixture.scenario.stored_session_id}/messages",
params={"profile": "default", "limit": 500, "offset": 0, "order": "asc"},
) as response:
self.assertEqual(200, response.status)
self.assertIsInstance((await response.json())["messages"], list)
await observer.close()
producer_frames += await self.frames_until(
producer,
lambda frame: frame.get("params", {}).get("type") == "message.complete",
)
self.assertIn(
"message.complete",
[frame.get("params", {}).get("type") for frame in producer_frames],
)
async with self.session.get(f"{base_url}/__fixture__/evidence") as response:
evidence = await response.json()
observer_methods = [
entry.get("method")
for entry in evidence["entries"]
if entry.get("kind") == "rpc" and entry.get("connection") == 2
]
self.assertEqual(["session.active_list"], observer_methods)
self.assertNotIn("session.interrupt", observer_methods)
async def test_rapid_chunks_tools_and_interims_keep_wire_order(self) -> None:
_, base_url = await self.start("rapid_tools_interims")
ws, _ = await self.connect(base_url)
@@ -306,7 +265,6 @@ class ScenarioTestCase(unittest.TestCase):
"active_status_lifecycle",
"active_status_profile_scope",
"active_status_unsupported",
"cross_client_observation",
"ordinary_turn",
"rapid_tools_interims",
"terminal_gap_activate",
@@ -346,10 +304,6 @@ class ScenarioTestCase(unittest.TestCase):
("gateway.settled_session_info",),
load_scenario("terminal_gap_session_info").contract_requirements,
)
self.assertEqual(
("gateway.message_complete", "gateway.session_active_list"),
load_scenario("cross_client_observation").contract_requirements,
)
def test_tls_arguments_must_be_paired(self) -> None:
with contextlib.redirect_stderr(io.StringIO()):
@@ -1,40 +0,0 @@
{
"name": "cross_client_observation",
"live_session_id": "fixture-desktop-live",
"stored_session_id": "fixture-shared-session",
"profile": "default",
"contract_requirements": [
"gateway.message_complete",
"gateway.session_active_list"
],
"initial_history": [],
"turns": [
{
"steps": [
{"op": "event", "type": "message.start"},
{"op": "event", "type": "message.delta", "payload": {"text": "Desktop still owns this turn."}},
{"op": "sleep", "milliseconds": 250},
{
"op": "persist",
"messages": [
{"id": 1, "role": "user", "content": "Desktop prompt.", "timestamp": 1.0},
{"id": 2, "role": "assistant", "content": "Desktop still owns this turn.", "timestamp": 2.0, "finish_reason": "stop"}
]
},
{"op": "set_running", "value": false},
{"op": "event", "type": "message.complete", "payload": {"text": "Desktop still owns this turn.", "status": "complete"}}
]
}
],
"active_list": {
"supported": true,
"snapshots": [
[
{"id": "fixture-desktop-live", "session_key": "fixture-shared-session", "status": "working", "current": false}
],
[
{"id": "fixture-desktop-live", "session_key": "fixture-shared-session", "status": "idle", "current": false}
]
]
}
}