mod liveness;
mod route;
mod sweeps;
pub mod worker;
use super::Roots;
use super::drift::{self, Drift};
use super::snapshot::{Growth, Snapshot};
use crate::binding::Workspace;
use crate::budgets::StepBill;
use crate::git_tree::{GitTree, ProbeStack};
use crate::opslog::OpRow;
use crate::projects::balls::Ball;
use crate::projects::join::JoinRow;
use crate::projects::runner::BlRunner;
use crate::state::{
DirtySet, SnapshotCell, WatchSetHandle, lock_watchset, new_watchset, publish_snapshot,
};
use crate::ui_state::Clock;
use crate::watch::Mark;
use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
pub struct Deriver {
roots: Roots,
clock: Arc<dyn Clock>,
probes: ProbeStack,
pub(super) watches: WatchSetHandle,
dirty: DirtySet,
schedule: super::dirty::Schedule,
pub(super) cadence: super::Cadence,
pub(super) fleet: std::collections::BTreeMap<String, crate::fleet::Policy>,
cell: SnapshotCell,
balls: Box<dyn BlRunner>,
pub(super) projects: Vec<PathBuf>,
pub(super) workspaces: Vec<Workspace>,
pub(super) trees: HashMap<PathBuf, GitTree>,
pub(super) bills: HashMap<PathBuf, Vec<StepBill>>,
pub(super) windows: std::collections::BTreeMap<String, u64>,
pub(super) balls_by_project: HashMap<PathBuf, Vec<Ball>>,
pub(super) closed_by_project: HashMap<PathBuf, Vec<Ball>>,
pub(super) join_rows: Vec<JoinRow>,
pub(super) ops: Vec<OpRow>,
pub(super) changed: bool,
late: bool,
growth: Vec<Growth>,
ui_bytes: Option<Vec<u8>>,
}
impl Deriver {
pub(super) fn new(
roots: Roots,
clock: Arc<dyn Clock>,
balls: Box<dyn BlRunner>,
dirty: DirtySet,
cell: SnapshotCell,
) -> Self {
let cadence = super::Cadence::default();
let schedule = super::dirty::Schedule::new(Arc::clone(&clock), cadence);
Self {
roots,
clock,
probes: ProbeStack::platform(),
watches: new_watchset(),
dirty,
schedule,
cadence,
fleet: std::collections::BTreeMap::new(),
cell,
balls,
projects: Vec::new(),
workspaces: Vec::new(),
trees: HashMap::new(),
bills: HashMap::new(),
windows: std::collections::BTreeMap::new(),
balls_by_project: HashMap::new(),
closed_by_project: HashMap::new(),
join_rows: Vec::new(),
ops: Vec::new(),
changed: false,
late: false,
growth: Vec::new(),
ui_bytes: None,
}
}
pub fn watchset_handle(&self) -> WatchSetHandle {
Arc::clone(&self.watches)
}
pub fn dirty_handle(&self) -> DirtySet {
self.dirty.clone()
}
pub(super) fn boot(&mut self) {
self.adopt_cadence();
self.adopt_windows();
self.workspaces = crate::binding::workspaces(&self.roots.yog_data, &self.roots.lernie_data);
lock_watchset(&self.watches)
.reconcile(&super::desired_watches(&self.roots, &self.workspaces));
let paths: Vec<PathBuf> = self.workspaces.iter().map(|w| w.path.clone()).collect();
for path in paths {
self.rederive(&path);
}
self.refresh_balls();
self.refresh_ops();
self.publish();
}
pub fn step(&mut self) -> bool {
let started = self.clock.now();
self.changed = false;
self.growth.clear();
self.ui_bytes = None;
let delivered = self.dirty.drain();
let mut found: Vec<Drift> = delivered
.iter()
.filter(|(_, mark)| **mark == Mark::Desync)
.map(|(root, _)| Drift::Desync(root.clone()))
.collect();
self.dispatch_dirty(delivered);
let sweep = self.schedule.sweep();
match sweep {
super::dirty::Sweep::Full => found.extend(self.full_sweep()),
super::dirty::Sweep::Cheap => found.extend(self.cheap_sweep()),
super::dirty::Sweep::None => {}
}
for (root, mark) in self.schedule.due() {
let baseline = self.trees.contains_key(&root);
if mark == Mark::Watch || mark == Mark::Desync {
self.refresh_liveness(&root);
}
if self.rederive(&root) && mark == Mark::Sweep && baseline {
found.push(Drift::Unannounced(root));
}
}
let late = drift::lateness(started, self.clock.now(), self.cadence.late_pass(sweep));
if let Some(secs) = drift::late_edge(late, self.late) {
found.push(Drift::Late(self.roots.yog_state.clone(), secs));
}
self.late = late.is_some();
self.report_drift(&found);
let publish = self.changed || sweep == super::dirty::Sweep::Full;
if publish {
self.publish();
}
publish
}
pub(super) fn publish(&mut self) {
publish_snapshot(
&self.cell,
Arc::new(Snapshot {
workspaces: self.workspaces.clone(),
projects: self.projects.clone(),
trees: self.trees.clone(),
bills: self.bills.clone(),
windows: self.windows.clone(),
balls_by_project: self.balls_by_project.clone(),
closed_by_project: self.closed_by_project.clone(),
join_rows: self.join_rows.clone(),
ops: self.ops.clone(),
growth: std::mem::take(&mut self.growth),
ui_bytes: self.ui_bytes.take(),
derived_at_unix: self.clock.unix(),
cadence: self.cadence,
fleet: self.fleet.clone(),
}),
);
}
}