use std::collections::{BTreeMap, BTreeSet};
use std::path::{Path, PathBuf};
use serde_json::{json, Value};
use crate::error::Result;
use crate::graph::{Graph, NodeStatus};
use crate::journal::{self, Journal};
use crate::ledger::{self, RunPaths};
pub const PREFIX: &str = "run:";
pub const SYNTAX: &str = "run:<run_id>#<node_id>";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Reference {
pub run: String,
pub node: String,
}
pub fn parse(dependency: &str) -> Option<Reference> {
let rest = dependency.strip_prefix(PREFIX)?;
let (run, node) = rest.split_once('#')?;
if node.is_empty() || !ledger::is_valid_run_id(run) {
return None;
}
Some(Reference {
run: run.to_string(),
node: node.to_string(),
})
}
pub fn is_reference(dependency: &str) -> bool {
parse(dependency).is_some()
}
pub fn is_malformed(dependency: &str) -> bool {
dependency.starts_with(PREFIX) && parse(dependency).is_none()
}
pub fn extent(root: &Path, run: &str) -> Option<u64> {
let paths = RunPaths::under(root, run);
if !paths.exists() {
return None;
}
Some(ledger::read_lines(&paths.journal()).len() as u64)
}
fn settled_status(root: &Path, reference: &Reference) -> Option<NodeStatus> {
let paths = RunPaths::under(root, &reference.run);
if !paths.exists() {
return None;
}
ledger::read_lines(&paths.journal())
.iter()
.filter_map(|line| serde_json::from_str::<Value>(line).ok())
.filter(|event| {
event.get("kind").and_then(Value::as_str)
== Some(journal::PipelineKind::NodeSettled.as_str())
})
.filter(|event| {
event
.get("labels")
.and_then(|l| l.get("node"))
.and_then(Value::as_str)
== Some(reference.node.as_str())
})
.filter_map(|event| {
event
.get("payload")
.and_then(|p| p.get("status"))
.and_then(Value::as_str)
.and_then(NodeStatus::parse)
})
.next_back()
}
pub fn edges(graph: &Graph) -> BTreeMap<String, Vec<String>> {
let mut edges: BTreeMap<String, Vec<String>> = BTreeMap::new();
for node in graph.iter() {
for dep in &node.deps {
if is_reference(dep) {
edges.entry(dep.clone()).or_default().push(node.id.clone());
}
}
}
edges
}
#[derive(Debug)]
pub struct Observer {
root: PathBuf,
baselines: BTreeMap<String, u64>,
reported: BTreeSet<(String, String)>,
}
impl Observer {
pub fn new(
root: &Path,
baselines: BTreeMap<String, u64>,
reported: BTreeSet<(String, String)>,
) -> Self {
Self {
root: root.to_path_buf(),
baselines,
reported,
}
}
pub fn of_run(paths: &RunPaths, state: &crate::projection::RunState) -> Self {
let root = paths
.dir
.parent()
.map_or_else(ledger::runs_root, Path::to_path_buf);
Self::new(
&root,
state.cross_dag_baselines.clone(),
state.cross_dag_reported.clone(),
)
}
pub fn resolve(
&mut self,
graph: &Graph,
paths: &RunPaths,
round: u64,
journal: &mut Journal,
) -> Result<BTreeMap<String, NodeStatus>> {
let mut resolved = BTreeMap::new();
for (dependency, consumers) in edges(graph) {
let Some(reference) = parse(&dependency) else {
continue;
};
let status = settled_status(&self.root, &reference);
if status != Some(NodeStatus::Done) {
resolved.insert(dependency, NodeStatus::Blocked);
continue;
}
resolved.insert(dependency.clone(), NodeStatus::Done);
let Some(extent) = extent(&self.root, &reference.run) else {
continue;
};
let baseline = match self.baselines.get(&dependency) {
Some(baseline) => *baseline,
None => {
self.baselines.insert(dependency.clone(), extent);
if let Some(first) = consumers.first() {
journal.emit(
journal::PipelineKind::CrossDagSatisfied,
journal::labels(&paths.run, Some(round), Some(first)),
journal::payload(&[
("dependency", json!(dependency)),
("last_seq", json!(extent)),
]),
)?;
}
extent
}
};
if extent <= baseline {
continue;
}
for consumer in consumers {
let pair = (dependency.clone(), consumer.clone());
if !self.reported.insert(pair) {
continue;
}
journal.emit(
journal::PipelineKind::UpstreamModified,
journal::labels(&paths.run, Some(round), Some(&consumer)),
journal::payload(&[
("dependency", json!(dependency)),
("captured_last_seq", json!(baseline)),
("observed_last_seq", json!(extent)),
]),
)?;
}
}
Ok(resolved)
}
}
pub fn resolve_quietly(root: &Path, graph: &Graph) -> BTreeMap<String, NodeStatus> {
edges(graph)
.into_keys()
.filter_map(|dependency| {
let reference = parse(&dependency)?;
let status = match settled_status(root, &reference) {
Some(NodeStatus::Done) => NodeStatus::Done,
_ => NodeStatus::Blocked,
};
Some((dependency, status))
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_reference_needs_both_halves() {
assert_eq!(
parse("run:other#build"),
Some(Reference {
run: "other".into(),
node: "build".into()
})
);
assert_eq!(
parse("run:other#ship/verify").map(|r| r.node),
Some("ship/verify".to_string())
);
for malformed in [
"run:other",
"run:#build",
"run:other#",
"run:",
"run:#",
"run:../elsewhere#build",
"run:../../elsewhere#build",
"run:a/b#build",
"run:/absolute#build",
"run:.#build",
"run:..#build",
] {
assert_eq!(parse(malformed), None, "{malformed} parsed");
assert!(is_malformed(malformed), "{malformed} is not reported wrong");
}
assert_eq!(parse("build"), None);
assert!(!is_malformed("build"));
}
#[test]
fn an_unknown_run_has_no_extent_and_no_status() {
let root =
std::env::temp_dir().join(format!("onepipeline-crossdag-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
std::fs::create_dir_all(&root).expect("a scratch root");
assert_eq!(extent(&root, "nobody"), None);
assert_eq!(
settled_status(
&root,
&Reference {
run: "nobody".into(),
node: "build".into()
}
),
None
);
let _ = std::fs::remove_dir_all(&root);
}
}