1use std::sync::mpsc::Receiver;
8
9fn 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#[derive(Clone, Debug, Default, PartialEq, Eq)]
24pub struct SyncInSummary {
25 pub ingested: Vec<String>,
26 pub suppressed: Vec<String>,
28 pub conflicted: Vec<String>,
30 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#[derive(Clone, Debug, Default, PartialEq, Eq)]
48pub struct SyncOutSummary {
49 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
66pub 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 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 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 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 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 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 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 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 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 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 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 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 {
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 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 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 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}