Skip to main content

somatize_runtime/cache/
local.rs

1//! [`LocalCache`] — filesystem-backed [`CacheStore`] with sharded
2//! directories and atomic writes.
3
4use 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
11/// Filesystem-based cache store.
12///
13/// Each entry is stored as a JSON file named by the cache key's hex.
14/// Suitable for persistent local caching across process restarts.
15pub struct LocalCache {
16    base_dir: PathBuf,
17}
18
19impl LocalCache {
20    /// Open (or create) a cache rooted at `base_dir`.
21    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        // Shard into subdirectories: first 2 chars / next 2 chars / full key
30        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    /// Number of cached entries (scans filesystem).
45    pub fn len(&self) -> usize {
46        walkdir_count(&self.base_dir)
47    }
48
49    /// Whether the cache holds no entries (scans the filesystem).
50    pub fn is_empty(&self) -> bool {
51        self.len() == 0
52    }
53
54    /// Remove all cached entries.
55    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    /// Write `data` to `path` atomically: temp file in the same directory
64    /// → fsync → rename. A crash mid-write leaves only an orphan temp
65    /// file (never read back: lookups go through the exact entry path),
66    /// so entry existence is the commit point. Renames are idempotent,
67    /// which makes concurrent same-key writers from multiple processes
68    /// safe on POSIX filesystems.
69    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        // Meta first, value last: the value file's existence commits the
98        // entry (`exists()`/`get()` check only the value path).
99        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        // Simulate a crash mid-write: a temp file exists, the entry doesn't.
279        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        // A subsequent put commits normally despite the orphan.
287        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        // "restart": create a new instance pointing to same dir
345        {
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}