use std::collections::{BTreeMap, BTreeSet};
use std::io::{Read, Write};
use std::path::Path;
use std::process::Stdio;
use std::sync::mpsc;
use std::time::{Duration, Instant};
use oneagentgraph::config::{ConfigRef, GraphConfig, JudgeSide, Member};
use oneagentgraph::resolve::{self, Resolver};
use crate::cli::DEFAULT_DISPATCH_ENV_HOOK_TIMEOUT_SECONDS;
use crate::error::{Error, Result};
use crate::hooks::{self, HOOK_ENV, POLL, RUN_ID_ENV, RUN_ROOT_ENV};
use crate::ledger::{LaunchRecord, RunPaths};
use crate::sys;
pub(crate) const HOOK: &str = "dispatch-env";
pub(crate) const NODE_ID_ENV: &str = "ONEPIPELINE_NODE_ID";
const DOCUMENT_VERSION: u64 = 1;
const DOCUMENT_MEMBERS: [&str; 2] = ["version", "env"];
const MAX_STDOUT_BYTES: u64 = 1 << 20;
const STDOUT_GRACE: Duration = Duration::from_secs(2);
pub(crate) struct Launching<'a> {
pub paths: &'a RunPaths,
pub record: &'a LaunchRecord,
pub node: &'a str,
pub graph: &'a ConfigRef,
pub sets: &'a [String],
pub own: &'a [&'a str],
}
pub(crate) fn run_and_check(launching: &Launching<'_>) -> Result<Vec<(String, String)>> {
let Some(command) = launching.record.dispatch_env_hook() else {
return Ok(Vec::new());
};
let node = launching.node;
let added = run(launching, command).map_err(|ending| {
Error::Refused(format!(
"the {HOOK} hook '{command}' refused the launch of node '{node}': {ending}; its \
stderr is kept in {}",
hooks::hook_log(launching.paths, HOOK).display()
))
})?;
let mut refreshed = process_env();
refreshed.extend(added.iter().cloned());
for name in launching.own {
refreshed.entry((*name).to_string()).or_default();
}
validate(launching.graph, launching.sets, &refreshed).map_err(|why| {
Error::Refused(format!(
"the launch of node '{node}' was refused after the {HOOK} hook '{command}' ran: \
{why}"
))
})?;
Ok(added)
}
fn process_env() -> BTreeMap<String, String> {
std::env::vars_os()
.filter_map(|(key, value)| Some((key.into_string().ok()?, value.into_string().ok()?)))
.collect()
}
fn run(
launching: &Launching<'_>,
command: &str,
) -> std::result::Result<Vec<(String, String)>, String> {
let paths = launching.paths;
let record = launching.record;
let log = hooks::hook_log(paths, HOOK);
let opened = log
.parent()
.map_or(Ok(()), std::fs::create_dir_all)
.and_then(|()| {
std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&log)
});
let mut stderr = match opened {
Ok(stderr) => stderr,
Err(error) => {
return Err(format!(
"could-not-start: its log {} could not be opened: {error}",
log.display()
))
}
};
let _ = writeln!(
stderr,
"onepipeline: {HOOK} hook for node '{}'",
launching.node
);
let mut spawning = std::process::Command::new(command);
if !record.dir.as_os_str().is_empty() {
spawning.current_dir(&record.dir);
}
let stderr = match stderr.try_clone() {
Ok(stderr) => stderr,
Err(error) => return Err(format!("could-not-start: {error}")),
};
spawning
.env(HOOK_ENV, HOOK)
.env(RUN_ID_ENV, &paths.run)
.env(RUN_ROOT_ENV, hooks::run_root(paths))
.env(NODE_ID_ENV, launching.node)
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(stderr);
let mut child = spawning
.spawn()
.map_err(|error| format!("could-not-start: {error}"))?;
let (printed_tx, printed_rx) = mpsc::channel();
let mut stdout = child.stdout.take();
let reading = std::thread::Builder::new()
.name(format!("{HOOK}-hook-stdout"))
.spawn(move || {
let mut printed = Vec::new();
let read = match stdout.as_mut() {
Some(stdout) => stdout
.take(MAX_STDOUT_BYTES + 1)
.read_to_end(&mut printed)
.and_then(|_| std::io::copy(stdout, &mut std::io::sink()))
.map(|_| ()),
None => Ok(()),
};
let _ = printed_tx.send(read.map(|()| printed));
});
if let Err(error) = reading {
let _ = sys::stop(child.id(), sys::Stop::Now);
let _ = child.kill();
let _ = child.wait();
return Err(format!(
"could-not-start: no thread could be started to read its stdout: {error}"
));
}
let timeout = record.dispatch_env_hook_timeout();
let deadline = hooks::deadline_after(timeout);
let ended = loop {
match child.try_wait() {
Ok(Some(status)) => break Ok(status),
Ok(None) if deadline.is_none_or(|deadline| Instant::now() < deadline) => {}
waited => {
let _ = sys::stop(child.id(), sys::Stop::Now);
let _ = child.kill();
let _ = child.wait();
break Err(match waited {
Ok(_) => format!(
"timeout: still running after {timeout} seconds, so its process tree \
was ended"
),
Err(error) => format!("it could not be waited for: {error}"),
});
}
}
std::thread::sleep(POLL);
};
let status = ended?;
if !status.success() {
return Err(match status.code() {
Some(code) => format!("exit {code}"),
None => "it ended with no exit code".to_string(),
});
}
let printed = match printed_rx.recv_timeout(STDOUT_GRACE) {
Ok(Ok(printed)) => printed,
Ok(Err(error)) => return Err(format!("its stdout could not be read: {error}")),
Err(_) => {
return Err(
"malformed: its stdout was still held open after it exited, so what it \
printed is not one whole document"
.to_string(),
)
}
};
if printed.len() as u64 > MAX_STDOUT_BYTES {
return Err(format!(
"malformed: it printed more than {MAX_STDOUT_BYTES} bytes, which is not one \
environment document"
));
}
Document::read(&printed).map_err(|what| format!("malformed: {what}"))
}
struct Document;
impl Document {
fn read(printed: &[u8]) -> std::result::Result<Vec<(String, String)>, String> {
let document: serde_json::Value = serde_json::from_slice(printed).map_err(|error| {
format!("its stdout is not one JSON document: {error}")
})?;
let Some(members) = document.as_object() else {
return Err("its document is not a JSON object".to_string());
};
for member in members.keys() {
if !DOCUMENT_MEMBERS.contains(&member.as_str()) {
return Err(format!(
"its document carries a top-level member `{member}`, and the only members \
are `version` and `env`"
));
}
}
match members.get("version").and_then(serde_json::Value::as_u64) {
Some(DOCUMENT_VERSION) => {}
Some(other) => {
return Err(format!(
"its document is version {other}, and the only version is {DOCUMENT_VERSION}"
))
}
None => {
return Err(format!(
"its document names no `version`, which is {DOCUMENT_VERSION}"
))
}
}
let Some(env) = members.get("env").and_then(serde_json::Value::as_object) else {
return Err("its document's `env` is not an object of names to values".to_string());
};
let mut added = Vec::with_capacity(env.len());
for (name, value) in env {
if !oneharness_core::domain::config::valid_env_name(name) {
return Err(format!(
"its document's `env` names `{name}`, which is not a valid environment \
variable name"
));
}
let Some(value) = value.as_str() else {
return Err(format!("its document's `env.{name}` is not a string"));
};
if value.contains('\0') {
return Err(format!(
"its document's `env.{name}` holds a NUL, which no environment can carry"
));
}
added.push((name.clone(), value.to_string()));
}
Ok(added)
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
struct EnvFromSource {
file: String,
variant: String,
source: String,
}
fn validate(
graph: &ConfigRef,
sets: &[String],
refreshed: &BTreeMap<String, String>,
) -> std::result::Result<(), String> {
let missing: Vec<String> = env_from_sources(graph, sets)?
.into_iter()
.filter(|named| !refreshed.contains_key(&named.source))
.map(|named| {
format!(
"{} names harness variant '{}' whose env_from source '{}' is not in the \
environment this dispatch would be launched with",
named.file, named.variant, named.source
)
})
.collect();
if missing.is_empty() {
Ok(())
} else {
Err(missing.join("; "))
}
}
fn env_from_sources(
graph: &ConfigRef,
sets: &[String],
) -> std::result::Result<BTreeSet<EnvFromSource>, String> {
let mut resolver = Resolver::new();
let document = resolver
.resolve(graph, None)
.map_err(|error| format!("the graph '{}' could not be read: {error}", graph.0))?
.clone();
let graph_dir = document.base_dir.clone();
let config = graph_with_overrides(&document.content, &graph.0, sets)?;
let mut refs: Vec<&ConfigRef> = Vec::new();
for member in config.members.values() {
match member {
Member::Oneharness(member) => refs.push(&member.oneharness_config),
Member::Onejudge(member) => {
refs.push(&member.agent.oneharness_config);
for side in &member.judge {
if let JudgeSide::Harness(harness) = side {
refs.push(&harness.oneharness_config);
}
}
}
}
}
let mut named = BTreeSet::new();
for reference in refs {
let file = config_file(reference, graph_dir.as_deref());
let resolved = resolver
.resolve(reference, graph_dir.as_deref())
.map_err(|error| format!("{file} could not be read: {error}"))?;
let config = oneharness_core::domain::config::parse(&resolved.content)
.map_err(|error| format!("{file} is not an oneharness config: {error}"))?;
for (id, harness) in &config.harness {
for (name, variant) in &harness.variant {
for source in variant.env_from.values() {
named.insert(EnvFromSource {
file: file.clone(),
variant: format!("{id}:{name}"),
source: source.clone(),
});
}
}
}
}
Ok(named)
}
fn graph_with_overrides(
content: &str,
origin: &str,
sets: &[String],
) -> std::result::Result<GraphConfig, String> {
let mut parsed: serde_json::Value = serde_norway::from_str(content)
.map_err(|error| format!("the graph '{origin}' could not be read: {error}"))?;
let overrides = sets
.iter()
.map(|value| oneagentgraph::run::parse_set(value).map_err(|error| error.to_string()))
.collect::<std::result::Result<Vec<_>, _>>()?;
oneagentgraph::run::apply_overrides(&mut parsed, &overrides)
.map_err(|error| error.to_string())?;
serde_norway::to_value(&parsed)
.and_then(serde_norway::from_value)
.map_err(|error| format!("the graph '{origin}' could not be read: {error}"))
}
fn config_file(reference: &ConfigRef, graph_dir: Option<&Path>) -> String {
if resolve::is_remote(reference) {
reference.0.clone()
} else {
resolve::local_path(reference, graph_dir)
.display()
.to_string()
}
}
pub(crate) fn refused_zero_timeout(spelling: &str) -> String {
format!(
"{spelling} names a {HOOK} hook timeout of zero seconds, which ends the hook before it \
has begun — give it a positive whole number of seconds, or leave it out to take \
{DEFAULT_DISPATCH_ENV_HOOK_TIMEOUT_SECONDS} seconds"
)
}
#[cfg(test)]
fn graph_at(path: &Path) -> ConfigRef {
ConfigRef(path.display().to_string())
}
#[cfg(test)]
mod tests {
use std::path::PathBuf;
use super::*;
fn scratch(name: &str) -> PathBuf {
let root = std::env::temp_dir().join(format!(
"onepipeline-dispatchenv-{name}-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&root);
std::fs::create_dir_all(&root).expect("a scratch root");
root
}
#[test]
fn a_malformed_document_is_named_by_its_member_and_never_by_a_value() {
let sentinel = "hunter2-sentinel";
for (printed, names) in [
(
r#"{"version": 1, "env": {"A": "b"}, "extra": 1}"#,
"`extra`",
),
(
&format!(r#"{{"version": 1, "env": {{"A": 5, "B": "{sentinel}"}}}}"#),
"`env.A` is not a string",
),
(r#"{"version": 2, "env": {}}"#, "version 2"),
(r#"{"env": {}}"#, "no `version`"),
(r#"{"version": 1}"#, "`env` is not an object"),
(r#"{"version": 1, "env": []}"#, "`env` is not an object"),
(r#"{"version": 1, "env": {"1BAD": "x"}}"#, "`1BAD`"),
(r#"{"version": 1, "env": {"A=B": "x"}}"#, "`A=B`"),
(r#"[1]"#, "not a JSON object"),
(
&format!(
r#"{{"version": 1, "env": {{}}}} {{"version": 1, "env": {{"S": "{sentinel}"}}}}"#
),
"not one JSON document",
),
(&format!("not json {sentinel}"), "not one JSON document"),
("", "not one JSON document"),
] {
let what = Document::read(printed.as_bytes())
.expect_err(&format!("{printed} was read as a document"));
assert!(what.contains(names), "{printed}: {what}");
assert!(
!what.contains(sentinel),
"the refusal of {printed} quoted a value: {what}"
);
}
let nul = format!(r#"{{"version": 1, "env": {{"A": "x\u0000{sentinel}"}}}}"#);
let what = Document::read(nul.as_bytes()).expect_err("a NUL was accepted");
assert!(what.contains("`env.A`") && what.contains("NUL"), "{what}");
assert!(!what.contains(sentinel), "{what}");
assert_eq!(
Document::read(br#"{"version": 1, "env": {"B": "2", "A": "1"}}"#).expect("it reads"),
vec![
("B".to_string(), "2".to_string()),
("A".to_string(), "1".to_string())
]
);
assert_eq!(
Document::read(b"{\"version\": 1, \"env\": {}}\n").expect("it reads"),
Vec::new()
);
}
#[test]
fn every_env_from_source_the_launch_would_read_is_checked_after_the_overrides() {
let root = scratch("sources");
let write = |name: &str, text: &str| {
std::fs::write(root.join(name), text).expect("the fixture is written");
};
write(
"worker.toml",
"run_mode = \"fallback\"\nharnesses = [\"claude-code:work\"]\n\
[harness.claude-code.variant.work.env_from]\nCLAUDE_CONFIG_DIR = \"HOST_CLAUDE_HOME\"\n",
);
write(
"judge.toml",
"harnesses = [\"codex:review\"]\n[harness.codex.variant.review.env_from]\n\
CODEX_HOME = \"HOST_CODEX_HOME\"\nOPENAI_API_KEY = \"HOST_OPENAI_KEY\"\n",
);
write(
"other.toml",
"harnesses = [\"claude-code\"]\n[harness.claude-code.variant.alt.env_from]\n\
CLAUDE_CONFIG_DIR = \"HOST_ALT_HOME\"\n",
);
write("base.yaml", "system_prompt: Do the work.\n");
write(
"graph.yaml",
"version: 1\nname: node-scope\nmembers:\n worker:\n kind: onejudge\n \
base_config: ./base.yaml\n agent:\n oneharness_config: ./worker.toml\n \
judge:\n - oneharness_config: ./judge.toml\n - kind: llmlint\n \
mode: bypass\n reporter:\n kind: oneharness\n oneharness_config: ./worker.toml\n",
);
let graph = graph_at(&root.join("graph.yaml"));
let file = |name: &str| root.join(name).display().to_string();
let named = env_from_sources(&graph, &[]).expect("the sources read");
let seen: Vec<(String, String, String)> = named
.iter()
.map(|named| {
(
named.file.clone(),
named.variant.clone(),
named.source.clone(),
)
})
.collect();
assert_eq!(
seen,
vec![
(
file("judge.toml"),
"codex:review".into(),
"HOST_CODEX_HOME".into()
),
(
file("judge.toml"),
"codex:review".into(),
"HOST_OPENAI_KEY".into()
),
(
file("worker.toml"),
"claude-code:work".into(),
"HOST_CLAUDE_HOME".into()
),
]
);
let overridden = env_from_sources(
&graph,
&["members.worker.agent.oneharness_config=./other.toml".to_string()],
)
.expect("the overridden sources read");
assert!(
overridden
.iter()
.any(|named| named.source == "HOST_ALT_HOME"),
"{overridden:?}"
);
let mut refreshed = BTreeMap::new();
for present in ["HOST_CODEX_HOME", "HOST_OPENAI_KEY", "HOST_CLAUDE_HOME"] {
refreshed.insert(present.to_string(), "set".to_string());
}
validate(&graph, &[], &refreshed).expect("every source is present");
refreshed.remove("HOST_CLAUDE_HOME");
refreshed.remove("HOST_OPENAI_KEY");
let why = validate(&graph, &[], &refreshed).expect_err("two sources are missing");
assert!(
why.contains(&format!(
"{} names harness variant 'claude-code:work' whose env_from source \
'HOST_CLAUDE_HOME' is not in the environment",
file("worker.toml")
)),
"{why}"
);
assert!(
why.contains(&format!(
"{} names harness variant 'codex:review' whose env_from source \
'HOST_OPENAI_KEY'",
file("judge.toml")
)),
"{why}"
);
assert!(!why.contains("HOST_CODEX_HOME"), "{why}");
write("graph-absent.yaml", "version: 1\nname: g\nmembers:\n worker:\n kind: oneharness\n oneharness_config: ./nowhere.toml\n");
let why = validate(&graph_at(&root.join("graph-absent.yaml")), &[], &refreshed)
.expect_err("an absent config is refused");
assert!(
why.contains(&file("nowhere.toml")) && why.contains("could not be read"),
"{why}"
);
write("bad.toml", "harnesses = [\"no-such-harness\"]\n");
write("graph-bad.yaml", "version: 1\nname: g\nmembers:\n worker:\n kind: oneharness\n oneharness_config: ./bad.toml\n");
let why = validate(&graph_at(&root.join("graph-bad.yaml")), &[], &refreshed)
.expect_err("a config oneharness refuses is refused");
assert!(
why.contains(&file("bad.toml")) && why.contains("not an oneharness config"),
"{why}"
);
let why = validate(&graph_at(&root.join("nowhere.yaml")), &[], &refreshed)
.expect_err("an absent graph is refused");
assert!(
why.contains("nowhere.yaml") && why.contains("could not be read"),
"{why}"
);
let why = validate(&graph, &["members.nobody.model=x".to_string()], &refreshed)
.expect_err("an override naming nothing is refused");
assert!(why.contains("nobody"), "{why}");
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_launch_naming_no_hook_adds_nothing_and_reads_no_graph() {
let root = scratch("unhooked");
let paths = RunPaths::under(&root, "demo");
let record: LaunchRecord = serde_json::from_value(serde_json::json!({
"run_id": "demo",
"project": "plans:demo",
"dir": root.display().to_string(),
}))
.expect("a record naming no hook, as an earlier build wrote one");
assert_eq!(record.dispatch_env_hook(), None);
let added = run_and_check(&Launching {
paths: &paths,
record: &record,
node: "build",
graph: &graph_at(&root.join("no-such-graph.yaml")),
sets: &[],
own: &[],
})
.expect("a launch naming no hook is not refused");
assert!(added.is_empty());
assert!(!hooks::hook_log(&paths, HOOK).exists());
let _ = std::fs::remove_dir_all(&root);
}
}