use std::path::Path;
use std::process::Command;
use std::sync::mpsc::Sender;
use std::thread::JoinHandle;
use std::time::{Duration, Instant};
use onevcs::rules::RuleMatch;
use onevcs::{IdentityOutcome, Scope, SlotOutcome, Span};
use serde::{Deserialize, Deserializer, Serialize};
use serde_json::{json, Map, Value};
use crate::engine::Message;
use crate::error::{Error, Result};
use crate::ledger::RunPaths;
pub const FLAG: &str = "--maintenance-config";
pub const KEY: &str = "maintenance_config";
pub const SCHEDULE_VERSION: u32 = 1;
pub const PACE_ENV: &str = "ONEPIPELINE_MAINTENANCE_PACE_SECONDS";
pub const DEFAULT_PACE_SECONDS: u64 = 600;
const RECHECK: Duration = Duration::from_secs(60);
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct MaintenanceConfig {
pub version: u32,
pub default: MaintenanceDefault,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub rules: Vec<MaintenanceRule>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct MaintenanceDefault {
#[serde(deserialize_with = "every")]
pub every: Span,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct MaintenanceRule {
#[serde(rename = "match", deserialize_with = "matching")]
pub r#match: RuleMatch,
#[serde(deserialize_with = "every")]
pub every: Span,
}
impl MaintenanceConfig {
pub fn load(spelling: &str, path: &Path) -> Result<Self> {
let text = std::fs::read_to_string(path).map_err(|source| Error::Ledger {
path: path.to_path_buf(),
source,
})?;
let named = |why: String| Error::Invalid(format!("{spelling} {}: {why}", path.display()));
let config: Self =
serde_norway::from_str(&text).map_err(|failure| named(failure.to_string()))?;
if config.version != SCHEDULE_VERSION {
return Err(named(format!(
"`version` is {}, and this build reads {SCHEDULE_VERSION} — set `version: \
{SCHEDULE_VERSION}`",
config.version
)));
}
Ok(config)
}
pub fn every_for(&self, identity: &str) -> std::result::Result<Span, String> {
let criteria: Vec<RuleMatch> = self.rules.iter().map(|rule| rule.r#match.clone()).collect();
let matched =
onevcs::first_matching(&criteria, identity).map_err(|failure| failure.to_string())?;
Ok(matched.map_or(self.default.every, |index| self.rules[index].every))
}
}
fn every<'de, D: Deserializer<'de>>(deserializer: D) -> std::result::Result<Span, D::Error> {
let value = Value::deserialize(deserializer)?;
let text = match &value {
Value::String(text) => text.as_str(),
other => {
return Err(serde::de::Error::custom(format!(
"`every` holds {other}, which is not a span — write digits then one of s, m, h \
or d, such as 7d"
)))
}
};
text.parse::<Span>()
.map_err(|why| serde::de::Error::custom(format!("`every`: {why}")))
}
fn matching<'de, D: Deserializer<'de>>(
deserializer: D,
) -> std::result::Result<RuleMatch, D::Error> {
let written = Map::<String, Value>::deserialize(deserializer)?;
let read: RuleMatch = serde_json::from_value(Value::Object(written.clone()))
.map_err(|why| serde::de::Error::custom(format!("`match`: {why}")))?;
let kept = serde_json::to_value(&read)
.ok()
.and_then(|value| value.as_object().cloned())
.unwrap_or_default();
if let Some(unknown) = written.keys().find(|key| !kept.contains_key(*key)) {
let known: Vec<String> = kept.keys().cloned().collect();
return Err(serde::de::Error::custom(format!(
"`match` names `{unknown}`, which is not a field a match has{}",
if known.is_empty() {
String::new()
} else {
format!("; it names {}", known.join(", "))
}
)));
}
if kept.is_empty() {
return Err(serde::de::Error::custom(
"`match` names no field — a rule matches on at least one of the identity's host, \
owner, name or path",
));
}
Ok(read)
}
fn pace() -> Duration {
Duration::from_secs(
std::env::var(PACE_ENV)
.ok()
.and_then(|value| value.parse().ok())
.filter(|seconds| *seconds > 0)
.unwrap_or(DEFAULT_PACE_SECONDS),
)
}
pub(crate) fn host_has_room() -> bool {
crate::executor::Executor::capacity(&crate::executor::LocalExecutor).slots_free > 0
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct Maintained {
pub(crate) identity: String,
pub(crate) every: Span,
pub(crate) outcome: std::result::Result<IdentityOutcome, String>,
}
impl Maintained {
fn is_recorded(&self) -> bool {
match &self.outcome {
Err(_) | Ok(IdentityOutcome::Claimed { .. }) => true,
Ok(IdentityOutcome::Slots(slots)) => slots
.iter()
.any(|slot| matches!(slot.outcome, SlotOutcome::Ran { .. })),
Ok(IdentityOutcome::NoMaintainCommand | IdentityOutcome::NoSlots) => false,
}
}
fn payload(&self) -> Value {
let mut entry = json!({
"identity": self.identity,
"every": self.every.to_string(),
});
let answer = match &self.outcome {
Ok(outcome) => serde_json::to_value(outcome)
.map_err(|why| format!("the sibling's outcome could not be recorded: {why}")),
Err(why) => Err(why.clone()),
};
match answer {
Ok(outcome) => entry["outcome"] = outcome,
Err(why) => entry["error"] = json!(crate::engine::bounded(&why)),
}
entry
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct Swept {
pub(crate) started_at: String,
pub(crate) identities: Vec<Maintained>,
pub(crate) failure: Option<String>,
}
impl Swept {
fn is_recorded(&self) -> bool {
self.failure.is_some() || self.identities.iter().any(Maintained::is_recorded)
}
fn payload(&self) -> Map<String, Value> {
let mut payload = crate::journal::payload(&[
("started_at", json!(self.started_at)),
(
"identities",
Value::Array(
self.identities
.iter()
.filter(|identity| identity.is_recorded())
.map(Maintained::payload)
.collect(),
),
),
]);
if let Some(failure) = &self.failure {
payload.insert("error".to_owned(), json!(crate::engine::bounded(failure)));
}
payload
}
}
fn identities() -> std::result::Result<Vec<String>, String> {
let binary = crate::destination::binary();
let output = Command::new(&binary)
.arg("repos")
.output()
.map_err(|error| {
format!(
"`{} repos` could not be run: {error} (set {} to an executable one)",
binary.to_string_lossy(),
crate::destination::BINARY_ENV
)
})?;
if !output.status.success() {
return Err(format!(
"`{} repos` refused: {}",
binary.to_string_lossy(),
String::from_utf8_lossy(&output.stderr).trim()
));
}
let listed = String::from_utf8(output.stdout).map_err(|error| {
format!(
"`{} repos` answered bytes that are not UTF-8: {error}",
binary.to_string_lossy()
)
})?;
let mut keys: Vec<String> = Vec::new();
for line in listed.lines() {
if line.is_empty() || line.starts_with(char::is_whitespace) {
continue;
}
if line == "no repositories registered" {
continue;
}
let Some((key, _gate)) = line.split_once('\t') else {
return Err(format!(
"`{} repos` printed a line this build does not read as an identity: {line:?}",
binary.to_string_lossy()
));
};
keys.push(key.to_owned());
} keys.sort_unstable();
keys.dedup();
Ok(keys)
}
fn sweep(config: &MaintenanceConfig) -> Swept {
let started_at = crate::sys::now_rfc3339();
let keys = match identities() {
Ok(keys) => keys,
Err(failure) => {
return Swept {
started_at,
identities: Vec::new(),
failure: Some(failure),
}
}
};
let identities = keys
.into_iter()
.map(|identity| {
let every = match config.every_for(&identity) {
Ok(every) => every,
Err(why) => {
return Maintained {
identity,
every: config.default.every,
outcome: Err(why),
}
}
};
let outcome = onevcs::pool_maintain(Scope::Repo(identity.clone()), Some(every))
.map_err(|failure| failure.to_string())
.and_then(|report| {
report
.identities
.into_iter()
.next()
.map(|maintained| maintained.outcome)
.ok_or_else(|| "the sibling answered for no identity".to_owned())
});
Maintained {
identity,
every,
outcome,
}
})
.collect();
Swept {
started_at,
identities,
failure: None,
}
}
pub(crate) struct Sweep {
handle: Option<JoinHandle<()>>,
paths: RunPaths,
}
impl Sweep {
fn start(config: MaintenanceConfig, paths: &RunPaths, tx: Sender<Message>) -> Option<Self> {
let started_at = crate::sys::now_rfc3339();
let _ = crate::ledger::write_json(
&paths.maintenance(),
&json!({"started_at": started_at, "pid": crate::sys::pid()}),
); let handle = std::thread::Builder::new()
.name("pool-maintenance".to_owned())
.spawn(move || {
let swept = sweep(&config);
let _ = tx.send(Message::Maintained(Box::new(swept)));
});
match handle {
Ok(handle) => Some(Self {
handle: Some(handle),
paths: paths.clone(),
}),
Err(error) => {
eprintln!("onepipeline: cannot start the pool-maintenance sweep: {error}");
let _ = std::fs::remove_file(paths.maintenance());
None
} }
}
fn join(&mut self) {
if let Some(handle) = self.handle.take() {
let _ = handle.join();
}
let _ = std::fs::remove_file(self.paths.maintenance());
}
}
impl Drop for Sweep {
fn drop(&mut self) {
self.join();
}
}
pub(crate) fn status_line(view: &crate::views::RunView) -> String {
if view.liveness() == crate::views::DriverLiveness::Driving {
if let Some(marker) = crate::ledger::read_json_opt::<Value>(&view.paths.maintenance()) {
let since = marker["started_at"]
.as_str()
.unwrap_or("an unrecorded time");
return format!(
" pool maintenance: a sweep of the host's worktree pools is in progress, \
started {since}\n"
);
}
}
String::new()
}
pub(crate) fn results_lines(view: &crate::views::RunView) -> String {
let Some(record) = view.events.iter().rev().find(|event| {
crate::journal::PipelineKind::from_wire(&event.kind)
== Some(crate::journal::PipelineKind::PoolMaintenance)
}) else {
return String::new();
};
let mut out = format!(
" pool maintenance: last sweep that did something started {}\n",
record
.payload
.get("started_at")
.and_then(Value::as_str)
.map_or_else(
|| "at an unrecorded time".to_owned(),
crate::views::one_line
)
);
if let Some(error) = record.payload.get("error").and_then(Value::as_str) {
out.push_str(&format!(
" the host's identities could not be enumerated: {}\n",
crate::views::one_line(error)
));
}
for entry in record
.payload
.get("identities")
.and_then(Value::as_array)
.into_iter()
.flatten()
{
let identity = entry
.get("identity")
.and_then(Value::as_str)
.map_or_else(|| "an unnamed identity".to_owned(), crate::views::one_line);
let every = entry
.get("every")
.and_then(Value::as_str)
.map_or_else(|| "?".to_owned(), crate::views::one_line);
let what = match (
entry.get("error").and_then(Value::as_str),
entry.get("outcome"),
) {
(Some(error), _) => format!("failed — {}", crate::views::one_line(error)),
(None, Some(outcome)) => {
match serde_json::from_value::<IdentityOutcome>(outcome.clone()) {
Ok(outcome) => identity_phrase(&outcome),
Err(_) => "an outcome this build does not read".to_owned(),
}
}
(None, None) => "no outcome recorded".to_owned(),
};
out.push_str(&format!(" {identity} (every {every}): {what}\n"));
}
out
}
fn identity_phrase(outcome: &IdentityOutcome) -> String {
match outcome {
IdentityOutcome::NoMaintainCommand => "no maintain command".to_owned(),
IdentityOutcome::NoSlots => "no slots".to_owned(),
IdentityOutcome::Claimed { by_pid } => {
format!("claimed — another pool maintain (pid {by_pid}) was maintaining it")
}
IdentityOutcome::Slots(slots) => slots
.iter()
.map(|slot| format!("slot {} {}", slot.number, slot_phrase(&slot.outcome)))
.collect::<Vec<_>>()
.join("; "),
}
}
fn slot_phrase(outcome: &SlotOutcome) -> String {
match outcome {
SlotOutcome::NotDue { last_maintained } => {
format!(
"not due, last maintained {}",
crate::views::one_line(last_maintained)
)
}
SlotOutcome::InUse { session } => {
format!(
"kept: session {} is working in it",
crate::views::one_line(&session.0)
)
}
SlotOutcome::Broken { reason } => format!("kept: {}", crate::views::one_line(reason)),
SlotOutcome::Ran {
outcome,
duration_ms,
log,
} => {
let ended = match outcome {
onevcs::MaintenanceOutcome::Succeeded => "succeeded".to_owned(),
onevcs::MaintenanceOutcome::Failed { exit: Some(exit) } => {
format!("failed (exit {exit})")
}
onevcs::MaintenanceOutcome::Failed { exit: None } => {
"failed (ended by a signal, or never started)".to_owned()
}
onevcs::MaintenanceOutcome::TimedOut => "timed out".to_owned(),
};
let log = log
.as_ref()
.map(|log| format!(", log {}", crate::views::one_line(&log.0)))
.unwrap_or_default();
format!("ran — {ended} in {duration_ms} ms{log}")
}
}
}
pub(crate) struct Maintenance {
config: Option<MaintenanceConfig>,
pace: Duration,
last_started: Option<Instant>,
sweep: Option<Sweep>,
}
impl Maintenance {
pub(crate) fn of_launch(config: Option<MaintenanceConfig>, paths: &RunPaths) -> Self {
let _ = std::fs::remove_file(paths.maintenance());
Self {
config,
pace: pace(),
last_started: None,
sweep: None,
}
}
fn is_scheduled(&self) -> bool {
self.config.is_some()
}
pub(crate) fn consider(&mut self, idle: bool, paths: &RunPaths, tx: &Sender<Message>) {
let Some(config) = &self.config else {
return;
};
if !idle || self.sweep.is_some() || !crate::engine::due(self.last_started, self.pace) {
return;
}
self.last_started = Some(Instant::now());
self.sweep = Sweep::start(config.clone(), paths, tx.clone());
}
pub(crate) fn next_due(&self) -> Duration {
if !self.is_scheduled() || self.sweep.is_some() {
return Duration::MAX;
}
let until = self.last_started.map_or(Duration::ZERO, |last| {
self.pace.saturating_sub(last.elapsed())
});
if until.is_zero() {
RECHECK
} else {
until
}
}
pub(crate) fn record(
&mut self,
paths: &RunPaths,
journal: &mut crate::journal::Journal,
swept: &Swept,
) -> Result<()> {
if let Some(mut sweep) = self.sweep.take() {
sweep.join();
}
if let Some(failure) = &swept.failure {
eprintln!("onepipeline: the pool-maintenance sweep could not enumerate this host's identities: {failure}");
}
if !swept.is_recorded() {
return Ok(());
}
journal.emit(
crate::journal::PipelineKind::PoolMaintenance,
crate::journal::labels(&paths.run, None),
swept.payload(),
)
}
pub(crate) fn close(
mut self,
paths: &RunPaths,
journal: &mut crate::journal::Journal,
rx: &std::sync::mpsc::Receiver<Message>,
) -> Result<()> {
let Some(mut sweep) = self.sweep.take() else {
return Ok(());
};
sweep.join();
drop(sweep);
let swept = rx.try_iter().find_map(|message| match message {
Message::Maintained(swept) => Some(swept),
_ => None,
});
match swept {
Some(swept) => self.record(paths, journal, &swept),
None => Ok(()),
}
}
#[cfg(test)]
pub(crate) fn is_sweeping(&self) -> bool {
self.sweep.is_some()
}
}
#[cfg(test)]
mod tests {
use super::*;
use onevcs::{IdentityMaintenance, MaintenanceOutcome, SlotMaintenance};
fn written(text: &str) -> Result<MaintenanceConfig> {
static NTH: std::sync::atomic::AtomicU32 = std::sync::atomic::AtomicU32::new(0);
let dir = std::env::temp_dir().join(format!(
"onepipeline-maintenance-{}-{}",
std::process::id(),
NTH.fetch_add(1, std::sync::atomic::Ordering::SeqCst)
));
std::fs::create_dir_all(&dir).expect("a scratch directory");
let path = dir.join("maintenance.yml");
std::fs::write(&path, text).expect("the schedule is written");
let read = MaintenanceConfig::load(FLAG, &path);
let _ = std::fs::remove_dir_all(&dir);
read
}
fn refused(text: &str) -> String {
match written(text) {
Err(Error::Invalid(why)) => why,
other => panic!("{text:?} was not refused as invalid: {other:?}"),
}
}
#[test]
fn the_schedule_reads_through_the_siblings_span_and_match_and_round_trips() {
let config = written(
"version: 1\ndefault:\n every: 7d\nrules:\n - match: {host: github.com, owner: \
nickderobertis, name: onevcs}\n every: 3d\n",
)
.expect("the schedule reads");
assert_eq!(config.version, SCHEDULE_VERSION);
assert_eq!(config.default.every.to_string(), "7d");
assert_eq!(config.rules.len(), 1);
assert_eq!(config.rules[0].every.to_string(), "3d");
assert_eq!(
config.rules[0].r#match,
RuleMatch {
host: Some("github.com".into()),
owner: Some("nickderobertis".into()),
name: Some("onevcs".into()),
path: None,
}
);
let text = serde_json::to_string(&config).expect("it serializes");
assert_eq!(
serde_json::from_str::<MaintenanceConfig>(&text).expect("it re-parses"),
config
);
let bare = written("version: 1\ndefault:\n every: 36h\n").expect("it reads");
assert!(bare.rules.is_empty());
assert_eq!(
serde_json::to_string(&bare).expect("serializes"),
r#"{"version":1,"default":{"every":"36h"}}"#
);
}
#[test]
fn a_schedule_is_refused_by_the_key_at_fault() {
let why = refused("version: 1\ndefault:\n every: 7d\nat: \"03:00\"\n");
assert!(why.starts_with(FLAG), "{why}");
assert!(why.contains("unknown field `at`"), "{why}");
let why = refused("version: 2\ndefault:\n every: 7d\n");
assert!(why.contains("`version` is 2"), "{why}");
assert!(why.contains("this build reads 1"), "{why}");
let why =
refused("version: 1\ndefault:\n every: 7d\nrules:\n - match: {}\n every: 1d\n");
assert!(why.contains("`match` names no field"), "{why}");
let why = refused(
"version: 1\ndefault:\n every: 7d\nrules:\n - match: {repo: onevcs}\n every: 1d\n",
);
assert!(why.contains("`match` names `repo`"), "{why}");
let why = refused("version: 1\ndefault:\n every: 7x\n");
assert!(why.contains("`every`"), "{why}");
assert!(why.contains("not a unit letter"), "{why}");
let why = refused("version: 1\ndefault:\n every: 7\n");
assert!(why.contains("`every` holds 7"), "{why}");
let why = refused("version: 1\ndefault:\n every: 7d\n at: never\n");
assert!(why.contains("unknown field `at`"), "{why}");
let why = refused("version: 1\ndefault:\n every: 7d\nrules:\n - every: 1d\n");
assert!(why.contains("missing field `match`"), "{why}");
}
#[test]
fn the_pace_is_read_from_the_environment() {
let _held = crate::vcs::scratch_home_held();
for unusable in ["", "0", "-1", "soon"] {
std::env::set_var(PACE_ENV, unusable);
assert_eq!(
pace(),
Duration::from_secs(DEFAULT_PACE_SECONDS),
"{PACE_ENV}={unusable:?}"
);
}
std::env::set_var(PACE_ENV, "7");
assert_eq!(pace(), Duration::from_secs(7));
std::env::remove_var(PACE_ENV);
assert_eq!(pace(), Duration::from_secs(DEFAULT_PACE_SECONDS));
}
fn ran(identity: &str) -> Maintained {
Maintained {
identity: identity.into(),
every: "7d".parse().expect("a span"),
outcome: Ok(IdentityOutcome::Slots(vec![
SlotMaintenance {
number: 1,
outcome: SlotOutcome::Ran {
outcome: MaintenanceOutcome::Succeeded,
duration_ms: 12,
log: None,
},
},
SlotMaintenance {
number: 2,
outcome: SlotOutcome::NotDue {
last_maintained: "2026-09-19T00:00:00.000Z".into(),
},
},
])),
}
}
fn quiet(identity: &str, outcome: IdentityOutcome) -> Maintained {
Maintained {
identity: identity.into(),
every: "7d".parse().expect("a span"),
outcome: Ok(outcome),
}
}
#[test]
fn a_sweep_is_recorded_only_where_something_ran_was_claimed_or_failed() {
let nothing = Swept {
started_at: "2026-09-20T00:00:00.000Z".into(),
identities: vec![
quiet("a", IdentityOutcome::NoMaintainCommand),
quiet("b", IdentityOutcome::NoSlots),
quiet(
"c",
IdentityOutcome::Slots(vec![SlotMaintenance {
number: 1,
outcome: SlotOutcome::NotDue {
last_maintained: "2026-09-19T00:00:00.000Z".into(),
},
}]),
),
],
failure: None,
};
assert!(!nothing.is_recorded());
let something = Swept {
started_at: "2026-09-20T00:00:00.000Z".into(),
identities: vec![
quiet("a", IdentityOutcome::NoMaintainCommand),
ran("b"),
quiet("c", IdentityOutcome::Claimed { by_pid: 42 }),
Maintained {
identity: "d".into(),
every: "1d".parse().expect("a span"),
outcome: Err("the registry could not be read".into()),
},
],
failure: None,
};
assert!(something.is_recorded());
let payload = Value::Object(something.payload());
assert_eq!(payload["started_at"], "2026-09-20T00:00:00.000Z");
assert!(payload.get("error").is_none());
let identities = payload["identities"].as_array().expect("an array");
assert_eq!(
identities
.iter()
.map(|entry| entry["identity"].as_str().expect("a key"))
.collect::<Vec<_>>(),
["b", "c", "d"],
"{payload}"
);
assert_eq!(identities[0]["every"], "7d");
assert_eq!(
identities[0]["outcome"]["slots"][0]["outcome"]["ran"]["outcome"],
"succeeded"
);
assert_eq!(
identities[0]["outcome"]["slots"][1]["outcome"]["not-due"]["last_maintained"],
"2026-09-19T00:00:00.000Z"
);
assert_eq!(identities[1]["outcome"]["claimed"]["by_pid"], 42);
assert_eq!(identities[2]["error"], "the registry could not be read");
assert!(identities[2].get("outcome").is_none());
let read: IdentityOutcome = serde_json::from_value(identities[0]["outcome"].clone())
.expect("the sibling reads its own outcome back");
assert_eq!(read, ran("b").outcome.expect("an outcome"));
let _ = IdentityMaintenance {
identity: "b".into(),
outcome: read,
};
let failed = Swept {
started_at: "2026-09-20T00:00:00.000Z".into(),
identities: Vec::new(),
failure: Some("`onevcs repos` could not be run".into()),
};
assert!(failed.is_recorded());
let payload = Value::Object(failed.payload());
assert_eq!(payload["error"], "`onevcs repos` could not be run");
assert_eq!(payload["identities"], json!([]));
}
#[test]
fn every_slot_outcome_has_a_phrase() {
let phrases = [
(
SlotOutcome::NotDue {
last_maintained: "2026-09-19T00:00:00.000Z".into(),
},
"not due, last maintained 2026-09-19T00:00:00.000Z",
),
(
SlotOutcome::InUse {
session: onevcs::SessionToken("s-abc".into()),
},
"kept: session s-abc is working in it",
),
(
SlotOutcome::Broken {
reason: "its record could not be read".into(),
},
"kept: its record could not be read",
),
(
SlotOutcome::Ran {
outcome: MaintenanceOutcome::Succeeded,
duration_ms: 12,
log: Some(onevcs::ArtifactId("a-1".into())),
},
"ran — succeeded in 12 ms, log a-1",
),
(
SlotOutcome::Ran {
outcome: MaintenanceOutcome::Failed { exit: Some(3) },
duration_ms: 12,
log: None,
},
"ran — failed (exit 3) in 12 ms",
),
(
SlotOutcome::Ran {
outcome: MaintenanceOutcome::Failed { exit: None },
duration_ms: 12,
log: None,
},
"ran — failed (ended by a signal, or never started) in 12 ms",
),
(
SlotOutcome::Ran {
outcome: MaintenanceOutcome::TimedOut,
duration_ms: 1_000,
log: None,
},
"ran — timed out in 1000 ms",
),
];
for (outcome, phrase) in phrases {
assert_eq!(slot_phrase(&outcome), phrase);
}
assert_eq!(
identity_phrase(&IdentityOutcome::Claimed { by_pid: 7 }),
"claimed — another pool maintain (pid 7) was maintaining it"
);
assert_eq!(
identity_phrase(&IdentityOutcome::NoMaintainCommand),
"no maintain command"
);
assert_eq!(identity_phrase(&IdentityOutcome::NoSlots), "no slots");
assert_eq!(
identity_phrase(&IdentityOutcome::Slots(vec![
SlotMaintenance {
number: 1,
outcome: SlotOutcome::NotDue {
last_maintained: "t".into()
},
},
SlotMaintenance {
number: 2,
outcome: SlotOutcome::Broken { reason: "r".into() },
},
])),
"slot 1 not due, last maintained t; slot 2 kept: r"
);
}
#[test]
fn the_pace_decides_the_wait_and_only_an_idle_pass_starts_a_sweep() {
let root = std::env::temp_dir().join(format!(
"onepipeline-maintenance-pace-{}",
std::process::id()
));
let paths = RunPaths::under(&root, "paced");
let (tx, _rx) = std::sync::mpsc::channel();
let none = Maintenance::of_launch(None, &paths);
assert_eq!(none.next_due(), Duration::MAX);
assert!(!none.is_sweeping());
let config = MaintenanceConfig {
version: SCHEDULE_VERSION,
default: MaintenanceDefault {
every: "7d".parse().expect("a span"),
},
rules: Vec::new(),
};
let mut scheduled = Maintenance::of_launch(Some(config), &paths);
scheduled.pace = Duration::from_secs(3_600);
assert_eq!(scheduled.next_due(), RECHECK);
scheduled.consider(false, &paths, &tx);
assert!(!scheduled.is_sweeping());
assert_eq!(scheduled.next_due(), RECHECK);
scheduled.last_started = Some(Instant::now());
scheduled.consider(true, &paths, &tx);
assert!(!scheduled.is_sweeping());
let until = scheduled.next_due();
assert!(
until > Duration::from_secs(3_500) && until <= Duration::from_secs(3_600),
"{until:?}"
);
let _ = std::fs::remove_dir_all(&root);
}
}