use std::io::{self, BufRead, BufReader, Read, Write};
use std::path::Path;
use std::process::{Command, Stdio};
use serde::Deserialize;
use crate::lifecycle::{Plugins, Sealed};
use crate::log::{Level, Log};
use crate::message::{self, PROTOCOL};
use crate::op::Phase;
use crate::registry::PluginRef;
use crate::verb::Verb;
use crate::wire::{OpContext, SealFacts};
pub const DEPTH_CAP: u32 = 8;
const BUSY_RETRIES: u32 = 6;
const BUSY_BACKOFF_MS: u64 = 2;
pub(crate) fn retry_busy<T>(mut exec: impl FnMut() -> io::Result<T>) -> io::Result<T> {
for _ in 0..BUSY_RETRIES {
match exec() {
Err(e) if e.kind() == io::ErrorKind::ExecutableFileBusy => {
std::thread::sleep(std::time::Duration::from_millis(BUSY_BACKOFF_MS));
}
other => return other,
}
}
exec()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Protocol {
pub protocol: Vec<u32>,
pub ops: Vec<String>,
}
#[derive(Deserialize)]
#[serde(untagged)]
enum Versions {
One(u32),
Many(Vec<u32>),
}
#[derive(Deserialize)]
struct RawProtocol {
protocol: Versions,
ops: Vec<String>,
}
impl Protocol {
#[must_use]
pub fn handles(&self, op: Verb) -> bool {
self.ops.iter().any(|o| o == op.token())
}
#[must_use]
pub fn speaks(&self, version: u32) -> bool {
self.protocol.contains(&version)
}
}
pub fn describe(bin: &Path) -> io::Result<Protocol> {
let out = retry_busy(|| Command::new(bin).arg("protocol").output())?;
if !out.status.success() {
return Err(io::Error::other(format!(
"plugin protocol self-describe exited {}",
out.status
)));
}
let raw: RawProtocol = serde_json::from_slice(&out.stdout).map_err(io::Error::other)?;
let protocol = match raw.protocol {
Versions::One(v) => vec![v],
Versions::Many(vs) => vs,
};
Ok(Protocol { protocol, ops: raw.ops })
}
pub struct Subprocess<'a> {
ctx: OpContext,
log: &'a Log,
depth: u32,
}
impl<'a> Subprocess<'a> {
#[must_use]
pub fn new(ctx: OpContext, log: &'a Log, depth: u32) -> Self {
Self { ctx, log, depth }
}
fn at_cap(&self) -> bool {
self.depth >= DEPTH_CAP
}
fn invoke(
&self,
plugin: &PluginRef,
op: Verb,
phase: Phase,
dir: &Path,
sealed: Option<&Sealed>,
rolling_back: Option<&str>,
) -> io::Result<()> {
let bin = plugin.bin.as_ref().ok_or_else(|| {
let n = &plugin.name;
io::Error::other(format!("plugin {n} referenced but bin/{n} missing — run bl install"))
})?;
let metadata = match sealed.and_then(|s| s.message) {
Some(m) => Some(message::parse(m)?),
None => None,
};
let facts = sealed.map(|s| SealFacts {
commit: s.commit,
previous_commit: s.previous_commit,
metadata: metadata.as_ref(),
});
let payload = self.ctx.wire(&plugin.name, op.token(), phase.token(), facts, rolling_back);
let json = serde_json::to_string(&payload).map_err(io::Error::other)?;
let status = self.spawn(bin, &plugin.name, op, phase, dir, &json)?;
if status.success() {
return Ok(());
}
let (name, o, p) = (&plugin.name, op.token(), phase.token());
let msg = match rolling_back {
Some(_) => format!("plugin {name} rollback failed ({status}) — its {o}.{p} side effects may not be unwound"),
None => format!("plugin {name} aborted the op ({status})"),
};
self.log.record(Level::Error, "core", Some(phase), &msg);
Err(io::Error::other(msg))
}
fn spawn(
&self,
bin: &Path,
name: &str,
op: Verb,
phase: Phase,
dir: &Path,
payload: &str,
) -> io::Result<std::process::ExitStatus> {
self.log.record(Level::Info, "core", Some(phase), &format!("invoke {name}"));
let depth = (self.depth + 1).to_string();
let mut child = retry_busy(|| {
Command::new(bin)
.arg(op.token())
.arg(phase.token())
.current_dir(dir)
.env("BALLS_PROTOCOL", PROTOCOL.to_string())
.env("BALLS_PLUGIN_NAME", name)
.env("BALLS_PLUGIN_DEPTH", &depth)
.stdin(Stdio::piped())
.stdout(Stdio::inherit())
.stderr(Stdio::piped())
.spawn()
})?;
child.stdin.take().expect("stdin was configured as a pipe").write_all(payload.as_bytes())?;
self.relay(name, phase, child.stderr.take().expect("stderr was configured as a pipe"));
child.wait()
}
fn relay(&self, name: &str, phase: Phase, stderr: std::process::ChildStderr) {
capped_lines(BufReader::new(stderr), RELAY_LINE_MAX, |line| {
self.log.record(Level::Info, name, Some(phase), line);
});
}
}
const RELAY_LINE_MAX: u64 = 1 << 20;
fn capped_lines(mut reader: impl BufRead, cap: u64, mut sink: impl FnMut(&str)) {
let mut buf = Vec::new();
while reader.by_ref().take(cap).read_until(b'\n', &mut buf).unwrap_or(0) != 0 {
if buf.last() == Some(&b'\n') {
buf.pop();
}
sink(&String::from_utf8_lossy(&buf));
buf.clear();
}
}
impl Plugins for Subprocess<'_> {
fn run(&self, plugin: &PluginRef, op: Verb, phase: Phase, dir: &Path, sealed: Option<&Sealed>) -> io::Result<()> {
if self.at_cap() {
let msg = format!(
"invocation-tree depth cap ({DEPTH_CAP}) reached at {}.{} — aborting before plugin {} (§6)",
op.token(),
phase.token(),
plugin.name
);
self.log.record(Level::Error, "core", Some(phase), &msg);
return Err(io::Error::other(msg));
}
self.invoke(plugin, op, phase, dir, sealed, None)
}
fn rollback(&self, plugin: &PluginRef, op: Verb, phase: Phase, dir: &Path, sealed: Option<&Sealed>) {
if self.at_cap() {
return;
}
let _ = self.invoke(plugin, op, phase, dir, sealed, Some(phase.token()));
}
}
#[cfg(test)]
#[path = "plugin_tests.rs"]
mod tests;