use crate::agent::message_router::{self, AgentJob, MessageKind};
use crate::agent::{run_agent, run_default_agent};
use crate::session::analyze_agent_id;
use crate::tools::Tool;
use crate::tools::analyze::DispatchMode;
use crate::{Role, Workspace};
use anyhow::{Result, anyhow};
use async_trait::async_trait;
use futures_util::FutureExt;
use serde_json::json;
pub struct ImplementTool {
dispatch_mode: DispatchMode,
pub caller_role: Role,
}
impl ImplementTool {
#[must_use]
pub const fn new(dispatch_mode: DispatchMode, caller_role: Role) -> Self {
Self {
dispatch_mode,
caller_role,
}
}
}
#[async_trait]
impl Tool for ImplementTool {
fn name(&self) -> &'static str {
"implement"
}
fn parameters_schema(&self) -> serde_json::Value {
super::tool_params_schema(
&json!({
"task": {
"type": "string",
"description": "The implementation task to delegate to the coder sub-agent"
}
}),
&["task"],
)
}
async fn execute(&self, ws: &Workspace, args: serde_json::Value) -> Result<String> {
let task = super::get_str(&args, "task")?;
if self.dispatch_mode.is_async() {
let ws = ws.clone();
let task = task.to_string();
let caller_role = self.caller_role;
let user_name = crate::agent::CURRENT_TOOL_USER_NAME
.try_with(String::clone)
.unwrap_or_default();
let channel = crate::agent::CURRENT_TOOL_CHANNEL
.try_with(String::clone)
.unwrap_or_default();
tokio::spawn(async move {
let round = std::panic::AssertUnwindSafe(async {
dispatch_durable_implement(
&ws,
&task,
caller_role,
user_name.clone(),
channel.clone(),
)
.await
})
.catch_unwind()
.await;
let envelope = match round {
Ok(Some(envelope)) => envelope,
Ok(None) => {
return;
}
Err(panic) => {
let panic = crate::util::panic_message(&*panic);
tracing::error!(panic = %panic, "implement round dispatch panicked");
AgentJob {
content: build_async_implement_message(&Err(anyhow!(
"implement round dispatch panicked: {panic}"
))),
workspace_name: ws.name.clone(),
user_name,
channel,
kind: MessageKind::ImplementResult,
role: caller_role,
reply_target: None,
pending_job_id: None,
}
}
};
if crate::shutdown::aborting() {
return;
}
message_router::route(&crate::jobs::envelope_target(&envelope), envelope);
});
return Ok("Sub-agent dispatched. Results will follow shortly.".to_string());
}
let agent_id = analyze_agent_id(&ws.name, Role::Coder.as_str());
let parent_key = crate::agent::CURRENT_TOOL_PARENT_KEY
.try_with(std::clone::Clone::clone)
.unwrap_or(None);
let parent_label = crate::agent::CURRENT_TOOL_PARENT_LABEL
.try_with(std::clone::Clone::clone)
.unwrap_or(None);
let (agent, response) = run_default_agent(
&agent_id,
Role::Coder,
ws,
task,
None,
parent_key,
parent_label,
)
.await;
if let Some(response) = response {
Ok(response)
} else if agent.is_cancelled() || crate::shutdown::aborting() {
anyhow::bail!("Sub-agent cancelled");
} else {
anyhow::bail!(
"Sub-agent failed: {}",
agent.failure.as_deref().unwrap_or("unknown error")
);
}
}
}
enum ImplementRunOutcome {
Result(anyhow::Result<String>),
DrainCut,
}
async fn dispatch_durable_implement(
ws: &Workspace,
task: &str,
caller_role: Role,
user_name: String,
channel: String,
) -> Option<AgentJob> {
let job_id = crate::generate_id();
let result =
match run_implement_with_job(ws, task, &job_id, caller_role, &user_name, &channel, false)
.await
{
Ok(ImplementRunOutcome::DrainCut) => {
tracing::info!(
job = %job_id,
"Implement round cut short by drain — job stays launched for boot resume",
);
return None;
}
Ok(ImplementRunOutcome::Result(result)) => result,
Err(e) => Err(e),
};
let envelope = crate::jobs::complete_durable_job(
&job_id,
build_async_implement_message(&result),
MessageKind::ImplementResult,
caller_role,
&user_name,
&channel,
&ws.name,
)
.await;
Some(envelope)
}
async fn run_implement_with_job(
ws: &Workspace,
task: &str,
job_id: &str,
caller_role: Role,
user_name: &str,
channel: &str,
resume: bool,
) -> anyhow::Result<ImplementRunOutcome> {
let (coder_agent_id, pre_done) = if resume {
let rows = crate::jobs::list_agents_for_job(&crate::session::store().conn, job_id).await?;
let row = rows
.first()
.ok_or_else(|| anyhow!("Implement resume: no coder roster row"))?;
let pre_done = (row.status == crate::jobs::RowStatus::Done.as_str())
.then(|| row.outcome.clone())
.flatten();
(row.agent_id.clone(), pre_done)
} else {
let suffix = crate::generate_suffix();
let coder_agent_id = analyze_agent_id(&ws.name, Role::Coder.as_str()) + &suffix;
let agents = vec![crate::jobs::NewAgent {
agent_id: coder_agent_id.clone(),
kind: crate::jobs::AgentKind::Coder,
idx: Some(0),
task: task.to_string(),
}];
crate::jobs::spawn_job(
&crate::session::store().conn,
job_id,
task,
&ws.name,
user_name,
channel,
caller_role,
&agents,
&crate::jobs::SpawnChild::Implement,
)
.await?;
(coder_agent_id, None)
};
if let Some(outcome) = pre_done {
return Ok(ImplementRunOutcome::Result(Ok(outcome)));
}
let parent_key = crate::agent::CURRENT_TOOL_PARENT_KEY
.try_with(std::clone::Clone::clone)
.unwrap_or(None);
let parent_label = crate::agent::CURRENT_TOOL_PARENT_LABEL
.try_with(std::clone::Clone::clone)
.unwrap_or(None);
let has_session = resume && crate::session::store().has_content(&coder_agent_id).await;
let (agent, response) = run_agent(
coder_agent_id.clone(),
Role::Coder,
ws,
None,
if has_session { "" } else { task },
user_name.to_string(),
channel.to_string(),
false,
None,
resume,
None,
parent_key,
parent_label,
)
.await;
let (status, outcome) = match &response {
Some(r) => (crate::jobs::RowStatus::Done, r.clone()),
None => (
crate::jobs::RowStatus::Failed,
agent.failure_reason("coder produced no response"),
),
};
if let Err(e) = crate::jobs::write_agent_outcome(
&crate::session::store().conn,
job_id,
&coder_agent_id,
status,
Some(&outcome),
)
.await
{
tracing::warn!(job = %job_id, error = %e, "Failed to checkpoint implement outcome");
}
if crate::shutdown::aborting() {
return Ok(ImplementRunOutcome::DrainCut);
}
Ok(ImplementRunOutcome::Result(match response {
Some(r) => Ok(r),
None => Err(anyhow!(
"Sub-agent failed: {}",
agent.failure_reason("unknown error")
)),
}))
}
pub(crate) async fn resume_implement_round(job_id: &str, ws: &Workspace) {
let Some((caller, caller_role)) = crate::jobs::resume_job_preamble(
&crate::session::store().conn,
job_id,
"Implement resume",
"Implement resume",
)
.await
else {
return;
};
let result = match run_implement_with_job(
ws,
&caller.task,
job_id,
caller_role,
&caller.user_name,
&caller.channel,
true,
)
.await
{
Ok(ImplementRunOutcome::Result(result)) => result,
Ok(ImplementRunOutcome::DrainCut) => {
tracing::info!(job = %job_id, "Implement resume aborted — job stays for next boot");
return;
}
Err(e) => Err(e),
};
if crate::shutdown::aborting() {
tracing::info!(job = %job_id, "Implement resume aborted — job stays for next boot");
return;
}
let envelope = crate::jobs::complete_durable_job(
job_id,
build_async_implement_message(&result),
MessageKind::ImplementResult,
caller_role,
&caller.user_name,
&caller.channel,
&ws.name,
)
.await;
message_router::route(&crate::jobs::envelope_target(&envelope), envelope);
}
fn build_async_implement_message(result: &anyhow::Result<String>) -> String {
crate::tools::analyze::build_async_result_envelope(result, "implement-tool-result")
}
#[cfg(test)]
mod tests {
use super::*;
use crate::workspace::test_ws;
use serde_json::json;
#[tokio::test]
async fn test_implement_missing_args() {
let tool = ImplementTool::new(DispatchMode::Sync, Role::Coder);
let ws = test_ws("/tmp/test_ws");
let result = tool.execute(&ws, json!({})).await;
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("Missing required field: task"),
"Should mention missing task"
);
}
}