use std::ffi::{OsStr, OsString};
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::sync::{Arc, Mutex};
use clap::{CommandFactory, Parser};
use crate::controls::NodeControls;
use crate::error::{Error, Result};
use crate::event::{Envelope, Labels};
use crate::executor::{CancellationToken, DispatchRequest, LocalExecutor, WorkspaceSpec};
use crate::ledger::RunPaths;
use crate::lifecycle::{Drafted, Undrafted, PR_AUTHOR_PERSONA};
pub(crate) const PR_AUTHOR_GRAPH_FLAG: &str = "--pr-author-graph";
pub(crate) const NO_DRAFT_FLAG: &str = "--no-draft";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Verb {
PublishBranch,
Recover,
}
impl Verb {
fn onevcs(self) -> &'static str {
match self {
Self::PublishBranch => "publish-branch",
Self::Recover => "recover",
}
}
}
#[derive(Debug, Default, PartialEq, Eq)]
struct Split {
forwarded: Vec<OsString>,
graph: Option<OsString>,
drafting: Drafting,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
enum Drafting {
#[default]
WhereWanted,
Declined,
}
pub(crate) fn land(verb: Verb, args: Vec<OsString>) -> Result<i32> {
let split = split(verb, args)?;
let argv = [OsString::from("onevcs"), OsString::from(verb.onevcs())]
.into_iter()
.chain(split.forwarded);
let mut cli = match onevcs::cli::Cli::try_parse_from(argv) {
Ok(cli) => cli,
Err(refused) => {
let _ = refused.print();
return Ok(refused.exit_code());
}
};
if let Some(body) = drafted_for(verb, &cli, &split.graph, split.drafting)? {
match &mut cli.command {
onevcs::cli::Command::PublishBranch(args) => args.body = Some(body),
onevcs::cli::Command::Recover(args) => args.body = Some(body),
_ => {}
}
}
Ok(i32::from(onevcs::run(&cli)))
}
fn split(verb: Verb, args: Vec<OsString>) -> Result<Split> {
let valued = valued_options(verb);
let mut split = Split::default();
let mut args = args.into_iter();
while let Some(arg) = args.next() {
let text = arg.to_str().unwrap_or_default();
if text == "--" {
split.forwarded.push(arg);
split.forwarded.extend(args);
break;
} else if text == NO_DRAFT_FLAG {
split.drafting = Drafting::Declined;
} else if text == PR_AUTHOR_GRAPH_FLAG {
split.graph = Some(args.next().ok_or_else(|| {
Error::Invalid(format!(
"{PR_AUTHOR_GRAPH_FLAG} was given no value: name the graph the body is \
drafted by, or pass {NO_DRAFT_FLAG}"
))
})?);
} else if let Some(graph) = text.strip_prefix(&format!("{PR_AUTHOR_GRAPH_FLAG}=")) {
split.graph = Some(graph.into());
} else if valued.iter().any(|option| option == text) {
split.forwarded.push(arg);
split.forwarded.extend(args.next());
} else {
split.forwarded.push(arg);
}
}
Ok(split)
}
fn valued_options(verb: Verb) -> Vec<String> {
let cli = onevcs::cli::Cli::command();
let Some(command) = cli.find_subcommand(verb.onevcs()) else {
return Vec::new();
};
command
.get_arguments()
.filter(|argument| !argument.is_positional() && argument.get_action().takes_values())
.flat_map(|argument| {
argument
.get_long()
.map(|long| format!("--{long}"))
.into_iter()
.chain(argument.get_short().map(|short| format!("-{short}")))
})
.collect()
}
fn drafted_for(
verb: Verb,
cli: &onevcs::cli::Cli,
graph: &Option<OsString>,
drafting: Drafting,
) -> Result<Option<String>> {
let (branch, repo, brought, narrowed) = match (&cli.command, verb) {
(onevcs::cli::Command::PublishBranch(args), Verb::PublishBranch) => (
&args.branch,
&args.repo,
args.body.is_some() || args.body_file.is_some(),
args.policy,
),
(onevcs::cli::Command::Recover(args), Verb::Recover) => (
&args.branch,
&args.repo,
args.body.is_some() || args.body_file.is_some(),
None,
),
_ => {
return Err(Error::Invalid(format!(
"`onevcs {}` parsed as a different command",
verb.onevcs()
)))
}
};
if drafting == Drafting::Declined || brought {
return Ok(None);
}
let named = repo.to_string_lossy();
let destination = crate::destination::resolve(&named);
let publication = narrowed.or_else(|| {
destination
.as_ref()
.ok()
.and_then(|destination| destination.publication.clone().ok())
});
if publication.is_some_and(|policy| !crate::destination::opens_a_change_request(policy)) {
return Ok(None);
}
let Some(graph) = graph else {
return Err(Error::Invalid(format!(
"`onepipeline {}` would draft the change request's body for {branch}, and no graph \
to draft it with was named: pass {PR_AUTHOR_GRAPH_FLAG} <PATH>, bring a body with \
--body or --body-file, or pass {NO_DRAFT_FLAG}",
match verb {
Verb::PublishBranch => "publish-branch",
Verb::Recover => "repo-recover",
}
)));
};
let graph = graph.to_str().ok_or_else(|| {
Error::Invalid(format!(
"{PR_AUTHOR_GRAPH_FLAG} names a path that is not UTF-8: {}",
graph.to_string_lossy()
))
})?;
let here = std::env::current_dir().map_err(|error| {
Error::Invalid(format!(
"the working directory {PR_AUTHOR_GRAPH_FLAG} resolves against cannot be read: \
{error}"
))
})?;
let graph = crate::driver::resolve_graph(graph, &here)?;
let checkout = if repo.is_dir() {
Ok(repo.clone())
} else {
destination.map(|destination| destination.resolved.publication_checkout)
};
let drafted = match checkout {
Ok(checkout) => draft(&graph, &checkout, branch),
Err(why) => Drafted::Undrafted(Undrafted::Dispatch(format!(
"the drafting dispatch could not start: {why}"
))),
};
Ok(match drafted {
Drafted::Body(body) => Some(body),
Drafted::Undrafted(ending) => {
eprintln!(
"onepipeline: {} ({}), so {branch} lands with no body",
ending.why(),
ending.ending()
);
None
}
})
}
fn draft(graph: &str, checkout: &Path, branch: &str) -> Drafted {
let could_not_start = |why: String| {
Drafted::Undrafted(Undrafted::Dispatch(format!(
"the drafting dispatch could not start: {why}"
)))
};
let scratch = match Scratch::new(checkout) {
Ok(scratch) => scratch,
Err(why) => return could_not_start(why),
};
let tree = match scratch.cut(branch) {
Ok(tree) => tree,
Err(why) => return could_not_start(why),
};
let paths = RunPaths::under(&scratch.dir, "draft");
if let Err(error) = std::fs::create_dir_all(&paths.dir) {
return could_not_start(format!(
"{} could not be created: {error}",
paths.dir.display()
));
}
let mut deaths = Vec::new();
let (drafted, handle) = crate::lifecycle::draft(
&LocalExecutor,
DispatchRequest {
graph: oneagentgraph::config::ConfigRef(graph.to_owned()),
task: crate::lifecycle::out_of_band_drafting_task(branch, &tree.base, &tree.commits),
labels: Labels {
persona: Some(PR_AUTHOR_PERSONA.to_owned()),
..Labels::default()
},
controls: NodeControls::default(),
workspace: WorkspaceSpec::Path(tree.worktree),
cancel: CancellationToken::new(),
attempt: std::num::NonZeroU32::MIN,
},
&paths,
&mut |envelope| deaths.extend(death(&envelope)),
);
drop(handle);
if matches!(drafted, Drafted::Undrafted(Undrafted::Dispatch(_))) {
for died in &deaths {
eprintln!("onepipeline: {died}");
}
}
drafted
}
fn death(envelope: &Envelope) -> Option<String> {
if envelope.source != crate::event::Source::Agentgraph
|| envelope.kind.0 != oneagentgraph::event::EventKind::MemberDied.as_str()
{
return None;
}
let member = envelope.labels.member.as_deref().unwrap_or("a member");
let stated: Vec<String> = ["rule", "cause", "detail"]
.into_iter()
.filter_map(|field| {
let value = envelope.payload.get(field)?.as_str()?;
(!value.is_empty()).then(|| format!("{field}={}", crate::views::one_line(value)))
})
.collect();
let truncated = envelope
.payload
.get("truncated")
.and_then(serde_json::Value::as_bool)
== Some(true);
Some(match (stated.is_empty(), truncated) {
(true, _) => {
format!("{member} died, and oneagentgraph named no rule, cause or detail for it")
}
(false, false) => format!("{member} died: {}", stated.join(" ")),
(false, true) => format!(
"{member} died: {} (oneagentgraph truncated this detail)",
stated.join(" ")
),
})
}
struct Scratch {
dir: PathBuf,
checkout: PathBuf,
left: Arc<Mutex<Leftover>>,
_interrupts: crate::sys::OnInterrupt,
}
struct Leftover {
checkout: PathBuf,
made: Made,
}
enum Made {
Nothing,
Dir(PathBuf),
Worktree { dir: PathBuf, worktree: PathBuf },
}
impl Leftover {
fn remove(&mut self) {
let dir = match std::mem::replace(&mut self.made, Made::Nothing) {
Made::Nothing => return,
Made::Dir(dir) => dir,
Made::Worktree { dir, worktree } => {
let removed = git_in(
&self.checkout,
[
OsStr::new("worktree"),
OsStr::new("remove"),
OsStr::new("--force"),
worktree.as_os_str(),
],
);
if let Err(why) = removed {
eprintln!(
"onepipeline: the drafting worktree at {} could not be removed from {}: \
{why}; remove it with `git -C {} worktree remove --force {}`",
worktree.display(),
self.checkout.display(),
self.checkout.display(),
worktree.display()
);
}
dir
}
};
if let Err(error) = std::fs::remove_dir_all(&dir) {
eprintln!(
"onepipeline: the drafting turn's directory {} could not be removed: \
{error}; remove it by hand",
dir.display()
);
}
}
}
fn lock_leftover(leftover: &Mutex<Leftover>) -> std::sync::MutexGuard<'_, Leftover> {
leftover
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
struct Tree {
base: String,
commits: String,
worktree: PathBuf,
}
impl Scratch {
fn new(checkout: &Path) -> std::result::Result<Self, String> {
static MINTED: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let leftover = Arc::new(Mutex::new(Leftover {
checkout: checkout.to_path_buf(),
made: Made::Nothing,
}));
let interrupts = crate::sys::on_interrupt({
let leftover = Arc::clone(&leftover);
move || {
eprintln!(
"onepipeline: interrupted while drafting; removing the drafting worktree"
);
lock_leftover(&leftover).remove();
}
})?;
let root = std::env::temp_dir();
for _ in 0..64 {
let dir = root.join(format!(
"onepipeline-draft-{}-{}",
crate::sys::pid(),
MINTED.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
));
let mut left = lock_leftover(&leftover);
match std::fs::create_dir(&dir) {
Ok(()) => {
left.made = Made::Dir(dir.clone());
drop(left);
return Ok(Self {
dir,
checkout: checkout.to_path_buf(),
left: leftover,
_interrupts: interrupts,
});
}
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {}
Err(error) => {
return Err(format!("{} could not be created: {error}", dir.display()))
}
}
}
Err(format!(
"no directory for the drafting turn could be created under {}",
root.display()
))
}
fn cut(&self, branch: &str) -> std::result::Result<Tree, String> {
let head = [
format!("refs/heads/{branch}"),
format!("refs/remotes/origin/{branch}"),
]
.into_iter()
.find(|reference| {
self.git(&["show-ref", "--verify", "--quiet", reference])
.is_ok()
})
.ok_or_else(|| {
format!(
"the checkout {} has no branch '{branch}', locally or on its origin",
self.checkout.display()
)
})?;
let base = self.base().ok_or_else(|| {
format!(
"the base {branch} lands on is unknown: {} names no default branch for its \
origin",
self.checkout.display()
)
})?;
let commits = self.git(&[
"log",
"--reverse",
"--format=commit %h%n%n%B",
&format!("{base}..{head}"),
"--",
])?;
let worktree = self.dir.join("branch");
let mut left = lock_leftover(&self.left);
git_in(
&self.checkout,
[
OsStr::new("worktree"),
OsStr::new("add"),
OsStr::new("--detach"),
worktree.as_os_str(),
OsStr::new(&head),
],
)?;
left.made = Made::Worktree {
dir: self.dir.clone(),
worktree: worktree.clone(),
};
drop(left);
Ok(Tree {
base,
commits,
worktree,
})
}
fn base(&self) -> Option<String> {
if let Ok(tracked) = self.git(&[
"symbolic-ref",
"--quiet",
"--short",
"refs/remotes/origin/HEAD",
]) {
return Some(tracked);
}
let advertised = self
.git(&["ls-remote", "--symref", "origin", "HEAD"])
.ok()
.and_then(|said| {
said.lines().find_map(|line| {
let (named, _) = line.strip_prefix("ref: refs/heads/")?.split_once('\t')?;
Some(format!("origin/{named}"))
})
})
.filter(|named| {
self.git(&[
"show-ref",
"--verify",
"--quiet",
&format!("refs/remotes/{named}"),
])
.is_ok()
});
if advertised.is_some() {
return advertised;
}
let tracked = self
.git(&[
"for-each-ref",
"--format=%(refname:short)",
"refs/remotes/origin",
])
.ok()?;
let mut candidates = tracked.lines().filter(|named| *named != "origin/HEAD");
match (candidates.next(), candidates.next()) {
(Some(only), None) => Some(only.to_owned()),
_ => None,
}
}
fn git(&self, args: &[&str]) -> std::result::Result<String, String> {
git_in(&self.checkout, args.iter().map(OsStr::new))
}
}
impl Drop for Scratch {
fn drop(&mut self) {
lock_leftover(&self.left).remove();
}
}
fn git_in<'a>(
dir: &Path,
args: impl IntoIterator<Item = &'a OsStr>,
) -> std::result::Result<String, String> {
let args: Vec<&OsStr> = args.into_iter().collect();
let named = || {
let words: Vec<String> = args
.iter()
.map(|arg| arg.to_string_lossy().into_owned())
.collect();
format!("`git {}` in {}", words.join(" "), dir.display())
};
let output = Command::new("git")
.args(&args)
.current_dir(dir)
.stdin(Stdio::null())
.output()
.map_err(|error| format!("{} could not be run: {error}", named()))?;
if !output.status.success() {
return Err(format!(
"{} exited {}: {}",
named(),
output.status.code().unwrap_or(-1),
crate::views::one_line(&String::from_utf8_lossy(&output.stderr))
));
}
String::from_utf8(output.stdout)
.map(|said| said.trim().to_owned())
.map_err(|error| format!("{} answered bytes that are not UTF-8: {error}", named()))
}
#[cfg(test)]
mod tests {
use super::*;
fn words(args: &[&str]) -> Vec<OsString> {
args.iter().map(OsString::from).collect()
}
#[test]
fn only_this_crates_two_flags_are_taken_out_and_the_rest_keep_their_order() {
let taken = split(
Verb::PublishBranch,
words(&[
"feature",
"--title",
"--no-draft",
"--repo=service",
"--no-draft",
"--policy",
"change-open",
"--pr-author-graph",
"graphs/pr-author.yaml",
]),
)
.expect("the line splits");
assert_eq!(
taken,
Split {
forwarded: words(&[
"feature",
"--title",
"--no-draft",
"--repo=service",
"--policy",
"change-open"
]),
graph: Some("graphs/pr-author.yaml".into()),
drafting: Drafting::Declined,
}
);
let after = split(
Verb::Recover,
words(&["--pr-author-graph=g.yaml", "--", "--no-draft"]),
)
.expect("the line splits");
assert_eq!(after.forwarded, words(&["--", "--no-draft"]));
assert_eq!(after.graph, Some("g.yaml".into()));
assert_eq!(after.drafting, Drafting::WhereWanted);
assert!(split(Verb::Recover, words(&["b", "--pr-author-graph"])).is_err());
}
#[test]
fn the_valued_options_are_read_off_the_siblings_own_parser() {
let valued = valued_options(Verb::PublishBranch);
for option in ["--repo", "--title", "--policy", "--body", "--body-file"] {
assert!(
valued.iter().any(|known| known == option),
"{option}: {valued:?}"
);
}
let valued = valued_options(Verb::Recover);
assert!(
!valued.iter().any(|known| known == "--policy"),
"{valued:?}"
);
}
#[test]
fn a_member_death_is_said_in_oneagentgraphs_own_classification() {
let envelope = |died: oneagentgraph::event::MemberDied| -> Envelope {
serde_json::from_value(serde_json::json!({
"v": 1, "ts": "2026-01-01T00:00:00Z", "stream": "s", "seq": 1,
"source": "agentgraph",
"kind": oneagentgraph::event::EventKind::MemberDied.as_str(),
"labels": {"member": "author"},
"payload": serde_json::to_value(died).expect("a death serialises"),
}))
.expect("an envelope")
};
let died = |rule: oneagentgraph::member::Rule,
cause: oneagentgraph::event::Cause,
detail: &str,
truncated: bool| oneagentgraph::event::MemberDied {
rule: rule.as_str().to_owned(),
cause,
detail: detail.to_owned(),
truncated,
exit_code: None,
disposition: None,
stderr_tail: None,
candidates: Vec::new(),
};
assert_eq!(
death(&envelope(died(
oneagentgraph::member::Rule::Unstartable,
oneagentgraph::event::Cause::Spawn,
"no such file",
false
))),
Some("author died: rule=unstartable cause=spawn detail=no such file".to_owned())
);
assert_eq!(
death(&envelope(died(
oneagentgraph::member::Rule::ProviderFailure,
oneagentgraph::event::Cause::Unclassified,
"cut",
true
))),
Some(
"author died: rule=provider-failure cause=unclassified detail=cut \
(oneagentgraph truncated this detail)"
.to_owned()
)
);
let mut bare = envelope(died(
oneagentgraph::member::Rule::Unstartable,
oneagentgraph::event::Cause::Spawn,
"",
false,
));
bare.payload.clear();
assert_eq!(
death(&bare),
Some("author died, and oneagentgraph named no rule, cause or detail for it".to_owned())
);
}
}