mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-09-14 23:22:51 +00:00
Add SubAgentStatus enum for explicit subagent lifecycle tracking
Retain agents in the HashMap after wait/close instead of removing them, enabling cached result retrieval, status queries, and disambiguated error messages (never spawned vs completed vs closed vs failed). Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
b91ae3df47
commit
72386166c9
3 changed files with 241 additions and 56 deletions
|
|
@ -45,7 +45,9 @@ pub use sandbox::{
|
|||
};
|
||||
pub use session::Session;
|
||||
pub use skills::Skill;
|
||||
pub use subagent::{SubAgent, SubAgentEventCallback, SubAgentManager, SubAgentResult};
|
||||
pub use subagent::{
|
||||
SubAgent, SubAgentEventCallback, SubAgentManager, SubAgentResult, SubAgentStatus,
|
||||
};
|
||||
pub use tool_registry::ToolRegistry;
|
||||
pub use tools::{
|
||||
make_edit_file_tool, make_glob_tool, make_grep_tool, make_read_file_tool, make_shell_tool,
|
||||
|
|
|
|||
|
|
@ -2713,8 +2713,11 @@ mod tests {
|
|||
let mut rx = session.subscribe();
|
||||
session.close();
|
||||
|
||||
// The subagent should have been cleaned up
|
||||
assert!(manager.lock().await.get(&agent_id).is_none());
|
||||
// The subagent should have been closed
|
||||
assert_eq!(
|
||||
manager.lock().await.status(&agent_id),
|
||||
Some(&crate::subagent::SubAgentStatus::Closed)
|
||||
);
|
||||
|
||||
// Verify event ordering: SubAgentClosed before SessionEnded
|
||||
let mut events = Vec::new();
|
||||
|
|
|
|||
|
|
@ -18,11 +18,21 @@ pub struct SubAgentResult {
|
|||
pub turns_used: usize,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub enum SubAgentStatus {
|
||||
Running,
|
||||
Completed,
|
||||
Failed,
|
||||
Closed,
|
||||
}
|
||||
|
||||
pub struct SubAgent {
|
||||
task: Option<tokio::task::JoinHandle<Result<SubAgentResult, AgentError>>>,
|
||||
followup_queue: Arc<Mutex<VecDeque<String>>>,
|
||||
cancel_token: CancellationToken,
|
||||
depth: usize,
|
||||
status: SubAgentStatus,
|
||||
cached_result: Option<Result<SubAgentResult, AgentError>>,
|
||||
}
|
||||
|
||||
pub struct SubAgentManager {
|
||||
|
|
@ -123,6 +133,8 @@ impl SubAgentManager {
|
|||
followup_queue,
|
||||
cancel_token,
|
||||
depth: depth + 1,
|
||||
status: SubAgentStatus::Running,
|
||||
cached_result: None,
|
||||
},
|
||||
);
|
||||
|
||||
|
|
@ -137,9 +149,30 @@ impl SubAgentManager {
|
|||
|
||||
pub fn send_input(&self, agent_id: &str, message: &str) -> Result<(), AgentError> {
|
||||
let agent = self.agents.get(agent_id).ok_or_else(|| {
|
||||
AgentError::InvalidState(format!("No agent found with id: {agent_id}"))
|
||||
AgentError::InvalidState(format!(
|
||||
"No agent found with id: {agent_id} (it was never spawned)"
|
||||
))
|
||||
})?;
|
||||
|
||||
match agent.status {
|
||||
SubAgentStatus::Running => {}
|
||||
SubAgentStatus::Completed => {
|
||||
return Err(AgentError::InvalidState(format!(
|
||||
"Agent {agent_id} has already completed"
|
||||
)));
|
||||
}
|
||||
SubAgentStatus::Failed => {
|
||||
return Err(AgentError::InvalidState(format!(
|
||||
"Agent {agent_id} has failed"
|
||||
)));
|
||||
}
|
||||
SubAgentStatus::Closed => {
|
||||
return Err(AgentError::InvalidState(format!(
|
||||
"Agent {agent_id} has been closed"
|
||||
)));
|
||||
}
|
||||
}
|
||||
|
||||
agent
|
||||
.followup_queue
|
||||
.lock()
|
||||
|
|
@ -150,61 +183,112 @@ impl SubAgentManager {
|
|||
}
|
||||
|
||||
pub async fn wait(&mut self, agent_id: &str) -> Result<SubAgentResult, AgentError> {
|
||||
let mut agent = self.agents.remove(agent_id).ok_or_else(|| {
|
||||
AgentError::InvalidState(format!("No agent found with id: {agent_id}"))
|
||||
})?;
|
||||
// Phase 1: Check existence and current status
|
||||
let agent = self.agents.get(agent_id);
|
||||
let (status, depth) = match agent {
|
||||
None => {
|
||||
return Err(AgentError::InvalidState(format!(
|
||||
"No agent found with id: {agent_id} (it was never spawned)"
|
||||
)));
|
||||
}
|
||||
Some(a) => (a.status.clone(), a.depth),
|
||||
};
|
||||
|
||||
let depth = agent.depth;
|
||||
|
||||
match agent.task.take() {
|
||||
Some(join_handle) => match join_handle.await {
|
||||
Ok(Ok(result)) => {
|
||||
self.emit_event(AgentEvent::SubAgentCompleted {
|
||||
agent_id: agent_id.to_string(),
|
||||
depth,
|
||||
success: result.success,
|
||||
turns_used: result.turns_used,
|
||||
});
|
||||
Ok(result)
|
||||
}
|
||||
Ok(Err(e)) => {
|
||||
self.emit_event(AgentEvent::SubAgentFailed {
|
||||
agent_id: agent_id.to_string(),
|
||||
depth,
|
||||
error: e.clone(),
|
||||
});
|
||||
Err(e)
|
||||
}
|
||||
Err(e) => {
|
||||
let error_msg = format!("Agent task panicked: {e}");
|
||||
self.emit_event(AgentEvent::SubAgentFailed {
|
||||
agent_id: agent_id.to_string(),
|
||||
depth,
|
||||
error: AgentError::InvalidState(error_msg.clone()),
|
||||
});
|
||||
Err(AgentError::InvalidState(error_msg))
|
||||
}
|
||||
},
|
||||
None => Err(AgentError::InvalidState(format!(
|
||||
"Agent {agent_id} has no running task"
|
||||
))),
|
||||
match status {
|
||||
SubAgentStatus::Closed => {
|
||||
return Err(AgentError::InvalidState(format!(
|
||||
"Agent {agent_id} has been closed"
|
||||
)));
|
||||
}
|
||||
SubAgentStatus::Completed | SubAgentStatus::Failed => {
|
||||
// Return cached result without emitting events again
|
||||
let agent = self.agents.get(agent_id).unwrap();
|
||||
return agent.cached_result.clone().unwrap();
|
||||
}
|
||||
SubAgentStatus::Running => {}
|
||||
}
|
||||
|
||||
// Phase 2: Take the JoinHandle (brief mutable borrow, no await)
|
||||
let join_handle = self
|
||||
.agents
|
||||
.get_mut(agent_id)
|
||||
.unwrap()
|
||||
.task
|
||||
.take()
|
||||
.ok_or_else(|| {
|
||||
AgentError::InvalidState(format!("Agent {agent_id} has no running task"))
|
||||
})?;
|
||||
|
||||
// Phase 3: Await the task (no borrow held)
|
||||
let task_result = match join_handle.await {
|
||||
Ok(result) => result,
|
||||
Err(e) => Err(AgentError::InvalidState(format!(
|
||||
"Agent task panicked: {e}"
|
||||
))),
|
||||
};
|
||||
|
||||
// Phase 4: Cache result and update status
|
||||
let new_status = if task_result.is_ok() {
|
||||
SubAgentStatus::Completed
|
||||
} else {
|
||||
SubAgentStatus::Failed
|
||||
};
|
||||
let agent = self.agents.get_mut(agent_id).unwrap();
|
||||
agent.status = new_status;
|
||||
agent.cached_result = Some(task_result.clone());
|
||||
|
||||
// Phase 5: Emit event
|
||||
match &task_result {
|
||||
Ok(result) => {
|
||||
self.emit_event(AgentEvent::SubAgentCompleted {
|
||||
agent_id: agent_id.to_string(),
|
||||
depth,
|
||||
success: result.success,
|
||||
turns_used: result.turns_used,
|
||||
});
|
||||
}
|
||||
Err(e) => {
|
||||
self.emit_event(AgentEvent::SubAgentFailed {
|
||||
agent_id: agent_id.to_string(),
|
||||
depth,
|
||||
error: e.clone(),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
task_result
|
||||
}
|
||||
|
||||
pub fn close(&mut self, agent_id: &str) -> Result<(), AgentError> {
|
||||
let agent = self.agents.remove(agent_id).ok_or_else(|| {
|
||||
AgentError::InvalidState(format!("No agent found with id: {agent_id}"))
|
||||
let agent = self.agents.get_mut(agent_id).ok_or_else(|| {
|
||||
AgentError::InvalidState(format!(
|
||||
"No agent found with id: {agent_id} (it was never spawned)"
|
||||
))
|
||||
})?;
|
||||
|
||||
agent.cancel_token.cancel();
|
||||
|
||||
if let Some(join_handle) = agent.task {
|
||||
join_handle.abort();
|
||||
match agent.status {
|
||||
SubAgentStatus::Closed => {
|
||||
return Err(AgentError::InvalidState(format!(
|
||||
"Agent {agent_id} is already closed"
|
||||
)));
|
||||
}
|
||||
SubAgentStatus::Running => {
|
||||
agent.cancel_token.cancel();
|
||||
if let Some(join_handle) = agent.task.take() {
|
||||
join_handle.abort();
|
||||
}
|
||||
}
|
||||
SubAgentStatus::Completed | SubAgentStatus::Failed => {
|
||||
// No task to cancel, just transition status
|
||||
}
|
||||
}
|
||||
|
||||
agent.status = SubAgentStatus::Closed;
|
||||
let depth = agent.depth;
|
||||
|
||||
self.emit_event(AgentEvent::SubAgentClosed {
|
||||
agent_id: agent_id.to_string(),
|
||||
depth: agent.depth,
|
||||
depth,
|
||||
});
|
||||
|
||||
Ok(())
|
||||
|
|
@ -212,9 +296,13 @@ impl SubAgentManager {
|
|||
|
||||
/// Close all active subagents, cancelling their tokens and aborting tasks.
|
||||
pub fn close_all(&mut self) {
|
||||
let ids: Vec<String> = self.agents.keys().cloned().collect();
|
||||
let ids: Vec<String> = self
|
||||
.agents
|
||||
.iter()
|
||||
.filter(|(_, a)| a.status == SubAgentStatus::Running)
|
||||
.map(|(id, _)| id.clone())
|
||||
.collect();
|
||||
for id in ids {
|
||||
// close() always succeeds for known IDs, ignore result
|
||||
let _ = self.close(&id);
|
||||
}
|
||||
}
|
||||
|
|
@ -224,6 +312,12 @@ impl SubAgentManager {
|
|||
pub fn get(&self, agent_id: &str) -> Option<&SubAgent> {
|
||||
self.agents.get(agent_id)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[must_use]
|
||||
pub fn status(&self, agent_id: &str) -> Option<&SubAgentStatus> {
|
||||
self.agents.get(agent_id).map(|a| &a.status)
|
||||
}
|
||||
}
|
||||
|
||||
pub fn make_spawn_agent_tool(
|
||||
|
|
@ -449,7 +543,7 @@ mod tests {
|
|||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn close_removes_agent() {
|
||||
async fn close_sets_closed_status() {
|
||||
let mut manager = SubAgentManager::new(3);
|
||||
let session = make_session(vec![text_response("Hello")]).await;
|
||||
let agent_id = manager.spawn(session, "Do something".into(), 0).unwrap();
|
||||
|
|
@ -457,7 +551,7 @@ mod tests {
|
|||
|
||||
let result = manager.close(&agent_id);
|
||||
assert!(result.is_ok());
|
||||
assert!(manager.get(&agent_id).is_none());
|
||||
assert_eq!(manager.status(&agent_id), Some(&SubAgentStatus::Closed));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
|
@ -488,7 +582,7 @@ mod tests {
|
|||
assert_eq!(agent_result.output, "Task completed successfully");
|
||||
assert!(agent_result.success);
|
||||
assert!(agent_result.turns_used > 0);
|
||||
assert!(manager.get(&agent_id).is_none());
|
||||
assert_eq!(manager.status(&agent_id), Some(&SubAgentStatus::Completed));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -633,7 +727,7 @@ mod tests {
|
|||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn close_all_removes_all_agents() {
|
||||
async fn close_all_closes_all_agents() {
|
||||
let mut manager = SubAgentManager::new(3);
|
||||
let session1 = make_session(vec![text_response("Hello")]).await;
|
||||
let session2 = make_session(vec![text_response("World")]).await;
|
||||
|
|
@ -644,8 +738,8 @@ mod tests {
|
|||
|
||||
manager.close_all();
|
||||
|
||||
assert!(manager.get(&id1).is_none());
|
||||
assert!(manager.get(&id2).is_none());
|
||||
assert_eq!(manager.status(&id1), Some(&SubAgentStatus::Closed));
|
||||
assert_eq!(manager.status(&id2), Some(&SubAgentStatus::Closed));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
|
@ -654,4 +748,90 @@ mod tests {
|
|||
manager.close_all(); // should not panic
|
||||
assert!(manager.agents.is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn wait_twice_returns_cached_result() {
|
||||
let mut manager = SubAgentManager::new(3);
|
||||
let session = make_session(vec![text_response("cached output")]).await;
|
||||
let agent_id = manager.spawn(session, "Do something".into(), 0).unwrap();
|
||||
|
||||
let result1 = manager.wait(&agent_id).await.unwrap();
|
||||
let result2 = manager.wait(&agent_id).await.unwrap();
|
||||
|
||||
assert_eq!(result1.output, "cached output");
|
||||
assert_eq!(result2.output, "cached output");
|
||||
assert_eq!(manager.status(&agent_id), Some(&SubAgentStatus::Completed));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn send_input_to_completed_agent_errors() {
|
||||
let mut manager = SubAgentManager::new(3);
|
||||
let session = make_session(vec![text_response("done")]).await;
|
||||
let agent_id = manager.spawn(session, "Do something".into(), 0).unwrap();
|
||||
let _ = manager.wait(&agent_id).await.unwrap();
|
||||
|
||||
let result = manager.send_input(&agent_id, "hello");
|
||||
assert!(result.is_err());
|
||||
assert!(result
|
||||
.unwrap_err()
|
||||
.to_string()
|
||||
.contains("already completed"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn send_input_to_closed_agent_errors() {
|
||||
let mut manager = SubAgentManager::new(3);
|
||||
let session = make_session(vec![text_response("Hello")]).await;
|
||||
let agent_id = manager.spawn(session, "Do something".into(), 0).unwrap();
|
||||
manager.close(&agent_id).unwrap();
|
||||
|
||||
let result = manager.send_input(&agent_id, "hello");
|
||||
assert!(result.is_err());
|
||||
assert!(result.unwrap_err().to_string().contains("has been closed"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn close_already_closed_agent_errors() {
|
||||
let mut manager = SubAgentManager::new(3);
|
||||
let session = make_session(vec![text_response("Hello")]).await;
|
||||
let agent_id = manager.spawn(session, "Do something".into(), 0).unwrap();
|
||||
manager.close(&agent_id).unwrap();
|
||||
|
||||
let result = manager.close(&agent_id);
|
||||
assert!(result.is_err());
|
||||
assert!(result.unwrap_err().to_string().contains("already closed"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn close_completed_agent_succeeds() {
|
||||
let mut manager = SubAgentManager::new(3);
|
||||
let session = make_session(vec![text_response("done")]).await;
|
||||
let agent_id = manager.spawn(session, "Do something".into(), 0).unwrap();
|
||||
let _ = manager.wait(&agent_id).await.unwrap();
|
||||
assert_eq!(manager.status(&agent_id), Some(&SubAgentStatus::Completed));
|
||||
|
||||
let result = manager.close(&agent_id);
|
||||
assert!(result.is_ok());
|
||||
assert_eq!(manager.status(&agent_id), Some(&SubAgentStatus::Closed));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn status_is_running_after_spawn() {
|
||||
let mut manager = SubAgentManager::new(3);
|
||||
let session = make_session(vec![text_response("Hello")]).await;
|
||||
let agent_id = manager.spawn(session, "Do something".into(), 0).unwrap();
|
||||
assert_eq!(manager.status(&agent_id), Some(&SubAgentStatus::Running));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn wait_on_closed_agent_errors() {
|
||||
let mut manager = SubAgentManager::new(3);
|
||||
let session = make_session(vec![text_response("Hello")]).await;
|
||||
let agent_id = manager.spawn(session, "Do something".into(), 0).unwrap();
|
||||
manager.close(&agent_id).unwrap();
|
||||
|
||||
let result = manager.wait(&agent_id).await;
|
||||
assert!(result.is_err());
|
||||
assert!(result.unwrap_err().to_string().contains("has been closed"));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue