Scope mobile agent activity to threads

Signed-off-by: kenny lopez <klopez4212@gmail.com>
This commit is contained in:
kenny lopez
2026-08-17 15:48:43 +01:00
parent 64ea2d430a
commit 7da997a567
14 changed files with 323 additions and 49 deletions
+5
View File
@@ -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<ObserverHandle> {
self.observer.clone()
+55 -13
View File
@@ -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<i64>| {
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,
+9
View File
@@ -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<String>,
/// NIP-10 thread root for the current turn, when thread-scoped.
pub thread_head_id: Option<String>,
/// ACP session ID associated with the current turn, once known.
pub session_id: Option<String>,
/// Local UUID for one prompt turn.
@@ -67,6 +69,9 @@ pub struct ObserverEvent {
pub agent_index: Option<usize>,
/// Buzz channel UUID for channel-scoped events.
pub channel_id: Option<String>,
/// NIP-10 thread root for thread-scoped events.
#[serde(skip_serializing_if = "Option::is_none")]
pub thread_head_id: Option<String>,
/// ACP session ID when known.
pub session_id: Option<String>,
/// 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<uuid::Uuid>,
thread_head_id: Option<String>,
session_id: Option<String>,
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),
+104 -8
View File
@@ -59,6 +59,8 @@ pub struct SuccessfulSteerDelivery {
pub struct TaskMeta {
pub agent_index: usize,
pub channel_id: Option<Uuid>,
/// NIP-10 thread root for panic recovery telemetry.
pub thread_head_id: Option<String>,
/// 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<String> {
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<observer::ObserverHandle>,
agent_index: Option<usize>,
channel_id: Option<uuid::Uuid>,
thread_head_id: Option<String>,
turn_id: String,
}
@@ -3987,12 +4001,14 @@ impl TurnCompletionGuard {
observer: Option<observer::ObserverHandle>,
agent_index: Option<usize>,
channel_id: Option<uuid::Uuid>,
thread_head_id: Option<String>,
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 {
@@ -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;
+9 -6
View File
@@ -85,21 +85,23 @@ The `content` field decrypts to an `ObserverEvent` JSON object:
"kind": "<frame_kind>",
"agentIndex": <integer> | null,
"channelId": "<channel_uuid>" | null,
"threadHeadId": "<nip10_root_event_id>" | null,
"sessionId": "<session_id>" | null,
"turnId": "<turn_id>" | 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": {
@@ -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<AgentTurnState> reduceAgentTurnStates(
turnsById[turnId] = AgentTurnState(
agentPubkey: agentPubkey,
channelId: channelId,
threadHeadId: frame.threadHeadId,
turnId: turnId,
startedAt: _safeStartedAt(frame, frameAt),
lastActivityAt: frameAt,
@@ -134,6 +139,7 @@ List<AgentTurnState> reduceAgentTurnStates(
AgentTurnState(
agentPubkey: agentPubkey,
channelId: channelId,
threadHeadId: frame.threadHeadId,
turnId: turnId,
startedAt: _safeStartedAt(frame, frameAt),
lastActivityAt: frameAt,
@@ -148,7 +154,12 @@ List<AgentTurnState> 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<AgentTurnState?>(
null,
(latest, turn) =>
@@ -189,6 +200,7 @@ List<AgentTurnState> 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<AgentTurnState> 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;
}
@@ -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;
@@ -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?,
@@ -372,6 +372,7 @@ class ObserverRelayNotifier extends Notifier<ObserverRelayState> {
kind: frame.kind,
agentIndex: frame.agentIndex,
channelId: frame.channelId,
threadHeadId: frame.threadHeadId,
sessionId: frame.sessionId,
turnId: frame.turnId,
startedAt: frame.startedAt,
@@ -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<ComposerActivityState, ComposerActivityKey>((ref, key) {
final currentPubkey = ref.watch(currentPubkeyProvider)?.toLowerCase();
@@ -81,7 +81,10 @@ final composerActivityStateProvider = Provider.autoDispose
final composerTurns = ref.watch(composerAgentTurnStatesProvider);
final activeByAgent = <String, AgentTurnState>{};
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 = <String, WorkingAgentSignal>{};
// 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 = <TypingEntry>[];
@@ -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 <String, dynamic>{},
@@ -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
@@ -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<String, dynamic> _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'],
@@ -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 <ChannelMember>[]),
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),