mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
Merge pull request #694 from fabro-sh/feat/reusable-subagent-sessions
Reuse completed subagent sessions
This commit is contained in:
commit
f212594875
14 changed files with 1284 additions and 450 deletions
|
|
@ -189,10 +189,13 @@ pub(super) enum ProgressEvent {
|
|||
LlmRequestFinished {
|
||||
stage_node_id: String,
|
||||
},
|
||||
SubagentSpawned {
|
||||
/// A subagent started work. Generation 1 is the spawn; later generations
|
||||
/// are further turns in the same child session.
|
||||
SubagentStarted {
|
||||
stage_node_id: String,
|
||||
agent_id: String,
|
||||
task: String,
|
||||
generation: u64,
|
||||
},
|
||||
SubagentCompleted {
|
||||
stage_node_id: String,
|
||||
|
|
@ -426,10 +429,17 @@ pub(super) fn from_run_event(stored: &RunEvent) -> Option<ProgressEvent> {
|
|||
stage_node_id: node_id,
|
||||
})
|
||||
}
|
||||
EventBody::AgentSubSpawned(props) => Some(ProgressEvent::SubagentSpawned {
|
||||
EventBody::AgentSubSpawned(props) => Some(ProgressEvent::SubagentStarted {
|
||||
stage_node_id: node_id,
|
||||
agent_id: props.agent_id.clone(),
|
||||
task: props.task.clone(),
|
||||
generation: props.generation,
|
||||
}),
|
||||
EventBody::AgentSubTurnStarted(props) => Some(ProgressEvent::SubagentStarted {
|
||||
stage_node_id: node_id,
|
||||
agent_id: props.agent_id.clone(),
|
||||
task: props.task.clone(),
|
||||
generation: props.generation,
|
||||
}),
|
||||
EventBody::AgentSubCompleted(props) => Some(ProgressEvent::SubagentCompleted {
|
||||
stage_node_id: node_id,
|
||||
|
|
|
|||
|
|
@ -363,13 +363,19 @@ impl ProgressUI {
|
|||
ProgressEvent::LlmRequestFinished { stage_node_id } => {
|
||||
self.stage.on_llm_request_finished(&stage_node_id);
|
||||
}
|
||||
ProgressEvent::SubagentSpawned {
|
||||
ProgressEvent::SubagentStarted {
|
||||
stage_node_id,
|
||||
agent_id,
|
||||
task,
|
||||
generation,
|
||||
} => {
|
||||
self.stage
|
||||
.on_subagent_spawned(renderer, &stage_node_id, &agent_id, &task);
|
||||
self.stage.on_subagent_started(
|
||||
renderer,
|
||||
&stage_node_id,
|
||||
&agent_id,
|
||||
&task,
|
||||
generation,
|
||||
);
|
||||
}
|
||||
ProgressEvent::SubagentCompleted {
|
||||
stage_node_id,
|
||||
|
|
@ -970,13 +976,15 @@ mod tests {
|
|||
},
|
||||
}),
|
||||
agent_event("code", AgentEvent::SubAgentSpawned {
|
||||
agent_id: "a1".into(),
|
||||
depth: 1,
|
||||
task: "review recent changes".into(),
|
||||
agent_id: "a1".into(),
|
||||
depth: 1,
|
||||
task: "review recent changes".into(),
|
||||
generation: 1,
|
||||
}),
|
||||
agent_event("code", AgentEvent::SubAgentCompleted {
|
||||
agent_id: "a1".into(),
|
||||
depth: 1,
|
||||
generation: 1,
|
||||
success: true,
|
||||
turns_used: 3,
|
||||
}),
|
||||
|
|
@ -1332,9 +1340,10 @@ mod tests {
|
|||
emit(
|
||||
&mut ui,
|
||||
agent_event("code", AgentEvent::SubAgentSpawned {
|
||||
agent_id: "a1".into(),
|
||||
depth: 1,
|
||||
task: "review recent changes".into(),
|
||||
agent_id: "a1".into(),
|
||||
depth: 1,
|
||||
task: "review recent changes".into(),
|
||||
generation: 1,
|
||||
}),
|
||||
);
|
||||
emit(
|
||||
|
|
@ -1342,10 +1351,30 @@ mod tests {
|
|||
agent_event("code", AgentEvent::SubAgentCompleted {
|
||||
agent_id: "a1".into(),
|
||||
depth: 1,
|
||||
generation: 1,
|
||||
success: true,
|
||||
turns_used: 3,
|
||||
}),
|
||||
);
|
||||
emit(
|
||||
&mut ui,
|
||||
agent_event("code", AgentEvent::SubAgentTurnStarted {
|
||||
agent_id: "a1".into(),
|
||||
depth: 1,
|
||||
task: "fix the review findings".into(),
|
||||
generation: 2,
|
||||
}),
|
||||
);
|
||||
emit(
|
||||
&mut ui,
|
||||
agent_event("code", AgentEvent::SubAgentCompleted {
|
||||
agent_id: "a1".into(),
|
||||
depth: 1,
|
||||
generation: 2,
|
||||
success: true,
|
||||
turns_used: 2,
|
||||
}),
|
||||
);
|
||||
emit(&mut ui, Event::SetupStarted { command_count: 1 });
|
||||
emit(&mut ui, Event::SetupCommandCompleted {
|
||||
command: "bun install".into(),
|
||||
|
|
@ -1363,6 +1392,8 @@ mod tests {
|
|||
⚠ retry: gpt-5-mini attempt 2 (busy, delay 1s)
|
||||
▸ subagent[a1] "review recent changes"
|
||||
✓ subagent[a1] (3 turns)
|
||||
↻ subagent[a1] turn 2 "fix the review findings"
|
||||
✓ subagent[a1] (2 turns)
|
||||
✓ [1/1] bun install 2s
|
||||
Setup: 1 command (2s)
|
||||
✓ Code 5s (1 turns, 0 tools, 1.5k toks)
|
||||
|
|
|
|||
|
|
@ -3,7 +3,7 @@ use std::convert::TryFrom;
|
|||
use std::time::Duration;
|
||||
|
||||
use chrono::{DateTime, Utc};
|
||||
use fabro_types::LlmOutputKind;
|
||||
use fabro_types::{INITIAL_SUBAGENT_GENERATION, LlmOutputKind};
|
||||
use fabro_workflow::outcome::{StageOutcome, format_cost};
|
||||
use indicatif::ProgressBar;
|
||||
|
||||
|
|
@ -606,17 +606,26 @@ impl StageDisplay {
|
|||
);
|
||||
}
|
||||
|
||||
pub(super) fn on_subagent_spawned(
|
||||
/// Show a subagent starting work. Generation 1 is the spawn; a later
|
||||
/// generation is another turn in the same child session, so it reads as a
|
||||
/// return to work rather than a new agent.
|
||||
pub(super) fn on_subagent_started(
|
||||
&mut self,
|
||||
renderer: &ProgressRenderer,
|
||||
stage_node_id: &str,
|
||||
agent_id: &str,
|
||||
task: &str,
|
||||
generation: u64,
|
||||
) {
|
||||
if !self.verbose {
|
||||
return;
|
||||
}
|
||||
|
||||
let (glyph, turn) = if generation > INITIAL_SUBAGENT_GENERATION {
|
||||
("\u{21bb}", format!("turn {generation} "))
|
||||
} else {
|
||||
("\u{25b8}", String::new())
|
||||
};
|
||||
self.insert_subagent_line_for_stage(
|
||||
renderer,
|
||||
stage_node_id,
|
||||
|
|
@ -624,7 +633,7 @@ impl StageDisplay {
|
|||
.styles()
|
||||
.dim
|
||||
.apply_to(format!(
|
||||
"\u{25b8} subagent[{agent_id}] \"{}\"",
|
||||
"{glyph} subagent[{agent_id}] {turn}\"{}\"",
|
||||
styles::truncate(task, 50)
|
||||
))
|
||||
.to_string(),
|
||||
|
|
|
|||
|
|
@ -692,8 +692,20 @@ pub async fn run_with_args_and_client_and_catalog(
|
|||
agent_id,
|
||||
depth,
|
||||
task,
|
||||
..
|
||||
generation,
|
||||
}
|
||||
| AgentEvent::SubAgentTurnStarted {
|
||||
agent_id,
|
||||
depth,
|
||||
task,
|
||||
generation,
|
||||
} => {
|
||||
let started =
|
||||
if matches!(event.event, AgentEvent::SubAgentSpawned { .. }) {
|
||||
"spawned"
|
||||
} else {
|
||||
"turn started"
|
||||
};
|
||||
let task_preview = if task.len() > 60 {
|
||||
&task[..task.floor_char_boundary(60)]
|
||||
} else {
|
||||
|
|
@ -702,40 +714,46 @@ pub async fn run_with_args_and_client_and_catalog(
|
|||
eprintln!(
|
||||
" {}",
|
||||
s.dim.apply_to(format!(
|
||||
"{child_prefix}\u{25b6} subagent {agent_id} spawned (depth={depth}) task={task_preview:?}"
|
||||
"{child_prefix}\u{25b6} subagent {agent_id} {started} (depth={depth}, generation={generation}) task={task_preview:?}"
|
||||
)),
|
||||
);
|
||||
}
|
||||
AgentEvent::SubAgentCompleted {
|
||||
agent_id,
|
||||
depth,
|
||||
generation,
|
||||
success,
|
||||
turns_used,
|
||||
} => {
|
||||
eprintln!(
|
||||
" {}",
|
||||
s.dim.apply_to(format!(
|
||||
"{child_prefix}\u{25a0} subagent {agent_id} completed (depth={depth}, success={success}, turns={turns_used})"
|
||||
"{child_prefix}\u{25a0} subagent {agent_id} completed (depth={depth}, generation={generation}, success={success}, turns={turns_used})"
|
||||
)),
|
||||
);
|
||||
}
|
||||
AgentEvent::SubAgentFailed {
|
||||
agent_id,
|
||||
depth,
|
||||
generation,
|
||||
error,
|
||||
} => {
|
||||
eprintln!(
|
||||
" {}",
|
||||
s.red.apply_to(format!(
|
||||
"{child_prefix}\u{2717} subagent {agent_id} failed (depth={depth}): {error}"
|
||||
"{child_prefix}\u{2717} subagent {agent_id} failed (depth={depth}, generation={generation}): {error}"
|
||||
)),
|
||||
);
|
||||
}
|
||||
AgentEvent::SubAgentClosed { agent_id, depth } => {
|
||||
AgentEvent::SubAgentClosed {
|
||||
agent_id,
|
||||
depth,
|
||||
generation,
|
||||
} => {
|
||||
eprintln!(
|
||||
" {}",
|
||||
s.dim.apply_to(format!(
|
||||
"{child_prefix}\u{25a0} subagent {agent_id} closed (depth={depth})"
|
||||
"{child_prefix}\u{25a0} subagent {agent_id} closed (depth={depth}, generation={generation})"
|
||||
)),
|
||||
);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -333,7 +333,7 @@ pub(crate) fn make_task_output_tool(supervisor: SubAgentSupervisor) -> Registere
|
|||
}
|
||||
|
||||
match supervisor.status(task_id) {
|
||||
Some(SubAgentStatus::Finished(result)) => {
|
||||
Some(SubAgentStatus::Finished { result, .. }) => {
|
||||
return finished_output(&supervisor, task_id, result);
|
||||
}
|
||||
Some(SubAgentStatus::Running) if !block => {
|
||||
|
|
@ -383,7 +383,7 @@ pub(crate) fn make_task_stop_tool(supervisor: SubAgentSupervisor) -> RegisteredT
|
|||
RegisteredTool {
|
||||
definition: definition(
|
||||
NativeTool::StopAgent,
|
||||
"Stop a running background agent by task ID.",
|
||||
"Stop a running or completed background agent by task ID.",
|
||||
serde_json::json!({
|
||||
"type": "object",
|
||||
"properties": {
|
||||
|
|
@ -416,7 +416,7 @@ pub(crate) fn make_send_message_tool(supervisor: SubAgentSupervisor) -> Register
|
|||
RegisteredTool {
|
||||
definition: definition(
|
||||
NativeTool::MessageAgent,
|
||||
"Send additional instructions to a running background agent by its task ID.",
|
||||
"Send additional instructions to a background agent by its task ID. A running agent receives them at a safe turn boundary. A completed agent starts another turn in the same session with its existing history.",
|
||||
serde_json::json!({
|
||||
"type": "object",
|
||||
"properties": {
|
||||
|
|
@ -590,11 +590,17 @@ mod tests {
|
|||
assert_schema(&make_task_stop_tool(supervisor.clone()), &["task_id"], &[
|
||||
"task_id",
|
||||
]);
|
||||
assert_schema(
|
||||
&make_send_message_tool(supervisor),
|
||||
&["message", "summary", "to"],
|
||||
&["message", "to"],
|
||||
let send_message = make_send_message_tool(supervisor);
|
||||
assert_schema(&send_message, &["message", "summary", "to"], &[
|
||||
"message", "to",
|
||||
]);
|
||||
assert!(
|
||||
send_message
|
||||
.definition
|
||||
.description
|
||||
.contains("completed agent")
|
||||
);
|
||||
assert!(send_message.definition.description.contains("same session"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
|
|
|||
File diff suppressed because it is too large
Load diff
|
|
@ -352,24 +352,38 @@ pub enum AgentEvent {
|
|||
phase: LlmRetryPhase,
|
||||
},
|
||||
SubAgentSpawned {
|
||||
agent_id: String,
|
||||
depth: usize,
|
||||
task: String,
|
||||
agent_id: String,
|
||||
depth: usize,
|
||||
task: String,
|
||||
#[serde(default = "fabro_types::initial_subagent_generation")]
|
||||
generation: u64,
|
||||
},
|
||||
SubAgentTurnStarted {
|
||||
agent_id: String,
|
||||
depth: usize,
|
||||
task: String,
|
||||
generation: u64,
|
||||
},
|
||||
SubAgentCompleted {
|
||||
agent_id: String,
|
||||
depth: usize,
|
||||
#[serde(default = "fabro_types::initial_subagent_generation")]
|
||||
generation: u64,
|
||||
success: bool,
|
||||
turns_used: usize,
|
||||
},
|
||||
SubAgentFailed {
|
||||
agent_id: String,
|
||||
depth: usize,
|
||||
error: Error,
|
||||
agent_id: String,
|
||||
depth: usize,
|
||||
#[serde(default = "fabro_types::initial_subagent_generation")]
|
||||
generation: u64,
|
||||
error: Error,
|
||||
},
|
||||
SubAgentClosed {
|
||||
agent_id: String,
|
||||
depth: usize,
|
||||
agent_id: String,
|
||||
depth: usize,
|
||||
#[serde(default = "fabro_types::initial_subagent_generation")]
|
||||
generation: u64,
|
||||
},
|
||||
McpServerReady {
|
||||
server_name: String,
|
||||
|
|
@ -586,35 +600,57 @@ impl AgentEvent {
|
|||
agent_id,
|
||||
depth,
|
||||
task,
|
||||
generation,
|
||||
} => {
|
||||
debug!(session_id, agent_id, depth, task, "Sub-agent spawned");
|
||||
debug!(
|
||||
session_id,
|
||||
agent_id, depth, generation, task, "Sub-agent spawned"
|
||||
);
|
||||
}
|
||||
Self::SubAgentTurnStarted {
|
||||
agent_id,
|
||||
depth,
|
||||
task,
|
||||
generation,
|
||||
} => {
|
||||
debug!(
|
||||
session_id,
|
||||
agent_id, depth, generation, task, "Sub-agent turn started"
|
||||
);
|
||||
}
|
||||
Self::SubAgentCompleted {
|
||||
agent_id,
|
||||
depth,
|
||||
generation,
|
||||
success,
|
||||
turns_used,
|
||||
} => {
|
||||
debug!(
|
||||
session_id,
|
||||
agent_id, depth, success, turns_used, "Sub-agent completed"
|
||||
agent_id, depth, generation, success, turns_used, "Sub-agent completed"
|
||||
);
|
||||
}
|
||||
Self::SubAgentFailed {
|
||||
agent_id,
|
||||
depth,
|
||||
generation,
|
||||
error,
|
||||
} => {
|
||||
warn!(
|
||||
session_id,
|
||||
agent_id,
|
||||
depth,
|
||||
generation,
|
||||
error = %error,
|
||||
"Sub-agent failed"
|
||||
);
|
||||
}
|
||||
Self::SubAgentClosed { agent_id, depth } => {
|
||||
debug!(session_id, agent_id, depth, "Sub-agent closed");
|
||||
Self::SubAgentClosed {
|
||||
agent_id,
|
||||
depth,
|
||||
generation,
|
||||
} => {
|
||||
debug!(session_id, agent_id, depth, generation, "Sub-agent closed");
|
||||
}
|
||||
Self::McpServerReady {
|
||||
server_name,
|
||||
|
|
@ -765,9 +801,10 @@ mod tests {
|
|||
#[test]
|
||||
fn subagent_spawned_constructible() {
|
||||
let event = AgentEvent::SubAgentSpawned {
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 1,
|
||||
task: "list files".into(),
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 1,
|
||||
task: "list files".into(),
|
||||
generation: 1,
|
||||
};
|
||||
assert!(matches!(event, AgentEvent::SubAgentSpawned {
|
||||
depth: 1,
|
||||
|
|
@ -780,6 +817,7 @@ mod tests {
|
|||
let event = AgentEvent::SubAgentCompleted {
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 1,
|
||||
generation: 1,
|
||||
success: true,
|
||||
turns_used: 5,
|
||||
};
|
||||
|
|
@ -793,9 +831,10 @@ mod tests {
|
|||
#[test]
|
||||
fn subagent_failed_constructible() {
|
||||
let event = AgentEvent::SubAgentFailed {
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
error: Error::ToolExecution("timeout".into()),
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
generation: 1,
|
||||
error: Error::ToolExecution("timeout".into()),
|
||||
};
|
||||
assert!(matches!(event, AgentEvent::SubAgentFailed { depth: 0, .. }));
|
||||
}
|
||||
|
|
@ -803,8 +842,9 @@ mod tests {
|
|||
#[test]
|
||||
fn subagent_closed_constructible() {
|
||||
let event = AgentEvent::SubAgentClosed {
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 2,
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 2,
|
||||
generation: 1,
|
||||
};
|
||||
assert!(matches!(event, AgentEvent::SubAgentClosed { depth: 2, .. }));
|
||||
}
|
||||
|
|
@ -813,29 +853,52 @@ mod tests {
|
|||
fn subagent_events_serde_round_trip() {
|
||||
let events = vec![
|
||||
AgentEvent::SubAgentSpawned {
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
task: "test".into(),
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
task: "test".into(),
|
||||
generation: 1,
|
||||
},
|
||||
AgentEvent::SubAgentTurnStarted {
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
task: "fix it".into(),
|
||||
generation: 2,
|
||||
},
|
||||
AgentEvent::SubAgentCompleted {
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
generation: 2,
|
||||
success: true,
|
||||
turns_used: 3,
|
||||
},
|
||||
AgentEvent::SubAgentFailed {
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
error: Error::ToolExecution("oops".into()),
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
generation: 2,
|
||||
error: Error::ToolExecution("oops".into()),
|
||||
},
|
||||
AgentEvent::SubAgentClosed {
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
generation: 2,
|
||||
},
|
||||
];
|
||||
let json = serde_json::to_string(&events).unwrap();
|
||||
let deserialized: Vec<AgentEvent> = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(deserialized.len(), 4);
|
||||
assert_eq!(deserialized.len(), 5);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn legacy_subagent_event_defaults_to_the_initial_generation() {
|
||||
let event: AgentEvent = serde_json::from_str(
|
||||
r#"{"SubAgentSpawned":{"agent_id":"sa-1","depth":0,"task":"test"}}"#,
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert!(matches!(event, AgentEvent::SubAgentSpawned {
|
||||
generation: 1,
|
||||
..
|
||||
}));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -1054,9 +1117,10 @@ mod tests {
|
|||
#[test]
|
||||
fn subagent_failed_carries_agent_error() {
|
||||
let event = AgentEvent::SubAgentFailed {
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
error: Error::ToolExecution("cmd failed".into()),
|
||||
agent_id: "sa-1".into(),
|
||||
depth: 0,
|
||||
generation: 1,
|
||||
error: Error::ToolExecution("cmd failed".into()),
|
||||
};
|
||||
let json = serde_json::to_string(&event).unwrap();
|
||||
let deserialized: AgentEvent = serde_json::from_str(&json).unwrap();
|
||||
|
|
|
|||
|
|
@ -689,37 +689,54 @@ impl RunProjectionReducer for RunProjection {
|
|||
status: SubAgentStatus::Running,
|
||||
});
|
||||
}
|
||||
// A reused subagent stays one projected row: the spawn task and
|
||||
// generation 1 identify it, and every later generation only moves
|
||||
// its status. The per-turn task and generation stay in the event
|
||||
// log for consumers that need each turn.
|
||||
EventBody::AgentSubTurnStarted(props) => {
|
||||
set_subagent_status(
|
||||
self,
|
||||
stored,
|
||||
props.visit,
|
||||
event.seq,
|
||||
&props.agent_id,
|
||||
SubAgentStatus::Running,
|
||||
);
|
||||
}
|
||||
EventBody::AgentSubCompleted(props) => {
|
||||
let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq)
|
||||
else {
|
||||
return Ok(());
|
||||
};
|
||||
if let Some(subagent) = subagent_mut(stage, &props.agent_id) {
|
||||
subagent.status = SubAgentStatus::Completed {
|
||||
set_subagent_status(
|
||||
self,
|
||||
stored,
|
||||
props.visit,
|
||||
event.seq,
|
||||
&props.agent_id,
|
||||
SubAgentStatus::Completed {
|
||||
success: props.success,
|
||||
turns_used: props.turns_used,
|
||||
};
|
||||
}
|
||||
},
|
||||
);
|
||||
}
|
||||
EventBody::AgentSubFailed(props) => {
|
||||
let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq)
|
||||
else {
|
||||
return Ok(());
|
||||
};
|
||||
if let Some(subagent) = subagent_mut(stage, &props.agent_id) {
|
||||
subagent.status = SubAgentStatus::Failed {
|
||||
set_subagent_status(
|
||||
self,
|
||||
stored,
|
||||
props.visit,
|
||||
event.seq,
|
||||
&props.agent_id,
|
||||
SubAgentStatus::Failed {
|
||||
error: props.error.clone(),
|
||||
};
|
||||
}
|
||||
},
|
||||
);
|
||||
}
|
||||
EventBody::AgentSubClosed(props) => {
|
||||
let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq)
|
||||
else {
|
||||
return Ok(());
|
||||
};
|
||||
if let Some(subagent) = subagent_mut(stage, &props.agent_id) {
|
||||
subagent.status = SubAgentStatus::Closed;
|
||||
}
|
||||
set_subagent_status(
|
||||
self,
|
||||
stored,
|
||||
props.visit,
|
||||
event.seq,
|
||||
&props.agent_id,
|
||||
SubAgentStatus::Closed,
|
||||
);
|
||||
}
|
||||
EventBody::AgentSkillsDiscovered(props) => {
|
||||
let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq)
|
||||
|
|
@ -889,6 +906,25 @@ fn apply_todo_deleted(stage: &mut StageProjection, props: &TodoDeletedProps) {
|
|||
}
|
||||
}
|
||||
|
||||
/// Move an already-projected subagent to a new lifecycle status. Every
|
||||
/// subagent event after the spawn updates the same row, so reuse shows one
|
||||
/// agent returning to running rather than a second agent appearing.
|
||||
fn set_subagent_status(
|
||||
state: &mut RunProjection,
|
||||
stored: &RunEvent,
|
||||
visit: u32,
|
||||
seq: u32,
|
||||
agent_id: &str,
|
||||
status: SubAgentStatus,
|
||||
) {
|
||||
let Some(stage) = stage_at_stored_or_visit(state, stored, visit, seq) else {
|
||||
return;
|
||||
};
|
||||
if let Some(subagent) = subagent_mut(stage, agent_id) {
|
||||
subagent.status = status;
|
||||
}
|
||||
}
|
||||
|
||||
fn subagent_mut<'a>(
|
||||
stage: &'a mut StageProjection,
|
||||
agent_id: &str,
|
||||
|
|
@ -1602,9 +1638,10 @@ mod tests {
|
|||
AgentSessionDeactivatedProps, AgentSessionEndedProps, AgentSessionStartedProps,
|
||||
AgentSkillActivatedProps, AgentSkillActivationSource, AgentSkillSummary,
|
||||
AgentSkillsDiscoveredProps, AgentSteeringInjectedProps, AgentSubClosedProps,
|
||||
AgentSubCompletedProps, AgentSubFailedProps, AgentSubSpawnedProps, AgentToolCategory,
|
||||
AgentToolSource, AgentToolStartedProps, AgentToolSummary, AgentToolsAvailableProps,
|
||||
CheckpointCompletedProps, InterviewCompletedProps, InterviewOption, InterviewStartedProps,
|
||||
AgentSubCompletedProps, AgentSubFailedProps, AgentSubSpawnedProps,
|
||||
AgentSubTurnStartedProps, AgentToolCategory, AgentToolSource, AgentToolStartedProps,
|
||||
AgentToolSummary, AgentToolsAvailableProps, CheckpointCompletedProps,
|
||||
InterviewCompletedProps, InterviewOption, InterviewStartedProps,
|
||||
ParallelBranchCompletedProps, ParallelBranchStartedProps, RunCompletedProps,
|
||||
RunControlEffectProps, StageCompletedProps, StageFailedProps, StagePromptProps,
|
||||
StageRetryingProps, StageStartedProps,
|
||||
|
|
@ -6425,10 +6462,11 @@ mod tests {
|
|||
.apply_event(&test_stage_event(
|
||||
1,
|
||||
EventBody::AgentSubSpawned(AgentSubSpawnedProps {
|
||||
agent_id: "sub-1".to_string(),
|
||||
depth: 1,
|
||||
task: "write tests".to_string(),
|
||||
visit: 1,
|
||||
agent_id: "sub-1".to_string(),
|
||||
depth: 1,
|
||||
task: "write tests".to_string(),
|
||||
generation: 1,
|
||||
visit: 1,
|
||||
}),
|
||||
stage_id.clone(),
|
||||
))
|
||||
|
|
@ -6446,6 +6484,7 @@ mod tests {
|
|||
EventBody::AgentSubCompleted(AgentSubCompletedProps {
|
||||
agent_id: "sub-1".to_string(),
|
||||
depth: 1,
|
||||
generation: 1,
|
||||
success: true,
|
||||
turns_used: 3,
|
||||
visit: 1,
|
||||
|
|
@ -6462,23 +6501,64 @@ mod tests {
|
|||
state
|
||||
.apply_event(&test_stage_event(
|
||||
3,
|
||||
EventBody::AgentSubTurnStarted(AgentSubTurnStartedProps {
|
||||
agent_id: "sub-1".to_string(),
|
||||
depth: 1,
|
||||
task: "fix the review findings".to_string(),
|
||||
generation: 2,
|
||||
visit: 1,
|
||||
}),
|
||||
stage_id.clone(),
|
||||
))
|
||||
.unwrap();
|
||||
let stage = state.stage(&stage_id).unwrap();
|
||||
assert_eq!(stage.subagents.len(), 1);
|
||||
assert_eq!(stage.subagents[0].task, "write tests");
|
||||
assert_eq!(stage.subagents[0].status, SubAgentStatus::Running);
|
||||
|
||||
state
|
||||
.apply_event(&test_stage_event(
|
||||
4,
|
||||
EventBody::AgentSubCompleted(AgentSubCompletedProps {
|
||||
agent_id: "sub-1".to_string(),
|
||||
depth: 1,
|
||||
generation: 2,
|
||||
success: true,
|
||||
turns_used: 5,
|
||||
visit: 1,
|
||||
}),
|
||||
stage_id.clone(),
|
||||
))
|
||||
.unwrap();
|
||||
let stage = state.stage(&stage_id).unwrap();
|
||||
assert_eq!(stage.subagents.len(), 1);
|
||||
assert_eq!(stage.subagents[0].status, SubAgentStatus::Completed {
|
||||
success: true,
|
||||
turns_used: 5,
|
||||
});
|
||||
|
||||
state
|
||||
.apply_event(&test_stage_event(
|
||||
5,
|
||||
EventBody::AgentSubSpawned(AgentSubSpawnedProps {
|
||||
agent_id: "sub-2".to_string(),
|
||||
depth: 2,
|
||||
task: "debug failure".to_string(),
|
||||
visit: 1,
|
||||
agent_id: "sub-2".to_string(),
|
||||
depth: 2,
|
||||
task: "debug failure".to_string(),
|
||||
generation: 1,
|
||||
visit: 1,
|
||||
}),
|
||||
stage_id.clone(),
|
||||
))
|
||||
.unwrap();
|
||||
state
|
||||
.apply_event(&test_stage_event(
|
||||
4,
|
||||
6,
|
||||
EventBody::AgentSubFailed(AgentSubFailedProps {
|
||||
agent_id: "sub-2".to_string(),
|
||||
depth: 2,
|
||||
error: json!({ "message": "boom" }),
|
||||
visit: 1,
|
||||
agent_id: "sub-2".to_string(),
|
||||
depth: 2,
|
||||
generation: 1,
|
||||
error: json!({ "message": "boom" }),
|
||||
visit: 1,
|
||||
}),
|
||||
stage_id.clone(),
|
||||
))
|
||||
|
|
@ -6490,11 +6570,12 @@ mod tests {
|
|||
|
||||
state
|
||||
.apply_event(&test_stage_event(
|
||||
5,
|
||||
7,
|
||||
EventBody::AgentSubClosed(AgentSubClosedProps {
|
||||
agent_id: "sub-2".to_string(),
|
||||
depth: 2,
|
||||
visit: 1,
|
||||
agent_id: "sub-2".to_string(),
|
||||
depth: 2,
|
||||
generation: 1,
|
||||
visit: 1,
|
||||
}),
|
||||
stage_id.clone(),
|
||||
))
|
||||
|
|
|
|||
|
|
@ -758,20 +758,36 @@ fn event_body_from_event(event: &Event) -> EventBody {
|
|||
agent_id,
|
||||
depth,
|
||||
task,
|
||||
generation,
|
||||
} => EventBody::AgentSubSpawned(fabro_types::AgentSubSpawnedProps {
|
||||
agent_id: agent_id.clone(),
|
||||
depth: *depth,
|
||||
task: task.clone(),
|
||||
visit: *visit,
|
||||
agent_id: agent_id.clone(),
|
||||
depth: *depth,
|
||||
task: task.clone(),
|
||||
generation: *generation,
|
||||
visit: *visit,
|
||||
}),
|
||||
AgentEvent::SubAgentTurnStarted {
|
||||
agent_id,
|
||||
depth,
|
||||
task,
|
||||
generation,
|
||||
} => EventBody::AgentSubTurnStarted(fabro_types::AgentSubTurnStartedProps {
|
||||
agent_id: agent_id.clone(),
|
||||
depth: *depth,
|
||||
task: task.clone(),
|
||||
generation: *generation,
|
||||
visit: *visit,
|
||||
}),
|
||||
AgentEvent::SubAgentCompleted {
|
||||
agent_id,
|
||||
depth,
|
||||
generation,
|
||||
success,
|
||||
turns_used,
|
||||
} => EventBody::AgentSubCompleted(fabro_types::AgentSubCompletedProps {
|
||||
agent_id: agent_id.clone(),
|
||||
depth: *depth,
|
||||
generation: *generation,
|
||||
success: *success,
|
||||
turns_used: *turns_used,
|
||||
visit: *visit,
|
||||
|
|
@ -779,18 +795,25 @@ fn event_body_from_event(event: &Event) -> EventBody {
|
|||
AgentEvent::SubAgentFailed {
|
||||
agent_id,
|
||||
depth,
|
||||
generation,
|
||||
error,
|
||||
} => EventBody::AgentSubFailed(fabro_types::AgentSubFailedProps {
|
||||
agent_id: agent_id.clone(),
|
||||
depth: *depth,
|
||||
error: serde_json::to_value(error).expect("agent Error derives Serialize with no custom logic that can fail"),
|
||||
visit: *visit,
|
||||
agent_id: agent_id.clone(),
|
||||
depth: *depth,
|
||||
generation: *generation,
|
||||
error: serde_json::to_value(error).expect("agent Error derives Serialize with no custom logic that can fail"),
|
||||
visit: *visit,
|
||||
}),
|
||||
AgentEvent::SubAgentClosed { agent_id, depth } => {
|
||||
AgentEvent::SubAgentClosed {
|
||||
agent_id,
|
||||
depth,
|
||||
generation,
|
||||
} => {
|
||||
EventBody::AgentSubClosed(fabro_types::AgentSubClosedProps {
|
||||
agent_id: agent_id.clone(),
|
||||
depth: *depth,
|
||||
visit: *visit,
|
||||
agent_id: agent_id.clone(),
|
||||
depth: *depth,
|
||||
generation: *generation,
|
||||
visit: *visit,
|
||||
})
|
||||
}
|
||||
AgentEvent::McpServerReady {
|
||||
|
|
|
|||
|
|
@ -86,6 +86,7 @@ pub fn event_name(event: &Event) -> &'static str {
|
|||
AgentEvent::CompactionCompleted { .. } => "agent.compaction.completed",
|
||||
AgentEvent::LlmRetry { .. } => "agent.llm.retry",
|
||||
AgentEvent::SubAgentSpawned { .. } => "agent.sub.spawned",
|
||||
AgentEvent::SubAgentTurnStarted { .. } => "agent.sub.turn.started",
|
||||
AgentEvent::SubAgentCompleted { .. } => "agent.sub.completed",
|
||||
AgentEvent::SubAgentFailed { .. } => "agent.sub.failed",
|
||||
AgentEvent::SubAgentClosed { .. } => "agent.sub.closed",
|
||||
|
|
@ -184,9 +185,10 @@ mod tests {
|
|||
stage: "code".to_string(),
|
||||
visit: 1,
|
||||
event: AgentEvent::SubAgentSpawned {
|
||||
agent_id: "a1".to_string(),
|
||||
depth: 1,
|
||||
task: "do it".to_string(),
|
||||
agent_id: "a1".to_string(),
|
||||
depth: 1,
|
||||
task: "do it".to_string(),
|
||||
generation: 1,
|
||||
},
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
|
|
@ -194,6 +196,22 @@ mod tests {
|
|||
}),
|
||||
"agent.sub.spawned"
|
||||
);
|
||||
assert_eq!(
|
||||
event_name(&Event::Agent {
|
||||
stage: "code".to_string(),
|
||||
visit: 1,
|
||||
event: AgentEvent::SubAgentTurnStarted {
|
||||
agent_id: "a1".to_string(),
|
||||
depth: 1,
|
||||
task: "fix it".to_string(),
|
||||
generation: 2,
|
||||
},
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
tool_call_id: None,
|
||||
}),
|
||||
"agent.sub.turn.started"
|
||||
);
|
||||
assert_eq!(
|
||||
event_name(&Event::Agent {
|
||||
stage: "code".to_string(),
|
||||
|
|
|
|||
|
|
@ -4079,8 +4079,9 @@ enabled = true
|
|||
|
||||
session.sub_agent_event_callback()(
|
||||
fabro_agent::subagent::SubAgentCallbackEvent::Lifecycle(AgentEvent::SubAgentClosed {
|
||||
agent_id: "child-1".to_string(),
|
||||
depth: 1,
|
||||
agent_id: "child-1".to_string(),
|
||||
depth: 1,
|
||||
generation: 1,
|
||||
}),
|
||||
);
|
||||
let session_id = session.id().to_string();
|
||||
|
|
|
|||
|
|
@ -114,10 +114,10 @@ pub use run_blob_id::RunBlobId;
|
|||
pub use run_event::{
|
||||
AgentMcpToolSummary, AgentMemoryFileProps, AgentSkillActivationSource, AgentSkillSummary,
|
||||
AgentToolCategory, AgentToolSource, AgentToolSummary, AgentToolsAvailableProps, EventBody,
|
||||
ExecOutputTail, FailoverProps, InterviewOption, LlmOutputKind, LlmRetryPhase,
|
||||
MetadataSnapshotFailureKind, MetadataSnapshotPhase, RunEvent, RunNoticeCode, RunNoticeLevel,
|
||||
RunPairEndedReason, RunPairFailedReason, RunRunnableSource, SessionCapability,
|
||||
TodoCreatedProps, TodoDeletedProps, TodoUpdatedProps,
|
||||
ExecOutputTail, FailoverProps, INITIAL_SUBAGENT_GENERATION, InterviewOption, LlmOutputKind,
|
||||
LlmRetryPhase, MetadataSnapshotFailureKind, MetadataSnapshotPhase, RunEvent, RunNoticeCode,
|
||||
RunNoticeLevel, RunPairEndedReason, RunPairFailedReason, RunRunnableSource, SessionCapability,
|
||||
TodoCreatedProps, TodoDeletedProps, TodoUpdatedProps, initial_subagent_generation,
|
||||
};
|
||||
pub use run_failure::RunFailure;
|
||||
pub use run_id::{RunId, fixtures};
|
||||
|
|
|
|||
|
|
@ -378,16 +378,29 @@ pub struct AgentLlmFirstOutputProps {
|
|||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct AgentSubSpawnedProps {
|
||||
pub agent_id: String,
|
||||
pub depth: usize,
|
||||
pub task: String,
|
||||
pub visit: u32,
|
||||
pub agent_id: String,
|
||||
pub depth: usize,
|
||||
pub task: String,
|
||||
#[serde(default = "initial_subagent_generation")]
|
||||
pub generation: u64,
|
||||
pub visit: u32,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct AgentSubTurnStartedProps {
|
||||
pub agent_id: String,
|
||||
pub depth: usize,
|
||||
pub task: String,
|
||||
pub generation: u64,
|
||||
pub visit: u32,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct AgentSubCompletedProps {
|
||||
pub agent_id: String,
|
||||
pub depth: usize,
|
||||
#[serde(default = "initial_subagent_generation")]
|
||||
pub generation: u64,
|
||||
pub success: bool,
|
||||
pub turns_used: usize,
|
||||
pub visit: u32,
|
||||
|
|
@ -395,17 +408,32 @@ pub struct AgentSubCompletedProps {
|
|||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct AgentSubFailedProps {
|
||||
pub agent_id: String,
|
||||
pub depth: usize,
|
||||
pub error: Value,
|
||||
pub visit: u32,
|
||||
pub agent_id: String,
|
||||
pub depth: usize,
|
||||
#[serde(default = "initial_subagent_generation")]
|
||||
pub generation: u64,
|
||||
pub error: Value,
|
||||
pub visit: u32,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct AgentSubClosedProps {
|
||||
pub agent_id: String,
|
||||
pub depth: usize,
|
||||
pub visit: u32,
|
||||
pub agent_id: String,
|
||||
pub depth: usize,
|
||||
#[serde(default = "initial_subagent_generation")]
|
||||
pub generation: u64,
|
||||
pub visit: u32,
|
||||
}
|
||||
|
||||
/// The generation of a subagent's first turn. Events stored before subagent
|
||||
/// session reuse existed carry no generation, so they read back as this.
|
||||
pub const INITIAL_SUBAGENT_GENERATION: u64 = 1;
|
||||
|
||||
/// Serde default for the generation of a stored subagent event. Public so
|
||||
/// crates with their own subagent event types share this one definition.
|
||||
#[must_use]
|
||||
pub const fn initial_subagent_generation() -> u64 {
|
||||
INITIAL_SUBAGENT_GENERATION
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
|
|
|
|||
|
|
@ -242,6 +242,8 @@ pub enum EventBody {
|
|||
AgentLlmRetry(AgentLlmRetryProps),
|
||||
#[serde(rename = "agent.sub.spawned")]
|
||||
AgentSubSpawned(AgentSubSpawnedProps),
|
||||
#[serde(rename = "agent.sub.turn.started")]
|
||||
AgentSubTurnStarted(AgentSubTurnStartedProps),
|
||||
#[serde(rename = "agent.sub.completed")]
|
||||
AgentSubCompleted(AgentSubCompletedProps),
|
||||
#[serde(rename = "agent.sub.failed")]
|
||||
|
|
@ -510,6 +512,7 @@ impl EventBody {
|
|||
Self::AgentLlmFirstOutput(_) => "agent.llm.first_output",
|
||||
Self::AgentLlmRetry(_) => "agent.llm.retry",
|
||||
Self::AgentSubSpawned(_) => "agent.sub.spawned",
|
||||
Self::AgentSubTurnStarted(_) => "agent.sub.turn.started",
|
||||
Self::AgentSubCompleted(_) => "agent.sub.completed",
|
||||
Self::AgentSubFailed(_) => "agent.sub.failed",
|
||||
Self::AgentSubClosed(_) => "agent.sub.closed",
|
||||
|
|
@ -682,6 +685,7 @@ fn is_known_event_name(event: &str) -> bool {
|
|||
| "agent.llm.first_output"
|
||||
| "agent.llm.retry"
|
||||
| "agent.sub.spawned"
|
||||
| "agent.sub.turn.started"
|
||||
| "agent.sub.completed"
|
||||
| "agent.sub.failed"
|
||||
| "agent.sub.closed"
|
||||
|
|
@ -2497,6 +2501,37 @@ mod tests {
|
|||
assert_eq!(parsed, body);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn subagent_generations_are_typed_and_legacy_events_default_to_one() {
|
||||
let started = EventBody::AgentSubTurnStarted(AgentSubTurnStartedProps {
|
||||
agent_id: "sub-1".to_string(),
|
||||
depth: 1,
|
||||
task: "fix the review findings".to_string(),
|
||||
generation: 2,
|
||||
visit: 1,
|
||||
});
|
||||
let value = serde_json::to_value(&started).unwrap();
|
||||
assert_eq!(value["event"], "agent.sub.turn.started");
|
||||
assert_eq!(value["properties"]["generation"], 2);
|
||||
assert_eq!(serde_json::from_value::<EventBody>(value).unwrap(), started);
|
||||
|
||||
let legacy: EventBody = serde_json::from_value(json!({
|
||||
"event": "agent.sub.completed",
|
||||
"properties": {
|
||||
"agent_id": "sub-1",
|
||||
"depth": 1,
|
||||
"success": true,
|
||||
"turns_used": 3,
|
||||
"visit": 1
|
||||
}
|
||||
}))
|
||||
.unwrap();
|
||||
let EventBody::AgentSubCompleted(props) = legacy else {
|
||||
panic!("expected subagent completion");
|
||||
};
|
||||
assert_eq!(props.generation, 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn agent_tool_source_and_category_use_public_json_shape() {
|
||||
assert_eq!(
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue