Compare commits

..
Author SHA1 Message Date
Bailey Dixon d80b3db329 fix(android): keep passive gateway observation read-only 2026-08-28 23:40:13 -04:00
42 changed files with 734 additions and 3185 deletions
+1 -1
View File
@@ -8,7 +8,6 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/), and this
### Added
- **Android can preview delegated agent work without leaving the parent chat.** The current-chat activity sheet shows bounded lifecycle, progress, and tool previews for concurrent children, opens vanilla Hermes child history read-only when the Gateway exposes it, and stays explicit when reconnect gaps or older routes leave details unavailable.
- **Android presents Relay Git as a first-class native workspace.** A compact optional Chat rail opens repository status, line totals, filters, diffs, branches, staging, commits, and remotes; the full workspace remains available from Settings when Chat controls are hidden.
- **Hermes-Relay Plugin provides a bounded Git workspace API for authenticated Dashboard clients.** Configured repository roots, path validation, tracked line totals, scoped write grants, and explicit confirmation protect repository reads and mutations.
@@ -18,6 +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.
- **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,6 +26,9 @@ 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,6 +239,60 @@ 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"
@@ -263,6 +317,9 @@ 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)
@@ -282,6 +339,18 @@ 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())
}
@@ -345,6 +414,15 @@ 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 }
@@ -353,4 +431,9 @@ 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"
}
}
@@ -10,16 +10,6 @@ package com.hermesandroid.relay.data
*/
object AgentDisplay {
const val SERVER_DEFAULT_PROFILE_KEY: String = "__server_default__"
private const val PROFILE_CONTEXT_SEPARATOR = "::"
data class ProfileContextIdentity(
val connectionId: String,
val profileKey: String,
) {
/** Null means the upstream request must inherit Server Default. */
val requestProfileName: String?
get() = profileRequestName(profileKey)
}
private val GENERIC_MODEL_ALIASES = setOf(
"hermes-agent",
"hermes_agent",
@@ -173,25 +163,7 @@ object AgentDisplay {
profileRequestName(profileName) ?: SERVER_DEFAULT_PROFILE_KEY
fun profileContextKey(connectionId: String?, profileName: String?): String =
"${connectionId.orEmpty()}$PROFILE_CONTEXT_SEPARATOR${profileSessionKey(profileName)}"
/**
* Parse the canonical profile/context identity used by persisted chat state.
*
* Legacy or malformed opaque keys deliberately return null: recovery may
* still use the exact key for ownership, but must not invent an upstream
* profile override from it. The first separator is authoritative so legal
* profile names containing `::` remain round-trippable.
*/
fun parseProfileContextKey(contextKey: String?): ProfileContextIdentity? {
val raw = contextKey?.trim().orEmpty()
val separator = raw.indexOf(PROFILE_CONTEXT_SEPARATOR)
if (separator <= 0 || separator + PROFILE_CONTEXT_SEPARATOR.length >= raw.length) return null
val connectionId = raw.substring(0, separator).trim()
val profileKey = raw.substring(separator + PROFILE_CONTEXT_SEPARATOR.length).trim()
if (connectionId.isEmpty() || profileKey.isEmpty()) return null
return ProfileContextIdentity(connectionId, profileKey)
}
"${connectionId.orEmpty()}::${profileSessionKey(profileName)}"
fun localDisplayAlias(value: String?): String? =
value
@@ -23,8 +23,6 @@ import kotlinx.serialization.json.Json
data class ChatTurnCheckpoint(
val schemaVersion: Int = CURRENT_SCHEMA,
val contextKey: String,
/** Explicit persisted profile identity; null only for legacy checkpoints. */
val profileKey: String? = null,
val sessionId: String,
val liveSessionId: String? = null,
val transport: String,
@@ -1006,56 +1006,6 @@ class ChatHandler {
}
}
/**
* Bound the ephemeral, read-only child-watch projection. This is stricter
* than the main transcript: system rows and tool results are not part of
* the preview contract, and one live child must not retain unbounded text.
*/
internal fun boundReadOnlyPreview(
maxMessages: Int = 100,
maxTotalChars: Int = 32_000,
maxFieldChars: Int = 8_000,
maxToolChars: Int = 1_000,
): Boolean {
var truncated = false
_messages.update { current ->
val visible = current.filterNot { it.role == MessageRole.SYSTEM }
if (visible.size != current.size || visible.size > maxMessages) truncated = true
var remaining = maxTotalChars
val kept = mutableListOf<ChatMessage>()
visible.takeLast(maxMessages).asReversed().forEach { message ->
if (remaining <= 0) {
truncated = true
return@forEach
}
fun bounded(value: String, limit: Int): String {
val allowed = minOf(limit, remaining)
val next = value.takeLast(allowed)
if (next.length != value.length) truncated = true
remaining -= next.length
return next
}
val content = bounded(message.content, maxFieldChars)
val thinking = bounded(message.thinkingContent, maxFieldChars)
val tools = message.toolCalls.takeLast(50).map { tool ->
if (message.toolCalls.size > 50) truncated = true
tool.copy(
args = tool.args?.let { bounded(it, maxToolChars) },
result = null,
error = tool.error?.let { bounded(it, maxToolChars) },
)
}
kept += message.copy(
content = content,
thinkingContent = thinking,
toolCalls = tools,
)
}
kept.asReversed()
}
return truncated
}
/**
* Rehydrate the last client-owned state of an unfinished turn.
*
@@ -3080,9 +3030,7 @@ class ChatHandler {
fun onSubagentEvent(messageId: String, event: GatewaySubagentEvent) {
val label = event.goal.trim().take(60).ifBlank { null }
when (event.phase) {
GatewaySubagentEvent.Phase.SPAWN_REQUESTED,
GatewaySubagentEvent.Phase.START,
-> {
GatewaySubagentEvent.Phase.START -> {
if (label != null) subagentLabels[event.taskIndex] = label
event.subagentId?.takeIf(String::isNotBlank)?.let {
subagentIds[event.taskIndex] = it
@@ -209,9 +209,6 @@ class GatewayChatClient(
private const val INBOUND_BIND_TIMEOUT_MS = 2_000L
private const val CANCELLED_TURN_SUBMIT_WAIT_MS = 2_000L
private const val MAX_RECOVERY_BUFFERED_EVENTS = 256
internal const val MAX_CHILD_WATCH_HISTORY_ITEMS = 200
internal const val MAX_CHILD_WATCH_HISTORY_CHARS = 64_000
private const val MAX_PENDING_CHILD_WATCH_EVENTS = 256
/** Distinct socket-loss (flap) events per turn we'll try to recover from. */
private const val MAX_TURN_REJOINS = 4
@@ -384,15 +381,6 @@ class GatewayChatClient(
private val prewarmRequestGeneration = AtomicLong(0)
private val pendingRpcs = ConcurrentHashMap<Long, CompletableDeferred<JsonObject>>()
/** Monotonic client-local fence for lazy child watch open/close races. */
private val childWatchGeneration = AtomicLong(0)
/** Live child runtime id -> exact watcher that owns its callbacks. */
private val childWatches = ConcurrentHashMap<String, ChildWatchRegistration>()
/** Events that race a lazy `session.resume` acknowledgement. */
private val pendingChildWatchOpens = ConcurrentHashMap<Long, PendingChildWatchOpen>()
/** Live (per-connection) session id ←→ the stored DB id it was resumed/created from. */
@Volatile
private var liveSessionId: String? = null
@@ -459,62 +447,6 @@ class GatewayChatClient(
@Volatile var pendingAsk: GatewayAsk? = null,
)
private class ChildWatchRegistration(
val storedSessionId: String,
val liveSessionId: String,
val profile: String?,
val generation: Long,
val callbacks: GatewayTurnCallbacks,
) {
lateinit var mapper: GatewayEventMapper
}
private data class ChildWatchEvent(
val sessionId: String,
val type: String,
val payload: JsonObject?,
)
private data class PendingChildWatchReplay(
val events: List<ChildWatchEvent>,
val truncated: Boolean,
)
private class PendingChildWatchOpen {
private val lock = Any()
private val events = mutableListOf<ChildWatchEvent>()
private var closed = false
private var truncated = false
fun capture(event: ChildWatchEvent): Boolean = synchronized(lock) {
if (closed) return@synchronized false
if (events.size >= MAX_PENDING_CHILD_WATCH_EVENTS) {
events.removeAt(0)
truncated = true
}
events += event
true
}
fun closeAndTake(sessionId: String): PendingChildWatchReplay = synchronized(lock) {
closed = true
PendingChildWatchReplay(
events = events.filter { it.sessionId == sessionId },
truncated = truncated,
).also { events.clear() }
}
fun close() = synchronized(lock) {
closed = true
events.clear()
}
}
private data class BoundedChildHistory(
val messages: List<MessageItem>,
val truncated: Boolean,
)
/**
* Upstream may emit the interrupted turn's tail and terminal event after
* `session.interrupt` returns. Keep a short exact-session tombstone so that
@@ -974,6 +906,30 @@ 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
@@ -1015,200 +971,6 @@ class GatewayChatClient(
return sessionReady
}
/**
* Open the vanilla-upstream child-session watcher advertised by
* `subagent.*.child_session_id`. This RPC deliberately does not mutate
* [liveSessionId], [storedSessionId], or [liveSessionProfile]: the parent
* conversation keeps owning the main mapper while the returned short live
* id routes a second, read-only event stream on the same socket.
*
* Returned history is bounded locally even when an upstream gateway sends
* the child's entire transcript in the resume acknowledgement. The server
* may still enforce its own larger resume safety limit before replying.
*/
suspend fun openChildWatch(
childSessionId: String,
profile: String? = currentSessionProfile(),
callbacks: GatewayTurnCallbacks,
historyLimit: Int = MAX_CHILD_WATCH_HISTORY_ITEMS,
): Result<GatewayChildWatch> = runCatching {
val storedChildId = childSessionId.trim()
require(storedChildId.isNotEmpty()) { "child session id is required" }
val requestedProfile = profile?.trim()?.takeIf(String::isNotEmpty)
connectMutex.withLock {
// Allocate and register under the same mutex as the resume RPC so
// concurrent opens complete in generation order; an older caller
// can never `put` after a newer one for the same live child id.
val generation = childWatchGeneration.incrementAndGet()
val pending = PendingChildWatchOpen()
pendingChildWatchOpens[generation] = pending
try {
ensureConnected()
val result = rpc(
"session.resume",
buildJsonObject {
put("session_id", storedChildId)
put("cols", DEFAULT_COLS)
put("source", sessionSource)
put("lazy", true)
put("close_on_disconnect", true)
requestedProfile?.let { put("profile", it) }
},
).getOrElse { error ->
throw GatewayPreflightException(
"child session resume failed: ${error.message}",
)
}
val liveChildId = result.stringField("session_id")?.takeIf(String::isNotBlank)
?: throw GatewayPreflightException(
"child session resume returned no live session id",
)
try {
requireConfirmedSessionProfile(result, requestedProfile)
} catch (error: GatewayPreflightException) {
// The wrong profile must not leave an unowned lazy watcher behind.
rpc(
"session.close",
buildJsonObject { put("session_id", liveChildId) },
)
throw error
}
val registration = ChildWatchRegistration(
storedSessionId = storedChildId,
liveSessionId = liveChildId,
profile = requestedProfile,
generation = generation,
callbacks = callbacks,
)
val dispatchedCallbacks = dispatchOn(callbacks) {
childWatches[liveChildId] === registration
}
registration.mapper = GatewayEventMapper(
dispatchedCallbacks,
dedupeAdjacentMessageStarts = true,
)
childWatches.put(liveChildId, registration)?.let { prior ->
if (prior.generation != generation) {
notifyChildWatchFailure(
prior,
"Child watch was replaced by a newer view",
)
}
}
// Replay only frames tagged with the exact live id returned by
// this resume. Unknown gateway sessions captured during the
// narrow ack race remain foreign and are discarded.
val replay = pending.closeAndTake(liveChildId)
if (replay.truncated) dispatchedCallbacks.onReconcileRequired()
replay.events.forEach { event ->
if (childWatches[liveChildId] === registration) {
registration.mapper.onEvent(event.type, event.payload)
}
}
val replayedTerminal = replay.events.any {
it.type == "message.complete" || it.type == "error"
}
val history = parseChildWatchMessages(result, historyLimit)
GatewayChildWatch(
storedSessionId = storedChildId,
liveSessionId = liveChildId,
profile = requestedProfile,
generation = generation,
messages = history.messages,
historyTruncated = history.truncated,
running = !replayedTerminal && result.booleanField("running") == true,
status = if (replayedTerminal) "idle" else result.stringField("status"),
)
} finally {
pendingChildWatchOpens.remove(generation, pending)
pending.close()
}
}
}
/**
* Close only the exact lazy watcher represented by [watch]. A stale handle
* is a no-op so it can never close a newer watcher whose live id was reused.
* The parent session and delegated child continue running server-side.
*/
suspend fun closeChildWatch(watch: GatewayChildWatch): Result<Unit> =
connectMutex.withLock {
val registration = childWatches[watch.liveSessionId]
?: return@withLock Result.success(Unit)
if (
registration.generation != watch.generation ||
registration.storedSessionId != watch.storedSessionId ||
registration.profile != watch.profile ||
!childWatches.remove(watch.liveSessionId, registration)
) {
return@withLock Result.success(Unit)
}
if (webSocket == null || readySignal?.isCompleted != true) {
return@withLock Result.success(Unit)
}
val result = rpc(
"session.close",
buildJsonObject { put("session_id", watch.liveSessionId) },
)
result.fold(
onSuccess = { Result.success(Unit) },
onFailure = { error ->
// Permit an exact-handle retry. Opens share connectMutex,
// so no newer registration can race this restoration.
childWatches.putIfAbsent(watch.liveSessionId, registration)
Result.failure(error)
},
)
}
private fun parseChildWatchMessages(
result: JsonObject,
requestedLimit: Int,
): BoundedChildHistory {
val limit = requestedLimit.coerceIn(1, MAX_CHILD_WATCH_HISTORY_ITEMS)
val all = (result["messages"] as? JsonArray).orEmpty()
val raw = all.takeLast(limit)
var retainedChars = 0
var truncated = all.size > raw.size
val newestFirst = raw.asReversed().mapNotNull { element ->
val message = element as? JsonObject ?: run {
truncated = true
return@mapNotNull null
}
// Gateway display history uses `text`; the shared session DTO uses
// `content`. Normalize only that projection boundary.
val normalized = JsonObject(message.toMutableMap().apply {
if (!containsKey("content")) {
put("content", message["text"] ?: message["context"] ?: JsonNull)
}
if (!containsKey("tool_name") && message.containsKey("name")) {
put("tool_name", message["name"] ?: JsonNull)
}
})
val serializedChars = normalized.toString().length
if (serializedChars > MAX_CHILD_WATCH_HISTORY_CHARS - retainedChars) {
truncated = true
return@mapNotNull null
}
val decoded = runCatching {
json.decodeFromJsonElement(MessageItem.serializer(), normalized)
}.onFailure {
Log.w(TAG, "child watch returned an unreadable history row", it)
}.getOrNull()
if (decoded == null) {
truncated = true
null
} else {
retainedChars += serializedChars
decoded
}
}
return BoundedChildHistory(newestFirst.asReversed(), truncated)
}
/**
* Obtain a session-scoped target before a model-selection `config.set`.
*
@@ -3284,9 +3046,6 @@ class GatewayChatClient(
val reason = payload?.stringField("reason")
val supportedReason = reason in setOf("idle_timeout", "lru_evict", "ws_orphan_reap")
if (!reclaimedLiveId.isNullOrBlank() && supportedReason) {
childWatches.remove(reclaimedLiveId)?.let { registration ->
notifyChildWatchFailure(registration, "Gateway reclaimed the child watch")
}
val background = backgroundTurns.remove(reclaimedLiveId)
if (background != null) {
callbackDispatcher {
@@ -3338,23 +3097,6 @@ class GatewayChatClient(
return
}
// A lazy child watcher is a second session on this shared socket. Route
// it before the main-session recovery/foreign-session gates and require
// the exact live id returned by its own session.resume acknowledgement.
val childWatch = eventSessionId?.let(childWatches::get)
if (childWatch != null) {
if (childWatches[eventSessionId] === childWatch) {
childWatch.mapper.onEvent(type, payload)
}
return
}
val capturedForPendingChildWatch = !eventSessionId.isNullOrBlank() &&
eventSessionId != liveSessionId &&
!backgroundTurns.containsKey(eventSessionId) &&
capturePendingChildWatchEvent(ChildWatchEvent(eventSessionId, type, payload))
if (capturedForPendingChildWatch) return
// A cold session.resume may schedule auto-continue before its RPC
// response reaches Android. The recovery buffer is an ownership gate,
// not an observational copy: an event is either claimed here for
@@ -3597,7 +3339,6 @@ class GatewayChatClient(
it.completeExceptionally(GatewayRpcException("gateway connection lost"))
}
pendingRpcs.clear()
failChildWatches("Child watch disconnected from the gateway")
val turn = activeTurn
if (turn == null) {
if (backgroundTurns.isNotEmpty() && !backgroundRejoinInProgress) {
@@ -3780,7 +3521,6 @@ class GatewayChatClient(
_activeSessionCapability.value = GatewayActiveSessionCapability.Unknown
_approvalModeCapability.value = GatewayApprovalModeCapability.Unknown
_connectionState.value = GatewayConnectionState.Idle
failChildWatches("Child watch closed with the gateway socket")
}
private fun scheduleBackgroundClose() {
@@ -3790,7 +3530,7 @@ class GatewayChatClient(
backgroundCloseJob?.cancel()
backgroundCloseJob = scope.launch {
delay(BACKGROUND_CLOSE_GRACE_MS)
if (activeTurn == null && backgroundTurns.isEmpty() && childWatches.isEmpty() &&
if (activeTurn == null && backgroundTurns.isEmpty() &&
!AppForegroundTracker.isForeground.value
) {
closeSocket("app backgrounded")
@@ -3798,28 +3538,6 @@ class GatewayChatClient(
}
}
private fun failChildWatches(message: String) {
if (childWatches.isEmpty()) return
val registrations = childWatches.values.toSet()
childWatches.clear()
registrations.forEach { notifyChildWatchFailure(it, message) }
}
private fun notifyChildWatchFailure(
registration: ChildWatchRegistration,
message: String,
) {
callbackDispatcher { registration.callbacks.onResumeFailure(message) }
}
private fun capturePendingChildWatchEvent(event: ChildWatchEvent): Boolean {
var captured = false
pendingChildWatchOpens.values.forEach { pending ->
if (pending.capture(event)) captured = true
}
return captured
}
// ------------------------------------------------------------------
// JSON-RPC
// ------------------------------------------------------------------
@@ -4368,55 +4086,44 @@ class GatewayChatClient(
}
/** Wrap callbacks so every invocation lands on the callback dispatcher (main thread). */
private fun dispatchOn(
callbacks: GatewayTurnCallbacks,
stillCurrent: () -> Boolean = { true },
) = GatewayTurnCallbacks(
onSessionId = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onSessionId(v) } },
onStart = { dispatchIfCurrent(stillCurrent) { callbacks.onStart() } },
onTextDelta = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onTextDelta(v) } },
private fun dispatchOn(callbacks: GatewayTurnCallbacks) = GatewayTurnCallbacks(
onSessionId = { v -> callbackDispatcher { callbacks.onSessionId(v) } },
onStart = { callbackDispatcher { callbacks.onStart() } },
onTextDelta = { v -> callbackDispatcher { callbacks.onTextDelta(v) } },
onInterimMessage = { text, alreadyStreamed ->
dispatchIfCurrent(stillCurrent) { callbacks.onInterimMessage(text, alreadyStreamed) }
callbackDispatcher { callbacks.onInterimMessage(text, alreadyStreamed) }
},
onInterimReconciled = { text ->
dispatchIfCurrent(stillCurrent) { callbacks.onInterimReconciled(text) }
callbackDispatcher { callbacks.onInterimReconciled(text) }
},
onThinkingDelta = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onThinkingDelta(v) } },
onThinkingDelta = { v -> callbackDispatcher { callbacks.onThinkingDelta(v) } },
onToolCallStart = { id, name, args ->
dispatchIfCurrent(stillCurrent) { callbacks.onToolCallStart(id, name, args) }
callbackDispatcher { callbacks.onToolCallStart(id, name, args) }
},
onToolCallDone = { a, b -> dispatchIfCurrent(stillCurrent) { callbacks.onToolCallDone(a, b) } },
onToolCallFailed = { a, b -> dispatchIfCurrent(stillCurrent) { callbacks.onToolCallFailed(a, b) } },
onToolOutputRisk = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onToolOutputRisk(v) } },
onTurnComplete = { dispatchIfCurrent(stillCurrent) { callbacks.onTurnComplete() } },
onReconcileRequired = { dispatchIfCurrent(stillCurrent) { callbacks.onReconcileRequired() } },
onComplete = { dispatchIfCurrent(stillCurrent) { callbacks.onComplete() } },
onUsage = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onUsage(v) } },
onError = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onError(v) } },
onToolGenerating = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onToolGenerating(v) } },
onSubagentEvent = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onSubagentEvent(v) } },
onMoaReference = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onMoaReference(v) } },
onInteractionRequest = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onInteractionRequest(v) } },
onInteractionExpired = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onInteractionExpired(v) } },
onResumeFailure = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onResumeFailure(v) } },
onFailure = { v -> dispatchIfCurrent(stillCurrent) { callbacks.onFailure(v) } },
onToolCallDone = { a, b -> callbackDispatcher { callbacks.onToolCallDone(a, b) } },
onToolCallFailed = { a, b -> callbackDispatcher { callbacks.onToolCallFailed(a, b) } },
onToolOutputRisk = { v -> callbackDispatcher { callbacks.onToolOutputRisk(v) } },
onTurnComplete = { callbackDispatcher { callbacks.onTurnComplete() } },
onReconcileRequired = { callbackDispatcher { callbacks.onReconcileRequired() } },
onComplete = { callbackDispatcher { callbacks.onComplete() } },
onUsage = { v -> callbackDispatcher { callbacks.onUsage(v) } },
onError = { v -> callbackDispatcher { callbacks.onError(v) } },
onToolGenerating = { v -> callbackDispatcher { callbacks.onToolGenerating(v) } },
onSubagentEvent = { v -> callbackDispatcher { callbacks.onSubagentEvent(v) } },
onMoaReference = { v -> callbackDispatcher { callbacks.onMoaReference(v) } },
onInteractionRequest = { v -> callbackDispatcher { callbacks.onInteractionRequest(v) } },
onInteractionExpired = { v -> callbackDispatcher { callbacks.onInteractionExpired(v) } },
onResumeFailure = { v -> callbackDispatcher { callbacks.onResumeFailure(v) } },
onFailure = { v -> callbackDispatcher { callbacks.onFailure(v) } },
// MUST be wrapped like every other member: GatewayTurnCallbacks gives
// onStatusUpdate a default no-op, so omitting it here silently swallows
// EVERY gateway status line — the ❌ terminal-error lifecycle update
// included. Without it markError never fires, the turn isn't badged
// "Error", and onComplete's history reload wipes the error bubble (the
// "reply appears then vanishes" bug).
onStatusUpdate = { kind, text ->
dispatchIfCurrent(stillCurrent) { callbacks.onStatusUpdate(kind, text) }
},
onStatusClear = { kind -> dispatchIfCurrent(stillCurrent) { callbacks.onStatusClear(kind) } },
onStatusUpdate = { kind, text -> callbackDispatcher { callbacks.onStatusUpdate(kind, text) } },
onStatusClear = { kind -> callbackDispatcher { callbacks.onStatusClear(kind) } },
)
private fun dispatchIfCurrent(stillCurrent: () -> Boolean, callback: () -> Unit) {
callbackDispatcher {
if (stillCurrent()) callback()
}
}
}
internal fun parseGatewayPersonalityOptions(result: JsonObject): List<String> =
@@ -281,12 +281,11 @@ class GatewayEventMapper(
callbacks.onError(payload.string("message") ?: "Gateway error")
}
"subagent.spawn_requested", "subagent.start", "subagent.thinking", "subagent.tool",
"subagent.start", "subagent.thinking", "subagent.tool",
"subagent.progress", "subagent.complete",
-> {
clearActivityStatuses()
val phase = when (type) {
"subagent.spawn_requested" -> GatewaySubagentEvent.Phase.SPAWN_REQUESTED
"subagent.start" -> GatewaySubagentEvent.Phase.START
"subagent.thinking" -> GatewaySubagentEvent.Phase.THINKING
"subagent.tool" -> GatewaySubagentEvent.Phase.TOOL
@@ -307,10 +306,6 @@ class GatewayEventMapper(
preview = payload.string("tool_preview") ?: payload.string("text"),
durationSeconds = payload.double("duration_seconds"),
subagentId = payload.string("subagent_id"),
childSessionId = payload.string("child_session_id"),
parentId = payload.string("parent_id"),
depth = payload.int("depth"),
model = payload.string("model"),
),
)
}
@@ -1,6 +1,5 @@
package com.hermesandroid.relay.network.upstream
import com.hermesandroid.relay.network.upstream.models.MessageItem
import com.hermesandroid.relay.network.upstream.models.UsageInfo
import kotlinx.serialization.json.JsonArray
import kotlinx.serialization.json.JsonObject
@@ -258,8 +257,7 @@ data class GatewayToolOutputRisk(
/**
* One `subagent.*` lifecycle event, emitted on the PARENT session. Lifecycle
* per task: SPAWN_REQUESTED → START → (THINKING | TOOL | PROGRESS)* →
* COMPLETE. Field
* per task: START → (THINKING | TOOL | PROGRESS)* → COMPLETE. Field
* availability varies by phase — [toolName]/[preview] ride TOOL,
* [status]/[summary]/[durationSeconds] ride COMPLETE — and older emitters
* omit everything beyond the three defaults-bearing fields.
@@ -275,36 +273,10 @@ data class GatewaySubagentEvent(
val preview: String? = null,
val durationSeconds: Double? = null,
val subagentId: String? = null,
/** Durable child session id accepted by `session.resume {lazy:true}`. */
val childSessionId: String? = null,
/** Owning subagent id for nested delegation; null for first-level children. */
val parentId: String? = null,
/** Zero-based depth used by the upstream spawn-tree renderer. */
val depth: Int? = null,
/** Effective child model, when the emitter exposes it. */
val model: String? = null,
) {
enum class Phase { SPAWN_REQUESTED, START, THINKING, TOOL, PROGRESS, COMPLETE }
enum class Phase { START, THINKING, TOOL, PROGRESS, COMPLETE }
}
/**
* One profile-pinned, read-only child-session watch opened through the vanilla
* upstream Gateway. [storedSessionId] is the durable child id from
* `subagent.*`; [liveSessionId] is the short runtime id that tags subsequent
* mirror events on this socket. The bounded [messages] snapshot is child-only.
*/
data class GatewayChildWatch(
val storedSessionId: String,
val liveSessionId: String,
val profile: String?,
val generation: Long,
val messages: List<MessageItem>,
/** True when Android retained only a bounded recent tail of the response. */
val historyTruncated: Boolean,
val running: Boolean,
val status: String?,
)
/**
* One session-owned background process returned by the upstream gateway's
* `process.list` RPC. The registry calls its process id `session_id`; Android
@@ -20,17 +20,13 @@ import androidx.compose.foundation.shape.CircleShape
import androidx.compose.foundation.shape.RoundedCornerShape
import androidx.compose.foundation.text.selection.SelectionContainer
import androidx.compose.material.icons.Icons
import androidx.compose.material.icons.filled.AccountTree
import androidx.compose.material.icons.filled.CheckCircle
import androidx.compose.material.icons.filled.Close
import androidx.compose.material.icons.filled.ErrorOutline
import androidx.compose.material.icons.filled.ExpandLess
import androidx.compose.material.icons.filled.ExpandMore
import androidx.compose.material.icons.filled.Refresh
import androidx.compose.material.icons.filled.PauseCircleOutline
import androidx.compose.material.icons.filled.Stop
import androidx.compose.material.icons.filled.Terminal
import androidx.compose.material.icons.filled.Visibility
import androidx.compose.material3.CircularProgressIndicator
import androidx.compose.material3.ExperimentalMaterial3Api
import androidx.compose.material3.HorizontalDivider
@@ -43,13 +39,9 @@ import androidx.compose.material3.Text
import androidx.compose.material3.TextButton
import androidx.compose.material3.rememberModalBottomSheetState
import androidx.compose.runtime.Composable
import androidx.compose.runtime.LaunchedEffect
import androidx.compose.runtime.derivedStateOf
import androidx.compose.runtime.getValue
import androidx.compose.runtime.mutableStateOf
import androidx.compose.runtime.remember
import androidx.compose.runtime.rememberCoroutineScope
import androidx.compose.runtime.snapshotFlow
import androidx.compose.runtime.setValue
import androidx.compose.ui.Alignment
import androidx.compose.ui.Modifier
@@ -57,7 +49,6 @@ import androidx.compose.ui.draw.clip
import androidx.compose.ui.graphics.Color
import androidx.compose.ui.res.stringResource
import androidx.compose.ui.semantics.contentDescription
import androidx.compose.ui.semantics.heading
import androidx.compose.ui.semantics.semantics
import androidx.compose.ui.semantics.stateDescription
import androidx.compose.ui.text.font.FontFamily
@@ -65,10 +56,6 @@ import androidx.compose.ui.text.style.TextOverflow
import androidx.compose.ui.unit.dp
import com.hermesandroid.relay.R
import com.hermesandroid.relay.network.upstream.GatewayProcess
import com.hermesandroid.relay.viewmodel.SubagentActivity
import com.hermesandroid.relay.viewmodel.SubagentActivityPhase
import com.hermesandroid.relay.viewmodel.SubagentChildPreview
import kotlinx.coroutines.launch
/**
* Composer-adjacent summary of upstream Hermes processes for the active chat.
@@ -76,10 +63,8 @@ import kotlinx.coroutines.launch
* in the global session/navigation drawer.
*/
@Composable
internal fun GatewayBackgroundProcessStrip(
fun GatewayBackgroundProcessStrip(
processes: List<GatewayProcess>,
subagentActivities: List<SubagentActivity>,
subagentPreviewVisibility: SubagentPreviewVisibility,
loading: Boolean,
onClick: () -> Unit,
modifier: Modifier = Modifier,
@@ -87,26 +72,16 @@ internal fun GatewayBackgroundProcessStrip(
// Initial/switch refreshes are silent. The strip appears only after the
// session actually owns a process, avoiding a transient "Checking" row on
// every ordinary chat open.
val visibleActivities = subagentActivities.takeIf { subagentPreviewVisibility.showLifecycle }.orEmpty()
if (processes.isEmpty() && visibleActivities.isEmpty()) return
if (processes.isEmpty()) return
val running = processes.count { it.isRunning }
val runningAgents = visibleActivities.count { !it.isTerminal }
val failed = processes.count { !it.isRunning && (it.exitCode ?: 0) != 0 }
val failedAgents = visibleActivities.count { it.phase == SubagentActivityPhase.FAILED }
val interruptedAgents = visibleActivities.count {
it.phase == SubagentActivityPhase.INTERRUPTED ||
it.phase == SubagentActivityPhase.ENDED_WITH_PARENT
}
val failureCount = failed + failedAgents
val displayedCount = if (running > 0) running else processes.size
val status = when {
runningAgents > 0 -> stringResource(R.string.subagent_lane_running_count, runningAgents)
running > 0 -> "$running ${stringResource(R.string.bg_processes_running)}"
failureCount > 0 -> "$failureCount ${stringResource(R.string.task_status_failed)}"
interruptedAgents > 0 -> stringResource(R.string.agent_activity_status_interrupted)
failed > 0 -> "$failed ${stringResource(R.string.task_status_failed)}"
else -> stringResource(R.string.task_status_complete)
}
val openDescription = stringResource(R.string.current_chat_activity_open)
Surface(
modifier = modifier
@@ -115,11 +90,11 @@ internal fun GatewayBackgroundProcessStrip(
.heightIn(min = 48.dp)
.semantics {
contentDescription =
"$openDescription, $status"
"Background processes, $status. Open current chat activity."
stateDescription = status
}
.clickable(
onClickLabel = openDescription,
onClickLabel = stringResource(R.string.bg_processes_open),
onClick = onClick,
),
shape = RoundedCornerShape(14.dp),
@@ -130,41 +105,37 @@ internal fun GatewayBackgroundProcessStrip(
modifier = Modifier.padding(horizontal = 12.dp, vertical = 9.dp),
verticalAlignment = Alignment.CenterVertically,
) {
if (running > 0 || runningAgents > 0 || loading) {
if (running > 0 || loading) {
CircularProgressIndicator(modifier = Modifier.size(16.dp), strokeWidth = 2.dp)
} else {
Icon(
imageVector = when {
failureCount > 0 -> Icons.Filled.ErrorOutline
interruptedAgents > 0 -> Icons.Filled.PauseCircleOutline
else -> Icons.Filled.CheckCircle
},
imageVector = if (failed > 0) Icons.Filled.ErrorOutline else Icons.Filled.CheckCircle,
contentDescription = null,
modifier = Modifier.size(17.dp),
tint = when {
failureCount > 0 -> MaterialTheme.colorScheme.error
interruptedAgents > 0 -> MaterialTheme.colorScheme.tertiary
else -> MaterialTheme.colorScheme.primary
tint = if (failed > 0) {
MaterialTheme.colorScheme.error
} else {
MaterialTheme.colorScheme.primary
},
)
}
Spacer(Modifier.width(9.dp))
Text(
text = stringResource(R.string.current_chat_activity_title),
text = stringResource(R.string.background_process_count, displayedCount),
style = MaterialTheme.typography.labelLarge,
modifier = Modifier.weight(1f),
)
Text(
text = status,
style = MaterialTheme.typography.labelMedium,
color = if (failureCount > 0 && running == 0) {
color = if (failed > 0 && running == 0) {
MaterialTheme.colorScheme.error
} else {
MaterialTheme.colorScheme.onSurfaceVariant
},
)
Icon(
imageVector = Icons.Filled.Visibility,
imageVector = Icons.Filled.ExpandLess,
contentDescription = null,
modifier = Modifier
.padding(start = 6.dp)
@@ -178,56 +149,18 @@ internal fun GatewayBackgroundProcessStrip(
/** Mobile analogue of Hermes Desktop's composer process stack + terminal viewer. */
@OptIn(ExperimentalMaterial3Api::class)
@Composable
internal fun GatewayBackgroundProcessSheet(
fun GatewayBackgroundProcessSheet(
processes: List<GatewayProcess>,
subagentActivities: List<SubagentActivity>,
subagentChildPreview: SubagentChildPreview?,
subagentPreviewVisibility: SubagentPreviewVisibility,
loading: Boolean,
stoppingProcessIds: Set<String>,
onRefresh: () -> Unit,
onStop: (String) -> Unit,
onDismissProcess: (String) -> Unit,
onOpenSubagentChild: (String) -> Unit,
onDismiss: () -> Unit,
) {
val sheetState = rememberModalBottomSheetState(skipPartiallyExpanded = false)
val listState = androidx.compose.foundation.lazy.rememberLazyListState()
val scope = rememberCoroutineScope()
var expandedAgentKeys by remember { mutableStateOf<Set<String>>(emptySet()) }
val running = processes.filter { it.isRunning }
val recent = processes.filterNot { it.isRunning }
val visibleActivities = subagentActivities.takeIf { subagentPreviewVisibility.showLifecycle }.orEmpty()
val followTarget = subagentActivityFollowTarget(
visibleActivities,
expandedAgentKeys,
subagentPreviewVisibility,
subagentChildPreview,
)
val activityRevision = visibleActivities.sumOf { it.revision } +
subagentChildPreview?.messages.orEmpty().sumOf { message ->
message.content.length + message.thinkingContent.length +
message.toolCalls.sumOf { (it.args?.length ?: 0) + it.name.length }
}
val nearActivityTail by remember(followTarget, listState) {
derivedStateOf {
val last = listState.layoutInfo.visibleItemsInfo.lastOrNull()?.index ?: 0
followTarget >= 0 && last in maxOf(0, followTarget - 2)..followTarget
}
}
var followAgentTail by remember { mutableStateOf(true) }
LaunchedEffect(listState, followTarget) {
snapshotFlow { listState.isScrollInProgress to nearActivityTail }.collect { (scrolling, nearTail) ->
if (scrolling) followAgentTail = nearTail
else if (nearTail) followAgentTail = true
}
}
LaunchedEffect(activityRevision, followTarget) {
if (followAgentTail && followTarget >= 0) {
listState.scrollToItem(followTarget)
}
}
ModalBottomSheet(
onDismissRequest = onDismiss,
@@ -245,48 +178,23 @@ internal fun GatewayBackgroundProcessSheet(
verticalAlignment = Alignment.CenterVertically,
) {
Column(modifier = Modifier.weight(1f)) {
Text(stringResource(R.string.background_processes_title), style = MaterialTheme.typography.titleLarge)
Text(
stringResource(R.string.current_chat_activity_title),
modifier = Modifier.semantics { heading() },
style = MaterialTheme.typography.titleLarge,
)
Text(
stringResource(R.string.current_chat_activity_subtitle),
stringResource(R.string.background_processes_subtitle),
style = MaterialTheme.typography.bodySmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
)
}
if (processes.isNotEmpty()) IconButton(onClick = onRefresh, enabled = !loading) {
IconButton(onClick = onRefresh, enabled = !loading) {
if (loading) {
CircularProgressIndicator(modifier = Modifier.size(20.dp), strokeWidth = 2.dp)
} else {
Icon(Icons.Filled.Refresh, contentDescription = stringResource(R.string.background_processes_refresh_a11y))
}
}
IconButton(onClick = onDismiss) {
Icon(
Icons.Filled.Close,
contentDescription = stringResource(R.string.current_chat_activity_close),
)
}
}
if (!followAgentTail && followTarget >= 0) {
Row(
modifier = Modifier
.fillMaxWidth()
.padding(horizontal = 12.dp),
horizontalArrangement = Arrangement.End,
) {
TextButton(onClick = {
followAgentTail = true
scope.launch { listState.animateScrollToItem(followTarget) }
}) {
Text(stringResource(R.string.current_chat_activity_latest))
}
}
}
if (processes.isEmpty() && visibleActivities.isEmpty() && !loading) {
if (processes.isEmpty() && !loading) {
Column(
modifier = Modifier
.fillMaxWidth()
@@ -294,47 +202,23 @@ internal fun GatewayBackgroundProcessSheet(
horizontalAlignment = Alignment.CenterHorizontally,
) {
Icon(
Icons.Filled.AccountTree,
Icons.Filled.Terminal,
contentDescription = null,
modifier = Modifier.size(30.dp),
tint = MaterialTheme.colorScheme.onSurfaceVariant,
)
Text(
stringResource(R.string.current_chat_activity_empty),
stringResource(R.string.bg_processes_empty),
modifier = Modifier.padding(top = 12.dp),
color = MaterialTheme.colorScheme.onSurfaceVariant,
)
}
} else {
LazyColumn(
state = listState,
modifier = Modifier
.fillMaxWidth()
.heightIn(max = 560.dp),
) {
subagentActivityItems(
activities = visibleActivities,
expandedKeys = expandedAgentKeys,
visibility = subagentPreviewVisibility,
childPreview = subagentChildPreview,
onToggle = { key ->
expandedAgentKeys = if (key in expandedAgentKeys) {
expandedAgentKeys - key
} else {
expandedAgentKeys + key
}
},
onOpenChild = onOpenSubagentChild,
)
if (visibleActivities.isNotEmpty() && processes.isNotEmpty()) {
item { HorizontalDivider(modifier = Modifier.padding(vertical = 6.dp)) }
item {
ProcessSectionLabel(
stringResource(R.string.background_processes_title),
processes.size,
)
}
}
if (running.isNotEmpty()) {
item { ProcessSectionLabel(stringResource(R.string.bg_processes_running), running.size) }
items(running, key = { it.id }) { process ->
@@ -1,440 +0,0 @@
package com.hermesandroid.relay.ui.components
import androidx.compose.foundation.clickable
import androidx.compose.foundation.layout.Column
import androidx.compose.foundation.layout.Row
import androidx.compose.foundation.layout.Spacer
import androidx.compose.foundation.layout.fillMaxWidth
import androidx.compose.foundation.layout.height
import androidx.compose.foundation.layout.heightIn
import androidx.compose.foundation.layout.padding
import androidx.compose.foundation.layout.size
import androidx.compose.foundation.layout.width
import androidx.compose.foundation.lazy.LazyListScope
import androidx.compose.material.icons.Icons
import androidx.compose.material.icons.filled.AccountTree
import androidx.compose.material.icons.filled.CheckCircle
import androidx.compose.material.icons.filled.ErrorOutline
import androidx.compose.material.icons.filled.ExpandLess
import androidx.compose.material.icons.filled.ExpandMore
import androidx.compose.material.icons.filled.HourglassTop
import androidx.compose.material.icons.filled.PauseCircleOutline
import androidx.compose.material3.Icon
import androidx.compose.material3.MaterialTheme
import androidx.compose.material3.Surface
import androidx.compose.material3.Text
import androidx.compose.runtime.Composable
import androidx.compose.ui.Alignment
import androidx.compose.ui.Modifier
import androidx.compose.ui.graphics.vector.ImageVector
import androidx.compose.ui.res.stringResource
import androidx.compose.ui.semantics.LiveRegionMode
import androidx.compose.ui.semantics.contentDescription
import androidx.compose.ui.semantics.liveRegion
import androidx.compose.ui.semantics.semantics
import androidx.compose.ui.semantics.stateDescription
import androidx.compose.ui.text.font.FontFamily
import androidx.compose.ui.text.style.TextOverflow
import androidx.compose.ui.unit.dp
import com.hermesandroid.relay.R
import com.hermesandroid.relay.data.ChatMessage
import com.hermesandroid.relay.data.MessageRole
import com.hermesandroid.relay.viewmodel.SubagentChildPreview
import com.hermesandroid.relay.viewmodel.SubagentActivity
import com.hermesandroid.relay.viewmodel.SubagentActivityEvent
import com.hermesandroid.relay.viewmodel.SubagentActivityEventKind
import com.hermesandroid.relay.viewmodel.SubagentActivityPhase
internal data class SubagentPreviewVisibility(
val showLifecycle: Boolean = true,
val showReasoning: Boolean = true,
val showToolNames: Boolean = true,
val showToolDetails: Boolean = true,
val showChildHistory: Boolean = true,
)
internal fun LazyListScope.subagentActivityItems(
activities: List<SubagentActivity>,
expandedKeys: Set<String>,
visibility: SubagentPreviewVisibility,
childPreview: SubagentChildPreview?,
onToggle: (String) -> Unit,
onOpenChild: (String) -> Unit,
) {
if (activities.isEmpty() || !visibility.showLifecycle) return
item(key = "subagent-section") {
Column(modifier = Modifier.padding(horizontal = 20.dp, vertical = 8.dp)) {
Text(
stringResource(R.string.agent_activity_section),
style = MaterialTheme.typography.labelLarge,
)
Text(
stringResource(R.string.agent_activity_disclosure),
style = MaterialTheme.typography.bodySmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
)
}
}
activities.forEach { activity ->
val key = activity.stableKey
item(key = "subagent-header-$key") {
SubagentActivityHeader(
activity = activity,
expanded = key in expandedKeys,
visibility = visibility,
onClick = {
if (key !in expandedKeys && visibility.showChildHistory) onOpenChild(key)
onToggle(key)
},
)
}
if (key in expandedKeys) {
if (activity.truncated) {
item(key = "subagent-truncated-$key") {
SubagentMetaRow(stringResource(R.string.agent_activity_older_omitted))
}
}
activity.events.forEach { event ->
item(key = "subagent-event-$key-${event.sequence}") {
SubagentEventRow(event, visibility)
}
}
if (activity.partialAfterGap) {
item(key = "subagent-gap-$key") {
SubagentMetaRow(stringResource(R.string.agent_activity_partial))
}
}
childPreview?.takeIf { visibility.showChildHistory && it.activityKey == key }?.let { preview ->
when (preview.childWatchAvailable) {
null -> item(key = "subagent-child-loading-$key") {
SubagentMetaRow(stringResource(R.string.agent_activity_child_loading))
}
false -> item(key = "subagent-child-unavailable-$key") {
SubagentMetaRow(
preview.error?.takeIf(String::isNotBlank)
?: stringResource(R.string.agent_activity_child_unavailable),
)
}
true -> {
item(key = "subagent-child-heading-$key") {
SubagentMetaRow(
if (preview.running) {
stringResource(R.string.agent_activity_child_live)
} else {
stringResource(R.string.agent_activity_child_history)
},
)
}
if (preview.historyTruncated) {
item(key = "subagent-child-truncated-$key") {
SubagentMetaRow(stringResource(R.string.agent_activity_child_truncated))
}
}
preview.messages.filterNot { it.role == MessageRole.SYSTEM }.forEach { message ->
item(key = "subagent-child-message-$key-${message.uiKey}") {
SubagentChildMessageRow(message, visibility)
}
}
preview.error?.takeIf(String::isNotBlank)?.let { error ->
item(key = "subagent-child-error-$key") { SubagentMetaRow(error) }
}
}
}
item(key = "subagent-child-tail-$key") {
Spacer(Modifier.height(1.dp))
}
}
}
}
}
internal fun subagentActivityItemCount(
activities: List<SubagentActivity>,
expandedKeys: Set<String>,
visibility: SubagentPreviewVisibility,
childPreview: SubagentChildPreview?,
): Int {
if (activities.isEmpty() || !visibility.showLifecycle) return 0
return 1 + activities.sumOf { activity ->
val expanded = activity.stableKey in expandedKeys
val preview = childPreview?.takeIf {
visibility.showChildHistory && it.activityKey == activity.stableKey
}
val previewRows = when (preview?.childWatchAvailable) {
null -> if (preview != null) 1 else 0
false -> 1
true -> 1 + preview.messages.count { it.role != MessageRole.SYSTEM } +
(if (preview.historyTruncated) 1 else 0) +
(if (preview.error.isNullOrBlank()) 0 else 1)
} + if (preview != null) 1 else 0 // explicit bottom anchor for growing rows
1 + if (!expanded) 0 else activity.events.size +
(if (activity.truncated) 1 else 0) +
(if (activity.partialAfterGap) 1 else 0) + previewRows
}
}
internal fun subagentActivityFollowTarget(
activities: List<SubagentActivity>,
expandedKeys: Set<String>,
visibility: SubagentPreviewVisibility,
childPreview: SubagentChildPreview?,
): Int {
if (activities.isEmpty() || !visibility.showLifecycle) return -1
var index = 0 // section heading
var selectedTarget = -1
activities.forEach { activity ->
index += 1 // lane header
if (activity.stableKey in expandedKeys) {
if (activity.truncated) index += 1
index += activity.events.size
if (activity.partialAfterGap) index += 1
if (visibility.showChildHistory && childPreview?.activityKey == activity.stableKey) {
index += when (childPreview.childWatchAvailable) {
null, false -> 1
true -> 1 + childPreview.messages.count { it.role != MessageRole.SYSTEM } +
(if (childPreview.historyTruncated) 1 else 0) +
(if (childPreview.error.isNullOrBlank()) 0 else 1)
}
index += 1 // explicit bottom anchor for growing child content
selectedTarget = index
}
}
}
return if (selectedTarget >= 0) selectedTarget else index
}
@Composable
private fun SubagentActivityHeader(
activity: SubagentActivity,
expanded: Boolean,
visibility: SubagentPreviewVisibility,
onClick: () -> Unit,
) {
val title = activity.goal.takeIf { visibility.showReasoning && it.isNotBlank() }
?: stringResource(R.string.agent_activity_fallback, activity.taskIndex + 1)
val phaseLabel = phaseLabel(activity.phase)
val description = stringResource(
R.string.agent_activity_lane_a11y,
title,
phaseLabel,
activity.taskIndex + 1,
activity.taskCount,
)
val icon: ImageVector = when (activity.phase) {
SubagentActivityPhase.COMPLETED -> Icons.Filled.CheckCircle
SubagentActivityPhase.FAILED -> Icons.Filled.ErrorOutline
SubagentActivityPhase.INTERRUPTED,
SubagentActivityPhase.ENDED_WITH_PARENT,
-> Icons.Filled.PauseCircleOutline
SubagentActivityPhase.STARTED,
SubagentActivityPhase.THINKING,
SubagentActivityPhase.TOOL,
SubagentActivityPhase.PROGRESS,
-> Icons.Filled.HourglassTop
}
val tint = when (activity.phase) {
SubagentActivityPhase.FAILED -> MaterialTheme.colorScheme.error
SubagentActivityPhase.COMPLETED -> MaterialTheme.colorScheme.primary
else -> MaterialTheme.colorScheme.tertiary
}
Surface(
modifier = Modifier
.fillMaxWidth()
.padding(horizontal = 16.dp, vertical = 3.dp)
.heightIn(min = 48.dp)
.semantics {
contentDescription = description
stateDescription = phaseLabel
liveRegion = LiveRegionMode.Polite
}
.clickable(onClick = onClick),
shape = MaterialTheme.shapes.medium,
color = MaterialTheme.colorScheme.surfaceVariant.copy(alpha = 0.5f),
) {
Row(
modifier = Modifier.padding(horizontal = 12.dp, vertical = 9.dp),
verticalAlignment = Alignment.CenterVertically,
) {
Icon(icon, contentDescription = null, tint = tint, modifier = Modifier.size(18.dp))
Spacer(Modifier.width(9.dp))
Column(modifier = Modifier.weight(1f)) {
Text(
title,
style = MaterialTheme.typography.labelLarge,
maxLines = 2,
overflow = TextOverflow.Ellipsis,
)
Text(
stringResource(
R.string.agent_activity_task_position,
activity.taskIndex + 1,
activity.taskCount,
phaseLabel,
),
style = MaterialTheme.typography.bodySmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
)
}
activity.durationSeconds?.takeIf { activity.isTerminal }?.let { seconds ->
Text(
stringResource(R.string.agent_activity_duration, seconds),
style = MaterialTheme.typography.labelSmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
)
Spacer(Modifier.width(6.dp))
}
Icon(
if (expanded) Icons.Filled.ExpandLess else Icons.Filled.ExpandMore,
contentDescription = stringResource(
if (expanded) R.string.cd_subagent_collapse else R.string.cd_subagent_expand,
),
)
}
}
}
@Composable
private fun SubagentEventRow(
event: SubagentActivityEvent,
visibility: SubagentPreviewVisibility,
) {
val label = when (event.kind) {
SubagentActivityEventKind.STARTED -> stringResource(R.string.agent_activity_event_started)
SubagentActivityEventKind.UPDATE -> stringResource(R.string.agent_activity_event_update)
SubagentActivityEventKind.TOOL -> stringResource(R.string.agent_activity_event_tool)
SubagentActivityEventKind.COMPLETED -> phaseLabel(event.phase)
}
val showText = when (event.kind) {
SubagentActivityEventKind.STARTED -> false
SubagentActivityEventKind.UPDATE,
SubagentActivityEventKind.COMPLETED,
-> visibility.showReasoning
SubagentActivityEventKind.TOOL -> visibility.showToolDetails
}
val toolName = event.toolName?.takeIf { visibility.showToolNames || visibility.showToolDetails }
Column(
modifier = Modifier
.fillMaxWidth()
.padding(start = 38.dp, end = 20.dp, top = 5.dp, bottom = 5.dp),
) {
Row(verticalAlignment = Alignment.CenterVertically) {
Icon(
Icons.Filled.AccountTree,
contentDescription = null,
modifier = Modifier.size(14.dp),
tint = MaterialTheme.colorScheme.onSurfaceVariant,
)
Spacer(Modifier.width(6.dp))
Text(label, style = MaterialTheme.typography.labelSmall)
toolName?.let {
Text(
" · $it",
style = MaterialTheme.typography.labelSmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
maxLines = 1,
overflow = TextOverflow.Ellipsis,
)
}
}
event.text?.takeIf { showText && it.isNotBlank() }?.let { text ->
Text(
text,
modifier = Modifier.padding(top = 3.dp),
style = MaterialTheme.typography.bodySmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
fontFamily = if (event.kind == SubagentActivityEventKind.TOOL) {
FontFamily.Monospace
} else {
FontFamily.Default
},
maxLines = 8,
overflow = TextOverflow.Ellipsis,
)
}
}
}
@Composable
private fun SubagentMetaRow(text: String) {
Text(
text,
modifier = Modifier.padding(start = 38.dp, end = 20.dp, top = 4.dp, bottom = 6.dp),
style = MaterialTheme.typography.bodySmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
)
}
@Composable
private fun SubagentChildMessageRow(
message: ChatMessage,
visibility: SubagentPreviewVisibility,
) {
if (message.role == MessageRole.SYSTEM) return
val role = when (message.role) {
MessageRole.USER -> stringResource(R.string.agent_activity_child_role_task)
MessageRole.ASSISTANT -> stringResource(R.string.agent_activity_child_role_agent)
MessageRole.SYSTEM -> stringResource(R.string.agent_activity_child_role_system)
}
Column(
modifier = Modifier
.fillMaxWidth()
.padding(start = 38.dp, end = 20.dp, top = 6.dp, bottom = 6.dp),
) {
Text(
role,
style = MaterialTheme.typography.labelSmall,
color = MaterialTheme.colorScheme.primary,
)
message.thinkingContent.takeIf { visibility.showReasoning && it.isNotBlank() }?.let { thought ->
Text(
thought,
modifier = Modifier.padding(top = 3.dp),
style = MaterialTheme.typography.bodySmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
maxLines = 10,
overflow = TextOverflow.Ellipsis,
)
}
message.content.takeIf(String::isNotBlank)?.let { content ->
Text(
content,
modifier = Modifier.padding(top = 3.dp),
style = MaterialTheme.typography.bodyMedium,
maxLines = 20,
overflow = TextOverflow.Ellipsis,
)
}
if (visibility.showToolNames || visibility.showToolDetails) {
message.toolCalls.forEach { tool ->
Text(
buildString {
append(tool.name)
if (visibility.showToolDetails) {
tool.args?.takeIf(String::isNotBlank)?.let { append(" · ").append(it.take(500)) }
}
},
modifier = Modifier.padding(top = 3.dp),
style = MaterialTheme.typography.bodySmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
fontFamily = FontFamily.Monospace,
maxLines = 5,
overflow = TextOverflow.Ellipsis,
)
}
}
}
}
@Composable
private fun phaseLabel(phase: SubagentActivityPhase): String = stringResource(
when (phase) {
SubagentActivityPhase.STARTED -> R.string.agent_activity_status_started
SubagentActivityPhase.THINKING -> R.string.agent_activity_status_thinking
SubagentActivityPhase.TOOL -> R.string.agent_activity_status_tool
SubagentActivityPhase.PROGRESS -> R.string.agent_activity_status_progress
SubagentActivityPhase.COMPLETED -> R.string.agent_activity_status_completed
SubagentActivityPhase.FAILED -> R.string.agent_activity_status_failed
SubagentActivityPhase.INTERRUPTED -> R.string.agent_activity_status_interrupted
SubagentActivityPhase.ENDED_WITH_PARENT -> R.string.agent_activity_status_unavailable
},
)
@@ -216,7 +216,6 @@ import com.hermesandroid.relay.ui.components.CHAT_PET_STEP_MESSAGE_MARKER
import com.hermesandroid.relay.ui.components.CHAT_PET_USER_MESSAGE_PERCH_PREFIX
import com.hermesandroid.relay.ui.components.GatewayBackgroundProcessSheet
import com.hermesandroid.relay.ui.components.GatewayBackgroundProcessStrip
import com.hermesandroid.relay.ui.components.SubagentPreviewVisibility
import com.hermesandroid.relay.ui.components.InjectedContextSheet
import com.hermesandroid.relay.ui.components.InlineAutocomplete
import com.hermesandroid.relay.ui.components.loadedContentTransform
@@ -929,8 +928,6 @@ fun ChatScreen(
val backgroundProcesses by chatViewModel.backgroundProcesses.collectAsState()
val backgroundProcessesLoading by chatViewModel.backgroundProcessesLoading.collectAsState()
val stoppingProcessIds by chatViewModel.stoppingProcessIds.collectAsState()
val subagentActivities by chatViewModel.subagentActivities.collectAsState()
val subagentChildPreview by chatViewModel.subagentChildPreview.collectAsState()
val isLoadingHistory by chatViewModel.isLoadingHistory.collectAsState()
val isLoadingSessions by chatViewModel.isLoadingSessions.collectAsState()
val selectedPersonality by chatViewModel.selectedPersonality.collectAsState()
@@ -1065,13 +1062,6 @@ fun ChatScreen(
supervisedVisibility.showToolNames -> "compact"
else -> "off"
}
val subagentPreviewVisibility = SubagentPreviewVisibility(
showLifecycle = !supervised || supervisedVisibility.showWorkingStatus,
showReasoning = showThinking,
showToolNames = toolDisplay == "compact" || toolDisplay == "detailed",
showToolDetails = toolDisplay == "detailed",
showChildHistory = !supervised,
)
val smoothAutoScroll by connectionViewModel.smoothAutoScroll.collectAsState()
val closeDrawerOnSend by connectionViewModel.closeDrawerOnSend.collectAsState()
val keepComposerFocusedOnSend by
@@ -1129,12 +1119,16 @@ fun ChatScreen(
}
// Recover any durable in-flight chat checkpoint whenever Chat returns to
// the foreground. On Gateway this also pre-warms/re-attaches the socket;
// sessions-SSE falls back to bounded persisted-history reconciliation.
// the foreground. setChatVisible owns that edge; an ordinary Gateway open
// warms only the observation socket and never attaches a saved session.
val appForeground by com.hermesandroid.relay.util.AppForegroundTracker.isForeground.collectAsState()
LaunchedEffect(isGatewayTransport, appForeground, chatReady) {
chatViewModel.setChatVisible(appForeground && chatReady)
if (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.prewarmGateway()
}
if (isGatewayTransport && appForeground && chatReady) {
@@ -1372,7 +1366,6 @@ fun ChatScreen(
// A process inventory is scoped to one gateway session. Never leave a
// sheet opened onto a different chat after a drawer/profile switch.
LaunchedEffect(currentSessionId, selectedProfile?.name, activeConnection?.id) {
chatViewModel.closeSubagentChildPreview()
showBackgroundProcesses = false
}
@@ -3965,8 +3958,6 @@ fun ChatScreen(
if (isGatewayTransport) {
GatewayBackgroundProcessStrip(
processes = backgroundProcesses,
subagentActivities = subagentActivities,
subagentPreviewVisibility = subagentPreviewVisibility,
loading = backgroundProcessesLoading,
onClick = { showBackgroundProcesses = true },
)
@@ -4876,19 +4867,12 @@ fun ChatScreen(
if (showBackgroundProcesses) {
GatewayBackgroundProcessSheet(
processes = backgroundProcesses,
subagentActivities = subagentActivities,
subagentChildPreview = subagentChildPreview,
subagentPreviewVisibility = subagentPreviewVisibility,
loading = backgroundProcessesLoading,
stoppingProcessIds = stoppingProcessIds,
onRefresh = chatViewModel::refreshBackgroundProcesses,
onStop = chatViewModel::stopBackgroundProcess,
onDismissProcess = chatViewModel::dismissBackgroundProcess,
onOpenSubagentChild = chatViewModel::openSubagentChildPreview,
onDismiss = {
chatViewModel.closeSubagentChildPreview()
showBackgroundProcesses = false
},
onDismiss = { showBackgroundProcesses = false },
)
}
@@ -408,6 +408,9 @@ 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
@@ -419,8 +422,10 @@ class ChatViewModel : ViewModel() {
_sessionDirectoryRefreshRequests.asSharedFlow()
private fun activityScope(contextKey: String? = activeProfileContextKey): SessionActivityScope? {
val identity = AgentDisplay.parseProfileContextKey(contextKey) ?: return null
return SessionActivityScope.of(identity.connectionId, identity.profileKey)
val raw = contextKey?.trim().orEmpty()
val separator = raw.lastIndexOf("::")
if (separator <= 0 || separator >= raw.lastIndex) return null
return SessionActivityScope.of(raw.substring(0, separator), raw.substring(separator + 2))
}
private fun activityOwner(
@@ -532,8 +537,7 @@ class ChatViewModel : ViewModel() {
private fun backgroundTurnKey(sessionId: String, profile: String?): TurnCheckpointKey? {
val profileKey = AgentDisplay.profileSessionKey(profile)
return backgroundTurnCheckpoints.keys.firstOrNull { key ->
key.sessionId == sessionId &&
AgentDisplay.parseProfileContextKey(key.contextKey)?.profileKey == profileKey
key.sessionId == sessionId && key.contextKey.substringAfterLast("::") == profileKey
}
}
private var checkpointWriteJob: Job? = null
@@ -546,7 +550,6 @@ class ChatViewModel : ViewModel() {
private data class ActiveTurnCheckpointSeed(
var contextKey: String?,
val profileKey: String?,
var sessionId: String,
var liveSessionId: String?,
var transport: String,
@@ -2078,10 +2081,6 @@ class ChatViewModel : ViewModel() {
private var chatVisible = false
private var gatewayProcessSource: GatewayProcessSource? = null
private val gatewayProcessController = GatewayProcessController(viewModelScope)
private val subagentActivityController = SubagentActivityController()
private val subagentChildPreviewController = SubagentChildPreviewController(viewModelScope)
internal val subagentChildPreview: StateFlow<SubagentChildPreview?> =
subagentChildPreviewController.state
/** Session-scoped upstream shell processes shown beside the composer. */
val backgroundProcesses: StateFlow<List<GatewayProcess>> = gatewayProcessController.processes
@@ -2093,33 +2092,6 @@ class ChatViewModel : ViewModel() {
val backgroundProcessesLoading: StateFlow<Boolean> = gatewayProcessController.loading
val stoppingProcessIds: StateFlow<Set<String>> = gatewayProcessController.stoppingProcessIds
/** Bounded parent-session lifecycle previews for the active chat's delegated work. */
internal val subagentActivities: StateFlow<List<SubagentActivity>> =
subagentActivityController.activities
fun openSubagentChildPreview(activityKey: String) {
val activity = subagentActivities.value.firstOrNull { it.stableKey == activityKey } ?: return
val parentSessionId = chatHandler?.currentSessionId?.value ?: return
val parentScopeKey = activeProfileContextKey
val client = gatewayClient
subagentChildPreviewController.open(
activity = activity,
client = client,
parentSessionId = parentSessionId,
parentScopeKey = parentScopeKey,
gatewayRouteActive = streamingEndpoint == "gateway",
stillOwnsParent = {
gatewayClient === client &&
chatHandler?.currentSessionId?.value == parentSessionId &&
activeProfileContextKey == parentScopeKey
},
)
}
fun closeSubagentChildPreview() {
subagentChildPreviewController.close()
}
private val _messageReactionsSupported = MutableStateFlow(true)
val messageReactionsSupported: StateFlow<Boolean> = _messageReactionsSupported.asStateFlow()
@@ -2283,6 +2255,8 @@ 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()
@@ -2303,6 +2277,43 @@ 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,
@@ -2364,6 +2375,18 @@ 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) }
@@ -2382,7 +2405,9 @@ class ChatViewModel : ViewModel() {
record.freshness == SessionActivityFreshness.Confirmed &&
record.phase(System.currentTimeMillis()) != SessionActivityPhase.Idle
}
val delayMs = if (hasConfirmedLiveWork) 1_500L else 30_000L
val delayMs = if (
hasConfirmedLiveWork || hasPassiveCurrentLiveWork || hasPassiveCatchupPending
) 1_500L else 30_000L
sessionActivityPollJob = viewModelScope.launch {
delay(delayMs)
if (gatewayClient === client && chatVisible) pollSessionActivity(client)
@@ -2455,6 +2480,10 @@ class ChatViewModel : ViewModel() {
clearProjectedBackgroundProcesses()
sessionActivityPollJob?.cancel()
sessionActivityPollJob = null
passiveGatewayHistoryRefreshJob?.cancel()
passiveGatewayHistoryRefreshJob = null
passivelyObservedGatewaySessionId = null
passiveObservationCatchupPendingSessionId = null
sessionActivityGeneration.incrementAndGet()
sessionActivityDirectory = emptySet()
lastLocalActivityOwner = null
@@ -2473,8 +2502,6 @@ class ChatViewModel : ViewModel() {
resetApprovalModeState()
_messageReactionsSupported.value = true
gatewayProcessSource = client?.let(::GatewayChatProcessSource)
closeSubagentChildPreview()
subagentActivityController.resetConnection()
gatewayProcessController.bind(
newSource = gatewayProcessSource,
sessionId = chatHandler?.currentSessionId?.value,
@@ -2576,8 +2603,9 @@ 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 visible durable session is
// resumed and its authoritative history reconciled.
// client transition so the observation socket is
// restored; only an exact Android checkpoint may
// resume/activate a live runtime.
prewarmGateway()
}
}
@@ -2637,7 +2665,7 @@ class ChatViewModel : ViewModel() {
val matching = backgroundTurnCheckpoints.keys.filter { key ->
val checkpoint = backgroundTurnCheckpoints[key]
key.sessionId == completion.storedSessionId &&
AgentDisplay.parseProfileContextKey(key.contextKey)?.profileKey == profileKey &&
key.contextKey.substringAfterLast("::") == profileKey &&
checkpoint?.liveSessionId == completion.liveSessionId
}
if (matching.isEmpty()) return
@@ -2686,14 +2714,11 @@ class ChatViewModel : ViewModel() {
queuedRecovery: QueuedRecoveryHandoff? = null,
): GatewayInboundTurnRegistration? {
val handler = chatHandler ?: return null
val eventScopeKey = activeProfileContextKey
val eventProfile = currentSessionProfileName()
fun matchesAdmissionContext(): Boolean =
gatewayClient === client &&
streamingEndpoint == "gateway" &&
chatHandler === handler &&
handler.currentSessionId.value == storedSessionId &&
activeProfileContextKey == eventScopeKey
handler.currentSessionId.value == storedSessionId
val messageId = "gateway-inbound-${UUID.randomUUID()}"
val queuedUserMessageId = "gateway-queued-user-${UUID.randomUUID()}"
@@ -2728,11 +2753,6 @@ class ChatViewModel : ViewModel() {
onStart = {
if (!started && ownsBoundTurn()) {
started = true
subagentActivityController.beginTurn(
storedSessionId,
eventScopeKey,
messageId,
)
cancelAnswerRecovery()
intentionallyCancelled = false
firstTokenNotified = false
@@ -2791,7 +2811,6 @@ class ChatViewModel : ViewModel() {
?.content
?.takeIf { it.isNotBlank() }
if (canWriteTranscript) {
subagentActivityController.endTurn(messageId)
val failed = handler.messages.value
.lastOrNull { it.id == messageId }
?.badges
@@ -2827,7 +2846,6 @@ class ChatViewModel : ViewModel() {
onError = { message ->
val canWriteTranscript = acceptsEvent()
if (canWriteTranscript) {
subagentActivityController.endTurn(messageId)
AppAnalytics.onStreamError()
handler.onStreamError(message)
emitError(Exception(message), context = "send_message")
@@ -2845,16 +2863,7 @@ class ChatViewModel : ViewModel() {
if (acceptsEvent()) handler.onToolGenerating(messageId, name)
},
onSubagentEvent = { event ->
if (acceptsEvent()) {
subagentActivityController.onEvent(
sessionId = storedSessionId,
eventScopeKey = eventScopeKey,
turnId = messageId,
event = event,
profile = eventProfile,
)
handler.onSubagentEvent(messageId, event)
}
if (acceptsEvent()) handler.onSubagentEvent(messageId, event)
},
onMoaReference = { event ->
if (acceptsEvent()) handler.onMoaReference(messageId, event)
@@ -3051,6 +3060,103 @@ 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 {
@@ -3061,11 +3167,11 @@ class ChatViewModel : ViewModel() {
}
/**
* 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.
* 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.
*/
fun prewarmGateway() {
val client = gatewayClient
@@ -3073,17 +3179,16 @@ class ChatViewModel : ViewModel() {
val sessionId = handler.currentSessionId.value
selectBackgroundProcessSession(sessionId)
if (sessionId == null) {
client?.prewarm(null)
client?.observe()
} 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 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)
// GatewayChatClient owns the socket IO scope, so the dial can
// progress while a paused UI dispatcher is being recreated.
observeGatewaySession(gateway, handler, sessionId)
return
}
if (activeStream == null && (streamRecovery == null || client != null)) {
@@ -3095,26 +3200,29 @@ class ChatViewModel : ViewModel() {
chatHandler === handler &&
handler.currentSessionId.value == sessionId
) {
if (client?.prewarmAwait(sessionId) == true) {
gatewayProcessController.sessionReady(sessionId)
}
observeGatewaySession(client, handler, sessionId)
}
checkpointRecoveryJob = null
}
return
}
// 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)
// 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()
}
}
} else {
observeGatewaySession(client, handler, sessionId)
}
}
}
@@ -3124,7 +3232,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) {
fun setChatVisible(visible: Boolean): Boolean {
val changed = chatVisible != visible
chatVisible = visible
if (visible && changed) {
@@ -3133,7 +3241,12 @@ class ChatViewModel : ViewModel() {
} else if (!visible) {
sessionActivityPollJob?.cancel()
sessionActivityPollJob = null
passiveGatewayHistoryRefreshJob?.cancel()
passiveGatewayHistoryRefreshJob = null
passivelyObservedGatewaySessionId = null
passiveObservationCatchupPendingSessionId = null
}
return changed
}
// === Gateway desktop-parity state ===
@@ -3424,13 +3537,7 @@ class ChatViewModel : ViewModel() {
sessionId: String?,
scopeKey: String? = activeProfileContextKey,
) {
subagentChildPreview.value?.let { preview ->
if (preview.parentSessionId != sessionId || preview.parentScopeKey != scopeKey) {
closeSubagentChildPreview()
}
}
gatewayProcessController.selectSession(sessionId, scopeKey)
subagentActivityController.selectSession(sessionId, scopeKey)
}
/**
@@ -3697,14 +3804,6 @@ class ChatViewModel : ViewModel() {
requestSessionActivityRefresh()
}
}
launch {
client.connectionState.collect { state ->
if (gatewayClient !== client) return@collect
subagentActivityController.onConnectionReady(
state == com.hermesandroid.relay.network.upstream.GatewayConnectionState.Ready,
)
}
}
launch {
client.serverPersonality.collect { value ->
if (gatewayClient !== client || value == null) return@collect
@@ -4531,7 +4630,9 @@ class ChatViewModel : ViewModel() {
)
if (stillCurrent()) {
handler.loadMessageHistory(messages)
if (streamingEndpoint == "gateway") gatewayClient?.prewarm(sessionId)
if (streamingEndpoint == "gateway") {
observeGatewaySession(gatewayClient, handler, sessionId)
}
}
}
} catch (e: kotlinx.coroutines.CancellationException) {
@@ -5042,7 +5143,9 @@ class ChatViewModel : ViewModel() {
handler.currentSessionId.value == sessionId
) {
handler.loadMessageHistory(messages)
if (streamingEndpoint == "gateway") gatewayClient?.prewarm(sessionId)
if (streamingEndpoint == "gateway") {
observeGatewaySession(gatewayClient, handler, sessionId)
}
}
} catch (e: kotlinx.coroutines.CancellationException) {
throw e
@@ -6329,16 +6432,6 @@ class ChatViewModel : ViewModel() {
checkpointWriteJob?.cancel()
activeTurnCheckpointSeed = ActiveTurnCheckpointSeed(
contextKey = activeProfileContextKey,
// Persist the explicit UI selection, not the effective sticky
// server profile used to bind the current live session. Server
// Default must survive restart as the sentinel even when Hermes
// currently resolves it to a named profile such as `victor`.
profileKey = AgentDisplay.profileSessionKey(
conversationBinding.value.let { binding ->
if (binding.hasExplicitOwner) binding.profileName
else selectedProfileProvider()?.name
},
),
sessionId = sessionId,
liveSessionId = null,
transport = transport,
@@ -6379,7 +6472,6 @@ class ChatViewModel : ViewModel() {
checkpointWriteJob?.cancel()
activeTurnCheckpointSeed = ActiveTurnCheckpointSeed(
contextKey = checkpoint.contextKey,
profileKey = checkpoint.profileKey,
sessionId = checkpoint.sessionId,
liveSessionId = checkpoint.liveSessionId,
transport = checkpoint.transport,
@@ -6440,7 +6532,6 @@ class ChatViewModel : ViewModel() {
val now = System.currentTimeMillis()
return ChatTurnCheckpoint(
contextKey = contextKey,
profileKey = seed.profileKey,
sessionId = sessionId,
liveSessionId = gatewayClient?.currentLiveSessionId(sessionId) ?: seed.liveSessionId,
transport = seed.transport,
@@ -6644,6 +6735,10 @@ 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
@@ -6888,11 +6983,6 @@ class ChatViewModel : ViewModel() {
queuedSuccessorPending: AtomicBoolean,
): GatewayTurnCallbacks {
val messageId = checkpoint.assistant.id
subagentActivityController.beginTurn(
checkpoint.sessionId,
checkpoint.contextKey,
messageId,
)
fun owns(): Boolean = ownsTurnCheckpoint(checkpoint, handler)
return GatewayTurnCallbacks(
onSessionId = { },
@@ -6943,7 +7033,6 @@ class ChatViewModel : ViewModel() {
onReconcileRequired = { },
onComplete = {
if (owns()) {
subagentActivityController.endTurn(messageId)
cancelAnswerRecovery(settleUi = false)
val failed = handler.messages.value
.lastOrNull { it.id == messageId }
@@ -6981,7 +7070,6 @@ class ChatViewModel : ViewModel() {
},
onError = { error ->
if (owns()) {
subagentActivityController.endTurn(messageId)
if (queuedSuccessorPending.get()) {
AppAnalytics.onStreamError()
handler.onStreamError(error)
@@ -7006,18 +7094,6 @@ class ChatViewModel : ViewModel() {
},
onSubagentEvent = { event ->
if (owns()) {
subagentActivityController.onEvent(
sessionId = checkpoint.sessionId,
eventScopeKey = checkpoint.contextKey,
turnId = messageId,
event = event,
profile = if (checkpoint.profileKey != null) {
AgentDisplay.profileRequestName(checkpoint.profileKey)
} else {
AgentDisplay.parseProfileContextKey(checkpoint.contextKey)
?.requestProfileName
},
)
handler.onSubagentEvent(messageId, event)
scheduleCheckpointWrite(immediate = true)
}
@@ -8790,7 +8866,6 @@ class ChatViewModel : ViewModel() {
// wins — stop the poller before finalizing so the turn can't
// finish twice.
cancelAnswerRecovery(settleUi = false)
subagentActivityController.endTurn(currentMessageId)
val completedTransport = dispatchedSseEndpoint
?: if (activeStreamIsGateway) "gateway" else streamingEndpoint
val turnErrored = handler.messages.value
@@ -8918,7 +8993,6 @@ class ChatViewModel : ViewModel() {
}
val onErrorCb = { errorMsg: String ->
markTransportFailed(errorMsg)
subagentActivityController.endTurn(currentMessageId)
stopImageActivityBridge()
flushAndReleaseStreamDeltas()
val errorSessionId = handler.currentSessionId.value
@@ -9319,13 +9393,6 @@ class ChatViewModel : ViewModel() {
// context rides the SSE systemMessage (invisible) + the on-demand
// android_phone_status tool instead. See PhoneStatusPromptBuilder.
_steerableTurn.value = true
handler.currentSessionId.value?.let { existingSessionId ->
subagentActivityController.beginTurn(
existingSessionId,
activeProfileContextKey,
currentMessageId,
)
}
gateway.sendTurn(
sessionId = handler.currentSessionId.value,
text = message,
@@ -9347,11 +9414,6 @@ class ChatViewModel : ViewModel() {
markSessionActivityStarting(sid)
updateTurnCheckpointSession(sid)
selectBackgroundProcessSession(sid)
subagentActivityController.beginTurn(
sid,
activeProfileContextKey,
currentMessageId,
)
gatewayProcessController.sessionReady(sid)
onSessionChanged?.invoke(sid)
// The brand-new chat now has a session — apply any
@@ -9396,13 +9458,6 @@ class ChatViewModel : ViewModel() {
onSubagentEvent = { event ->
ensurePostInterimMessage()
streamDeltas.flushNow()
subagentActivityController.onEvent(
sessionId = handler.currentSessionId.value,
eventScopeKey = activeProfileContextKey,
turnId = currentMessageId,
event = event,
profile = currentSessionProfileName(),
)
handler.onSubagentEvent(currentMessageId, event)
scheduleCheckpointWrite(immediate = true)
},
@@ -1,276 +0,0 @@
package com.hermesandroid.relay.viewmodel
import com.hermesandroid.relay.network.upstream.GatewaySubagentEvent
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
internal enum class SubagentActivityPhase {
STARTED,
THINKING,
TOOL,
PROGRESS,
COMPLETED,
FAILED,
INTERRUPTED,
ENDED_WITH_PARENT,
}
internal enum class SubagentActivityEventKind { STARTED, UPDATE, TOOL, COMPLETED }
internal data class SubagentActivityEvent(
val sequence: Long,
val kind: SubagentActivityEventKind,
val text: String? = null,
val toolName: String? = null,
val phase: SubagentActivityPhase,
val observedAtMillis: Long,
)
/**
* A bounded, ephemeral projection of parent-session `subagent.*` events.
*
* This is intentionally not a child transcript. Upstream currently exposes no
* durable child-session key or child-history route, so the projection is owned
* by the exact profile-scoped parent session and parent turn that emitted it.
*/
internal data class SubagentActivity(
val laneId: Long,
val turnId: String,
val taskIndex: Int,
val taskCount: Int,
val goal: String,
val subagentId: String? = null,
val childSessionId: String? = null,
val parentId: String? = null,
val depth: Int? = null,
val model: String? = null,
val profile: String? = null,
val phase: SubagentActivityPhase,
val summary: String? = null,
val durationSeconds: Double? = null,
val events: List<SubagentActivityEvent> = emptyList(),
val truncated: Boolean = false,
val partialAfterGap: Boolean = false,
val revision: Long = 0L,
) {
val stableKey: String
get() = "$turnId:$laneId"
val isTerminal: Boolean
get() = phase in setOf(
SubagentActivityPhase.COMPLETED,
SubagentActivityPhase.FAILED,
SubagentActivityPhase.INTERRUPTED,
SubagentActivityPhase.ENDED_WITH_PARENT,
)
}
/**
* Keeps live child activity isolated from the unrelated `process.list`
* registry. All text is control-sanitized and bounded before entering UI state.
*/
internal class SubagentActivityController(
private val clock: () -> Long = System::currentTimeMillis,
) {
companion object {
internal const val MAX_EVENTS_PER_CHILD = 50
internal const val MAX_CHARS_PER_CHILD = 32_000
internal const val MAX_GOAL_CHARS = 500
internal const val MAX_EVENT_TEXT_CHARS = 2_000
internal const val MAX_TOOL_NAME_CHARS = 160
}
private val _activities = MutableStateFlow<List<SubagentActivity>>(emptyList())
val activities: StateFlow<List<SubagentActivity>> = _activities.asStateFlow()
private var storedSessionId: String? = null
private var scopeKey: String? = null
private var activeTurnId: String? = null
private var sequence = 0L
private var laneSequence = 0L
private var connectionWasReady = false
private var pendingGap = false
fun selectSession(sessionId: String?, newScopeKey: String?) {
if (storedSessionId == sessionId && scopeKey == newScopeKey) return
storedSessionId = sessionId
scopeKey = newScopeKey
activeTurnId = null
sequence = 0L
laneSequence = 0L
connectionWasReady = false
pendingGap = false
_activities.value = emptyList()
}
fun resetConnection() {
activeTurnId = null
sequence = 0L
laneSequence = 0L
connectionWasReady = false
pendingGap = false
_activities.value = emptyList()
}
fun onConnectionReady(ready: Boolean) {
if (connectionWasReady && !ready && _activities.value.any { !it.isTerminal }) {
pendingGap = true
}
if (ready && pendingGap) {
_activities.value = _activities.value.map { activity ->
if (activity.isTerminal) activity else activity.copy(
partialAfterGap = true,
revision = activity.revision + 1,
)
}
pendingGap = false
}
connectionWasReady = ready
}
fun beginTurn(sessionId: String?, eventScopeKey: String?, turnId: String) {
if (sessionId == null || sessionId != storedSessionId || eventScopeKey != scopeKey) return
if (activeTurnId == turnId) return
activeTurnId = turnId
sequence = 0L
laneSequence = 0L
_activities.value = emptyList()
}
fun onEvent(
sessionId: String?,
eventScopeKey: String?,
turnId: String,
event: GatewaySubagentEvent,
profile: String? = null,
) {
if (sessionId == null || sessionId != storedSessionId || eventScopeKey != scopeKey) return
if (activeTurnId != turnId) return
val taskIndex = event.taskIndex.coerceAtLeast(0)
val eventIdentity = event.subagentId?.takeIf(String::isNotBlank)
?: event.childSessionId?.takeIf(String::isNotBlank)
val identityMatch = eventIdentity?.let { identity ->
_activities.value.firstOrNull {
it.subagentId == identity || it.childSessionId == identity
}
}
val compatibleIndexMatches = _activities.value.filter { activity ->
activity.taskIndex == taskIndex &&
(event.subagentId.isNullOrBlank() || activity.subagentId.isNullOrBlank() ||
event.subagentId == activity.subagentId) &&
(event.childSessionId.isNullOrBlank() || activity.childSessionId.isNullOrBlank() ||
event.childSessionId == activity.childSessionId) &&
(event.parentId.isNullOrBlank() || activity.parentId.isNullOrBlank() ||
event.parentId == activity.parentId) &&
(event.depth == null || activity.depth == null || event.depth == activity.depth)
}
val current = identityMatch ?: compatibleIndexMatches.singleOrNull()
if (
current?.isTerminal == true &&
event.phase != GatewaySubagentEvent.Phase.SPAWN_REQUESTED &&
event.phase != GatewaySubagentEvent.Phase.START
) return
val base = if (current?.isTerminal == true) null else current
val phase = event.toActivityPhase()
val goal = sanitize(event.goal, MAX_GOAL_CHARS)
val preview = sanitize(event.preview, MAX_EVENT_TEXT_CHARS).ifBlank { null }
val summary = sanitize(event.summary, MAX_EVENT_TEXT_CHARS).ifBlank { null }
val toolName = sanitize(event.toolName, MAX_TOOL_NAME_CHARS).ifBlank { null }
val eventRow = SubagentActivityEvent(
sequence = sequence++,
kind = when (event.phase) {
GatewaySubagentEvent.Phase.SPAWN_REQUESTED,
GatewaySubagentEvent.Phase.START,
-> SubagentActivityEventKind.STARTED
GatewaySubagentEvent.Phase.THINKING,
GatewaySubagentEvent.Phase.PROGRESS,
-> SubagentActivityEventKind.UPDATE
GatewaySubagentEvent.Phase.TOOL -> SubagentActivityEventKind.TOOL
GatewaySubagentEvent.Phase.COMPLETE -> SubagentActivityEventKind.COMPLETED
},
text = if (event.phase == GatewaySubagentEvent.Phase.COMPLETE) summary else preview,
toolName = toolName,
phase = phase,
observedAtMillis = clock(),
)
val priorEvents = base?.events.orEmpty()
val coalesced = eventRow.kind == SubagentActivityEventKind.UPDATE &&
priorEvents.lastOrNull()?.let { previous ->
previous.kind == eventRow.kind && previous.text == eventRow.text
} == true
val appended = if (coalesced) priorEvents else priorEvents + eventRow
val (boundedEvents, truncated) = boundEvents(appended)
val next = SubagentActivity(
laneId = base?.laneId ?: laneSequence++,
turnId = turnId,
taskIndex = taskIndex,
taskCount = maxOf(1, event.taskCount, base?.taskCount ?: 1),
goal = goal.ifBlank { base?.goal.orEmpty() },
subagentId = event.subagentId?.takeIf(String::isNotBlank) ?: base?.subagentId,
childSessionId = event.childSessionId?.takeIf(String::isNotBlank) ?: base?.childSessionId,
parentId = event.parentId?.takeIf(String::isNotBlank) ?: base?.parentId,
depth = event.depth ?: base?.depth,
model = event.model?.takeIf(String::isNotBlank) ?: base?.model,
profile = profile?.takeIf(String::isNotBlank) ?: base?.profile,
phase = phase,
summary = summary ?: base?.summary,
durationSeconds = event.durationSeconds ?: base?.durationSeconds,
events = boundedEvents,
truncated = base?.truncated == true || truncated,
partialAfterGap = base?.partialAfterGap == true,
revision = (base?.revision ?: 0L) + 1,
)
_activities.value = (_activities.value.filterNot { it.stableKey == next.stableKey } + next)
.sortedWith(compareBy<SubagentActivity> { it.isTerminal }.thenBy { it.taskIndex })
}
fun endTurn(turnId: String) {
if (activeTurnId != turnId) return
_activities.value = _activities.value.map { activity ->
if (activity.isTerminal) activity else activity.copy(
phase = SubagentActivityPhase.ENDED_WITH_PARENT,
partialAfterGap = true,
revision = activity.revision + 1,
)
}
}
private fun boundEvents(
events: List<SubagentActivityEvent>,
): Pair<List<SubagentActivityEvent>, Boolean> {
val bounded = events.toMutableList()
var truncated = false
fun charCount(): Int = bounded.sumOf { (it.text?.length ?: 0) + (it.toolName?.length ?: 0) }
while (bounded.size > MAX_EVENTS_PER_CHILD || charCount() > MAX_CHARS_PER_CHILD) {
if (bounded.size <= 1) break
bounded.removeAt(if (bounded.first().kind == SubagentActivityEventKind.STARTED) 1 else 0)
truncated = true
}
return bounded to truncated
}
}
private fun GatewaySubagentEvent.toActivityPhase(): SubagentActivityPhase = when (phase) {
GatewaySubagentEvent.Phase.SPAWN_REQUESTED,
GatewaySubagentEvent.Phase.START -> SubagentActivityPhase.STARTED
GatewaySubagentEvent.Phase.THINKING -> SubagentActivityPhase.THINKING
GatewaySubagentEvent.Phase.TOOL -> SubagentActivityPhase.TOOL
GatewaySubagentEvent.Phase.PROGRESS -> SubagentActivityPhase.PROGRESS
GatewaySubagentEvent.Phase.COMPLETE -> when (status?.trim()?.lowercase()) {
"failed", "error" -> SubagentActivityPhase.FAILED
"interrupted", "cancelled", "canceled" -> SubagentActivityPhase.INTERRUPTED
else -> SubagentActivityPhase.COMPLETED
}
}
private val ANSI_ESCAPE = Regex("\\u001B(?:\\[[0-?]*[ -/]*[@-~]|\\][^\\u0007]*(?:\\u0007|\\u001B\\\\))")
private fun sanitize(value: String?, maxChars: Int): String = value.orEmpty()
.replace(ANSI_ESCAPE, "")
.filter { it == '\n' || it == '\t' || it >= ' ' }
.trim()
.take(maxChars)
@@ -1,317 +0,0 @@
package com.hermesandroid.relay.viewmodel
import com.hermesandroid.relay.data.ChatMessage
import com.hermesandroid.relay.data.MessageRole
import com.hermesandroid.relay.network.upstream.ChatHandler
import com.hermesandroid.relay.network.upstream.GatewayChatClient
import com.hermesandroid.relay.network.upstream.GatewayChildWatch
import com.hermesandroid.relay.network.upstream.GatewayTurnCallbacks
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.launch
import java.util.concurrent.atomic.AtomicLong
internal data class SubagentChildPreview(
val activityKey: String,
val parentSessionId: String,
val parentScopeKey: String?,
/** null while opening, false for a truthful parent-event fallback. */
val childWatchAvailable: Boolean? = null,
val messages: List<ChatMessage> = emptyList(),
val running: Boolean = false,
val status: String? = null,
val historyTruncated: Boolean = false,
val partialAfterGap: Boolean = false,
val error: String? = null,
)
internal class SubagentChildPreviewController(
private val scope: CoroutineScope,
private val openWatch: suspend (
GatewayChatClient,
String,
String?,
GatewayTurnCallbacks,
) -> Result<GatewayChildWatch> = { client, sessionId, profile, callbacks ->
client.openChildWatch(sessionId, profile, callbacks)
},
private val closeWatch: suspend (GatewayChatClient, GatewayChildWatch) -> Result<Unit> =
{ client, watch -> client.closeChildWatch(watch) },
) {
private class WatchContext(
val activity: SubagentActivity,
val client: GatewayChatClient,
val parentSessionId: String,
val parentScopeKey: String?,
val generation: Long,
val stillOwnsParent: () -> Boolean,
) {
val handler = ChatHandler()
var messageOrdinal = 0
var messageId = "child-watch-$generation-0"
var contentTruncated = false
var initialized = false
var pendingOverflow = false
val pendingCallbacks = mutableListOf<() -> Unit>()
}
private val _state = MutableStateFlow<SubagentChildPreview?>(null)
val state: StateFlow<SubagentChildPreview?> = _state.asStateFlow()
private val generation = AtomicLong(0)
private var watch: GatewayChildWatch? = null
private var watchClient: GatewayChatClient? = null
fun open(
activity: SubagentActivity,
client: GatewayChatClient?,
parentSessionId: String,
parentScopeKey: String?,
gatewayRouteActive: Boolean,
stillOwnsParent: () -> Boolean,
) {
if (isAlreadyOpen(activity.stableKey, parentSessionId, parentScopeKey)) return
close(clearState = false)
if (activity.childSessionId.isNullOrBlank() || client == null || !gatewayRouteActive) {
_state.value = fallbackState(activity, parentSessionId, parentScopeKey)
return
}
val context = WatchContext(
activity = activity,
client = client,
parentSessionId = parentSessionId,
parentScopeKey = parentScopeKey,
generation = generation.incrementAndGet(),
stillOwnsParent = stillOwnsParent,
)
_state.value = baseState(context)
// Once upstream creates a lazy watcher, only the resume acknowledgement
// reveals the live id needed to close it. Let a dismissed open finish;
// generation invalidation makes [acceptOpenedWatch] close the late handle.
scope.launch { openWatch(context) }
}
fun close() = close(clearState = true)
private suspend fun openWatch(context: WatchContext) {
openWatch(
context.client,
context.activity.childSessionId.orEmpty(),
context.activity.profile,
callbacks(context),
).fold(
onSuccess = { opened -> acceptOpenedWatch(context, opened) },
onFailure = { error -> publishOpenFailure(context, error.message) },
)
}
private suspend fun acceptOpenedWatch(context: WatchContext, opened: GatewayChildWatch) {
if (!owns(context)) {
closeWatch(context.client, opened)
return
}
watch = opened
watchClient = context.client
context.handler.setSessionId(opened.storedSessionId)
context.handler.loadMessageHistory(opened.messages)
context.contentTruncated = context.handler.boundReadOnlyPreview()
val pending = synchronized(context.pendingCallbacks) {
context.initialized = true
context.pendingCallbacks.toList().also { context.pendingCallbacks.clear() }
}
publish(
context = context,
running = opened.running,
status = opened.status,
historyTruncated = opened.historyTruncated || context.contentTruncated,
partial = context.activity.partialAfterGap || context.pendingOverflow,
)
pending.forEach { callback -> if (owns(context)) callback() }
}
private fun callbacks(context: WatchContext) = GatewayTurnCallbacks(
onSessionId = { },
onStart = { runOrQueue(context) { startMessage(context) } },
onTextDelta = { delta ->
runOrQueue(context) { mutate(context) { onTextDelta(context.messageId, delta) } }
},
onThinkingDelta = { delta ->
runOrQueue(context) { mutate(context) { onThinkingDelta(context.messageId, delta) } }
},
onToolCallStart = { id, name, preview ->
runOrQueue(context) {
mutate(context) { onToolCallStart(context.messageId, id, name, preview) }
}
},
onToolCallDone = { id, preview ->
runOrQueue(context) {
mutate(context, ensureMessage = false) {
onToolCallComplete(context.messageId, id, preview)
}
}
},
onToolCallFailed = { id, error ->
runOrQueue(context) {
mutate(context, ensureMessage = false) {
onToolCallFailed(context.messageId, id, error)
}
}
},
onTurnComplete = {
runOrQueue(context) {
if (owns(context)) context.handler.onTurnComplete(context.messageId)
}
},
onReconcileRequired = { runOrQueue(context) { publish(context, partial = true) } },
onComplete = { runOrQueue(context) { complete(context) } },
onUsage = { },
onError = { message ->
runOrQueue(context) {
publish(context, running = false, partial = true, error = message)
}
},
onToolGenerating = { },
onSubagentEvent = { event ->
runOrQueue(context) { mutate(context) { onSubagentEvent(context.messageId, event) } }
},
onMoaReference = { },
onInteractionRequest = { },
onInteractionExpired = { },
onResumeFailure = { message ->
runOrQueue(context) {
publish(context, running = false, partial = true, error = message)
}
},
)
private fun runOrQueue(context: WatchContext, callback: () -> Unit) {
if (!owns(context)) return
val runNow = synchronized(context.pendingCallbacks) {
if (context.initialized) {
true
} else {
if (context.pendingCallbacks.size >= 256) {
context.pendingCallbacks.removeAt(0)
context.pendingOverflow = true
}
context.pendingCallbacks += callback
false
}
}
if (runNow) callback()
}
private fun startMessage(context: WatchContext) {
if (!owns(context)) return
context.messageId = "child-watch-${context.generation}-${context.messageOrdinal++}"
ensureLiveMessage(context)
publish(context)
}
private inline fun mutate(
context: WatchContext,
ensureMessage: Boolean = true,
mutation: ChatHandler.() -> Unit,
) {
if (!owns(context)) return
if (ensureMessage) ensureLiveMessage(context)
context.handler.mutation()
context.contentTruncated = context.handler.boundReadOnlyPreview() || context.contentTruncated
publish(context)
}
private fun complete(context: WatchContext) {
if (!owns(context)) return
context.handler.onStreamComplete(context.messageId)
// The child mirror's message.complete omits failed/interrupted status.
// Keep this neutral; the parent activity lane is authoritative.
publish(context, running = false)
}
private fun ensureLiveMessage(context: WatchContext) {
if (context.handler.messages.value.any { it.id == context.messageId }) return
context.handler.addPlaceholderMessage(
ChatMessage(
id = context.messageId,
role = MessageRole.ASSISTANT,
content = "",
timestamp = System.currentTimeMillis(),
isStreaming = true,
),
)
}
private fun publish(
context: WatchContext,
running: Boolean = true,
status: String? = _state.value?.status,
historyTruncated: Boolean =
_state.value?.historyTruncated == true || context.contentTruncated,
partial: Boolean = _state.value?.partialAfterGap == true,
error: String? = null,
) {
if (!owns(context)) return
_state.value = baseState(context).copy(
childWatchAvailable = true,
messages = context.handler.messages.value.takeLast(200),
running = running,
status = status,
historyTruncated = historyTruncated,
partialAfterGap = partial,
error = error,
)
}
private fun publishOpenFailure(context: WatchContext, message: String?) {
if (!owns(context)) return
_state.value = fallbackState(
context.activity,
context.parentSessionId,
context.parentScopeKey,
).copy(error = message)
}
private fun owns(context: WatchContext): Boolean =
generation.get() == context.generation && context.stillOwnsParent()
private fun isAlreadyOpen(key: String, sessionId: String, scopeKey: String?): Boolean =
_state.value?.let {
it.activityKey == key &&
it.parentSessionId == sessionId &&
it.parentScopeKey == scopeKey &&
it.error.isNullOrBlank() &&
it.childWatchAvailable != false
} == true
private fun baseState(context: WatchContext) = SubagentChildPreview(
activityKey = context.activity.stableKey,
parentSessionId = context.parentSessionId,
parentScopeKey = context.parentScopeKey,
)
private fun fallbackState(
activity: SubagentActivity,
parentSessionId: String,
parentScopeKey: String?,
) = SubagentChildPreview(
activityKey = activity.stableKey,
parentSessionId = parentSessionId,
parentScopeKey = parentScopeKey,
childWatchAvailable = false,
partialAfterGap = activity.partialAfterGap,
)
private fun close(clearState: Boolean) {
generation.incrementAndGet()
val closingWatch = watch
val closingClient = watchClient
watch = null
watchClient = null
if (closingWatch != null && closingClient != null) {
scope.launch { closeWatch(closingClient, closingWatch) }
}
if (clearState) _state.value = null
}
}
@@ -4280,38 +4280,4 @@
<plurals name="chat_git_change_count"><item quantity="one">%1$d alteração</item><item quantity="other">%1$d alterações</item></plurals>
<string name="settings_git_workspace">Espaço de trabalho Git</string>
<string name="settings_git_workspace_desc">Revise alterações, branches, commits e remotos</string>
<string name="current_chat_activity_title">Atividade atual do chat</string>
<string name="current_chat_activity_subtitle">Detalhes ao vivo e somente leitura deste chat</string>
<string name="current_chat_activity_open">Visualizar atividade atual do chat</string>
<string name="current_chat_activity_summary">%1$d agentes · %2$d processos</string>
<string name="current_chat_activity_close">Fechar prévia da atividade</string>
<string name="current_chat_activity_empty">Nenhuma atividade atual neste chat</string>
<string name="current_chat_activity_latest">Mais recente</string>
<string name="agent_activity_section">Atividade ao vivo dos agentes</string>
<string name="agent_activity_disclosure">Atualizações recebidas por este chat. O histórico completo do agente filho pode não estar disponível.</string>
<string name="agent_activity_fallback">Agente %d</string>
<string name="agent_activity_task_position">%1$d de %2$d · %3$s</string>
<string name="agent_activity_lane_a11y">%1$s, %2$s, agente %3$d de %4$d</string>
<string name="agent_activity_duration">%1$.1fs</string>
<string name="agent_activity_older_omitted">Atividade anterior omitida</string>
<string name="agent_activity_partial">Atualizações ao vivo retomadas. Pode faltar atividade do período offline.</string>
<string name="agent_activity_event_started">Iniciado</string>
<string name="agent_activity_event_update">Atualização</string>
<string name="agent_activity_event_tool">Prévia da ferramenta</string>
<string name="agent_activity_status_started">Iniciando</string>
<string name="agent_activity_status_thinking">Pensando</string>
<string name="agent_activity_status_tool">Usando uma ferramenta</string>
<string name="agent_activity_status_progress">Trabalhando</string>
<string name="agent_activity_status_completed">Concluído</string>
<string name="agent_activity_status_failed">Falhou</string>
<string name="agent_activity_status_interrupted">Interrompido</string>
<string name="agent_activity_status_unavailable">Estado final indisponível</string>
<string name="agent_activity_child_loading">Abrindo histórico filho somente leitura…</string>
<string name="agent_activity_child_unavailable">O histórico filho não está disponível nesta versão ou rota do Hermes. A atividade da sessão principal é mostrada acima.</string>
<string name="agent_activity_child_live">Histórico filho · atualizações ao vivo</string>
<string name="agent_activity_child_history">Histórico filho · somente leitura</string>
<string name="agent_activity_child_truncated">Atividade filha recente exibida · detalhes antigos ou muito grandes foram omitidos</string>
<string name="agent_activity_child_role_task">Tarefa</string>
<string name="agent_activity_child_role_agent">Agente</string>
<string name="agent_activity_child_role_system">Sistema</string>
</resources>
@@ -4362,38 +4362,4 @@
<plurals name="chat_git_change_count"><item quantity="other">%1$d 个更改</item></plurals>
<string name="settings_git_workspace">Git 工作区</string>
<string name="settings_git_workspace_desc">查看更改、分支、提交和远程仓库</string>
<string name="current_chat_activity_title">当前聊天活动</string>
<string name="current_chat_activity_subtitle">此聊天中的只读实时详情</string>
<string name="current_chat_activity_open">预览当前聊天活动</string>
<string name="current_chat_activity_summary">%1$d 个代理 · %2$d 个进程</string>
<string name="current_chat_activity_close">关闭活动预览</string>
<string name="current_chat_activity_empty">此聊天当前没有活动</string>
<string name="current_chat_activity_latest">最新</string>
<string name="agent_activity_section">实时代理活动</string>
<string name="agent_activity_disclosure">此聊天接收到的更新。可能无法获取子代理的完整历史记录。</string>
<string name="agent_activity_fallback">代理 %d</string>
<string name="agent_activity_task_position">第 %1$d 个,共 %2$d 个 · %3$s</string>
<string name="agent_activity_lane_a11y">%1$s,%2$s,第 %3$d 个代理,共 %4$d 个</string>
<string name="agent_activity_duration">%1$.1f 秒</string>
<string name="agent_activity_older_omitted">已省略较早的活动</string>
<string name="agent_activity_partial">实时更新已恢复。离线期间的活动可能缺失。</string>
<string name="agent_activity_event_started">已开始</string>
<string name="agent_activity_event_update">更新</string>
<string name="agent_activity_event_tool">工具预览</string>
<string name="agent_activity_status_started">正在启动</string>
<string name="agent_activity_status_thinking">正在思考</string>
<string name="agent_activity_status_tool">正在使用工具</string>
<string name="agent_activity_status_progress">正在工作</string>
<string name="agent_activity_status_completed">已完成</string>
<string name="agent_activity_status_failed">失败</string>
<string name="agent_activity_status_interrupted">已中断</string>
<string name="agent_activity_status_unavailable">最终状态不可用</string>
<string name="agent_activity_child_loading">正在打开只读子历史记录…</string>
<string name="agent_activity_child_unavailable">此 Hermes 版本或路由不提供子历史记录。上方显示父会话活动。</string>
<string name="agent_activity_child_live">子历史记录 · 实时更新</string>
<string name="agent_activity_child_history">子历史记录 · 只读</string>
<string name="agent_activity_child_truncated">显示近期子活动 · 已省略较早或过大的详情</string>
<string name="agent_activity_child_role_task">任务</string>
<string name="agent_activity_child_role_agent">代理</string>
<string name="agent_activity_child_role_system">系统</string>
</resources>
-34
View File
@@ -4437,38 +4437,4 @@
<plurals name="chat_git_change_count"><item quantity="one">%1$d Änderung</item><item quantity="other">%1$d Änderungen</item></plurals>
<string name="settings_git_workspace">Git-Arbeitsbereich</string>
<string name="settings_git_workspace_desc">Änderungen, Branches, Commits und Remotes prüfen</string>
<string name="current_chat_activity_title">Aktuelle Chat-Aktivität</string>
<string name="current_chat_activity_subtitle">Schreibgeschützte Live-Details aus diesem Chat</string>
<string name="current_chat_activity_open">Aktuelle Chat-Aktivität ansehen</string>
<string name="current_chat_activity_summary">%1$d Agenten · %2$d Prozesse</string>
<string name="current_chat_activity_close">Aktivitätsvorschau schließen</string>
<string name="current_chat_activity_empty">Keine aktuelle Aktivität in diesem Chat</string>
<string name="current_chat_activity_latest">Neueste</string>
<string name="agent_activity_section">Live-Agentenaktivität</string>
<string name="agent_activity_disclosure">Von diesem Chat empfangene Updates. Der vollständige Verlauf des untergeordneten Agenten ist möglicherweise nicht verfügbar.</string>
<string name="agent_activity_fallback">Agent %d</string>
<string name="agent_activity_task_position">%1$d von %2$d · %3$s</string>
<string name="agent_activity_lane_a11y">%1$s, %2$s, Agent %3$d von %4$d</string>
<string name="agent_activity_duration">%1$.1fs</string>
<string name="agent_activity_older_omitted">Ältere Aktivität ausgelassen</string>
<string name="agent_activity_partial">Live-Updates fortgesetzt. Aktivität während der Offlinezeit kann fehlen.</string>
<string name="agent_activity_event_started">Gestartet</string>
<string name="agent_activity_event_update">Update</string>
<string name="agent_activity_event_tool">Werkzeugvorschau</string>
<string name="agent_activity_status_started">Wird gestartet</string>
<string name="agent_activity_status_thinking">Denkt nach</string>
<string name="agent_activity_status_tool">Verwendet ein Werkzeug</string>
<string name="agent_activity_status_progress">Arbeitet</string>
<string name="agent_activity_status_completed">Abgeschlossen</string>
<string name="agent_activity_status_failed">Fehlgeschlagen</string>
<string name="agent_activity_status_interrupted">Unterbrochen</string>
<string name="agent_activity_status_unavailable">Endstatus nicht verfügbar</string>
<string name="agent_activity_child_loading">Schreibgeschützter untergeordneter Verlauf wird geöffnet…</string>
<string name="agent_activity_child_unavailable">Der untergeordnete Verlauf ist in dieser Hermes-Version oder Route nicht verfügbar. Die Aktivität der übergeordneten Sitzung wird oben angezeigt.</string>
<string name="agent_activity_child_live">Untergeordneter Verlauf · Live-Updates</string>
<string name="agent_activity_child_history">Untergeordneter Verlauf · schreibgeschützt</string>
<string name="agent_activity_child_truncated">Neueste untergeordnete Aktivität angezeigt · ältere oder zu große Details ausgelassen</string>
<string name="agent_activity_child_role_task">Aufgabe</string>
<string name="agent_activity_child_role_agent">Agent</string>
<string name="agent_activity_child_role_system">System</string>
</resources>
-34
View File
@@ -4128,38 +4128,4 @@
<plurals name="chat_git_change_count"><item quantity="one">%1$d cambio</item><item quantity="other">%1$d cambios</item></plurals>
<string name="settings_git_workspace">Espacio de Git</string>
<string name="settings_git_workspace_desc">Revisa cambios, ramas, commits y remotos</string>
<string name="current_chat_activity_title">Actividad actual del chat</string>
<string name="current_chat_activity_subtitle">Detalles en vivo de solo lectura de este chat</string>
<string name="current_chat_activity_open">Ver la actividad actual del chat</string>
<string name="current_chat_activity_summary">%1$d agentes · %2$d procesos</string>
<string name="current_chat_activity_close">Cerrar vista previa de actividad</string>
<string name="current_chat_activity_empty">No hay actividad actual en este chat</string>
<string name="current_chat_activity_latest">Más reciente</string>
<string name="agent_activity_section">Actividad de agentes en vivo</string>
<string name="agent_activity_disclosure">Actualizaciones recibidas por este chat. Es posible que el historial completo del agente secundario no esté disponible.</string>
<string name="agent_activity_fallback">Agente %d</string>
<string name="agent_activity_task_position">%1$d de %2$d · %3$s</string>
<string name="agent_activity_lane_a11y">%1$s, %2$s, agente %3$d de %4$d</string>
<string name="agent_activity_duration">%1$.1fs</string>
<string name="agent_activity_older_omitted">Se omitió la actividad anterior</string>
<string name="agent_activity_partial">Se reanudaron las actualizaciones en vivo. Puede faltar actividad mientras estaba sin conexión.</string>
<string name="agent_activity_event_started">Iniciado</string>
<string name="agent_activity_event_update">Actualización</string>
<string name="agent_activity_event_tool">Vista previa de herramienta</string>
<string name="agent_activity_status_started">Iniciando</string>
<string name="agent_activity_status_thinking">Pensando</string>
<string name="agent_activity_status_tool">Usando una herramienta</string>
<string name="agent_activity_status_progress">Trabajando</string>
<string name="agent_activity_status_completed">Completado</string>
<string name="agent_activity_status_failed">Falló</string>
<string name="agent_activity_status_interrupted">Interrumpido</string>
<string name="agent_activity_status_unavailable">Estado final no disponible</string>
<string name="agent_activity_child_loading">Abriendo el historial secundario de solo lectura…</string>
<string name="agent_activity_child_unavailable">El historial secundario no está disponible en esta versión o ruta de Hermes. La actividad de la sesión principal se muestra arriba.</string>
<string name="agent_activity_child_live">Historial secundario · actualizaciones en vivo</string>
<string name="agent_activity_child_history">Historial secundario · solo lectura</string>
<string name="agent_activity_child_truncated">Se muestra la actividad secundaria reciente · se omitieron detalles anteriores o demasiado grandes</string>
<string name="agent_activity_child_role_task">Tarea</string>
<string name="agent_activity_child_role_agent">Agente</string>
<string name="agent_activity_child_role_system">Sistema</string>
</resources>
-34
View File
@@ -4433,38 +4433,4 @@
<plurals name="chat_git_change_count"><item quantity="other">%1$d 件の変更</item></plurals>
<string name="settings_git_workspace">Git ワークスペース</string>
<string name="settings_git_workspace_desc">変更、ブランチ、コミット、リモートを確認</string>
<string name="current_chat_activity_title">現在のチャットのアクティビティ</string>
<string name="current_chat_activity_subtitle">このチャットからの読み取り専用ライブ詳細</string>
<string name="current_chat_activity_open">現在のチャットのアクティビティを表示</string>
<string name="current_chat_activity_summary">エージェント %1$d · プロセス %2$d</string>
<string name="current_chat_activity_close">アクティビティのプレビューを閉じる</string>
<string name="current_chat_activity_empty">このチャットに現在のアクティビティはありません</string>
<string name="current_chat_activity_latest">最新</string>
<string name="agent_activity_section">エージェントのライブアクティビティ</string>
<string name="agent_activity_disclosure">このチャットが受信した更新です。子エージェントの完全な履歴は利用できない場合があります。</string>
<string name="agent_activity_fallback">エージェント %d</string>
<string name="agent_activity_task_position">%2$d 件中 %1$d 件目 · %3$s</string>
<string name="agent_activity_lane_a11y">%1$s、%2$s、%4$d 件中 %3$d 件目のエージェント</string>
<string name="agent_activity_duration">%1$.1f秒</string>
<string name="agent_activity_older_omitted">古いアクティビティは省略されました</string>
<string name="agent_activity_partial">ライブ更新を再開しました。オフライン中のアクティビティが欠けている場合があります。</string>
<string name="agent_activity_event_started">開始</string>
<string name="agent_activity_event_update">更新</string>
<string name="agent_activity_event_tool">ツールのプレビュー</string>
<string name="agent_activity_status_started">開始中</string>
<string name="agent_activity_status_thinking">思考中</string>
<string name="agent_activity_status_tool">ツールを使用中</string>
<string name="agent_activity_status_progress">作業中</string>
<string name="agent_activity_status_completed">完了</string>
<string name="agent_activity_status_failed">失敗</string>
<string name="agent_activity_status_interrupted">中断</string>
<string name="agent_activity_status_unavailable">最終状態を確認できません</string>
<string name="agent_activity_child_loading">読み取り専用の子履歴を開いています…</string>
<string name="agent_activity_child_unavailable">この Hermes のバージョンまたはルートでは子履歴を利用できません。親セッションのアクティビティは上に表示されます。</string>
<string name="agent_activity_child_live">子履歴 · ライブ更新</string>
<string name="agent_activity_child_history">子履歴 · 読み取り専用</string>
<string name="agent_activity_child_truncated">最近の子アクティビティを表示 · 古い詳細または大きすぎる詳細は省略されました</string>
<string name="agent_activity_child_role_task">タスク</string>
<string name="agent_activity_child_role_agent">エージェント</string>
<string name="agent_activity_child_role_system">システム</string>
</resources>
-34
View File
@@ -4174,38 +4174,4 @@
<plurals name="chat_git_change_count"><item quantity="one">%1$d изменение</item><item quantity="few">%1$d изменения</item><item quantity="many">%1$d изменений</item><item quantity="other">%1$d изменения</item></plurals>
<string name="settings_git_workspace">Рабочая область Git</string>
<string name="settings_git_workspace_desc">Изменения, ветки, коммиты и удалённые репозитории</string>
<string name="current_chat_activity_title">Текущая активность чата</string>
<string name="current_chat_activity_subtitle">Доступные только для чтения сведения в реальном времени из этого чата</string>
<string name="current_chat_activity_open">Просмотреть текущую активность чата</string>
<string name="current_chat_activity_summary">Агенты: %1$d · процессы: %2$d</string>
<string name="current_chat_activity_close">Закрыть просмотр активности</string>
<string name="current_chat_activity_empty">В этом чате сейчас нет активности</string>
<string name="current_chat_activity_latest">Последнее</string>
<string name="agent_activity_section">Активность агентов в реальном времени</string>
<string name="agent_activity_disclosure">Обновления, полученные этим чатом. Полная история дочернего агента может быть недоступна.</string>
<string name="agent_activity_fallback">Агент %d</string>
<string name="agent_activity_task_position">%1$d из %2$d · %3$s</string>
<string name="agent_activity_lane_a11y">%1$s, %2$s, агент %3$d из %4$d</string>
<string name="agent_activity_duration">%1$.1f с</string>
<string name="agent_activity_older_omitted">Более ранняя активность опущена</string>
<string name="agent_activity_partial">Обновления возобновлены. Активность во время отсутствия подключения может быть пропущена.</string>
<string name="agent_activity_event_started">Запущено</string>
<string name="agent_activity_event_update">Обновление</string>
<string name="agent_activity_event_tool">Предпросмотр инструмента</string>
<string name="agent_activity_status_started">Запуск</string>
<string name="agent_activity_status_thinking">Размышляет</string>
<string name="agent_activity_status_tool">Использует инструмент</string>
<string name="agent_activity_status_progress">Работает</string>
<string name="agent_activity_status_completed">Завершено</string>
<string name="agent_activity_status_failed">Ошибка</string>
<string name="agent_activity_status_interrupted">Прервано</string>
<string name="agent_activity_status_unavailable">Итоговое состояние недоступно</string>
<string name="agent_activity_child_loading">Открывается дочерняя история только для чтения…</string>
<string name="agent_activity_child_unavailable">Дочерняя история недоступна в этой версии или маршруте Hermes. Активность родительского сеанса показана выше.</string>
<string name="agent_activity_child_live">Дочерняя история · обновления в реальном времени</string>
<string name="agent_activity_child_history">Дочерняя история · только чтение</string>
<string name="agent_activity_child_truncated">Показана недавняя дочерняя активность · более ранние или слишком большие сведения опущены</string>
<string name="agent_activity_child_role_task">Задача</string>
<string name="agent_activity_child_role_agent">Агент</string>
<string name="agent_activity_child_role_system">Система</string>
</resources>
-34
View File
@@ -3735,40 +3735,6 @@
<string name="tool_failed_a11y">Failed</string>
<string name="background_process_count">Background · %1$d</string>
<string name="background_processes_title">Background processes</string>
<string name="current_chat_activity_title">Current chat activity</string>
<string name="current_chat_activity_subtitle">Read-only live details from this chat</string>
<string name="current_chat_activity_open">Preview current chat activity</string>
<string name="current_chat_activity_summary">%1$d agents · %2$d processes</string>
<string name="current_chat_activity_close">Close activity preview</string>
<string name="current_chat_activity_empty">No current activity in this chat</string>
<string name="current_chat_activity_latest">Latest</string>
<string name="agent_activity_section">Live agent activity</string>
<string name="agent_activity_disclosure">Updates received by this chat. Full child history may be unavailable.</string>
<string name="agent_activity_fallback">Agent %d</string>
<string name="agent_activity_task_position">%1$d of %2$d · %3$s</string>
<string name="agent_activity_lane_a11y">%1$s, %2$s, agent %3$d of %4$d</string>
<string name="agent_activity_duration">%1$.1fs</string>
<string name="agent_activity_older_omitted">Older activity omitted</string>
<string name="agent_activity_partial">Live updates resumed. Activity while offline may be missing.</string>
<string name="agent_activity_event_started">Started</string>
<string name="agent_activity_event_update">Update</string>
<string name="agent_activity_event_tool">Tool preview</string>
<string name="agent_activity_status_started">Starting</string>
<string name="agent_activity_status_thinking">Thinking</string>
<string name="agent_activity_status_tool">Using a tool</string>
<string name="agent_activity_status_progress">Working</string>
<string name="agent_activity_status_completed">Completed</string>
<string name="agent_activity_status_failed">Failed</string>
<string name="agent_activity_status_interrupted">Interrupted</string>
<string name="agent_activity_status_unavailable">Final state unavailable</string>
<string name="agent_activity_child_loading">Opening read-only child history…</string>
<string name="agent_activity_child_unavailable">Child transcript unavailable on this Hermes version or route. Parent-session activity is shown above.</string>
<string name="agent_activity_child_live">Child history · live updates</string>
<string name="agent_activity_child_history">Child history · read only</string>
<string name="agent_activity_child_truncated">Recent child activity shown · older or oversized details omitted</string>
<string name="agent_activity_child_role_task">Task</string>
<string name="agent_activity_child_role_agent">Agent</string>
<string name="agent_activity_child_role_system">System</string>
<string name="background_processes_refresh_a11y">Refresh processes</string>
<string name="background_processes_subtitle">Current chat · live output and recent results</string>
<string name="background_processes_stop">Stop</string>
@@ -277,40 +277,4 @@ class AgentDisplayTest {
assertEquals("conn::__server_default__", AgentDisplay.profileContextKey("conn", null))
assertEquals("conn::mizu", AgentDisplay.profileContextKey("conn", "mizu"))
}
@Test
fun parseProfileContextKey_preservesRequestIdentityAndConnectionScope() {
val serverDefault = AgentDisplay.parseProfileContextKey(
AgentDisplay.profileContextKey("connection-a", null),
)
assertEquals("connection-a", serverDefault?.connectionId)
assertEquals(AgentDisplay.SERVER_DEFAULT_PROFILE_KEY, serverDefault?.profileKey)
assertNull(serverDefault?.requestProfileName)
val literalDefault = AgentDisplay.parseProfileContextKey(
AgentDisplay.profileContextKey("connection-a", "default"),
)
assertEquals("default", literalDefault?.profileKey)
assertEquals("default", literalDefault?.requestProfileName)
val named = AgentDisplay.parseProfileContextKey(
AgentDisplay.profileContextKey("connection-b", "mizu"),
)
assertEquals("connection-b", named?.connectionId)
assertEquals("mizu", named?.requestProfileName)
val delimitedProfile = AgentDisplay.parseProfileContextKey(
AgentDisplay.profileContextKey("connection-c", "team::writer"),
)
assertEquals("connection-c", delimitedProfile?.connectionId)
assertEquals("team::writer", delimitedProfile?.requestProfileName)
}
@Test
fun parseProfileContextKey_failsClosedForLegacyOrMalformedKeys() {
assertNull(AgentDisplay.parseProfileContextKey("connection/profile-default"))
assertNull(AgentDisplay.parseProfileContextKey("connection-a::"))
assertNull(AgentDisplay.parseProfileContextKey("::default"))
assertNull(AgentDisplay.parseProfileContextKey(null))
}
}
@@ -87,20 +87,6 @@ class ChatTurnCheckpointStoreTest {
assertEquals(checkpoint, store.read())
}
@Test
fun checkpointWithoutProfileKey_remainsReadableAsLegacyIdentity() = runTest {
val current = sampleCheckpoint().copy(profileKey = "default")
val legacyJson = Json.encodeToString(current)
.replace("\"profileKey\":\"default\",", "")
dataStore.edit { preferences ->
preferences[stringPreferencesKey("chat_inflight_turn_checkpoint_v1")] = legacyJson
}
val restored = store.read()
assertEquals(current.contextKey, restored?.contextKey)
assertNull(restored?.profileKey)
}
@Test
fun multipleRunningSessions_mergeAndRemoveIndependently() = runTest {
val first = sampleCheckpoint()
@@ -120,8 +106,7 @@ class ChatTurnCheckpointStoreTest {
}
private fun sampleCheckpoint() = ChatTurnCheckpoint(
contextKey = AgentDisplay.profileContextKey("connection-a", null),
profileKey = AgentDisplay.SERVER_DEFAULT_PROFILE_KEY,
contextKey = "connection-a/profile-default",
sessionId = "stored-42",
liveSessionId = "live-42",
transport = "gateway",
@@ -74,12 +74,6 @@ class GatewayClientHarness(
@Volatile
var recoveryAssistant = ""
@Volatile
var recoveryMessages: JsonArray = JsonArray(emptyList())
@Volatile
var resumeEventsBeforeAck: List<Pair<String, JsonObject?>> = emptyList()
@Volatile
var recoveryInflightStreaming: Boolean? = null
var recoveryInflightError: String? = null
@@ -237,11 +231,14 @@ 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)
webSocket.send(eventFrame("gateway.ready", null, null))
if (!suppressGatewayReady) sendGatewayReady(webSocket)
}
override fun onMessage(webSocket: WebSocket, text: String) {
@@ -545,16 +542,6 @@ class GatewayClientHarness(
put("error", buildJsonObject { put("message", "$method refused") })
}
}
if (method == "session.resume" && result != null) {
val liveId = (result["session_id"] as? JsonPrimitive)?.contentOrNull
val events = resumeEventsBeforeAck
resumeEventsBeforeAck = emptyList()
if (!liveId.isNullOrBlank()) {
events.forEach { (type, payload) ->
webSocket.send(eventFrame(type, payload, liveId))
}
}
}
webSocket.send(reply.toString())
}
}
@@ -565,7 +552,6 @@ class GatewayClientHarness(
put("session_id", sessionId)
put("running", recoveryRunning)
put("status", if (recoveryRunning) "streaming" else "idle")
put("messages", recoveryMessages)
if (!omitSessionProfileMetadata || recoveryProject != null) {
put("info", buildJsonObject {
if (!omitSessionProfileMetadata) {
@@ -644,6 +630,10 @@ 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) {
@@ -777,14 +767,13 @@ class GatewayChatClientTest {
rpcTimeoutMs: Long = 15_000L,
promptSubmitTimeoutMs: Long = 1_800_000L,
turnIdleTimeoutMs: Long = 180_000L,
callbackDispatcher: (block: () -> Unit) -> Unit = { it() },
) = GatewayChatClient(
initialDashboardClient = DashboardApiClient(
baseUrl = harness.server.url("/").toString().trimEnd('/'),
okHttpClient = OkHttpClient(),
),
okHttpClient = OkHttpClient(),
callbackDispatcher = callbackDispatcher,
callbackDispatcher = { it() },
onGatewayUnsupported = { unsupportedMarked = true },
scope = scope,
// Keep the mid-turn reconnect window short so `failed rejoin`
@@ -1425,6 +1414,27 @@ 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"
@@ -2200,196 +2210,6 @@ class GatewayChatClientTest {
assertEquals("focus on Android", (params["text"] as? JsonPrimitive)?.contentOrNull)
}
@Test
fun `child watch is profile pinned bounded and isolated from main session`() = runBlocking {
harness.sessionProfileOverride = "operator"
harness.resumeLiveSessionIds["parent-session"] = "live-parent"
harness.resumeLiveSessionIds["child-session"] = "live-child"
client.sessionProfileProvider = { "operator" }
assertTrue(client.prewarmAwait("parent-session"))
val serverWs = harness.awaitServerSocket()
assertEquals("live-parent", client.currentLiveSessionId("parent-session"))
harness.recoveryRunning = true
harness.recoveryMessages = JsonArray(listOf(
buildJsonObject { put("role", "user"); put("text", "old") },
buildJsonObject { put("role", "assistant"); put("text", "recent") },
buildJsonObject { put("role", "assistant"); put("text", "newest") },
))
val childRecorder = Recorder()
val watch = client.openChildWatch(
childSessionId = "child-session",
profile = "operator",
callbacks = childRecorder.callbacks,
historyLimit = 2,
).getOrThrow()
val resume = harness.awaitRpcCount("session.resume", 2).last()
assertEquals("child-session", (resume["session_id"] as? JsonPrimitive)?.contentOrNull)
assertEquals("operator", (resume["profile"] as? JsonPrimitive)?.contentOrNull)
assertEquals(true, (resume["lazy"] as? JsonPrimitive)?.booleanOrNull)
assertEquals(true, (resume["close_on_disconnect"] as? JsonPrimitive)?.booleanOrNull)
assertEquals("live-child", watch.liveSessionId)
assertTrue(watch.running)
assertTrue(watch.historyTruncated)
assertEquals(listOf("recent", "newest"), watch.messages.map { it.contentText })
assertEquals("live-parent", client.currentLiveSessionId("parent-session"))
serverWs.send(harness.eventFrame("message.start", null, "live-child"))
serverWs.send(
harness.eventFrame(
"reasoning.delta",
buildJsonObject { put("text", "checking") },
"live-child",
),
)
serverWs.send(
harness.eventFrame(
"message.delta",
buildJsonObject { put("text", "working") },
"live-child",
),
)
serverWs.send(
harness.eventFrame(
"message.complete",
buildJsonObject { put("text", "done") },
"live-child",
),
)
assertTrue(childRecorder.completeLatch.await(5, TimeUnit.SECONDS))
assertEquals(listOf("checking"), childRecorder.thinkingDeltas.toList())
assertTrue(childRecorder.textDeltas.contains("working"))
assertTrue(childRecorder.errors.isEmpty())
client.closeChildWatch(watch).getOrThrow()
val close = harness.awaitRpc("session.close")
assertEquals("live-child", (close["session_id"] as? JsonPrimitive)?.contentOrNull)
assertEquals("live-parent", client.currentLiveSessionId("parent-session"))
}
@Test
fun `concurrent child opens keep newest generation and stale close is harmless`() = runBlocking {
harness.resumeLiveSessionIds["child-session"] = "live-child"
val recorders = listOf(Recorder(), Recorder())
val opens = recorders.map { recorder ->
async(Dispatchers.IO) {
client.openChildWatch(
"child-session",
callbacks = recorder.callbacks,
).getOrThrow()
}
}
val watches = opens.map { it.await() }
harness.awaitServerSocket()
val stale = watches.minBy { it.generation }
val newest = watches.maxBy { it.generation }
client.closeChildWatch(stale).getOrThrow()
assertTrue(harness.rpcLog.none { it.first == "session.close" })
client.closeChildWatch(newest).getOrThrow()
val close = harness.awaitRpc("session.close")
assertEquals("live-child", (close["session_id"] as? JsonPrimitive)?.contentOrNull)
assertEquals(1, recorders.sumOf { it.resumeFailures.size })
}
@Test
fun `child watch replays terminal event that arrives before resume ack`() = runBlocking {
harness.resumeLiveSessionIds["child-session"] = "live-child"
harness.recoveryRunning = true
harness.resumeEventsBeforeAck = listOf(
"message.start" to null,
"message.delta" to buildJsonObject { put("text", "pre-ack") },
"message.complete" to buildJsonObject { put("text", "done") },
)
val recorder = Recorder()
val watch = client.openChildWatch(
"child-session",
callbacks = recorder.callbacks,
).getOrThrow()
harness.awaitServerSocket()
assertEquals("live-child", watch.liveSessionId)
assertFalse(watch.running)
assertTrue(recorder.completeLatch.await(5, TimeUnit.SECONDS))
assertTrue(recorder.textDeltas.contains("pre-ack"))
assertTrue(recorder.errors.isEmpty())
}
@Test
fun `failed child watch close can be retried`() = runBlocking {
harness.resumeLiveSessionIds["child-session"] = "live-child"
val watch = client.openChildWatch(
"child-session",
callbacks = Recorder().callbacks,
).getOrThrow()
harness.awaitServerSocket()
harness.rpcErrors["session.close"] = 5000 to "busy"
assertTrue(client.closeChildWatch(watch).isFailure)
harness.rpcErrors.remove("session.close")
client.closeChildWatch(watch).getOrThrow()
val closes = harness.awaitRpcCount("session.close", 2)
assertEquals(2, closes.size)
assertTrue(closes.all {
(it["session_id"] as? JsonPrimitive)?.contentOrNull == "live-child"
})
}
@Test
fun `queued child callback is dropped after exact watch closes`() = runBlocking {
client.shutdown()
scope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
val queuedCallbacks = ConcurrentLinkedQueue<() -> Unit>()
client = buildClient(callbackDispatcher = { queuedCallbacks += it })
harness.resumeLiveSessionIds["child-session"] = "live-child"
val recorder = Recorder()
val watch = client.openChildWatch(
"child-session",
callbacks = recorder.callbacks,
).getOrThrow()
val serverWs = harness.awaitServerSocket()
serverWs.send(
harness.eventFrame(
"message.delta",
buildJsonObject { put("text", "stale") },
"live-child",
),
)
awaitCondition { queuedCallbacks.isNotEmpty() }
client.closeChildWatch(watch).getOrThrow()
while (true) queuedCallbacks.poll()?.invoke() ?: break
assertTrue(recorder.textDeltas.isEmpty())
}
@Test
fun `child watch history enforces total character bound`() = runBlocking {
harness.resumeLiveSessionIds["child-session"] = "live-child"
harness.recoveryMessages = JsonArray(listOf(
buildJsonObject { put("role", "assistant"); put("text", "kept") },
buildJsonObject {
put("role", "assistant")
put("text", "x".repeat(GatewayChatClient.MAX_CHILD_WATCH_HISTORY_CHARS + 1))
},
))
val watch = client.openChildWatch(
"child-session",
callbacks = Recorder().callbacks,
).getOrThrow()
harness.awaitServerSocket()
assertTrue(watch.historyTruncated)
assertEquals(listOf("kept"), watch.messages.map { it.contentText })
}
@Test
fun `compress session uses dedicated rpc and parses authoritative messages`() {
val r = Recorder()
@@ -595,10 +595,6 @@ class GatewayEventMapperTest {
fun `subagent lifecycle maps phases and fields`() {
val r = Recorder()
val mapper = mapperWith(r)
mapper.onEvent(
"subagent.spawn_requested",
obj("""{"goal":"research topic","task_index":1,"task_count":3,"subagent_id":"child-17","child_session_id":"session-17","parent_id":"parent-child","depth":2,"model":"hermes-4"}"""),
)
mapper.onEvent("subagent.start", obj("""{"goal":"research topic","task_index":1,"task_count":3,"subagent_id":"child-17"}"""))
mapper.onEvent("subagent.thinking", obj("""{"goal":"research topic","task_index":1,"task_count":3,"text":"hmm"}"""))
mapper.onEvent(
@@ -613,7 +609,6 @@ class GatewayEventMapperTest {
assertEquals(
listOf(
GatewaySubagentEvent.Phase.SPAWN_REQUESTED,
GatewaySubagentEvent.Phase.START,
GatewaySubagentEvent.Phase.THINKING,
GatewaySubagentEvent.Phase.TOOL,
@@ -622,22 +617,17 @@ class GatewayEventMapperTest {
),
r.subagentEvents.map { it.phase },
)
val spawn = r.subagentEvents[0]
assertEquals("session-17", spawn.childSessionId)
assertEquals("parent-child", spawn.parentId)
assertEquals(2, spawn.depth)
assertEquals("hermes-4", spawn.model)
val start = r.subagentEvents[1]
val start = r.subagentEvents[0]
assertEquals(1, start.taskIndex)
assertEquals(3, start.taskCount)
assertEquals("research topic", start.goal)
assertEquals("child-17", start.subagentId)
assertEquals("hmm", r.subagentEvents[2].preview)
val tool = r.subagentEvents[3]
assertEquals("hmm", r.subagentEvents[1].preview)
val tool = r.subagentEvents[2]
assertEquals("web_search", tool.toolName)
assertEquals("searching docs", tool.preview)
assertEquals("halfway", r.subagentEvents[4].preview)
val complete = r.subagentEvents[5]
assertEquals("halfway", r.subagentEvents[3].preview)
val complete = r.subagentEvents[4]
assertEquals("completed", complete.status)
assertEquals("found it", complete.summary)
assertEquals(12.5, complete.durationSeconds!!, 0.001)
@@ -1067,7 +1057,7 @@ class GatewayEventMapperTest {
"message.complete", "error", "clarify.request", "approval.request",
"sudo.request", "secret.request", "reasoning.available",
"clarify.expire", "sudo.expire", "secret.expire", "approval.expire",
"tool.generating", "subagent.spawn_requested", "subagent.start", "subagent.thinking",
"tool.generating", "subagent.start", "subagent.thinking",
"subagent.tool", "subagent.progress", "subagent.complete",
"tool.output_risk", "moa.reference", "moa.progress", "moa.phase", "moa.aggregating",
).forEach { type ->
@@ -1,45 +0,0 @@
package com.hermesandroid.relay.network.upstream
import com.hermesandroid.relay.data.ChatMessage
import com.hermesandroid.relay.data.MessageRole
import com.hermesandroid.relay.data.ToolCall
import org.junit.Assert.assertEquals
import org.junit.Assert.assertFalse
import org.junit.Assert.assertTrue
import org.junit.Test
class ReadOnlyPreviewBoundsTest {
@Test
fun `child preview drops system rows results and oversized live content`() {
val handler = ChatHandler()
handler.addPlaceholderMessage(
ChatMessage("system", MessageRole.SYSTEM, "private system context", 1L),
)
handler.addPlaceholderMessage(
ChatMessage(
id = "child",
role = MessageRole.ASSISTANT,
content = "x".repeat(20_000),
timestamp = 2L,
thinkingContent = "y".repeat(20_000),
toolCalls = listOf(
ToolCall(
name = "read_file",
args = "a".repeat(5_000),
result = "secret result",
success = true,
),
),
),
)
assertTrue(handler.boundReadOnlyPreview(maxTotalChars = 4_000, maxFieldChars = 2_000))
val messages = handler.messages.value
assertEquals(listOf("child"), messages.map(ChatMessage::id))
assertTrue(messages.sumOf { it.content.length + it.thinkingContent.length } <= 4_000)
assertTrue(messages.single().toolCalls.single().args.orEmpty().length <= 1_000)
assertEquals(null, messages.single().toolCalls.single().result)
assertFalse(messages.any { it.role == MessageRole.SYSTEM })
}
}
@@ -1,105 +0,0 @@
package com.hermesandroid.relay.screenshots
import androidx.compose.ui.test.junit4.createComposeRule
import androidx.compose.ui.test.onNodeWithContentDescription
import androidx.compose.ui.test.onNodeWithText
import androidx.compose.ui.test.onRoot
import androidx.compose.ui.test.performClick
import androidx.test.ext.junit.runners.AndroidJUnit4
import com.github.takahirom.roborazzi.captureRoboImage
import com.hermesandroid.relay.data.ChatMessage
import com.hermesandroid.relay.data.MessageRole
import com.hermesandroid.relay.ui.components.GatewayBackgroundProcessSheet
import com.hermesandroid.relay.ui.components.SubagentPreviewVisibility
import com.hermesandroid.relay.ui.theme.HermesRelayTheme
import com.hermesandroid.relay.viewmodel.SubagentActivity
import com.hermesandroid.relay.viewmodel.SubagentActivityEvent
import com.hermesandroid.relay.viewmodel.SubagentActivityEventKind
import com.hermesandroid.relay.viewmodel.SubagentActivityPhase
import com.hermesandroid.relay.viewmodel.SubagentChildPreview
import org.junit.Rule
import org.junit.Test
import org.junit.runner.RunWith
import org.robolectric.annotation.Config
import org.robolectric.annotation.GraphicsMode
@RunWith(AndroidJUnit4::class)
@GraphicsMode(GraphicsMode.Mode.NATIVE)
@Config(qualifiers = "w360dp-h720dp-xhdpi")
class SubagentActivitySheetScreenshotTest {
@get:Rule val compose = createComposeRule()
@Test
fun concurrentLiveAgentsRenderAsReadOnlyActivity() {
val activities = listOf(
activity(0, 0, "Inspect Android event handling", SubagentActivityPhase.PROGRESS),
activity(1, 1, "Review privacy boundaries", SubagentActivityPhase.INTERRUPTED),
)
val preview = SubagentChildPreview(
activityKey = activities.first().stableKey,
parentSessionId = "parent",
parentScopeKey = "scope",
childWatchAvailable = true,
messages = listOf(
ChatMessage("task", MessageRole.USER, "Trace the upstream child watch contract.", 1L),
ChatMessage("answer", MessageRole.ASSISTANT, "The child-only history is available read-only.", 2L),
),
running = true,
status = "streaming",
)
compose.setContent {
HermesRelayTheme(appThemeId = "hermes-relay", themePreference = "dark") {
GatewayBackgroundProcessSheet(
processes = emptyList(),
subagentActivities = activities,
subagentChildPreview = preview,
subagentPreviewVisibility = SubagentPreviewVisibility(),
loading = false,
stoppingProcessIds = emptySet(),
onRefresh = {},
onStop = {},
onDismissProcess = {},
onOpenSubagentChild = {},
onDismiss = {},
)
}
}
compose.onNodeWithContentDescription(
"Inspect Android event handling, Working, agent 1 of 2",
).performClick()
compose.onNodeWithText("Child history · live updates").assertExists()
compose.onNodeWithText("The child-only history is available read-only.").assertExists()
compose.onNodeWithText("Stop").assertDoesNotExist()
compose.onRoot().captureRoboImage("build/ui-regression/subagent-activity-sheet.png")
}
private fun activity(
laneId: Long,
taskIndex: Int,
goal: String,
phase: SubagentActivityPhase,
) = SubagentActivity(
laneId = laneId,
turnId = "turn",
taskIndex = taskIndex,
taskCount = 2,
goal = goal,
phase = phase,
childSessionId = "child-$taskIndex",
profile = "default",
events = listOf(
SubagentActivityEvent(
sequence = laneId,
kind = SubagentActivityEventKind.UPDATE,
text = if (phase == SubagentActivityPhase.INTERRUPTED) {
"Stopped safely"
} else {
"Mapping Gateway events"
},
phase = phase,
observedAtMillis = 1,
),
),
)
}
@@ -1,86 +0,0 @@
package com.hermesandroid.relay.ui.components
import com.hermesandroid.relay.data.ChatMessage
import com.hermesandroid.relay.data.MessageRole
import com.hermesandroid.relay.viewmodel.SubagentActivity
import com.hermesandroid.relay.viewmodel.SubagentActivityEvent
import com.hermesandroid.relay.viewmodel.SubagentActivityEventKind
import com.hermesandroid.relay.viewmodel.SubagentActivityPhase
import com.hermesandroid.relay.viewmodel.SubagentChildPreview
import org.junit.Assert.assertEquals
import org.junit.Assert.assertTrue
import org.junit.Test
class SubagentActivityPreviewTest {
@Test
fun `supervised visibility excludes child history rows`() {
val activity = activity(laneId = 0, taskIndex = 0)
val preview = preview(activity, messageCount = 2)
val expanded = setOf(activity.stableKey)
val full = subagentActivityItemCount(
listOf(activity),
expanded,
SubagentPreviewVisibility(showChildHistory = true),
preview,
)
val supervised = subagentActivityItemCount(
listOf(activity),
expanded,
SubagentPreviewVisibility(showChildHistory = false),
preview,
)
assertTrue(full > supervised)
assertEquals(4, full - supervised) // heading, two child messages, and tail anchor
}
@Test
fun `follow target stays with selected child instead of final concurrent lane`() {
val first = activity(laneId = 0, taskIndex = 0)
val second = activity(laneId = 1, taskIndex = 1)
val preview = preview(first, messageCount = 1)
val target = subagentActivityFollowTarget(
activities = listOf(first, second),
expandedKeys = setOf(first.stableKey, second.stableKey),
visibility = SubagentPreviewVisibility(),
childPreview = preview,
)
val total = subagentActivityItemCount(
listOf(first, second),
setOf(first.stableKey, second.stableKey),
SubagentPreviewVisibility(),
preview,
)
assertTrue(target < total - 1)
}
private fun activity(laneId: Long, taskIndex: Int) = SubagentActivity(
laneId = laneId,
turnId = "turn",
taskIndex = taskIndex,
taskCount = 2,
goal = "Task $taskIndex",
phase = SubagentActivityPhase.PROGRESS,
events = listOf(
SubagentActivityEvent(
sequence = 0,
kind = SubagentActivityEventKind.UPDATE,
text = "Working",
phase = SubagentActivityPhase.PROGRESS,
observedAtMillis = 1,
),
),
)
private fun preview(activity: SubagentActivity, messageCount: Int) = SubagentChildPreview(
activityKey = activity.stableKey,
parentSessionId = "parent",
parentScopeKey = "scope",
childWatchAvailable = true,
messages = List(messageCount) { index ->
ChatMessage("message-$index", MessageRole.ASSISTANT, "Text", index.toLong())
},
)
}
@@ -992,11 +992,12 @@ 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(
@@ -1297,120 +1298,6 @@ class ChatViewModelGatewayInboundTurnTest {
)
}
@Test
fun recoveredServerDefaultSubagentWatchOmitsProfileOverride() {
assertRecoveredSubagentWatchProfile(
contextKey = AgentDisplay.profileContextKey("connection-a", null),
persistedProfileKey = AgentDisplay.SERVER_DEFAULT_PROFILE_KEY,
expectedProfile = null,
)
}
@Test
fun recoveredLiteralDefaultSubagentWatchKeepsExplicitProfile() {
assertRecoveredSubagentWatchProfile(
contextKey = AgentDisplay.profileContextKey("connection-a", "default"),
persistedProfileKey = "default",
expectedProfile = "default",
)
}
@Test
fun recoveredNamedSubagentWatchKeepsOwningProfileAcrossConnectionScope() {
assertRecoveredSubagentWatchProfile(
contextKey = AgentDisplay.profileContextKey("connection-b", "team::writer"),
persistedProfileKey = "team::writer",
expectedProfile = "team::writer",
)
}
@Test
fun recoveredLegacyCheckpointFailsClosedWithoutInventingProfile() {
assertRecoveredSubagentWatchProfile(
contextKey = "connection-a/profile-default",
persistedProfileKey = null,
expectedProfile = null,
)
}
@Test
fun currentServerDefaultCheckpointPersistsExplicitSentinel() {
assertCurrentCheckpointProfileKey(
profileName = null,
effectiveSessionProfileName = "victor",
expectedProfileKey = AgentDisplay.SERVER_DEFAULT_PROFILE_KEY,
)
}
@Test
fun currentLiteralDefaultCheckpointPersistsNamedProfile() {
assertCurrentCheckpointProfileKey(
profileName = "default",
effectiveSessionProfileName = "default",
expectedProfileKey = "default",
)
}
@Test
fun explicitConversationOwnerWinsAmbientSelectorInCurrentCheckpoint() {
val checkpointStore = MemoryCheckpointStore()
val global = Profile(name = "global", model = "global-model")
val writer = Profile(name = "writer", model = "writer-model")
viewModel.setSelectedProfileProvider { global }
viewModel.setSessionProfileNameProvider { global.name }
viewModel.setProfileMessageLoader { Result.success(emptyList()) }
viewModel.setChatTurnCheckpointStore(checkpointStore)
assertTrue(
viewModel.openProfileSession(
profileName = writer.name,
profile = writer,
contextKey = AgentDisplay.profileContextKey("connection-a", writer.name),
sessionId = STORED_SESSION_ID,
),
)
awaitCondition { viewModel.conversationBinding.value.hasExplicitOwner }
viewModel.sendMessage("Persist explicit owner")
gatewayHarness.awaitRpc("prompt.submit")
awaitCondition { checkpointStore.checkpoint?.profileKey == writer.name }
assertEquals(writer.name, checkpointStore.checkpoint?.profileKey)
}
@Test
fun liveNonRecoveredSubagentWatchKeepsLiteralDefaultProfile() {
val profile = Profile(name = "default", model = "model")
viewModel.setSelectedProfileProvider { profile }
viewModel.setSessionProfileNameProvider { profile.name }
viewModel.switchProfileContext(
AgentDisplay.profileContextKey("connection-a", profile.name),
STORED_SESSION_ID,
)
gatewayHarness.resumeLiveSessionIds["live-child-stored"] = "live-child"
viewModel.sendMessage("Delegate live work")
gatewayHarness.awaitRpc("prompt.submit")
serverWs.send(
gatewayHarness.eventFrame(
"subagent.start",
buildJsonObject {
put("goal", "Inspect live path")
put("task_index", 0)
put("task_count", 1)
put("subagent_id", "live-child-agent")
put("child_session_id", "live-child-stored")
},
"live-resumed",
),
)
awaitCondition { viewModel.subagentActivities.value.size == 1 }
viewModel.openSubagentChildPreview(viewModel.subagentActivities.value.single().stableKey)
val resume = gatewayHarness.awaitRpcCount("session.resume", 2).last()
assertEquals(JsonPrimitive("default"), resume["profile"])
}
@Test
fun explicitApprovalActionAloneEmitsResponseAndCollapsesCard() {
viewModel.sendMessage("Run the guarded command")
@@ -1938,11 +1825,12 @@ 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(
@@ -2145,13 +2033,16 @@ 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)
}
@@ -2162,16 +2053,19 @@ 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 trigger an exact-session reattach.
// callback must itself restore the observation socket and catch up
// history without attaching the live session.
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)
}
@@ -2480,36 +2374,189 @@ class ChatViewModelGatewayInboundTurnTest {
}
@Test
fun reconnectAfterMissedStartRecoversOnExactSessionCompletion() {
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)
}
}
serverWs.close(1012, "missed start")
awaitCondition { gatewayClient.connectionState.value == GatewayConnectionState.Idle }
viewModel.setChatVisible(true)
viewModel.prewarmGateway()
serverWs = gatewayHarness.awaitServerSocket()
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",
),
)
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.
persistedHistory = persistedAnswerHistory()
serverWs.send(
gatewayHarness.eventFrame(
"message.complete",
buildJsonObject { put("text", BACKGROUND_ANSWER) },
"live-resumed",
),
)
releaseFirstRead.complete(Unit)
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)
@@ -2579,88 +2626,6 @@ class ChatViewModelGatewayInboundTurnTest {
)
}
private fun assertRecoveredSubagentWatchProfile(
contextKey: String,
persistedProfileKey: String?,
expectedProfile: String?,
) {
val now = System.currentTimeMillis()
val checkpointStore = MemoryCheckpointStore(
ChatTurnCheckpoint(
contextKey = contextKey,
profileKey = persistedProfileKey,
sessionId = STORED_SESSION_ID,
liveSessionId = "live-resumed",
transport = "gateway",
user = ChatTurnUserCheckpoint("pending-user", "Delegate this", now - 2_000L),
assistant = ChatTurnAssistantCheckpoint(
id = "pending-assistant",
content = "Partial",
timestamp = now - 1_900L,
),
priorUserMessageCount = 0,
baselineAssistantCount = 0,
startedAt = now - 1_900L,
updatedAt = now,
),
)
gatewayHarness.recoveryRunning = true
gatewayHarness.recoveryAssistant = "Partial"
gatewayHarness.resumeLiveSessionIds["child-stored"] = "child-live"
viewModel.setChatTurnCheckpointStore(checkpointStore)
handler.setSessionId(null)
viewModel.switchProfileContext(contextKey, STORED_SESSION_ID)
viewModel.prewarmGateway()
gatewayHarness.awaitRpc("session.activate")
awaitCondition { handler.isStreaming.value }
serverWs.send(
gatewayHarness.eventFrame(
"subagent.start",
buildJsonObject {
put("goal", "Inspect recovery")
put("task_index", 0)
put("task_count", 1)
put("subagent_id", "child-agent")
put("child_session_id", "child-stored")
},
"live-resumed",
),
)
awaitCondition { viewModel.subagentActivities.value.size == 1 }
viewModel.openSubagentChildPreview(viewModel.subagentActivities.value.single().stableKey)
val resume = gatewayHarness.awaitRpcCount("session.resume", 2).last()
if (expectedProfile == null) {
assertFalse(resume.containsKey("profile"))
} else {
assertEquals(JsonPrimitive(expectedProfile), resume["profile"])
}
}
private fun assertCurrentCheckpointProfileKey(
profileName: String?,
effectiveSessionProfileName: String?,
expectedProfileKey: String,
) {
val checkpointStore = MemoryCheckpointStore()
val profile = profileName?.let { Profile(name = it, model = "model") }
viewModel.setSelectedProfileProvider { profile }
viewModel.setSessionProfileNameProvider { effectiveSessionProfileName }
viewModel.setChatTurnCheckpointStore(checkpointStore)
viewModel.switchProfileContext(
AgentDisplay.profileContextKey("connection-a", profileName),
STORED_SESSION_ID,
)
viewModel.sendMessage("Persist profile identity")
gatewayHarness.awaitRpc("prompt.submit")
awaitCondition { checkpointStore.checkpoint?.profileKey == expectedProfileKey }
assertEquals(expectedProfileKey, checkpointStore.checkpoint?.profileKey)
}
private fun activeSessionPayload(status: String) = buildJsonObject {
put("sessions", buildJsonArray {
add(buildJsonObject {
@@ -1,159 +0,0 @@
package com.hermesandroid.relay.viewmodel
import com.hermesandroid.relay.network.upstream.GatewaySubagentEvent
import org.junit.Assert.assertEquals
import org.junit.Assert.assertFalse
import org.junit.Assert.assertTrue
import org.junit.Test
class SubagentActivityControllerTest {
private var now = 1_000L
private val controller = SubagentActivityController { now++ }
@Test
fun `interleaved children retain independent lifecycle previews`() {
controller.selectSession("parent", "connection::default")
controller.beginTurn("parent", "connection::default", "turn-1")
controller.onEvent("parent", "connection::default", "turn-1", event(0, GatewaySubagentEvent.Phase.START, goal = "Research"))
controller.onEvent("parent", "connection::default", "turn-1", event(1, GatewaySubagentEvent.Phase.START, goal = "Review"))
controller.onEvent("parent", "connection::default", "turn-1", event(0, GatewaySubagentEvent.Phase.PROGRESS, preview = "Halfway"))
controller.onEvent("parent", "connection::default", "turn-1", event(1, GatewaySubagentEvent.Phase.TOOL, preview = "file.kt", tool = "read_file"))
controller.onEvent("parent", "connection::default", "turn-1", event(0, GatewaySubagentEvent.Phase.COMPLETE, status = "complete", summary = "Done"))
val activities = controller.activities.value.sortedBy { it.taskIndex }
assertEquals(listOf("Research", "Review"), activities.map { it.goal })
assertEquals(SubagentActivityPhase.COMPLETED, activities[0].phase)
assertEquals("Done", activities[0].summary)
assertEquals(SubagentActivityPhase.TOOL, activities[1].phase)
assertEquals("read_file", activities[1].events.last().toolName)
}
@Test
fun `profile session and newer turn fence stale events`() {
controller.selectSession("shared", "connection::alpha")
controller.beginTurn("shared", "connection::alpha", "turn-old")
controller.onEvent("shared", "connection::alpha", "turn-old", event(0, GatewaySubagentEvent.Phase.START, goal = "Old"))
controller.selectSession("shared", "connection::beta")
controller.onEvent("shared", "connection::alpha", "turn-old", event(0, GatewaySubagentEvent.Phase.PROGRESS, preview = "stale"))
assertTrue(controller.activities.value.isEmpty())
controller.beginTurn("shared", "connection::beta", "turn-new")
controller.onEvent("shared", "connection::beta", "turn-new", event(0, GatewaySubagentEvent.Phase.START, goal = "New"))
controller.beginTurn("shared", "connection::beta", "turn-newer")
controller.onEvent("shared", "connection::beta", "turn-newer", event(0, GatewaySubagentEvent.Phase.START, goal = "Newest"))
controller.onEvent("shared", "connection::beta", "turn-new", event(0, GatewaySubagentEvent.Phase.PROGRESS, preview = "late"))
assertEquals(listOf("Newest"), controller.activities.value.map { it.goal })
}
@Test
fun `terminal truth distinguishes failure interruption and missing terminal`() {
controller.selectSession("parent", "scope")
controller.beginTurn("parent", "scope", "turn-1")
controller.onEvent("parent", "scope", "turn-1", event(0, GatewaySubagentEvent.Phase.START))
controller.onEvent("parent", "scope", "turn-1", event(0, GatewaySubagentEvent.Phase.COMPLETE, status = "interrupted"))
assertEquals(SubagentActivityPhase.INTERRUPTED, controller.activities.value.single().phase)
controller.beginTurn("parent", "scope", "turn-2")
controller.onEvent("parent", "scope", "turn-2", event(0, GatewaySubagentEvent.Phase.START))
controller.onEvent("parent", "scope", "turn-2", event(0, GatewaySubagentEvent.Phase.COMPLETE, status = "failed"))
assertEquals(SubagentActivityPhase.FAILED, controller.activities.value.single().phase)
controller.beginTurn("parent", "scope", "turn-3")
controller.onEvent("parent", "scope", "turn-3", event(0, GatewaySubagentEvent.Phase.START))
controller.endTurn("turn-3")
assertEquals(SubagentActivityPhase.ENDED_WITH_PARENT, controller.activities.value.single().phase)
assertTrue(controller.activities.value.single().partialAfterGap)
}
@Test
fun `reconnect marks only live activity partial and late events do not reopen terminal child`() {
controller.selectSession("parent", "scope")
controller.beginTurn("parent", "scope", "turn")
controller.onConnectionReady(true)
controller.onEvent("parent", "scope", "turn", event(0, GatewaySubagentEvent.Phase.START))
controller.onConnectionReady(false)
controller.onConnectionReady(true)
assertTrue(controller.activities.value.single().partialAfterGap)
controller.onEvent("parent", "scope", "turn", event(0, GatewaySubagentEvent.Phase.COMPLETE, status = "complete"))
val revision = controller.activities.value.single().revision
controller.onEvent("parent", "scope", "turn", event(0, GatewaySubagentEvent.Phase.PROGRESS, preview = "late"))
assertEquals(revision, controller.activities.value.single().revision)
}
@Test
fun `event history is sanitized coalesced and bounded`() {
controller.selectSession("parent", "scope")
controller.beginTurn("parent", "scope", "turn")
controller.onEvent("parent", "scope", "turn", event(0, GatewaySubagentEvent.Phase.START, goal = "\u001B[31mSecret\u0000"))
repeat(SubagentActivityController.MAX_EVENTS_PER_CHILD + 10) { index ->
controller.onEvent(
"parent",
"scope",
"turn",
event(0, GatewaySubagentEvent.Phase.PROGRESS, preview = "update-$index"),
)
}
controller.onEvent("parent", "scope", "turn", event(0, GatewaySubagentEvent.Phase.PROGRESS, preview = "same"))
controller.onEvent("parent", "scope", "turn", event(0, GatewaySubagentEvent.Phase.PROGRESS, preview = "same"))
val activity = controller.activities.value.single()
assertEquals("Secret", activity.goal)
assertTrue(activity.truncated)
assertTrue(activity.events.size <= SubagentActivityController.MAX_EVENTS_PER_CHILD)
assertEquals(1, activity.events.count { it.text == "same" })
assertFalse(activity.events.any { it.text?.contains('\u0000') == true })
}
@Test
fun `identity enrichment keeps one lane while conflicting child stays separate`() {
controller.selectSession("parent", "scope")
controller.beginTurn("parent", "scope", "turn")
controller.onEvent(
"parent", "scope", "turn",
event(0, GatewaySubagentEvent.Phase.SPAWN_REQUESTED, childId = "session-1"),
)
controller.onEvent(
"parent", "scope", "turn",
event(0, GatewaySubagentEvent.Phase.START, subagentId = "agent-1"),
)
assertEquals(1, controller.activities.value.size)
assertEquals("session-1", controller.activities.value.single().childSessionId)
assertEquals("agent-1", controller.activities.value.single().subagentId)
controller.onEvent(
"parent", "scope", "turn",
event(
0,
GatewaySubagentEvent.Phase.START,
subagentId = "agent-2",
childId = "session-2",
),
)
assertEquals(2, controller.activities.value.size)
assertEquals(2, controller.activities.value.map { it.stableKey }.distinct().size)
}
private fun event(
index: Int,
phase: GatewaySubagentEvent.Phase,
goal: String = "",
preview: String? = null,
tool: String? = null,
status: String? = null,
summary: String? = null,
subagentId: String? = null,
childId: String? = null,
) = GatewaySubagentEvent(
phase = phase,
taskIndex = index,
taskCount = 2,
goal = goal,
preview = preview,
toolName = tool,
status = status,
summary = summary,
subagentId = subagentId,
childSessionId = childId,
)
}
@@ -1,152 +0,0 @@
package com.hermesandroid.relay.viewmodel
import com.hermesandroid.relay.network.upstream.GatewayChatClient
import com.hermesandroid.relay.network.upstream.GatewayChildWatch
import com.hermesandroid.relay.network.upstream.models.MessageItem
import io.mockk.mockk
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.test.advanceUntilIdle
import kotlinx.coroutines.test.runCurrent
import kotlinx.coroutines.test.runTest
import kotlinx.serialization.json.JsonPrimitive
import org.junit.Assert.assertEquals
import org.junit.Assert.assertFalse
import org.junit.Assert.assertNull
import org.junit.Assert.assertTrue
import org.junit.Test
@OptIn(ExperimentalCoroutinesApi::class)
class SubagentChildPreviewControllerTest {
@Test
fun `pre-ack live callbacks append after hydrated history`() = runTest {
val client = mockk<GatewayChatClient>()
val controller = SubagentChildPreviewController(
scope = this,
openWatch = { _, _, _, callbacks ->
callbacks.onStart()
callbacks.onTextDelta("live")
callbacks.onComplete()
Result.success(watch(messages = listOf(message("history")), running = true))
},
)
controller.open(activity(), client, "parent", "scope", true) { true }
advanceUntilIdle()
val state = controller.state.value!!
assertTrue(state.messages.any { it.content == "history" })
assertTrue(state.messages.any { it.content == "live" })
assertFalse(state.running)
assertFalse(state.status == "completed")
}
@Test
fun `dismissed in-flight open closes the late exact watch without publishing`() = runTest {
val client = mockk<GatewayChatClient>()
val acknowledgement = CompletableDeferred<GatewayChildWatch>()
var closed: GatewayChildWatch? = null
val controller = SubagentChildPreviewController(
scope = this,
openWatch = { _, _, _, _ -> Result.success(acknowledgement.await()) },
closeWatch = { _, watch ->
closed = watch
Result.success(Unit)
},
)
controller.open(activity(), client, "parent", "scope", true) { true }
runCurrent()
controller.close()
val late = watch()
acknowledgement.complete(late)
advanceUntilIdle()
assertEquals(late, closed)
assertNull(controller.state.value)
}
@Test
fun `stale open generation closes late watch and cannot replace newer profile`() = runTest {
val client = mockk<GatewayChatClient>()
val oldAcknowledgement = CompletableDeferred<GatewayChildWatch>()
val closed = mutableListOf<GatewayChildWatch>()
val controller = SubagentChildPreviewController(
scope = this,
openWatch = { _, sessionId, profile, _ ->
if (sessionId == "child-old") {
Result.success(oldAcknowledgement.await())
} else {
Result.success(watch(sessionId, "live-new", profile))
}
},
closeWatch = { _, watch ->
closed += watch
Result.success(Unit)
},
)
controller.open(activity("child-old", "old", laneId = 0), client, "parent", "scope", true) { true }
runCurrent()
controller.open(activity("child-new", "default", laneId = 1), client, "parent", "scope", true) { true }
advanceUntilIdle()
assertEquals("default", controller.state.value?.status)
val late = watch("child-old", "live-old", "old")
oldAcknowledgement.complete(late)
advanceUntilIdle()
assertTrue(late in closed)
assertEquals("default", controller.state.value?.status)
assertEquals("turn:1", controller.state.value?.activityKey)
}
private fun activity(
childSessionId: String = "child-stored",
profile: String = "default",
laneId: Long = 0,
) = SubagentActivity(
laneId = laneId,
turnId = "turn",
taskIndex = 0,
taskCount = 1,
goal = "Inspect",
phase = SubagentActivityPhase.PROGRESS,
childSessionId = childSessionId,
profile = profile,
)
private fun message(text: String) = MessageItem(
role = "assistant",
content = JsonPrimitive(text),
)
private fun watch(
messages: List<MessageItem> = emptyList(),
running: Boolean = false,
) = GatewayChildWatch(
storedSessionId = "child-stored",
liveSessionId = "child-live",
profile = "default",
generation = 1,
messages = messages,
historyTruncated = false,
running = running,
status = if (running) "streaming" else "idle",
)
private fun watch(
storedSessionId: String,
liveSessionId: String,
status: String?,
) = GatewayChildWatch(
storedSessionId = storedSessionId,
liveSessionId = liveSessionId,
profile = status,
generation = if (storedSessionId == "child-old") 1 else 2,
messages = emptyList(),
historyTruncated = false,
running = true,
status = status,
)
}
+33
View File
@@ -3977,3 +3977,36 @@ 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,6 +41,7 @@ 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": "e98464fc5fd7283a6cfd0053155c20ca2306f360f7f3dac52c0e56343850e61a",
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"sideload": "4abff4f1069091ec2de735c3037a7ec7d77699cb4321e8511a622437bceaf7c2"
},
"surfaces": {
@@ -48,7 +48,7 @@
"verification": "ai-translated",
"review_refs": [],
"source_sha256": {
"main": "e98464fc5fd7283a6cfd0053155c20ca2306f360f7f3dac52c0e56343850e61a",
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"sideload": "4abff4f1069091ec2de735c3037a7ec7d77699cb4321e8511a622437bceaf7c2"
},
"surfaces": {
@@ -72,7 +72,7 @@
"verification": "ai-translated",
"review_refs": [],
"source_sha256": {
"main": "e98464fc5fd7283a6cfd0053155c20ca2306f360f7f3dac52c0e56343850e61a",
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"sideload": "4abff4f1069091ec2de735c3037a7ec7d77699cb4321e8511a622437bceaf7c2"
},
"surfaces": {
@@ -96,7 +96,7 @@
"verification": "ai-translated",
"review_refs": [],
"source_sha256": {
"main": "e98464fc5fd7283a6cfd0053155c20ca2306f360f7f3dac52c0e56343850e61a",
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"sideload": "4abff4f1069091ec2de735c3037a7ec7d77699cb4321e8511a622437bceaf7c2"
},
"surfaces": {
@@ -120,7 +120,7 @@
"verification": "ai-translated",
"review_refs": [],
"source_sha256": {
"main": "e98464fc5fd7283a6cfd0053155c20ca2306f360f7f3dac52c0e56343850e61a",
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"sideload": "4abff4f1069091ec2de735c3037a7ec7d77699cb4321e8511a622437bceaf7c2"
},
"surfaces": {
@@ -135,7 +135,7 @@
"verification": "ai-translated",
"review_refs": [],
"source_sha256": {
"main": "e98464fc5fd7283a6cfd0053155c20ca2306f360f7f3dac52c0e56343850e61a",
"main": "28d58a3b9803968124ea6581fd9dbb3a0ce2a9ee1c946228bb6f0b79f987a24a",
"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 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.
- **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.
- **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.
@@ -30,7 +30,6 @@ GATEWAY_SETTLED_INFO = "gateway.settled_session_info"
SESSION_ACTIVATE = "gateway.session_activate_live"
SESSION_RESUME = "gateway.session_resume_durable"
SESSION_ACTIVE_LIST = "gateway.session_active_list"
SUBAGENT_CHILD_WATCH = "gateway.subagent_child_watch"
API_BOUNDARY = "api.fallback_boundary"
ALL_CONTRACTS = (
GATEWAY_TERMINAL,
@@ -38,7 +37,6 @@ ALL_CONTRACTS = (
SESSION_ACTIVATE,
SESSION_RESUME,
SESSION_ACTIVE_LIST,
SUBAGENT_CHILD_WATCH,
API_BOUNDARY,
)
@@ -377,38 +375,6 @@ def _check_api_boundary(api: SourceFile) -> CheckResult:
return CheckResult(contract, False, (), str(exc))
def _check_subagent_child_watch(server: SourceFile, methods: SourceFile) -> CheckResult:
contract = SUBAGENT_CHILD_WATCH
try:
resume = methods.method_handler("session.resume")
resume_segment = methods.segment(resume)
resume_strings = _string_constants(resume)
server_strings = _string_constants(server.tree)
missing_resume = sorted(
{"lazy", "close_on_disconnect"} - resume_strings
)
if missing_resume:
raise ValueError("lazy child resume field(s) missing: " + ", ".join(missing_resume))
if "include_ancestors" not in resume_segment:
raise ValueError("lazy child resume does not declare child-only history")
missing_events = sorted(
{"child_session_id", "subagent.text", "reasoning.delta", "message.delta"}
- server_strings
)
if missing_events:
raise ValueError("child watch mirror event(s) missing: " + ", ".join(missing_events))
return CheckResult(
contract,
True,
(
methods.evidence(resume, "session.resume supports a lazy child-only watch"),
"tui_gateway/server.py: child_session_id routes child mirror events",
),
)
except ValueError as exc:
return CheckResult(contract, False, (), str(exc))
def load_requirements(manifest: Path | None) -> tuple[str, ...]:
if manifest is None:
return ALL_CONTRACTS
@@ -450,7 +416,6 @@ def audit_sources(root: Path, requirements: Iterable[str]) -> list[CheckResult]:
SESSION_ACTIVATE: lambda: _check_activate(server, methods),
SESSION_RESUME: lambda: _check_resume(methods),
SESSION_ACTIVE_LIST: lambda: _check_active_list(server, methods),
SUBAGENT_CHILD_WATCH: lambda: _check_subagent_child_watch(server, methods),
API_BOUNDARY: lambda: _check_api_boundary(api),
}
return [checks[requirement]() for requirement in requirements]
@@ -56,12 +56,6 @@ def _session_live_item(sid, session, current_sid=""):
"session_key": session.get("session_key", sid),
"status": _session_live_status(sid, session),
}
def _mirror_subagent_child(event):
child = event.get("child_session_id")
if event.get("type") == "subagent.text":
return (child, "reasoning.delta", "message.delta")
return child
'''
METHODS_SOURCE = '''
@@ -71,9 +65,6 @@ def method(name):
@method("session.resume")
def _(rid, params):
target = params.get("session_id", "")
lazy = bool(params.get("lazy"))
close_on_disconnect = bool(params.get("close_on_disconnect"))
include_ancestors = not lazy
found = db.get_session(target)
if not found:
return _err(rid, 4007, "session not found")
-5
View File
@@ -88,11 +88,6 @@ The initial catalog covers ordinary streaming, rapid chunks/reasoning/tool
events, queued turns, scoped and foreign/unscoped inputs, persisted history,
and both issue #365 terminal-gap forms:
- `subagent_child_preview`: interleaved concurrent child lifecycle events carry
stable child/session identity, thinking/progress/tool previews, and distinct
completed/interrupted terminal states. Its upstream requirement also proves
the vanilla lazy child-session watch contract used by read-only clients.
- `active_status_lifecycle`: one successful live snapshot contains starting,
working, waiting, and idle rows; the next successful snapshot is empty so a
client can prove a complete, unambiguously resolved snapshot clears prior
@@ -100,6 +100,47 @@ 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)
@@ -265,9 +306,9 @@ class ScenarioTestCase(unittest.TestCase):
"active_status_lifecycle",
"active_status_profile_scope",
"active_status_unsupported",
"cross_client_observation",
"ordinary_turn",
"rapid_tools_interims",
"subagent_child_preview",
"terminal_gap_activate",
"terminal_gap_session_info",
"queued_follow_up",
@@ -305,6 +346,10 @@ 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()):
@@ -0,0 +1,40 @@
{
"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}
]
]
}
}
@@ -1,28 +0,0 @@
{
"name": "subagent_child_preview",
"live_session_id": "fixture-live-1",
"stored_session_id": "20260821_120000_fixture",
"profile": "default",
"contract_requirements": [
"gateway.message_complete",
"gateway.subagent_child_watch"
],
"turns": [
{
"steps": [
{"op": "event", "type": "message.start"},
{"op": "event", "type": "subagent.spawn_requested", "payload": {"goal": "Inspect Android", "task_index": 0, "task_count": 2, "subagent_id": "child-a", "child_session_id": "child-session-a", "depth": 0}},
{"op": "event", "type": "subagent.start", "payload": {"goal": "Inspect Android", "task_index": 0, "task_count": 2, "subagent_id": "child-a", "child_session_id": "child-session-a", "depth": 0}},
{"op": "event", "type": "subagent.start", "payload": {"goal": "Review privacy", "task_index": 1, "task_count": 2, "subagent_id": "child-b", "child_session_id": "child-session-b", "depth": 0}},
{"op": "event", "type": "subagent.thinking", "payload": {"task_index": 0, "task_count": 2, "subagent_id": "child-a", "child_session_id": "child-session-a", "text": "Mapping events"}},
{"op": "event", "type": "subagent.tool", "payload": {"task_index": 1, "task_count": 2, "subagent_id": "child-b", "child_session_id": "child-session-b", "tool_name": "read_file", "tool_preview": "policy.md"}},
{"op": "event", "type": "subagent.progress", "payload": {"task_index": 0, "task_count": 2, "subagent_id": "child-a", "child_session_id": "child-session-a", "text": "One tool complete"}},
{"op": "event", "type": "subagent.complete", "payload": {"task_index": 1, "task_count": 2, "subagent_id": "child-b", "child_session_id": "child-session-b", "status": "interrupted", "summary": "Stopped safely"}},
{"op": "event", "type": "subagent.complete", "payload": {"task_index": 0, "task_count": 2, "subagent_id": "child-a", "child_session_id": "child-session-a", "status": "completed", "summary": "Mapped Android events", "duration_seconds": 2.5}},
{"op": "persist", "messages": [{"id": 1, "role": "user", "content": "Exercise child previews.", "timestamp": 1.0}, {"id": 2, "role": "assistant", "content": "Delegation complete.", "timestamp": 2.0}]},
{"op": "set_running", "value": false},
{"op": "event", "type": "message.complete", "payload": {"text": "Delegation complete.", "status": "complete"}}
]
}
]
}