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