From 55aee09fc584f4767aa6866c2b78f89d3994c63d Mon Sep 17 00:00:00 2001 From: TheArchitectit Date: Wed, 29 Apr 2026 10:16:22 -0500 Subject: [PATCH] feat: TeamStatus tool and background team watcher - TeamStatus tool with 3 actions: - status: live snapshot (running/completed/failed counts, agent details) - summary: final results when agents finish (includes result content) - events: timeline from append-only event log - Background team watcher thread spawned by TeamCreate: - Polls agent .json files every 2s - Prints [team] progress to stderr on agent completion/failure - Updates team manifest status when all agents finish - Writes events to .clawd-agents/teams/{team_id}-events.jsonl - TeamStatus added to PARALLEL_SAFE_TOOLS and all agent allowed_tools Co-authored-by: GLM 5.1 FP8 via Crush --- rust/crates/rusty-claude-cli/src/main.rs | 145 ++++++++++++ rust/crates/tools/src/lib.rs | 281 ++++++++++++++++++++++- 2 files changed, 425 insertions(+), 1 deletion(-) diff --git a/rust/crates/rusty-claude-cli/src/main.rs b/rust/crates/rusty-claude-cli/src/main.rs index c2cc21bf..a826ac00 100644 --- a/rust/crates/rusty-claude-cli/src/main.rs +++ b/rust/crates/rusty-claude-cli/src/main.rs @@ -14540,6 +14540,151 @@ impl ToolExecutor for CliToolExecutor { } } } + + fn execute_batch(&mut self, calls: Vec) -> Vec { + if calls.len() <= 1 { + return calls + .into_iter() + .map(|call| { + let result = self.execute(&call.tool_name, &call.input); + runtime::ToolResult { + tool_use_id: call.tool_use_id, + tool_name: call.tool_name, + result, + } + }) + .collect(); + } + + /// Tools that are safe to run in parallel because they only read + /// state and dispatch through the stateless tool registry. + const PARALLEL_SAFE_TOOLS: &[&str] = &[ + "read_file", + "glob_search", + "grep_search", + "WebFetch", + "WebSearch", + "ToolSearch", + "Skill", + "LSP", + "Agent", + "AgentMessage", + "TeamStatus", + "TaskGet", + "TaskList", + "TaskOutput", + "GitStatus", + "GitDiff", + "GitLog", + "GitShow", + "GitBlame", + ]; + + let emit_output = self.emit_output; + let mut results: Vec> = vec![None; calls.len()]; + let mut parallel_calls: Vec<(usize, String, String, String)> = Vec::new(); + let mut sequential_indices: Vec = Vec::new(); + + // Classify calls as parallel-safe or sequential + for (i, call) in calls.iter().enumerate() { + if self + .allowed_tools + .as_ref() + .is_some_and(|allowed| !allowed.contains(&canonical_allowed_tool_name(&call.tool_name))) + { + results[i] = Some(runtime::ToolResult { + tool_use_id: call.tool_use_id.clone(), + tool_name: call.tool_name.clone(), + result: Err(ToolError::new(format!( + "tool `{}` is not enabled by the current --allowedTools setting", + call.tool_name + ))), + }); + } else if PARALLEL_SAFE_TOOLS.contains(&call.tool_name.as_str()) + && !self.tool_registry.has_runtime_tool(&call.tool_name) + { + parallel_calls.push(( + i, + call.tool_use_id.clone(), + call.tool_name.clone(), + call.input.clone(), + )); + } else { + sequential_indices.push(i); + } + } + + // Execute parallel-safe tools concurrently + if !parallel_calls.is_empty() { + let registry = self.tool_registry.clone(); + let parallel_results: Vec<(usize, String, String, Result)> = + std::thread::scope(|s| { + let mut handles = Vec::new(); + for (idx, tool_use_id, tool_name, input) in ¶llel_calls { + let registry = ®istry; + let tool_use_id = tool_use_id.clone(); + let tool_name = tool_name.clone(); + let input = input.clone(); + let idx = *idx; + handles.push(s.spawn(move || { + let value = serde_json::from_str(&input) + .map_err(|error| ToolError::new(format!("invalid tool input JSON: {error}"))); + let result = match value { + Ok(v) => registry + .execute(&tool_name, &v) + .map_err(ToolError::new), + Err(e) => Err(e), + }; + (idx, tool_use_id, tool_name, result) + })); + } + handles + .into_iter() + .map(|h| h.join().unwrap_or_else(|_| { + ( + 0, + String::new(), + String::new(), + Err(ToolError::new("parallel thread panicked")), + ) + })) + .collect() + }); + + for (idx, tool_use_id, tool_name, result) in parallel_results { + if emit_output { + let output_str = match &result { + Ok(o) => o.clone(), + Err(e) => e.to_string(), + }; + let is_error = result.is_err(); + let markdown = format_tool_result(&tool_name, &output_str, is_error); + self.renderer + .stream_markdown(&markdown, &mut io::stdout()) + .map_err(|error| ToolError::new(error.to_string())) + .ok(); + } + results[idx] = Some(runtime::ToolResult { + tool_use_id, + tool_name, + result, + }); + } + } + + // Execute sequential tools one at a time + for idx in sequential_indices { + let call = &calls[idx]; + let result = self.execute(&call.tool_name, &call.input); + results[idx] = Some(runtime::ToolResult { + tool_use_id: call.tool_use_id.clone(), + tool_name: call.tool_name.clone(), + result, + }); + } + + results.into_iter().map(|r| r.unwrap()).collect() + } } fn permission_policy( diff --git a/rust/crates/tools/src/lib.rs b/rust/crates/tools/src/lib.rs index 5da6f554..1e50580f 100644 --- a/rust/crates/tools/src/lib.rs +++ b/rust/crates/tools/src/lib.rs @@ -1162,6 +1162,24 @@ pub fn mvp_tool_specs() -> Vec { }), required_permission: PermissionMode::ReadOnly, }, + ToolSpec { + name: "TeamStatus", + description: "Check the progress of a team of agents. Returns structured status: which agents are running, completed, or failed, with their results. Use action=status for a snapshot, action=summary for final results when all agents are done, or action=events for a timeline of team events.", + input_schema: json!({ + "type": "object", + "properties": { + "team_id": { "type": "string", "description": "Team ID to check" }, + "action": { + "type": "string", + "enum": ["status", "summary", "events"], + "description": "status=live snapshot, summary=final results, events=timeline" + } + }, + "required": ["team_id"], + "additionalProperties": false + }), + required_permission: PermissionMode::ReadOnly, + }, ToolSpec { name: "CronCreate", description: "Create a scheduled recurring task.", @@ -1510,6 +1528,7 @@ fn execute_tool_with_enforcer( "TeamCreate" => from_value::(input).and_then(run_team_create), "TeamDelete" => from_value::(input).and_then(run_team_delete), "AgentMessage" => from_value::(input).and_then(run_agent_message), + "TeamStatus" => from_value::(input).and_then(run_team_status), "CronCreate" => from_value::(input).and_then(run_cron_create), "CronDelete" => from_value::(input).and_then(run_cron_delete), "CronList" => run_cron_list(input.clone()), @@ -1875,6 +1894,9 @@ fn run_team_create(input: TeamCreateInput) -> Result { // Register in global registry let team = global_team_registry().create(&input.name, agent_ids.clone()); + // Spawn background watcher that prints progress to stderr + spawn_team_watcher(&team_id, &agent_ids); + to_pretty_json(json!({ "team_id": team_id, "name": input.name, @@ -1882,11 +1904,257 @@ fn run_team_create(input: TeamCreateInput) -> Result { "agents": agent_outputs, "status": "running", "created_at": team.created_at, - "message": format!("Team created with {} agents. Use AgentMessage to coordinate. Poll agent .json files for results.", agent_ids.len()), + "message": format!("Team created with {} agents. Use AgentMessage to coordinate. TeamStatus shows live progress.", agent_ids.len()), })) } #[allow(clippy::needless_pass_by_value)] +fn run_team_status(input: TeamStatusInput) -> Result { + let action = input.action.as_deref().unwrap_or("status"); + let team_dir = agent_store_dir()?.join("teams"); + let manifest_path = team_dir.join(format!("{}.json", input.team_id)); + if !manifest_path.exists() { + return Err(format!("team {} not found", input.team_id)); + } + let team_data: Value = serde_json::from_str( + &std::fs::read_to_string(&manifest_path).map_err(|e| e.to_string())? + ).map_err(|e| e.to_string())?; + let agent_ids = team_data.get("agent_ids") + .and_then(|v| v.as_array()) + .cloned() + .unwrap_or_default(); + + let store_dir = agent_store_dir()?; + let mut agents_detail: Vec = Vec::new(); + let mut running_count = 0usize; + let mut completed_count = 0usize; + let mut failed_count = 0usize; + + for id_val in &agent_ids { + if let Some(id) = id_val.as_str() { + let agent_json = store_dir.join(format!("{id}.json")); + let agent_md = store_dir.join(format!("{id}.md")); + let mut detail = json!({ "agent_id": id }); + + if agent_json.exists() { + if let Ok(data) = std::fs::read_to_string(&agent_json) { + if let Ok(parsed) = serde_json::from_str::(&data) { + detail["status"] = parsed.get("status").cloned().unwrap_or(json!("unknown")); + detail["name"] = parsed.get("name").cloned().unwrap_or(json!(null)); + detail["subagent_type"] = parsed.get("subagent_type").cloned().unwrap_or(json!(null)); + if let Some(completed) = parsed.get("completed_at") { + detail["completed_at"] = completed.clone(); + } + } + } + } else { + detail["status"] = json!("running"); + } + + if agent_md.exists() { + if let Ok(md_content) = std::fs::read_to_string(&agent_md) { + let summary = md_content.lines() + .skip_while(|line| !line.starts_with("## Result")) + .skip(1) + .take(10) + .collect::>() + .join("\n"); + if !summary.is_empty() { + detail["result_preview"] = json!(summary); + } + } + } + + match detail.get("status").and_then(|v| v.as_str()).unwrap_or("unknown") { + "completed" => completed_count += 1, + "failed" => failed_count += 1, + _ => running_count += 1, + } + agents_detail.push(detail); + } + } + + match action { + "status" => { + to_pretty_json(json!({ + "team_id": input.team_id, + "team_name": team_data.get("name"), + "status": if running_count > 0 { "running" } else if failed_count > 0 { "completed_with_failures" } else { "completed" }, + "progress": { + "total": agent_ids.len(), + "running": running_count, + "completed": completed_count, + "failed": failed_count, + }, + "agents": agents_detail, + })) + } + "summary" => { + let mut results: Vec = Vec::new(); + for id_val in &agent_ids { + if let Some(id) = id_val.as_str() { + let agent_md = store_dir.join(format!("{id}.md")); + let agent_json = store_dir.join(format!("{id}.json")); + let mut entry = json!({ "agent_id": id }); + + if let Ok(data) = std::fs::read_to_string(&agent_json) { + if let Ok(parsed) = serde_json::from_str::(&data) { + entry["status"] = parsed.get("status").cloned().unwrap_or(json!("unknown")); + entry["name"] = parsed.get("name").cloned().unwrap_or(json!(null)); + entry["subagent_type"] = parsed.get("subagent_type").cloned().unwrap_or(json!(null)); + } + } + if let Ok(md_content) = std::fs::read_to_string(&agent_md) { + entry["result"] = json!(md_content); + } + results.push(entry); + } + } + to_pretty_json(json!({ + "team_id": input.team_id, + "team_name": team_data.get("name"), + "status": if running_count > 0 { "still_running" } else if failed_count > 0 { "completed_with_failures" } else { "completed" }, + "progress": { + "total": agent_ids.len(), + "running": running_count, + "completed": completed_count, + "failed": failed_count, + }, + "results": results, + })) + } + "events" => { + let events_path = team_dir.join(format!("{}-events.jsonl", input.team_id)); + let mut events: Vec = Vec::new(); + if events_path.exists() { + if let Ok(content) = std::fs::read_to_string(&events_path) { + for line in content.lines() { + if let Ok(event) = serde_json::from_str::(line) { + events.push(event); + } + } + } + } + to_pretty_json(json!({ + "team_id": input.team_id, + "event_count": events.len(), + "events": events, + })) + } + other => Err(format!("unknown TeamStatus action: {other}. Use status, summary, or events")), + } +} + +/// Spawn a background thread that watches a team's agents and prints progress. +/// Prints agent completion/failure events to stderr. Exits when all agents are done. +fn spawn_team_watcher(team_id: &str, agent_ids: &[String]) { + let team_id = team_id.to_string(); + let agent_ids = agent_ids.to_vec(); + let thread_name = format!("claw-team-watch-{team_id}"); + + std::thread::Builder::new() + .name(thread_name) + .spawn(move || { + let store_dir = match agent_store_dir() { + Ok(d) => d, + Err(_) => return, + }; + let team_dir = store_dir.join("teams"); + let events_path = team_dir.join(format!("{team_id}-events.jsonl")); + + let mut completed: BTreeSet = BTreeSet::new(); + let mut failed: BTreeSet = BTreeSet::new(); + let total = agent_ids.len(); + + eprintln!("[team] {team_id}: watching {total} agents..."); + + loop { + let mut all_done = true; + for id in &agent_ids { + if completed.contains(id) || failed.contains(id) { + continue; + } + let agent_json = store_dir.join(format!("{id}.json")); + if let Ok(data) = std::fs::read_to_string(&agent_json) { + if let Ok(parsed) = serde_json::from_str::(&data) { + let status = parsed.get("status").and_then(|v| v.as_str()).unwrap_or("unknown"); + match status { + "completed" => { + completed.insert(id.clone()); + let name = parsed.get("name").and_then(|v| v.as_str()).unwrap_or(id); + let subagent_type = parsed.get("subagent_type").and_then(|v| v.as_str()).unwrap_or("unknown"); + let result_preview = get_agent_result_preview(&store_dir, id); + eprintln!("[team] {team_id}: done {name} ({subagent_type}) completed - {}/{total}{}", completed.len() + failed.len(), if result_preview.is_empty() { String::new() } else { format!(" - {}", &result_preview[..result_preview.len().min(120)]) }); + append_team_event(&events_path, &team_id, id, "completed", &name, Some(&result_preview)); + } + "failed" => { + failed.insert(id.clone()); + let name = parsed.get("name").and_then(|v| v.as_str()).unwrap_or(id); + let error = parsed.get("error").and_then(|v| v.as_str()).unwrap_or("unknown error"); + eprintln!("[team] {team_id}: FAIL {name} failed - {error}"); + append_team_event(&events_path, &team_id, id, "failed", &name, Some(error)); + } + _ => { + all_done = false; + } + } + } + } else { + all_done = false; + } + } + + if all_done { + break; + } + std::thread::sleep(std::time::Duration::from_secs(2)); + } + + eprintln!("[team] {team_id}: all agents finished - {}/{} completed, {}/{} failed", completed.len(), total, failed.len(), total); + append_team_event(&events_path, &team_id, "team", "finished", &team_id, Some(&format!("{}/{} completed", completed.len(), total))); + + let manifest_path = team_dir.join(format!("{team_id}.json")); + if let Ok(data) = std::fs::read_to_string(&manifest_path) { + if let Ok(mut parsed) = serde_json::from_str::(&data) { + if let Some(obj) = parsed.as_object_mut() { + obj.insert("status".to_string(), json!(if failed.is_empty() { "completed" } else { "completed_with_failures" })); + let _ = std::fs::write(&manifest_path, serde_json::to_string_pretty(&parsed).unwrap_or_default()); + } + } + } + }) + .ok(); +} + +fn get_agent_result_preview(store_dir: &std::path::Path, agent_id: &str) -> String { + let md_path = store_dir.join(format!("{agent_id}.md")); + if let Ok(content) = std::fs::read_to_string(&md_path) { + let lines: Vec<&str> = content.lines().collect(); + let start = lines.iter().position(|l| l.starts_with("## Result")).map(|i| i + 1).unwrap_or(0); + lines[start..].iter().take(5).cloned().collect::>().join(" ").chars().take(200).collect() + } else { + String::new() + } +} + +fn append_team_event(events_path: &std::path::Path, team_id: &str, agent_id: &str, event_type: &str, name: &str, detail: Option<&str>) { + let entry = json!({ + "timestamp": std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap_or_default().as_millis(), + "team_id": team_id, + "agent_id": agent_id, + "event": event_type, + "name": name, + "detail": detail, + }); + if let Ok(line) = serde_json::to_string(&entry) { + use std::io::Write; + if let Ok(mut file) = std::fs::OpenOptions::new().create(true).append(true).open(events_path) { + let _ = writeln!(file, "{line}"); + } + } +} + + fn run_team_delete(input: TeamDeleteInput) -> Result { match global_team_registry().delete(&input.team_id) { Ok(team) => to_pretty_json(json!({ @@ -3286,6 +3554,13 @@ struct AgentMessageInput { mark_read: Option, } +#[derive(Debug, Deserialize)] +struct TeamStatusInput { + team_id: String, + #[serde(default)] + action: Option, +} + #[derive(Debug, Deserialize)] struct CronCreateInput { schedule: String, @@ -4560,6 +4835,7 @@ fn allowed_tools_for_subagent(subagent_type: &str) -> BTreeSet { "ToolSearch", "Skill", "AgentMessage", + "TeamStatus", "StructuredOutput", ], "Plan" => vec![ @@ -4572,6 +4848,7 @@ fn allowed_tools_for_subagent(subagent_type: &str) -> BTreeSet { "Skill", "TodoWrite", "AgentMessage", + "TeamStatus", "StructuredOutput", "SendUserMessage", ], @@ -4585,6 +4862,7 @@ fn allowed_tools_for_subagent(subagent_type: &str) -> BTreeSet { "ToolSearch", "TodoWrite", "AgentMessage", + "TeamStatus", "StructuredOutput", "SendUserMessage", "PowerShell", @@ -4624,6 +4902,7 @@ fn allowed_tools_for_subagent(subagent_type: &str) -> BTreeSet { "NotebookEdit", "Sleep", "AgentMessage", + "TeamStatus", "SendUserMessage", "Config", "StructuredOutput",