use super::super::{assembler, child_result, drain, driver, terminal, tool_step};
use super::AdvanceOutcome;
use crate::config::Event;
use crate::prompt::inbox::{self, Epitaph, ExecutorLock};
use crate::prompt::resolve::WorkerConfig;
use crate::prompt::step::{next_step_seq, step_dir_rel};
use crate::prompt::workflow_actions;
use crate::prompt::{Deps, Error};
use crate::workspace::hold;
use brazen::{Content, Message, Role};
use std::path::Path;
pub(super) enum Resumption {
Done(AdvanceOutcome),
Stale(ExecutorLock),
}
pub(super) fn resume(
workspace: &Path,
agent_id: &str,
mark: &hold::Held,
lock: ExecutorLock,
deps: &Deps<'_>,
resolve: &mut dyn FnMut() -> Result<WorkerConfig, Error>,
) -> Result<Resumption, Error> {
let seen = drain::seen_all(&inbox::inbox_dir(workspace, agent_id))?;
let worktree = crate::workspace::agent_worktree(workspace, agent_id);
if !worktree.exists() {
driver::release_then_reprobe(lock, workspace, agent_id, &seen, deps.launcher);
return Ok(Resumption::Done(AdvanceOutcome::NothingToDo));
}
let tail = assembler::transcript(&worktree)?;
let Some(window) = open_window(&tail).filter(|w| w.unpaired(&mark.tool_use_id)) else {
hold::clear(workspace, agent_id, deps.git).map_err(|source| Error::Git {
op: "stale hold mark clear",
source,
})?;
return Ok(Resumption::Stale(lock));
};
let cfg = resolve()?;
let resolved = cfg.as_resolved();
let step_seq = next_step_seq(workspace, agent_id)?.saturating_sub(1).max(1);
let step_dir_rel_str = step_dir_rel(agent_id, step_seq);
match tool_step::run_tool_calls(
workspace,
&worktree,
agent_id,
&resolved,
&step_dir_rel_str,
window.content,
deps,
)? {
tool_step::ToolWindow::Held => {
driver::release_then_reprobe(lock, workspace, agent_id, &seen, deps.launcher);
Ok(Resumption::Done(AdvanceOutcome::Held))
}
tool_step::ToolWindow::Stopped => {
terminal::conclude(
workspace,
agent_id,
Epitaph::Stopped,
&cfg.workflow,
lock,
&seen,
deps,
)?;
Ok(Resumption::Done(AdvanceOutcome::Terminal))
}
tool_step::ToolWindow::Completed => {
workflow_actions::run_step_hook(
&cfg.workflow,
Event::OnToolReturn,
&worktree,
agent_id,
deps.git,
)?;
child_result::run_flush(workspace, agent_id, &worktree, &cfg.workflow, deps)?;
Ok(Resumption::Done(AdvanceOutcome::ToolsPending(lock)))
}
}
}
struct OpenWindow<'a> {
content: &'a [Content],
resolved_ids: Vec<&'a str>,
}
impl OpenWindow<'_> {
fn unpaired(&self, id: &str) -> bool {
self.content
.iter()
.any(|b| matches!(b, Content::ToolUse { id: block_id, .. } if block_id == id))
&& !self.resolved_ids.contains(&id)
}
}
fn open_window(tail: &[Message]) -> Option<OpenWindow<'_>> {
let at = tail.iter().rposition(|m| m.role == Role::Assistant)?;
let resolved_ids = tail[at + 1..]
.iter()
.flat_map(|m| m.content.iter())
.filter_map(|b| match b {
Content::ToolResult { tool_use_id, .. } => Some(tool_use_id.as_str()),
_ => None,
})
.collect();
Some(OpenWindow {
content: &tail[at].content,
resolved_ids,
})
}
#[cfg(test)]
mod tests {
use super::*;
fn msg(role: Role, content: Vec<Content>) -> Message {
Message { role, content }
}
#[test]
fn open_window_is_total_over_arbitrary_tails() {
assert!(open_window(&[msg(Role::User, vec![Content::Text("hi".into())])]).is_none());
let tail = [
msg(
Role::Assistant,
vec![
Content::ToolUse {
id: "t1".into(),
name: "bash".into(),
input: serde_json::json!({}),
signature: None,
},
Content::ToolUse {
id: "t2".into(),
name: "bash".into(),
input: serde_json::json!({}),
signature: None,
},
],
),
msg(
Role::Tool,
vec![Content::ToolResult {
tool_use_id: "t1".into(),
content: vec![],
is_error: false,
}],
),
msg(Role::User, vec![Content::Text("mail".into())]),
];
let window = open_window(&tail).unwrap();
assert!(!window.unpaired("t1"), "t1 has a committed result");
assert!(window.unpaired("t2"), "t2 is the open frontier");
assert!(!window.unpaired("t9"), "t9 is not of this window");
}
}