diff --git a/apps/fabro-web/app/lib/run-events.test.tsx b/apps/fabro-web/app/lib/run-events.test.tsx index 1c48947e0..a440fe5fe 100644 --- a/apps/fabro-web/app/lib/run-events.test.tsx +++ b/apps/fabro-web/app/lib/run-events.test.tsx @@ -68,6 +68,13 @@ describe("queryKeysForRunEvent", () => { queryKeys.runs.stageEvents("run-1", "agent@1"), ]); }); + + test("stage-scoped interrupt injection invalidates run events and stage events", () => { + expect(queryKeysForRunEvent("run-1", "agent.interrupt.injected", "nap@1")).toEqual([ + queryKeys.runs.events("run-1", 1000), + queryKeys.runs.stageEvents("run-1", "nap@1"), + ]); + }); }); describe("subscribeToRunEvents", () => { diff --git a/apps/fabro-web/app/lib/run-events.ts b/apps/fabro-web/app/lib/run-events.ts index 71bf73479..2d94f48ca 100644 --- a/apps/fabro-web/app/lib/run-events.ts +++ b/apps/fabro-web/app/lib/run-events.ts @@ -67,6 +67,8 @@ export const STAGE_ACTIVITY_EVENT_TYPES = [ "agent.message", "agent.tool.started", "agent.tool.completed", + "agent.steering.injected", + "agent.interrupt.injected", "command.started", "command.completed", ] as const; @@ -82,6 +84,7 @@ const STEERING_EVENTS = new Set([ "run.interrupt", "run.steer", "agent.steering.injected", + "agent.interrupt.injected", "agent.session.activated", "agent.session.deactivated", "agent.steer.buffered", @@ -134,10 +137,6 @@ export function queryKeysForRunEvent( return keys; } - if (STAGE_ACTIVITY_EVENTS.has(event)) { - return stageId ? [queryKeys.runs.stageEvents(runId, stageId)] : []; - } - if (STEERING_EVENTS.has(event)) { const keys: SseKey[] = [queryKeys.runs.events(runId, 1000)]; if (stageId) { @@ -146,6 +145,10 @@ export function queryKeysForRunEvent( return keys; } + if (STAGE_ACTIVITY_EVENTS.has(event)) { + return stageId ? [queryKeys.runs.stageEvents(runId, stageId)] : []; + } + return []; } diff --git a/apps/fabro-web/app/routes/run-stages.test.ts b/apps/fabro-web/app/routes/run-stages.test.ts index b82264d6c..3ca8e54f7 100644 --- a/apps/fabro-web/app/routes/run-stages.test.ts +++ b/apps/fabro-web/app/routes/run-stages.test.ts @@ -168,6 +168,64 @@ describe("eventsToActivity", () => { } }); + test("renders injected steering as a transcript turn for the matching stage", () => { + const events: EventEnvelope[] = [ + envelope(1, { + event: "run.steer", + properties: { text: "say hello" }, + }), + envelope(2, { + event: "agent.steering.injected", + stage_id: "nap@1", + node_id: "nap", + properties: { text: "say hello", visit: 1 }, + }), + envelope(3, { + event: "agent.steering.injected", + stage_id: "other@1", + node_id: "other", + properties: { text: "wrong stage", visit: 1 }, + }), + ]; + + expect(eventsToActivity(events, "nap@1")).toEqual([ + { + kind: "steer", + ts: "2026-04-09T12:00:00Z", + content: "say hello", + }, + ]); + }); + + test("renders injected interrupt as a transcript turn for the matching stage", () => { + const events: EventEnvelope[] = [ + envelope(1, { + event: "run.interrupt", + properties: {}, + }), + envelope(2, { + event: "agent.interrupt.injected", + stage_id: "nap@1", + node_id: "nap", + properties: { visit: 1 }, + }), + envelope(3, { + event: "agent.interrupt.injected", + stage_id: "other@1", + node_id: "other", + properties: { visit: 1 }, + }), + ]; + + expect(eventsToActivity(events, "nap@1")).toEqual([ + { + kind: "interrupt", + ts: "2026-04-09T12:00:00Z", + content: "Agent interrupted", + }, + ]); + }); + test("extractStageModel pulls model from agent.session.activated, ignoring other stages", () => { const events: EventEnvelope[] = [ envelope(1, { diff --git a/lib/crates/fabro-api/tests/run_event_round_trip.rs b/lib/crates/fabro-api/tests/run_event_round_trip.rs index 1ac43d302..b9d150fe3 100644 --- a/lib/crates/fabro-api/tests/run_event_round_trip.rs +++ b/lib/crates/fabro-api/tests/run_event_round_trip.rs @@ -78,6 +78,26 @@ fn run_event_round_trips_run_steer() { assert_run_event_round_trip(value); } +#[test] +fn run_event_round_trips_agent_interrupt_injected() { + let value = json!({ + "id": "evt_interrupt_injected", + "ts": "2026-04-29T12:00:00Z", + "run_id": fixtures::RUN_1, + "event": "agent.interrupt.injected", + "node_id": "code", + "node_label": "code", + "stage_id": "code@2", + "session_id": "ses_1", + "actor": { "kind": "system", "system_kind": "engine" }, + "properties": { + "visit": 2 + } + }); + + assert_run_event_round_trip(value); +} + #[test] fn run_event_round_trips_stage_started() { let value = json!({ diff --git a/lib/crates/fabro-types/src/run_event/agent.rs b/lib/crates/fabro-types/src/run_event/agent.rs index da5856167..73f5b85f9 100644 --- a/lib/crates/fabro-types/src/run_event/agent.rs +++ b/lib/crates/fabro-types/src/run_event/agent.rs @@ -109,6 +109,11 @@ pub struct AgentSteeringInjectedProps { pub visit: u32, } +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct AgentInterruptInjectedProps { + pub visit: u32, +} + #[allow( clippy::empty_structs_with_brackets, reason = "This type must serialize as {} rather than null." diff --git a/lib/crates/fabro-types/src/run_event/mod.rs b/lib/crates/fabro-types/src/run_event/mod.rs index 339c8c22d..91bac42d7 100644 --- a/lib/crates/fabro-types/src/run_event/mod.rs +++ b/lib/crates/fabro-types/src/run_event/mod.rs @@ -178,6 +178,8 @@ pub enum EventBody { AgentTurnLimitReached(AgentTurnLimitReachedProps), #[serde(rename = "agent.steering.injected")] AgentSteeringInjected(AgentSteeringInjectedProps), + #[serde(rename = "agent.interrupt.injected")] + AgentInterruptInjected(AgentInterruptInjectedProps), #[serde(rename = "agent.steer.buffered")] AgentSteerBuffered(AgentSteerBufferedProps), #[serde(rename = "agent.steer.dropped")] @@ -412,6 +414,7 @@ impl EventBody { Self::AgentLoopDetected(_) => "agent.loop.detected", Self::AgentTurnLimitReached(_) => "agent.turn.limit", Self::AgentSteeringInjected(_) => "agent.steering.injected", + Self::AgentInterruptInjected(_) => "agent.interrupt.injected", Self::AgentSteerBuffered(_) => "agent.steer.buffered", Self::AgentSteerDropped(_) => "agent.steer.dropped", Self::AgentCompactionStarted(_) => "agent.compaction.started", @@ -552,6 +555,7 @@ fn is_known_event_name(event: &str) -> bool { | "agent.loop.detected" | "agent.turn.limit" | "agent.steering.injected" + | "agent.interrupt.injected" | "agent.steer.buffered" | "agent.steer.dropped" | "agent.compaction.started" @@ -964,6 +968,29 @@ mod tests { assert_eq!(parsed.to_value().unwrap(), line); } + #[test] + fn agent_interrupt_injected_round_trips_with_stage_session_and_actor() { + let line = json!({ + "id": "evt_interrupt_injected", + "ts": "2026-04-04T12:00:00Z", + "run_id": fixtures::RUN_1, + "event": "agent.interrupt.injected", + "node_id": "code", + "node_label": "code", + "stage_id": "code@2", + "session_id": "ses_1", + "actor": { "kind": "system", "system_kind": "engine" }, + "properties": { "visit": 2 } + }); + + let parsed = RunEvent::from_value(line.clone()).unwrap(); + assert!(matches!( + &parsed.body, + EventBody::AgentInterruptInjected(props) if props.visit == 2 + )); + assert_eq!(parsed.to_value().unwrap(), line); + } + #[test] fn run_interrupt_then_steer_is_not_a_known_persisted_event() { let line = json!({ diff --git a/lib/crates/fabro-workflow/src/event/convert.rs b/lib/crates/fabro-workflow/src/event/convert.rs index 0fb72c9aa..86348dbfb 100644 --- a/lib/crates/fabro-workflow/src/event/convert.rs +++ b/lib/crates/fabro-workflow/src/event/convert.rs @@ -1034,6 +1034,11 @@ fn event_body_from_event(event: &Event) -> EventBody { Event::AgentSessionEnded { .. } => { EventBody::AgentSessionEnded(fabro_types::AgentSessionEndedProps {}) } + Event::AgentInterruptInjected { visit, .. } => { + EventBody::AgentInterruptInjected(fabro_types::AgentInterruptInjectedProps { + visit: *visit, + }) + } Event::AgentSteerBuffered { .. } => { EventBody::AgentSteerBuffered(fabro_types::AgentSteerBufferedProps::default()) } @@ -1584,6 +1589,30 @@ mod tests { ); } + #[test] + fn agent_interrupt_injected_populates_stage_session_and_actor() { + let actor = Principal::System { + system_kind: SystemActorKind::Engine, + }; + let stored = to_run_event(&fixtures::RUN_1, &Event::AgentInterruptInjected { + node_id: "code".to_string(), + visit: 3, + session_id: "ses_1".to_string(), + actor: Some(actor.clone()), + }); + + assert_eq!(stored.event_name(), "agent.interrupt.injected"); + assert_eq!(stored.node_id.as_deref(), Some("code")); + assert_eq!(stored.node_label.as_deref(), Some("code")); + assert_eq!(stored.stage_id, Some(StageId::new("code", 3))); + assert_eq!(stored.session_id.as_deref(), Some("ses_1")); + assert_eq!(stored.actor, Some(actor)); + match stored.body { + EventBody::AgentInterruptInjected(props) => assert_eq!(props.visit, 3), + other => panic!("unexpected body: {other:?}"), + } + } + #[test] fn stage_scope_populates_stage_id_on_non_stage_events() { // Events tied to a concrete stage execution but lacking scope in their diff --git a/lib/crates/fabro-workflow/src/event/events.rs b/lib/crates/fabro-workflow/src/event/events.rs index 16f43f710..cf31c0df6 100644 --- a/lib/crates/fabro-workflow/src/event/events.rs +++ b/lib/crates/fabro-workflow/src/event/events.rs @@ -571,6 +571,15 @@ pub enum Event { #[serde(default, skip_serializing_if = "Option::is_none")] parent_session_id: Option, }, + /// A run-level interrupt was delivered to a concrete API-mode agent + /// session/stage. + AgentInterruptInjected { + node_id: String, + visit: u32, + session_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + actor: Option, + }, /// A steer arrived with no active session and was parked in the run-wide /// pending buffer. The actor (steer author) is lifted to top-level. AgentSteerBuffered { @@ -1355,6 +1364,14 @@ impl Event { Self::AgentSessionEnded { session_id, .. } => { debug!(session_id, "Agent session ended"); } + Self::AgentInterruptInjected { + node_id, + visit, + session_id, + .. + } => { + debug!(node_id, visit, session_id, "Agent interrupt injected"); + } Self::AgentSteerBuffered { .. } => { debug!("Steer buffered (no active session)"); } diff --git a/lib/crates/fabro-workflow/src/event/names.rs b/lib/crates/fabro-workflow/src/event/names.rs index f91290fa2..69ae5ef63 100644 --- a/lib/crates/fabro-workflow/src/event/names.rs +++ b/lib/crates/fabro-workflow/src/event/names.rs @@ -122,6 +122,7 @@ pub fn event_name(event: &Event) -> &'static str { Event::AgentSessionActivated { .. } => "agent.session.activated", Event::AgentSessionDeactivated { .. } => "agent.session.deactivated", Event::AgentSessionEnded { .. } => "agent.session.ended", + Event::AgentInterruptInjected { .. } => "agent.interrupt.injected", Event::AgentSteerBuffered { .. } => "agent.steer.buffered", Event::AgentSteerDropped { .. } => "agent.steer.dropped", Event::AgentCliCancelled { .. } => "agent.cli.cancelled", diff --git a/lib/crates/fabro-workflow/src/event/stored_fields.rs b/lib/crates/fabro-workflow/src/event/stored_fields.rs index 92527aee3..59a68248b 100644 --- a/lib/crates/fabro-workflow/src/event/stored_fields.rs +++ b/lib/crates/fabro-workflow/src/event/stored_fields.rs @@ -156,6 +156,23 @@ fn stored_event_fields_for_variant(event: &Event) -> StoredEventFields { ..StoredEventFields::default() } } + Event::AgentInterruptInjected { + node_id, + visit, + session_id, + actor, + } => { + let node_id_str = node_id.clone(); + let node_label = default_node_label(Some(&node_id_str), None); + StoredEventFields { + session_id: Some(session_id.clone()), + node_id: Some(node_id_str.clone()), + node_label, + stage_id: Some(StageId::new(node_id_str, *visit)), + actor: actor.clone(), + ..StoredEventFields::default() + } + } Event::AgentSteerDropped { actor, node_id, diff --git a/lib/crates/fabro-workflow/src/steering_hub.rs b/lib/crates/fabro-workflow/src/steering_hub.rs index 7ef2b3c97..3bb22ccb3 100644 --- a/lib/crates/fabro-workflow/src/steering_hub.rs +++ b/lib/crates/fabro-workflow/src/steering_hub.rs @@ -226,8 +226,14 @@ impl SteeringHub { self.emitter.emit(&Event::RunInterrupt { actor: actor.cloned(), }); - for entry in active.values() { + for (stage_id, entry) in active.iter() { entry.handle.interrupt(actor.cloned()); + self.emitter.emit(&Event::AgentInterruptInjected { + node_id: stage_id.node_id().to_string(), + visit: stage_id.visit(), + session_id: entry.session_id.clone(), + actor: actor.cloned(), + }); } } @@ -261,6 +267,12 @@ impl SteeringHub { visit: Some(stage_id.visit()), }); } + self.emitter.emit(&Event::AgentInterruptInjected { + node_id: stage_id.node_id().to_string(), + visit: stage_id.visit(), + session_id: entry.session_id.clone(), + actor: actor.cloned(), + }); } } @@ -312,7 +324,7 @@ mod tests { use std::sync::{Arc, Mutex}; use fabro_agent::SessionControlHandle; - use fabro_types::{Principal, RunId, StageId, SystemActorKind}; + use fabro_types::{Principal, RunEvent, RunId, StageId, SystemActorKind}; use super::SteeringHub; use crate::event::Emitter; @@ -330,6 +342,16 @@ mod tests { (Arc::new(SteeringHub::new(emitter)), names) } + fn hub_with_events() -> (Arc, Arc>>) { + let emitter = Arc::new(Emitter::new(RunId::new())); + let events = Arc::new(Mutex::new(Vec::new())); + let events_for_listener = Arc::clone(&events); + emitter.on_event(move |event| { + events_for_listener.lock().unwrap().push(event.clone()); + }); + (Arc::new(SteeringHub::new(emitter)), events) + } + #[test] fn deliver_with_no_active_buffers_message() { let (hub, names) = hub_with_event_names(); @@ -461,7 +483,7 @@ mod tests { #[test] fn pure_interrupt_marks_active_sessions_waiting_without_queueing_text() { - let (hub, names) = hub_with_event_names(); + let (hub, events) = hub_with_events(); let stage = StageId::new("a", 1); let handle = SessionControlHandle::new(); assert!(hub.attach_handle(&stage, "session-a", &handle)); @@ -472,15 +494,23 @@ mod tests { assert!(handle.is_waiting_for_steer()); assert_eq!(handle.queue_len(), 0); assert_eq!(hub.pending_len(), 0); - assert_eq!(names.lock().unwrap().as_slice(), [ + let events = events.lock().unwrap(); + let names = events.iter().map(RunEvent::event_name).collect::>(); + assert_eq!(names, [ "run.interrupt", - "run.interrupt" + "agent.interrupt.injected", + "run.interrupt", + "agent.interrupt.injected", ]); + assert_eq!(events[1].stage_id, Some(stage.clone())); + assert_eq!(events[1].session_id.as_deref(), Some("session-a")); + assert_eq!(events[3].stage_id, Some(stage)); + assert_eq!(events[3].session_id.as_deref(), Some("session-a")); } #[test] fn interrupt_then_steer_cancels_and_queues_text() { - let (hub, names) = hub_with_event_names(); + let (hub, events) = hub_with_events(); let stage = StageId::new("a", 1); let handle = SessionControlHandle::new(); assert!(hub.attach_handle(&stage, "session-a", &handle)); @@ -490,10 +520,15 @@ mod tests { assert!(!handle.is_waiting_for_steer()); assert_eq!(handle.queue_len(), 1); assert_eq!(hub.pending_len(), 0); - assert_eq!(names.lock().unwrap().as_slice(), [ + let events = events.lock().unwrap(); + let names = events.iter().map(RunEvent::event_name).collect::>(); + assert_eq!(names, [ "run.interrupt", - "run.steer" + "run.steer", + "agent.interrupt.injected", ]); + assert_eq!(events[2].stage_id, Some(stage)); + assert_eq!(events[2].session_id.as_deref(), Some("session-a")); } #[test]