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::{
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}