pub mod advance;
mod assembler;
mod child_result;
mod drain;
pub mod driver;
mod model_call;
mod result_deposit;
mod staging;
mod step_commit;
pub mod stop_signal;
mod terminal;
mod tool_step;
mod tools;
mod transcript;
mod transfer;
pub use model_call::{RealSleeper, Sleeper};
pub(crate) use step_commit::remove_control_files;
pub use stop_signal::{flag as stop_flag, install as install_stop_handler};
use super::inbox::{self, Epitaph};
use super::step::{RESPONSE_FILE, STAGING_FILE, StepMeta, step_dir_rel};
use super::{Deps, Error};
use crate::config::manifest::RoleRules;
use crate::config::{Budgets, Model, RetryConfig, Workflow};
use assembler::assemble;
use brazen::Content;
use model_call::ModelCall;
use std::ffi::OsString;
use std::path::Path;
use step_commit::{
commit_dispatch, prepend_goal, read_branch_tip, spawn_branch, write_dispatch_files, write_meta,
write_request,
};
use tool_step::run_tool_calls;
const DEFAULT_MAX_TOKENS: u32 = 4096;
pub(super) struct Resolved<'a> {
pub(super) role: &'a str,
pub(super) model: &'a Model,
pub(super) provider_row: &'a str,
pub(super) tools: &'a [String],
pub(super) soul: String,
pub(super) binary: OsString,
pub(super) retry: RetryConfig,
pub(super) budgets: Budgets,
pub(super) workflow: &'a Workflow,
pub(super) manifest: Option<&'a RoleRules>,
pub(super) expect_handshake: bool,
}
pub(super) fn run_exchange(
repo: &Path,
user_message: &str,
resolved: &Resolved<'_>,
deps: &Deps<'_>,
) -> Result<String, Error> {
let ts = deps.clock.now_compact();
let short_id = deps.id_gen.short();
let conv_id = format!("{ts}-{short_id}");
let branch_name = conv_id.clone();
let worktree_path = crate::workspace::agent_worktree(repo, &conv_id);
let inbox = inbox::inbox_dir(repo, &conv_id);
let executor_lock = match inbox::try_acquire(&inbox).map_err(|source| Error::ExecutorLock {
path: inbox.clone(),
source,
})? {
Some(guard) => guard,
None => return Ok(branch_name),
};
spawn_branch(repo, &worktree_path, &conv_id, deps)?;
inbox::deposit(repo, &conv_id, inbox::USER_SENDER, user_message, deps.clock)?;
let system_with_goal = prepend_goal(user_message, &resolved.soul);
let call = ModelCall {
adapter: deps.adapter,
sleeper: deps.sleeper,
binary: &resolved.binary,
provider_row: resolved.provider_row,
retry: resolved.retry,
stop: deps.stop,
expect_handshake: resolved.expect_handshake,
};
let mut step_seq: u32 = 1;
let mut exhausted = false;
let mut stopped = false;
let mut seen;
loop {
if step_seq == 1 {
write_dispatch_files(&worktree_path, user_message, &resolved.soul)?;
commit_dispatch(&worktree_path, &conv_id, deps)?;
}
seen = drain::drain(&worktree_path, &inbox, &conv_id, deps.git)?.left;
child_result::interpret_pending(repo, &conv_id, &worktree_path, resolved.workflow, deps)?;
let commit_sha = read_branch_tip(&worktree_path, deps)?;
if stop_signal::stopped(deps.stop) {
stopped = true;
break;
}
if terminal::budget_exhausted(
repo,
&conv_id,
&branch_name,
&worktree_path,
&resolved.budgets,
deps,
)? {
exhausted = true;
break;
}
let messages = assemble(&worktree_path, resolved.manifest)?;
let tools = tools::compose(&worktree_path, resolved.role, resolved.tools, &messages)?;
let request = model_call::build_request(
&resolved.model.model_id,
&system_with_goal,
messages,
tools,
DEFAULT_MAX_TOKENS,
);
let request_value =
serde_json::to_value(&request).expect("CanonicalRequest is always serializable");
let step_dir_rel_str = step_dir_rel(&conv_id, step_seq);
write_request(repo, &step_dir_rel_str, &request_value)?;
let request_bytes =
serde_json::to_vec(&request).expect("CanonicalRequest is always serializable");
let started_at = deps.clock.now_iso8601();
let response_path = repo.join(&step_dir_rel_str).join(RESPONSE_FILE);
let call_outcome = model_call::run(&call, &request_bytes, &response_path);
if stop_signal::stopped(deps.stop) {
stopped = true;
break;
}
call_outcome?;
let ended_at = deps.clock.now_iso8601();
write_meta(
repo,
&step_dir_rel_str,
&StepMeta {
commit: commit_sha,
started_at,
ended_at,
},
)?;
let staging_path = repo.join(&step_dir_rel_str).join(STAGING_FILE);
let assistant_content = transcript::commit_assistant(
&worktree_path,
&conv_id,
&resolved.model.model_id,
&staging_path,
deps.git,
)?;
if !assistant_content
.iter()
.any(|b| matches!(b, Content::ToolUse { .. }))
{
let response = result_deposit::terminal_text(&assistant_content);
result_deposit::deposit_terminal(
repo,
&conv_id,
&worktree_path,
Epitaph::FinalResponse,
response.as_deref(),
deps,
)?;
break;
}
if run_tool_calls(
repo,
&worktree_path,
&conv_id,
resolved.role,
&step_dir_rel_str,
&assistant_content,
deps,
)? {
stopped = true;
break;
}
child_result::run_flush(repo, &conv_id, &worktree_path, resolved.workflow, deps)?;
step_seq += 1;
}
let epitaph = match (stopped, exhausted) {
(true, _) => Epitaph::Stopped,
(false, true) => Epitaph::BudgetExhausted,
(false, false) => Epitaph::FinalResponse,
};
terminal::conclude(
repo,
&conv_id,
epitaph,
resolved.workflow,
executor_lock,
&seen,
deps,
)?;
Ok(branch_name)
}