somatize_runtime/cache/
local.rs1use chrono::Utc;
5use somatize_core::cache::{CacheKey, CacheStore, EntryMeta, Origin};
6use somatize_core::error::{Result, SomaError};
7use somatize_core::value::Value;
8use std::fs;
9use std::path::{Path, PathBuf};
10
11pub struct LocalCache {
16 base_dir: PathBuf,
17}
18
19impl LocalCache {
20 pub fn new(base_dir: impl Into<PathBuf>) -> Result<Self> {
22 let base_dir = base_dir.into();
23 fs::create_dir_all(&base_dir)?;
24 Ok(Self { base_dir })
25 }
26
27 fn key_path(&self, key: &CacheKey) -> PathBuf {
28 let hex = key.to_hex();
29 self.base_dir
31 .join(&hex[..2])
32 .join(&hex[2..4])
33 .join(format!("{hex}.json"))
34 }
35
36 fn meta_path(&self, key: &CacheKey) -> PathBuf {
37 let hex = key.to_hex();
38 self.base_dir
39 .join(&hex[..2])
40 .join(&hex[2..4])
41 .join(format!("{hex}.meta.json"))
42 }
43
44 pub fn len(&self) -> usize {
46 walkdir_count(&self.base_dir)
47 }
48
49 pub fn is_empty(&self) -> bool {
51 self.len() == 0
52 }
53
54 pub fn clear(&self) -> Result<()> {
56 if self.base_dir.exists() {
57 fs::remove_dir_all(&self.base_dir)?;
58 fs::create_dir_all(&self.base_dir)?;
59 }
60 Ok(())
61 }
62
63 fn write_atomic(path: &Path, data: &[u8]) -> Result<()> {
70 use std::io::Write;
71 use std::sync::atomic::{AtomicU64, Ordering};
72 static WRITE_SEQ: AtomicU64 = AtomicU64::new(0);
73
74 let parent = path
75 .parent()
76 .ok_or_else(|| SomaError::Cache("cache path has no parent".into()))?;
77 fs::create_dir_all(parent)?;
78 let seq = WRITE_SEQ.fetch_add(1, Ordering::Relaxed);
79 let tmp = path.with_extension(format!("tmp-{}-{seq}", std::process::id()));
80 {
81 let mut f = fs::File::create(&tmp)?;
82 f.write_all(data)?;
83 f.sync_all()?;
84 }
85 if let Err(e) = fs::rename(&tmp, path) {
86 let _ = fs::remove_file(&tmp);
87 return Err(e.into());
88 }
89 Ok(())
90 }
91
92 fn write_entry(&self, key: &CacheKey, value: &Value, origin: &Origin) -> Result<()> {
93 let data = serde_json::to_string(value)
94 .map_err(|e| SomaError::Cache(format!("serialize error: {e}")))?;
95 let size = data.len() as u64;
96
97 let meta = EntryMeta {
100 key: key.clone(),
101 size_bytes: size,
102 created_at: Utc::now(),
103 last_accessed: Utc::now(),
104 ttl: None,
105 origin: origin.clone(),
106 };
107 let meta_data = serde_json::to_string(&meta)
108 .map_err(|e| SomaError::Cache(format!("meta serialize error: {e}")))?;
109 Self::write_atomic(&self.meta_path(key), meta_data.as_bytes())?;
110 Self::write_atomic(&self.key_path(key), data.as_bytes())?;
111 Ok(())
112 }
113}
114
115impl CacheStore for LocalCache {
116 fn tier(&self) -> somatize_core::cache::CacheTier {
117 somatize_core::cache::CacheTier::Local
118 }
119
120 fn get(&self, key: &CacheKey) -> Result<Option<Value>> {
121 let path = self.key_path(key);
122 if !path.exists() {
123 return Ok(None);
124 }
125 let data = fs::read_to_string(&path)?;
126 let value: Value = serde_json::from_str(&data)
127 .map_err(|e| SomaError::Cache(format!("deserialize error: {e}")))?;
128 Ok(Some(value))
129 }
130
131 fn put(&self, key: &CacheKey, value: &Value) -> Result<()> {
132 self.write_entry(
133 key,
134 value,
135 &Origin::Ingested {
136 source: "unknown".into(),
137 },
138 )
139 }
140
141 fn put_with_origin(&self, key: &CacheKey, value: &Value, origin: &Origin) -> Result<()> {
142 self.write_entry(key, value, origin)
143 }
144
145 fn exists(&self, key: &CacheKey) -> Result<bool> {
146 Ok(self.key_path(key).exists())
147 }
148
149 fn remove(&self, key: &CacheKey) -> Result<()> {
150 let path = self.key_path(key);
151 if path.exists() {
152 fs::remove_file(&path)?;
153 }
154 let meta = self.meta_path(key);
155 if meta.exists() {
156 fs::remove_file(&meta)?;
157 }
158 Ok(())
159 }
160
161 fn metadata(&self, key: &CacheKey) -> Result<Option<EntryMeta>> {
162 let path = self.meta_path(key);
163 if !path.exists() {
164 return Ok(None);
165 }
166 let data = fs::read_to_string(&path)?;
167 let meta: EntryMeta = serde_json::from_str(&data)
168 .map_err(|e| SomaError::Cache(format!("meta deserialize error: {e}")))?;
169 Ok(Some(meta))
170 }
171}
172
173fn walkdir_count(dir: &Path) -> usize {
174 if !dir.exists() {
175 return 0;
176 }
177 fs::read_dir(dir)
178 .map(|entries| {
179 entries
180 .filter_map(|e| e.ok())
181 .map(|e| {
182 if e.path().is_dir() {
183 walkdir_count(&e.path())
184 } else if e.path().extension().is_some_and(|ext| ext == "json")
185 && !e.path().to_string_lossy().contains(".meta.")
186 {
187 1
188 } else {
189 0
190 }
191 })
192 .sum()
193 })
194 .unwrap_or(0)
195}
196
197#[cfg(test)]
198mod tests {
199 use super::*;
200 use serde_json::json;
201 use std::env;
202
203 use std::sync::atomic::{AtomicU64, Ordering};
204 static COUNTER: AtomicU64 = AtomicU64::new(0);
205
206 fn temp_dir() -> PathBuf {
207 let id = COUNTER.fetch_add(1, Ordering::Relaxed);
208 let dir = env::temp_dir().join(format!("soma_test_cache_{}_{id}", std::process::id()));
209 let _ = fs::remove_dir_all(&dir);
210 dir
211 }
212
213 #[test]
214 fn put_and_get() {
215 let dir = temp_dir();
216 let cache = LocalCache::new(&dir).unwrap();
217 let key = CacheKey::hash_data(b"test");
218 let value = Value::tensor(vec![1.0, 2.0, 3.0], vec![3]);
219
220 cache.put(&key, &value).unwrap();
221 let retrieved = cache.get(&key).unwrap().unwrap();
222 assert_eq!(retrieved, value);
223
224 cache.clear().unwrap();
225 }
226
227 #[test]
228 fn get_missing() {
229 let dir = temp_dir();
230 let cache = LocalCache::new(&dir).unwrap();
231 assert!(cache.get(&CacheKey::hash_data(b"nope")).unwrap().is_none());
232 cache.clear().unwrap();
233 }
234
235 #[test]
236 fn exists_check() {
237 let dir = temp_dir();
238 let cache = LocalCache::new(&dir).unwrap();
239 let key = CacheKey::hash_data(b"test");
240 assert!(!cache.exists(&key).unwrap());
241 cache.put(&key, &Value::Empty).unwrap();
242 assert!(cache.exists(&key).unwrap());
243 cache.clear().unwrap();
244 }
245
246 #[test]
247 fn remove_entry() {
248 let dir = temp_dir();
249 let cache = LocalCache::new(&dir).unwrap();
250 let key = CacheKey::hash_data(b"test");
251 cache.put(&key, &Value::json(json!(42))).unwrap();
252 assert!(cache.exists(&key).unwrap());
253 cache.remove(&key).unwrap();
254 assert!(!cache.exists(&key).unwrap());
255 cache.clear().unwrap();
256 }
257
258 #[test]
259 fn metadata_persists() {
260 let dir = temp_dir();
261 let cache = LocalCache::new(&dir).unwrap();
262 let key = CacheKey::hash_data(b"test");
263 cache
264 .put(&key, &Value::tensor(vec![1.0; 50], vec![50]))
265 .unwrap();
266
267 let meta = cache.metadata(&key).unwrap().unwrap();
268 assert!(meta.size_bytes > 0);
269 cache.clear().unwrap();
270 }
271
272 #[test]
273 fn orphan_tmp_files_are_not_entries() {
274 let dir = temp_dir();
275 let cache = LocalCache::new(&dir).unwrap();
276 let key = CacheKey::hash_data(b"torn");
277
278 let entry_path = cache.key_path(&key);
280 fs::create_dir_all(entry_path.parent().unwrap()).unwrap();
281 fs::write(entry_path.with_extension("tmp-999-0"), b"garbage{{{").unwrap();
282
283 assert!(!cache.exists(&key).unwrap());
284 assert!(cache.get(&key).unwrap().is_none());
285
286 cache.put(&key, &Value::tensor(vec![1.0], vec![1])).unwrap();
288 assert!(cache.exists(&key).unwrap());
289 cache.clear().unwrap();
290 }
291
292 #[test]
293 fn concurrent_same_key_puts_are_safe() {
294 let dir = temp_dir();
295 let cache = std::sync::Arc::new(LocalCache::new(&dir).unwrap());
296 let key = CacheKey::hash_data(b"contended");
297 let value = Value::tensor(vec![1.0; 256], vec![256]);
298
299 std::thread::scope(|s| {
300 for _ in 0..8 {
301 let cache = cache.clone();
302 let key = key.clone();
303 let value = value.clone();
304 s.spawn(move || cache.put(&key, &value).unwrap());
305 }
306 });
307
308 assert_eq!(cache.get(&key).unwrap().unwrap(), value);
309 cache.clear().unwrap();
310 }
311
312 #[test]
313 fn put_with_origin_records_provenance() {
314 let dir = temp_dir();
315 let cache = LocalCache::new(&dir).unwrap();
316 let key = CacheKey::hash_data(b"prov");
317 let origin = Origin::Computed {
318 node_id: "scaler".into(),
319 run_id: "run_42".into(),
320 };
321 cache.put_with_origin(&key, &Value::Empty, &origin).unwrap();
322
323 let meta = cache.metadata(&key).unwrap().unwrap();
324 match meta.origin {
325 Origin::Computed { node_id, run_id } => {
326 assert_eq!(node_id, "scaler");
327 assert_eq!(run_id, "run_42");
328 }
329 other => panic!("expected Computed origin, got {other:?}"),
330 }
331 cache.clear().unwrap();
332 }
333
334 #[test]
335 fn survives_restart() {
336 let dir = temp_dir();
337 let key = CacheKey::hash_data(b"persist");
338 let value = Value::tensor(vec![42.0], vec![1]);
339
340 {
341 let cache = LocalCache::new(&dir).unwrap();
342 cache.put(&key, &value).unwrap();
343 }
344 {
346 let cache = LocalCache::new(&dir).unwrap();
347 let retrieved = cache.get(&key).unwrap().unwrap();
348 assert_eq!(retrieved, value);
349 }
350
351 let _ = fs::remove_dir_all(&dir);
352 }
353}