use std::collections::HashMap;
use std::path::{Path, PathBuf};
use super::heal_select::{claim_key, Claim};
const CLAIMS_FILE: &str = "heal-claims.json";
const PRUNE_AFTER_MS: u64 = 7 * 24 * 60 * 60 * 1000;
pub const BACKOFF_BASE_MS: u64 = 60 * 60 * 1000;
pub const MAX_ATTEMPTS: u32 = 5;
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct Attempts {
pub count: u32,
pub last_failed_ms: u64,
pub last_reason: String,
}
impl Attempts {
pub fn next_eligible_ms(&self) -> Option<u64> {
if self.count >= MAX_ATTEMPTS {
return None;
}
let shift = self.count.saturating_sub(1).min(10);
Some(self.last_failed_ms.saturating_add(BACKOFF_BASE_MS << shift))
}
pub fn ready(&self, now_ms: u64) -> bool {
match self.next_eligible_ms() {
None => false,
Some(at) => now_ms >= at,
}
}
}
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct ClaimStore {
claims: HashMap<String, Claim>,
attempts: HashMap<String, Attempts>,
}
impl ClaimStore {
pub fn new() -> Self {
Self::default()
}
pub fn load(dir: &Path) -> Self {
let path = dir.join(CLAIMS_FILE);
let Ok(raw) = std::fs::read_to_string(&path) else {
return Self::new();
};
match serde_json::from_str::<StoredLedger>(&raw) {
Ok(l) => Self {
claims: l
.claims
.into_iter()
.map(|(k, v)| {
(
k,
Claim {
run_id: v.run_id,
claimed_ms: v.claimed_ms,
},
)
})
.collect(),
attempts: l.attempts,
},
Err(e) => {
tracing::warn!(
path = %path.display(),
error = %e,
"unreadable heal-claim ledger; continuing with an empty one"
);
Self::new()
}
}
}
pub fn save(&self, dir: &Path, now_ms: u64) -> Result<(), String> {
std::fs::create_dir_all(dir).map_err(|e| format!("create {dir:?}: {e}"))?;
let keep: HashMap<&String, StoredClaim> = self
.claims
.iter()
.filter(|(_, c)| now_ms.saturating_sub(c.claimed_ms) < PRUNE_AFTER_MS)
.map(|(k, c)| {
(
k,
StoredClaim {
run_id: c.run_id.clone(),
claimed_ms: c.claimed_ms,
},
)
})
.collect();
let attempts: HashMap<&String, &Attempts> = self
.attempts
.iter()
.filter(|(_, a)| now_ms.saturating_sub(a.last_failed_ms) < PRUNE_AFTER_MS)
.collect();
let json = serde_json::to_string_pretty(&StoredLedgerRef {
claims: keep,
attempts,
})
.map_err(|e| format!("serialize heal claims: {e}"))?;
let path = dir.join(CLAIMS_FILE);
let tmp = tmp_path(&path);
std::fs::write(&tmp, json).map_err(|e| format!("write {tmp:?}: {e}"))?;
std::fs::rename(&tmp, &path).map_err(|e| format!("replace {path:?}: {e}"))?;
Ok(())
}
pub fn as_map(&self) -> &HashMap<String, Claim> {
&self.claims
}
pub fn attempts(&self) -> &HashMap<String, Attempts> {
&self.attempts
}
pub fn record_failure(&mut self, repo: &str, number: u64, reason: &str, now_ms: u64) {
let e = self.attempts.entry(claim_key(repo, number)).or_default();
e.count = e.count.saturating_add(1);
e.last_failed_ms = now_ms;
e.last_reason = reason.to_string();
}
pub fn record_permanent_failure(&mut self, repo: &str, number: u64, reason: &str, now_ms: u64) {
let e = self.attempts.entry(claim_key(repo, number)).or_default();
e.count = MAX_ATTEMPTS;
e.last_failed_ms = now_ms;
e.last_reason = reason.to_string();
}
pub fn clear_failures(&mut self, repo: &str, number: u64) {
self.attempts.remove(&claim_key(repo, number));
}
pub fn claim(
&mut self,
repo: &str,
number: u64,
run_id: &str,
now_ms: u64,
) -> Result<(), ClaimRefused> {
let key = claim_key(repo, number);
if let Some(existing) = self.claims.get(&key) {
let live =
now_ms.saturating_sub(existing.claimed_ms) < super::heal_select::CLAIM_TTL_MS;
if live && existing.run_id != run_id {
return Err(ClaimRefused {
held_by: existing.run_id.clone(),
});
}
}
self.claims.insert(
key,
Claim {
run_id: run_id.to_string(),
claimed_ms: now_ms,
},
);
Ok(())
}
pub fn release(&mut self, repo: &str, number: u64, run_id: &str) {
let key = claim_key(repo, number);
if self.claims.get(&key).is_some_and(|c| c.run_id == run_id) {
self.claims.remove(&key);
}
}
pub fn held_by(&self, repo: &str, number: u64) -> Option<&str> {
self.claims
.get(&claim_key(repo, number))
.map(|c| c.run_id.as_str())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ClaimRefused {
pub held_by: String,
}
impl std::fmt::Display for ClaimRefused {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "already claimed by run `{}`", self.held_by)
}
}
#[derive(serde::Serialize, serde::Deserialize)]
struct StoredClaim {
run_id: String,
claimed_ms: u64,
}
#[derive(serde::Deserialize, Default)]
struct StoredLedger {
#[serde(default)]
claims: HashMap<String, StoredClaim>,
#[serde(default)]
attempts: HashMap<String, Attempts>,
}
#[derive(serde::Serialize)]
struct StoredLedgerRef<'a> {
claims: HashMap<&'a String, StoredClaim>,
attempts: HashMap<&'a String, &'a Attempts>,
}
fn tmp_path(path: &Path) -> PathBuf {
let mut name = path.file_name().unwrap_or_default().to_os_string();
name.push(".tmp");
path.with_file_name(name)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::coder::heal_select::CLAIM_TTL_MS;
#[test]
fn claiming_then_reading_back_survives_a_round_trip() {
let dir = tempfile::tempdir().unwrap();
let mut s = ClaimStore::new();
s.claim("acme/widgets", 3, "run-1", 1_000).unwrap();
s.save(dir.path(), 1_000).unwrap();
let loaded = ClaimStore::load(dir.path());
assert_eq!(loaded.held_by("acme/widgets", 3), Some("run-1"));
}
#[test]
fn a_missing_ledger_is_an_empty_one_not_an_error() {
let dir = tempfile::tempdir().unwrap();
assert_eq!(ClaimStore::load(dir.path()), ClaimStore::new());
}
#[test]
fn a_corrupt_ledger_does_not_stop_the_loop() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join(CLAIMS_FILE), "{not json").unwrap();
assert_eq!(ClaimStore::load(dir.path()), ClaimStore::new());
}
#[test]
fn a_second_run_cannot_take_a_live_claim() {
let mut s = ClaimStore::new();
s.claim("acme/widgets", 3, "run-1", 1_000).unwrap();
let err = s
.claim("acme/widgets", 3, "run-2", 1_000 + CLAIM_TTL_MS - 1)
.unwrap_err();
assert_eq!(err.held_by, "run-1");
}
#[test]
fn a_second_run_may_take_an_expired_claim() {
let mut s = ClaimStore::new();
s.claim("acme/widgets", 3, "run-1", 1_000).unwrap();
s.claim("acme/widgets", 3, "run-2", 1_000 + CLAIM_TTL_MS)
.expect("expired claims do not hold");
assert_eq!(s.held_by("acme/widgets", 3), Some("run-2"));
}
#[test]
fn re_claiming_your_own_work_is_idempotent() {
let mut s = ClaimStore::new();
s.claim("acme/widgets", 3, "run-1", 1_000).unwrap();
s.claim("acme/widgets", 3, "run-1", 2_000)
.expect("same run re-claims");
assert_eq!(s.as_map()[&claim_key("acme/widgets", 3)].claimed_ms, 2_000);
}
#[test]
fn only_the_holder_may_release() {
let mut s = ClaimStore::new();
s.claim("acme/widgets", 3, "run-1", 1_000).unwrap();
s.release("acme/widgets", 3, "run-2");
assert_eq!(
s.held_by("acme/widgets", 3),
Some("run-1"),
"a stranger's release is a no-op"
);
s.release("acme/widgets", 3, "run-1");
assert_eq!(s.held_by("acme/widgets", 3), None);
}
#[test]
fn long_dead_claims_are_pruned_on_save() {
let dir = tempfile::tempdir().unwrap();
let mut s = ClaimStore::new();
s.claim("acme/widgets", 1, "old-run", 0).unwrap();
s.claim("acme/widgets", 2, "new-run", PRUNE_AFTER_MS)
.unwrap();
s.save(dir.path(), PRUNE_AFTER_MS).unwrap();
let loaded = ClaimStore::load(dir.path());
assert_eq!(loaded.held_by("acme/widgets", 1), None, "pruned");
assert_eq!(loaded.held_by("acme/widgets", 2), Some("new-run"));
}
#[test]
fn an_expired_claim_is_kept_until_the_prune_horizon() {
let dir = tempfile::tempdir().unwrap();
let mut s = ClaimStore::new();
s.claim("acme/widgets", 1, "run-1", 0).unwrap();
s.save(dir.path(), CLAIM_TTL_MS + 1).unwrap();
assert_eq!(
ClaimStore::load(dir.path()).held_by("acme/widgets", 1),
Some("run-1")
);
}
#[test]
fn the_store_feeds_selection_directly() {
let mut s = ClaimStore::new();
s.claim("acme/widgets", 9, "run-1", 500).unwrap();
let map = s.as_map();
assert!(map.contains_key(&claim_key("acme/widgets", 9)));
}
#[test]
fn a_partial_write_cannot_be_observed() {
let dir = tempfile::tempdir().unwrap();
let mut s = ClaimStore::new();
s.claim("acme/widgets", 1, "run-1", 0).unwrap();
s.save(dir.path(), 0).unwrap();
let entries: Vec<_> = std::fs::read_dir(dir.path())
.unwrap()
.filter_map(|e| e.ok())
.map(|e| e.file_name().to_string_lossy().to_string())
.collect();
assert!(
!entries.iter().any(|n| n.ends_with(".tmp")),
"temp file left behind: {entries:?}"
);
}
#[test]
fn a_failed_item_is_not_retried_immediately() {
let mut s = ClaimStore::new();
s.record_failure("acme/w", 1, "panel rejected", 1_000);
let a = &s.attempts()[&claim_key("acme/w", 1)];
assert_eq!(a.count, 1);
assert!(!a.ready(1_000), "not immediately");
assert!(a.ready(1_000 + BACKOFF_BASE_MS), "but eventually");
}
#[test]
fn backoff_doubles_with_each_failure() {
let mut s = ClaimStore::new();
s.record_failure("acme/w", 1, "x", 0);
s.record_failure("acme/w", 1, "x", 0);
let a = &s.attempts()[&claim_key("acme/w", 1)];
assert_eq!(a.next_eligible_ms(), Some(BACKOFF_BASE_MS * 2));
}
#[test]
fn an_item_that_keeps_failing_is_eventually_left_to_a_human() {
let mut s = ClaimStore::new();
for _ in 0..MAX_ATTEMPTS {
s.record_failure("acme/w", 1, "x", 0);
}
let a = &s.attempts()[&claim_key("acme/w", 1)];
assert_eq!(a.next_eligible_ms(), None);
assert!(!a.ready(u64::MAX), "never ready again");
}
#[test]
fn success_clears_the_history() {
let mut s = ClaimStore::new();
s.record_failure("acme/w", 1, "x", 0);
s.clear_failures("acme/w", 1);
assert!(s.attempts().is_empty());
}
#[test]
fn attempt_history_survives_a_round_trip() {
let dir = tempfile::tempdir().unwrap();
let mut s = ClaimStore::new();
s.record_failure("acme/w", 1, "panel rejected", 1_000);
s.claim("acme/w", 2, "run-1", 1_000).unwrap();
s.save(dir.path(), 1_000).unwrap();
let loaded = ClaimStore::load(dir.path());
assert_eq!(loaded.attempts()[&claim_key("acme/w", 1)].count, 1);
assert_eq!(loaded.held_by("acme/w", 2), Some("run-1"));
}
#[test]
fn a_backoff_shift_cannot_overflow() {
let a = Attempts {
count: 3,
last_failed_ms: u64::MAX - 10,
last_reason: String::new(),
};
assert!(a.next_eligible_ms().is_some());
assert!(!a.ready(0));
}
}