use std::io::{self, Write};
use std::path::Path;
use std::process::{Command, Stdio};
use crate::checkout;
use crate::config::{self, EffectiveConfig};
use crate::edge::Edge;
use crate::hooks::Hooks;
use crate::log::{Level, Log};
use crate::message::{Metadata, PROTOCOL};
use crate::plugin::{retry_busy, DEPTH_CAP};
use crate::registry::Registry;
use crate::verb::Verb;
use crate::wire::OpContext;
pub(crate) fn fold(edge: &Edge, store: &Path, verb: Verb, id: Option<&str>, cfg: &EffectiveConfig, log: &Log) -> String {
if edge.depth >= DEPTH_CAP {
return String::new();
}
let landing = edge.xdg.clone_dir(&edge.invocation_path).landing();
let Ok(hooks) = Hooks::effective(&landing, &edge.xdg.user_config()) else {
return String::new();
};
let refs = hooks.resolve_read(&Registry::at(&landing), verb.token());
if refs.is_empty() {
return String::new();
}
let binding_path = edge.xdg.clone_dir(&edge.invocation_path).binding();
let (remote, stealth) =
config::remote_ladder(None, &landing, &binding_path, &edge.xdg.user_config()).unwrap_or((None, false));
let binding = checkout::binding(&landing, store, &edge.invocation_path, remote, stealth, cfg.tasks_branch.clone());
let ctx = OpContext { actor: edge.default_actor.clone(), binding, command: None, before: None };
let metadata = id.map_or_else(Metadata::new, |id| Metadata::from([("bl-id".to_string(), vec![id.to_string()])]));
let mut out = String::new();
for plugin in refs {
let Some(bin) = plugin.bin else { continue };
let payload = ctx.read_wire(&plugin.name, verb.token(), &metadata);
let line = serde_json::to_string(&payload)
.map_err(io::Error::other)
.and_then(|json| capture(&bin, &plugin.name, verb.token(), edge.depth, store, &json, log));
match line {
Ok(line) => out.push_str(&line),
Err(e) => log.record(Level::Error, "core", None, &e.to_string()),
}
}
out
}
fn capture(bin: &Path, name: &str, op: &str, depth: u32, store: &Path, payload: &str, log: &Log) -> io::Result<String> {
log.record(Level::Debug, "core", None, &format!("invoke {name}"));
let mut child = retry_busy(|| {
Command::new(bin)
.args([op, "read"])
.current_dir(store)
.env("BALLS_PROTOCOL", PROTOCOL.to_string())
.env("BALLS_PLUGIN_NAME", name)
.env("BALLS_PLUGIN_DEPTH", (depth + 1).to_string())
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
})?;
child.stdin.take().expect("stdin was configured as a pipe").write_all(payload.as_bytes())?;
let out = child.wait_with_output()?;
for line in String::from_utf8_lossy(&out.stderr).lines() {
log.record(Level::Info, name, None, line);
}
if !out.status.success() {
return Err(io::Error::other(format!("plugin {name} failed the {op} read dispatch")));
}
Ok(String::from_utf8_lossy(&out.stdout).into_owned())
}
#[cfg(test)]
#[path = "readop_tests.rs"]
mod tests;