use std::collections::{BTreeMap, BTreeSet};
use std::sync::mpsc::Sender;
use onevcs::SessionRequest;
use crate::engine::{self, Message, Settlement};
use crate::event::{Envelope, Labels};
use crate::executor::{DispatchRequest, Executor, WorkspaceSpec};
use crate::graph::NodeStatus;
use crate::plan::{Node, NodeKind, Step};
pub const PR_AUTHOR_PERSONA: &str = "pr-author";
pub fn execute(
executor: &dyn Executor,
run: &str,
round: u64,
default_graph: &str,
node: &Node,
cancel: &crate::executor::CancellationToken,
tx: &Sender<Message>,
) -> Settlement {
let Some(request) = crate::vcs::request_for(node) else {
return Settlement {
detail: Some("a lifecycle node needs a repo".into()),
..Settlement::plain(&node.id, NodeStatus::Failed, Some("invalid-node"))
};
};
let declared_steps = node.steps.is_some();
let steps = match ordered_steps(node) {
Ok(steps) => steps,
Err(reason) => {
return Settlement {
detail: Some(reason),
..Settlement::plain(&node.id, NodeStatus::Failed, Some("invalid-node"))
}
}
};
let mut session: Option<String> = None;
let mut stream: Option<crate::vcs::Follower> = None;
let mut worktree: Option<std::path::PathBuf> = None;
let whose = engine::dispatch_labels(run, round, &node.id, None, node.persona.as_deref());
let mut branch: Option<String> = node.branch.clone();
let mut completed: Vec<String> = node
.resume
.as_ref()
.map(|resume| resume.completed_steps.clone())
.unwrap_or_default();
for step in &steps {
if declared_steps && completed.iter().any(|id| id == &step.id) {
continue;
}
if step.kind == NodeKind::Human {
return Settlement {
branch,
completed_steps: completed,
..Settlement::plain(&node.id, NodeStatus::Waiting, None)
};
}
if step.expects_no_diff {
continue;
}
let request = SessionRequest {
branch: branch.clone().or_else(|| request.branch.clone()),
..request.clone()
};
let workspace = match &worktree {
Some(dir) => WorkspaceSpec::Path(dir.clone()),
None => WorkspaceSpec::VcsSession(request.clone()),
};
let graph = engine::node_graph(
step.agent_graph.as_ref().or(node.agent_graph.as_ref()),
default_graph,
);
let build = || DispatchRequest {
graph: graph.clone(),
task: step.rendered_task(node.context.as_deref()),
labels: engine::dispatch_labels(
run,
round,
&node.id,
declared_steps.then_some(step.id.as_str()),
step.persona.as_deref(),
),
workspace: workspace.clone(),
cancel: cancel.clone(),
};
let drained = engine::attempt(executor, &node.id, engine::Role::Worker, cancel, tx, &build);
session = drained.session.or(session);
branch = drained.branch.or(branch);
if stream.is_none() {
if let Some(token) = &session {
worktree = crate::vcs::worktree_of(token);
stream = crate::vcs::follow(token, relay_into(tx, whose.clone()));
}
}
if drained.settlement.status != NodeStatus::Done {
end_session(stream, tx, session.as_deref(), &whose);
return Settlement {
branch,
completed_steps: completed,
..drained.settlement
};
}
if declared_steps {
completed.push(step.id.clone());
}
}
let Some(token) = session else {
return Settlement {
branch,
..Settlement::plain(&node.id, NodeStatus::Done, Some("no-changes"))
};
};
let settlement = publish(
executor,
run,
round,
default_graph,
node,
worktree.as_deref(),
cancel,
tx,
&token,
branch,
);
end_session(stream, tx, Some(&token), &whose);
settlement
}
#[allow(
clippy::too_many_arguments,
reason = "publication needs the dispatch context (executor, run, round, node, cancel, \
stream) as well as what the steps left behind (the session token, its branch, \
and the worktree they worked in); the first six are the node's own dispatch \
identity and bundling them would only move the same list one indirection away"
)]
fn publish(
executor: &dyn Executor,
run: &str,
round: u64,
default_graph: &str,
node: &Node,
worktree: Option<&std::path::Path>,
cancel: &crate::executor::CancellationToken,
tx: &Sender<Message>,
token: &str,
branch: Option<String>,
) -> Settlement {
let title = node.title.clone().unwrap_or_else(|| {
draft_title(
executor,
run,
round,
default_graph,
node,
worktree,
cancel,
tx,
)
});
match crate::vcs::publish(token, node.merge_policy, Some(&title)) {
Ok(published) => {
if let onevcs::PublishOutcome::Failed { reason, .. } = &published.outcome {
return Settlement {
branch,
detail: Some(format!("onevcs: {reason}")),
..Settlement::plain(&node.id, NodeStatus::Failed, Some("publication-failed"))
};
}
let labels =
engine::dispatch_labels(run, round, &node.id, None, node.persona.as_deref());
let _ = tx.send(Message::Event(Box::new(crate::vcs::published_event(
&published, &labels,
))));
Settlement {
branch: branch.or_else(|| Some(published.branch.clone())),
change_url: crate::vcs::change_url(&published.outcome),
outcome: Some(crate::vcs::outcome_of(&published.outcome).to_owned()),
..Settlement::plain(&node.id, NodeStatus::Done, None)
}
}
Err(error) => Settlement {
branch,
detail: Some(error.to_string()),
..Settlement::plain(&node.id, NodeStatus::Failed, Some("publication-failed"))
},
}
}
#[allow(
clippy::too_many_arguments,
reason = "the draft is a dispatch inside one lifecycle execution and needs that execution's \
executor, labels, resolved graph, workspace, cancellation, and event stream"
)]
fn draft_title(
executor: &dyn Executor,
run: &str,
round: u64,
default_graph: &str,
node: &Node,
worktree: Option<&std::path::Path>,
cancel: &crate::executor::CancellationToken,
tx: &Sender<Message>,
) -> String {
let fallback = deterministic_title(node);
let Some(request) = crate::vcs::request_for(node) else {
return fallback;
};
let workspace = match worktree {
Some(dir) => WorkspaceSpec::Path(dir.to_path_buf()),
None => WorkspaceSpec::VcsSession(SessionRequest {
branch: node.branch.clone().or(request.branch.clone()),
..request
}),
};
let dispatch = executor.dispatch(DispatchRequest {
graph: engine::node_graph(None, default_graph),
task: format!(
"Read this branch's diff and write the change request's title and body, \
following the repository's own template. The task this branch delivered:\n\n{}",
node.rendered_task()
),
labels: engine::dispatch_labels(run, round, &node.id, None, Some(PR_AUTHOR_PERSONA)),
workspace,
cancel: cancel.clone(),
});
let Ok(mut handle) = dispatch else {
return fallback;
};
let mut drafted = None;
for envelope in handle.events() {
let Ok(envelope) = envelope else { continue };
if let Some(title) = envelope.payload.get("title").and_then(|v| v.as_str()) {
drafted = Some(title.to_string());
}
let _ = tx.send(Message::Event(Box::new(envelope)));
}
match handle.wait() {
Ok(outcome) if outcome.succeeded => drafted.unwrap_or(fallback),
_ => fallback,
}
}
pub fn deterministic_title(node: &Node) -> String {
node.title
.clone()
.unwrap_or_else(|| format!("chore: {}", node.id))
}
fn relay_into(tx: &Sender<Message>, node: Labels) -> Box<dyn Fn(Envelope) + Send> {
let tx = tx.clone();
Box::new(move |mut envelope| {
stamp(&mut envelope.labels, &node);
let _ = tx.send(Message::Event(Box::new(envelope)));
})
}
fn stamp(labels: &mut Labels, known: &Labels) {
labels.run_id = labels.run_id.take().or_else(|| known.run_id.clone());
labels.round = labels.round.or(known.round);
labels.node = labels.node.take().or_else(|| known.node.clone());
labels.persona = labels.persona.take().or_else(|| known.persona.clone());
}
fn end_session(
stream: Option<crate::vcs::Follower>,
tx: &Sender<Message>,
token: Option<&str>,
node: &Labels,
) {
close(token);
let followed_through = stream.and_then(crate::vcs::Follower::finish);
relay_session_events(tx, token, node, followed_through);
}
fn relay_session_events(
tx: &Sender<Message>,
token: Option<&str>,
node: &Labels,
followed_through: Option<u64>,
) {
let Some(token) = token else { return };
let relay = relay_into(tx, node.clone());
for envelope in beyond(crate::vcs::events(token), followed_through) {
relay(envelope);
}
}
fn beyond(envelopes: Vec<Envelope>, followed_through: Option<u64>) -> Vec<Envelope> {
envelopes
.into_iter()
.filter(|envelope| !followed_through.is_some_and(|seq| envelope.seq <= seq))
.collect()
}
fn close(token: Option<&str>) {
if let Some(token) = token {
let _ = crate::vcs::session_close(token);
}
}
pub fn ordered_steps(node: &Node) -> std::result::Result<Vec<Step>, String> {
let Some(steps) = &node.steps else {
return Ok(vec![Step {
id: node.id.clone(),
kind: node.kind,
task: node.task.clone(),
persona: node.persona.clone(),
deps: Vec::new(),
done_when: node.done_when.clone(),
max_turns: node.max_turns,
expects_no_diff: node.expects_no_diff,
executor: node.executor.clone(),
agent_graph: node.agent_graph.clone(),
}]);
};
let by_id: BTreeMap<&str, &Step> = steps.iter().map(|s| (s.id.as_str(), s)).collect();
let mut settled: BTreeSet<&str> = BTreeSet::new();
let mut order: Vec<Step> = Vec::new();
while order.len() < steps.len() {
let mut progressed = false;
for step in steps {
if settled.contains(step.id.as_str()) {
continue;
}
if step
.deps
.iter()
.all(|dep| settled.contains(dep.as_str()) || !by_id.contains_key(dep.as_str()))
{
settled.insert(step.id.as_str());
order.push(step.clone());
progressed = true;
}
}
if !progressed {
return Err(format!(
"node '{}': its steps have a dependency cycle",
node.id
));
}
}
Ok(order)
}
#[cfg(test)]
mod tests {
use super::*;
fn step(id: &str, deps: &[&str]) -> Step {
Step {
id: id.into(),
persona: Some("engineer".into()),
task: Some("## What\ndo it".into()),
deps: deps.iter().map(|d| (*d).to_string()).collect(),
..Step::default()
}
}
fn lifecycle(steps: Option<Vec<Step>>) -> Node {
Node {
id: "service".into(),
repo: Some("owner/repo".into()),
persona: steps.is_none().then(|| "engineer".into()),
task: steps.is_none().then(|| "## What\nship".into()),
steps,
..Node::default()
}
}
#[test]
fn steps_run_serially_in_topological_order() {
let node = lifecycle(Some(vec![
step("publish", &["review"]),
step("implement", &[]),
step("review", &["implement"]),
]));
let order: Vec<String> = ordered_steps(&node)
.expect("the steps order")
.into_iter()
.map(|s| s.id)
.collect();
assert_eq!(order, vec!["implement", "review", "publish"]);
}
#[test]
fn steps_with_a_cycle_are_reported_rather_than_run_in_some_order() {
let node = lifecycle(Some(vec![step("a", &["b"]), step("b", &["a"])]));
let message = ordered_steps(&node).unwrap_err();
assert!(message.contains("dependency cycle"), "{message}");
}
#[test]
fn a_lifecycle_node_with_no_steps_is_one_implicit_step() {
let node = lifecycle(None);
let steps = ordered_steps(&node).expect("one implicit step");
assert_eq!(steps.len(), 1);
assert_eq!(steps[0].id, "service");
assert_eq!(steps[0].persona.as_deref(), Some("engineer"));
}
#[test]
fn a_deterministic_title_is_the_same_every_time() {
let node = lifecycle(None);
assert_eq!(deterministic_title(&node), "chore: service");
let titled = Node {
title: Some("feat: ship the thing".into()),
..node
};
assert_eq!(deterministic_title(&titled), "feat: ship the thing");
}
#[test]
fn relays_only_what_the_follow_did_not() {
let wrote = |seq: u64| Envelope {
v: crate::event::ENVELOPE_VERSION,
ts: "2026-01-01T00:00:00.000Z".into(),
stream: "s-1".into(),
seq,
source: crate::event::Source::Vcs,
kind: crate::event::EventKind("session-closed".into()),
labels: Labels::default(),
payload: serde_json::Map::new(),
artifacts: Vec::new(),
};
let stream: Vec<Envelope> = (1..=4).map(wrote).collect();
let tail = beyond(stream.clone(), Some(3));
assert_eq!(tail.iter().map(|e| e.seq).collect::<Vec<_>>(), vec![4]);
assert!(beyond(stream.clone(), Some(4)).is_empty());
assert_eq!(
beyond(stream, None)
.iter()
.map(|e| e.seq)
.collect::<Vec<_>>(),
vec![1, 2, 3, 4]
);
}
#[test]
fn a_step_ordering_ignores_a_dependency_on_something_outside_the_node() {
let node = lifecycle(Some(vec![step("only", &["elsewhere"])]));
let steps = ordered_steps(&node).expect("an outside reference does not deadlock");
assert_eq!(steps.len(), 1);
}
}