use anyhow::{Context, Result, anyhow};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use std::collections::{HashMap, HashSet, VecDeque};
use std::fs;
use std::io::Read;
use std::path::{Path, PathBuf};
use crate::store::atomic_write::atomic_write;
pub mod migration;
const DEFAULT_MAX_HASHES: usize = 50_000;
const MAX_STATE_JSON_BYTES: usize = 128 * 1024 * 1024;
fn read_state_json_validated(path: &Path) -> Result<String> {
let validated = crate::sanitize::validate_read_path(path)?;
let file = fs::OpenOptions::new()
.read(true)
.open(&validated)
.map_err(|e| anyhow!("Failed to open '{}': {}", validated.display(), e))?;
let metadata = file.metadata().map_err(|e| {
anyhow!(
"Failed to stat opened fd for '{}': {}",
validated.display(),
e
)
})?;
if metadata.len() > MAX_STATE_JSON_BYTES as u64 {
return Err(anyhow!(
"State file '{}' exceeds {} bytes (actual: {})",
validated.display(),
MAX_STATE_JSON_BYTES,
metadata.len()
));
}
let mut reader = file.take(MAX_STATE_JSON_BYTES.saturating_add(1) as u64);
let mut bytes = Vec::new();
reader
.read_to_end(&mut bytes)
.map_err(|e| anyhow!("Failed to read '{}': {}", validated.display(), e))?;
if bytes.len() > MAX_STATE_JSON_BYTES {
return Err(anyhow!(
"State file '{}' exceeds {} bytes (actual: {})",
validated.display(),
MAX_STATE_JSON_BYTES,
bytes.len()
));
}
String::from_utf8(bytes).map_err(|e| anyhow!("Failed to read '{}': {}", validated.display(), e))
}
#[derive(Debug, Clone, Default)]
pub struct SeenHashSet {
order: VecDeque<String>,
set: HashSet<String>,
}
impl SeenHashSet {
pub fn len(&self) -> usize {
self.set.len()
}
pub fn is_empty(&self) -> bool {
self.set.is_empty()
}
pub fn contains(&self, hash: &str) -> bool {
self.set.contains(hash)
}
pub fn insert(&mut self, hash: String) {
if self.set.remove(&hash) {
self.order.retain(|existing| *existing != hash);
}
self.set.insert(hash.clone());
self.order.push_back(hash);
}
pub(crate) fn extend_from(&mut self, other: SeenHashSet) {
for hash in other.order {
self.insert(hash);
}
}
pub fn prune_oldest(&mut self, limit: usize) {
while self.set.len() > limit {
let Some(hash) = self.order.pop_front() else {
break;
};
self.set.remove(&hash);
}
}
}
impl Serialize for SeenHashSet {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
self.order.serialize(serializer)
}
}
impl<'de> Deserialize<'de> for SeenHashSet {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let hashes = Vec::<String>::deserialize(deserializer)?;
let mut out = SeenHashSet::default();
for hash in hashes {
out.insert(hash);
}
Ok(out)
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RunRecord {
pub timestamp: DateTime<Utc>,
pub entries_added: usize,
pub sources: Vec<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StateManager {
pub last_processed: HashMap<String, DateTime<Utc>>,
pub seen_hashes: HashMap<String, SeenHashSet>,
pub runs: Vec<RunRecord>,
#[serde(default)]
pub hash_algorithm: String,
}
#[derive(Debug, Deserialize)]
struct LegacySiphashStateManager {
#[serde(default)]
last_processed: HashMap<String, DateTime<Utc>>,
#[serde(default)]
seen_hashes: HashMap<String, Vec<u64>>,
#[serde(default)]
runs: Vec<RunRecord>,
hash_algorithm: String,
}
impl LegacySiphashStateManager {
fn into_current_state(self) -> StateManager {
let seen_hashes = self
.seen_hashes
.into_iter()
.map(|(project, hashes)| {
let mut set = SeenHashSet::default();
for hash in hashes {
set.insert(hash.to_string());
}
(project, set)
})
.collect();
StateManager {
last_processed: self.last_processed,
seen_hashes,
runs: self.runs,
hash_algorithm: self.hash_algorithm,
}
}
}
impl Default for StateManager {
fn default() -> Self {
Self {
last_processed: HashMap::new(),
seen_hashes: HashMap::new(),
runs: Vec::new(),
hash_algorithm: migration::BLAKE3_128_ALGORITHM.to_string(),
}
}
}
pub(crate) fn state_path_for(home: &Path) -> PathBuf {
home.join("state.json")
}
impl StateManager {
fn state_path() -> Result<PathBuf> {
Ok(state_path_for(&crate::store::store_base_dir()?))
}
pub fn load() -> Result<Self> {
Self::load_from_path(&Self::state_path()?)
}
fn load_from_path(path: &Path) -> Result<Self> {
Self::load_from_path_with_legacy_warning(path, |message| eprintln!("{message}"))
}
fn load_from_path_with_legacy_warning<W>(path: &Path, mut warn_legacy: W) -> Result<Self>
where
W: FnMut(&str),
{
if !path.exists() {
return Ok(Self::default());
}
let backup_path = Self::backup_path(path);
let contents = read_state_json_validated(path)
.with_context(|| format!("Failed to read state file: {}", path.display()))?;
let state: Self =
match Self::deserialize_and_migrate_contents(path, &contents, &mut warn_legacy) {
Ok(state) => state,
Err(err) => {
tracing::warn!(
path = %path.display(),
error = %err,
"state.json parse failed"
);
if backup_path.exists() {
let backup =
read_state_json_validated(&backup_path).with_context(|| {
format!("Failed to read state backup: {}", backup_path.display())
})?;
let recovered = Self::deserialize_and_migrate_contents(
&backup_path,
&backup,
&mut warn_legacy,
)
.map_err(|backup_err| {
anyhow!(
"state.json malformed AND backup unreadable: {err} / {backup_err}"
)
})?;
recovered.save_recovered_backup_to_primary(path)?;
recovered
} else {
return Err(anyhow!(
"state.json corrupted, no backup; manual recovery needed: {}",
path.display()
));
}
}
};
Ok(state)
}
fn save_recovered_backup_to_primary(&self, path: &Path) -> Result<()> {
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)
.with_context(|| format!("Failed to create config dir: {}", parent.display()))?;
}
let json = serde_json::to_string_pretty(self).context("Failed to serialize state")?;
atomic_write(path, json.as_bytes()).with_context(|| {
format!(
"Failed to self-heal state file from backup: {}",
path.display()
)
})?;
Ok(())
}
fn deserialize_and_migrate_contents<W>(
path: &Path,
contents: &str,
warn_legacy: &mut W,
) -> std::result::Result<Self, serde_json::Error>
where
W: FnMut(&str),
{
let (mut state, loaded_legacy_u64_shape) =
Self::deserialize_current_or_legacy_siphash(contents)?;
let previous_hash_algorithm = state.hash_algorithm.clone();
let report = state.apply_load_migrations();
Self::emit_load_migration_warning(
path,
&previous_hash_algorithm,
loaded_legacy_u64_shape,
&report,
warn_legacy,
);
Ok(state)
}
fn deserialize_current_or_legacy_siphash(
contents: &str,
) -> std::result::Result<(Self, bool), serde_json::Error> {
match serde_json::from_str::<Self>(contents) {
Ok(state) => Ok((state, false)),
Err(strict_err) => {
let value: serde_json::Value =
serde_json::from_str(contents).map_err(|_| strict_err)?;
if !Self::is_legacy_siphash_state_value(&value) {
return Err(serde_json::from_str::<Self>(contents).unwrap_err());
}
serde_json::from_value::<LegacySiphashStateManager>(value)
.map(|legacy| (legacy.into_current_state(), true))
.map_err(|_| serde_json::from_str::<Self>(contents).unwrap_err())
}
}
}
fn is_legacy_siphash_state_value(value: &serde_json::Value) -> bool {
value
.get("hash_algorithm")
.and_then(|algorithm| algorithm.as_str())
.is_some_and(migration::is_legacy_siphash13_algorithm)
}
fn emit_load_migration_warning<W>(
path: &Path,
previous_hash_algorithm: &str,
loaded_legacy_u64_shape: bool,
report: &migration::StateMigrationReport,
warn_legacy: &mut W,
) where
W: FnMut(&str),
{
if !report.hash_algorithm_changed {
return;
}
tracing::warn!(
path = %path.display(),
previous_hash_algorithm = %previous_hash_algorithm,
current_hash_algorithm = %migration::BLAKE3_128_ALGORITHM,
cleared_seen_hashes = report.cleared_seen_hashes,
legacy_u64_shape = loaded_legacy_u64_shape,
"state.json migrated from legacy hash algorithm"
);
let previous = if previous_hash_algorithm.trim().is_empty() {
"missing"
} else {
previous_hash_algorithm.trim()
};
warn_legacy(&format!(
"Warning: state.json migrated from legacy hash algorithm {previous} to {}; cleared {} legacy seen_hashes",
migration::BLAKE3_128_ALGORITHM,
report.cleared_seen_hashes
));
}
pub fn save(&self) -> Result<()> {
let path = Self::state_path()?;
self.save_to_path(&path)
}
fn save_to_path(&self, path: &Path) -> Result<()> {
self.save_to_path_with_writer(path, atomic_write)
}
fn save_to_path_with_writer<W>(&self, path: &Path, write_atomic: W) -> Result<()>
where
W: Fn(&Path, &[u8]) -> std::io::Result<()>,
{
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)
.with_context(|| format!("Failed to create config dir: {}", parent.display()))?;
}
let json = serde_json::to_string_pretty(self).context("Failed to serialize state")?;
match fs::OpenOptions::new().read(true).open(path) {
Ok(mut previous_file) => {
let mut previous = Vec::new();
previous_file
.read_to_end(&mut previous)
.with_context(|| format!("Failed to read state file: {}", path.display()))?;
let backup_path = Self::backup_path(path);
write_atomic(&backup_path, &previous).with_context(|| {
format!("Failed to write state backup: {}", backup_path.display())
})?;
}
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {}
Err(err) => {
return Err(err)
.with_context(|| format!("Failed to open state file: {}", path.display()));
}
}
write_atomic(path, json.as_bytes())
.with_context(|| format!("Failed to write state file: {}", path.display()))?;
Ok(())
}
fn backup_path(path: &Path) -> PathBuf {
path.with_file_name("state.json.bak")
}
pub fn content_hash(agent: &str, timestamp: i64, message: &str) -> String {
let timestamp = timestamp.to_le_bytes();
Self::stable_hash_fields(&[agent.as_bytes(), ×tamp, message.as_bytes()])
}
pub fn overlap_hash(timestamp: i64, message: &str) -> String {
let bucket = timestamp / 60; let bucket = bucket.to_le_bytes();
Self::stable_hash_fields(&[&bucket, message.as_bytes()])
}
fn stable_hash_fields(fields: &[&[u8]]) -> String {
let mut hasher = blake3::Hasher::new();
for field in fields {
hasher.update(&(field.len() as u64).to_le_bytes());
hasher.update(field);
}
hex::encode(&hasher.finalize().as_bytes()[..16])
}
pub fn is_new(&self, project: &str, hash: &str) -> bool {
let project = migration::canonical_state_bucket(project);
self.seen_hashes
.get(&project)
.is_none_or(|set| !set.contains(hash))
}
pub fn mark_seen(&mut self, project: &str, hash: String) {
let project = migration::canonical_state_bucket(project);
self.seen_hashes.entry(project).or_default().insert(hash);
}
pub fn get_watermark(&self, source: &str) -> Option<DateTime<Utc>> {
self.last_processed.get(source).copied()
}
pub fn migrate_watermark_aliases(&mut self, canonical: &str, aliases: &[String]) -> bool {
if self.last_processed.contains_key(canonical) {
return false;
}
let migrated = aliases
.iter()
.filter_map(|alias| self.last_processed.get(alias).copied())
.max();
if let Some(ts) = migrated {
self.last_processed.insert(canonical.to_string(), ts);
true
} else {
false
}
}
pub fn update_watermark(&mut self, source: &str, ts: DateTime<Utc>) {
let entry = self.last_processed.entry(source.to_string()).or_insert(ts);
if ts > *entry {
*entry = ts;
}
}
pub fn record_run(&mut self, entries: usize, sources: Vec<String>) {
self.runs.push(RunRecord {
timestamp: Utc::now(),
entries_added: entries,
sources,
});
}
pub fn prune_old_hashes(&mut self, max_per_project: usize) {
let limit = if max_per_project == 0 {
DEFAULT_MAX_HASHES
} else {
max_per_project
};
for set in self.seen_hashes.values_mut() {
set.prune_oldest(limit);
}
}
pub fn reset_project(&mut self, project: &str) {
self.seen_hashes
.remove(&migration::canonical_state_bucket(project));
}
pub fn reset_all(&mut self) {
self.seen_hashes.clear();
}
pub fn total_hashes(&self) -> usize {
self.seen_hashes.values().map(|s| s.len()).sum()
}
fn apply_load_migrations(&mut self) -> migration::StateMigrationReport {
migration::migrate_loaded_state(self)
}
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::TimeZone;
use std::path::{Path, PathBuf};
use std::time::{SystemTime, UNIX_EPOCH};
fn unique_state_path(label: &str) -> PathBuf {
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
let dir =
std::env::temp_dir().join(format!("aicx-state-{label}-{}-{nanos}", std::process::id()));
fs::create_dir_all(&dir).unwrap();
dir.join("state.json")
}
fn cleanup_state_path(path: &Path) {
if let Some(parent) = path.parent() {
let _ = fs::remove_dir_all(parent);
}
}
fn ts(seconds: i64) -> DateTime<Utc> {
Utc.timestamp_opt(seconds, 0).single().unwrap()
}
fn state_with_marker(source: &str, watermark_seconds: i64, project_hash: u64) -> StateManager {
let mut state = StateManager::default();
state.update_watermark(source, ts(watermark_seconds));
state.mark_seen("project", project_hash.to_string());
state
}
#[test]
fn test_default_state_is_empty() {
let state = StateManager::default();
assert!(state.last_processed.is_empty());
assert!(state.seen_hashes.is_empty());
assert!(state.runs.is_empty());
}
#[test]
fn test_content_hash_deterministic() {
let h1 = StateManager::content_hash("claude", 1700000000, "hello world");
let h2 = StateManager::content_hash("claude", 1700000000, "hello world");
assert_eq!(h1, h2);
}
#[test]
fn test_content_hash_varies_with_input() {
let h1 = StateManager::content_hash("claude", 1700000000, "hello");
let h2 = StateManager::content_hash("claude", 1700000000, "world");
assert_ne!(h1, h2, "different message → different hash");
let h3 = StateManager::content_hash("codex", 1700000000, "hello");
assert_ne!(h1, h3, "different agent → different hash");
let h5 = StateManager::content_hash("claude", 1700000001, "hello");
assert_ne!(h1, h5, "different timestamp → different hash");
}
#[test]
fn test_content_hash_separates_legacy_concat_collision() {
fn legacy_content_hash(agent: &str, timestamp: i64, message: &str) -> String {
let mut data = Vec::new();
data.extend_from_slice(agent.as_bytes());
data.extend_from_slice(×tamp.to_le_bytes());
data.extend_from_slice(message.as_bytes());
migration::stable_blake3_128(&data)
}
let left = ("a\0\0\0\0\0\0\0\0", 0, "bc");
let right = ("a", 0, "\0\0\0\0\0\0\0\0bc");
assert_eq!(
legacy_content_hash(left.0, left.1, left.2),
legacy_content_hash(right.0, right.1, right.2),
"legacy raw concatenation should demonstrate the split collision"
);
assert_ne!(
StateManager::content_hash(left.0, left.1, left.2),
StateManager::content_hash(right.0, right.1, right.2),
"field boundaries should split the triples"
);
}
#[test]
fn test_overlap_hash_ignores_agent() {
let prompt = "Deploy the new auth module to staging and run integration tests";
let ts = 1700000000i64;
let h_claude = StateManager::overlap_hash(ts, prompt);
let h_codex = StateManager::overlap_hash(ts, prompt);
assert_eq!(
h_claude, h_codex,
"same message + same bucket → SAME overlap hash"
);
}
#[test]
fn test_overlap_hash_buckets_60s() {
let prompt = "Identical broadcast prompt";
let base = 1700000040i64; let same_bucket = base + 19;
let h1 = StateManager::overlap_hash(base, prompt);
let h2 = StateManager::overlap_hash(same_bucket, prompt);
assert_eq!(h1, h2, "within same 60s bucket → SAME hash");
let next_bucket = base - (base % 60) + 60; let h3 = StateManager::overlap_hash(next_bucket, prompt);
assert_ne!(h1, h3, "different 60s bucket → different hash");
}
#[test]
fn test_overlap_hash_different_message() {
let ts = 1700000000i64;
let h1 = StateManager::overlap_hash(ts, "prompt A");
let h2 = StateManager::overlap_hash(ts, "prompt B");
assert_ne!(h1, h2, "different message → different overlap hash");
}
#[test]
fn test_overlap_hash_uses_field_boundary_between_bucket_and_message() {
let timestamp = 0i64;
let message = "broadcast prompt";
let mut legacy_data = Vec::new();
legacy_data.extend_from_slice(&(timestamp / 60).to_le_bytes());
legacy_data.extend_from_slice(message.as_bytes());
assert_ne!(
StateManager::overlap_hash(timestamp, message),
migration::stable_blake3_128(&legacy_data),
"overlap hash should move off the legacy raw-concat byte stream"
);
}
#[test]
fn test_is_new_and_mark_seen_per_project() {
let mut state = StateManager::default();
let hash = StateManager::content_hash("claude", 100, "msg");
assert!(state.is_new("projA", &hash));
assert!(state.is_new("projB", &hash));
state.mark_seen("projA", hash.clone());
assert!(!state.is_new("projA", &hash));
assert!(!state.is_new("proja", &hash));
assert!(state.is_new("projB", &hash));
state.mark_seen("projB", hash.clone());
assert!(!state.is_new("projB", &hash));
}
#[test]
fn test_watermark_none_for_unknown_source() {
let state = StateManager::default();
assert_eq!(state.get_watermark("nonexistent"), None);
}
#[test]
fn test_watermark_update_only_if_newer() {
let mut state = StateManager::default();
let t1 = Utc.with_ymd_and_hms(2026, 1, 1, 10, 0, 0).unwrap();
let t2 = Utc.with_ymd_and_hms(2026, 1, 1, 12, 0, 0).unwrap();
let t0 = Utc.with_ymd_and_hms(2026, 1, 1, 8, 0, 0).unwrap();
state.update_watermark("claude:CodeScribe", t1);
assert_eq!(state.get_watermark("claude:CodeScribe"), Some(t1));
state.update_watermark("claude:CodeScribe", t2);
assert_eq!(state.get_watermark("claude:CodeScribe"), Some(t2));
state.update_watermark("claude:CodeScribe", t0);
assert_eq!(state.get_watermark("claude:CodeScribe"), Some(t2));
}
#[test]
fn test_record_run() {
let mut state = StateManager::default();
assert!(state.runs.is_empty());
state.record_run(
42,
vec!["claude:Proj".to_string(), "codex:global".to_string()],
);
assert_eq!(state.runs.len(), 1);
assert_eq!(state.runs[0].entries_added, 42);
assert_eq!(state.runs[0].sources, vec!["claude:Proj", "codex:global"]);
}
#[test]
fn test_save_uses_atomic_write() {
let path = unique_state_path("atomic-write");
let old = state_with_marker("claude:test", 10, 101);
old.save_to_path(&path).unwrap();
let old_contents = fs::read(&path).unwrap();
let new = state_with_marker("claude:test", 20, 202);
let target_path = path.clone();
let err = new
.save_to_path_with_writer(&path, move |target, content| {
if target == target_path.as_path() {
return Err(std::io::Error::other("mock atomic_write failure"));
}
crate::store::atomic_write::atomic_write(target, content)
})
.expect_err("mocked final atomic write should fail");
assert!(err.to_string().contains("Failed to write state file"));
assert_eq!(fs::read(&path).unwrap(), old_contents);
let loaded = StateManager::load_from_path(&path).unwrap();
assert_eq!(loaded.get_watermark("claude:test"), Some(ts(10)));
assert!(!loaded.is_new("project", "101"));
cleanup_state_path(&path);
}
#[test]
fn test_load_malformed_returns_error_not_default() {
let path = unique_state_path("malformed");
fs::write(&path, b"{ this is not json").unwrap();
let err = StateManager::load_from_path(&path)
.expect_err("malformed state without backup must not default");
assert!(err.to_string().contains("state.json corrupted"));
cleanup_state_path(&path);
}
#[test]
fn test_load_recovers_from_backup_when_main_corrupt() {
let path = unique_state_path("backup-recovery");
let backup_path = StateManager::backup_path(&path);
let backup_state = state_with_marker("claude:test", 20, 202);
fs::write(&path, b"{ this is not json").unwrap();
fs::write(
&backup_path,
serde_json::to_string_pretty(&backup_state).unwrap(),
)
.unwrap();
let loaded = StateManager::load_from_path(&path).unwrap();
assert_eq!(loaded.get_watermark("claude:test"), Some(ts(20)));
assert!(!loaded.is_new("project", "202"));
let repaired_primary: StateManager =
serde_json::from_str(&fs::read_to_string(&path).unwrap()).unwrap();
assert_eq!(repaired_primary.get_watermark("claude:test"), Some(ts(20)));
assert!(!repaired_primary.is_new("project", "202"));
let repaired = StateManager::load_from_path(&path).unwrap();
assert_eq!(repaired.get_watermark("claude:test"), Some(ts(20)));
assert!(!repaired.is_new("project", "202"));
cleanup_state_path(&path);
}
#[test]
fn test_load_migrates_legacy_siphash_u64_state() {
let path = unique_state_path("legacy-siphash-u64");
let legacy_state = serde_json::json!({
"last_processed": {
"claude:test": "2026-05-20T12:00:00Z"
},
"seen_hashes": {
"Vista": [101_u64, 202_u64]
},
"runs": [
{
"timestamp": "2026-05-20T12:30:00Z",
"entries_added": 2,
"sources": ["claude:test"]
}
],
"hash_algorithm": migration::SIPHASH13_ALGORITHM
});
fs::write(&path, serde_json::to_vec_pretty(&legacy_state).unwrap()).unwrap();
let mut warnings = Vec::new();
let loaded = StateManager::load_from_path_with_legacy_warning(&path, |message| {
warnings.push(message.to_string());
})
.unwrap();
assert_eq!(loaded.hash_algorithm, migration::BLAKE3_128_ALGORITHM);
assert!(loaded.seen_hashes.is_empty());
assert_eq!(loaded.total_hashes(), 0);
assert_eq!(loaded.get_watermark("claude:test"), Some(ts(1779278400)));
assert_eq!(loaded.runs.len(), 1);
assert_eq!(loaded.runs[0].entries_added, 2);
assert_eq!(warnings.len(), 1);
assert!(warnings[0].contains("migrated from legacy hash algorithm"));
assert!(warnings[0].contains(migration::SIPHASH13_ALGORITHM));
assert!(warnings[0].contains(migration::BLAKE3_128_ALGORITHM));
assert!(warnings[0].contains("cleared 2 legacy seen_hashes"));
loaded.save_to_path(&path).unwrap();
let persisted: serde_json::Value =
serde_json::from_str(&fs::read_to_string(&path).unwrap()).unwrap();
assert_eq!(
persisted["hash_algorithm"],
serde_json::Value::String(migration::BLAKE3_128_ALGORITHM.to_string())
);
assert_eq!(persisted["seen_hashes"].as_object().unwrap().len(), 0);
cleanup_state_path(&path);
}
#[test]
fn test_load_rejects_non_legacy_schema_mismatch() {
let path = unique_state_path("non-legacy-schema-mismatch");
let invalid_current_state = serde_json::json!({
"last_processed": {},
"seen_hashes": {
"Vista": [101_u64]
},
"runs": [],
"hash_algorithm": migration::BLAKE3_128_ALGORITHM
});
fs::write(
&path,
serde_json::to_vec_pretty(&invalid_current_state).unwrap(),
)
.unwrap();
let mut warnings = Vec::new();
let err = StateManager::load_from_path_with_legacy_warning(&path, |message| {
warnings.push(message.to_string());
})
.expect_err("non-legacy u64 hashes must still be rejected");
assert!(err.to_string().contains("state.json corrupted"));
assert!(warnings.is_empty());
cleanup_state_path(&path);
}
#[test]
fn test_save_creates_backup_before_overwrite() {
let path = unique_state_path("backup-overwrite");
let backup_path = StateManager::backup_path(&path);
let old = state_with_marker("claude:test", 10, 101);
old.save_to_path(&path).unwrap();
let old_contents = fs::read_to_string(&path).unwrap();
let new = state_with_marker("claude:test", 20, 202);
new.save_to_path(&path).unwrap();
assert_eq!(fs::read_to_string(&backup_path).unwrap(), old_contents);
let loaded = StateManager::load_from_path(&path).unwrap();
assert_eq!(loaded.get_watermark("claude:test"), Some(ts(20)));
assert!(!loaded.is_new("project", "202"));
cleanup_state_path(&path);
}
#[test]
fn test_prune_old_hashes_below_limit() {
let mut state = StateManager::default();
for i in 0..10u64 {
state.mark_seen("proj", i.to_string());
}
state.prune_old_hashes(100);
assert_eq!(state.seen_hashes["proj"].len(), 10);
}
#[test]
fn test_prune_old_hashes_above_limit() {
let mut state = StateManager::default();
for i in 0..100u64 {
state.mark_seen("proj", i.to_string());
}
state.prune_old_hashes(30);
assert_eq!(state.seen_hashes["proj"].len(), 30);
}
#[test]
fn lru_evicts_oldest_first() {
let mut state = StateManager::default();
for i in 0..10u64 {
state.mark_seen("proj", i.to_string());
}
state.prune_old_hashes(5);
for old in 0..5u64 {
assert!(
state.is_new("proj", &old.to_string()),
"old hash {old} should be evicted"
);
}
for fresh in 5..10u64 {
assert!(
!state.is_new("proj", &fresh.to_string()),
"fresh hash {fresh} should remain"
);
}
}
#[test]
fn watermark_migration_carries_timestamp_forward() {
let mut state = StateManager::default();
let ts = Utc.with_ymd_and_hms(2026, 5, 6, 11, 0, 0).unwrap();
state.update_watermark("claude+codex+gemini:all", ts);
let migrated = state.migrate_watermark_aliases(
"claude+codex+gemini+junie:all",
&["claude+codex+gemini:all".to_string()],
);
assert!(migrated);
assert_eq!(
state.get_watermark("claude+codex+gemini+junie:all"),
Some(ts)
);
}
#[test]
fn test_prune_old_hashes_default_limit() {
let mut state = StateManager::default();
state.prune_old_hashes(0);
assert_eq!(state.total_hashes(), 0);
}
#[test]
fn test_reset_project() {
let mut state = StateManager::default();
state.mark_seen("projA", "1".to_string());
state.mark_seen("projA", "2".to_string());
state.mark_seen("projB", "3".to_string());
state.reset_project("projA");
assert!(state.is_new("projA", "1"));
assert!(!state.is_new("projB", "3"));
}
#[test]
fn test_reset_all() {
let mut state = StateManager::default();
state.mark_seen("projA", "1".to_string());
state.mark_seen("projB", "2".to_string());
state.reset_all();
assert!(state.is_new("projA", "1"));
assert!(state.is_new("projB", "2"));
assert_eq!(state.total_hashes(), 0);
}
#[test]
fn test_serialization_roundtrip() {
let mut state = StateManager::default();
let t = Utc.with_ymd_and_hms(2026, 1, 20, 15, 30, 0).unwrap();
state.update_watermark("claude:TestProject", t);
state.mark_seen("myproj", "123456789".to_string());
state.mark_seen("myproj", "987654321".to_string());
state.record_run(5, vec!["claude:TestProject".to_string()]);
let json = serde_json::to_string_pretty(&state).unwrap();
let restored: StateManager = serde_json::from_str(&json).unwrap();
assert_eq!(restored.get_watermark("claude:TestProject"), Some(t));
assert!(!restored.is_new("myproj", "123456789"));
assert!(!restored.is_new("myproj", "987654321"));
assert!(restored.is_new("myproj", "111111111"));
assert!(restored.is_new("other", "123456789")); assert_eq!(restored.runs.len(), 1);
assert_eq!(restored.runs[0].entries_added, 5);
}
#[test]
fn pre_siphash_state_clears_seen_hashes_once() {
let mut state = StateManager::default();
state.hash_algorithm.clear();
state
.seen_hashes
.entry("Vista".to_string())
.or_default()
.insert("42".to_string());
let report = migration::migrate_loaded_state(&mut state);
assert!(report.hash_algorithm_changed);
assert_eq!(report.cleared_seen_hashes, 1);
assert_eq!(state.hash_algorithm, migration::BLAKE3_128_ALGORITHM);
assert_eq!(state.total_hashes(), 0);
}
#[test]
fn current_state_lowercases_and_merges_seen_hash_buckets() {
let mut state = StateManager::default();
state
.seen_hashes
.entry("Vista".to_string())
.or_default()
.insert("1".to_string());
state
.seen_hashes
.entry("vista".to_string())
.or_default()
.insert("2".to_string());
let report = migration::migrate_loaded_state(&mut state);
assert!(!report.hash_algorithm_changed);
assert_eq!(report.lowercased_seen_hash_buckets, 1);
assert_eq!(report.merged_seen_hash_buckets, 1);
assert_eq!(state.seen_hashes.len(), 1);
assert!(!state.is_new("vista", "1"));
assert!(!state.is_new("Vista", "2"));
}
#[test]
fn case_bucket_merge_script_is_reviewable_and_lowercase_targeted() {
let root = std::env::temp_dir().join(format!(
"aicx-case-bucket-script-{}",
Utc::now().timestamp_nanos_opt().unwrap()
));
let _ = std::fs::remove_dir_all(&root);
std::fs::create_dir_all(root.join("VetCoders").join("Vista")).unwrap();
std::fs::create_dir_all(root.join("vetcoders").join("vista")).unwrap();
let plan = migration::generate_case_bucket_merge_script(&root).unwrap();
assert_eq!(plan.merges.len(), 2);
assert!(plan.script.contains("Review before running"));
assert!(plan.script.contains("vetcoders"));
assert!(plan.script.contains("VetCoders"));
assert!(root.join("VetCoders").join("Vista").exists());
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn test_state_path_for_lives_under_home_and_named_state_json() {
let home = PathBuf::from("/tmp/test-aicx-state");
let state = super::state_path_for(&home);
assert!(
state.starts_with(&home),
"state_path_for should live under home; got {state:?}"
);
assert_eq!(
state.file_name().and_then(|n| n.to_str()),
Some("state.json"),
"state_path_for should end with state.json; got {state:?}"
);
}
}