mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-09-07 08:27:12 +00:00
fix(events): expose stage interrupt events in transcript
Emit agent.interrupt.injected when run interrupts reach active agent sessions, persist the stage/session fields, and refresh/render those rows in the Transcript tab.
This commit is contained in:
parent
e47e738e8b
commit
893903ee04
11 changed files with 231 additions and 12 deletions
|
|
@ -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", () => {
|
||||
|
|
|
|||
|
|
@ -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 [];
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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, {
|
||||
|
|
|
|||
|
|
@ -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!({
|
||||
|
|
|
|||
|
|
@ -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."
|
||||
|
|
|
|||
|
|
@ -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!({
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -571,6 +571,15 @@ pub enum Event {
|
|||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
parent_session_id: Option<String>,
|
||||
},
|
||||
/// 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<Principal>,
|
||||
},
|
||||
/// 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)");
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<SteeringHub>, Arc<Mutex<Vec<RunEvent>>>) {
|
||||
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::<Vec<_>>();
|
||||
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::<Vec<_>>();
|
||||
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]
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue