Skip to main content

kimun_notes/server_client/
sync.rs

1//! Orchestration: turn observed changes and the vault's authoritative state into
2//! server pushes/deletes. [`RagSync`] wires the observer; `drain` flushes the
3//! dirty-set (the fast path); `reconcile` is the correctness backbone.
4//! All server I/O goes through [`RagTransport`], so this logic is
5//! tested with a fake against a real vault.
6
7use std::collections::HashMap;
8use std::sync::Arc;
9
10use kimun_core::{IndexObserver, NoteVault, error::VaultError, nfs::VaultPath};
11
12use crate::server_client::dto::{WireDoc, WireSection};
13use crate::server_client::{
14    DirtyOp, DirtySet, RagClient, RagError, RagObserver, RagTransport, hash_string, reconcile_diff,
15};
16
17/// What a reachable server can do, derived from `/health`: search
18/// needs an embedder, question-answering needs an embedder AND an LLM.
19#[derive(Debug, Clone, Copy, PartialEq, Eq)]
20pub enum ServerCapability {
21    /// No embedder configured — nothing works server-side; the client must not
22    /// push or reconcile (every call would 503).
23    Unconfigured,
24    /// Embedder, no LLM: search and sync work, question-answering does not.
25    SemanticOnly,
26    /// Embedder and LLM: everything works.
27    Full,
28}
29
30impl ServerCapability {
31    /// Derives the capability from a health probe's fields.
32    pub fn from_health(health: &crate::server_client::dto::Health) -> Self {
33        match (health.embedder.is_some(), health.llm_provider.is_some()) {
34            (false, _) => ServerCapability::Unconfigured,
35            (true, false) => ServerCapability::SemanticOnly,
36            (true, true) => ServerCapability::Full,
37        }
38    }
39
40    /// Whether question-answering is usable.
41    pub fn llm_available(self) -> bool {
42        matches!(self, ServerCapability::Full)
43    }
44}
45
46/// A server update the user should hear about, derived from `/health`.
47#[derive(Debug, Clone, PartialEq, Eq)]
48pub enum ServerUpdate {
49    /// The server's update check found this newer release.
50    Newer(String),
51    /// The server reports no `version` at all: it predates version reporting,
52    /// so it is behind every release that has it, but which one is unknown.
53    Legacy,
54}
55
56impl ServerUpdate {
57    /// Derives the notice from a health probe: a server without `version` is
58    /// [`Legacy`](ServerUpdate::Legacy); otherwise it is whatever newer release
59    /// the server itself reported, if any.
60    pub fn from_health(health: &crate::server_client::dto::Health) -> Option<Self> {
61        match &health.version {
62            None => Some(ServerUpdate::Legacy),
63            Some(_) => health.latest_version.clone().map(ServerUpdate::Newer),
64        }
65    }
66}
67
68/// One `/health` round-trip's worth of facts: what the server can do, and
69/// whether it gates its API behind a bearer token. `/health` itself is
70/// un-gated, so a client with a missing/wrong token still probes fine —
71/// `auth_required` lets it report "unauthorized" up front instead of
72/// discovering a 401 on the first sync call. `server_update` is the passive
73/// footer hint that the server is out of date.
74#[derive(Debug, Clone, PartialEq, Eq)]
75pub struct ServerProbe {
76    pub capability: ServerCapability,
77    pub auth_required: bool,
78    pub server_update: Option<ServerUpdate>,
79}
80
81/// Bundles a vault, its dirty-set, and the server client, and drives sync. The
82/// caller (the TUI) owns the schedule: [`probe`](RagSync::probe) to gate
83/// features, and call [`tick`](RagSync::tick) periodically to keep the server in
84/// step.
85pub struct RagSync {
86    vault: Arc<NoteVault>,
87    dirty: Arc<DirtySet>,
88    /// The exact observer this sync registered, kept so `Drop` can deregister
89    /// *only* ours (by identity) and never a newer one that replaced it.
90    observer: Arc<dyn IndexObserver>,
91    client: RagClient,
92}
93
94impl RagSync {
95    /// Registers the observer on `vault` and returns a handle over it. Construct
96    /// **one** `RagSync` per vault: the observer is zero-or-one, so a second
97    /// `RagSync` for the same vault replaces the first's observer and strands its
98    /// dirty-set (its `tick` still self-heals via reconcile, but its drain fast
99    /// path goes silent).
100    pub fn new(vault: Arc<NoteVault>, client: RagClient) -> Self {
101        let dirty = Arc::new(DirtySet::default());
102        let observer: Arc<dyn IndexObserver> = Arc::new(RagObserver::new(dirty.clone()));
103        vault.set_index_observer(observer.clone());
104        Self {
105            vault,
106            dirty,
107            observer,
108            client,
109        }
110    }
111
112    /// Probe reachability, capability, and auth in one `/health` request:
113    /// `None` = offline; otherwise a [`ServerProbe`] carrying the server's
114    /// [`ServerCapability`] — `Unconfigured` (no embedder: don't sync),
115    /// `SemanticOnly` (search, no Q&A), or `Full` — and whether the API
116    /// requires a bearer token.
117    pub async fn probe(&self) -> Option<ServerProbe> {
118        self.client.health().await.ok().map(|h| ServerProbe {
119            capability: ServerCapability::from_health(&h),
120            auth_required: h.auth_required,
121            server_update: ServerUpdate::from_health(&h),
122        })
123    }
124
125    /// Whether the local index is filled and safe to sync from. `false` while
126    /// a healed/rebuilding index is still empty — syncing then would read "no
127    /// notes" and tear the server collection down (see [`reconcile`]).
128    pub fn index_ready(&self) -> bool {
129        self.vault.index_ready()
130    }
131
132    /// One sync pass: flush pending changes, then reconcile to repair drift.
133    /// Returns `false` when either half was skipped because the local index
134    /// is not ready yet — call again once it is.
135    pub async fn tick(&self) -> Result<bool, RagError> {
136        let drained = drain(&self.vault, &self.dirty, &self.client).await?;
137        let reconciled = reconcile(&self.vault, &self.client).await?;
138        Ok(drained && reconciled)
139    }
140
141    /// Flush pending changes only — the cheap fast path (touches only dirty
142    /// notes). Run this often; run [`tick`](Self::tick)/[`reconcile`](Self::reconcile)
143    /// occasionally as the safety net. Returns `false` when skipped because
144    /// the local index is not ready.
145    pub async fn drain(&self) -> Result<bool, RagError> {
146        drain(&self.vault, &self.dirty, &self.client).await
147    }
148
149    /// Full hash-diff reconciliation only — an index-wide read + a full-collection
150    /// hash fetch. The periodic backbone; not needed on every tick. Returns
151    /// `false` when skipped because the local index is not ready.
152    pub async fn reconcile(&self) -> Result<bool, RagError> {
153        reconcile(&self.vault, &self.client).await
154    }
155
156    /// The underlying client, for queries (search / ask).
157    pub fn client(&self) -> &RagClient {
158        &self.client
159    }
160}
161
162impl Drop for RagSync {
163    fn drop(&mut self) {
164        // Deregister *our* observer so a superseded/aborted sync doesn't leave
165        // the vault feeding a dirty-set nobody drains — but only if it's still
166        // ours, so we never wipe a newer sync that has replaced it.
167        self.vault.clear_index_observer_if(&self.observer);
168    }
169}
170
171/// Builds the wire document for a note: its canonical path, content hash, and
172/// heading sections pulled from the index. Returns `None` when the note has no
173/// indexable sections — an empty note is not RAG content, so it is never pushed
174/// (this keeps both backends from perpetually re-pushing chunkless notes, since
175/// only one of them records a hash for them server-side).
176pub async fn build_doc(
177    vault: &NoteVault,
178    path: &VaultPath,
179    hash: u64,
180) -> Result<Option<WireDoc>, VaultError> {
181    let chunks = vault.get_note_chunks(path).await?;
182    let sections: Vec<WireSection> = chunks
183        .into_values()
184        .flatten()
185        .map(|c| WireSection {
186            title: c.get_breadcrumb().to_string(),
187            text: c.get_text().to_string(),
188        })
189        .collect();
190    if sections.is_empty() {
191        return Ok(None);
192    }
193    Ok(Some(WireDoc {
194        path: path.to_string(),
195        hash: hash_string(hash),
196        sections,
197    }))
198}
199
200/// Flushes the dirty-set to the server. Failed operations are re-queued so the
201/// next drain (or a reconcile) retries them.
202///
203/// Returns `false` (doing nothing) when the local index is not ready: an
204/// unready (healed/rebuilding) index reads as empty, so build_doc would find
205/// no chunks and turn queued upserts into server-side deletes. The dirty-set
206/// stays queued until the index is filled; callers should retry then.
207pub async fn drain<T: RagTransport>(
208    vault: &NoteVault,
209    dirty: &DirtySet,
210    transport: &T,
211) -> Result<bool, RagError> {
212    if !vault.index_ready() {
213        return Ok(false);
214    }
215    let ops = dirty.drain();
216    if ops.is_empty() {
217        return Ok(true);
218    }
219
220    let mut upserts: Vec<(VaultPath, u64)> = Vec::new();
221    let mut deletes: Vec<String> = Vec::new();
222    for (path, op) in ops {
223        match op {
224            DirtyOp::Upsert(hash) => upserts.push((path, hash)),
225            DirtyOp::Delete => deletes.push(path.to_string()),
226        }
227    }
228
229    // Build the docs for upserts; a note that can't be read right now is
230    // re-queued rather than dropped.
231    let mut docs = Vec::new();
232    let mut built: Vec<(VaultPath, u64)> = Vec::new();
233    for (path, hash) in upserts {
234        match build_doc(vault, &path, hash).await {
235            Ok(Some(doc)) => {
236                docs.push(doc);
237                built.push((path, hash));
238            }
239            // An emptied note has no chunks to index — delete it server-side so
240            // its old chunks don't linger (and so /hashes stops reporting it).
241            Ok(None) => deletes.push(path.to_string()),
242            Err(_) => dirty.requeue([(path, DirtyOp::Upsert(hash))]),
243        }
244    }
245
246    let mut first_err: Option<RagError> = None;
247    if !docs.is_empty()
248        && let Err(e) = transport.push_docs(docs).await
249    {
250        dirty.requeue(built.into_iter().map(|(p, h)| (p, DirtyOp::Upsert(h))));
251        first_err = Some(e);
252    }
253    if !deletes.is_empty() {
254        let paths_for_requeue: Vec<VaultPath> = deletes.iter().map(VaultPath::new).collect();
255        if let Err(e) = transport.delete_paths(deletes).await {
256            dirty.requeue(paths_for_requeue.into_iter().map(|p| (p, DirtyOp::Delete)));
257            first_err = first_err.or(Some(e));
258        }
259    }
260
261    match first_err {
262        Some(e) => Err(e),
263        None => Ok(true),
264    }
265}
266
267/// Reconciles the server with the vault: diff hash sets, then push/delete only
268/// the differences. Self-healing — repairs anything the drain path missed.
269///
270/// Returns `false` (doing nothing) when the local index is not ready: a
271/// healed/rebuilding index reads as an empty vault, and diffing against that
272/// snapshot would put every server doc in `to_delete` — wiping the collection.
273/// Callers should retry once the index is filled.
274pub async fn reconcile<T: RagTransport>(
275    vault: &NoteVault,
276    transport: &T,
277) -> Result<bool, RagError> {
278    if !vault.index_ready() {
279        return Ok(false);
280    }
281    let notes = vault
282        .get_all_notes()
283        .await
284        .map_err(|e| RagError::Protocol(format!("read vault notes: {e}")))?;
285
286    let local_hashes: HashMap<String, u64> = notes
287        .into_iter()
288        .map(|(entry, content)| (entry.path.to_string(), content.hash))
289        .collect();
290    let local_str: HashMap<String, String> = local_hashes
291        .iter()
292        .map(|(p, h)| (p.clone(), hash_string(*h)))
293        .collect();
294
295    let server = transport.server_hashes().await?;
296    let plan = reconcile_diff(&local_str, &server);
297
298    let mut docs = Vec::new();
299    let mut to_delete = plan.to_delete;
300    for path_str in &plan.to_push {
301        let hash = local_hashes[path_str];
302        match build_doc(vault, &VaultPath::new(path_str), hash)
303            .await
304            .map_err(|e| RagError::Protocol(format!("build doc {path_str}: {e}")))?
305        {
306            Some(doc) => docs.push(doc),
307            // Empty note: it carries no chunks. If the server still has it (it
308            // was emptied), delete it; if not, it's already converged (nothing
309            // to push). Either way it never becomes a stale server entry.
310            None => {
311                if server.contains_key(path_str) {
312                    to_delete.push(path_str.clone());
313                }
314            }
315        }
316    }
317    if !docs.is_empty() {
318        transport.push_docs(docs).await?;
319    }
320    if !to_delete.is_empty() {
321        transport.delete_paths(to_delete).await?;
322    }
323    Ok(true)
324}
325
326#[cfg(test)]
327mod tests {
328    use super::*;
329    use async_trait::async_trait;
330
331    #[test]
332    fn server_update_from_health_fields() {
333        use crate::server_client::dto::Health;
334        let h = |version: Option<&str>, latest: Option<&str>| Health {
335            status: "ok".into(),
336            reranker: false,
337            embedder: None,
338            llm_provider: None,
339            auth_required: false,
340            version: version.map(str::to_string),
341            latest_version: latest.map(str::to_string),
342        };
343        assert_eq!(ServerUpdate::from_health(&h(Some("0.4.4"), None)), None);
344        assert_eq!(
345            ServerUpdate::from_health(&h(Some("0.4.4"), Some("0.4.5"))),
346            Some(ServerUpdate::Newer("0.4.5".into()))
347        );
348        // A 0.4.3-or-earlier server's /health has no version at all.
349        let legacy: Health = serde_json::from_str(
350            r#"{"status":"ok","reranker":false,"embedder":"fastembed","llm_provider":null,"auth_required":false,"degraded":null}"#,
351        )
352        .unwrap();
353        assert_eq!(
354            ServerUpdate::from_health(&legacy),
355            Some(ServerUpdate::Legacy)
356        );
357    }
358
359    #[test]
360    fn capability_from_health_fields() {
361        use crate::server_client::dto::Health;
362        let h = |embedder: Option<&str>, llm: Option<&str>| Health {
363            status: "ok".into(),
364            reranker: false,
365            embedder: embedder.map(str::to_string),
366            llm_provider: llm.map(str::to_string),
367            auth_required: false,
368            version: Some("0.4.4".into()),
369            latest_version: None,
370        };
371        assert_eq!(
372            ServerCapability::from_health(&h(None, None)),
373            ServerCapability::Unconfigured
374        );
375        assert_eq!(
376            // An LLM without an embedder still can't answer — retrieval is dead.
377            ServerCapability::from_health(&h(None, Some("gemini"))),
378            ServerCapability::Unconfigured
379        );
380        assert_eq!(
381            ServerCapability::from_health(&h(Some("fastembed"), None)),
382            ServerCapability::SemanticOnly
383        );
384        assert_eq!(
385            ServerCapability::from_health(&h(Some("fastembed"), Some("gemini"))),
386            ServerCapability::Full
387        );
388    }
389    use kimun_core::VaultConfig;
390    use std::sync::Mutex;
391    use tempfile::TempDir;
392
393    #[derive(Default)]
394    struct FakeTransport {
395        pushed: Mutex<Vec<WireDoc>>,
396        deleted: Mutex<Vec<String>>,
397        server: Mutex<HashMap<String, String>>,
398        fail_push: Mutex<bool>,
399    }
400
401    #[async_trait]
402    impl RagTransport for FakeTransport {
403        async fn push_docs(&self, docs: Vec<WireDoc>) -> Result<(), RagError> {
404            if *self.fail_push.lock().unwrap() {
405                return Err(RagError::Protocol("boom".into()));
406            }
407            self.pushed.lock().unwrap().extend(docs);
408            Ok(())
409        }
410        async fn delete_paths(&self, paths: Vec<String>) -> Result<(), RagError> {
411            self.deleted.lock().unwrap().extend(paths);
412            Ok(())
413        }
414        async fn server_hashes(&self) -> Result<HashMap<String, String>, RagError> {
415            Ok(self.server.lock().unwrap().clone())
416        }
417    }
418
419    /// Test-only observer wiring feeding a bare dirty-set, for exercising the
420    /// free `drain`/`reconcile` fns directly. Production always goes through
421    /// [`RagSync::new`], which additionally keeps the observer handle so its
422    /// `Drop` can deregister by identity.
423    fn register(vault: &NoteVault) -> Arc<DirtySet> {
424        let dirty = Arc::new(DirtySet::default());
425        vault.set_index_observer(Arc::new(RagObserver::new(dirty.clone())));
426        dirty
427    }
428
429    /// A vault root as a `SystemPath`. Built from `kimun_core` rather than the
430    /// crate's own test helper: this module must not name anything outside
431    /// itself, so it stays extractable as a crate (adr/0042).
432    fn sys(path: impl AsRef<std::path::Path>) -> kimun_core::SystemPath {
433        kimun_core::SystemPath::try_absolute(path).expect("test path must be absolute")
434    }
435
436    async fn vault(dir: &std::path::Path) -> NoteVault {
437        let vault = NoteVault::new(VaultConfig::new(sys(dir))).await.unwrap();
438        // Fill the freshly-healed index so index_ready() holds — the state the
439        // drain/reconcile gates require (mirrors the app's validate_and_init).
440        vault.validate_and_init().await.unwrap();
441        vault
442    }
443
444    #[tokio::test]
445    async fn drain_pushes_created_note_and_deletes_removed() {
446        let dir = TempDir::new().unwrap();
447        let vault = vault(dir.path()).await;
448        let dirty = register(&vault);
449        let transport = FakeTransport::default();
450
451        vault
452            .create_note(&VaultPath::new("a.md"), "# Title\n\nbody")
453            .await
454            .unwrap();
455        drain(&vault, &dirty, &transport).await.unwrap();
456
457        // Block-scoped: clippy's await_holding_lock tracks lexical scope, so an
458        // explicit drop() before the awaits below wouldn't silence it.
459        {
460            let pushed = transport.pushed.lock().unwrap();
461            assert_eq!(pushed.len(), 1);
462            assert_eq!(pushed[0].path, "/a.md"); // canonical
463            assert!(!pushed[0].sections.is_empty());
464            assert!(dirty.is_empty());
465        }
466
467        vault.delete_note(&VaultPath::new("a.md")).await.unwrap();
468        drain(&vault, &dirty, &transport).await.unwrap();
469        assert_eq!(
470            *transport.deleted.lock().unwrap(),
471            vec!["/a.md".to_string()]
472        );
473    }
474
475    #[tokio::test]
476    async fn failed_push_requeues() {
477        let dir = TempDir::new().unwrap();
478        let vault = vault(dir.path()).await;
479        let dirty = register(&vault);
480        let transport = FakeTransport::default();
481        *transport.fail_push.lock().unwrap() = true;
482
483        vault
484            .create_note(&VaultPath::new("a.md"), "body")
485            .await
486            .unwrap();
487        assert!(drain(&vault, &dirty, &transport).await.is_err());
488        // The op survived for a later retry.
489        assert_eq!(dirty.len(), 1);
490    }
491
492    #[tokio::test]
493    async fn reconcile_pushes_missing_and_deletes_stale() {
494        let dir = TempDir::new().unwrap();
495        let vault = vault(dir.path()).await;
496        let _dirty = register(&vault);
497        let transport = FakeTransport::default();
498
499        vault
500            .create_note(&VaultPath::new("keep.md"), "kept")
501            .await
502            .unwrap();
503        // Server already has a stale note the vault no longer contains.
504        transport
505            .server
506            .lock()
507            .unwrap()
508            .insert("/gone.md".to_string(), "oldhash".to_string());
509
510        assert!(reconcile(&vault, &transport).await.unwrap());
511
512        let pushed = transport.pushed.lock().unwrap();
513        assert!(pushed.iter().any(|d| d.path == "/keep.md"));
514        assert_eq!(
515            *transport.deleted.lock().unwrap(),
516            vec!["/gone.md".to_string()]
517        );
518    }
519
520    #[tokio::test]
521    async fn reconcile_skipped_while_index_not_ready() {
522        let dir = TempDir::new().unwrap();
523        // No validate_and_init: the fresh index is healed-but-empty, exactly
524        // the state where a reconcile would read "no local notes" and delete
525        // the whole server collection.
526        let vault = NoteVault::new(VaultConfig::new(sys(dir.path())))
527            .await
528            .unwrap();
529        assert!(!vault.index_ready());
530        let transport = FakeTransport::default();
531        transport
532            .server
533            .lock()
534            .unwrap()
535            .insert("/precious.md".to_string(), "hash".to_string());
536
537        assert!(!reconcile(&vault, &transport).await.unwrap());
538        assert!(transport.deleted.lock().unwrap().is_empty());
539        assert!(transport.pushed.lock().unwrap().is_empty());
540
541        // Drain likewise holds queued ops instead of misreading the empty
542        // index (an upsert would otherwise become a server-side delete),
543        // and reports the skip so callers don't claim the vault is synced.
544        let dirty = register(&vault);
545        dirty.record(&kimun_core::NoteChange::Upsert {
546            path: VaultPath::new("precious.md"),
547            hash: 1,
548        });
549        assert!(!drain(&vault, &dirty, &transport).await.unwrap());
550        assert_eq!(dirty.len(), 1);
551        assert!(transport.deleted.lock().unwrap().is_empty());
552    }
553}