Skip to main content

omgbase_sync/
coordinator.rs

1//! The coordinator (`spec/sync/README.md` §6): drives a source against an
2//! engine client with no reconciliation logic of its own — whole-file bytes
3//! in, the change feed out. Loop safety: the engine's echo gate makes a
4//! written file that comes back an echo, and observed commits are never
5//! exported.
6
7use std::sync::mpsc::Receiver;
8
9/// The storage form of a path the engine reported (`spec/surface` §1 "Paths",
10/// 2.0: every path a tool returns is `/`-rooted; a source speaks the
11/// repo-relative form). Every leading `/` is stripped.
12fn storage_path(path: &str) -> String {
13    path.trim_start_matches('/').to_owned()
14}
15
16use serde_json::Value;
17
18use crate::engine::{EngineClient, FileBytes};
19use crate::error::Result;
20use crate::source::{SyncSource, WatchEvent};
21
22/// What a `sync_in`/`reconcile` did, by path.
23#[derive(Clone, Debug, Default, PartialEq, Eq)]
24pub struct SyncInSummary {
25    pub ingested: Vec<String>,
26    /// Echoes: the bytes already matched, no commit.
27    pub suppressed: Vec<String>,
28    /// Ingested but carrying git conflict markers.
29    pub conflicted: Vec<String>,
30    /// Gone paths whose live doc was tombstoned.
31    pub deleted: Vec<String>,
32}
33
34impl SyncInSummary {
35    #[must_use]
36    pub fn to_json(&self) -> Value {
37        serde_json::json!({
38            "ingested": self.ingested,
39            "suppressed": self.suppressed,
40            "conflicted": self.conflicted,
41            "deleted": self.deleted,
42        })
43    }
44}
45
46/// What a `sync_out` did.
47#[derive(Clone, Debug, Default, PartialEq, Eq)]
48pub struct SyncOutSummary {
49    /// The final feed cursor.
50    pub cursor: i64,
51    pub written: Vec<String>,
52    pub removed: Vec<String>,
53}
54
55impl SyncOutSummary {
56    #[must_use]
57    pub fn to_json(&self) -> Value {
58        serde_json::json!({
59            "cursor": self.cursor,
60            "written": self.written,
61            "removed": self.removed,
62        })
63    }
64}
65
66/// A source paired with an engine.
67pub struct Coordinator<'a> {
68    engine: &'a mut dyn EngineClient,
69    source: &'a mut dyn SyncSource,
70}
71
72impl<'a> Coordinator<'a> {
73    pub fn new(engine: &'a mut dyn EngineClient, source: &'a mut dyn SyncSource) -> Self {
74        Self { engine, source }
75    }
76
77    /// source → engine: the source's full current scope.
78    pub fn sync_in(&mut self) -> Result<SyncInSummary> {
79        let paths: Vec<String> = self
80            .source
81            .enumerate()?
82            .into_iter()
83            .map(|e| e.path)
84            .collect();
85        self.reconcile(&paths)
86    }
87
88    /// source → engine for a set of paths: fetch each; present items go to
89    /// `observe_many` in one call, gone paths to `observe_delete` one by one.
90    pub fn reconcile(&mut self, paths: &[String]) -> Result<SyncInSummary> {
91        let mut files: Vec<FileBytes> = Vec::new();
92        let mut gone: Vec<String> = Vec::new();
93        for path in paths {
94            match self.source.fetch(path)? {
95                Some(item) => files.push(FileBytes {
96                    path: path.clone(),
97                    content: item.content,
98                }),
99                None => gone.push(path.clone()),
100            }
101        }
102        let mut summary = SyncInSummary::default();
103        if !files.is_empty() {
104            // The engine answers in the surface's reference form (`/a.md`,
105            // `spec/surface` §1 "Paths", 2.0); a source speaks storage paths.
106            for r in self.engine.observe_many(&files)? {
107                let path = storage_path(&r.path);
108                if r.echo {
109                    summary.suppressed.push(path);
110                } else if r.conflicted {
111                    summary.conflicted.push(path);
112                } else {
113                    summary.ingested.push(path);
114                }
115            }
116        }
117        for path in gone {
118            if self.engine.observe_delete(&path)?.deleted() {
119                summary.deleted.push(path);
120            }
121        }
122        Ok(summary)
123    }
124
125    /// engine → source: page `changes_since(cursor)`; for every digest whose
126    /// origin is not `observed`, every revision's doc is re-read by path and
127    /// written, or removed when it no longer reads; follows `truncated`
128    /// pages; returns the final cursor. A source without write-through
129    /// exports nothing and returns `cursor` unchanged.
130    pub fn sync_out(&mut self, cursor: i64) -> Result<SyncOutSummary> {
131        let mut summary = SyncOutSummary {
132            cursor,
133            ..SyncOutSummary::default()
134        };
135        if !self.source.capabilities().write_through {
136            return Ok(summary);
137        }
138        let mut cur = cursor;
139        loop {
140            let page = self.engine.changes_since(cur, None, None)?;
141            for digest in &page.digests {
142                if digest.origin == "observed" {
143                    continue;
144                }
145                for rev in &digest.revisions {
146                    // `changes_since` reports the reference form; the source takes storage paths.
147                    let path = storage_path(&rev.path);
148                    match self.engine.read_doc(&path)? {
149                        Some(doc) => {
150                            self.source.write(&path, &doc.content)?;
151                            summary.written.push(path);
152                        }
153                        None => {
154                            self.source.remove(&path)?;
155                            summary.removed.push(path);
156                        }
157                    }
158                }
159            }
160            cur = page.cursor;
161            if !page.truncated {
162                break;
163            }
164        }
165        summary.cursor = cur;
166        Ok(summary)
167    }
168
169    /// Live source → engine: subscribe when the source can watch; `None`
170    /// otherwise. The stream yields [`WatchEvent::Ready`] once the feed is
171    /// primed (a host waits for it — [`crate::wait_ready`] — before its
172    /// priming sweep, §5), then batches. Drive it with
173    /// [`Coordinator::handle_batches`] or call [`Coordinator::reconcile`] per
174    /// received batch.
175    pub fn watch_in(&mut self) -> Result<Option<Receiver<WatchEvent>>> {
176        if !self.source.capabilities().watch {
177            return Ok(None);
178        }
179        Ok(Some(self.source.watch()?))
180    }
181
182    /// Reconcile every batch the stream yields until it closes (the adapter
183    /// exited or `unwatch` ran), reporting each summary or error. `Ready` is
184    /// not a batch and is skipped.
185    pub fn handle_batches(
186        &mut self,
187        events: &Receiver<WatchEvent>,
188        mut on_summary: impl FnMut(SyncInSummary),
189        mut on_error: impl FnMut(crate::error::Error),
190    ) {
191        for event in events.iter() {
192            let WatchEvent::Batch(paths) = event else {
193                continue;
194            };
195            match self.reconcile(&paths) {
196                Ok(s) => on_summary(s),
197                Err(e) => on_error(e),
198            }
199        }
200    }
201
202    /// Stop the watch stream.
203    pub fn stop_watch(&mut self) -> Result<()> {
204        self.source.unwatch()
205    }
206}
207
208#[cfg(test)]
209mod tests {
210    use super::*;
211    use crate::engine::InProcessEngineClient;
212    use crate::source::{MemSource, SourceCapabilities};
213    use omgbase_store::{NullDocStore, SequentialMinter, Store};
214
215    const TS: &str = "2026-09-26T10:00:00.000Z";
216
217    #[test]
218    fn sync_in_reconcile_and_sync_out_round_trip() {
219        let mut store =
220            Store::open_in_memory_with_minter(Box::new(SequentialMinter::new())).unwrap();
221        let repo = store.create_repo("fixture").unwrap();
222        let mut source = MemSource::with_files(&[("a.md", "# A\n\nOne.\n"), ("b.md", "# B\n")]);
223        {
224            let mut engine = InProcessEngineClient::new(&mut store, &repo).at(TS);
225            let mut co = Coordinator::new(&mut engine, &mut source);
226            let s = co.sync_in().unwrap();
227            assert_eq!(s.ingested, ["a.md", "b.md"]);
228            assert!(s.suppressed.is_empty() && s.deleted.is_empty() && s.conflicted.is_empty());
229            // Again: everything echoes.
230            let s = co.sync_in().unwrap();
231            assert_eq!(s.suppressed, ["a.md", "b.md"]);
232            assert_eq!(
233                s.to_json()["suppressed"],
234                serde_json::json!(["a.md", "b.md"])
235            );
236            // Observed commits are never exported.
237            let out = co.sync_out(0).unwrap();
238            assert_eq!(out.cursor, 2);
239            assert!(out.written.is_empty() && out.removed.is_empty());
240        }
241        // An engine-authored (api) commit exports; a gone path deletes.
242        let ctx = omgbase_store::DocOpContext {
243            repo_id: repo.clone(),
244            actor: Some("agent:test".into()),
245            ts: TS.into(),
246        };
247        store
248            .docs_create(&ctx, &mut NullDocStore, "authored.md", "# Authored\n", None)
249            .unwrap();
250        source.files.remove("b.md");
251        source.set("c.md", "<<<<<<< a\nx\n=======\ny\n>>>>>>> b\n");
252        let out = {
253            let mut engine = InProcessEngineClient::new(&mut store, &repo).at(TS);
254            let mut co = Coordinator::new(&mut engine, &mut source);
255            let s = co
256                .reconcile(&["b.md".to_owned(), "c.md".to_owned(), "zzz.md".to_owned()])
257                .unwrap();
258            assert_eq!(s.deleted, ["b.md"]);
259            assert_eq!(s.conflicted, ["c.md"]);
260            co.sync_out(2).unwrap()
261        };
262        assert_eq!(out.written, ["authored.md"]);
263        assert!(out.removed.is_empty());
264        assert_eq!(out.cursor, 5);
265        assert_eq!(source.files["authored.md"], "# Authored\n");
266        assert_eq!(source.log.len(), 1);
267        assert_eq!(out.to_json()["cursor"], 5);
268        // The written file comes back as an echo: the loop terminates.
269        {
270            let mut engine = InProcessEngineClient::new(&mut store, &repo).at(TS);
271            let mut co = Coordinator::new(&mut engine, &mut source);
272            let s = co.reconcile(&["authored.md".to_owned()]).unwrap();
273            assert_eq!(s.suppressed, ["authored.md"]);
274            assert_eq!(co.sync_out(5).unwrap().cursor, 5);
275        }
276    }
277
278    #[test]
279    fn sync_out_removes_tombstoned_docs_and_pages() {
280        let mut store =
281            Store::open_in_memory_with_minter(Box::new(SequentialMinter::new())).unwrap();
282        let repo = store.create_repo("fixture").unwrap();
283        let mut source = MemSource::with_files(&[]);
284        // 60 api commits via docs_create so paging (limit 50) is exercised.
285        let ctx = omgbase_store::DocOpContext {
286            repo_id: repo.clone(),
287            actor: Some("t".into()),
288            ts: TS.into(),
289        };
290        for i in 0..60 {
291            store
292                .docs_create(
293                    &ctx,
294                    &mut NullDocStore,
295                    &format!("n{i:02}.md"),
296                    "# N\n",
297                    None,
298                )
299                .unwrap();
300        }
301        store
302            .docs_delete(&ctx, &mut NullDocStore, "n00.md")
303            .unwrap();
304        let out = {
305            let mut engine = InProcessEngineClient::new(&mut store, &repo).at(TS);
306            let mut co = Coordinator::new(&mut engine, &mut source);
307            co.sync_out(0).unwrap()
308        };
309        assert_eq!(out.cursor, 61);
310        assert_eq!(out.written.len(), 59);
311        assert_eq!(out.removed, ["n00.md"]);
312        assert_eq!(source.files.len(), 59);
313        // Read-only source: nothing exported, cursor unchanged.
314        let mut ro = MemSource::new(SourceCapabilities::default());
315        let mut engine = InProcessEngineClient::new(&mut store, &repo).at(TS);
316        let mut co = Coordinator::new(&mut engine, &mut ro);
317        assert_eq!(
318            co.sync_out(3).unwrap(),
319            SyncOutSummary {
320                cursor: 3,
321                ..SyncOutSummary::default()
322            }
323        );
324        assert!(co.watch_in().unwrap().is_none());
325    }
326
327    #[test]
328    fn watch_in_reconciles_batches() {
329        let mut store = Store::open_in_memory().unwrap();
330        let repo = store.create_repo("fixture").unwrap();
331        let mut source = MemSource::with_files(&[("a.md", "# A\n")]);
332        let rx = {
333            let mut engine = InProcessEngineClient::new(&mut store, &repo);
334            let mut co = Coordinator::new(&mut engine, &mut source);
335            co.watch_in().unwrap().expect("watching source")
336        };
337        source.emit(&["a.md"]);
338        source.emit(&["gone.md"]);
339        // Stopping the watch drops the sender, so the stream ends.
340        source.unwatch().unwrap();
341        let mut engine = InProcessEngineClient::new(&mut store, &repo);
342        let mut co = Coordinator::new(&mut engine, &mut source);
343        let mut summaries = Vec::new();
344        co.handle_batches(&rx, |s| summaries.push(s), |e| panic!("{e}"));
345        assert_eq!(summaries.len(), 2, "the leading `Ready` is not a batch");
346        assert_eq!(summaries[0].ingested, ["a.md"]);
347        assert!(
348            summaries[1].deleted.is_empty(),
349            "nothing was live at gone.md"
350        );
351        co.stop_watch().unwrap();
352    }
353}