unifier-cli 0.2.0

Filesystem postbox for inter-process communication via a Unix tree
Documentation
//! ACID tick staging: read previous tick at start, write current tick in memory, commit at end.

use std::collections::{BTreeMap, BTreeSet};
use std::fs;

use chrono::Utc;
use serde::Serialize;
use uuid::Uuid;

use crate::constants::{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>,
}

#[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()?;
        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()?;
        }

        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(),
        }
    }

    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 post_event(&mut self, payload: &str) -> Result<Uuid> {
        serde_json::from_str::<serde_json::Value>(payload)
            .map_err(|e| Error::msg(format!("event payload must be valid JSON: {e}")))?;
        let id = Uuid::new_v4();
        self.events.insert(
            id,
            EventState {
                body: payload.to_string(),
                dirty: true,
            },
        );
        Ok(id)
    }

    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<()> {
        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)?;
                write_text(
                    &root.join(format!("{}.json", id.hyphenated())),
                    &event.body,
                )?;
                event.dirty = false;
            }
        }
        Ok(())
    }

    fn begin_tick(&mut self) -> Result<u64> {
        let number = self.committed_tick.saturating_add(1);
        let read_snapshot = self.committed_key_snapshot();
        self.active_tick = Some(ActiveTick {
            number,
            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;
        };
        events.insert(
            id,
            EventState {
                body: read_text(&entry.path())?,
                dirty: false,
            },
        );
    }
    Ok(events)
}

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());
    }
}