use std::io::{BufRead, BufReader};
use std::path::PathBuf;
use std::process::{Child, Command, Stdio};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
use onevcs::{MergePolicy, SessionRequest};
use serde::Deserialize;
use crate::error::{Error, Result};
use crate::event::Envelope;
pub const BINARY_ENV: &str = "ONEPIPELINE_ONEVCS_BIN";
pub const DEFAULT_BINARY: &str = "onevcs";
pub fn binary() -> String {
std::env::var(BINARY_ENV)
.ok()
.filter(|value| !value.is_empty())
.unwrap_or_else(|| DEFAULT_BINARY.to_string())
}
fn sibling(message: impl Into<String>) -> Error {
Error::Sibling {
tool: "onevcs",
message: message.into(),
}
}
#[derive(Debug, Clone, PartialEq, Eq, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct OpenSession {
pub token: String,
pub worktree: PathBuf,
pub branch: String,
pub base: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Published {
#[serde(default)]
pub url: Option<String>,
#[serde(default)]
pub id: Option<String>,
#[serde(default)]
pub outcome: Option<String>,
}
fn run_json<T: serde::de::DeserializeOwned>(command: &mut Command, what: &str) -> Result<T> {
let output = command
.stdin(Stdio::null())
.output()
.map_err(|e| sibling(format!("cannot start `{} {what}`: {e}", binary())))?;
if !output.status.success() {
return Err(sibling(format!(
"{what} exited {}: {}",
output.status.code().unwrap_or(-1),
String::from_utf8_lossy(&output.stderr).trim()
)));
}
let stdout = String::from_utf8_lossy(&output.stdout);
serde_json::from_str(stdout.trim())
.map_err(|e| sibling(format!("{what} printed something unreadable: {e}")))
}
pub fn session_open(request: &SessionRequest) -> Result<OpenSession> {
let mut command = Command::new(binary());
command.arg("session").arg("open").arg(&request.repo);
if let Some(branch) = &request.branch {
command.arg("--branch").arg(branch);
}
if let Some(base) = &request.base {
command.arg("--base").arg(base);
}
if let Some(checkout) = &request.execution_checkout {
command.arg("--execution-checkout").arg(checkout);
}
run_json(&mut command, "session open")
}
pub fn publish(token: &str, policy: Option<MergePolicy>, title: Option<&str>) -> Result<Published> {
let mut command = Command::new(binary());
command.arg("publish").arg(token);
if let Some(policy) = policy {
command.arg("--policy").arg(policy_arg(policy));
}
if let Some(title) = title {
command.arg("--title").arg(title);
}
run_json(&mut command, "publish")
}
pub fn policy_arg(policy: MergePolicy) -> &'static str {
match policy {
MergePolicy::LocalDirect => "local-direct",
MergePolicy::ChangeOpen => "change-open",
MergePolicy::ChangeAuto => "change-auto",
MergePolicy::ChangeDirect => "change-direct",
}
}
pub fn session_close(token: &str) -> Result<()> {
let output = Command::new(binary())
.arg("session")
.arg("close")
.arg(token)
.stdin(Stdio::null())
.output()
.map_err(|e| sibling(format!("cannot start `{} session close`: {e}", binary())))?;
if output.status.success() {
return Ok(());
}
Err(sibling(format!(
"session close {token} exited {}: {}",
output.status.code().unwrap_or(-1),
String::from_utf8_lossy(&output.stderr).trim()
)))
}
pub fn events(token: &str) -> Vec<Envelope> {
let read = Command::new(binary())
.arg("events")
.arg(token)
.stdin(Stdio::null())
.output();
let output = match read {
Ok(output) if output.status.success() => output,
Ok(output) => {
eprintln!(
"onepipeline: cannot read session {token}'s events: {}",
String::from_utf8_lossy(&output.stderr).trim()
);
return Vec::new();
}
Err(error) => {
eprintln!("onepipeline: cannot read session {token}'s events: {error}");
return Vec::new();
}
};
let text = String::from_utf8_lossy(&output.stdout);
let lines: Vec<&str> = text
.lines()
.filter(|line| !line.trim().is_empty())
.collect();
let envelopes: Vec<Envelope> = lines
.iter()
.filter_map(|line| serde_json::from_str::<Envelope>(line).ok())
.collect();
report_skipped("onevcs", lines.len() - envelopes.len());
envelopes
}
const FOLLOW_GRACE: Duration = Duration::from_secs(5);
const FOLLOW_POLL: Duration = Duration::from_millis(20);
pub fn follow(token: &str, sink: Box<dyn Fn(Envelope) + Send>) -> Option<Follower> {
let started = Command::new(binary())
.arg("events")
.arg(token)
.arg("--follow")
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::inherit())
.spawn();
let mut child = match started {
Ok(child) => child,
Err(error) => {
eprintln!("onepipeline: cannot follow session {token}'s events: {error}");
return None;
}
};
let Some(stdout) = child.stdout.take() else {
stop(&mut child);
eprintln!("onepipeline: cannot read session {token}'s events as they are written");
return None;
};
let relayed = Arc::new(AtomicU64::new(0));
let counted = Arc::clone(&relayed);
let reader = std::thread::Builder::new()
.name(format!("{}-events", binary()))
.spawn(move || {
let mut skipped = 0usize;
for line in BufReader::new(stdout)
.lines()
.map_while(std::io::Result::ok)
{
if line.trim().is_empty() {
continue;
}
match serde_json::from_str::<Envelope>(&line) {
Ok(envelope) => {
counted.fetch_add(1, Ordering::SeqCst);
sink(envelope);
}
Err(_) => skipped += 1,
}
}
report_skipped("onevcs", skipped);
});
match reader {
Ok(reader) => Some(Follower {
child,
reader: Some(reader),
relayed,
}),
Err(error) => {
stop(&mut child);
eprintln!("onepipeline: cannot follow session {token}'s events: {error}");
None
}
}
}
#[derive(Debug)]
pub struct Follower {
child: Child,
reader: Option<std::thread::JoinHandle<()>>,
relayed: Arc<AtomicU64>,
}
impl Drop for Follower {
fn drop(&mut self) {
stop(&mut self.child);
if let Some(reader) = self.reader.take() {
let _ = reader.join();
}
}
}
impl Follower {
pub fn finish(mut self) -> bool {
let deadline = Instant::now() + FOLLOW_GRACE;
let ended = loop {
match self.child.try_wait() {
Ok(Some(status)) => break status.success(),
Err(_) => break false,
Ok(None) if Instant::now() >= deadline => {
stop(&mut self.child);
break false;
}
Ok(None) => std::thread::sleep(FOLLOW_POLL),
}
};
if let Some(reader) = self.reader.take() {
let _ = reader.join();
}
ended || self.relayed.load(Ordering::SeqCst) > 0
}
}
fn stop(child: &mut Child) {
let _ = child.kill();
let _ = child.wait();
}
pub fn report_skipped(tool: &str, skipped: usize) {
if skipped > 0 {
eprintln!("onepipeline: skipped {skipped} {tool} line(s) this build cannot read");
}
}
pub fn session_opened_event(session: &OpenSession, labels: &crate::event::Labels) -> Envelope {
Envelope {
v: crate::event::ENVELOPE_VERSION,
ts: crate::sys::now_rfc3339(),
stream: format!("onevcs-{}", session.token),
seq: 0,
source: crate::event::Source::Vcs,
kind: crate::event::EventKind("session-opened".into()),
labels: labels.clone(),
payload: crate::journal::payload(&[
("token", serde_json::json!(session.token)),
("branch", serde_json::json!(session.branch)),
("base", serde_json::json!(session.base)),
("worktree", serde_json::json!(session.worktree)),
]),
artifacts: Vec::new(),
}
}
pub fn published_event(
published: &Published,
branch: &str,
labels: &crate::event::Labels,
) -> Envelope {
Envelope {
v: crate::event::ENVELOPE_VERSION,
ts: crate::sys::now_rfc3339(),
stream: format!("onevcs-{branch}"),
seq: 1,
source: crate::event::Source::Vcs,
kind: crate::event::EventKind("published".into()),
labels: labels.clone(),
payload: crate::journal::payload(&[
("branch", serde_json::json!(branch)),
("url", serde_json::json!(published.url)),
("id", serde_json::json!(published.id)),
("outcome", serde_json::json!(published.outcome)),
]),
artifacts: Vec::new(),
}
}
pub fn request_for(node: &crate::plan::Node) -> Option<SessionRequest> {
Some(SessionRequest {
repo: node.repo.clone()?,
branch: node.branch.clone(),
base: node.base_branch.clone(),
execution_checkout: node.execution_checkout.clone(),
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::plan::Node;
#[test]
fn every_merge_policy_has_one_spelling_on_the_command_line() {
assert_eq!(policy_arg(MergePolicy::LocalDirect), "local-direct");
assert_eq!(policy_arg(MergePolicy::ChangeOpen), "change-open");
assert_eq!(policy_arg(MergePolicy::ChangeAuto), "change-auto");
assert_eq!(policy_arg(MergePolicy::ChangeDirect), "change-direct");
}
#[test]
fn a_lifecycle_node_asks_for_the_session_its_fields_describe() {
let node = Node {
id: "service".into(),
repo: Some("owner/repo".into()),
branch: Some("feature".into()),
base_branch: Some("main".into()),
execution_checkout: Some("primary".into()),
persona: Some("engineer".into()),
task: Some("## What\nship".into()),
..Node::default()
};
let request = request_for(&node).expect("a lifecycle node asks for a session");
assert_eq!(request.repo, "owner/repo");
assert_eq!(request.branch.as_deref(), Some("feature"));
assert_eq!(request.base.as_deref(), Some("main"));
assert_eq!(request.execution_checkout.as_deref(), Some("primary"));
}
#[test]
fn a_direct_agent_node_asks_for_no_session() {
let node = Node {
id: "build".into(),
persona: Some("engineer".into()),
task: Some("## What\ndo it".into()),
..Node::default()
};
assert!(request_for(&node).is_none());
}
#[test]
fn the_binary_comes_from_the_environment_or_falls_back() {
assert_eq!(
std::env::var(BINARY_ENV)
.ok()
.filter(|v| !v.is_empty())
.unwrap_or_else(|| DEFAULT_BINARY.to_string()),
binary()
);
}
}