1use 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#[derive(Clone, Debug, Default, PartialEq, Eq)]
18pub struct RecoveryResult {
19 pub healed: Vec<String>,
21 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
32pub 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 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 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}