Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
40fcb80a7d | ||
|
|
ff64548085 | ||
|
|
042f02b852 |
@@ -19,6 +19,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/), and this
|
||||
- **Android keeps completed chat text visible when Dashboard sign-in expires.** Generic and reason-coded history `401` responses settle the local turn, preserve its transcript, and surface the existing sign-in recovery without reading another profile's API history.
|
||||
- **Android keeps long-running context compaction alive.** A client-visible compaction status extends and refreshes the Gateway turn watchdog instead of interrupting healthy compression after the ordinary idle window. (Supersedes #484.)
|
||||
- **Android Bot Chats render loaded history immediately.** Route-owned chat screens observe their own handler state from first composition, including fast history loads that settle before another frame. (Supersedes #453.)
|
||||
- **Android Chat settles an owned Gateway turn when its terminal frame is lost.** An exact idle `session.active_list` snapshot now completes the matching local stream, reconciles durable history, and drains its queued follow-up without interrupting or claiming Desktop/TUI work.
|
||||
- **Supervised Gateway setup stays parent-owned.** Add Gateway is single-flight and checks live parent authority before allocating a draft, relock/back cancels the exact pending setup, and the locked Chat footer no longer attempts protected navigation.
|
||||
- **Generated images stay visible and use their intended Chat animation.** Completed image media survives a marker-lagging history refresh, and both the built-in `image_generate` tool and profile tools ending in `_create_image` use the image-generation presentation.
|
||||
|
||||
|
||||
@@ -53,8 +53,11 @@ and older Gateways without `session.active_list`. Before calling the status
|
||||
model device-certified:
|
||||
|
||||
- Exercise working, quiet tool-heavy work, each pending-input surface, normal
|
||||
completion, Stop, reconnect, app restart, and process recreation against
|
||||
current vanilla upstream.
|
||||
completion, a lost terminal followed by an exact active-list Idle row, Stop,
|
||||
reconnect, app restart, and process recreation against current vanilla
|
||||
upstream. Confirm the lost-terminal path preserves the partial transcript,
|
||||
settles composer/steering state, and drains or cancels queued corrections
|
||||
exactly once according to the owning turn outcome.
|
||||
- Verify All Profiles with duplicate session ids across two profiles and two
|
||||
saved connections; no late snapshot or old socket generation may mark the
|
||||
wrong row live.
|
||||
|
||||
+40
@@ -258,6 +258,46 @@ class GatewayForegroundRecoveryInstrumentedTest {
|
||||
assertEquals(0, fixture.requestsTo("/v1/chat/completions"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun terminalGapActiveList_settlesExactOwnedTurnAndRendersAuthoritativeHistory() {
|
||||
viewModel.sendMessage("Run an Android-owned task")
|
||||
fixture.awaitRpc("prompt.submit")
|
||||
serverSocket.send(fixture.event("message.start", null, LIVE_SESSION_ID))
|
||||
serverSocket.send(
|
||||
fixture.event(
|
||||
"message.delta",
|
||||
buildJsonObject { put("text", PARTIAL_ANSWER) },
|
||||
LIVE_SESSION_ID,
|
||||
),
|
||||
)
|
||||
compose.waitUntil(5_000) { handler.isStreaming.value }
|
||||
compose.onNodeWithTag("stream-state").assertTextEquals("STREAMING")
|
||||
|
||||
persistedHistory = listOf(
|
||||
MessageItem(
|
||||
id = PERSISTED_ANSWER_ID,
|
||||
sessionId = STORED_SESSION_ID,
|
||||
role = "assistant",
|
||||
content = JsonPrimitive(AUTHORITATIVE_ANSWER),
|
||||
),
|
||||
)
|
||||
fixture.activeSessionStatus = "idle"
|
||||
runBlocking { gatewayClient.listActiveSessions() }
|
||||
|
||||
compose.waitUntil(5_000) {
|
||||
!handler.isStreaming.value &&
|
||||
!gatewayClient.hasActiveTurn() &&
|
||||
handler.messages.value.singleOrNull()?.id == PERSISTED_ANSWER_ID
|
||||
}
|
||||
compose.onNodeWithTag("stream-state").assertTextEquals("IDLE")
|
||||
compose.onNodeWithTag("message-$PERSISTED_ANSWER_ID")
|
||||
.assertTextEquals("${MessageRole.ASSISTANT.name}:$AUTHORITATIVE_ANSWER")
|
||||
assertEquals(1, fixture.rpcCount("prompt.submit"))
|
||||
assertEquals(0, fixture.rpcCount("session.interrupt"))
|
||||
assertEquals(0, fixture.rpcCount("session.activate"))
|
||||
assertEquals(0, fixture.requestsTo("/v1/chat/completions"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun desktopOwnedTurn_remainsReadOnlyAcrossAndroidForegroundLifecycle() {
|
||||
viewModel.setChatVisible(false)
|
||||
|
||||
@@ -563,6 +563,28 @@ class GatewayChatClient(
|
||||
val terminalRequired: Boolean,
|
||||
)
|
||||
|
||||
/**
|
||||
* One exact turn settled from authoritative session state may still receive
|
||||
* the terminal frame that was already in flight. Consume only that terminal
|
||||
* so it cannot be reported as a second unmatched completion. A subsequent
|
||||
* message.start clears the drain because it establishes the next turn on
|
||||
* the same live runtime.
|
||||
*/
|
||||
@Volatile
|
||||
private var settledTurnDrain: SettledTurnDrain? = null
|
||||
|
||||
private data class SettledTurnDrain(
|
||||
val storedSessionId: String,
|
||||
val liveSessionId: String,
|
||||
)
|
||||
|
||||
private data class ActiveTurnLivenessProbe(
|
||||
val turn: GatewayTurn,
|
||||
val storedSessionId: String,
|
||||
val liveSessionId: String,
|
||||
val progressGeneration: Long,
|
||||
)
|
||||
|
||||
/**
|
||||
* Creates UI callbacks when the server starts a turn that has no matching
|
||||
* [sendTurn] call (for example a background-process completion). The
|
||||
@@ -692,6 +714,7 @@ class GatewayChatClient(
|
||||
): ActiveTurnHandle {
|
||||
val turn = GatewayTurn(
|
||||
callbacks = dispatchOn(callbacks),
|
||||
androidOwned = true,
|
||||
onTransportAccepted = onTransportAccepted,
|
||||
)
|
||||
// Warm = the connection-establish phases are skipped this turn (socket
|
||||
@@ -739,6 +762,10 @@ class GatewayChatClient(
|
||||
cleanupStagedAttachments(stagedImagePaths)
|
||||
return@launch
|
||||
}
|
||||
// A newly accepted Android send is a distinct generation on
|
||||
// this runtime. Its terminal must never be consumed by the
|
||||
// prior turn's optional late-terminal drain.
|
||||
settledTurnDrain = null
|
||||
activeTurn = turn
|
||||
turn.armWatchdog()
|
||||
// Generic `file.attach` uploads are staged artifacts, not
|
||||
@@ -1378,6 +1405,7 @@ class GatewayChatClient(
|
||||
callbacks = dispatchOn(callbacks),
|
||||
dedupeAdjacentMessageStarts = true,
|
||||
deferEvents = true,
|
||||
androidOwned = true,
|
||||
).also { turn ->
|
||||
turn.markRecoveredStarted()
|
||||
activeTurn = turn
|
||||
@@ -1501,6 +1529,7 @@ class GatewayChatClient(
|
||||
boundTurn = GatewayTurn(
|
||||
callbacks = dispatchOn(callbacks),
|
||||
dedupeAdjacentMessageStarts = true,
|
||||
androidOwned = true,
|
||||
).also { turn ->
|
||||
turn.markRecoveredStarted()
|
||||
activeTurn = turn
|
||||
@@ -1534,6 +1563,7 @@ class GatewayChatClient(
|
||||
val queuedTurn = GatewayTurn(
|
||||
callbacks = dispatchOn(registration.callbacks),
|
||||
dedupeAdjacentMessageStarts = true,
|
||||
androidOwned = true,
|
||||
)
|
||||
// recoverTurn is resumed on its caller's coroutine context;
|
||||
// ChatViewModel calls it from Main, so this admission runs
|
||||
@@ -2429,6 +2459,7 @@ class GatewayChatClient(
|
||||
} catch (error: Exception) {
|
||||
return GatewayActiveSessionsResult.TransientFailure(error)
|
||||
}
|
||||
val livenessProbe = captureActiveTurnLivenessProbe()
|
||||
val result = rpc(
|
||||
"session.active_list",
|
||||
buildJsonObject {
|
||||
@@ -2448,12 +2479,58 @@ class GatewayChatClient(
|
||||
val payload = result.getOrThrow()
|
||||
val rows = payload["sessions"] as? JsonArray
|
||||
?: throw GatewayRpcException("session.active_list returned no sessions array")
|
||||
GatewayActiveSessionsResult.Success(rows.map(::parseGatewayActiveSession))
|
||||
val sessions = rows.map(::parseGatewayActiveSession)
|
||||
reconcileActiveTurnFromSnapshot(livenessProbe, sessions)
|
||||
GatewayActiveSessionsResult.Success(sessions)
|
||||
} catch (parseError: Exception) {
|
||||
GatewayActiveSessionsResult.TransientFailure(parseError)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Capture only a locally submitted/recovered turn. Merely observing an
|
||||
* exact session through the shared Gateway socket never grants Android
|
||||
* authority to settle Desktop/TUI work.
|
||||
*/
|
||||
private fun captureActiveTurnLivenessProbe(): ActiveTurnLivenessProbe? {
|
||||
val turn = activeTurn ?: return null
|
||||
val generation = turn.captureLivenessGeneration() ?: return null
|
||||
val storedId = storedSessionId ?: return null
|
||||
val liveId = liveSessionId ?: return null
|
||||
return ActiveTurnLivenessProbe(turn, storedId, liveId, generation)
|
||||
}
|
||||
|
||||
/**
|
||||
* `session.active_list` is process-wide, but a row naming both identifiers
|
||||
* already owned by this client is authoritative for that exact runtime.
|
||||
* Fence the delayed snapshot by turn identity and progress generation so an
|
||||
* old idle result cannot settle a newer turn or race newer live events.
|
||||
*/
|
||||
private fun reconcileActiveTurnFromSnapshot(
|
||||
probe: ActiveTurnLivenessProbe?,
|
||||
sessions: List<GatewayActiveSession>,
|
||||
) {
|
||||
probe ?: return
|
||||
if (activeTurn !== probe.turn ||
|
||||
storedSessionId != probe.storedSessionId ||
|
||||
liveSessionId != probe.liveSessionId
|
||||
) return
|
||||
val exact = sessions.singleOrNull { row ->
|
||||
row.runtimeSessionId == probe.liveSessionId &&
|
||||
row.storedSessionId == probe.storedSessionId
|
||||
} ?: return
|
||||
if (exact.status != GatewayActiveSessionStatus.Idle) return
|
||||
if (probe.turn.settleFromAuthoritativeSessionState(
|
||||
running = false,
|
||||
source = "session.active_list",
|
||||
expectedProgressGeneration = probe.progressGeneration,
|
||||
)
|
||||
) {
|
||||
if (activeTurn === probe.turn) activeTurn = null
|
||||
if (!AppForegroundTracker.isForeground.value) scheduleBackgroundClose()
|
||||
}
|
||||
}
|
||||
|
||||
/** Stop one process owned by the current live gateway session. */
|
||||
suspend fun killProcess(processId: String): Result<Unit> {
|
||||
if (processId.isBlank()) {
|
||||
@@ -2846,6 +2923,7 @@ class GatewayChatClient(
|
||||
activeTurn = null
|
||||
backgroundTurns.clear()
|
||||
cancelledTurnDrain = null
|
||||
settledTurnDrain = null
|
||||
unsolicitedTurnProvider = null
|
||||
coldPrewarmSessionReadyListener = null
|
||||
unmatchedTurnCompleteListener = null
|
||||
@@ -3282,6 +3360,7 @@ class GatewayChatClient(
|
||||
val requestedProfile = currentSessionProfile()
|
||||
if (requestedStoredId != null && requestedStoredId != storedSessionId) {
|
||||
cancelledTurnDrain = null
|
||||
settledTurnDrain = null
|
||||
}
|
||||
if (
|
||||
liveSessionId != null &&
|
||||
@@ -3351,6 +3430,7 @@ class GatewayChatClient(
|
||||
storedSessionId = stored
|
||||
liveSessionProfile = requestedProfile
|
||||
if (cancelledTurnDrain?.storedSessionId != stored) cancelledTurnDrain = null
|
||||
if (settledTurnDrain?.storedSessionId != stored) settledTurnDrain = null
|
||||
turn.callbacks.onSessionId(stored)
|
||||
}
|
||||
|
||||
@@ -3728,6 +3808,7 @@ class GatewayChatClient(
|
||||
}
|
||||
dispatchProcessEvent(type, payload, eventSessionId)
|
||||
if (consumeCancelledTurnEvent(type, eventSessionId)) return
|
||||
if (consumeSettledTurnTerminal(type, eventSessionId)) return
|
||||
var turn = activeTurn
|
||||
if (turn == null && type == "message.start") {
|
||||
// Unsolicited turns are accepted only with an explicit exact live-
|
||||
@@ -4254,6 +4335,7 @@ class GatewayChatClient(
|
||||
val callbacks: GatewayTurnCallbacks,
|
||||
dedupeAdjacentMessageStarts: Boolean = false,
|
||||
deferEvents: Boolean = false,
|
||||
private val androidOwned: Boolean = false,
|
||||
private val onTransportAccepted: () -> Unit = { },
|
||||
) : ActiveTurnHandle {
|
||||
private val mapper = GatewayEventMapper(callbacks, dedupeAdjacentMessageStarts)
|
||||
@@ -4280,6 +4362,7 @@ class GatewayChatClient(
|
||||
|
||||
private val rejoinAttempts = java.util.concurrent.atomic.AtomicInteger(0)
|
||||
private val transportAccepted = AtomicBoolean(false)
|
||||
private val progressGeneration = java.util.concurrent.atomic.AtomicLong(0L)
|
||||
|
||||
fun markTransportAccepted() {
|
||||
if (transportAccepted.compareAndSet(false, true)) {
|
||||
@@ -4356,6 +4439,7 @@ class GatewayChatClient(
|
||||
if (settledWithoutTerminalFrame) return
|
||||
if (type != "session.info") {
|
||||
started = true
|
||||
progressGeneration.incrementAndGet()
|
||||
markTransportAccepted()
|
||||
}
|
||||
tracer.mark("ttfe")
|
||||
@@ -4385,10 +4469,20 @@ class GatewayChatClient(
|
||||
* exact turn has proved it went live. A pre-start `running=false`
|
||||
* heartbeat can race `prompt.submit` and is not a completion boundary.
|
||||
*/
|
||||
fun settleFromAuthoritativeSessionState(running: Boolean?, source: String): Boolean {
|
||||
fun captureLivenessGeneration(): Long? =
|
||||
if (androidOwned && started && !ended) progressGeneration.get() else null
|
||||
|
||||
fun settleFromAuthoritativeSessionState(
|
||||
running: Boolean?,
|
||||
source: String,
|
||||
expectedProgressGeneration: Long? = null,
|
||||
): Boolean {
|
||||
if (running != false || !started) return false
|
||||
val settled = synchronized(deferredEventLock) {
|
||||
if (ended) {
|
||||
if (ended ||
|
||||
(expectedProgressGeneration != null &&
|
||||
progressGeneration.get() != expectedProgressGeneration)
|
||||
) {
|
||||
false
|
||||
} else {
|
||||
settledWithoutTerminalFrame = true
|
||||
@@ -4399,6 +4493,7 @@ class GatewayChatClient(
|
||||
if (!settled) return false
|
||||
|
||||
disarmWatchdog()
|
||||
armSettledTurnDrain()
|
||||
Log.i(TAG, "Gateway turn settled from $source after missing terminal frame")
|
||||
callbacks.onReconcileRequired()
|
||||
callbacks.onComplete()
|
||||
@@ -4419,6 +4514,7 @@ class GatewayChatClient(
|
||||
callbacks = dispatchOn(registration.callbacks),
|
||||
dedupeAdjacentMessageStarts = true,
|
||||
deferEvents = true,
|
||||
androidOwned = true,
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -4546,6 +4642,26 @@ class GatewayChatClient(
|
||||
)
|
||||
}
|
||||
|
||||
private fun armSettledTurnDrain() {
|
||||
val storedId = storedSessionId ?: return
|
||||
val liveId = liveSessionId ?: return
|
||||
settledTurnDrain = SettledTurnDrain(storedId, liveId)
|
||||
}
|
||||
|
||||
/** Consume one late terminal from a turn already settled by session state. */
|
||||
private fun consumeSettledTurnTerminal(type: String, eventSessionId: String?): Boolean {
|
||||
val drain = settledTurnDrain ?: return false
|
||||
if (eventSessionId != drain.liveSessionId) return false
|
||||
if (type == "message.start") {
|
||||
if (settledTurnDrain === drain) settledTurnDrain = null
|
||||
return false
|
||||
}
|
||||
if (type != "message.complete" && type != "error") return false
|
||||
if (settledTurnDrain === drain) settledTurnDrain = null
|
||||
Log.d(TAG, "Ignored late terminal for gateway turn settled from session state")
|
||||
return true
|
||||
}
|
||||
|
||||
private fun updateCancelledDrainLiveSession(storedId: String, liveId: String) {
|
||||
val drain = cancelledTurnDrain ?: return
|
||||
if (drain.storedSessionId == storedId) {
|
||||
|
||||
+191
-1
@@ -772,6 +772,7 @@ class GatewayChatClientTest {
|
||||
val moaReferences = ConcurrentLinkedQueue<GatewayMoaReference>()
|
||||
val usages = ConcurrentLinkedQueue<UsageInfo>()
|
||||
val reconcileRequests = AtomicInteger(0)
|
||||
val completions = AtomicInteger(0)
|
||||
val completeLatch = CountDownLatch(1)
|
||||
val preflightFailures = ConcurrentLinkedQueue<String>()
|
||||
|
||||
@@ -785,7 +786,7 @@ class GatewayChatClientTest {
|
||||
onToolCallFailed = { _, _ -> },
|
||||
onTurnComplete = { },
|
||||
onReconcileRequired = { reconcileRequests.incrementAndGet() },
|
||||
onComplete = { completeLatch.countDown() },
|
||||
onComplete = { completions.incrementAndGet(); completeLatch.countDown() },
|
||||
onUsage = { it?.let(usages::add) },
|
||||
onError = { errors += it; completeLatch.countDown() },
|
||||
onToolGenerating = { toolGenerating += it ?: "" },
|
||||
@@ -839,6 +840,21 @@ class GatewayChatClientTest {
|
||||
assertTrue("condition did not settle within ${timeoutMs}ms", condition())
|
||||
}
|
||||
|
||||
private fun exactActiveSessionPayload(
|
||||
status: String,
|
||||
liveSessionId: String = "live-1",
|
||||
storedSessionId: String = "20260612_120000_abc123",
|
||||
): JsonObject = buildJsonObject {
|
||||
put("sessions", buildJsonArray {
|
||||
add(buildJsonObject {
|
||||
put("id", liveSessionId)
|
||||
put("session_key", storedSessionId)
|
||||
put("status", status)
|
||||
put("last_active", 1_777_000_000.0)
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Swap in a client with shortened timeout seams. Mints a FRESH scope:
|
||||
* shutdown() cancels the scope's Job, and the replacement client must
|
||||
@@ -4714,6 +4730,180 @@ class GatewayChatClientTest {
|
||||
assertTrue(harness.ticketMints.get() >= 2)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `active session idle settles exact Android turn without interrupt`() = runBlocking {
|
||||
val recorder = Recorder()
|
||||
client.sendTurn(null, "finish without terminal", null, recorder.callbacks) {
|
||||
recorder.preflightFailures += it
|
||||
}
|
||||
val serverWs = harness.awaitServerSocket()
|
||||
harness.awaitRpc("prompt.submit")
|
||||
serverWs.send(harness.eventFrame("message.start", null, "live-1"))
|
||||
serverWs.send(
|
||||
harness.eventFrame(
|
||||
"message.delta",
|
||||
buildJsonObject { put("text", "durable partial") },
|
||||
"live-1",
|
||||
),
|
||||
)
|
||||
awaitCondition { recorder.textDeltas.isNotEmpty() }
|
||||
harness.activeSessionListPayload = exactActiveSessionPayload("idle")
|
||||
|
||||
assertTrue(client.listActiveSessions() is GatewayActiveSessionsResult.Success)
|
||||
assertTrue("idle snapshot did not settle turn", recorder.completeLatch.await(5, TimeUnit.SECONDS))
|
||||
assertEquals(1, recorder.completions.get())
|
||||
assertEquals(1, recorder.reconcileRequests.get())
|
||||
assertTrue(recorder.errors.isEmpty())
|
||||
assertTrue(harness.rpcLog.none { it.first == "session.interrupt" })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `active session idle never settles passively observed turn`() = runBlocking {
|
||||
val recorder = Recorder()
|
||||
client.setUnsolicitedTurnProvider {
|
||||
GatewayInboundTurnRegistration(recorder.callbacks) { true }
|
||||
}
|
||||
assertTrue(client.prewarmAwait("stored-session"))
|
||||
val serverWs = harness.awaitServerSocket()
|
||||
serverWs.send(harness.eventFrame("message.start", null, "live-resumed"))
|
||||
serverWs.send(
|
||||
harness.eventFrame(
|
||||
"message.delta",
|
||||
buildJsonObject { put("text", "desktop-owned") },
|
||||
"live-resumed",
|
||||
),
|
||||
)
|
||||
awaitCondition { recorder.textDeltas.isNotEmpty() }
|
||||
harness.activeSessionListPayload = exactActiveSessionPayload(
|
||||
status = "idle",
|
||||
liveSessionId = "live-resumed",
|
||||
storedSessionId = "stored-session",
|
||||
)
|
||||
|
||||
client.listActiveSessions()
|
||||
assertFalse(recorder.completeLatch.await(250, TimeUnit.MILLISECONDS))
|
||||
assertTrue(harness.rpcLog.none { it.first == "session.interrupt" })
|
||||
|
||||
serverWs.send(
|
||||
harness.eventFrame(
|
||||
"message.complete",
|
||||
buildJsonObject { put("text", "desktop-owned") },
|
||||
"live-resumed",
|
||||
),
|
||||
)
|
||||
assertTrue(recorder.completeLatch.await(5, TimeUnit.SECONDS))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `stale active session snapshot cannot settle newer turn generation`() = runBlocking {
|
||||
val first = Recorder()
|
||||
client.sendTurn(null, "first", null, first.callbacks) { first.preflightFailures += it }
|
||||
val serverWs = harness.awaitServerSocket()
|
||||
harness.awaitRpc("prompt.submit")
|
||||
serverWs.send(harness.eventFrame("message.start", null, "live-1"))
|
||||
serverWs.send(
|
||||
harness.eventFrame("message.delta", buildJsonObject { put("text", "first") }, "live-1"),
|
||||
)
|
||||
awaitCondition { first.textDeltas.isNotEmpty() }
|
||||
|
||||
harness.suppressAckMethods += "session.active_list"
|
||||
val staleSnapshot = scope.async { client.listActiveSessions() }
|
||||
val staleAck = harness.awaitPendingAck()
|
||||
serverWs.send(
|
||||
harness.eventFrame(
|
||||
"message.complete",
|
||||
buildJsonObject { put("text", "first") },
|
||||
"live-1",
|
||||
),
|
||||
)
|
||||
assertTrue(first.completeLatch.await(5, TimeUnit.SECONDS))
|
||||
|
||||
val second = Recorder()
|
||||
client.sendTurn(
|
||||
"20260612_120000_abc123",
|
||||
"second",
|
||||
null,
|
||||
second.callbacks,
|
||||
) { second.preflightFailures += it }
|
||||
harness.awaitRpcCount("prompt.submit", 2)
|
||||
serverWs.send(harness.eventFrame("message.start", null, "live-1"))
|
||||
serverWs.send(
|
||||
harness.eventFrame("message.delta", buildJsonObject { put("text", "second") }, "live-1"),
|
||||
)
|
||||
awaitCondition { second.textDeltas.isNotEmpty() }
|
||||
|
||||
harness.releaseAck(staleAck, exactActiveSessionPayload("idle"))
|
||||
assertTrue(staleSnapshot.await() is GatewayActiveSessionsResult.Success)
|
||||
assertFalse(second.completeLatch.await(250, TimeUnit.MILLISECONDS))
|
||||
assertEquals(0, second.reconcileRequests.get())
|
||||
|
||||
serverWs.send(
|
||||
harness.eventFrame(
|
||||
"message.complete",
|
||||
buildJsonObject { put("text", "second") },
|
||||
"live-1",
|
||||
),
|
||||
)
|
||||
assertTrue(second.completeLatch.await(5, TimeUnit.SECONDS))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `cancellation wins over delayed active session idle snapshot`() = runBlocking {
|
||||
val recorder = Recorder()
|
||||
val handle = client.sendTurn(null, "cancel me", null, recorder.callbacks) {
|
||||
recorder.preflightFailures += it
|
||||
}
|
||||
val serverWs = harness.awaitServerSocket()
|
||||
harness.awaitRpc("prompt.submit")
|
||||
serverWs.send(harness.eventFrame("message.start", null, "live-1"))
|
||||
serverWs.send(
|
||||
harness.eventFrame("message.delta", buildJsonObject { put("text", "partial") }, "live-1"),
|
||||
)
|
||||
awaitCondition { recorder.textDeltas.isNotEmpty() }
|
||||
|
||||
harness.suppressAckMethods += "session.active_list"
|
||||
val delayedSnapshot = scope.async { client.listActiveSessions() }
|
||||
val delayedAck = harness.awaitPendingAck()
|
||||
handle.cancel()
|
||||
harness.awaitRpc("session.interrupt")
|
||||
harness.releaseAck(delayedAck, exactActiveSessionPayload("idle"))
|
||||
|
||||
assertTrue(delayedSnapshot.await() is GatewayActiveSessionsResult.Success)
|
||||
assertEquals(0, recorder.completions.get())
|
||||
assertEquals(0, recorder.reconcileRequests.get())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `late terminal after active session settle is consumed once`() = runBlocking {
|
||||
val recorder = Recorder()
|
||||
val unmatched = ConcurrentLinkedQueue<GatewayBackgroundTurnCompletion>()
|
||||
client.setUnmatchedTurnCompleteListener(unmatched::add)
|
||||
client.sendTurn(null, "late terminal", null, recorder.callbacks) {
|
||||
recorder.preflightFailures += it
|
||||
}
|
||||
val serverWs = harness.awaitServerSocket()
|
||||
harness.awaitRpc("prompt.submit")
|
||||
serverWs.send(harness.eventFrame("message.start", null, "live-1"))
|
||||
serverWs.send(
|
||||
harness.eventFrame("message.delta", buildJsonObject { put("text", "done") }, "live-1"),
|
||||
)
|
||||
awaitCondition { recorder.textDeltas.isNotEmpty() }
|
||||
harness.activeSessionListPayload = exactActiveSessionPayload("idle")
|
||||
client.listActiveSessions()
|
||||
assertTrue(recorder.completeLatch.await(5, TimeUnit.SECONDS))
|
||||
|
||||
serverWs.send(
|
||||
harness.eventFrame(
|
||||
"message.complete",
|
||||
buildJsonObject { put("text", "done") },
|
||||
"live-1",
|
||||
),
|
||||
)
|
||||
Thread.sleep(150)
|
||||
assertEquals(1, recorder.completions.get())
|
||||
assertTrue(unmatched.isEmpty())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `idle watchdog does not fire while events keep arriving slowly`() {
|
||||
rebuildClient(turnIdleTimeoutMs = 1_000L)
|
||||
|
||||
+45
@@ -3014,6 +3014,51 @@ class ChatViewModelGatewayInboundTurnTest {
|
||||
assertTrue(viewModel.queuedMessages.value.isEmpty())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun queuedCorrectionDrainsOnceAfterOwnedTurnSettlesFromActiveSessionIdle() = runBlocking {
|
||||
viewModel.switchProfileContext(PROFILE_CONTEXT, STORED_SESSION_ID)
|
||||
gatewayHarness.redirectStatus = "rejected"
|
||||
viewModel.sendMessage("Original Android turn")
|
||||
gatewayHarness.awaitRpc("prompt.submit")
|
||||
serverWs.send(gatewayHarness.eventFrame("message.start", null, "live-resumed"))
|
||||
serverWs.send(
|
||||
gatewayHarness.eventFrame(
|
||||
"message.delta",
|
||||
buildJsonObject { put("text", "Answer without terminal") },
|
||||
"live-resumed",
|
||||
),
|
||||
)
|
||||
awaitCondition { handler.isStreaming.value }
|
||||
|
||||
viewModel.sendMessage("Queued correction")
|
||||
gatewayHarness.awaitRpc("session.redirect")
|
||||
awaitCondition { viewModel.queuedMessages.value == listOf("Queued correction") }
|
||||
persistedHistory = persistedAnswerHistory("Answer without terminal", "settled-answer")
|
||||
gatewayHarness.activeSessionListPayload = buildJsonObject {
|
||||
put("sessions", buildJsonArray {
|
||||
add(buildJsonObject {
|
||||
put("id", "live-resumed")
|
||||
put("session_key", STORED_SESSION_ID)
|
||||
put("status", "idle")
|
||||
put("last_active", 1_777_000_000.0)
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
gatewayClient.listActiveSessions()
|
||||
shadowOf(Looper.getMainLooper()).idle()
|
||||
|
||||
awaitCondition {
|
||||
gatewayHarness.rpcLog.count { (method, params) ->
|
||||
method == "prompt.submit" &&
|
||||
params["text"] == JsonPrimitive("Queued correction") &&
|
||||
params["queued"] == JsonPrimitive(true)
|
||||
} == 1
|
||||
}
|
||||
assertTrue(viewModel.queuedMessages.value.isEmpty())
|
||||
assertTrue(viewModel.steerableTurn.value)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun multipleQueuedMessagesDrainAsAnOwnedRunChain() {
|
||||
viewModel.switchProfileContext(
|
||||
|
||||
+14
-9
@@ -3710,15 +3710,20 @@ live even when upstream had already persisted the final answer. The earlier
|
||||
visible-chat Idle reattach fix covered a socket that closed after foreground
|
||||
prewarm; it did not cover this already-active turn state.
|
||||
|
||||
**Decision.** A Gateway turn may settle from `session.activate` or exact-session
|
||||
`session.info` only when the turn has already received turn-scoped activity and
|
||||
upstream reports `running=false`. Pre-start idle snapshots are ignored because
|
||||
they can race prompt admission. This backstop is a successful server-owned
|
||||
settle, not cancellation or transport failure: Android completes the local
|
||||
stream, keeps the durable session identity, performs bounded identity-fenced
|
||||
history reconciliation, and never resubmits through API fallback. Cold open
|
||||
continues to use `session.resume`; an authoritative resume rejection remains
|
||||
visible and cannot create or switch to a replacement context.
|
||||
**Decision.** A Gateway turn may settle from `session.activate`, exact-session
|
||||
`session.info`, or an exact live/durable `session.active_list` row only when the
|
||||
turn is Android-owned, has already received turn-scoped activity, and upstream
|
||||
reports `running=false` or Idle. The active-list request captures the exact turn
|
||||
and its progress generation; a session/profile switch, cancellation, newer turn,
|
||||
or intervening live event rejects the delayed snapshot. Pre-start idle snapshots
|
||||
and passively observed Desktop/TUI turns are never eligible. This backstop is a
|
||||
successful server-owned settle, not cancellation or transport failure: Android
|
||||
completes the local stream, keeps the durable session identity, performs bounded
|
||||
identity-fenced history reconciliation, consumes one late terminal without
|
||||
double-completion, and never resubmits through API fallback. A locally queued
|
||||
correction drains once through its existing owner chain after settlement. Cold
|
||||
open continues to use `session.resume`; an authoritative resume rejection
|
||||
remains visible and cannot create or switch to a replacement context.
|
||||
|
||||
The recovery writes one bounded content-free diagnostic containing only route,
|
||||
missing-terminal phase, and reconciliation action. It records no prompt or
|
||||
|
||||
@@ -73,6 +73,7 @@ the upstream contract identifiers it depends on.
|
||||
| `queued_follow_up` | Two explicitly owned turns and ordered queue drainage |
|
||||
| `scope_rejection_inputs` | Exact, foreign, and unscoped event inputs |
|
||||
| `terminal_gap_activate` | Socket closes after live output; replacement `session.activate` reports `running=false`; history is authoritative |
|
||||
| `terminal_gap_active_list` | An exact Android-owned turn receives deltas but no terminal; `session.active_list` reports the same live/durable owner idle; history is authoritative |
|
||||
| `terminal_gap_session_info` | Scoped `session.info {running:false}` settles a turn without `message.complete` |
|
||||
| `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 |
|
||||
|
||||
+1
-1
@@ -610,7 +610,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. A profile switch marks the replacement list loading before clearing the previous profile's rows and keeps that state until the exact-profile fetch settles, so an empty-state claim never flashes before server truth arrives. 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.
|
||||
- **Cold profile hydration** — a persisted named profile scopes its Dashboard session directory and last-session restore immediately, before `/api/profiles` metadata is available. Server-default selection waits for the lightweight active-profile scope. Roster, avatars, pets, skills, and model metadata never precede the first directory result. The startup sphere releases after route selection; Chat keeps identity and cached rows mounted while its existing animated status surfaces show Gateway wake, session restore, and directory loading.
|
||||
- **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 through exact foreground/detached ownership already held by that client, explicit profile metadata if a future upstream sends it, or the currently selected passive session when its durable id has exactly one owner in the current connection directory. Duplicate same-id owners across profiles remain unresolved and 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, ambiguous bare session ids, and delayed snapshots cannot revive newer settled state.
|
||||
- **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**; an exact terminal, `session.info {running:false}`, or an exact live/durable active-list row reporting Idle can settle only the matching Android-owned turn and progress generation. Because active-list rows normally have no profile metadata, Android assigns a row through exact foreground/detached ownership already held by that client, explicit profile metadata if a future upstream sends it, or the currently selected passive session when its durable id has exactly one owner in the current connection directory. Duplicate same-id owners across profiles remain unresolved and 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, ambiguous bare session ids, delayed snapshots, and snapshots crossed by newer turn events cannot settle or revive a newer generation.
|
||||
- **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 Direct API compatibility chat 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.
|
||||
|
||||
@@ -42,7 +42,7 @@ Verified upstream source snapshot:
|
||||
| `/v1/skills`, `/v1/toolsets` | Upstream API server | No | Discovery | Authenticated read-only API-server skill/toolset inventory; Android Diagnostics summarizes enabled toolsets and Relay tool visibility. |
|
||||
| Dashboard `/api/status`, `/api/auth/me` | Upstream dashboard | No | Manage auth and post-selection diagnostics | Dashboard cookie/session path; separate from API bearer. Optional status diagnostics include Nous bootstrap validity, resource pressure, and profile/gateway topology; these do not gate transport selection. |
|
||||
| Dashboard `/api/auth/ws-ticket`, `/api/ws` | Upstream dashboard/tui_gateway | No | Preferred chat transport | Vanilla Hermes gateway chat path with live reasoning/thinking events. `message.complete` is the ordinary terminal event; `session.info {running:false}` is the authoritative settle backstop when a replacement socket missed that terminal frame. A reconnect reactivates the exact live runtime with `session.activate`; durable `session.resume` remains the cold-open path and an explicit rejection never creates a replacement context. |
|
||||
| Gateway `session.active_list` | Upstream tui_gateway | No | Authoritative process-wide live activity | Returns attachable runtimes across the Gateway process, with live `id`, durable `session_key`, and `starting`, `working`, `waiting`, or `idle`. The only optional selector is `current_session_id`; rows normally carry no profile metadata. Android attributes a row from exact foreground/detached ownership already held by that client, explicit profile metadata if a future upstream sends it, or a unique match to the currently selected passive session in the current connection directory. Duplicate same-id owners across profiles remain unresolved. Unresolved rows stay unattributed, and absence settles a scope only after a complete, unambiguously resolved successful snapshot. Method-not-found or refresh failure is Unavailable, not Idle. Pending input outranks running work. |
|
||||
| Gateway `session.active_list` | Upstream tui_gateway | No | Authoritative process-wide live activity | Returns attachable runtimes across the Gateway process, with live `id`, durable `session_key`, and `starting`, `working`, `waiting`, or `idle`. The only optional selector is `current_session_id`; rows normally carry no profile metadata. Android attributes a row from exact foreground/detached ownership already held by that client, explicit profile metadata if a future upstream sends it, or a unique match to the currently selected passive session in the current connection directory. Duplicate same-id owners across profiles remain unresolved. An exact live/durable Idle row may settle only the same Android-owned turn and unchanged progress generation when a terminal frame is missing; it never claims a passively observed Desktop/TUI turn. Unresolved rows stay unattributed, and absence settles a scope only after a complete, unambiguously resolved successful snapshot. Method-not-found or refresh failure is Unavailable, not Idle. Pending input outranks running work. |
|
||||
| Dashboard `model.options` / `/api/model/*` | Upstream dashboard/tui_gateway | No | Provider/model inventory and selection | Source of truth for coherent provider/model identities. A reasoning boolean or exact effort list is consumed when present; clients do not infer provider identity from a model string alone. |
|
||||
| Gateway `pet.info`, `pet.gallery`, `pet.select`, `pet.disable` | Upstream tui_gateway | No | Profile-scoped animated companion | `pet.info` supplies bounded PNG/WebP sheet bytes, revision, geometry, real frame counts, loop timing, scale, and row taxonomy. Android passes `knownRevision` to avoid duplicate sheet transfer, renders the active pet through its native activity-aware companion, and keeps phone-local pet packs separate. All four RPCs carry the effective profile. |
|
||||
| Dashboard `/api/audio/transcribe`, `/api/audio/speak-stream`, `/api/audio/speak` | Upstream dashboard | No | Vanilla Hermes voice | Manage sign-in unlocks Vanilla Hermes voice. Assistant text streams into upstream speech when available; older hosts fall back to whole-request speech before audio starts. API server has no `/v1/audio/*` route today. |
|
||||
|
||||
@@ -239,6 +239,37 @@ class FixtureTestCase(unittest.IsolatedAsyncioTestCase):
|
||||
self.assertEqual(fixture.scenario.live_session_id, events[-1]["session_id"])
|
||||
self.assertNotIn("message.complete", [event["type"] for event in events])
|
||||
|
||||
async def test_active_list_idle_settles_without_message_complete(self) -> None:
|
||||
fixture, base_url = await self.start("terminal_gap_active_list")
|
||||
ws, _ = await self.connect(base_url)
|
||||
await self.rpc(ws, 1, "prompt.submit", {"text": "fixture"})
|
||||
frames = await self.frames_until(
|
||||
ws,
|
||||
lambda frame: frame.get("params", {}).get("type") == "message.delta",
|
||||
)
|
||||
for _ in range(50):
|
||||
async with self.session.get(f"{base_url}/__fixture__/state") as response:
|
||||
state = await response.json()
|
||||
if not state["running"]:
|
||||
break
|
||||
await asyncio.sleep(0.01)
|
||||
self.assertFalse(state["running"])
|
||||
|
||||
await self.rpc(
|
||||
ws,
|
||||
2,
|
||||
"session.active_list",
|
||||
{"current_session_id": fixture.scenario.live_session_id},
|
||||
)
|
||||
snapshot = (await ws.receive_json())["result"]["sessions"]
|
||||
self.assertEqual("idle", snapshot[0]["status"])
|
||||
self.assertEqual(fixture.scenario.live_session_id, snapshot[0]["id"])
|
||||
self.assertEqual(fixture.scenario.stored_session_id, snapshot[0]["session_key"])
|
||||
self.assertNotIn(
|
||||
"message.complete",
|
||||
[frame.get("params", {}).get("type") for frame in frames],
|
||||
)
|
||||
|
||||
async def test_queued_follow_up_runs_after_first_turn(self) -> None:
|
||||
_, base_url = await self.start("queued_follow_up")
|
||||
ws, _ = await self.connect(base_url)
|
||||
@@ -353,6 +384,7 @@ class ScenarioTestCase(unittest.TestCase):
|
||||
"rapid_tools_interims",
|
||||
"subagent_child_preview",
|
||||
"terminal_gap_activate",
|
||||
"terminal_gap_active_list",
|
||||
"terminal_gap_session_info",
|
||||
"queued_follow_up",
|
||||
"scope_rejection_inputs",
|
||||
@@ -389,6 +421,10 @@ class ScenarioTestCase(unittest.TestCase):
|
||||
),
|
||||
load_scenario("terminal_gap_activate").contract_requirements,
|
||||
)
|
||||
self.assertEqual(
|
||||
("gateway.session_active_list",),
|
||||
load_scenario("terminal_gap_active_list").contract_requirements,
|
||||
)
|
||||
self.assertEqual(
|
||||
("gateway.settled_session_info",),
|
||||
load_scenario("terminal_gap_session_info").contract_requirements,
|
||||
|
||||
@@ -0,0 +1,33 @@
|
||||
{
|
||||
"name": "terminal_gap_active_list",
|
||||
"live_session_id": "fixture-live-1",
|
||||
"stored_session_id": "20260831_112003_fixture",
|
||||
"profile": "default",
|
||||
"contract_requirements": [
|
||||
"gateway.session_active_list"
|
||||
],
|
||||
"turns": [
|
||||
{
|
||||
"steps": [
|
||||
{"op": "event", "type": "message.start"},
|
||||
{"op": "event", "type": "message.delta", "payload": {"text": "Persisted without a terminal frame."}},
|
||||
{
|
||||
"op": "persist",
|
||||
"messages": [
|
||||
{"id": 1, "role": "user", "content": "Exercise active-list settlement.", "timestamp": 1.0},
|
||||
{"id": 2, "role": "assistant", "content": "Persisted without a terminal frame.", "timestamp": 2.0, "finish_reason": "stop"}
|
||||
]
|
||||
},
|
||||
{"op": "set_running", "value": false}
|
||||
]
|
||||
}
|
||||
],
|
||||
"active_list": {
|
||||
"supported": true,
|
||||
"snapshots": [
|
||||
[
|
||||
{"id": "fixture-live-1", "session_key": "20260831_112003_fixture", "status": "idle", "current": true}
|
||||
]
|
||||
]
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user