Skip to main content

mesofact_core/proxy/
source_gen.rs

1//! Source generation tokens for cache-key input 6.
2//!
3//! The proxy folds each read source's *current generation* into the Mode 2
4//! cache key so a backend bump is an automatic miss (no manual purge). Per
5//! §"Cache-key composition", generations come from:
6//!
7//! | source | token | refresh |
8//! |---|---|---|
9//! | `sqlite` (global) | file mtime | cached 1s |
10//! | `r2` | bucket/object `Last-Modified` | cached 5s |
11//! | `pg` / `rpc` | LSN / roster token | (post-MVP) |
12//!
13//! MVP implements the sqlite mtime path (the P9 slice's source). Other kinds
14//! return a stable placeholder until their poll lands — a stable token is safe
15//! (it just means generation never advances on its own; TTL still expires).
16//!
17//! The proxy caches each token with a 1s TTL so a burst of misses doesn't
18//! amplify into N filesystem stats / backend pings.
19
20use serde::Deserialize;
21use std::collections::HashMap;
22use std::path::Path;
23use std::sync::Mutex;
24use std::time::{Duration, Instant, UNIX_EPOCH};
25
26const GENERATION_TTL: Duration = Duration::from_secs(1);
27
28/// One declared source, parsed from `[sources.<name>]` in `mesofact.config.toml`.
29/// Only the fields the generation poll needs are kept.
30#[derive(Debug, Clone)]
31pub struct SourceDef {
32    pub kind: String,
33    pub path: Option<String>,
34}
35
36#[derive(Debug, Deserialize)]
37struct RawConfig {
38    #[serde(default)]
39    sources: HashMap<String, RawSource>,
40}
41
42#[derive(Debug, Deserialize)]
43struct RawSource {
44    kind: String,
45    #[serde(default)]
46    path: Option<String>,
47}
48
49/// Generation provider: maps source name → current generation token, with a
50/// 1s memo so repeated cache-key composition is cheap.
51pub struct Generations {
52    defs: HashMap<String, SourceDef>,
53    memo: Mutex<HashMap<String, (String, Instant)>>,
54}
55
56impl Generations {
57    /// Empty provider — every source resolves to the placeholder token. Used
58    /// when no `mesofact.config.toml` is configured.
59    pub fn empty() -> Self {
60        Self { defs: HashMap::new(), memo: Mutex::new(HashMap::new()) }
61    }
62
63    /// Parse `[sources.*]` from a `mesofact.config.toml` string.
64    pub fn from_config_str(toml_str: &str) -> Result<Self, toml::de::Error> {
65        let raw: RawConfig = toml::from_str(toml_str)?;
66        let defs = raw
67            .sources
68            .into_iter()
69            .map(|(name, s)| (name, SourceDef { kind: s.kind, path: s.path }))
70            .collect();
71        Ok(Self { defs, memo: Mutex::new(HashMap::new()) })
72    }
73
74    /// Load from a config file path. A missing file yields an empty provider —
75    /// a project with no scoped sources legitimately ships no config.
76    pub fn from_config_file(path: &Path) -> anyhow::Result<Self> {
77        match std::fs::read_to_string(path) {
78            Ok(s) => Ok(Self::from_config_str(&s)?),
79            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(Self::empty()),
80            Err(e) => Err(e.into()),
81        }
82    }
83
84    /// Current generation token for `name`, memoized for 1s. Unknown sources
85    /// (not in config) resolve to the stable placeholder.
86    pub fn token(&self, name: &str) -> String {
87        let now = Instant::now();
88        {
89            let memo = self.memo.lock().unwrap();
90            if let Some((tok, at)) = memo.get(name) {
91                if now.saturating_duration_since(*at) < GENERATION_TTL {
92                    return tok.clone();
93                }
94            }
95        }
96        let fresh = self.compute(name);
97        self.memo.lock().unwrap().insert(name.to_string(), (fresh.clone(), now));
98        fresh
99    }
100
101    fn compute(&self, name: &str) -> String {
102        let Some(def) = self.defs.get(name) else {
103            return PLACEHOLDER.to_string();
104        };
105        match def.kind.as_str() {
106            "sqlite" => def
107                .path
108                .as_deref()
109                .map(mtime_token)
110                .unwrap_or_else(|| PLACEHOLDER.to_string()),
111            // r2/pg/rpc polling is post-MVP — a stable token keeps the key
112            // correct (TTL still drives expiry); it just never self-advances.
113            _ => PLACEHOLDER.to_string(),
114        }
115    }
116}
117
118const PLACEHOLDER: &str = "0";
119
120/// File mtime as a nanosecond token, or `"missing"` when the file is absent.
121/// A non-existent DB resolves stably; once the file appears its mtime token
122/// changes, busting the key.
123fn mtime_token(path: &str) -> String {
124    match std::fs::metadata(path).and_then(|m| m.modified()) {
125        Ok(t) => t
126            .duration_since(UNIX_EPOCH)
127            .map(|d| d.as_nanos().to_string())
128            .unwrap_or_else(|_| "0".to_string()),
129        Err(_) => "missing".to_string(),
130    }
131}
132
133/// Convenience: a `SystemTime` formatted the same way `mtime_token` does, for
134/// tests that want to assert against a known mtime.
135#[cfg(test)]
136fn token_of(t: std::time::SystemTime) -> String {
137    t.duration_since(UNIX_EPOCH).map(|d| d.as_nanos().to_string()).unwrap_or_default()
138}
139
140#[cfg(test)]
141mod tests {
142    use super::*;
143    use std::io::Write;
144    use std::time::SystemTime;
145
146    #[test]
147    fn unknown_source_is_placeholder() {
148        let g = Generations::empty();
149        assert_eq!(g.token("whatever"), PLACEHOLDER);
150    }
151
152    #[test]
153    fn parses_sources_and_keeps_kind_and_path() {
154        let g = Generations::from_config_str(
155            r#"
156            [sources.project_db]
157            kind = "sqlite"
158            scope = "global"
159            path = "/tmp/x.db"
160
161            [sources.assets]
162            kind = "r2"
163            scope = "global"
164            bucket = "b"
165            endpoint_env = "R2_ENDPOINT"
166            "#,
167        )
168        .unwrap();
169        assert_eq!(g.defs.get("project_db").unwrap().kind, "sqlite");
170        assert_eq!(g.defs.get("project_db").unwrap().path.as_deref(), Some("/tmp/x.db"));
171        // r2 resolves to placeholder for MVP.
172        assert_eq!(g.token("assets"), PLACEHOLDER);
173    }
174
175    #[test]
176    fn sqlite_token_tracks_file_mtime() {
177        let dir = tempfile::tempdir().unwrap();
178        let db = dir.path().join("project.db");
179        let mut f = std::fs::File::create(&db).unwrap();
180        f.write_all(b"v1").unwrap();
181        f.sync_all().unwrap();
182
183        let cfg = format!(
184            "[sources.project_db]\nkind = \"sqlite\"\nscope = \"global\"\npath = \"{}\"\n",
185            db.display()
186        );
187        let g = Generations::from_config_str(&cfg).unwrap();
188
189        let t1 = g.token("project_db");
190        assert_ne!(t1, "missing");
191        // Same instant within the 1s memo window → identical token.
192        assert_eq!(g.token("project_db"), t1);
193
194        // Bump the mtime past the memo window and confirm the token advances.
195        let later = SystemTime::now() + Duration::from_secs(5);
196        let f2 = std::fs::OpenOptions::new().write(true).open(&db).unwrap();
197        f2.set_modified(later).unwrap();
198        std::thread::sleep(Duration::from_millis(1100));
199        let t2 = g.token("project_db");
200        assert_ne!(t1, t2, "token should follow the new mtime");
201        assert_eq!(t2, token_of(later));
202    }
203
204    #[test]
205    fn missing_sqlite_file_is_stable_missing_token() {
206        let g = Generations::from_config_str(
207            "[sources.project_db]\nkind = \"sqlite\"\npath = \"/no/such/file.db\"\n",
208        )
209        .unwrap();
210        assert_eq!(g.token("project_db"), "missing");
211    }
212}