mod execution;
mod handlers;
mod policy;
pub use policy::{BuildPolicy, IncrementalPolicy, ProductAction};
use anyhow::Result;
use indicatif::ProgressBar;
use parking_lot::Mutex;
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use crate::color;
use crate::display::DisplayOptions;
use crate::errors;
use crate::graph::BuildGraph;
use crate::object_store::{ExplainAction, ObjectStore};
use crate::processors::ProcessorMap;
use crate::stats::ProcessStats;
enum PreCheckResult {
Handled,
NeedsExecution,
}
enum RestoreOutcome {
Restored,
Failed,
NotRestorable,
}
struct WorkItem {
product_id: usize,
input_checksum: String,
needs_rebuild: bool,
}
struct HandlerContext<'b> {
product: &'b crate::graph::Product,
id: usize,
input_checksum: &'b str,
proc_name: &'b str,
keep_going: bool,
shared: &'b SharedState,
pb: &'b ProgressBar,
}
struct LevelWork {
batch_groups: HashMap<String, Vec<WorkItem>>,
non_batch_items: Vec<WorkItem>,
}
#[derive(Debug)]
pub struct ExecutorOptions {
pub parallel: usize,
pub verbose: bool,
pub display_opts: DisplayOptions,
pub batch_size: Option<usize>,
pub explain: bool,
pub retry: usize,
}
#[derive(Debug)]
struct SharedState {
stats: Arc<Mutex<HashMap<String, ProcessStats>>>,
errors: Arc<Mutex<Vec<anyhow::Error>>>,
failed_products: Arc<Mutex<HashSet<usize>>>,
failed_messages: Arc<Mutex<Vec<String>>>,
failed_processors: Arc<Mutex<HashSet<String>>>,
global_current: Arc<AtomicUsize>,
global_total: usize,
}
pub struct ClassifiedProduct {
pub id: usize,
pub action: ProductAction,
pub input_checksum: String,
}
pub struct Classification {
pub skip_count: usize,
pub restore_count: usize,
pub build_count: usize,
pub products: Vec<ClassifiedProduct>,
}
pub fn classify_products(
ctx: &crate::build_context::BuildContext,
policy: &dyn BuildPolicy,
graph: &BuildGraph,
order: &[usize],
object_store: &ObjectStore,
force: bool,
) -> Classification {
let mut skip_count = 0;
let mut restore_count = 0;
let mut build_count = 0;
let mut will_change: HashSet<usize> = HashSet::new();
let mut products: Vec<ClassifiedProduct> = Vec::with_capacity(order.len());
for &id in order {
let product = graph.get_product(id).expect(errors::INVALID_PRODUCT_ID);
let dep_changed = graph
.get_dependencies(id)
.iter()
.any(|d| will_change.contains(d));
let Ok(input_checksum) = crate::checksum::combined_input_checksum(ctx, &product.inputs)
else {
build_count += 1;
will_change.insert(id);
products.push(ClassifiedProduct {
id,
action: ProductAction::Build,
input_checksum: String::new(),
});
continue;
};
let action = policy.classify(
ctx,
product,
object_store,
&input_checksum,
dep_changed,
force,
);
match action {
ProductAction::Skip => {
skip_count += 1;
}
ProductAction::Restore => {
restore_count += 1;
will_change.insert(id);
}
ProductAction::Build => {
build_count += 1;
will_change.insert(id);
}
}
products.push(ClassifiedProduct {
id,
action,
input_checksum,
});
}
Classification {
skip_count,
restore_count,
build_count,
products,
}
}
pub fn unlink_pending_outputs(
graph: &BuildGraph,
object_store: &ObjectStore,
classification: &Classification,
) -> Result<()> {
for c in &classification.products {
if matches!(c.action, ProductAction::Skip) {
continue;
}
let product = graph.get_product(c.id).expect(errors::INVALID_PRODUCT_ID);
execution::remove_stale_outputs(product, object_store, &c.input_checksum)?;
}
Ok(())
}
pub struct Executor<'a> {
processors: &'a ProcessorMap,
build_ctx: &'a crate::build_context::BuildContext,
policy: &'a dyn BuildPolicy,
parallel: usize,
verbose: bool,
display_opts: DisplayOptions,
batch_size: Option<usize>,
explain: bool,
retry: usize,
}
impl<'a> Executor<'a> {
pub fn new(
processors: &'a ProcessorMap,
build_ctx: &'a crate::build_context::BuildContext,
policy: &'a dyn BuildPolicy,
opts: ExecutorOptions,
) -> Self {
Self {
processors,
build_ctx,
policy,
parallel: opts.parallel,
verbose: opts.verbose && crate::json_output::human_output_enabled(),
display_opts: opts.display_opts,
batch_size: opts.batch_size,
explain: opts.explain,
retry: opts.retry,
}
}
fn is_interrupted(&self) -> bool {
self.build_ctx.is_interrupted()
}
fn product_display(&self, product: &crate::graph::Product) -> String {
product.display(self.display_opts)
}
fn inc_global(shared: &SharedState) {
shared.global_current.fetch_add(1, Ordering::SeqCst);
}
fn inc_progress(pb: &ProgressBar, shared: &SharedState) {
Self::inc_global(shared);
pb.inc(1);
}
fn print_explain(&self, product: &crate::graph::Product, action: &ExplainAction) {
let styled = match action {
ExplainAction::Skip => color::dim("SKIP"),
ExplainAction::Restore(_) => color::cyan("RESTORE"),
ExplainAction::Rebuild(_) => color::yellow("BUILD"),
};
crate::output::info(&format!(
"[{}] {} {} ({})",
product.processor,
styled,
self.product_display(product),
action
));
}
pub fn clean(&self, graph: &BuildGraph, verbose: bool) -> Result<HashMap<String, usize>> {
let mut stats: HashMap<String, usize> = HashMap::new();
for product in graph.products() {
if self.is_interrupted() {
return Err(crate::exit_code::interrupted());
}
if let Some(processor) = self.processors.get(&product.processor) {
let count = processor.clean(product, verbose)?;
if count > 0 {
*stats.entry(product.processor.clone()).or_default() += count;
}
}
}
Ok(stats)
}
}
pub fn has_failed_dependency(graph: &BuildGraph, id: usize, failed: &HashSet<usize>) -> bool {
for &dep_id in graph.get_dependencies(id) {
if failed.contains(&dep_id) {
return true;
}
}
false
}
pub fn compute_parallel_levels(graph: &BuildGraph, order: &[usize]) -> Vec<Vec<usize>> {
let mut levels: Vec<Vec<usize>> = Vec::new();
let mut product_level: HashMap<usize, usize> = HashMap::new();
for &id in order {
let max_dep_level = graph
.get_dependencies(id)
.iter()
.filter_map(|&dep_id| product_level.get(&dep_id))
.max()
.copied()
.unwrap_or(0);
let my_level = if graph.get_dependencies(id).is_empty() {
0
} else {
max_dep_level + 1
};
product_level.insert(id, my_level);
while levels.len() <= my_level {
levels.push(Vec::new());
}
levels[my_level].push(id);
}
levels
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parallel_levels_diamond() {
let mut graph = BuildGraph::new();
let top = graph
.add_product(vec!["a.src".into()], vec!["a.o".into()], "cc", None)
.unwrap();
let left = graph
.add_product(vec!["a.o".into()], vec!["b.o".into()], "cc", None)
.unwrap();
let right = graph
.add_product(vec!["a.o".into()], vec!["c.o".into()], "cc", None)
.unwrap();
let bottom = graph
.add_product(
vec!["b.o".into(), "c.o".into()],
vec!["d.o".into()],
"cc",
None,
)
.unwrap();
let lone = graph
.add_product(vec!["x.src".into()], vec!["x.o".into()], "cc", None)
.unwrap();
graph.resolve_dependencies();
let order = graph.topological_sort().unwrap();
let levels = compute_parallel_levels(&graph, &order);
let sorted = |mut ids: Vec<usize>| {
ids.sort_unstable();
ids
};
assert_eq!(
levels.len(),
3,
"diamond plus a free node is three levels: {levels:?}"
);
assert_eq!(sorted(levels[0].clone()), sorted(vec![top, lone]));
assert_eq!(sorted(levels[1].clone()), sorted(vec![left, right]));
assert_eq!(levels[2], vec![bottom]);
}
#[test]
fn parallel_levels_cover_all_products() {
let mut g = BuildGraph::new();
g.add_product(vec!["a.src".into()], vec!["a.o".into()], "cc", None)
.unwrap();
g.add_product(vec!["a.o".into()], vec!["b.o".into()], "cc", None)
.unwrap();
g.add_product(vec!["free.src".into()], vec!["free.o".into()], "cc", None)
.unwrap();
g.resolve_dependencies();
let order = g.topological_sort().unwrap();
let levels = compute_parallel_levels(&g, &order);
let mut all: Vec<usize> = levels.into_iter().flatten().collect();
all.sort_unstable();
assert_eq!(all, vec![0, 1, 2]);
}
#[test]
fn failed_dependency_is_direct_only() {
let mut g = BuildGraph::new();
let a = g
.add_product(vec!["a.src".into()], vec!["a.o".into()], "cc", None)
.unwrap();
let b = g
.add_product(vec!["a.o".into()], vec!["b.o".into()], "cc", None)
.unwrap();
let c = g
.add_product(vec!["b.o".into()], vec!["c.o".into()], "cc", None)
.unwrap();
g.resolve_dependencies();
let failed: HashSet<usize> = [a].into();
assert!(
has_failed_dependency(&g, b, &failed),
"b directly depends on failed a"
);
assert!(
!has_failed_dependency(&g, c, &failed),
"c depends on a only through b; direct check must not see it"
);
assert!(
!has_failed_dependency(&g, a, &failed),
"a has no dependencies"
);
}
}