Skip to main content

unifier/
tick.rs

1//! ACID tick staging: read previous tick at start, write current tick in memory, commit at end.
2
3use std::collections::{BTreeMap, BTreeSet};
4use std::fs;
5
6use chrono::Utc;
7use serde::Serialize;
8use uuid::Uuid;
9
10use crate::constants::{EVENTS, KEYS, TICKS, TICK_CURRENT};
11use crate::error::{Error, Result};
12use crate::fs_text::{read_text, write_text, write_text_atomic};
13use crate::home::UnifierHome;
14use crate::store::{validate_key, ActiveTick, EventState, HotStore, KeyState, StagingValue};
15
16#[derive(Debug, Clone, PartialEq, Eq)]
17pub enum TickStartOutcome {
18    Started { tick: u64 },
19    Queued { position: usize, label: String },
20}
21
22#[derive(Debug, Clone, PartialEq, Eq)]
23pub struct TickStatus {
24    pub committed_tick: u64,
25    pub active_tick: Option<u64>,
26    pub queued: usize,
27    pub locked_keys: Vec<String>,
28}
29
30#[derive(Serialize)]
31struct TickMeta {
32    tick: u64,
33    committed_at: String,
34    keys_written: usize,
35}
36
37impl HotStore {
38    pub fn tick_start(&mut self, label: &str) -> Result<TickStartOutcome> {
39        if self.active_tick.is_some() {
40            self.tick_queue.push_back(label.to_string());
41            return Ok(TickStartOutcome::Queued {
42                position: self.tick_queue.len(),
43                label: label.to_string(),
44            });
45        }
46        let tick = self.begin_tick()?;
47        Ok(TickStartOutcome::Started { tick })
48    }
49
50    pub fn tick_end(&mut self, home: &UnifierHome) -> Result<u64> {
51        let Some(active) = self.active_tick.take() else {
52            return Err(Error::msg("no active tick"));
53        };
54
55        for (key, staged) in &active.staging {
56            match staged {
57                StagingValue::Present(value) => {
58                    self.keys.insert(
59                        key.clone(),
60                        KeyState::Present {
61                            value: value.clone(),
62                            dirty: true,
63                        },
64                    );
65                }
66                StagingValue::Deleted => {
67                    if self.keys.contains_key(key) {
68                        self.keys
69                            .insert(key.clone(), KeyState::Deleted { dirty: true });
70                    }
71                }
72            }
73        }
74
75        self.committed_tick = active.number;
76        commit_tick_version(home, active.number, &active.staging)?;
77        write_committed_tick(home, active.number)?;
78
79        if let Some(next) = self.tick_queue.pop_front() {
80            eprintln!(
81                "tick queue: starting queued tick {:?} ({} remaining)",
82                next,
83                self.tick_queue.len()
84            );
85            self.begin_tick()?;
86        }
87
88        Ok(active.number)
89    }
90
91    pub fn tick_status(&self) -> TickStatus {
92        TickStatus {
93            committed_tick: self.committed_tick,
94            active_tick: self.active_tick.as_ref().map(|t| t.number),
95            queued: self.tick_queue.len(),
96            locked_keys: self
97                .active_tick
98                .as_ref()
99                .map(|t| t.locks.iter().cloned().collect())
100                .unwrap_or_default(),
101        }
102    }
103
104    pub fn tick_lock(&mut self, key: &str) -> Result<()> {
105        let Some(tick) = &mut self.active_tick else {
106            return Err(Error::msg("no active tick"));
107        };
108        validate_key(key)?;
109        tick.locks.insert(key.to_string());
110        Ok(())
111    }
112
113    pub fn tick_unlock(&mut self, key: &str) -> Result<bool> {
114        let Some(tick) = &mut self.active_tick else {
115            return Err(Error::msg("no active tick"));
116        };
117        Ok(tick.locks.remove(key))
118    }
119
120    pub fn post_event(&mut self, payload: &str) -> Result<Uuid> {
121        serde_json::from_str::<serde_json::Value>(payload)
122            .map_err(|e| Error::msg(format!("event payload must be valid JSON: {e}")))?;
123        let id = Uuid::new_v4();
124        self.events.insert(
125            id,
126            EventState {
127                body: payload.to_string(),
128                dirty: true,
129            },
130        );
131        Ok(id)
132    }
133
134    pub fn send_agent_message(&mut self, from: &str, to: &str, payload: &str) -> Result<Uuid> {
135        serde_json::from_str::<serde_json::Value>(payload)
136            .map_err(|e| Error::msg(format!("message payload must be valid JSON: {e}")))?;
137        self.send_from(from, to, payload)
138    }
139
140    pub(crate) fn flush_events(&mut self, home: &UnifierHome) -> Result<()> {
141        let root = home.path().join(EVENTS);
142        for (id, event) in self.events.iter_mut() {
143            if event.dirty {
144                if let Some(p) = root.parent() {
145                    let _ = fs::create_dir_all(p);
146                }
147                fs::create_dir_all(&root)?;
148                write_text(&root.join(format!("{}.json", id.hyphenated())), &event.body)?;
149                event.dirty = false;
150            }
151        }
152        Ok(())
153    }
154
155    fn begin_tick(&mut self) -> Result<u64> {
156        let number = self.committed_tick.saturating_add(1);
157        let read_snapshot = self.committed_key_snapshot();
158        self.active_tick = Some(ActiveTick {
159            number,
160            read_snapshot,
161            staging: BTreeMap::new(),
162            locks: BTreeSet::new(),
163        });
164        Ok(number)
165    }
166
167    fn committed_key_snapshot(&self) -> BTreeMap<String, String> {
168        self.keys
169            .iter()
170            .filter_map(|(k, state)| match state {
171                KeyState::Present { value, .. } => Some((k.clone(), value.clone())),
172                KeyState::Deleted { .. } => None,
173            })
174            .collect()
175    }
176}
177
178pub(crate) fn load_committed_tick(home: &UnifierHome) -> Result<u64> {
179    let path = home.path().join(TICKS).join(TICK_CURRENT);
180    if !path.is_file() {
181        return Ok(0);
182    }
183    let text = read_text(&path)?;
184    text.parse::<u64>()
185        .map_err(|_| Error::msg("invalid ticks/current"))
186}
187
188pub(crate) fn load_events(home: &UnifierHome) -> Result<BTreeMap<Uuid, EventState>> {
189    let mut events = BTreeMap::new();
190    let root = home.path().join(EVENTS);
191    if !root.is_dir() {
192        return Ok(events);
193    }
194    for entry in fs::read_dir(root)? {
195        let entry = entry?;
196        if !entry.file_type()?.is_file() {
197            continue;
198        }
199        let name = entry.file_name().to_string_lossy().into_owned();
200        let Some(stem) = name.strip_suffix(".json") else {
201            continue;
202        };
203        let Ok(id) = Uuid::parse_str(stem) else {
204            continue;
205        };
206        events.insert(
207            id,
208            EventState {
209                body: read_text(&entry.path())?,
210                dirty: false,
211            },
212        );
213    }
214    Ok(events)
215}
216
217fn write_committed_tick(home: &UnifierHome, tick: u64) -> Result<()> {
218    let dir = home.path().join(TICKS);
219    fs::create_dir_all(&dir)?;
220    write_text_atomic(&dir.join(TICK_CURRENT), &tick.to_string())
221}
222
223fn commit_tick_version(
224    home: &UnifierHome,
225    tick: u64,
226    staging: &BTreeMap<String, StagingValue>,
227) -> Result<()> {
228    let tick_root = home.path().join(TICKS).join(tick.to_string());
229    let tick_keys = tick_root.join(KEYS);
230    fs::create_dir_all(&tick_keys)?;
231
232    for (key, staged) in staging {
233        match staged {
234            StagingValue::Present(value) => {
235                let dest = tick_keys.join(key);
236                if let Some(p) = dest.parent() {
237                    fs::create_dir_all(p)?;
238                }
239                write_text_atomic(&dest, value)?;
240            }
241            StagingValue::Deleted => {}
242        }
243    }
244
245    let meta = TickMeta {
246        tick,
247        committed_at: Utc::now().to_rfc3339(),
248        keys_written: staging.len(),
249    };
250    write_text(
251        &tick_root.join("meta.json"),
252        &serde_json::to_string_pretty(&meta)?,
253    )?;
254    Ok(())
255}
256
257#[cfg(test)]
258mod tests {
259    use super::*;
260    use crate::home::UnifierHome;
261    use tempfile::tempdir;
262
263    #[test]
264    fn tick_reads_previous_committed_state_only() {
265        let tmp = tempdir().unwrap();
266        let home = UnifierHome::resolve(Some(tmp.path().to_path_buf()), None).unwrap();
267        let mut store = HotStore::load(&home).unwrap();
268
269        store.put_key("counter", "1").unwrap();
270        store.tick_start("t1").unwrap();
271        assert_eq!(store.get_key("counter").unwrap(), Some("1".into()));
272        store.put_key("counter", "2").unwrap();
273        assert_eq!(store.get_key("counter").unwrap(), Some("2".into()));
274
275        store.tick_end(&home).unwrap();
276        assert_eq!(store.get_key("counter").unwrap(), Some("2".into()));
277        assert!(home.path().join("ticks/1/meta.json").is_file());
278    }
279
280    #[test]
281    fn locked_key_rejects_put() {
282        let tmp = tempdir().unwrap();
283        let home = UnifierHome::resolve(Some(tmp.path().to_path_buf()), None).unwrap();
284        let mut store = HotStore::load(&home).unwrap();
285        store.tick_start("t1").unwrap();
286        store.tick_lock("config/x").unwrap();
287        assert!(store.put_key("config/x", "y").is_err());
288    }
289}