Skip to main content

omgbase_sync/
recovery.rs

1//! Startup recovery (`spec/sync/README.md` §4.3 "Recovery"): every live
2//! doc's file is checked against its current revision; a divergent file is
3//! re-ingested as an observed commit, one document at a time (§9).
4
5use std::path::Path;
6
7use omgbase_format::hash::sha256;
8use omgbase_reconcile::Config;
9use omgbase_store::{Origin, Store};
10use rusqlite::params;
11use serde_json::Value;
12
13use crate::error::Result;
14use crate::fs::FileSystem;
15
16/// What recovery found.
17#[derive(Clone, Debug, Default, PartialEq, Eq)]
18pub struct RecoveryResult {
19    /// Paths re-ingested.
20    pub healed: Vec<String>,
21    /// Live docs whose file is gone.
22    pub missing: Vec<String>,
23}
24
25impl RecoveryResult {
26    #[must_use]
27    pub fn to_json(&self) -> Value {
28        serde_json::json!({ "healed": self.healed, "missing": self.missing })
29    }
30}
31
32/// For every live doc (in row order): a missing file is `missing`; a file
33/// whose hash differs from the current revision's `rendered_hash` (or
34/// `file_hash` when there is no revision) is re-ingested with the reconciling
35/// resolver and listed `healed`.
36pub fn recover_repo(
37    store: &mut Store,
38    repo_id: &str,
39    fs: &dyn FileSystem,
40    root: &Path,
41    ts: &str,
42    config: &Config,
43) -> Result<RecoveryResult> {
44    type Row = (String, Option<Vec<u8>>, Option<Vec<u8>>);
45    let docs: Vec<Row> = {
46        let mut stmt = store.conn().prepare(
47            "SELECT d.path, d.file_hash, r.rendered_hash
48             FROM docs d LEFT JOIN revisions r ON r.rev_id = d.current_rev
49             WHERE d.repo_id = ?1 AND d.deleted_commit IS NULL ORDER BY d.rowid",
50        )?;
51        let rows = stmt.query_map(params![repo_id], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?;
52        rows.collect::<std::result::Result<_, _>>()?
53    };
54    let mut result = RecoveryResult::default();
55    for (path, file_hash, rendered_hash) in docs {
56        let Some(content) = fs.read(root, &path)? else {
57            result.missing.push(path);
58            continue;
59        };
60        let on_disk = sha256(content.as_bytes());
61        let recorded = rendered_hash.or(file_hash);
62        if recorded.is_none_or(|r| r[..] != on_disk[..]) {
63            store.reconciling_ingest(
64                repo_id,
65                &path,
66                &content,
67                ts,
68                Origin::Observed,
69                None,
70                None,
71                config,
72            )?;
73            result.healed.push(path);
74        }
75    }
76    Ok(result)
77}
78
79#[cfg(test)]
80mod tests {
81    use super::*;
82    use crate::fs::MemFileSystem;
83    use omgbase_store::{BatchItem, SequentialMinter};
84
85    const TS: &str = "2026-09-26T10:00:00.000Z";
86
87    #[test]
88    fn heals_divergent_files_and_reports_missing() {
89        let mut store =
90            Store::open_in_memory_with_minter(Box::new(SequentialMinter::new())).unwrap();
91        let repo = store.create_repo("fixture").unwrap();
92        let cfg = Config::default();
93        let items = [
94            BatchItem::observed("a.md", "# A\n\nOne.\n"),
95            BatchItem::observed("b.md", "# B\n"),
96            BatchItem::observed("c.md", "# C\n"),
97        ];
98        store.observe_batch(&repo, &items, TS, &cfg).unwrap();
99        let mut fs = MemFileSystem::new();
100        fs.set("a.md", "# A\n\nOne, edited.\n", 1);
101        fs.set("b.md", "# B\n", 2);
102        let root = Path::new("/r");
103        let r = recover_repo(&mut store, &repo, &fs, root, TS, &cfg).unwrap();
104        assert_eq!(r.healed, ["a.md"]);
105        assert_eq!(r.missing, ["c.md"]);
106        assert_eq!(r.to_json()["healed"], serde_json::json!(["a.md"]));
107        let commits: i64 = store
108            .conn()
109            .query_row(
110                "SELECT count(*) FROM commits WHERE repo_id = ?1",
111                params![repo],
112                |r| r.get(0),
113            )
114            .unwrap();
115        assert_eq!(commits, 4);
116        let origin: String = store
117            .conn()
118            .query_row("SELECT origin FROM commits WHERE seq = 4", [], |r| r.get(0))
119            .unwrap();
120        assert_eq!(origin, "observed");
121        let doc: String = store
122            .conn()
123            .query_row("SELECT doc_id FROM docs WHERE path = 'a.md'", [], |r| {
124                r.get(0)
125            })
126            .unwrap();
127        assert_eq!(
128            store.reconstruct(&doc).unwrap().as_deref(),
129            Some("# A\n\nOne, edited.\n")
130        );
131        // The heading kept its id: the reconciling resolver was used.
132        let carried: i64 = store
133            .conn()
134            .query_row(
135                "SELECT count(*) FROM dispositions WHERE commit_id = 'c_3' AND kind = 'same'",
136                [],
137                |r| r.get(0),
138            )
139            .unwrap();
140        assert_eq!(carried, 1);
141        // Idempotent.
142        let again = recover_repo(&mut store, &repo, &fs, root, TS, &cfg).unwrap();
143        assert!(again.healed.is_empty());
144        assert_eq!(again.missing, ["c.md"]);
145    }
146}