1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
20pub enum ServerCapability {
21 Unconfigured,
24 SemanticOnly,
26 Full,
28}
29
30impl ServerCapability {
31 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 pub fn llm_available(self) -> bool {
42 matches!(self, ServerCapability::Full)
43 }
44}
45
46#[derive(Debug, Clone, PartialEq, Eq)]
48pub enum ServerUpdate {
49 Newer(String),
51 Legacy,
54}
55
56impl ServerUpdate {
57 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#[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
81pub struct RagSync {
86 vault: Arc<NoteVault>,
87 dirty: Arc<DirtySet>,
88 observer: Arc<dyn IndexObserver>,
91 client: RagClient,
92}
93
94impl RagSync {
95 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 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 pub fn index_ready(&self) -> bool {
129 self.vault.index_ready()
130 }
131
132 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 pub async fn drain(&self) -> Result<bool, RagError> {
146 drain(&self.vault, &self.dirty, &self.client).await
147 }
148
149 pub async fn reconcile(&self) -> Result<bool, RagError> {
153 reconcile(&self.vault, &self.client).await
154 }
155
156 pub fn client(&self) -> &RagClient {
158 &self.client
159 }
160}
161
162impl Drop for RagSync {
163 fn drop(&mut self) {
164 self.vault.clear_index_observer_if(&self.observer);
168 }
169}
170
171pub 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
200pub 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 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 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
267pub 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 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 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 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 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 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 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 {
460 let pushed = transport.pushed.lock().unwrap();
461 assert_eq!(pushed.len(), 1);
462 assert_eq!(pushed[0].path, "/a.md"); 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 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 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 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 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}