1use 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}