use std::collections::{BTreeMap, BTreeSet};
use std::path::Path;
use serde::{Deserialize, Serialize};
use crate::Engine;
use crate::binding::BuildMode;
use crate::pipeline_store::BindingConfigs;
use super::cursor::source_moved;
use super::resolve::{ResolvedIngest, resolve_binding_run};
pub const MAX_SKIP_LEVEL: u32 = 10;
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct BackoffEntry {
#[serde(default)]
pub skip_remaining: u32,
#[serde(default)]
pub skip_level: u32,
#[serde(default)]
pub snapshot: String,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct Cursor {
#[serde(default)]
pub last: Option<String>,
}
pub fn apply_backoff(entry: &mut BackoffEntry, current: &str) -> bool {
if !entry.snapshot.is_empty() && current != entry.snapshot {
entry.skip_remaining = 0;
entry.skip_level = 0;
entry.snapshot = current.to_string();
return false;
}
if entry.skip_remaining > 0 {
entry.skip_remaining -= 1;
return true;
}
if !entry.snapshot.is_empty() && current == entry.snapshot {
entry.skip_level = (entry.skip_level + 1).min(MAX_SKIP_LEVEL);
entry.skip_remaining = entry.skip_level;
}
entry.snapshot = current.to_string();
false
}
pub fn should_skip(
mode: BuildMode,
source_moved: bool,
entry: &mut BackoffEntry,
current: &str,
) -> bool {
match mode {
BuildMode::OneShot => return false,
BuildMode::Discovery => {}
}
if source_moved {
return false;
}
apply_backoff(entry, current)
}
fn read_json<T: Default + for<'de> Deserialize<'de>>(cache_root: &Path, name: &str) -> T {
std::fs::read(cache_root.join(name))
.ok()
.and_then(|b| serde_json::from_slice(&b).ok())
.unwrap_or_default()
}
fn write_json<T: Serialize>(cache_root: &Path, name: &str, value: &T) {
let _ = std::fs::create_dir_all(cache_root);
if let Ok(bytes) = serde_json::to_vec(value) {
let _ = std::fs::write(cache_root.join(name), bytes);
}
}
fn read_one_shot_runs(cache_root: &Path) -> BTreeSet<String> {
let map: BTreeMap<String, bool> = read_json(cache_root, "ingest-one-shot-runs.json");
map.into_iter()
.filter(|(_, v)| *v)
.map(|(k, _)| k)
.collect()
}
pub fn select_next_due(
engine: &Engine,
workspace_root: &Path,
configs: &BindingConfigs,
) -> Option<String> {
let cache_root = workspace_root.join(".memstead.cache").join("ingest");
let one_shot_ran = read_one_shot_runs(&cache_root);
let mut eligible: Vec<ResolvedIngest> = configs
.bindings
.iter()
.filter_map(|r| {
resolve_binding_run(configs, &format!("{}/{}", r.mem, r.name), &r.config).ok()
})
.filter(|ri| !(ri.mode == BuildMode::OneShot && one_shot_ran.contains(&ri.name)))
.collect();
eligible.sort_by(|a, b| a.name.cmp(&b.name));
let n = eligible.len();
if n == 0 {
return None;
}
let mut cursor: Cursor = read_json(&cache_root, "ingest-cursor.json");
let start = cursor
.last
.as_ref()
.and_then(|last| eligible.iter().position(|ri| &ri.name == last))
.map_or(0, |i| (i + 1) % n);
cursor.last = Some(eligible[start].name.clone());
write_json(&cache_root, "ingest-cursor.json", &cursor);
let mut backoff: BTreeMap<String, BackoffEntry> = read_json(&cache_root, "ingest-backoff.json");
let mut selected = None;
for offset in 0..n {
let ingest = &eligible[(start + offset) % n];
let current = engine
.mem_head_sha(&ingest.destination_mem)
.ok()
.flatten()
.unwrap_or_default();
let moved = source_moved(engine, ingest, workspace_root);
let entry = backoff.entry(ingest.name.clone()).or_default();
if !should_skip(ingest.mode, moved, entry, ¤t) {
selected = Some(ingest.name.clone());
break;
}
}
write_json(&cache_root, "ingest-backoff.json", &backoff);
selected
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn backoff_ramps_and_resets() {
let mut e = BackoffEntry::default();
assert!(!apply_backoff(&mut e, "sha1"));
assert_eq!(e.snapshot, "sha1");
assert_eq!(e.skip_level, 0);
assert!(!apply_backoff(&mut e, "sha1"));
assert_eq!(e.skip_level, 1);
assert_eq!(e.skip_remaining, 1);
assert!(apply_backoff(&mut e, "sha1"));
assert_eq!(e.skip_remaining, 0);
assert!(!apply_backoff(&mut e, "sha1"));
assert_eq!(e.skip_level, 2);
assert_eq!(e.skip_remaining, 2);
assert!(!apply_backoff(&mut e, "sha2"));
assert_eq!(e.skip_level, 0);
assert_eq!(e.skip_remaining, 0);
assert_eq!(e.snapshot, "sha2");
}
#[test]
fn backoff_caps_at_max_level() {
let mut e = BackoffEntry {
skip_level: MAX_SKIP_LEVEL,
skip_remaining: 0,
snapshot: "s".to_string(),
};
assert!(!apply_backoff(&mut e, "s")); assert_eq!(e.skip_level, MAX_SKIP_LEVEL, "capped");
assert_eq!(e.skip_remaining, MAX_SKIP_LEVEL);
}
#[test]
fn should_skip_honours_mode_and_source_movement() {
let mut e = BackoffEntry {
skip_remaining: 3,
skip_level: 3,
snapshot: "s".to_string(),
};
assert!(!should_skip(BuildMode::OneShot, false, &mut e.clone(), "s"));
let mut e2 = e.clone();
assert!(!should_skip(BuildMode::Discovery, true, &mut e2, "s"));
assert_eq!(e2.skip_remaining, 3, "moved source does not touch backoff");
assert!(should_skip(BuildMode::Discovery, false, &mut e, "s"));
}
}