From 7da997a567e7217a8d405c9cd0c791724e044c84 Mon Sep 17 00:00:00 2001 From: kenny lopez Date: Mon, 17 Aug 2026 15:48:43 +0100 Subject: [PATCH] Scope mobile agent activity to threads Signed-off-by: kenny lopez --- crates/buzz-acp/src/acp.rs | 5 + crates/buzz-acp/src/lib.rs | 68 +++++++++-- crates/buzz-acp/src/observer.rs | 9 ++ crates/buzz-acp/src/pool.rs | 112 ++++++++++++++++-- .../features/agents/ui/agentSessionTypes.ts | 1 + docs/nips/NIP-AO.md | 15 ++- .../agent_activity/active_agent_turns.dart | 16 ++- .../composer_agent_activity_indicator.dart | 2 + .../agent_activity/observer_models.dart | 3 + .../agent_activity/observer_subscription.dart | 1 + .../agent_activity/working_bots_provider.dart | 41 ++++--- .../active_agent_turns_test.dart | 42 +++++++ .../observer_subscription_test.dart | 4 + .../working_bots_provider_test.dart | 53 +++++++++ 14 files changed, 323 insertions(+), 49 deletions(-) diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index f04b8eeec..1c44e4afa 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -577,6 +577,11 @@ impl AcpClient { self.observer_context = context; } + /// Return the observer metadata for the current turn. + pub(crate) fn observer_context(&self) -> ObserverContext { + self.observer_context.clone() + } + /// Return a clone of the observer handle, if attached. pub(crate) fn observer_handle(&self) -> Option { self.observer.clone() diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 7fd40b83d..0e78ece92 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -591,6 +591,7 @@ fn batch_envelope(events: &[observer::ObserverEvent]) -> observer::ObserverEvent kind: OBSERVER_BATCH_KIND.to_string(), agent_index: last.agent_index, channel_id: last.channel_id.clone(), + thread_head_id: last.thread_head_id.clone(), session_id: last.session_id.clone(), turn_id: last.turn_id.clone(), started_at: last.started_at.clone(), @@ -1259,6 +1260,7 @@ fn emit_project_owner_control_result( None, &observer::ObserverContext { channel_id: None, + thread_head_id: None, session_id: None, turn_id: None, started_at: None, @@ -1296,6 +1298,7 @@ fn handle_cancel_turn_control( None, &observer::ObserverContext { channel_id: Some(channel_id.to_string()), + thread_head_id: None, session_id: None, turn_id: None, started_at: None, @@ -1373,6 +1376,7 @@ fn handle_switch_model_control( None, &observer::ObserverContext { channel_id: Some(channel_id.to_string()), + thread_head_id: None, session_id: None, turn_id: None, started_at: None, @@ -3765,6 +3769,7 @@ fn dispatch_pending( pool::TaskMeta { agent_index, channel_id: Some(channel_id), + thread_head_id: typing_scope.root_event_id.clone(), turn_id, recoverable_batch, control_tx: Some(control_tx), @@ -4029,11 +4034,16 @@ fn handle_prompt_result( .to_string(); let harness_pid = std::process::id(); - let channel_id = match &result.source { - PromptSource::Channel(ch) => Some(*ch), - PromptSource::Heartbeat => None, - }; - let turn_id = result.turn_id.clone(); + let mut observer_context = result.agent.acp.observer_context(); + if observer_context.channel_id.is_none() { + observer_context.channel_id = match &result.source { + PromptSource::Channel(channel_id) => Some(channel_id.to_string()), + PromptSource::Heartbeat => None, + }; + } + if observer_context.turn_id.is_none() { + observer_context.turn_id = Some(result.turn_id.clone()); + } let emit_turn_error = |error_msg: &str, error_code: Option| { if let Some(ref observer) = observer { let mut payload = serde_json::json!({ @@ -4043,12 +4053,7 @@ fn handle_prompt_result( if let Some(code) = error_code { payload["code"] = serde_json::json!(code); } - observer.emit( - "turn_error", - Some(agent_index), - &observer::context_for(channel_id, None, Some(turn_id.clone())), - payload, - ); + observer.emit("turn_error", Some(agent_index), &observer_context, payload); } }; @@ -4275,7 +4280,13 @@ fn recover_panicked_agent( observer.emit( "agent_panic", Some(i), - &observer::context_for(meta.channel_id, None, Some(meta.turn_id)), + &observer::ObserverContext { + channel_id: meta.channel_id.map(|channel_id| channel_id.to_string()), + thread_head_id: meta.thread_head_id, + session_id: None, + turn_id: Some(meta.turn_id), + started_at: None, + }, serde_json::json!({ "outcome": "panic", "error": format!("Agent task panicked: {join_error}"), @@ -4401,6 +4412,7 @@ fn dispatch_heartbeat( pool::TaskMeta { agent_index, channel_id: None, + thread_head_id: None, turn_id, recoverable_batch: None, control_tx: None, @@ -5195,6 +5207,7 @@ mod owner_control_command_tests { pool::TaskMeta { agent_index: 0, channel_id: Some(channel_id), + thread_head_id: None, turn_id: "test-turn-id".to_string(), recoverable_batch: None, control_tx: Some(control_tx), @@ -5758,6 +5771,7 @@ mod observer_publish_queue_tests { kind: kind.to_string(), agent_index: Some(0), channel_id: channel.map(ToOwned::to_owned), + thread_head_id: None, session_id: Some("session-1".to_string()), turn_id: Some("turn-1".to_string()), started_at: None, @@ -6626,6 +6640,7 @@ mod observer_chunk_coalescer_tests { kind: "acp_read".to_string(), agent_index: Some(0), channel_id: Some("channel-1".to_string()), + thread_head_id: None, session_id: Some("session-1".to_string()), turn_id: Some("turn-1".to_string()), started_at: None, @@ -6654,6 +6669,7 @@ mod observer_chunk_coalescer_tests { kind: "turn_started".to_string(), agent_index: Some(0), channel_id: Some("channel-1".to_string()), + thread_head_id: None, session_id: Some("session-1".to_string()), turn_id: Some("turn-1".to_string()), started_at: None, @@ -7049,6 +7065,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: Some(channel_id), + thread_head_id: None, turn_id: "test-turn-id".into(), recoverable_batch: None, control_tx: None, @@ -7121,6 +7138,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: Some(channel_id), + thread_head_id: None, turn_id: "test-turn-id".into(), recoverable_batch: None, control_tx: None, @@ -7236,6 +7254,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: Some(channel_id), + thread_head_id: None, turn_id: "test-turn-id".into(), recoverable_batch: None, control_tx: None, @@ -7288,7 +7307,11 @@ mod error_outcome_emission_tests { /// Drive one error outcome through `handle_prompt_result` and return how /// many `turn_error` events it emitted to the observer feed. async fn turn_errors_emitted_for(outcome: PromptOutcome) -> usize { - let agent = dummy_agent(0).await; + let mut agent = dummy_agent(0).await; + let mut observer_context = + crate::observer::context_for(None, None, Some("test-turn-id".to_string())); + observer_context.thread_head_id = Some("thread-1".to_string()); + agent.acp.set_observer_context(observer_context); let mut pool = AgentPool::from_slots(vec![None]); // `handle_prompt_result` asserts it removes exactly one in-flight task @@ -7301,6 +7324,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: None, + thread_head_id: None, turn_id: "test-turn-id".to_string(), recoverable_batch: None, control_tx: None, @@ -7355,6 +7379,12 @@ mod error_outcome_emission_tests { .all(|event| event.turn_id.as_deref() == Some("test-turn-id")), "turn_error must retain the completed turn id" ); + assert!( + turn_errors + .iter() + .all(|event| event.thread_head_id.as_deref() == Some("thread-1")), + "turn_error must retain the completed turn thread" + ); turn_errors.len() } @@ -7378,6 +7408,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: Some(channel_id), + thread_head_id: Some("thread-1".into()), turn_id: "panic-turn-id".to_string(), recoverable_batch: None, control_tx: None, @@ -7427,6 +7458,7 @@ mod error_outcome_emission_tests { Some(channel_id.to_string().as_str()) ); assert_eq!(panic.turn_id.as_deref(), Some("panic-turn-id")); + assert_eq!(panic.thread_head_id.as_deref(), Some("thread-1")); } #[tokio::test] @@ -7471,6 +7503,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: None, + thread_head_id: None, turn_id: "test-turn-id".to_string(), recoverable_batch: None, control_tx: None, @@ -7563,6 +7596,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: None, + thread_head_id: None, turn_id: "test-turn-id".to_string(), recoverable_batch: None, control_tx: None, @@ -7669,6 +7703,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: None, + thread_head_id: None, turn_id: "test-turn-id".to_string(), recoverable_batch: None, control_tx: None, @@ -7746,6 +7781,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: None, + thread_head_id: None, turn_id: "test-turn-id".to_string(), recoverable_batch: None, control_tx: None, @@ -7841,6 +7877,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: None, + thread_head_id: None, turn_id: "test-turn-id".to_string(), recoverable_batch: None, control_tx: None, @@ -7958,6 +7995,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: None, + thread_head_id: None, turn_id: "test-turn-id".to_string(), recoverable_batch: None, control_tx: None, @@ -8098,6 +8136,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: None, + thread_head_id: None, turn_id: "test-turn-id".to_string(), recoverable_batch: None, control_tx: None, @@ -8287,6 +8326,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: None, + thread_head_id: None, turn_id: "test-turn-id".to_string(), recoverable_batch: None, control_tx: None, @@ -8373,6 +8413,7 @@ mod error_outcome_emission_tests { crate::pool::TaskMeta { agent_index: 0, channel_id: None, + thread_head_id: None, turn_id: "test-turn-id".to_string(), recoverable_batch: None, control_tx: None, @@ -8437,6 +8478,7 @@ mod observer_payload_trim_tests { kind: kind.to_string(), agent_index: Some(0), channel_id: Some("11111111-1111-1111-1111-111111111111".to_string()), + thread_head_id: None, session_id: Some("sess-1".to_string()), turn_id: Some("turn-1".to_string()), started_at: None, diff --git a/crates/buzz-acp/src/observer.rs b/crates/buzz-acp/src/observer.rs index 7029e5af6..e4579524b 100644 --- a/crates/buzz-acp/src/observer.rs +++ b/crates/buzz-acp/src/observer.rs @@ -22,6 +22,8 @@ const OBSERVER_BUFFER_CAP: usize = 1_000; pub struct ObserverContext { /// Buzz channel UUID for the current turn, when channel-scoped. pub channel_id: Option, + /// NIP-10 thread root for the current turn, when thread-scoped. + pub thread_head_id: Option, /// ACP session ID associated with the current turn, once known. pub session_id: Option, /// Local UUID for one prompt turn. @@ -67,6 +69,9 @@ pub struct ObserverEvent { pub agent_index: Option, /// Buzz channel UUID for channel-scoped events. pub channel_id: Option, + /// NIP-10 thread root for thread-scoped events. + #[serde(skip_serializing_if = "Option::is_none")] + pub thread_head_id: Option, /// ACP session ID when known. pub session_id: Option, /// Local UUID for one prompt turn. @@ -114,6 +119,7 @@ impl ObserverHandle { kind: kind.into(), agent_index, channel_id: context.channel_id.clone(), + thread_head_id: context.thread_head_id.clone(), session_id: context.session_id.clone(), turn_id: context.turn_id.clone(), started_at: context.started_at.clone(), @@ -144,6 +150,7 @@ pub fn context_for( ) -> ObserverContext { ObserverContext { channel_id: channel_id.map(|id| id.to_string()), + thread_head_id: None, session_id, turn_id, started_at: None, @@ -153,12 +160,14 @@ pub fn context_for( /// Attach the authoritative start timestamp to every observer frame for a turn. pub fn context_for_turn( channel_id: Option, + thread_head_id: Option, session_id: Option, turn_id: String, started_at: String, ) -> ObserverContext { ObserverContext { channel_id: channel_id.map(|id| id.to_string()), + thread_head_id, session_id, turn_id: Some(turn_id), started_at: Some(started_at), diff --git a/crates/buzz-acp/src/pool.rs b/crates/buzz-acp/src/pool.rs index 923d19a3e..45a1f2fbe 100644 --- a/crates/buzz-acp/src/pool.rs +++ b/crates/buzz-acp/src/pool.rs @@ -59,6 +59,8 @@ pub struct SuccessfulSteerDelivery { pub struct TaskMeta { pub agent_index: usize, pub channel_id: Option, + /// NIP-10 thread root for panic recovery telemetry. + pub thread_head_id: Option, /// Identifies terminal events when the task panics before returning a result. pub turn_id: String, /// Clone of batch for Queue mode panic recovery. @@ -1462,6 +1464,12 @@ fn send_prompt_result( }); } +fn observer_thread_head_id(batch: Option<&FlushBatch>) -> Option { + batch + .and_then(|batch| batch.events.last()) + .and_then(|event| crate::queue::parse_thread_tags(&event.event).root_event_id) +} + /// Core async function spawned for each prompt. /// /// Lifecycle: @@ -1492,9 +1500,11 @@ pub async fn run_prompt_task( PromptSource::Channel(channel_id) => Some(*channel_id), PromptSource::Heartbeat => None, }; + let observer_thread_head_id = observer_thread_head_id(batch.as_ref()); let turn_started_at = chrono::Utc::now().to_rfc3339(); agent.acp.set_observer_context(observer::context_for_turn( observer_channel_id, + observer_thread_head_id.clone(), None, turn_id.clone(), turn_started_at.clone(), @@ -1523,6 +1533,7 @@ pub async fn run_prompt_task( agent.acp.observer_handle(), agent.acp.observer_agent_index(), observer_channel_id, + observer_thread_head_id.clone(), turn_id.clone(), ); @@ -1544,6 +1555,7 @@ pub async fn run_prompt_task( agent.acp.observer_agent_index(), observer::context_for_turn( observer_channel_id, + observer_thread_head_id.clone(), None, turn_id.clone(), turn_started_at.clone(), @@ -1814,6 +1826,7 @@ pub async fn run_prompt_task( }; agent.acp.set_observer_context(observer::context_for_turn( observer_channel_id, + observer_thread_head_id, Some(session_id.clone()), turn_id.clone(), turn_started_at, @@ -3979,6 +3992,7 @@ struct TurnCompletionGuard { observer: Option, agent_index: Option, channel_id: Option, + thread_head_id: Option, turn_id: String, } @@ -3987,12 +4001,14 @@ impl TurnCompletionGuard { observer: Option, agent_index: Option, channel_id: Option, + thread_head_id: Option, turn_id: String, ) -> Self { Self { observer, agent_index, channel_id, + thread_head_id, turn_id, } } @@ -4001,7 +4017,9 @@ impl TurnCompletionGuard { impl Drop for TurnCompletionGuard { fn drop(&mut self) { if let Some(observer) = self.observer.take() { - let context = observer::context_for(self.channel_id, None, Some(self.turn_id.clone())); + let mut context = + observer::context_for(self.channel_id, None, Some(self.turn_id.clone())); + context.thread_head_id = self.thread_head_id.clone(); observer.emit( "turn_completed", self.agent_index, @@ -6517,6 +6535,39 @@ printf '%s\n' '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}'" } } + #[test] + fn test_observer_thread_scope_uses_triggering_thread_root() { + let channel_id = Uuid::new_v4(); + let root = "a".repeat(64); + let root_tag = Tag::parse(vec![ + "e".to_string(), + root.clone(), + String::new(), + "root".to_string(), + ]) + .unwrap(); + let event = EventBuilder::new(Kind::Custom(9), "thread reply") + .tags([root_tag]) + .sign_with_keys(&Keys::generate()) + .unwrap(); + let batch = FlushBatch { + channel_id, + events: vec![crate::queue::BatchEvent { + event, + prompt_tag: "test".into(), + received_at: std::time::Instant::now(), + }], + cancelled_events: vec![], + cancel_reason: None, + }; + + assert_eq!( + observer_thread_head_id(Some(&batch)).as_deref(), + Some(root.as_str()) + ); + assert_eq!(observer_thread_head_id(None), None); + } + #[test] fn test_requeue_cancelled_batch_maps_control_signal_to_cancel_reason() { let cases = [ @@ -6745,11 +6796,37 @@ printf '%s\n' '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}'" })) } + #[test] + fn test_completion_guard_preserves_thread_scope() { + let observer = observer::ObserverHandle::in_process(); + { + let _guard = TurnCompletionGuard::new( + Some(observer.clone()), + Some(0), + Some(Uuid::new_v4()), + Some("thread-1".into()), + "turn-1".into(), + ); + } + + let event = observer + .snapshot() + .into_iter() + .find(|event| event.kind == "turn_completed") + .expect("completion frame"); + assert_eq!(event.thread_head_id.as_deref(), Some("thread-1")); + } + #[tokio::test(start_paused = true)] async fn test_liveness_stops_before_completion_frame() { let observer = observer::ObserverHandle::in_process(); - let context = - observer::context_for_turn(None, None, "t-1".into(), "2026-07-14T21:00:00Z".into()); + let context = observer::context_for_turn( + None, + None, + None, + "t-1".into(), + "2026-07-14T21:00:00Z".into(), + ); let completion_context = observer::context_for(None, None, Some("t-1".into())); let completion_observer = observer.clone(); let completion_handle = tokio::spawn(async move { @@ -6802,7 +6879,13 @@ printf '%s\n' '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}'" async fn test_liveness_fires_until_guard_drops() { let observer = observer::ObserverHandle::in_process(); let started_at = "2026-07-14T21:00:00Z".to_string(); - let context = observer::context_for_turn(None, None, "t-1".into(), started_at.clone()); + let context = observer::context_for_turn( + None, + Some("thread-1".into()), + None, + "t-1".into(), + started_at.clone(), + ); let state = open_liveness_state(); let guard = LivenessGuard::new( tokio::spawn(run_turn_liveness( @@ -6832,6 +6915,9 @@ printf '%s\n' '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}'" assert!(pings .iter() .all(|event| event.started_at.as_deref() == Some(&started_at))); + assert!(pings + .iter() + .all(|event| event.thread_head_id.as_deref() == Some("thread-1"))); assert!(pings .iter() .all(|event| { event.payload == serde_json::json!({ "livenessIntervalSecs": 10 }) })); @@ -6853,8 +6939,13 @@ printf '%s\n' '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}'" #[tokio::test(start_paused = true)] async fn test_liveness_backfills_session_id_after_resolution() { let observer = observer::ObserverHandle::in_process(); - let context = - observer::context_for_turn(None, None, "t-1".into(), "2026-07-14T21:00:00Z".into()); + let context = observer::context_for_turn( + None, + None, + None, + "t-1".into(), + "2026-07-14T21:00:00Z".into(), + ); let state = open_liveness_state(); let guard = LivenessGuard::new( tokio::spawn(run_turn_liveness( @@ -6953,8 +7044,13 @@ printf '%s\n' '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}'" #[tokio::test(start_paused = true)] async fn test_liveness_emits_nothing_once_closed_flag_is_set() { let observer = observer::ObserverHandle::in_process(); - let context = - observer::context_for_turn(None, None, "t-1".into(), "2026-07-14T21:00:00Z".into()); + let context = observer::context_for_turn( + None, + None, + None, + "t-1".into(), + "2026-07-14T21:00:00Z".into(), + ); // Set directly, bypassing `LivenessGuard` — isolates the read side of // the contract: the check under the lock must gate the emit on its own. let state = Arc::new(Mutex::new(LivenessState { diff --git a/desktop/src/features/agents/ui/agentSessionTypes.ts b/desktop/src/features/agents/ui/agentSessionTypes.ts index 578f98076..cb73922d9 100644 --- a/desktop/src/features/agents/ui/agentSessionTypes.ts +++ b/desktop/src/features/agents/ui/agentSessionTypes.ts @@ -6,6 +6,7 @@ export type ObserverEvent = { kind: string; agentIndex: number | null; channelId: string | null; + threadHeadId?: string | null; sessionId: string | null; turnId: string | null; startedAt?: string | null; diff --git a/docs/nips/NIP-AO.md b/docs/nips/NIP-AO.md index 36adea048..340506986 100644 --- a/docs/nips/NIP-AO.md +++ b/docs/nips/NIP-AO.md @@ -85,21 +85,23 @@ The `content` field decrypts to an `ObserverEvent` JSON object: "kind": "", "agentIndex": | null, "channelId": "" | null, + "threadHeadId": "" | null, "sessionId": "" | null, "turnId": "" | null, "payload": { ... } } ``` -`seq`, `timestamp`, `kind`, and `payload` are REQUIRED. `agentIndex`, `channelId`, `sessionId`, -and `turnId` are OPTIONAL — they MAY be `null` when the value is not yet known -(e.g., `sessionId` before session establishment). Clients MUST handle `null` values -gracefully. +`seq`, `timestamp`, `kind`, and `payload` are REQUIRED. `agentIndex`, `channelId`, +`threadHeadId`, `sessionId`, and `turnId` are OPTIONAL — they MAY be `null` when +the value is not yet known (e.g., `sessionId` before session establishment). +Clients MUST handle `null` values gracefully. `seq` is monotonically increasing per session (drop detection). `timestamp` is an RFC 3339 datetime string with sub-second precision (e.g., `"2026-04-29T12:00:41.500Z"`). -`agentIndex` identifies the agent in multi-agent scenarios. `sessionId`/`turnId` -correlate frames across a session and turn. `payload` is kind-specific (MAY be `{}`). +`agentIndex` identifies the agent in multi-agent scenarios. `threadHeadId` is the +NIP-10 root event ID for a thread-scoped turn. `sessionId`/`turnId` correlate +frames across a session and turn. `payload` is kind-specific (MAY be `{}`). Unknown `kind` values MUST be ignored. ### Frame Kinds @@ -254,6 +256,7 @@ of decrypted payloads and MUST NOT log it at INFO level or above. "kind": "acp_write", "agentIndex": 0, "channelId": "52a85618-0f8f-4542-94ec-599e6e1c6f2e", + "threadHeadId": "9f0d...c2a1", "sessionId": "a1b2c3d4", "turnId": "e5f6g7h8", "payload": { diff --git a/mobile/lib/features/channels/agent_activity/active_agent_turns.dart b/mobile/lib/features/channels/agent_activity/active_agent_turns.dart index 0e178efa8..f34181f71 100644 --- a/mobile/lib/features/channels/agent_activity/active_agent_turns.dart +++ b/mobile/lib/features/channels/agent_activity/active_agent_turns.dart @@ -21,6 +21,7 @@ enum AgentTurnPhase { working, finished, error } class AgentTurnState { final String agentPubkey; final String channelId; + final String? threadHeadId; final String turnId; final DateTime startedAt; final DateTime lastActivityAt; @@ -33,6 +34,7 @@ class AgentTurnState { const AgentTurnState({ required this.agentPubkey, required this.channelId, + this.threadHeadId, required this.turnId, required this.startedAt, required this.lastActivityAt, @@ -49,6 +51,7 @@ class AgentTurnState { AgentTurnState( agentPubkey: agentPubkey, channelId: channelId, + threadHeadId: threadHeadId, turnId: turnId, startedAt: startedAt, lastActivityAt: at, @@ -64,6 +67,7 @@ class AgentTurnState { }) => AgentTurnState( agentPubkey: agentPubkey, channelId: channelId, + threadHeadId: threadHeadId, turnId: turnId, startedAt: startedAt, lastActivityAt: at, @@ -100,6 +104,7 @@ List reduceAgentTurnStates( turnsById[turnId] = AgentTurnState( agentPubkey: agentPubkey, channelId: channelId, + threadHeadId: frame.threadHeadId, turnId: turnId, startedAt: _safeStartedAt(frame, frameAt), lastActivityAt: frameAt, @@ -134,6 +139,7 @@ List reduceAgentTurnStates( AgentTurnState( agentPubkey: agentPubkey, channelId: channelId, + threadHeadId: frame.threadHeadId, turnId: turnId, startedAt: _safeStartedAt(frame, frameAt), lastActivityAt: frameAt, @@ -148,7 +154,12 @@ List reduceAgentTurnStates( final channelId = frame.channelId; if (channelId == null) continue; final matching = turnsById.values - .where((turn) => turn.channelId == channelId && turn.isWorking) + .where( + (turn) => + turn.channelId == channelId && + turn.threadHeadId == frame.threadHeadId && + turn.isWorking, + ) .fold( null, (latest, turn) => @@ -189,6 +200,7 @@ List reduceAgentTurnStates( turnsById[turnId] = AgentTurnState( agentPubkey: agentPubkey, channelId: channelId, + threadHeadId: frame.threadHeadId, turnId: turnId, startedAt: _safeStartedAt(frame, frameAt), lastActivityAt: frameAt, @@ -222,6 +234,7 @@ AgentTurnState? latestAgentTurnState( Iterable states, { required String agentPubkey, required String channelId, + required String? threadHeadId, String? turnId, }) { final normalizedAgent = agentPubkey.toLowerCase(); @@ -229,6 +242,7 @@ AgentTurnState? latestAgentTurnState( for (final state in states) { if (state.agentPubkey != normalizedAgent || state.channelId != channelId || + state.threadHeadId != threadHeadId || (turnId != null && state.turnId != turnId)) { continue; } diff --git a/mobile/lib/features/channels/agent_activity/composer_agent_activity_indicator.dart b/mobile/lib/features/channels/agent_activity/composer_agent_activity_indicator.dart index 8ab33933f..b4c195d58 100644 --- a/mobile/lib/features/channels/agent_activity/composer_agent_activity_indicator.dart +++ b/mobile/lib/features/channels/agent_activity/composer_agent_activity_indicator.dart @@ -252,6 +252,7 @@ class ComposerAgentActivityIndicator extends HookConsumerWidget { turnStates, agentPubkey: effectiveSelectedAgent, channelId: channelId, + threadHeadId: threadHeadId, turnId: effectiveTurnId, ); final ObserverState? observerState; @@ -358,6 +359,7 @@ class ComposerAgentActivityIndicator extends HookConsumerWidget { turnStates, agentPubkey: pubkey, channelId: channelId, + threadHeadId: threadHeadId, ); selectedAgent.value = pubkey; pinnedTurnId.value = signal?.turnId ?? latestTurn?.turnId; diff --git a/mobile/lib/features/channels/agent_activity/observer_models.dart b/mobile/lib/features/channels/agent_activity/observer_models.dart index 4ad8f5f39..53c478822 100644 --- a/mobile/lib/features/channels/agent_activity/observer_models.dart +++ b/mobile/lib/features/channels/agent_activity/observer_models.dart @@ -14,6 +14,7 @@ class ObserverFrame { final String kind; final int? agentIndex; final String? channelId; + final String? threadHeadId; final String? sessionId; final String? turnId; final String? startedAt; @@ -26,6 +27,7 @@ class ObserverFrame { required this.kind, this.agentIndex, this.channelId, + this.threadHeadId, this.sessionId, this.turnId, this.startedAt, @@ -42,6 +44,7 @@ class ObserverFrame { kind: json['kind'] as String? ?? '', agentIndex: json['agentIndex'] as int?, channelId: json['channelId'] as String?, + threadHeadId: json['threadHeadId'] as String?, sessionId: json['sessionId'] as String?, turnId: json['turnId'] as String?, startedAt: json['startedAt'] as String?, diff --git a/mobile/lib/features/channels/agent_activity/observer_subscription.dart b/mobile/lib/features/channels/agent_activity/observer_subscription.dart index 8a7c94bb9..7cb582ecc 100644 --- a/mobile/lib/features/channels/agent_activity/observer_subscription.dart +++ b/mobile/lib/features/channels/agent_activity/observer_subscription.dart @@ -372,6 +372,7 @@ class ObserverRelayNotifier extends Notifier { kind: frame.kind, agentIndex: frame.agentIndex, channelId: frame.channelId, + threadHeadId: frame.threadHeadId, sessionId: frame.sessionId, turnId: frame.turnId, startedAt: frame.startedAt, diff --git a/mobile/lib/features/channels/agent_activity/working_bots_provider.dart b/mobile/lib/features/channels/agent_activity/working_bots_provider.dart index 1a27d885a..f3cdb6c09 100644 --- a/mobile/lib/features/channels/agent_activity/working_bots_provider.dart +++ b/mobile/lib/features/channels/agent_activity/working_bots_provider.dart @@ -51,9 +51,9 @@ class ComposerActivityState { }); } -/// Unified composer state. Observer activity is authoritative in channels; -/// kind:20002 typing fills gaps and is the only thread-scoped signal because -/// observer frames do not carry a thread id. +/// Unified composer state. Scoped observer activity is authoritative; +/// kind:20002 typing fills gaps, including for legacy observer frames that do +/// not carry a thread id. final composerActivityStateProvider = Provider.autoDispose .family((ref, key) { final currentPubkey = ref.watch(currentPubkeyProvider)?.toLowerCase(); @@ -81,7 +81,10 @@ final composerActivityStateProvider = Provider.autoDispose final composerTurns = ref.watch(composerAgentTurnStatesProvider); final activeByAgent = {}; for (final turn in composerTurns) { - if (turn.channelId != key.channelId) continue; + if (turn.channelId != key.channelId || + turn.threadHeadId != key.threadHeadId) { + continue; + } final existing = activeByAgent[turn.agentPubkey]; if (existing == null || _compareComposerTurnRecency(turn, existing) > 0) { @@ -97,23 +100,19 @@ final composerActivityStateProvider = Provider.autoDispose } final signals = {}; - // A channel can trust observer activity directly. A thread cannot: the - // observer protocol has no thread id, so a thread requires typing first. - if (key.threadHeadId == null) { - for (final entry in activeByAgent.entries) { - if (!channelAgents.contains(entry.key)) continue; - final turn = entry.value; - if (!turn.isWorking && !canView(entry.key)) continue; - signals[entry.key] = WorkingAgentSignal( - pubkey: entry.key, - source: AgentWorkingSource.observer, - canViewActivity: canView(entry.key), - isWorking: turn.isWorking, - phase: turn.phase, - turnId: turn.turnId, - startedAt: turn.startedAt, - ); - } + for (final entry in activeByAgent.entries) { + if (!channelAgents.contains(entry.key)) continue; + final turn = entry.value; + if (!turn.isWorking && !canView(entry.key)) continue; + signals[entry.key] = WorkingAgentSignal( + pubkey: entry.key, + source: AgentWorkingSource.observer, + canViewActivity: canView(entry.key), + isWorking: turn.isWorking, + phase: turn.phase, + turnId: turn.turnId, + startedAt: turn.startedAt, + ); } final humans = []; diff --git a/mobile/test/features/channels/agent_activity/active_agent_turns_test.dart b/mobile/test/features/channels/agent_activity/active_agent_turns_test.dart index a76b3ef73..a2b6b95b2 100644 --- a/mobile/test/features/channels/agent_activity/active_agent_turns_test.dart +++ b/mobile/test/features/channels/agent_activity/active_agent_turns_test.dart @@ -10,6 +10,7 @@ void main() { seq: 1, second: 1, kind: 'turn_started', + threadHeadId: 'thread-1', receivedSecond: 10, payload: { 'triggeringEventIds': ['message-1'], @@ -22,6 +23,7 @@ void main() { expect(turns, hasLength(1)); expect(turns.single.agentPubkey, 'agent-a'); expect(turns.single.phase, AgentTurnPhase.working); + expect(turns.single.threadHeadId, 'thread-1'); expect(turns.single.triggeringEventId, 'message-1'); expect(turns.single.lastActivityAt, DateTime.utc(2026, 8, 16, 12, 0, 20)); }); @@ -204,6 +206,44 @@ void main() { expect(turns.single.errorMessage, 'Process exited'); }); + test('terminal without a turn id stays within its observer thread scope', () { + final turns = reduceAgentTurnStates({ + 'agent-a': [ + _frame( + seq: 1, + second: 1, + kind: 'turn_started', + turnId: 'turn-a', + threadHeadId: 'thread-a', + ), + _frame( + seq: 2, + second: 2, + kind: 'turn_started', + turnId: 'turn-b', + threadHeadId: 'thread-b', + ), + _frame( + seq: 3, + second: 3, + kind: 'agent_panic', + turnId: null, + threadHeadId: 'thread-b', + payload: {'error': 'Process exited'}, + ), + ], + }, now: DateTime.utc(2026, 8, 16, 12, 0, 20)); + + expect( + turns.singleWhere((turn) => turn.turnId == 'turn-a').phase, + AgentTurnPhase.working, + ); + expect( + turns.singleWhere((turn) => turn.turnId == 'turn-b').phase, + AgentTurnPhase.error, + ); + }); + test( 'retains terminal outcomes beside the composer for a bounded window', () { @@ -246,6 +286,7 @@ ObserverFrame _frame({ required String kind, String? turnId = 'turn-1', String channelId = 'channel-1', + String? threadHeadId, int? receivedSecond, String? startedAt, dynamic payload = const {}, @@ -255,6 +296,7 @@ ObserverFrame _frame({ timestamp: DateTime.utc(2026, 8, 16, 12, 0, second).toIso8601String(), kind: kind, channelId: channelId, + threadHeadId: threadHeadId, turnId: turnId, startedAt: startedAt, receivedAt: receivedSecond == null diff --git a/mobile/test/features/channels/agent_activity/observer_subscription_test.dart b/mobile/test/features/channels/agent_activity/observer_subscription_test.dart index 2a775a109..e343aaf75 100644 --- a/mobile/test/features/channels/agent_activity/observer_subscription_test.dart +++ b/mobile/test/features/channels/agent_activity/observer_subscription_test.dart @@ -361,6 +361,7 @@ void main() { final earlierFrame = _observerFrameJson( seq: 1, channelId: channelId, + threadHeadId: 'thread-1', turnId: 'turn-1', ); relaySession.emit( @@ -388,6 +389,7 @@ void main() { isTrue, reason: 'batch receipt order must follow timestamp and sequence', ); + expect(frames[0].threadHeadId, 'thread-1'); final state = container.read(observerSubscriptionProvider(key)); expect(state.connection, ObserverConnectionState.open); @@ -594,12 +596,14 @@ ObserverFrame _turnMessageFrame({ Map _observerFrameJson({ required int seq, required String channelId, + String? threadHeadId, required String turnId, }) => { 'seq': seq, 'timestamp': '2026-04-30T12:00:0$seq.000Z', 'kind': 'turn_started', 'channelId': channelId, + 'threadHeadId': ?threadHeadId, 'turnId': turnId, 'payload': { 'triggeringEventIds': ['$seq'], diff --git a/mobile/test/features/channels/agent_activity/working_bots_provider_test.dart b/mobile/test/features/channels/agent_activity/working_bots_provider_test.dart index 7d42c7848..7d37c95be 100644 --- a/mobile/test/features/channels/agent_activity/working_bots_provider_test.dart +++ b/mobile/test/features/channels/agent_activity/working_bots_provider_test.dart @@ -164,6 +164,57 @@ void main() { expect(signal.startedAt, isNull); }); + test( + 'keeps a scoped terminal outcome reachable after thread typing stops', + () { + final failedTurn = _turn( + 'agent-a', + phase: AgentTurnPhase.error, + threadHeadId: 'thread-1', + ); + final container = ProviderContainer( + overrides: [ + currentPubkeyProvider.overrideWith((ref) => 'owner'), + channelMembersProvider( + _channelId, + ).overrideWith((ref) async => const []), + channelTypingProvider( + _channelId, + ).overrideWith(() => _FakeTypingNotifier(const [])), + agentMentionPubkeysProvider( + _channelId, + ).overrideWith((ref) => const {'agent-a'}), + agentOwnersProvider.overrideWithValue( + const AsyncData({'agent-a': 'owner'}), + ), + userCacheProvider.overrideWith(_FakeUserCacheNotifier.new), + observerRelayProvider.overrideWith( + () => _FakeObserverRelayNotifier({ + 'agent-a': [_observerFrame('agent-a')], + }), + ), + composerAgentTurnStatesProvider.overrideWithValue([failedTurn]), + ], + ); + addTearDown(container.dispose); + + final signal = container + .read( + composerActivityStateProvider(( + channelId: _channelId, + threadHeadId: 'thread-1', + )), + ) + .agents + .single; + + expect(signal.source, AgentWorkingSource.observer); + expect(signal.isWorking, isFalse); + expect(signal.phase, AgentTurnPhase.error); + expect(signal.turnId, failedTurn.turnId); + }, + ); + test('keeps a recent owned error reachable after typing stops', () { final failedTurn = _turn('agent-a', phase: AgentTurnPhase.error); final container = ProviderContainer( @@ -267,11 +318,13 @@ AgentTurnState _turn( String pubkey, { AgentTurnPhase phase = AgentTurnPhase.working, String? turnId, + String? threadHeadId, DateTime? startedAt, DateTime? lastActivityAt, }) => AgentTurnState( agentPubkey: pubkey, channelId: _channelId, + threadHeadId: threadHeadId, turnId: turnId ?? 'turn-$pubkey', startedAt: startedAt ?? DateTime.utc(2026, 8, 16, 12), lastActivityAt: lastActivityAt ?? DateTime.utc(2026, 8, 16, 12),