use crate::agent::message_router::{self, AgentJob, MessageKind};
use crate::agent::run_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 description(&self) -> String {
let base = crate::prompt::load_prompt(&format!("tool/{}.md", self.name()));
if self.dispatch_mode.is_async() {
let async_note = crate::prompt::load_prompt("tool/implement_async.md");
format!("{base}\n\n{async_note}")
} else {
base
}
}
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());
}
run_sync_implement(ws, task, self.caller_role).await
}
}
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,
crate::tools::CoreJobArgs {
job_id: &job_id,
caller_role,
user_name: &user_name,
channel: &channel,
resume: false,
caller_agent_id: None,
fail_on_checkpoint_error: false,
},
)
.await
{
Ok(crate::tools::SyncCoreOutcome::DrainCut) => {
tracing::info!(
job = %job_id,
"Implement round cut short by drain — job stays launched for boot resume",
);
return None;
}
Ok(crate::tools::SyncCoreOutcome::Terminal(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_sync_implement(ws: &Workspace, task: &str, caller_role: Role) -> Result<String> {
crate::tools::SyncDurableCore::Implement
.run_sync_dispatch(ws, task, caller_role)
.await
}
#[expect(clippy::too_many_lines)]
pub(crate) async fn run_implement_with_job(
ws: &Workspace,
task: &str,
args: crate::tools::CoreJobArgs<'_>,
) -> anyhow::Result<crate::tools::SyncCoreOutcome> {
let crate::tools::CoreJobArgs {
job_id,
caller_role,
user_name,
channel,
resume,
caller_agent_id,
fail_on_checkpoint_error,
} = args;
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,
caller_agent_id,
)
.await?;
(coder_agent_id, None)
};
if let Some(outcome) = pre_done {
return Ok(crate::tools::SyncCoreOutcome::Terminal(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
{
if fail_on_checkpoint_error {
return Err(e.context("failed to checkpoint implement outcome"));
}
tracing::warn!(job = %job_id, error = %e, "Failed to checkpoint implement outcome");
}
if crate::shutdown::aborting() {
return Ok(crate::tools::SyncCoreOutcome::DrainCut);
}
Ok(crate::tools::SyncCoreOutcome::Terminal(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 (caller_role, caller, result) = match crate::tools::SyncDurableCore::Implement
.resume_sync_core(ws, job_id, false)
.await
{
Ok(crate::jobs::SyncResumeOutcome::Terminal(caller_role, caller, result)) => {
(caller_role, caller, result)
}
Ok(crate::jobs::SyncResumeOutcome::DrainCut | crate::jobs::SyncResumeOutcome::Gone) => {
return;
}
Err(e) => {
let Some((caller, caller_role)) = crate::jobs::resume_job_preamble(
&crate::session::store().conn,
job_id,
"Implement resume",
"Implement resume",
)
.await
else {
return;
};
let envelope = crate::jobs::complete_durable_job(
job_id,
build_async_implement_message(&Err(e)),
MessageKind::ImplementResult,
caller_role,
&caller.user_name,
&caller.channel,
&ws.name,
)
.await;
message_router::route(&crate::jobs::envelope_target(&envelope), envelope);
return;
}
};
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"
);
}
#[test]
fn implement_description_mode_keyed() {
let sync = ImplementTool::new(DispatchMode::Sync, Role::Coder);
assert!(
!sync.description().contains("dispatched asynchronously"),
"sync description must not carry the async note"
);
let async_tool = ImplementTool::new(DispatchMode::Async, Role::Coder);
let async_desc = async_tool.description();
assert!(
async_desc.contains("dispatched asynchronously"),
"async description must carry the async note"
);
assert!(
async_desc.contains("<implement-tool-result>"),
"async note names the result envelope"
);
}
#[tokio::test]
#[expect(clippy::await_holding_lock)] #[expect(clippy::too_many_lines)] async fn sync_implement_draincut_and_completion_resumes_durable_job() {
let _lock = crate::util::test::retry_tests_lock();
crate::util::test::init_management_test_stores().await;
let ws = test_ws("/tmp/test_ws_sync_implement");
let pin = "sync_implement_pin";
let conn = &crate::session::store().conn;
crate::util::test::seed_session_row(conn, pin, "user", "implement this").await;
crate::shutdown::drain_begin();
let tool = ImplementTool::new(DispatchMode::Sync, crate::Role::Engineer);
let res = crate::agent::CURRENT_TOOL_AGENT_ID
.scope(Some(pin.to_string()), async {
tool.execute(&ws, json!({"task": "implement task"})).await
})
.await;
crate::shutdown::drain_clear();
let err = res.expect_err("drain must cut the sync implement dispatch");
assert!(
err.downcast_ref::<crate::tools::CallSuspended>().is_some(),
"CallSuspended carrier expected: {err:#}"
);
let jobs = conn
.query(
"SELECT id, status, caller_agent_id FROM jobs WHERE caller_agent_id = ?1 AND kind = 'implement'",
crate::db::params![pin],
)
.await
.unwrap();
assert_eq!(jobs.len(), 1, "one launched implement job");
assert_eq!(jobs[0].get::<String>(1).unwrap(), "launched");
assert_eq!(
jobs[0].get::<String>(2).unwrap(),
pin,
"job is caller-owned by the session pin"
);
let job_id = jobs[0].get::<String>(0).unwrap();
let roster = crate::jobs::list_agents_for_job(conn, &job_id)
.await
.unwrap();
let coder_id = roster[0].agent_id.clone();
crate::jobs::write_agent_outcome(
conn,
&job_id,
&coder_id,
crate::jobs::RowStatus::Done,
Some("CODER_RESPONSE"),
)
.await
.unwrap();
let frame = crate::providers::reasoning::assistant_replay_payload(
Some(""),
&[crate::ToolCall {
id: "call_implement_h".to_string(),
name: "implement".to_string(),
arguments: json!({"task": "implement task"}),
}],
None,
)
.to_string();
crate::util::test::seed_session_row(conn, pin, "assistant", &frame).await;
let mut session = crate::session::Session::default();
session.init(pin).await.unwrap();
let pending = session
.pending_tool_calls()
.expect("dangling implement call");
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].id, "call_implement_h");
let outcome = crate::tools::SyncDurableCore::Implement
.resume_sync_core(&ws, &job_id, true)
.await
.unwrap();
let crate::jobs::SyncResumeOutcome::Terminal(_, _, result) = outcome else {
panic!("expected a terminal resume outcome");
};
let text = result.expect("resumed implement result");
assert_eq!(text, "CODER_RESPONSE");
crate::jobs::terminalize_job(conn, &job_id).await.unwrap();
session
.settle_tool_results(pin, &[("call_implement_h".to_string(), text.clone())], &[])
.await
.unwrap();
let jobs = conn
.query(
"SELECT id FROM jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(jobs.is_empty(), "resumed job must be terminalized");
let rows = conn
.query(
"SELECT id, role, content FROM sessions WHERE agent_id = ?1 ORDER BY id",
crate::db::params![pin],
)
.await
.unwrap();
assert_eq!(rows.len(), 3);
assert_eq!(rows[1].get::<String>(1).unwrap(), "assistant");
assert_eq!(rows[2].get::<String>(1).unwrap(), "tool");
let payload: crate::ToolResultPayload =
serde_json::from_str(&rows[2].get::<String>(2).unwrap()).unwrap();
assert_eq!(payload.tool_call_id, "call_implement_h");
assert_eq!(payload.content, "CODER_RESPONSE");
}
}