use std::collections::{BTreeMap, BTreeSet};
use chrono::{DateTime, TimeDelta, Utc};
use crate::collect::agents::Agents;
use crate::collect::run::FailureKind;
use crate::collect::tracker::{OpenFailure, Trackers};
use crate::collect::worktree;
use crate::config::{Config, Project};
use crate::model::join::{self, Listed, ProjectRows};
use crate::model::snapshot::{
self, AgentProvider, Collected, FailedProject, Filter, ProviderState, Session, SessionState,
Snapshot, TrackerFailure, TrackerState, Tree,
};
use crate::model::types::Pane;
use super::tracker::{open_failure, refresh_project, ProjectWork, ReadAt, Refresh};
struct Read {
at: DateTime<Utc>,
work: Result<ProjectWork, TrackerFailure>,
taken_at: Option<ReadAt>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Wanted {
Everything,
Project(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Asked {
Read(Wanted),
Reloaded(Box<Config>),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Awaited {
pub wanted: Wanted,
pub asked_at: DateTime<Utc>,
pub patience: TimeDelta,
}
impl Awaited {
pub fn unanswered_at(&self, now: DateTime<Utc>) -> bool {
now - self.asked_at >= self.patience
}
}
impl Wanted {
pub fn names(&self, project: &str) -> bool {
match self {
Wanted::Everything => true,
Wanted::Project(named) => named == project,
}
}
}
fn placed(pane: Pane) -> Pane {
let main_tree = worktree::in_the_main_working_tree(&pane.cwd);
pane.with_cwd_in_the_main_working_tree(main_tree)
}
#[derive(Default)]
pub struct Collection {
read: BTreeMap<String, Read>,
panes_last_answered: BTreeMap<String, BTreeSet<String>>,
}
impl Collection {
pub fn collect(
&mut self,
cfg: &Config,
agents: &dyn Agents,
trackers: &dyn Trackers,
wanted: &Wanted,
filter: Filter,
now: DateTime<Utc>,
) -> Snapshot {
let (panes, provider, out_of_reach) = self.every_pane(agents);
for (project, answer) in self.refresh_together(cfg, trackers, wanted, &panes, now) {
match answer {
Ok(Refresh::Unchanged) => {
if let Some(standing) = self.read.get_mut(&project.name) {
standing.at = now;
}
}
Ok(Refresh::Read { at, work }) => {
self.read.insert(
project.name.clone(),
Read {
at: now,
work: Ok(*work),
taken_at: at.map(|at| *at),
},
);
}
Err(failure) => {
self.read.insert(
project.name.clone(),
Read {
at: now,
work: Err(open_failure(&failure)),
taken_at: None,
},
);
}
}
}
self.draw(cfg, &panes, &out_of_reach, provider, filter, now)
}
fn refresh_together<'a>(
&self,
cfg: &'a Config,
trackers: &dyn Trackers,
wanted: &Wanted,
panes: &[Pane],
now: DateTime<Utc>,
) -> Vec<(&'a Project, Result<Refresh, OpenFailure>)> {
std::thread::scope(|reads| {
let reading: Vec<_> = cfg
.read()
.filter(|p| wanted.names(&p.name))
.map(|project| {
let standing = self
.read
.get(&project.name)
.and_then(|read| read.taken_at.clone());
reads.spawn(move || {
let answer =
refresh_project(trackers, project, cfg, panes, standing.as_ref(), now);
(project, answer)
})
})
.collect();
reading
.into_iter()
.map(|read| {
read.join()
.unwrap_or_else(|panicked| std::panic::resume_unwind(panicked))
})
.collect()
})
}
fn draw(
&self,
cfg: &Config,
panes: &[Pane],
out_of_reach: &BTreeSet<String>,
agents: AgentProvider,
filter: Filter,
now: DateTime<Utc>,
) -> Snapshot {
let rows: Vec<ProjectRows<'_>> = self
.that_answered(cfg)
.flat_map(|(project, work)| {
work.roots.iter().filter_map(move |(_, read)| {
read.as_ref().ok().map(|assembled| ProjectRows {
project,
rows: &assembled.beads,
})
})
})
.collect();
let joined = &join::resolve(
&rows,
Listed {
panes,
out_of_reach,
},
cfg,
);
let trees = self
.that_answered(cfg)
.flat_map(|(project, work)| {
work.roots.iter().map(move |(root, read)| match read {
Ok(assembled) => snapshot::build_tree(
project,
assembled,
joined,
&work.readiness,
&work.relations,
agents.state,
cfg,
now,
),
Err(why) => Tree::unread(project, root, TrackerState::from(*why)),
})
})
.collect();
let failed_projects = self
.standing(cfg)
.filter_map(|(project, read)| {
read.work.as_ref().err().map(|failure| FailedProject {
project: project.to_string(),
tracker: *failure,
})
})
.collect();
let read_at = self
.standing(cfg)
.map(|(project, read)| (project.to_string(), read.at))
.collect();
snapshot::build(
Collected {
trees,
failed_projects,
read_at,
},
panes,
joined,
cfg,
agents,
filter,
now,
)
}
fn standing<'a>(&'a self, cfg: &'a Config) -> impl Iterator<Item = (&'a str, &'a Read)> {
cfg.read()
.filter_map(|p| Some((p.name.as_str(), self.read.get(&p.name)?)))
}
fn that_answered<'a>(
&'a self,
cfg: &'a Config,
) -> impl Iterator<Item = (&'a str, &'a ProjectWork)> {
self.standing(cfg)
.filter_map(|(project, read)| Some((project, read.work.as_ref().ok()?)))
}
fn every_pane(&mut self, agents: &dyn Agents) -> (Vec<Pane>, AgentProvider, BTreeSet<String>) {
let sessions = match agents.sessions() {
Ok(sessions) => sessions,
Err(failure) => {
return (
Vec::new(),
AgentProvider {
provider: agents.name(),
state: unlistable(failure.kind),
sessions: Vec::new(),
},
BTreeSet::new(),
)
}
};
self.panes_last_answered
.retain(|session, _| sessions.contains(session));
let mut panes = Vec::new();
let mut read = Vec::with_capacity(sessions.len());
let mut out_of_reach = BTreeSet::new();
for name in sessions {
let state = match agents.list(&name) {
Ok(listed) => {
self.panes_last_answered.insert(
name.clone(),
listed.iter().map(|pane| pane.pane_id.clone()).collect(),
);
panes.extend(listed.into_iter().map(placed));
SessionState::Answering
}
Err(_) => {
if let Some(held) = self.panes_last_answered.get(&name) {
out_of_reach.extend(held.iter().cloned());
}
SessionState::NotAnswering
}
};
read.push(Session { name, state });
}
(
panes,
AgentProvider::answering(agents.name(), read),
out_of_reach,
)
}
}
fn unlistable(kind: FailureKind) -> ProviderState {
match kind {
FailureKind::NotInstalled => ProviderState::Absent,
FailureKind::Unstartable
| FailureKind::InstalledUnstartable
| FailureKind::Auth
| FailureKind::Unavailable
| FailureKind::Gone
| FailureKind::Busy
| FailureKind::Parse
| FailureKind::Unsupported
| FailureKind::UnknownFlag => ProviderState::NotAnswering,
}
}
pub fn run(
cfg: &Config,
agents: &dyn Agents,
trackers: &dyn Trackers,
filter: Filter,
now: DateTime<Utc>,
) -> Snapshot {
Collection::default().collect(cfg, agents, trackers, &Wanted::Everything, filter, now)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::app::fixtures::*;
use crate::collect::agents::testing::{
in_session, named, pane, titled, Asked as AskedOfTheProvider, Fake as Provider, THE_FAKE,
};
use crate::collect::run::{FailureKind, RunFailure};
use crate::collect::tracker::testing::{Asked, Fake, Fakes};
use crate::collect::tracker::Tracker;
use crate::collect::worktree::testing::a_linked_worktree_git_made;
use crate::model::anomaly::Anomaly;
use crate::model::join::{BeadKey, Conflict};
use crate::model::snapshot::LoosePane;
use crate::model::types::testing::{key, A_SESSION};
use crate::model::types::PaneStatus;
use pretty_assertions::assert_eq;
use std::collections::BTreeSet;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Condvar, Mutex};
use std::time::Duration;
fn every_wanted() -> [&'static str; 2] {
match Wanted::Everything {
Wanted::Everything | Wanted::Project(_) => ["Wanted::Everything", "Wanted::Project"],
}
}
fn source_outside_tests() -> String {
let root = PathBuf::from(env!("CARGO_MANIFEST_DIR"));
let itself = root.join(file!());
let mut walking = vec![root.join("src")];
let mut read = String::new();
let mut left_itself_out = false;
while let Some(path) = walking.pop() {
for entry in std::fs::read_dir(&path).expect("the crate's own source") {
let found = entry.expect("a directory entry").path();
if found.is_dir() {
walking.push(found);
} else if found.extension().is_some_and(|kind| kind == "rs") {
if found == itself {
left_itself_out = true;
continue;
}
let text = std::fs::read_to_string(&found).expect("a source file");
read.push_str(text.split("\n#[cfg(test)]\n").next().unwrap_or_default());
}
}
}
assert!(
left_itself_out,
"{} was never met while walking the source, so this file was \
read into its own assertion and the check below proves nothing",
itself.display()
);
read
}
#[test]
fn something_that_is_not_a_test_asks_for_each_kind_of_collection() {
let source = source_outside_tests();
for wanted in every_wanted() {
assert!(
source.contains(wanted),
"{wanted} is built nowhere but in tests"
);
}
}
const UNSTAFFED_TREE: &str = r#"[
{"id":"orb-7","title":"lift the ground station","status":"blocked",
"priority":1,"issue_type":"epic"},
{"id":"orb-7.1","title":"re-point the dish","status":"open","parent":"orb-7",
"dependencies":[{"depends_on_id":"orb-7","type":"parent-child"}],
"priority":2,"issue_type":"task"},
{"id":"orb-7.2","title":"lay the feeder cable","status":"open","parent":"orb-7",
"dependencies":[{"depends_on_id":"orb-7","type":"parent-child"}],
"priority":2,"issue_type":"task"}
]"#;
fn tree_of<'a>(snap: &'a Snapshot, project: &str) -> &'a Tree {
snap.trees
.iter()
.find(|t| t.project == project)
.unwrap_or_else(|| panic!("{project} has a tree"))
}
#[test]
fn a_provider_that_will_not_answer_is_said_and_every_tree_still_draws() {
let no_session = Provider::unlistable(RunFailure {
kind: FailureKind::Unavailable,
program: "a provider".to_string(),
detail: "no such session".to_string(),
});
let snap = run(
&one_project(),
&no_session,
&orbital(),
Filter::LiveAgents,
now(),
);
assert_eq!(snap.agents.state, ProviderState::NotAnswering);
assert_eq!(snap.trees.len(), 1, "trees draw without liveness");
assert!(snap.trees[0].beads.iter().all(|n| n.agent.is_none()));
assert!(snap.unattributed.is_empty());
}
#[test]
fn a_provider_that_was_never_installed_is_absent_and_every_tree_still_draws() {
let nothing = Provider::unlistable(RunFailure::not_installed(
THE_FAKE,
"No such file or directory",
));
let snap = run(
&one_project(),
¬hing,
&orbital(),
Filter::LiveAgents,
now(),
);
assert_eq!(snap.agents.state, ProviderState::Absent);
assert_eq!(snap.agents.provider, THE_FAKE);
assert_eq!(snap.trees.len(), 1, "trees draw with no provider at all");
assert!(snap.trees[0].beads.iter().all(|n| n.agent.is_none()));
assert!(snap.unattributed.is_empty());
}
#[test]
fn a_provider_that_is_installed_and_will_not_start_is_a_finding_not_an_absence() {
let broken = Provider::unlistable(RunFailure::unstartable(
THE_FAKE,
"Permission denied (os error 13)",
));
let snap = run(
&one_project(),
&broken,
&orbital(),
Filter::LiveAgents,
now(),
);
assert_eq!(snap.agents.state, ProviderState::NotAnswering);
assert_eq!(snap.trees.len(), 1, "trees draw without liveness");
}
#[test]
fn a_provider_holding_no_pane_has_still_answered() {
let snap = run(&one_project(), &no_panes(), &orbital(), Filter::All, now());
assert_eq!(snap.agents.state, ProviderState::Answering);
}
fn orphaned(snap: &Snapshot) -> Vec<&str> {
snap.trees
.iter()
.flat_map(|tree| tree.beads.iter())
.filter(|node| {
node.anomalies
.iter()
.any(|fired| matches!(fired, Anomaly::OrphanClaim { .. }))
})
.map(|node| node.id.as_str())
.collect()
}
#[test]
fn a_claim_is_only_orphaned_against_a_provider_that_answered() {
let answered = run(&one_project(), &no_panes(), &orbital(), Filter::All, now());
assert_eq!(orphaned(&answered), ["orb-7", "orb-7.1"]);
for silent in [
Provider::unlistable(RunFailure::not_installed(
THE_FAKE,
"No such file or directory",
)),
Provider::unlistable(RunFailure::unstartable(
THE_FAKE,
"Permission denied (os error 13)",
)),
] {
let snap = run(&one_project(), &silent, &orbital(), Filter::All, now());
assert_eq!(
orphaned(&snap),
Vec::<&str>::new(),
"{:?} knows no more about panes than the other",
snap.agents.state
);
}
}
#[test]
fn a_tree_with_no_live_agent_is_reported_rather_than_dropped() {
let nobody = no_panes();
let trackers = orbital_with(
Fake::holding(beads(UNSTAFFED_TREE))
.ready(["orb-7.2"])
.blocked("orb-7", &["orb-9"]),
);
let filtered = run(
&one_project(),
&nobody,
&trackers,
Filter::LiveAgents,
now(),
);
assert!(filtered.trees.is_empty());
assert_eq!(filtered.hidden_trees.len(), 1);
assert_eq!(filtered.hidden_trees[0].root, "orb-7");
let all = run(&one_project(), &nobody, &trackers, Filter::All, now());
assert_eq!(all.trees.len(), 1);
assert!(all.hidden_trees.is_empty());
}
#[test]
fn a_pane_on_no_bead_is_reported_under_the_project_it_sits_in() {
let snap = run(
&one_project(),
&panes(),
&orbital(),
Filter::LiveAgents,
now(),
);
let loose: Vec<&str> = snap
.unattributed
.iter()
.map(|p| p.pane.id.as_str())
.collect();
assert_eq!(loose, vec!["w:p9"]);
assert_eq!(snap.unattributed[0].project, "orbital");
}
#[test]
fn a_pane_joins_only_the_project_its_directory_sits_in() {
let panes = Provider::holding(vec![named(
pane("w:p1", ORBITAL, PaneStatus::Working),
"x-1.1",
)]);
let snap = run(
&two_projects(),
&panes,
&colliding_trackers(),
Filter::All,
now(),
);
assert!(
node(tree_of(&snap, "orbital"), "x-1.1").agent.is_some(),
"the pane's own project"
);
assert!(
node(tree_of(&snap, "ferry"), "x-1.1").agent.is_none(),
"the same id in a tracker the pane is nowhere near"
);
}
fn panes_in_both() -> Provider {
Provider::holding(vec![
named(pane("w:p1", ORBITAL, PaneStatus::Working), "x-1.1"),
pane("w:p2", FERRY, PaneStatus::Idle),
])
}
const ALONE: Duration = Duration::from_secs(5);
struct Meeting {
inner: Fakes,
holds: Vec<(String, String)>,
arrived: Mutex<BTreeSet<String>>,
someone_arrived: Condvar,
waited_alone: Mutex<Vec<String>>,
}
impl Meeting {
fn at(inner: Fakes) -> Self {
Self {
inner,
holds: Vec::new(),
arrived: Mutex::new(BTreeSet::new()),
someone_arrived: Condvar::new(),
waited_alone: Mutex::new(Vec::new()),
}
}
fn holding(mut self, held: &str, until: &str) -> Self {
self.holds.push((held.to_string(), until.to_string()));
self
}
fn waited_alone(&self) -> Vec<String> {
self.waited_alone.lock().unwrap().clone()
}
fn arrived(&self) -> BTreeSet<String> {
self.arrived.lock().unwrap().clone()
}
fn met_by(&self, project: &str) {
let mut arrived = self.arrived.lock().unwrap();
arrived.insert(project.to_string());
self.someone_arrived.notify_all();
for (_, until) in self.holds.iter().filter(|(held, _)| held == project) {
let (still, waited) = self
.someone_arrived
.wait_timeout_while(arrived, ALONE, |arrived| !arrived.contains(until))
.unwrap();
arrived = still;
if waited.timed_out() {
self.waited_alone.lock().unwrap().push(project.to_string());
}
}
}
}
impl Trackers for Meeting {
fn of(&self, project: &Project) -> Result<Box<dyn Tracker + '_>, OpenFailure> {
Ok(Box::new(Held {
meeting: self,
project: project.name.clone(),
inner: self.inner.of(project)?,
}))
}
}
struct Held<'m> {
meeting: &'m Meeting,
project: String,
inner: Box<dyn Tracker + 'm>,
}
impl Tracker for Held<'_> {
fn fingerprint(&self) -> Option<Result<String, RunFailure>> {
self.meeting.met_by(&self.project);
self.inner.fingerprint()
}
fn all(&self) -> Result<Vec<crate::model::types::Bead>, RunFailure> {
self.inner.all()
}
fn ready(&self) -> Result<BTreeSet<String>, RunFailure> {
self.inner.ready()
}
fn blocked(&self) -> Result<BTreeMap<String, Vec<String>>, RunFailure> {
self.inner.blocked()
}
}
#[test]
fn the_projects_named_are_read_together_rather_than_in_turn() {
let trackers = Meeting::at(colliding_trackers())
.holding("orbital", "ferry")
.holding("ferry", "orbital");
collect(
&mut Collection::default(),
&panes_in_both(),
&trackers,
&Wanted::Everything,
);
assert_eq!(
trackers.arrived(),
BTreeSet::from(["orbital".to_string(), "ferry".to_string()]),
"both trackers were asked for a fingerprint, so a wait spent alone would have been recorded"
);
assert_eq!(
trackers.waited_alone(),
Vec::<String>::new(),
"a fingerprint that waited alone was asked for after the other project's read had finished"
);
}
fn collect(
collection: &mut Collection,
agents: &dyn Agents,
trackers: &dyn Trackers,
wanted: &Wanted,
) -> Snapshot {
collection.collect(
&two_projects(),
agents,
trackers,
wanted,
Filter::All,
now(),
)
}
fn orbital_alone() -> Wanted {
Wanted::Project("orbital".to_string())
}
fn asked_of(trackers: &Fakes, project: &str) -> usize {
trackers.tracker(project).asked().len()
}
const AGEING_TREE: &str = r#"[
{"id":"orb-7","title":"lift the ground station","status":"in_progress",
"priority":1,"issue_type":"epic","metadata":{"agent_pane":"w:p1"}},
{"id":"orb-7.1","title":"re-point the dish","status":"in_progress",
"dependencies":[{"depends_on_id":"orb-7","type":"parent-child"}],
"priority":2,"issue_type":"task",
"metadata":{"agent_pane":"w:p1"},
"updated_at":"2026-08-01T12:00:00Z"},
{"id":"orb-7.2","title":"lay the feeder cable","status":"open",
"dependencies":[{"depends_on_id":"orb-7","type":"parent-child"}],
"priority":2,"issue_type":"task"}
]"#;
fn a_claim_is_drawn_as_stale(snap: &Snapshot) -> bool {
snap.trees.iter().flat_map(|tree| &tree.beads).any(|node| {
node.anomalies
.iter()
.any(|fired| matches!(fired, Anomaly::StaleClaim { .. }))
})
}
const DEFERRED_TREE: &str = r#"[
{"id":"orb-7","title":"lift the ground station","status":"in_progress",
"priority":1,"issue_type":"epic"},
{"id":"orb-7.1","title":"re-point the dish","status":"in_progress",
"dependencies":[{"depends_on_id":"orb-7","type":"parent-child"}],
"priority":2,"issue_type":"task",
"metadata":{"agent_pane":"w:p1"}},
{"id":"orb-7.2","title":"lay the feeder cable","status":"open",
"dependencies":[{"depends_on_id":"orb-7","type":"parent-child"}],
"priority":2,"issue_type":"task",
"defer_until":"2026-08-30T14:00:00Z"}
]"#;
fn when_it_is_due() -> DateTime<Utc> {
"2026-08-30T14:00:00Z".parse().expect("the instant parses")
}
fn pane_on_a_bead() -> Provider {
Provider::holding(vec![named(
pane("w:p1", ORBITAL, PaneStatus::Working),
"orb-7.2",
)])
}
#[test]
fn a_project_whose_tracker_has_not_moved_is_asked_once() {
let trackers = orbital();
let cfg = one_project();
let mut standing = Collection::default();
standing.collect(
&cfg,
&panes(),
&trackers,
&Wanted::Everything,
Filter::All,
now(),
);
let first = asked_of(&trackers, "orbital");
standing.collect(
&cfg,
&panes(),
&trackers,
&orbital_alone(),
Filter::All,
now(),
);
assert_eq!(first, 4, "a project read for the first time costs both");
assert_eq!(
asked_of(&trackers, "orbital") - first,
1,
"and a project that has not moved since costs the fingerprint alone"
);
}
#[test]
fn a_project_whose_tracker_has_moved_is_read_in_full() {
let cfg = one_project();
let mut standing = Collection::default();
standing.collect(
&cfg,
&panes(),
&orbital(),
&Wanted::Everything,
Filter::All,
now(),
);
let moved = orbital_with(orbital_tracker().moved());
let after = standing.collect(&cfg, &panes(), &moved, &orbital_alone(), Filter::All, now());
assert_eq!(asked_of(&moved, "orbital"), 4);
assert!(
!trees_of(&after, "orbital").is_empty(),
"and everything it read is drawn"
);
}
#[test]
fn a_project_the_reader_has_rewritten_is_read_again_though_its_tracker_has_not_moved() {
let trackers = orbital();
let mut standing = Collection::default();
standing.collect(
&one_project(),
&panes(),
&trackers,
&Wanted::Everything,
Filter::All,
now(),
);
let first = asked_of(&trackers, "orbital");
let after = standing.collect(
&one_project_reached_with_a_credential(),
&panes(),
&trackers,
&Wanted::Everything,
Filter::All,
now(),
);
assert_eq!(
asked_of(&trackers, "orbital") - first,
4,
"the cascade ran again rather than being skipped against a read taken \
under the settings the reader has just replaced"
);
assert!(
!trees_of(&after, "orbital").is_empty(),
"and what it read under the new settings is drawn"
);
}
#[test]
fn a_cascade_that_failed_leaves_nothing_for_the_next_interval_to_skip_against() {
let cfg = one_project();
let mut standing = Collection::default();
standing.collect(
&cfg,
&panes(),
&orbital(),
&Wanted::Everything,
Filter::All,
now(),
);
let refused = orbital_with(
orbital_tracker()
.moved()
.failing(Asked::All, failing(FailureKind::Auth)),
);
let failed = standing.collect(
&cfg,
&panes(),
&refused,
&orbital_alone(),
Filter::All,
now(),
);
assert_eq!(failed.failed_projects.len(), 1, "the cascade failed");
let recovered = orbital_with(orbital_tracker().moved());
let after = standing.collect(
&cfg,
&panes(),
&recovered,
&orbital_alone(),
Filter::All,
now(),
);
assert_eq!(
asked_of(&recovered, "orbital"),
4,
"the cascade ran again rather than being skipped against the fingerprint the failure was taken at"
);
assert_eq!(after.failed_projects, vec![], "so the project recovered");
}
#[test]
fn a_tracker_that_cannot_answer_its_fingerprint_is_read_in_full_every_interval() {
let cfg = one_project();
let blind = orbital_with(
orbital_tracker().failing(Asked::Fingerprint, failing(FailureKind::Unavailable)),
);
let mut standing = Collection::default();
standing.collect(
&cfg,
&panes(),
&blind,
&Wanted::Everything,
Filter::All,
now(),
);
let first = asked_of(&blind, "orbital");
let after = standing.collect(&cfg, &panes(), &blind, &orbital_alone(), Filter::All, now());
assert_eq!(
first, 4,
"the fingerprint was asked for and the cascade ran anyway"
);
assert_eq!(
asked_of(&blind, "orbital") - first,
4,
"and again, rather than settling into a skip against a fingerprint nobody established"
);
assert!(
!trees_of(&after, "orbital").is_empty(),
"a tracker blind to its fingerprint still draws its trees"
);
}
#[test]
fn a_tracker_with_no_fingerprint_is_read_in_full_every_interval() {
let cfg = one_project();
let unfingerprinted = orbital_with(orbital_tracker().without_a_fingerprint());
let mut standing = Collection::default();
standing.collect(
&cfg,
&panes(),
&unfingerprinted,
&Wanted::Everything,
Filter::All,
now(),
);
let first = asked_of(&unfingerprinted, "orbital");
let after = standing.collect(
&cfg,
&panes(),
&unfingerprinted,
&orbital_alone(),
Filter::All,
now(),
);
assert_eq!(first, 4);
assert_eq!(asked_of(&unfingerprinted, "orbital") - first, 4);
assert!(!trees_of(&after, "orbital").is_empty());
}
#[test]
fn a_tracker_that_stops_answering_its_fingerprint_is_read_in_full_again() {
let cfg = one_project();
let mut standing = Collection::default();
standing.collect(
&cfg,
&panes(),
&orbital(),
&Wanted::Everything,
Filter::All,
now(),
);
let blind = orbital_with(
orbital_tracker().failing(Asked::Fingerprint, failing(FailureKind::Unavailable)),
);
standing.collect(&cfg, &panes(), &blind, &orbital_alone(), Filter::All, now());
let first = asked_of(&blind, "orbital");
standing.collect(&cfg, &panes(), &blind, &orbital_alone(), Filter::All, now());
assert_eq!(
first, 4,
"the fingerprint went unanswered, so the cascade ran rather than the standing one being kept"
);
assert_eq!(
asked_of(&blind, "orbital") - first,
4,
"and the read it just took left nothing for the next interval to skip against either"
);
}
#[test]
fn a_skipped_read_is_as_fresh_as_the_collection_that_skipped_it() {
let cfg = one_project();
let trackers = orbital();
let earlier = now();
let later = earlier + chrono::Duration::seconds(30);
let mut standing = Collection::default();
standing.collect(
&cfg,
&panes(),
&trackers,
&Wanted::Everything,
Filter::All,
earlier,
);
let after = standing.collect(
&cfg,
&panes(),
&trackers,
&orbital_alone(),
Filter::All,
later,
);
assert_eq!(after.read_at["orbital"], later);
assert!(
!trees_of(&after, "orbital").is_empty(),
"and everything the skipped read stood on is still drawn"
);
}
#[test]
fn a_skipped_read_still_ages_what_the_screen_says_about_it() {
let cfg = one_project();
let trackers = orbital_with(orbital_holding(AGEING_TREE));
let earlier = now();
let later = earlier + chrono::Duration::days(2);
let mut standing = Collection::default();
let before = standing.collect(
&cfg,
&panes(),
&trackers,
&Wanted::Everything,
Filter::All,
earlier,
);
let after = standing.collect(
&cfg,
&panes(),
&trackers,
&orbital_alone(),
Filter::All,
later,
);
assert!(
!a_claim_is_drawn_as_stale(&before),
"29 days is inside the window at the first collection"
);
assert!(
a_claim_is_drawn_as_stale(&after),
"and 31 days is outside it at the second, which read nothing"
);
assert_eq!(
asked_of(&trackers, "orbital"),
5,
"the second collection cost the fingerprint alone, so the ageing is the draw's and not the read's"
);
}
#[test]
fn a_read_stops_speaking_for_the_tracker_once_a_held_bead_is_due() {
let cfg = one_project();
let trackers = orbital_with(orbital_holding(DEFERRED_TREE));
let read_at = now();
let mut standing = Collection::default();
standing.collect(
&cfg,
&panes(),
&trackers,
&Wanted::Everything,
Filter::All,
read_at,
);
let first = asked_of(&trackers, "orbital");
let hour = chrono::Duration::hours(1);
standing.collect(
&cfg,
&panes(),
&trackers,
&orbital_alone(),
Filter::All,
read_at + hour,
);
let while_held = asked_of(&trackers, "orbital");
standing.collect(
&cfg,
&panes(),
&trackers,
&orbital_alone(),
Filter::All,
read_at + hour * 3,
);
let once_due = asked_of(&trackers, "orbital");
standing.collect(
&cfg,
&panes(),
&trackers,
&orbital_alone(),
Filter::All,
read_at + hour * 4,
);
assert_eq!(
while_held - first,
1,
"the fingerprint alone while the tracker is still holding the bead back"
);
assert_eq!(
once_due - while_held,
4,
"and the whole cascade at the first refresh past the instant it is due"
);
assert_eq!(
asked_of(&trackers, "orbital") - once_due,
1,
"after which nothing is held back, so the fingerprint alone again rather than for ever"
);
}
#[test]
fn a_read_has_stopped_speaking_at_the_instant_a_held_bead_is_due_rather_than_after_it() {
let cfg = one_project();
let trackers = orbital_with(orbital_holding(DEFERRED_TREE));
let mut standing = Collection::default();
standing.collect(
&cfg,
&panes(),
&trackers,
&Wanted::Everything,
Filter::All,
now(),
);
let first = asked_of(&trackers, "orbital");
standing.collect(
&cfg,
&panes(),
&trackers,
&orbital_alone(),
Filter::All,
when_it_is_due(),
);
assert_eq!(asked_of(&trackers, "orbital") - first, 4);
}
#[test]
fn a_bead_due_as_the_read_was_taken_leaves_nothing_to_turn_over() {
let cfg = one_project();
let trackers = orbital_with(orbital_holding(DEFERRED_TREE));
let mut standing = Collection::default();
standing.collect(
&cfg,
&panes(),
&trackers,
&Wanted::Everything,
Filter::All,
when_it_is_due(),
);
let first = asked_of(&trackers, "orbital");
let hour = chrono::Duration::hours(1);
standing.collect(
&cfg,
&panes(),
&trackers,
&orbital_alone(),
Filter::All,
when_it_is_due() + hour,
);
assert_eq!(asked_of(&trackers, "orbital") - first, 1);
}
#[test]
fn a_pane_that_starts_naming_a_bead_has_the_project_read_again() {
let cfg = one_project();
let mut standing = Collection::default();
standing.collect(
&cfg,
&panes(),
&orbital(),
&Wanted::Everything,
Filter::All,
now(),
);
let again = orbital();
standing.collect(
&cfg,
&pane_on_a_bead(),
&again,
&orbital_alone(),
Filter::All,
now(),
);
assert_eq!(
asked_of(&again, "orbital"),
4,
"the tracker had not moved, but what the panes name had"
);
}
fn orbital_refusing() -> Fakes {
Fakes::default()
.with(
"orbital",
colliding_tracker()
.failing(Asked::Fingerprint, failing(FailureKind::Auth))
.failing(Asked::All, failing(FailureKind::Auth)),
)
.with("ferry", colliding_tracker())
}
fn trees_of<'a>(snap: &'a Snapshot, project: &str) -> Vec<&'a Tree> {
snap.trees
.iter()
.filter(|t| t.project == project)
.map(Arc::as_ref)
.collect()
}
#[test]
fn a_pane_under_a_project_the_scope_left_out_is_not_reported() {
let scoped = two_projects()
.scoped_to(&["orbital".to_string()])
.expect("orbital is configured");
let snap = Collection::default().collect(
&scoped,
&panes_in_both(),
&colliding_trackers(),
&Wanted::Everything,
Filter::All,
now(),
);
assert_eq!(snap.unconfigured, vec![]);
assert!(
!snap.unattributed.iter().any(|pane| pane.pane.id == "w:p2"),
"ferry's pane was reported by a run reading orbital: {:#?}",
snap.unattributed
);
let named: Vec<&str> = scoped.projects.iter().map(|p| p.name.as_str()).collect();
assert_eq!(
named,
["orbital", "ferry"],
"the config as written is still reachable"
);
}
fn orbital_and_ferry_at(checkout: &Path) -> Config {
Config::from_toml(&format!(
r#"
[[projects]]
name = "orbital"
path = "{ORBITAL}"
[[projects]]
name = "ferry"
path = "{}"
"#,
checkout.display()
))
.expect("the config parses")
}
fn a_pane_sitting_in(cwd: &Path) -> Provider {
Provider::holding(vec![pane(
"w:p2",
&cwd.display().to_string(),
PaneStatus::Idle,
)])
}
#[test]
fn a_pane_in_a_linked_worktree_is_reported_in_the_project_its_main_working_tree_is_under() {
let fixture = a_linked_worktree_git_made("collection-linked-worktree");
let cfg = orbital_and_ferry_at(&fixture.checkout);
let snap = run(
&cfg,
&a_pane_sitting_in(&fixture.linked),
&colliding_trackers(),
Filter::All,
now(),
);
assert_eq!(
snap.unattributed,
vec![LoosePane {
pane: key("w:p2"),
project: "ferry".to_string(),
cwd: fixture.linked.display().to_string(),
pane_status: PaneStatus::Idle,
display_agent: None,
title: None,
claim_refused: false,
}],
"reported where it sits, placed by where its main working tree is"
);
assert_eq!(snap.unconfigured, vec![]);
}
#[test]
fn a_pane_in_an_excluded_projects_linked_worktree_is_placed_without_running_anything() {
let fixture = a_linked_worktree_git_made("collection-linked-worktree-excluded");
let cfg = orbital_and_ferry_at(&fixture.checkout)
.scoped_to(&["orbital".to_string()])
.expect("orbital is configured");
let provider = a_pane_sitting_in(&fixture.linked);
let snap = run(&cfg, &provider, &colliding_trackers(), Filter::All, now());
assert_eq!(snap.unattributed, vec![]);
assert_eq!(snap.unconfigured, vec![]);
assert_eq!(provider.asked(), one_session_read());
}
fn one_session_read() -> Vec<AskedOfTheProvider> {
vec![
AskedOfTheProvider::Sessions,
AskedOfTheProvider::List {
session: A_SESSION.to_string(),
},
]
}
#[test]
fn every_projects_trees_arrive_together_and_in_the_order_the_config_names() {
let snap = collect(
&mut Collection::default(),
&panes_in_both(),
&colliding_trackers(),
&Wanted::Everything,
);
let runs: Vec<&str> = snap
.trees
.chunk_by(|a, b| a.project == b.project)
.map(|run| run[0].project.as_str())
.collect();
let distinct: BTreeSet<&str> = runs.iter().copied().collect();
assert_eq!(runs, ["orbital", "ferry"], "{:#?}", snap.trees);
assert_eq!(
runs.len(),
distinct.len(),
"a project in two runs draws two project lines: {:#?}",
snap.trees
);
}
#[test]
fn refreshing_one_project_gives_the_snapshot_a_whole_rebuild_would_have() {
let panes = panes_in_both();
let trackers = colliding_trackers();
let mut standing = Collection::default();
collect(&mut standing, &panes, &trackers, &Wanted::Everything);
let refreshed = collect(&mut standing, &panes, &trackers, &orbital_alone());
let rebuilt = collect(
&mut Collection::default(),
&panes,
&trackers,
&Wanted::Everything,
);
assert_eq!(refreshed, rebuilt);
}
#[test]
fn refreshing_one_project_asks_no_other_projects_tracker() {
let panes = panes_in_both();
let trackers = colliding_trackers();
let mut standing = Collection::default();
collect(&mut standing, &panes, &trackers, &Wanted::Everything);
let (orbital, ferry) = (asked_of(&trackers, "orbital"), asked_of(&trackers, "ferry"));
collect(&mut standing, &panes, &trackers, &orbital_alone());
assert_eq!(
asked_of(&trackers, "ferry"),
ferry,
"ferry was not named, so its tracker was not asked again"
);
assert!(
asked_of(&trackers, "orbital") > orbital,
"orbital was named, so it was read"
);
}
#[test]
fn a_collection_asks_the_provider_for_its_sessions_once_and_each_session_once() {
let panes = Provider::holding(vec![
named(pane("w:p1", ORBITAL, PaneStatus::Working), "x-1.1"),
in_session(pane("w:p2", FERRY, PaneStatus::Idle), "beacon"),
]);
let trackers = colliding_trackers();
let mut standing = Collection::default();
let both_sessions_read = vec![
AskedOfTheProvider::Sessions,
AskedOfTheProvider::List {
session: A_SESSION.to_string(),
},
AskedOfTheProvider::List {
session: "beacon".to_string(),
},
];
collect(&mut standing, &panes, &trackers, &Wanted::Everything);
assert_eq!(
panes.asked(),
both_sessions_read,
"two projects were read, and each session on the machine was asked once"
);
collect(&mut standing, &panes, &trackers, &orbital_alone());
assert_eq!(
panes.asked(),
[both_sessions_read.clone(), both_sessions_read].concat(),
"a refresh naming one project reads every session on the machine, once each"
);
}
#[test]
fn a_session_that_will_not_answer_is_named_and_the_others_still_draw() {
let panes = Provider::holding(vec![named(
in_session(pane("w:p1", ORBITAL, PaneStatus::Working), "beacon"),
"x-1.1",
)])
.not_answering_for("standing-agents", wedged());
let snap = run(
&two_projects(),
&panes,
&colliding_trackers(),
Filter::All,
now(),
);
assert_eq!(snap.agents.state, ProviderState::Answering);
assert_eq!(
snap.agents.sessions,
vec![
Session {
name: A_SESSION.to_string(),
state: SessionState::Answering,
},
Session {
name: "beacon".to_string(),
state: SessionState::Answering,
},
Session {
name: "standing-agents".to_string(),
state: SessionState::NotAnswering,
},
]
);
assert_eq!(
snap.agents.unanswered().collect::<Vec<_>>(),
["standing-agents"]
);
let seat = node(tree_of(&snap, "orbital"), "x-1.1")
.agent
.as_ref()
.expect("the seat in beacon is on its bead");
assert_eq!(seat.pane.session, "beacon");
}
fn wedged() -> RunFailure {
RunFailure {
kind: FailureKind::Unavailable,
program: THE_FAKE.to_string(),
detail: "no socket".to_string(),
}
}
const SEATED_TREE: &str = r#"[
{"id":"orb-7","title":"lift the ground station","status":"open",
"priority":1,"issue_type":"epic"},
{"id":"orb-7.1","title":"re-point the dish","status":"in_progress","parent":"orb-7",
"dependencies":[{"depends_on_id":"orb-7","type":"parent-child"}],
"priority":2,"issue_type":"task","metadata":{"agent_pane":"w:p1"}},
{"id":"orb-7.2","title":"lay the feeder cable","status":"in_progress","parent":"orb-7",
"dependencies":[{"depends_on_id":"orb-7","type":"parent-child"}],
"priority":2,"issue_type":"task","metadata":{"agent_pane":"w:p2"}},
{"id":"orb-7.3","title":"tune the receiver","status":"in_progress","parent":"orb-7",
"dependencies":[{"depends_on_id":"orb-7","type":"parent-child"}],
"priority":2,"issue_type":"task","metadata":{"agent_pane":"w:p3"}}
]"#;
#[test]
fn a_claim_whose_seat_is_in_a_session_that_went_quiet_is_not_orphaned_and_the_rest_still_are() {
let cfg = one_project();
let trackers = orbital_with(orbital_holding(SEATED_TREE));
let mut standing = Collection::default();
let seated = Provider::holding(vec![
pane("w:p1", ORBITAL, PaneStatus::Working),
in_session(pane("w:p2", ORBITAL, PaneStatus::Working), "beacon"),
pane("w:p3", ORBITAL, PaneStatus::Working),
]);
let before = standing.collect(
&cfg,
&seated,
&trackers,
&Wanted::Everything,
Filter::All,
now(),
);
assert_eq!(
orphaned(&before),
Vec::<&str>::new(),
"every seat is on its pane while both sessions answer"
);
let quiet = Provider::holding(vec![pane("w:p1", ORBITAL, PaneStatus::Working)])
.not_answering_for("beacon", wedged());
let after = standing.collect(
&cfg,
&quiet,
&trackers,
&Wanted::Everything,
Filter::All,
now(),
);
assert_eq!(
after.agents.unanswered().collect::<Vec<_>>(),
["beacon"],
"the session that went quiet is still a finding of its own"
);
assert_eq!(
orphaned(&after),
["orb-7.3"],
"the seat that died in the session that answered is reported, and \
the one in the session that did not is not"
);
assert!(
node(tree_of(&after, "orbital"), "orb-7.1").agent.is_some(),
"the seat that answered is still drawn on its bead"
);
}
#[test]
fn a_claim_in_a_session_that_has_never_answered_is_orphaned_as_before() {
let quiet = Provider::holding(vec![pane("w:p1", ORBITAL, PaneStatus::Working)])
.not_answering_for("beacon", wedged());
let snap = run(
&one_project(),
&quiet,
&orbital_with(orbital_holding(SEATED_TREE)),
Filter::All,
now(),
);
assert_eq!(orphaned(&snap), ["orb-7.2", "orb-7.3"]);
}
#[test]
fn a_session_the_provider_has_stopped_running_takes_what_it_was_holding() {
let cfg = one_project();
let trackers = orbital_with(orbital_holding(SEATED_TREE));
let mut standing = Collection::default();
let both = Provider::holding(vec![
pane("w:p1", ORBITAL, PaneStatus::Working),
in_session(pane("w:p2", ORBITAL, PaneStatus::Working), "beacon"),
]);
standing.collect(
&cfg,
&both,
&trackers,
&Wanted::Everything,
Filter::All,
now(),
);
let alone = Provider::holding(vec![pane("w:p1", ORBITAL, PaneStatus::Working)]);
standing.collect(
&cfg,
&alone,
&trackers,
&Wanted::Everything,
Filter::All,
now(),
);
let quiet = Provider::holding(vec![pane("w:p1", ORBITAL, PaneStatus::Working)])
.not_answering_for("beacon", wedged());
let after = standing.collect(
&cfg,
&quiet,
&trackers,
&Wanted::Everything,
Filter::All,
now(),
);
assert_eq!(
orphaned(&after),
["orb-7.2", "orb-7.3"],
"the pane list of the session that went is not the new one's"
);
}
#[test]
fn a_pane_id_two_sessions_hold_is_awarded_to_nobody_and_reported() {
let panes = Provider::holding(vec![
titled(pane("w:p1", ORBITAL, PaneStatus::Working), "in default"),
titled(
in_session(pane("w:p1", ORBITAL, PaneStatus::Idle), "beacon"),
"in beacon",
),
]);
let snap = run(&one_project(), &panes, &orbital(), Filter::All, now());
let claimed = node(tree_of(&snap, "orbital"), "orb-7.1");
assert_eq!(claimed.agent, None, "neither pane is awarded");
assert_eq!(
claimed.anomalies,
vec![Anomaly::OrphanClaim {
refused: Some(Conflict::PaneIdInSeveralSessions {
bead: BeadKey {
project: "orbital".to_string(),
id: "orb-7.1".to_string(),
},
pane_id: "w:p1".to_string(),
sessions: vec![A_SESSION.to_string(), "beacon".to_string()],
}),
}]
);
assert_eq!(snap.conflicts.len(), 1);
let loose: Vec<(&str, &str)> = snap
.unattributed
.iter()
.map(|pane| (pane.pane.session.as_str(), pane.pane.id.as_str()))
.collect();
assert_eq!(
loose,
vec![(A_SESSION, "w:p1"), ("beacon", "w:p1")],
"both panes are still drawn, each under its session"
);
}
#[test]
fn a_project_a_scope_left_out_has_its_tracker_unasked() {
let trackers = colliding_trackers();
let scoped = two_projects()
.scoped_to(&["orbital".to_string()])
.expect("orbital is configured");
Collection::default().collect(
&scoped,
&panes_in_both(),
&trackers,
&Wanted::Everything,
Filter::All,
now(),
);
assert_eq!(
asked_of(&trackers, "ferry"),
0,
"ferry was scoped out, so nothing should have gone near its tracker"
);
assert!(
asked_of(&trackers, "orbital") > 0,
"orbital was scoped in, so it was read"
);
}
#[test]
fn a_tracker_that_fails_while_one_project_refreshes_leaves_the_others_drawn() {
let panes = panes_in_both();
let mut standing = Collection::default();
let before = collect(
&mut standing,
&panes,
&colliding_trackers(),
&Wanted::Everything,
);
let after = collect(&mut standing, &panes, &orbital_refusing(), &orbital_alone());
assert_eq!(
after.failed_projects,
vec![FailedProject {
project: "orbital".to_string(),
tracker: TrackerFailure::Auth,
}]
);
assert!(
trees_of(&after, "orbital").is_empty(),
"orbital's trees went with the tracker that could not be read"
);
assert_eq!(
trees_of(&after, "ferry"),
trees_of(&before, "ferry"),
"ferry is exactly what it was before orbital's outage"
);
}
#[test]
fn a_refresh_naming_one_project_still_reads_the_provider() {
let trackers = colliding_trackers();
let mut standing = Collection::default();
let before = collect(&mut standing, &no_panes(), &trackers, &Wanted::Everything);
assert!(node(tree_of(&before, "ferry"), "x-1.1").agent.is_none());
let arrived = Provider::holding(vec![named(
pane("w:p2", FERRY, PaneStatus::Working),
"x-1.1",
)]);
let after = collect(&mut standing, &arrived, &trackers, &orbital_alone());
assert!(
node(tree_of(&after, "ferry"), "x-1.1").agent.is_some(),
"the pane reached the project the refresh did not name"
);
}
#[test]
fn a_refresh_naming_one_project_dates_that_project_and_leaves_the_rest_alone() {
let panes = panes_in_both();
let trackers = colliding_trackers();
let mut standing = Collection::default();
let cfg = two_projects();
let earlier = now();
let later = earlier + chrono::Duration::seconds(30);
standing.collect(
&cfg,
&panes,
&trackers,
&Wanted::Everything,
Filter::All,
earlier,
);
let after = standing.collect(
&cfg,
&panes,
&trackers,
&orbital_alone(),
Filter::All,
later,
);
assert_eq!(
after.read_at,
BTreeMap::from([
("orbital".to_string(), later),
("ferry".to_string(), earlier),
]),
"only the project the refresh named was read again"
);
assert_eq!(
after.generated_at, later,
"the snapshot is still drawn at the later instant"
);
}
#[test]
fn a_read_that_failed_is_still_dated_by_the_attempt_that_failed() {
let panes = panes_in_both();
let mut standing = Collection::default();
let cfg = two_projects();
let earlier = now();
let later = earlier + chrono::Duration::seconds(30);
standing.collect(
&cfg,
&panes,
&colliding_trackers(),
&Wanted::Everything,
Filter::All,
earlier,
);
let after = standing.collect(
&cfg,
&panes,
&orbital_refusing(),
&orbital_alone(),
Filter::All,
later,
);
assert!(
trees_of(&after, "orbital").is_empty(),
"nothing orbital's earlier read produced is still drawn"
);
assert_eq!(after.read_at["orbital"], later);
}
}