use std::collections::{BTreeMap, BTreeSet};
use std::fs;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use chrono::Utc;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use uuid::Uuid;
use crate::constants::{DEFAULT_EVENT_TTL_SECS, EVENTS, KEYS, TICKS, TICK_CURRENT};
use crate::error::{Error, Result};
use crate::fs_text::{read_text, write_text, write_text_atomic};
use crate::home::UnifierHome;
use crate::store::{validate_key, ActiveTick, EventState, HotStore, KeyState, StagingValue};
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TickStartOutcome {
Started { tick: u64 },
Queued { position: usize, label: String },
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TickStatus {
pub committed_tick: u64,
pub active_tick: Option<u64>,
pub queued: usize,
pub locked_keys: Vec<String>,
pub phase: Option<String>,
pub label: Option<String>,
}
#[derive(Serialize)]
struct TickMeta {
tick: u64,
committed_at: String,
keys_written: usize,
}
impl HotStore {
pub fn tick_start(&mut self, label: &str) -> Result<TickStartOutcome> {
if self.active_tick.is_some() {
self.tick_queue.push_back(label.to_string());
return Ok(TickStartOutcome::Queued {
position: self.tick_queue.len(),
label: label.to_string(),
});
}
let tick = self.begin_tick(label)?;
Ok(TickStartOutcome::Started { tick })
}
pub fn tick_end(&mut self, home: &UnifierHome) -> Result<u64> {
let Some(active) = self.active_tick.take() else {
return Err(Error::msg("no active tick"));
};
for (key, staged) in &active.staging {
match staged {
StagingValue::Present(value) => {
self.keys.insert(
key.clone(),
KeyState::Present {
value: value.clone(),
dirty: true,
},
);
}
StagingValue::Deleted => {
if self.keys.contains_key(key) {
self.keys
.insert(key.clone(), KeyState::Deleted { dirty: true });
}
}
}
}
self.committed_tick = active.number;
commit_tick_version(home, active.number, &active.staging)?;
write_committed_tick(home, active.number)?;
if let Some(next) = self.tick_queue.pop_front() {
eprintln!(
"tick queue: starting queued tick {:?} ({} remaining)",
next,
self.tick_queue.len()
);
self.begin_tick(&next)?;
}
Ok(active.number)
}
pub fn tick_status(&self) -> TickStatus {
TickStatus {
committed_tick: self.committed_tick,
active_tick: self.active_tick.as_ref().map(|t| t.number),
queued: self.tick_queue.len(),
locked_keys: self
.active_tick
.as_ref()
.map(|t| t.locks.iter().cloned().collect())
.unwrap_or_default(),
phase: self.active_tick.as_ref().map(|t| t.phase.clone()),
label: self.active_tick.as_ref().map(|t| t.label.clone()),
}
}
pub fn tick_lock(&mut self, key: &str) -> Result<()> {
let Some(tick) = &mut self.active_tick else {
return Err(Error::msg("no active tick"));
};
validate_key(key)?;
tick.locks.insert(key.to_string());
Ok(())
}
pub fn tick_unlock(&mut self, key: &str) -> Result<bool> {
let Some(tick) = &mut self.active_tick else {
return Err(Error::msg("no active tick"));
};
Ok(tick.locks.remove(key))
}
pub fn tick_phase(&mut self, phase: &str) -> Result<(u64, String)> {
let phase = phase.trim();
if phase.is_empty() {
return Err(Error::msg("tick phase must be non-empty"));
}
if phase.contains(|c: char| c.is_whitespace()) {
return Err(Error::msg("tick phase must not contain whitespace"));
}
let Some(tick) = &mut self.active_tick else {
return Err(Error::msg("no active tick"));
};
tick.phase = phase.to_string();
Ok((tick.number, tick.label.clone()))
}
pub fn post_event(&mut self, payload: &str, ttl_secs: Option<u64>) -> Result<Uuid> {
let value: Value = serde_json::from_str(payload)
.map_err(|e| Error::msg(format!("event payload must be valid JSON: {e}")))?;
let created_at = Utc::now().to_rfc3339();
let expires_at = resolve_event_expiry(&value, ttl_secs)?;
let id = Uuid::new_v4();
self.events.insert(
id,
EventState {
body: payload.to_string(),
created_at,
expires_at,
dirty: true,
},
);
Ok(id)
}
pub fn expire_events(&mut self, home: &UnifierHome) -> Result<usize> {
let mut gone = Vec::new();
for (id, event) in &self.events {
if event_is_expired(event.expires_at.as_deref()) {
gone.push(*id);
}
}
let n = gone.len();
if n == 0 {
return Ok(0);
}
let root = home.path().join(EVENTS);
for id in gone {
self.events.remove(&id);
let path = root.join(format!("{}.json", id.hyphenated()));
if path.is_file() {
fs::remove_file(path)?;
}
}
Ok(n)
}
pub fn send_agent_message(&mut self, from: &str, to: &str, payload: &str) -> Result<Uuid> {
serde_json::from_str::<serde_json::Value>(payload)
.map_err(|e| Error::msg(format!("message payload must be valid JSON: {e}")))?;
self.send_from(from, to, payload)
}
pub(crate) fn flush_events(&mut self, home: &UnifierHome) -> Result<()> {
self.expire_events(home)?;
let root = home.path().join(EVENTS);
for (id, event) in self.events.iter_mut() {
if event.dirty {
if let Some(p) = root.parent() {
let _ = fs::create_dir_all(p);
}
fs::create_dir_all(&root)?;
let record = EventRecord {
payload: serde_json::from_str(&event.body)
.unwrap_or(Value::String(event.body.clone())),
created_at: event.created_at.clone(),
expires_at: event.expires_at.clone(),
};
write_text(
&root.join(format!("{}.json", id.hyphenated())),
&serde_json::to_string(&record)?,
)?;
event.dirty = false;
}
}
Ok(())
}
fn begin_tick(&mut self, label: &str) -> Result<u64> {
let number = self.committed_tick.saturating_add(1);
let read_snapshot = self.committed_key_snapshot();
self.active_tick = Some(ActiveTick {
number,
label: label.to_string(),
phase: "start".to_string(),
read_snapshot,
staging: BTreeMap::new(),
locks: BTreeSet::new(),
});
Ok(number)
}
fn committed_key_snapshot(&self) -> BTreeMap<String, String> {
self.keys
.iter()
.filter_map(|(k, state)| match state {
KeyState::Present { value, .. } => Some((k.clone(), value.clone())),
KeyState::Deleted { .. } => None,
})
.collect()
}
}
pub(crate) fn load_committed_tick(home: &UnifierHome) -> Result<u64> {
let path = home.path().join(TICKS).join(TICK_CURRENT);
if !path.is_file() {
return Ok(0);
}
let text = read_text(&path)?;
text.parse::<u64>()
.map_err(|_| Error::msg("invalid ticks/current"))
}
pub(crate) fn load_events(home: &UnifierHome) -> Result<BTreeMap<Uuid, EventState>> {
let mut events = BTreeMap::new();
let root = home.path().join(EVENTS);
if !root.is_dir() {
return Ok(events);
}
for entry in fs::read_dir(&root)? {
let entry = entry?;
if !entry.file_type()?.is_file() {
continue;
}
let name = entry.file_name().to_string_lossy().into_owned();
let Some(stem) = name.strip_suffix(".json") else {
continue;
};
let Ok(id) = Uuid::parse_str(stem) else {
continue;
};
let path = entry.path();
let text = read_text(&path)?;
let Some(state) = parse_event_file(&text, &path)? else {
let _ = fs::remove_file(&path);
continue;
};
if event_is_expired(state.expires_at.as_deref()) {
let _ = fs::remove_file(&path);
continue;
}
events.insert(id, state);
}
Ok(events)
}
#[derive(Serialize, Deserialize)]
struct EventRecord {
payload: Value,
created_at: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
expires_at: Option<String>,
}
fn parse_event_file(text: &str, path: &std::path::Path) -> Result<Option<EventState>> {
let value: Value = match serde_json::from_str(text) {
Ok(v) => v,
Err(_) => {
return Ok(Some(EventState {
body: text.to_string(),
created_at: Utc::now().to_rfc3339(),
expires_at: expires_from_mtime(path),
dirty: false,
}));
}
};
if is_event_record(&value) {
let record: EventRecord = serde_json::from_value(value)?;
return Ok(Some(EventState {
body: serde_json::to_string(&record.payload)?,
created_at: record.created_at,
expires_at: record.expires_at,
dirty: false,
}));
}
Ok(Some(EventState {
body: text.trim().to_string(),
created_at: Utc::now().to_rfc3339(),
expires_at: expires_from_mtime(path),
dirty: false,
}))
}
fn is_event_record(value: &Value) -> bool {
value.get("payload").is_some() && value.get("created_at").is_some()
}
pub fn default_event_ttl_secs() -> u64 {
std::env::var("UNIFIER_EVENT_TTL_SECS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(DEFAULT_EVENT_TTL_SECS)
}
fn resolve_event_expiry(payload: &Value, ttl_flag: Option<u64>) -> Result<Option<String>> {
if let Some(secs) = ttl_flag {
return Ok(ttl_to_expires_at(secs));
}
if let Some(exp) = payload.get("expires_at") {
if exp.is_null() {
return Ok(None);
}
if let Some(s) = exp.as_str() {
chrono::DateTime::parse_from_rfc3339(s)
.map_err(|e| Error::msg(format!("invalid expires_at: {e}")))?;
return Ok(Some(s.to_string()));
}
return Err(Error::msg("expires_at must be an RFC3339 string or null"));
}
if let Some(ttl) = payload.get("ttl") {
let secs = ttl
.as_u64()
.ok_or_else(|| Error::msg("ttl must be a non-negative integer"))?;
return Ok(ttl_to_expires_at(secs));
}
Ok(ttl_to_expires_at(default_event_ttl_secs()))
}
fn ttl_to_expires_at(secs: u64) -> Option<String> {
if secs == 0 {
return None;
}
Some(rfc3339_from_system(
SystemTime::now() + Duration::from_secs(secs),
))
}
fn expires_from_mtime(path: &std::path::Path) -> Option<String> {
let ttl = default_event_ttl_secs();
if ttl == 0 {
return None;
}
let modified = fs::metadata(path)
.ok()
.and_then(|m| m.modified().ok())
.unwrap_or_else(SystemTime::now);
Some(rfc3339_from_system(modified + Duration::from_secs(ttl)))
}
fn rfc3339_from_system(t: SystemTime) -> String {
let secs = t.duration_since(UNIX_EPOCH).unwrap_or_default().as_secs() as i64;
chrono::DateTime::from_timestamp(secs, 0)
.unwrap_or_else(Utc::now)
.to_rfc3339()
}
fn event_is_expired(expires_at: Option<&str>) -> bool {
let Some(exp) = expires_at else {
return false;
};
match chrono::DateTime::parse_from_rfc3339(exp) {
Ok(dt) => dt.with_timezone(&Utc) < Utc::now(),
Err(_) => false,
}
}
fn write_committed_tick(home: &UnifierHome, tick: u64) -> Result<()> {
let dir = home.path().join(TICKS);
fs::create_dir_all(&dir)?;
write_text_atomic(&dir.join(TICK_CURRENT), &tick.to_string())
}
fn commit_tick_version(
home: &UnifierHome,
tick: u64,
staging: &BTreeMap<String, StagingValue>,
) -> Result<()> {
let tick_root = home.path().join(TICKS).join(tick.to_string());
let tick_keys = tick_root.join(KEYS);
fs::create_dir_all(&tick_keys)?;
for (key, staged) in staging {
match staged {
StagingValue::Present(value) => {
let dest = tick_keys.join(key);
if let Some(p) = dest.parent() {
fs::create_dir_all(p)?;
}
write_text_atomic(&dest, value)?;
}
StagingValue::Deleted => {}
}
}
let meta = TickMeta {
tick,
committed_at: Utc::now().to_rfc3339(),
keys_written: staging.len(),
};
write_text(
&tick_root.join("meta.json"),
&serde_json::to_string_pretty(&meta)?,
)?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::home::UnifierHome;
use tempfile::tempdir;
#[test]
fn tick_reads_previous_committed_state_only() {
let tmp = tempdir().unwrap();
let home = UnifierHome::resolve(Some(tmp.path().to_path_buf()), None).unwrap();
let mut store = HotStore::load(&home).unwrap();
store.put_key("counter", "1").unwrap();
store.tick_start("t1").unwrap();
assert_eq!(store.get_key("counter").unwrap(), Some("1".into()));
store.put_key("counter", "2").unwrap();
assert_eq!(store.get_key("counter").unwrap(), Some("2".into()));
store.tick_end(&home).unwrap();
assert_eq!(store.get_key("counter").unwrap(), Some("2".into()));
assert!(home.path().join("ticks/1/meta.json").is_file());
}
#[test]
fn locked_key_rejects_put() {
let tmp = tempdir().unwrap();
let home = UnifierHome::resolve(Some(tmp.path().to_path_buf()), None).unwrap();
let mut store = HotStore::load(&home).unwrap();
store.tick_start("t1").unwrap();
store.tick_lock("config/x").unwrap();
assert!(store.put_key("config/x", "y").is_err());
}
#[test]
fn post_event_defaults_to_24h_expiry() {
let tmp = tempdir().unwrap();
let home = UnifierHome::resolve(Some(tmp.path().to_path_buf()), None).unwrap();
let mut store = HotStore::load(&home).unwrap();
let id = store.post_event(r#"{"name":"ping"}"#, None).unwrap();
let event = store.events.get(&id).unwrap();
let exp = event.expires_at.as_ref().unwrap();
let dt = chrono::DateTime::parse_from_rfc3339(exp).unwrap();
let delta = dt.with_timezone(&Utc) - Utc::now();
assert!(delta.num_seconds() > 23 * 3600);
assert!(delta.num_seconds() <= 24 * 3600);
}
#[test]
fn post_event_ttl_zero_never_expires() {
let tmp = tempdir().unwrap();
let home = UnifierHome::resolve(Some(tmp.path().to_path_buf()), None).unwrap();
let mut store = HotStore::load(&home).unwrap();
let id = store.post_event(r#"{"name":"keep"}"#, Some(0)).unwrap();
assert_eq!(store.events.get(&id).unwrap().expires_at, None);
}
#[test]
fn payload_expires_at_is_honored() {
let v = serde_json::json!({"name":"x","expires_at":"2099-01-01T00:00:00Z"});
let exp = resolve_event_expiry(&v, None).unwrap();
assert_eq!(exp.as_deref(), Some("2099-01-01T00:00:00Z"));
}
#[test]
fn expired_events_are_deleted_on_load() {
let tmp = tempdir().unwrap();
let home = UnifierHome::resolve(Some(tmp.path().to_path_buf()), None).unwrap();
let dir = home.path().join("events");
std::fs::create_dir_all(&dir).unwrap();
let id = Uuid::new_v4();
let record = EventRecord {
payload: serde_json::json!({"name":"stale"}),
created_at: "2000-01-01T00:00:00Z".into(),
expires_at: Some("2000-01-02T00:00:00Z".into()),
};
std::fs::write(
dir.join(format!("{}.json", id.hyphenated())),
serde_json::to_string(&record).unwrap(),
)
.unwrap();
let store = HotStore::load(&home).unwrap();
assert!(store.events.is_empty());
assert!(!dir.join(format!("{}.json", id.hyphenated())).exists());
}
#[test]
fn flush_drops_expired_events() {
let tmp = tempdir().unwrap();
let home = UnifierHome::resolve(Some(tmp.path().to_path_buf()), None).unwrap();
let mut store = HotStore::load(&home).unwrap();
let id = store.post_event(r#"{"name":"soon"}"#, Some(3600)).unwrap();
store.events.get_mut(&id).unwrap().expires_at = Some("2000-01-01T00:00:00Z".into());
store.flush(&home).unwrap();
assert!(store.events.is_empty());
}
#[test]
fn expire_events_drops_past_due() {
let tmp = tempdir().unwrap();
let home = UnifierHome::resolve(Some(tmp.path().to_path_buf()), None).unwrap();
let mut store = HotStore::load(&home).unwrap();
let id = store.post_event(r#"{"name":"soon"}"#, Some(3600)).unwrap();
store.flush(&home).unwrap();
store.events.get_mut(&id).unwrap().expires_at = Some("2000-01-01T00:00:00Z".into());
assert_eq!(store.expire_events(&home).unwrap(), 1);
assert!(store.events.is_empty());
assert!(!home
.path()
.join("events")
.join(format!("{}.json", id.hyphenated()))
.exists());
}
}