diff --git a/rust/crates/api/src/lib.rs b/rust/crates/api/src/lib.rs index 389212c5..d15731e2 100644 --- a/rust/crates/api/src/lib.rs +++ b/rust/crates/api/src/lib.rs @@ -29,7 +29,7 @@ pub use providers::openai_compat::{ pub use providers::{ detect_provider_kind, max_tokens_for_model, max_tokens_for_model_with_override, model_family_identity_for, model_family_identity_for_kind, model_token_limit, - ModelTokenLimit, provider_diagnostics_for_model, resolve_model_alias, ProviderDiagnostics, + provider_diagnostics_for_model, resolve_model_alias, ModelTokenLimit, ProviderDiagnostics, ProviderKind, }; pub use sse::{parse_frame, SseParser}; diff --git a/rust/crates/api/src/providers/mod.rs b/rust/crates/api/src/providers/mod.rs index c74429d1..8f846cbe 100644 --- a/rust/crates/api/src/providers/mod.rs +++ b/rust/crates/api/src/providers/mod.rs @@ -680,7 +680,8 @@ pub fn model_token_limit(model: &str) -> Option { context_window_tokens: 256_000, }), // Qwen models via DashScope / OpenAI-compat - "qwen3.6-35b-fast" | "qwen3-235b-a22b" | "qwen-max" | "qwen-plus" | "qwen-turbo" | "qwen-qwq" => Some(ModelTokenLimit { + "qwen3.6-35b-fast" | "qwen3-235b-a22b" | "qwen-max" | "qwen-plus" | "qwen-turbo" + | "qwen-qwq" => Some(ModelTokenLimit { max_output_tokens: 16_384, context_window_tokens: 131_072, }), diff --git a/rust/crates/commands/src/lib.rs b/rust/crates/commands/src/lib.rs index 2e5a9f55..a679a45d 100644 --- a/rust/crates/commands/src/lib.rs +++ b/rust/crates/commands/src/lib.rs @@ -815,7 +815,7 @@ const SLASH_COMMAND_SPECS: &[SlashCommandSpec] = &[ name: "team", aliases: &[], summary: "Manage agent teams", - argument_hint: Some("[list|create|delete]"), + argument_hint: Some("[on|off|status]"), resume_supported: true, }, SlashCommandSpec { @@ -5545,9 +5545,9 @@ mod tests { Ok(Some(SlashCommand::Team { action: None })) ); assert_eq!( - SlashCommand::parse("/team list"), + SlashCommand::parse("/team on"), Ok(Some(SlashCommand::Team { - action: Some("list".to_string()) + action: Some("on".to_string()) })) ); assert_eq!( diff --git a/rust/crates/runtime/src/config.rs b/rust/crates/runtime/src/config.rs index 7fdbda66..c076e6dc 100644 --- a/rust/crates/runtime/src/config.rs +++ b/rust/crates/runtime/src/config.rs @@ -1734,7 +1734,11 @@ fn parse_optional_model(root: &JsonValue) -> Option { /// falls back to the default model. fn parse_optional_subagent_model(root: &JsonValue) -> Option { root.as_object() - .and_then(|object| object.get("subagentModel").or_else(|| object.get("subagent_model"))) + .and_then(|object| { + object + .get("subagentModel") + .or_else(|| object.get("subagent_model")) + }) .and_then(JsonValue::as_str) .filter(|s| !s.trim().is_empty()) .map(|s| s.trim().to_string()) diff --git a/rust/crates/runtime/src/conversation.rs b/rust/crates/runtime/src/conversation.rs index 60b8f8b1..0ea0ee0a 100644 --- a/rust/crates/runtime/src/conversation.rs +++ b/rust/crates/runtime/src/conversation.rs @@ -284,10 +284,7 @@ where self } - pub fn with_turn_progress_reporter( - mut self, - reporter: Box, - ) -> Self { + pub fn with_turn_progress_reporter(mut self, reporter: Box) -> Self { self.turn_progress_reporter = Some(reporter); self } @@ -578,7 +575,13 @@ where // Phase 3: Post-hooks and session updates (sequential, original order). for p in &pending { // Capture progress data for the reporter. - let (progress_tool_name, progress_input, progress_output, progress_is_error, result_message) = if p.allowed { + let ( + progress_tool_name, + progress_input, + progress_output, + progress_is_error, + result_message, + ) = if p.allowed { let batch_result = &batch_results[batch_index]; batch_index += 1; let (mut output, mut is_error) = match &batch_result.result { @@ -621,16 +624,32 @@ where output, is_error, ); - (p.tool_name.clone(), p.effective_input.clone(), progress_output, is_error, result_message) + ( + p.tool_name.clone(), + p.effective_input.clone(), + progress_output, + is_error, + result_message, + ) } else { - let denied_output = merge_hook_feedback(&p.pre_hook_messages, p.deny_reason.clone().unwrap_or_default(), true); + let denied_output = merge_hook_feedback( + &p.pre_hook_messages, + p.deny_reason.clone().unwrap_or_default(), + true, + ); let result_message = ConversationMessage::tool_result( p.tool_use_id.clone(), p.tool_name.clone(), denied_output, true, ); - (p.tool_name.clone(), String::new(), String::new(), true, result_message) + ( + p.tool_name.clone(), + String::new(), + String::new(), + true, + result_message, + ) }; self.session .push_message(result_message.clone()) @@ -642,7 +661,13 @@ where } else { Ok(progress_output.as_str()) }; - reporter.on_tool_result(iterations, self.max_iterations, &progress_tool_name, &progress_input, report_result); + reporter.on_tool_result( + iterations, + self.max_iterations, + &progress_tool_name, + &progress_input, + report_result, + ); } tool_results.push(result_message); } diff --git a/rust/crates/rusty-claude-cli/src/main.rs b/rust/crates/rusty-claude-cli/src/main.rs index d6c0dbaa..25697604 100644 --- a/rust/crates/rusty-claude-cli/src/main.rs +++ b/rust/crates/rusty-claude-cli/src/main.rs @@ -7173,16 +7173,6 @@ fn run_repl( let _ = cli.set_model(Some(new_model)); } } - input::ReadOutcome::TeamToggle => { - let current = std::env::var("CLAWD_AGENT_TEAMS").unwrap_or_default(); - if current == "1" { - std::env::set_var("CLAWD_AGENT_TEAMS", "0"); - eprintln!("[team] Agent teams disabled"); - } else { - std::env::set_var("CLAWD_AGENT_TEAMS", "1"); - eprintln!("[team] Agent teams enabled"); - } - } } } @@ -8641,7 +8631,9 @@ impl LiveCli { if current == "1" { eprintln!("[team] Agent teams: ENABLED"); } else { - eprintln!("[team] Agent teams: DISABLED (use /team on or Ctrl+T to enable)"); + eprintln!( + "[team] Agent teams: DISABLED (use /team on or Ctrl+T to enable)" + ); } } "" => { @@ -8655,7 +8647,9 @@ impl LiveCli { eprintln!("[team] Agent teams enabled (TeamCreate now available)"); } } - other => eprintln!("[team] unknown action: {other}. Use: /team [on|off|status]"), + other => { + eprintln!("[team] unknown action: {other}. Use: /team [on|off|status]") + } } false } @@ -8664,7 +8658,10 @@ impl LiveCli { // Reload the model from config after wizard saves let cwd = std::env::current_dir().unwrap_or_default(); let config = runtime::ConfigLoader::default_for(&cwd).load().ok(); - if let Some(new_model) = config.as_ref().and_then(|c| c.provider().model().map(str::to_string)) { + if let Some(new_model) = config + .as_ref() + .and_then(|c| c.provider().model().map(str::to_string)) + { self.set_model(Some(new_model))?; } false @@ -14573,10 +14570,11 @@ impl ToolExecutor for CliToolExecutor { "LSP", "Agent", "AgentMessage", - "TeamStatus", + "TeamStatus", "TaskClaim", "AgentSuggestion", - "ContextRequest", "TaskGet", + "ContextRequest", + "TaskGet", "TaskList", "TaskOutput", "GitStatus", @@ -14593,11 +14591,9 @@ impl ToolExecutor for CliToolExecutor { // 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))) - { + 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(), @@ -14633,12 +14629,11 @@ impl ToolExecutor for CliToolExecutor { 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 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), + Ok(v) => registry.execute(&tool_name, &v).map_err(ToolError::new), Err(e) => Err(e), }; (idx, tool_use_id, tool_name, result) @@ -14646,14 +14641,16 @@ impl ToolExecutor for CliToolExecutor { } handles .into_iter() - .map(|h| h.join().unwrap_or_else(|_| { - ( - 0, - String::new(), - String::new(), - Err(ToolError::new("parallel thread panicked")), - ) - })) + .map(|h| { + h.join().unwrap_or_else(|_| { + ( + 0, + String::new(), + String::new(), + Err(ToolError::new("parallel thread panicked")), + ) + }) + }) .collect() }); diff --git a/rust/crates/tools/src/lib.rs b/rust/crates/tools/src/lib.rs index 9c65c32c..f65919df 100644 --- a/rust/crates/tools/src/lib.rs +++ b/rust/crates/tools/src/lib.rs @@ -7,10 +7,10 @@ use aspect_macros::aspect; use aspect_std::LoggingAspect; use api::{ - max_tokens_for_model, model_family_identity_for, model_token_limit, ModelTokenLimit, - resolve_model_alias, ApiError, ContentBlockDelta, InputContentBlock, InputMessage, - MessageRequest, MessageResponse, OutputContentBlock, ProviderClient, - StreamEvent as ApiStreamEvent, ToolChoice, ToolDefinition, ToolResultContentBlock, + max_tokens_for_model, model_family_identity_for, model_token_limit, resolve_model_alias, + ApiError, ContentBlockDelta, InputContentBlock, InputMessage, MessageRequest, MessageResponse, + OutputContentBlock, ProviderClient, StreamEvent as ApiStreamEvent, ToolChoice, ToolDefinition, + ToolResultContentBlock, }; use plugins::PluginTool; use reqwest::blocking::Client; @@ -1593,7 +1593,9 @@ fn execute_tool_with_enforcer( "AgentMessage" => from_value::(input).and_then(run_agent_message), "TeamStatus" => from_value::(input).and_then(run_team_status), "TaskClaim" => from_value::(input).and_then(run_task_claim), - "AgentSuggestion" => from_value::(input).and_then(run_agent_suggestion), + "AgentSuggestion" => { + from_value::(input).and_then(run_agent_suggestion) + } "ContextRequest" => from_value::(input).and_then(run_context_request), "CronCreate" => from_value::(input).and_then(run_cron_create), "CronDelete" => from_value::(input).and_then(run_cron_delete), @@ -1888,10 +1890,13 @@ fn run_team_create(input: TeamCreateInput) -> Result { return Err("Agent teams is disabled. Use /team on or Ctrl+T to enable.".to_string()); } - let team_id = format!("team-{}", std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap_or_default() - .as_nanos()); + let team_id = format!( + "team-{}", + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_nanos() + ); let output_dir = agent_store_dir()?; let team_dir = output_dir.join("teams"); std::fs::create_dir_all(&team_dir).map_err(|e| e.to_string())?; @@ -1899,9 +1904,23 @@ fn run_team_create(input: TeamCreateInput) -> Result { // Expand mode preset into tasks, or use manual tasks. // Default to "2x" when neither mode nor tasks are provided. let tasks = if let Some(mode) = &input.mode { - expand_team_mode(mode, input.prompt.as_deref().unwrap_or("Explore the codebase and report findings"), &team_id)? + expand_team_mode( + mode, + input + .prompt + .as_deref() + .unwrap_or("Explore the codebase and report findings"), + &team_id, + )? } else if input.tasks.is_empty() { - expand_team_mode("2x", input.prompt.as_deref().unwrap_or("Explore the codebase and report findings"), &team_id)? + expand_team_mode( + "2x", + input + .prompt + .as_deref() + .unwrap_or("Explore the codebase and report findings"), + &team_id, + )? } else { input.tasks.clone() }; @@ -1911,7 +1930,10 @@ fn run_team_create(input: TeamCreateInput) -> Result { for (i, task) in tasks.iter().enumerate() { let prompt = task.get("prompt").and_then(|v| v.as_str()).unwrap_or(""); - let description = task.get("description").and_then(|v| v.as_str()).unwrap_or(&input.name); + let description = task + .get("description") + .and_then(|v| v.as_str()) + .unwrap_or(&input.name); let subagent_type = task.get("subagent_type").and_then(|v| v.as_str()); let model_override = task.get("model").and_then(|v| v.as_str()); @@ -1919,12 +1941,19 @@ fn run_team_create(input: TeamCreateInput) -> Result { continue; } - let task_id = task.get("task_id").and_then(|v| v.as_str()).map(|s| s.to_string()); + let task_id = task + .get("task_id") + .and_then(|v| v.as_str()) + .map(|s| s.to_string()); let agent_input = AgentInput { description: description.to_string(), prompt: prompt.to_string(), subagent_type: subagent_type.map(|s| s.to_string()), - name: Some(format!("{}-agent-{}", slugify_agent_name(&input.name), i + 1)), + name: Some(format!( + "{}-agent-{}", + slugify_agent_name(&input.name), + i + 1 + )), model: model_override.map(|s| s.to_string()), team_id: Some(team_id.clone()), task_id, @@ -1961,8 +1990,11 @@ fn run_team_create(input: TeamCreateInput) -> Result { "created_at": iso8601_now(), }); let manifest_path = team_dir.join(format!("{team_id}.json")); - std::fs::write(&manifest_path, serde_json::to_string_pretty(&team_manifest).map_err(|e| e.to_string())?) - .map_err(|e| e.to_string())?; + std::fs::write( + &manifest_path, + serde_json::to_string_pretty(&team_manifest).map_err(|e| e.to_string())?, + ) + .map_err(|e| e.to_string())?; // Register in global registry let team = global_team_registry().create(&input.name, agent_ids.clone()); @@ -1989,10 +2021,11 @@ fn run_team_status(input: TeamStatusInput) -> Result { 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") + 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(); @@ -2012,9 +2045,11 @@ fn run_team_status(input: TeamStatusInput) -> Result { 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["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)); + 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(); } @@ -2026,7 +2061,8 @@ fn run_team_status(input: TeamStatusInput) -> Result { if agent_md.exists() { if let Ok(md_content) = std::fs::read_to_string(&agent_md) { - let summary = md_content.lines() + let summary = md_content + .lines() .skip_while(|line| !line.starts_with("## Result")) .skip(1) .take(10) @@ -2038,7 +2074,11 @@ fn run_team_status(input: TeamStatusInput) -> Result { } } - match detail.get("status").and_then(|v| v.as_str()).unwrap_or("unknown") { + match detail + .get("status") + .and_then(|v| v.as_str()) + .unwrap_or("unknown") + { "completed" => completed_count += 1, "failed" => failed_count += 1, _ => running_count += 1, @@ -2287,14 +2327,33 @@ fn get_agent_result_preview(store_dir: &std::path::Path, agent_id: &str) -> Stri 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() + 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>) { +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, @@ -2305,13 +2364,16 @@ fn append_team_event(events_path: &std::path::Path, team_id: &str, agent_id: &st }); 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) { + 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 { // Delete from disk-based team storage let team_dir = agent_store_dir()?.join("teams"); @@ -2381,8 +2443,11 @@ fn run_agent_message(input: AgentMessageInput) -> Result { "message": msg, "timestamp": ts, }); - std::fs::write(&msg_file, serde_json::to_string_pretty(&entry).map_err(|e| e.to_string())?) - .map_err(|e| e.to_string())?; + std::fs::write( + &msg_file, + serde_json::to_string_pretty(&entry).map_err(|e| e.to_string())?, + ) + .map_err(|e| e.to_string())?; to_pretty_json(json!({ "action": "sent", "to": target, @@ -2404,9 +2469,7 @@ fn run_agent_message(input: AgentMessageInput) -> Result { let entries: Vec<_> = std::fs::read_dir(&inbox_dir) .map_err(|e| e.to_string())? .filter_map(|e| e.ok()) - .filter(|e| { - e.path().extension().is_some_and(|ext| ext == "json") - }) + .filter(|e| e.path().extension().is_some_and(|ext| ext == "json")) .collect(); for entry in &entries { if let Ok(content) = std::fs::read_to_string(entry.path()) { @@ -2443,9 +2506,11 @@ fn run_agent_message(input: AgentMessageInput) -> Result { return Err(format!("team {team_id} not found")); } 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") + &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(); @@ -2457,7 +2522,9 @@ fn run_agent_message(input: AgentMessageInput) -> Result { let mut sent_to: Vec = Vec::new(); for id_val in agent_ids { if let Some(id) = id_val.as_str() { - if id == sender { continue; } + if id == sender { + continue; + } let inbox_dir = mailbox_dir.join(id); std::fs::create_dir_all(&inbox_dir).map_err(|e| e.to_string())?; let msg_file = inbox_dir.join(format!("msg-{ts}-{sender}.json")); @@ -2467,8 +2534,11 @@ fn run_agent_message(input: AgentMessageInput) -> Result { "timestamp": ts, "team_id": team_id, }); - std::fs::write(&msg_file, serde_json::to_string_pretty(&entry).map_err(|e| e.to_string())?) - .map_err(|e| e.to_string())?; + std::fs::write( + &msg_file, + serde_json::to_string_pretty(&entry).map_err(|e| e.to_string())?, + ) + .map_err(|e| e.to_string())?; sent_to.push(id.to_string()); } } @@ -2535,7 +2605,9 @@ fn run_task_claim(input: TaskClaimInput) -> Result { "count": claims.len(), })) } - other => Err(format!("unknown TaskClaim action: {other}. Use claim, release, or list")), + other => Err(format!( + "unknown TaskClaim action: {other}. Use claim, release, or list" + )), } } @@ -2557,8 +2629,11 @@ fn run_agent_suggestion(input: AgentSuggestionInput) -> Result { }); let filename = format!("suggestion-{agent_id}-{ts}.json"); let path = suggestions_dir.join(&filename); - std::fs::write(&path, serde_json::to_string_pretty(&entry).map_err(|e| e.to_string())?) - .map_err(|e| e.to_string())?; + std::fs::write( + &path, + serde_json::to_string_pretty(&entry).map_err(|e| e.to_string())?, + ) + .map_err(|e| e.to_string())?; to_pretty_json(json!({ "action": "suggestion", "file": filename, @@ -2579,7 +2654,11 @@ fn run_context_request(input: ContextRequestInput) -> Result { match std::fs::read_to_string(&path) { Ok(content) => { let truncated = if content.len() > 10000 { - format!("{}... (truncated, {} bytes total)", &content[..10000], content.len()) + format!( + "{}... (truncated, {} bytes total)", + &content[..10000], + content.len() + ) } else { content }; @@ -2610,7 +2689,15 @@ fn run_context_request(input: ContextRequestInput) -> Result { for symbol in &input.symbols { let cwd = std::env::current_dir().map_err(|e| e.to_string())?; let output = std::process::Command::new("grep") - .args(["-rn", "--include=*.rs", "--include=*.ts", "--include=*.js", "--include=*.py", symbol, "."]) + .args([ + "-rn", + "--include=*.rs", + "--include=*.ts", + "--include=*.js", + "--include=*.py", + symbol, + ".", + ]) .current_dir(&cwd) .output() .map_err(|e| e.to_string())?; @@ -2646,7 +2733,11 @@ fn expand_team_mode(mode: &str, base_prompt: &str, team_id: &str) -> Result 4, "5x" | "xlarge" => 5, "6x" | "mega" => 6, - other => return Err(format!("unknown team mode '{other}'. Use 1x-6x or tiny/small/medium/large/xlarge/mega")), + other => { + return Err(format!( + "unknown team mode '{other}'. Use 1x-6x or tiny/small/medium/large/xlarge/mega" + )) + } }; let short_team_id = &team_id[team_id.len().saturating_sub(8)..]; let roles: &[&str] = &["Explore", "Plan", "Verification"]; @@ -2721,8 +2812,11 @@ fn claim_task(task_id: &str, agent_id: &str, team_id: &str) -> Result Ok(true), Err(_) => { @@ -2755,7 +2849,8 @@ fn list_claims(team_id: Option<&str>) -> Vec { if path.extension().map_or(false, |e| e == "lock") { if let Ok(content) = std::fs::read_to_string(&path) { if let Ok(v) = serde_json::from_str::(&content) { - if team_id.map_or(true, |tid| v.get("team_id").map_or(false, |t| t == tid)) { + if team_id.map_or(true, |tid| v.get("team_id").map_or(false, |t| t == tid)) + { claims.push(v); } } @@ -2778,7 +2873,12 @@ impl TeamInboxReporter { fn new(team_id: String, agent_id: String, agent_name: String) -> Self { let inbox_dir = agent_mailbox_dir().join("team").join(&team_id); let _ = std::fs::create_dir_all(&inbox_dir); - Self { team_id, agent_id, agent_name, inbox_dir } + Self { + team_id, + agent_id, + agent_name, + inbox_dir, + } } } @@ -2812,10 +2912,9 @@ impl TurnProgressReporter for TeamInboxReporter { "max_iterations": max_iterations, "timestamp": ts, }); - let msg_file = self.inbox_dir.join(format!( - "tp-{}-{}-{ts}.json", - self.agent_id, iteration - )); + let msg_file = self + .inbox_dir + .join(format!("tp-{}-{}-{ts}.json", self.agent_id, iteration)); if let Ok(line) = serde_json::to_string(&entry) { let _ = std::fs::write(&msg_file, line); } @@ -2830,13 +2929,19 @@ impl TurnProgressReporter for TeamInboxReporter { .output(); if diff_check.map_or(true, |o| !o.status.success()) { let _ = std::process::Command::new("git") - .args(["commit", "-m", &format!("agent {} progress: iteration {iteration}", self.agent_id)]) + .args([ + "commit", + "-m", + &format!("agent {} progress: iteration {iteration}", self.agent_id), + ]) .output(); } } // Check for kill signal from team lead - for entry in std::fs::read_dir(&self.inbox_dir).unwrap_or_else(|_| std::fs::read_dir(".").unwrap()) { + for entry in + std::fs::read_dir(&self.inbox_dir).unwrap_or_else(|_| std::fs::read_dir(".").unwrap()) + { if let Ok(e) = entry { let name = e.file_name(); let name_str = name.to_string_lossy(); @@ -5352,14 +5457,25 @@ fn setup_agent_worktree(agent_id: &str) -> Result { // Create worktree on a new branch let branch_name = format!("agent/{agent_id}"); let output = std::process::Command::new("git") - .args(["worktree", "add", worktree_dir.to_str().unwrap_or(""), "-b", &branch_name]) + .args([ + "worktree", + "add", + worktree_dir.to_str().unwrap_or(""), + "-b", + &branch_name, + ]) .current_dir(&cwd) .output() .map_err(|e| e.to_string())?; if !output.status.success() { // Branch might already exist, try with existing branch let output2 = std::process::Command::new("git") - .args(["worktree", "add", worktree_dir.to_str().unwrap_or(""), &branch_name]) + .args([ + "worktree", + "add", + worktree_dir.to_str().unwrap_or(""), + &branch_name, + ]) .current_dir(&cwd) .output() .map_err(|e| e.to_string())?; @@ -5390,7 +5506,12 @@ fn teardown_agent_worktree(agent_id: &str, worktree_path: &std::path::Path) -> R } // Remove the worktree let _ = std::process::Command::new("git") - .args(["worktree", "remove", worktree_path.to_str().unwrap_or(""), "--force"]) + .args([ + "worktree", + "remove", + worktree_path.to_str().unwrap_or(""), + "--force", + ]) .current_dir(&cwd) .output(); // Delete the branch @@ -5405,8 +5526,7 @@ fn run_agent_job(job: &AgentJob) -> Result<(), String> { // Claim task if task_id is set (prevents duplicate work) if let Some(ref task_id) = job.task_id { if let Some(ref team_id) = job.team_id { - let claimed = claim_task(task_id, &job.manifest.agent_id, team_id) - .unwrap_or(false); + let claimed = claim_task(task_id, &job.manifest.agent_id, team_id).unwrap_or(false); if !claimed { return Err(format!("task {task_id} already claimed by another agent")); } @@ -5732,8 +5852,11 @@ fn post_agent_completion_to_team_inbox( "timestamp": ts, }); let msg_file = mailbox_dir.join(format!("{}-{ts}.json", manifest.agent_id)); - std::fs::write(&msg_file, serde_json::to_string_pretty(&entry).map_err(|e| e.to_string())?) - .map_err(|e| e.to_string()) + std::fs::write( + &msg_file, + serde_json::to_string_pretty(&entry).map_err(|e| e.to_string())?, + ) + .map_err(|e| e.to_string()) } const MIN_LANE_SUMMARY_WORDS: usize = 7; @@ -10038,8 +10161,8 @@ mod tests { subagent_type: Some("Explore".to_string()), name: Some("ship-audit".to_string()), model: None, - team_id: None, - task_id: None, + team_id: None, + task_id: None, }, move |job| { *captured_for_spawn @@ -10121,8 +10244,8 @@ mod tests { subagent_type: Some("Explore".to_string()), name: Some("complete-task".to_string()), model: Some("claude-sonnet-4-6".to_string()), - team_id: None, - task_id: None, + team_id: None, + task_id: None, }, |job| { persist_agent_terminal_state( @@ -10180,8 +10303,8 @@ mod tests { subagent_type: Some("Verification".to_string()), name: Some("fail-task".to_string()), model: None, - team_id: None, - task_id: None, + team_id: None, + task_id: None, }, |job| { persist_agent_terminal_state( @@ -10229,8 +10352,8 @@ mod tests { subagent_type: Some("Explore".to_string()), name: Some("summary-floor".to_string()), model: None, - team_id: None, - task_id: None, + team_id: None, + task_id: None, }, |job| { persist_agent_terminal_state( @@ -10326,8 +10449,8 @@ mod tests { subagent_type: Some("Verification".to_string()), name: Some("review-lane".to_string()), model: None, - team_id: None, - task_id: None, + team_id: None, + task_id: None, }, |job| { persist_agent_terminal_state( @@ -10488,8 +10611,8 @@ mod tests { subagent_type: Some("Explore".to_string()), name: Some("cron-closeout".to_string()), model: None, - team_id: None, - task_id: None, + team_id: None, + task_id: None, }, |job| { persist_agent_terminal_state( @@ -10531,8 +10654,8 @@ mod tests { subagent_type: None, name: Some("spawn-error".to_string()), model: None, - team_id: None, - task_id: None, + team_id: None, + task_id: None, }, |_| Err(String::from("thread creation failed")), )