From 1b76dd250e9e05ac01a8cf77961541b2eafd37bb Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 21 Jul 2026 23:34:29 +0000 Subject: [PATCH] feat(buzz): live observer telemetry tab in agent console Adds a third Observer tab streaming ephemeral agent telemetry frames (ObserverFrameEvent, kind:24200, p=owner) live from every Buzz-dialect relay via subscribeAsFlow. Frames are decrypted (frame:telemetry only), deduped across relays, sorted newest-first, and bounded to a 200-row ring. The live REQ is opened only while the Observer tab is on screen (started in a DisposableEffect, cancelled on dispose and in onCleared) because observer frames are never persisted by relays or LocalCache. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01J8KBSw6smQRyXLiWHeDsZ8 --- .../loggedIn/buzz/AgentConsoleScreen.kt | 65 ++++++++++++- .../loggedIn/buzz/AgentConsoleViewModel.kt | 96 +++++++++++++++++++ 2 files changed, 159 insertions(+), 2 deletions(-) diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/buzz/AgentConsoleScreen.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/buzz/AgentConsoleScreen.kt index 53a185748e..6c5bb6a685 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/buzz/AgentConsoleScreen.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/buzz/AgentConsoleScreen.kt @@ -39,6 +39,7 @@ import androidx.compose.material3.Tab import androidx.compose.material3.TabRow import androidx.compose.material3.Text import androidx.compose.runtime.Composable +import androidx.compose.runtime.DisposableEffect import androidx.compose.runtime.getValue import androidx.compose.runtime.mutableIntStateOf import androidx.compose.runtime.remember @@ -81,10 +82,19 @@ fun AgentConsoleScreen( val metrics by viewModel.metrics.collectAsStateWithLifecycle() val personas by viewModel.personas.collectAsStateWithLifecycle() + val observerFrames by viewModel.observerFrames.collectAsStateWithLifecycle() val isLoading by viewModel.isLoading.collectAsStateWithLifecycle() var selectedTab by rememberSaveable { mutableIntStateOf(0) } - val tabs = remember { listOf("Costs", "Personas") } + val tabs = remember { listOf("Costs", "Personas", "Observer") } + + // The observer stream is a live REQ; only keep it open while its tab is on screen. + if (selectedTab == 2) { + DisposableEffect(Unit) { + viewModel.startObserving() + onDispose { viewModel.stopObserving() } + } + } Scaffold( topBar = { TopBarWithBackButton("Agent Console", nav) }, @@ -103,7 +113,8 @@ fun AgentConsoleScreen( Box(modifier = Modifier.fillMaxSize()) { when (selectedTab) { 0 -> CostsTab(metrics) - else -> PersonasTab(personas) + 1 -> PersonasTab(personas) + else -> ObserverTab(observerFrames) } if (isLoading) { @@ -244,6 +255,56 @@ private fun PersonaCardView(persona: AgentConsoleViewModel.PersonaCard) { } } +@Composable +private fun ObserverTab(frames: List) { + if (frames.isEmpty()) { + EmptyState("Listening for live agent telemetry (kind:24200). Frames appear here while your agents are running.") + return + } + + LazyColumn( + modifier = Modifier.fillMaxSize().padding(horizontal = 16.dp, vertical = 8.dp), + verticalArrangement = Arrangement.spacedBy(4.dp), + ) { + items(frames) { frame -> + ObserverFrameRow(frame) + HorizontalDivider() + } + } +} + +@Composable +private fun ObserverFrameRow(frame: AgentConsoleViewModel.ObserverRow) { + Column(modifier = Modifier.fillMaxWidth().padding(vertical = 6.dp), verticalArrangement = Arrangement.spacedBy(2.dp)) { + Row(modifier = Modifier.fillMaxWidth(), horizontalArrangement = Arrangement.SpaceBetween) { + Text( + text = frame.kind, + style = MaterialTheme.typography.titleSmall, + fontFamily = FontFamily.Monospace, + fontWeight = FontWeight.Bold, + ) + Text( + text = frame.timestamp, + style = MaterialTheme.typography.bodySmall, + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + } + val context = + buildString { + append(shortKey(frame.agentPubKey)) + frame.sessionId?.let { append(" · session ${it.take(8)}") } + frame.turnId?.let { append(" · turn ${it.take(8)}") } + append(" · #${frame.seq}") + } + Text( + text = context, + style = MaterialTheme.typography.bodySmall, + fontFamily = FontFamily.Monospace, + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + } +} + @Composable private fun MetaRow( label: String, diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/buzz/AgentConsoleViewModel.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/buzz/AgentConsoleViewModel.kt index 3a85b073ec..6c45c3d884 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/buzz/AgentConsoleViewModel.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/buzz/AgentConsoleViewModel.kt @@ -31,17 +31,23 @@ import com.vitorpamplona.amethyst.model.LocalCache import com.vitorpamplona.amethyst.model.filter import com.vitorpamplona.quartz.buzz.amTurnMetrics.AgentTurnMetricEvent import com.vitorpamplona.quartz.buzz.amTurnMetrics.AgentTurnMetricPayload +import com.vitorpamplona.quartz.buzz.aoObserver.ObserverFrameEvent +import com.vitorpamplona.quartz.buzz.aoObserver.tags.FrameTag import com.vitorpamplona.quartz.buzz.apPersonas.PersonaEvent import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPagesFromPool +import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.subscribeAsFlow import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow +import kotlinx.coroutines.flow.update import kotlinx.coroutines.launch import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock +import java.util.Collections /** * Backing ViewModel for the [AgentConsoleScreen] — the workspace owner's read-only @@ -75,6 +81,14 @@ class AgentConsoleViewModel : ViewModel() { private val _isLoading = MutableStateFlow(false) val isLoading: StateFlow = _isLoading.asStateFlow() + private val _observerFrames = MutableStateFlow>(emptyList()) + val observerFrames: StateFlow> = _observerFrames.asStateFlow() + + private var observerJob: Job? = null + + /** Dedups ephemeral frames across relays and re-emissions (accessed off multiple readers). */ + private val observerSeen = Collections.synchronizedSet(HashSet()) + fun bindAccountIfMissing(account: Account) { if (this.account != null) return this.account = account @@ -157,6 +171,83 @@ class AgentConsoleViewModel : ViewModel() { _metrics.value = AgentFleetAggregator.aggregate(decrypted) } + /** + * Opens a live subscription to the owner's ephemeral observer telemetry frames + * ([ObserverFrameEvent], `kind:24200`, `p` = owner) across every Buzz-dialect relay, + * decrypts the `frame:telemetry` bodies, and pushes them newest-first into + * [observerFrames] (bounded to [MAX_OBSERVER_ROWS]). Idempotent; call [stopObserving] + * to tear the subscription down. Observer frames are never stored by relays or + * [LocalCache], so this live REQ is the only way to see them. + */ + fun startObserving() { + val account = account ?: return + if (observerJob != null) return + + val relays = BuzzRelayDialect.flow.value + if (relays.isEmpty()) return + + val signer = account.signer + val myPubkey = account.userProfile().pubkeyHex + val filter = Filter(kinds = listOf(ObserverFrameEvent.KIND), tags = mapOf("p" to listOf(myPubkey))) + + observerJob = + viewModelScope.launch(Dispatchers.IO) { + relays.forEach { relay -> + launch { + account.client.subscribeAsFlow(relay, filter).collect { events -> + val fresh = + events + .filterIsInstance() + .filter { observerSeen.add(it.id) } + if (fresh.isEmpty()) return@collect + + val rows = + fresh.mapNotNull { frame -> + if (frame.frame() != FrameTag.TELEMETRY) return@mapNotNull null + val payload = frame.decryptTelemetryOrNull(signer) ?: return@mapNotNull null + ObserverRow( + seq = payload.seq, + timestamp = payload.timestamp, + kind = payload.kind, + agentPubKey = frame.agentPubKey() ?: frame.pubKey, + sessionId = payload.sessionId, + turnId = payload.turnId, + ) + } + if (rows.isNotEmpty()) { + _observerFrames.update { existing -> + (rows + existing) + .sortedByDescending { it.timestamp } + .take(MAX_OBSERVER_ROWS) + } + } + } + } + } + } + } + + fun stopObserving() { + observerJob?.cancel() + observerJob = null + } + + override fun onCleared() { + stopObserving() + super.onCleared() + } + + /** One decrypted observer telemetry frame rendered on the Observer tab. */ + @Immutable + data class ObserverRow( + val seq: Long, + val timestamp: String, + val kind: String, + val agentPubKey: HexKey, + val sessionId: String?, + val turnId: String?, + ) + /** A persona rendered on the Personas tab; a flattened projection of [PersonaEvent]. */ @Immutable data class PersonaCard( @@ -167,4 +258,9 @@ class AgentConsoleViewModel : ViewModel() { val provider: String?, val systemPrompt: String?, ) + + companion object { + /** Ring size for the ephemeral observer stream — enough to scroll recent activity. */ + const val MAX_OBSERVER_ROWS = 200 + } }