1use std::fs::{File, OpenOptions};
10use std::io::Write;
11use std::path::{Path, PathBuf};
12use std::time::{Duration, Instant};
13
14use super::search_manifest::{
15 SearchArtifactError, SearchArtifactKey, SearchIndexKind, SearchManifest, SearchSourceSnapshot,
16};
17
18const CURRENT_FILE: &str = "current.json";
19const MANIFEST_FILE: &str = "manifest.json";
20const VERSIONS_DIR: &str = "versions";
21const BUILD_PREFIX: &str = ".build-";
22const VERSION_PREFIX: &str = "version-";
23const MAX_CURRENT_BYTES: usize = 4096;
24
25#[derive(Clone, Copy, Debug)]
27pub struct SearchCoordinationLimits {
28 pub lock_timeout: Duration,
30 pub lock_poll_interval: Duration,
32 pub cleanup_entries: usize,
34}
35
36impl Default for SearchCoordinationLimits {
37 fn default() -> Self {
38 Self {
39 lock_timeout: Duration::from_secs(30),
40 lock_poll_interval: Duration::from_millis(25),
41 cleanup_entries: 10_000,
42 }
43 }
44}
45
46#[derive(Clone, Copy, Debug, PartialEq, Eq)]
48pub enum SearchPublicationMode {
49 ReuseFresh,
51 Replace,
53}
54
55#[derive(Clone, Copy, Debug)]
57pub struct SearchPublicationPlan<'a> {
58 pub key: &'a SearchArtifactKey,
60 pub backend_version: &'a str,
62 pub contract_version: &'a str,
64 pub dimension: Option<u32>,
66 pub mode: SearchPublicationMode,
68}
69
70#[derive(Clone, Debug, PartialEq, Eq)]
72pub struct PublishedSearchArtifact {
73 pub path: PathBuf,
75 pub manifest: SearchManifest,
77}
78
79#[derive(Clone, Debug, PartialEq, Eq)]
81pub enum SearchPublicationOutcome {
82 Reused(PublishedSearchArtifact),
84 Published {
86 artifact: PublishedSearchArtifact,
88 attempts: u8,
90 },
91}
92
93#[derive(Clone, Copy, Debug, PartialEq, Eq)]
95pub enum SearchUpdateBuild {
96 ReuseCurrent,
98 Publish,
100}
101
102pub fn coordinate_search_publication<S, V, B, C>(
117 project_dir: &Path,
118 plan: SearchPublicationPlan<'_>,
119 limits: SearchCoordinationLimits,
120 mut snapshot: S,
121 mut validate_current: V,
122 mut build: B,
123 mut checkpoint: C,
124) -> Result<SearchPublicationOutcome, SearchArtifactError>
125where
126 S: FnMut() -> Result<SearchSourceSnapshot, SearchArtifactError>,
127 V: FnMut(&PublishedSearchArtifact) -> Result<(), SearchArtifactError>,
128 B: FnMut(&Path, &SearchSourceSnapshot) -> Result<(), SearchArtifactError>,
129 C: FnMut() -> Result<(), SearchArtifactError>,
130{
131 let root = plan.key.artifact_root(project_dir);
132 std::fs::create_dir_all(&root).map_err(|source| io("create artifact root", &root, source))?;
133 let _writer = SearchWriterLock::acquire(&root, limits, &mut checkpoint)?;
134 checkpoint()?;
135
136 let initial = snapshot()?;
137 match current_search_artifact(project_dir, plan.key) {
138 Ok(Some(artifact)) if plan.mode == SearchPublicationMode::ReuseFresh => {
139 match artifact.manifest.verify_fresh(
140 plan.key,
141 plan.backend_version,
142 plan.contract_version,
143 plan.dimension,
144 &initial,
145 ) {
146 Ok(()) => match validate_current(&artifact) {
147 Ok(()) => return Ok(SearchPublicationOutcome::Reused(artifact)),
148 Err(error)
149 if plan.key.kind() == SearchIndexKind::Text
150 && rebuildable_metadata(&error) => {}
151 Err(error) if plan.key.kind() == SearchIndexKind::Vector => {
152 return Err(primary_vector_error(root, error));
153 }
154 Err(error) => return Err(error),
155 },
156 Err(SearchArtifactError::Stale { .. })
157 if plan.key.kind() == SearchIndexKind::Text => {}
158 Err(error) => return Err(error),
159 }
160 }
161 Ok(Some(_) | None) => {}
162 Err(error) if plan.key.kind() == SearchIndexKind::Text && rebuildable_metadata(&error) => {}
163 Err(error) if plan.key.kind() == SearchIndexKind::Vector => {
164 return Err(primary_vector_error(root, error));
165 }
166 Err(error) => return Err(error),
167 }
168
169 let mut before = initial;
170 for attempt in 1_u8..=2 {
171 checkpoint()?;
172 let publication = PendingPublication::new(&root)?;
173 build(publication.path(), &before)?;
174 checkpoint()?;
175 let after = snapshot()?;
176 if before != after {
177 if attempt == 2 {
178 return Err(SearchArtifactError::ConcurrentMutation);
179 }
180 before = after;
181 continue;
182 }
183 let manifest = SearchManifest::for_key(
184 plan.key,
185 plan.backend_version,
186 plan.contract_version,
187 plan.dimension,
188 &before,
189 true,
190 )?;
191 let artifact = publication.publish(&manifest)?;
192 return Ok(SearchPublicationOutcome::Published {
193 artifact,
194 attempts: attempt,
195 });
196 }
197 unreachable!("the bounded publication loop returns on both terminal paths")
198}
199
200pub fn coordinate_search_update<S, B, C>(
218 project_dir: &Path,
219 plan: SearchPublicationPlan<'_>,
220 limits: SearchCoordinationLimits,
221 mut snapshot: S,
222 mut build: B,
223 mut checkpoint: C,
224) -> Result<SearchPublicationOutcome, SearchArtifactError>
225where
226 S: FnMut() -> Result<SearchSourceSnapshot, SearchArtifactError>,
227 B: FnMut(
228 Option<&PublishedSearchArtifact>,
229 &Path,
230 &SearchSourceSnapshot,
231 &mut C,
232 ) -> Result<SearchUpdateBuild, SearchArtifactError>,
233 C: FnMut() -> Result<(), SearchArtifactError>,
234{
235 if plan.mode != SearchPublicationMode::Replace {
236 return Err(SearchArtifactError::Build(
237 "atomic search updates require replacement mode".to_owned(),
238 ));
239 }
240 let root = plan.key.artifact_root(project_dir);
241 std::fs::create_dir_all(&root).map_err(|source| io("create artifact root", &root, source))?;
242 let _writer = SearchWriterLock::acquire(&root, limits, &mut checkpoint)?;
243 checkpoint()?;
244 let current = match current_search_artifact(project_dir, plan.key) {
245 Ok(current) => current,
246 Err(error) if plan.key.kind() == SearchIndexKind::Vector => {
247 return Err(primary_vector_error(root, error));
248 }
249 Err(error) => return Err(error),
250 };
251
252 let mut before = snapshot()?;
253 for attempt in 1_u8..=2 {
254 checkpoint()?;
255 let publication = PendingPublication::new(&root)?;
256 let decision = build(
257 current.as_ref(),
258 publication.path(),
259 &before,
260 &mut checkpoint,
261 )?;
262 checkpoint()?;
263 let after = snapshot()?;
264 if before != after {
265 if attempt == 2 {
266 return Err(SearchArtifactError::ConcurrentMutation);
267 }
268 before = after;
269 continue;
270 }
271 match decision {
272 SearchUpdateBuild::ReuseCurrent => {
273 return current
274 .clone()
275 .map(SearchPublicationOutcome::Reused)
276 .ok_or_else(|| {
277 SearchArtifactError::Build(
278 "update requested reuse without a current artifact".to_owned(),
279 )
280 });
281 }
282 SearchUpdateBuild::Publish => {
283 let manifest = SearchManifest::for_key(
284 plan.key,
285 plan.backend_version,
286 plan.contract_version,
287 plan.dimension,
288 &before,
289 true,
290 )?;
291 let artifact = publication.publish(&manifest)?;
292 return Ok(SearchPublicationOutcome::Published {
293 artifact,
294 attempts: attempt,
295 });
296 }
297 }
298 }
299 unreachable!("the bounded update loop returns on both terminal paths")
300}
301
302pub fn current_search_artifact(
310 project_dir: &Path,
311 key: &SearchArtifactKey,
312) -> Result<Option<PublishedSearchArtifact>, SearchArtifactError> {
313 let root = key.artifact_root(project_dir);
314 let pointer = root.join(CURRENT_FILE);
315 let bytes = match std::fs::read(&pointer) {
316 Ok(bytes) => bytes,
317 Err(source) if source.kind() == std::io::ErrorKind::NotFound => return Ok(None),
318 Err(source) => return Err(io("read current pointer", &pointer, source)),
319 };
320 if bytes.len() > MAX_CURRENT_BYTES {
321 return Err(SearchArtifactError::ResourceExhausted {
322 resource: "current_pointer_bytes",
323 limit: MAX_CURRENT_BYTES as u64,
324 });
325 }
326 let value: serde_json::Value =
327 serde_json::from_slice(&bytes).map_err(|error| SearchArtifactError::CorruptManifest {
328 path: pointer.clone(),
329 reason: error.to_string(),
330 })?;
331 let version = value
332 .as_object()
333 .and_then(|object| object.get("version"))
334 .and_then(serde_json::Value::as_str)
335 .filter(|name| valid_owned_name(name, VERSION_PREFIX))
336 .ok_or_else(|| SearchArtifactError::CorruptManifest {
337 path: pointer,
338 reason: "expected a safe version pointer".to_owned(),
339 })?;
340 let artifact_dir = root.join(VERSIONS_DIR).join(version);
341 let manifest_path = artifact_dir.join(MANIFEST_FILE);
342 let manifest_bytes = match std::fs::read(&manifest_path) {
343 Ok(bytes) => bytes,
344 Err(source) if source.kind() == std::io::ErrorKind::NotFound => {
345 return Err(SearchArtifactError::Missing {
346 path: manifest_path,
347 });
348 }
349 Err(source) => return Err(io("read published manifest", &manifest_path, source)),
350 };
351 let manifest = SearchManifest::from_json(&manifest_path, &manifest_bytes)?;
352 if !manifest.completed {
353 return Err(SearchArtifactError::CorruptManifest {
354 path: manifest_path,
355 reason: "published manifest is not completed".to_owned(),
356 });
357 }
358 Ok(Some(PublishedSearchArtifact {
359 path: artifact_dir,
360 manifest,
361 }))
362}
363
364pub fn cleanup_abandoned_search_builds(
376 project_dir: &Path,
377 max_entries: usize,
378) -> Result<usize, SearchArtifactError> {
379 let roots = [
380 project_dir.join("indexes").join("search"),
381 project_dir.join("embeddings"),
382 ];
383 let mut stack = roots.to_vec();
384 let mut inspected = 0_usize;
385 let mut remove = Vec::new();
386 while let Some(directory) = stack.pop() {
387 let entries = match std::fs::read_dir(&directory) {
388 Ok(entries) => entries,
389 Err(source) if source.kind() == std::io::ErrorKind::NotFound => continue,
390 Err(source) => return Err(io("scan abandoned builds", &directory, source)),
391 };
392 for entry in entries {
393 inspected = inspected
394 .checked_add(1)
395 .ok_or(SearchArtifactError::ResourceExhausted {
396 resource: "cleanup_entries",
397 limit: max_entries as u64,
398 })?;
399 if inspected > max_entries {
400 return Err(SearchArtifactError::ResourceExhausted {
401 resource: "cleanup_entries",
402 limit: max_entries as u64,
403 });
404 }
405 let entry = entry.map_err(|source| io("read cleanup entry", &directory, source))?;
406 let file_type = entry
407 .file_type()
408 .map_err(|source| io("read cleanup file type", &entry.path(), source))?;
409 let name = entry.file_name();
410 let name = name.to_string_lossy();
411 if file_type.is_dir() && valid_owned_name(&name, BUILD_PREFIX) {
412 remove.push((entry.path(), true));
413 } else if file_type.is_dir() && !file_type.is_symlink() {
414 stack.push(entry.path());
415 } else if file_type.is_file() && valid_pointer_temp(&name) {
416 remove.push((entry.path(), false));
417 }
418 }
419 }
420
421 remove.sort_unstable_by(|left, right| right.0.cmp(&left.0));
422 for (path, directory) in &remove {
423 if *directory {
424 std::fs::remove_dir_all(path)
425 .map_err(|source| io("remove abandoned build", path, source))?;
426 } else {
427 std::fs::remove_file(path)
428 .map_err(|source| io("remove abandoned pointer", path, source))?;
429 }
430 }
431 Ok(remove.len())
432}
433
434struct SearchWriterLock {
435 file: File,
436}
437
438impl SearchWriterLock {
439 fn acquire<C>(
440 root: &Path,
441 limits: SearchCoordinationLimits,
442 checkpoint: &mut C,
443 ) -> Result<Self, SearchArtifactError>
444 where
445 C: FnMut() -> Result<(), SearchArtifactError>,
446 {
447 let path = root.join(".writer.lock");
448 let file = OpenOptions::new()
449 .read(true)
450 .write(true)
451 .create(true)
452 .truncate(false)
453 .open(&path)
454 .map_err(|source| SearchArtifactError::Lock {
455 path: path.clone(),
456 reason: source.to_string(),
457 })?;
458 let started = Instant::now();
459 loop {
460 match file.try_lock() {
461 Ok(()) => return Ok(Self { file }),
462 Err(std::fs::TryLockError::WouldBlock) => {
463 checkpoint()?;
464 if started.elapsed() >= limits.lock_timeout {
465 return Err(SearchArtifactError::Lock {
466 path,
467 reason: format!(
468 "timed out after {} ms",
469 limits.lock_timeout.as_millis()
470 ),
471 });
472 }
473 std::thread::sleep(limits.lock_poll_interval);
474 }
475 Err(std::fs::TryLockError::Error(source)) => {
476 return Err(SearchArtifactError::Lock {
477 path,
478 reason: source.to_string(),
479 });
480 }
481 }
482 }
483 }
484}
485
486impl Drop for SearchWriterLock {
487 fn drop(&mut self) {
488 let _ = self.file.unlock();
489 }
490}
491
492struct PendingPublication {
493 root: PathBuf,
494 temp: tempfile::TempDir,
495}
496
497impl PendingPublication {
498 fn new(root: &Path) -> Result<Self, SearchArtifactError> {
499 let temp = tempfile::Builder::new()
500 .prefix(BUILD_PREFIX)
501 .tempdir_in(root)
502 .map_err(|source| io("create build directory", root, source))?;
503 Ok(Self {
504 root: root.to_path_buf(),
505 temp,
506 })
507 }
508
509 fn path(&self) -> &Path {
510 self.temp.path()
511 }
512
513 fn publish(
514 self,
515 manifest: &SearchManifest,
516 ) -> Result<PublishedSearchArtifact, SearchArtifactError> {
517 let manifest_path = self.temp.path().join(MANIFEST_FILE);
518 let manifest_bytes = manifest.to_canonical_json()?;
519 write_synced_file(&manifest_path, &manifest_bytes)?;
520 sync_tree(self.temp.path())?;
521
522 let token = self
523 .temp
524 .path()
525 .file_name()
526 .and_then(|name| name.to_str())
527 .and_then(|name| name.strip_prefix(BUILD_PREFIX))
528 .filter(|token| {
529 !token.is_empty() && token.bytes().all(|byte| byte.is_ascii_alphanumeric())
530 })
531 .ok_or_else(|| {
532 SearchArtifactError::Build("temporary build name is invalid".to_owned())
533 })?;
534 let version_name = format!("{VERSION_PREFIX}{token}");
535 let versions = self.root.join(VERSIONS_DIR);
536 std::fs::create_dir_all(&versions)
537 .map_err(|source| io("create versions directory", &versions, source))?;
538 let version_path = versions.join(&version_name);
539 let temp_path = self.temp.keep();
540 if let Err(source) = std::fs::rename(&temp_path, &version_path) {
541 let _ = std::fs::remove_dir_all(&temp_path);
542 return Err(io("publish immutable version", &version_path, source));
543 }
544 sync_directory(&versions)?;
545
546 let pointer = serde_json::to_vec(&serde_json::json!({ "version": version_name }))
547 .map_err(|error| SearchArtifactError::Build(error.to_string()))?;
548 persist_synced_pointer(&self.root.join(CURRENT_FILE), &pointer)?;
549 sync_directory(&self.root)?;
550 Ok(PublishedSearchArtifact {
551 path: version_path,
552 manifest: manifest.clone(),
553 })
554 }
555}
556
557fn rebuildable_metadata(error: &SearchArtifactError) -> bool {
558 matches!(
559 error,
560 SearchArtifactError::Missing { .. }
561 | SearchArtifactError::CorruptManifest { .. }
562 | SearchArtifactError::CorruptDerivedIndex { .. }
563 | SearchArtifactError::IncompatibleManifest { .. }
564 | SearchArtifactError::Stale { .. }
565 | SearchArtifactError::ResourceExhausted {
566 resource: "manifest_bytes" | "current_pointer_bytes",
567 ..
568 }
569 )
570}
571
572fn primary_vector_error(root: PathBuf, error: SearchArtifactError) -> SearchArtifactError {
573 match error {
574 error @ SearchArtifactError::CorruptPrimaryVectors { .. } => error,
575 error => SearchArtifactError::CorruptPrimaryVectors {
576 path: root,
577 reason: error.to_string(),
578 },
579 }
580}
581
582fn write_synced_file(path: &Path, bytes: &[u8]) -> Result<(), SearchArtifactError> {
583 let mut file = OpenOptions::new()
584 .create(true)
585 .truncate(true)
586 .write(true)
587 .open(path)
588 .map_err(|source| io("create publication file", path, source))?;
589 file.write_all(bytes)
590 .map_err(|source| io("write publication file", path, source))?;
591 file.sync_all()
592 .map_err(|source| io("sync publication file", path, source))
593}
594
595fn persist_synced_pointer(path: &Path, bytes: &[u8]) -> Result<(), SearchArtifactError> {
596 let parent = path
597 .parent()
598 .ok_or_else(|| SearchArtifactError::Build("current pointer has no parent".to_owned()))?;
599 let mut temp = tempfile::Builder::new()
600 .prefix("current.json.")
601 .suffix(".tmp")
602 .tempfile_in(parent)
603 .map_err(|source| io("create current pointer temp", path, source))?;
604 temp.write_all(bytes)
605 .map_err(|source| io("write current pointer temp", path, source))?;
606 temp.as_file()
607 .sync_all()
608 .map_err(|source| io("sync current pointer temp", path, source))?;
609 temp.persist(path)
610 .map_err(|error| io("publish current pointer", path, error.error))?;
611 Ok(())
612}
613
614fn sync_tree(root: &Path) -> Result<(), SearchArtifactError> {
615 let mut directories = vec![root.to_path_buf()];
616 let mut files = Vec::new();
617 let mut cursor = 0;
618 while cursor < directories.len() {
619 let directory = directories[cursor].clone();
620 cursor += 1;
621 let entries = std::fs::read_dir(&directory)
622 .map_err(|source| io("scan build for sync", &directory, source))?;
623 for entry in entries {
624 let entry = entry.map_err(|source| io("read build entry", &directory, source))?;
625 let file_type = entry
626 .file_type()
627 .map_err(|source| io("read build file type", &entry.path(), source))?;
628 if file_type.is_symlink() {
629 return Err(SearchArtifactError::Build(format!(
630 "search build must not contain symlink {}",
631 entry.path().display()
632 )));
633 }
634 if file_type.is_dir() {
635 directories.push(entry.path());
636 } else if file_type.is_file() {
637 files.push(entry.path());
638 }
639 }
640 }
641 files.sort_unstable();
642 for path in files {
643 sync_file(&path)?;
644 }
645 directories.sort_unstable_by_key(|path| std::cmp::Reverse(path.components().count()));
646 for directory in directories {
647 sync_directory(&directory)?;
648 }
649 Ok(())
650}
651
652fn sync_file(path: &Path) -> Result<(), SearchArtifactError> {
653 OpenOptions::new()
654 .read(true)
655 .write(true)
656 .open(path)
657 .and_then(|file| file.sync_all())
658 .map_err(|source| io("sync build file", path, source))
659}
660
661#[cfg(unix)]
662fn sync_directory(path: &Path) -> Result<(), SearchArtifactError> {
663 File::open(path)
664 .and_then(|file| file.sync_all())
665 .map_err(|source| io("sync directory", path, source))
666}
667
668#[cfg(not(unix))]
669fn sync_directory(_path: &Path) -> Result<(), SearchArtifactError> {
670 Ok(())
671}
672
673fn valid_owned_name(name: &str, prefix: &str) -> bool {
674 name.strip_prefix(prefix).is_some_and(|token| {
675 !token.is_empty() && token.bytes().all(|byte| byte.is_ascii_alphanumeric())
676 })
677}
678
679fn valid_pointer_temp(name: &str) -> bool {
680 name.strip_prefix("current.json.")
681 .and_then(|name| name.strip_suffix(".tmp"))
682 .is_some_and(|token| {
683 !token.is_empty() && token.bytes().all(|byte| byte.is_ascii_alphanumeric())
684 })
685}
686
687fn io(operation: &'static str, path: &Path, source: std::io::Error) -> SearchArtifactError {
688 SearchArtifactError::Io {
689 operation,
690 path: path.to_path_buf(),
691 source,
692 }
693}
694
695#[cfg(test)]
696mod tests {
697 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
698 use std::sync::{Arc, Barrier};
699
700 use tempfile::TempDir;
701
702 use super::*;
703
704 fn key() -> SearchArtifactKey {
705 SearchArtifactKey::text("Person", ["name"]).unwrap()
706 }
707
708 fn snapshot(generation: u64) -> SearchSourceSnapshot {
709 SearchSourceSnapshot {
710 generation,
711 fingerprint: format!("gf-fnv1a256:{generation:064x}"),
712 }
713 }
714
715 fn plan<'a>(
716 key: &'a SearchArtifactKey,
717 mode: SearchPublicationMode,
718 ) -> SearchPublicationPlan<'a> {
719 SearchPublicationPlan {
720 key,
721 backend_version: "tantivy-0.25",
722 contract_version: "graphforge_text_v1",
723 dimension: None,
724 mode,
725 }
726 }
727
728 #[test]
729 fn publish_is_invisible_until_current_pointer_swap() {
730 let dir = TempDir::new().unwrap();
731 let key = key();
732 let root = key.artifact_root(dir.path());
733 std::fs::create_dir_all(&root).unwrap();
734 let pending = PendingPublication::new(&root).unwrap();
735 std::fs::write(pending.path().join("index"), b"complete").unwrap();
736 assert!(current_search_artifact(dir.path(), &key).unwrap().is_none());
737
738 let manifest = SearchManifest::for_key(
739 &key,
740 "tantivy-0.25",
741 "graphforge_text_v1",
742 None,
743 &snapshot(1),
744 true,
745 )
746 .unwrap();
747 let published = pending.publish(&manifest).unwrap();
748 assert_eq!(
749 current_search_artifact(dir.path(), &key).unwrap(),
750 Some(published)
751 );
752 }
753
754 #[test]
755 fn sync_tree_flushes_regular_files_without_mutating_contents() {
756 let dir = TempDir::new().unwrap();
757 let nested = dir.path().join("nested");
758 std::fs::create_dir(&nested).unwrap();
759 let first = dir.path().join("first");
760 let second = nested.join("second");
761 std::fs::write(&first, b"one").unwrap();
762 std::fs::write(&second, b"two").unwrap();
763
764 sync_tree(dir.path()).unwrap();
765
766 assert_eq!(std::fs::read(first).unwrap(), b"one");
767 assert_eq!(std::fs::read(second).unwrap(), b"two");
768 }
769
770 #[test]
771 fn forced_rebuild_atomically_replaces_pointer_and_keeps_old_version() {
772 let dir = TempDir::new().unwrap();
773 let key = key();
774 let first = coordinate_search_publication(
775 dir.path(),
776 plan(&key, SearchPublicationMode::Replace),
777 SearchCoordinationLimits::default(),
778 || Ok(snapshot(1)),
779 |_| Ok(()),
780 |path, _| {
781 std::fs::write(path.join("index"), b"first")
782 .map_err(|error| SearchArtifactError::Build(error.to_string()))
783 },
784 || Ok(()),
785 )
786 .unwrap();
787 let first_path = match first {
788 SearchPublicationOutcome::Published { artifact, .. } => artifact.path,
789 SearchPublicationOutcome::Reused(_) => panic!("forced build reused"),
790 };
791 let second = coordinate_search_publication(
792 dir.path(),
793 plan(&key, SearchPublicationMode::Replace),
794 SearchCoordinationLimits::default(),
795 || Ok(snapshot(1)),
796 |_| Ok(()),
797 |path, _| {
798 std::fs::write(path.join("index"), b"second")
799 .map_err(|error| SearchArtifactError::Build(error.to_string()))
800 },
801 || Ok(()),
802 )
803 .unwrap();
804 let second_path = match second {
805 SearchPublicationOutcome::Published { artifact, .. } => artifact.path,
806 SearchPublicationOutcome::Reused(_) => panic!("forced build reused"),
807 };
808 assert_ne!(first_path, second_path);
809 assert!(
810 first_path.exists(),
811 "old readers retain an immutable version"
812 );
813 assert_eq!(
814 current_search_artifact(dir.path(), &key)
815 .unwrap()
816 .unwrap()
817 .path,
818 second_path
819 );
820 }
821
822 #[test]
823 fn failed_replacement_keeps_the_previous_publication() {
824 let dir = TempDir::new().unwrap();
825 let key = key();
826 coordinate_search_publication(
827 dir.path(),
828 plan(&key, SearchPublicationMode::Replace),
829 SearchCoordinationLimits::default(),
830 || Ok(snapshot(1)),
831 |_| Ok(()),
832 |path, _| {
833 std::fs::write(path.join("index"), b"committed")
834 .map_err(|error| SearchArtifactError::Build(error.to_string()))
835 },
836 || Ok(()),
837 )
838 .unwrap();
839 let previous = current_search_artifact(dir.path(), &key).unwrap().unwrap();
840
841 let error = coordinate_search_publication(
842 dir.path(),
843 plan(&key, SearchPublicationMode::Replace),
844 SearchCoordinationLimits::default(),
845 || Ok(snapshot(1)),
846 |_| Ok(()),
847 |path, _| {
848 std::fs::write(path.join("index"), b"partial").unwrap();
849 Err(SearchArtifactError::Build("injected failure".to_owned()))
850 },
851 || Ok(()),
852 )
853 .unwrap_err();
854 assert!(matches!(error, SearchArtifactError::Build(_)));
855 assert_eq!(
856 current_search_artifact(dir.path(), &key).unwrap().unwrap(),
857 previous
858 );
859 assert_eq!(
860 std::fs::read(previous.path.join("index")).unwrap(),
861 b"committed"
862 );
863 }
864
865 #[test]
866 fn fresh_lazy_request_reuses_without_running_builder() {
867 let dir = TempDir::new().unwrap();
868 let key = key();
869 let builds = AtomicUsize::new(0);
870 let run = |mode| {
871 coordinate_search_publication(
872 dir.path(),
873 plan(&key, mode),
874 SearchCoordinationLimits::default(),
875 || Ok(snapshot(1)),
876 |_| Ok(()),
877 |path, _| {
878 builds.fetch_add(1, Ordering::SeqCst);
879 std::fs::write(path.join("index"), b"data")
880 .map_err(|error| SearchArtifactError::Build(error.to_string()))
881 },
882 || Ok(()),
883 )
884 };
885 assert!(matches!(
886 run(SearchPublicationMode::ReuseFresh).unwrap(),
887 SearchPublicationOutcome::Published { .. }
888 ));
889 assert!(matches!(
890 run(SearchPublicationMode::ReuseFresh).unwrap(),
891 SearchPublicationOutcome::Reused(_)
892 ));
893 assert_eq!(builds.load(Ordering::SeqCst), 1);
894 }
895
896 #[test]
897 fn corrupt_derived_backend_is_rebuilt_before_reuse() {
898 let dir = TempDir::new().unwrap();
899 let key = key();
900 let builds = AtomicUsize::new(0);
901 let run = || {
902 coordinate_search_publication(
903 dir.path(),
904 plan(&key, SearchPublicationMode::ReuseFresh),
905 SearchCoordinationLimits::default(),
906 || Ok(snapshot(1)),
907 |artifact| {
908 let path = artifact.path.join("index");
909 let bytes = std::fs::read(&path).map_err(|error| {
910 SearchArtifactError::CorruptDerivedIndex {
911 path: path.clone(),
912 reason: error.to_string(),
913 }
914 })?;
915 if bytes == b"valid" {
916 Ok(())
917 } else {
918 Err(SearchArtifactError::CorruptDerivedIndex {
919 path,
920 reason: "backend validation failed".to_owned(),
921 })
922 }
923 },
924 |path, _| {
925 builds.fetch_add(1, Ordering::SeqCst);
926 std::fs::write(path.join("index"), b"valid")
927 .map_err(|error| SearchArtifactError::Build(error.to_string()))
928 },
929 || Ok(()),
930 )
931 };
932 assert!(matches!(
933 run().unwrap(),
934 SearchPublicationOutcome::Published { .. }
935 ));
936 let current = current_search_artifact(dir.path(), &key).unwrap().unwrap();
937 std::fs::write(current.path.join("index"), b"corrupt").unwrap();
938 assert!(matches!(
939 run().unwrap(),
940 SearchPublicationOutcome::Published { .. }
941 ));
942 assert_eq!(builds.load(Ordering::SeqCst), 2);
943 }
944
945 #[test]
946 fn mutation_retries_once_and_second_mutation_fails_closed() {
947 let dir = TempDir::new().unwrap();
948 let key = key();
949 let reads = AtomicUsize::new(0);
950 let outcome = coordinate_search_publication(
951 dir.path(),
952 plan(&key, SearchPublicationMode::Replace),
953 SearchCoordinationLimits::default(),
954 || {
955 let call = reads.fetch_add(1, Ordering::SeqCst);
956 Ok(snapshot(u64::from(call >= 1)))
957 },
958 |_| Ok(()),
959 |path, source| {
960 std::fs::write(path.join("source"), source.generation.to_string())
961 .map_err(|error| SearchArtifactError::Build(error.to_string()))
962 },
963 || Ok(()),
964 )
965 .unwrap();
966 assert!(matches!(
967 outcome,
968 SearchPublicationOutcome::Published { attempts: 2, .. }
969 ));
970
971 let reads = AtomicUsize::new(0);
972 let error = coordinate_search_publication(
973 dir.path(),
974 plan(&key, SearchPublicationMode::Replace),
975 SearchCoordinationLimits::default(),
976 || Ok(snapshot(reads.fetch_add(1, Ordering::SeqCst) as u64)),
977 |_| Ok(()),
978 |path, _| {
979 std::fs::write(path.join("index"), b"data")
980 .map_err(|error| SearchArtifactError::Build(error.to_string()))
981 },
982 || Ok(()),
983 )
984 .unwrap_err();
985 assert!(matches!(error, SearchArtifactError::ConcurrentMutation));
986 }
987
988 #[test]
989 fn atomic_update_requires_replace_and_cannot_reuse_missing_publication() {
990 let dir = TempDir::new().unwrap();
991 let key = key();
992 let error = coordinate_search_update(
993 dir.path(),
994 plan(&key, SearchPublicationMode::ReuseFresh),
995 SearchCoordinationLimits::default(),
996 || Ok(snapshot(1)),
997 |_, _, _, _| Ok(SearchUpdateBuild::Publish),
998 || Ok(()),
999 )
1000 .unwrap_err();
1001 assert!(
1002 matches!(error, SearchArtifactError::Build(reason) if reason.contains("replacement mode"))
1003 );
1004
1005 let error = coordinate_search_update(
1006 dir.path(),
1007 plan(&key, SearchPublicationMode::Replace),
1008 SearchCoordinationLimits::default(),
1009 || Ok(snapshot(1)),
1010 |current, _, _, _| {
1011 assert!(current.is_none());
1012 Ok(SearchUpdateBuild::ReuseCurrent)
1013 },
1014 || Ok(()),
1015 )
1016 .unwrap_err();
1017 assert!(
1018 matches!(error, SearchArtifactError::Build(reason) if reason.contains("without a current artifact"))
1019 );
1020 assert!(current_search_artifact(dir.path(), &key).unwrap().is_none());
1021 }
1022
1023 #[test]
1024 fn same_key_requests_serialize_and_share_one_lazy_build() {
1025 let dir = Arc::new(TempDir::new().unwrap());
1026 let key = Arc::new(key());
1027 let builds = Arc::new(AtomicUsize::new(0));
1028 let barrier = Arc::new(Barrier::new(2));
1029 let mut handles = Vec::new();
1030 for _ in 0..2 {
1031 let dir = Arc::clone(&dir);
1032 let key = Arc::clone(&key);
1033 let builds = Arc::clone(&builds);
1034 let barrier = Arc::clone(&barrier);
1035 handles.push(std::thread::spawn(move || {
1036 barrier.wait();
1037 coordinate_search_publication(
1038 dir.path(),
1039 plan(&key, SearchPublicationMode::ReuseFresh),
1040 SearchCoordinationLimits::default(),
1041 || Ok(snapshot(1)),
1042 |_| Ok(()),
1043 |path, _| {
1044 builds.fetch_add(1, Ordering::SeqCst);
1045 std::thread::sleep(Duration::from_millis(75));
1046 std::fs::write(path.join("index"), b"data")
1047 .map_err(|error| SearchArtifactError::Build(error.to_string()))
1048 },
1049 || Ok(()),
1050 )
1051 }));
1052 }
1053 let outcomes = handles
1054 .into_iter()
1055 .map(|handle| handle.join().unwrap().unwrap())
1056 .collect::<Vec<_>>();
1057 assert_eq!(builds.load(Ordering::SeqCst), 1);
1058 assert!(
1059 outcomes
1060 .iter()
1061 .any(|outcome| matches!(outcome, SearchPublicationOutcome::Reused(_)))
1062 );
1063 }
1064
1065 #[test]
1066 fn cancellation_while_waiting_does_not_publish() {
1067 let dir = TempDir::new().unwrap();
1068 let key = key();
1069 let root = key.artifact_root(dir.path());
1070 std::fs::create_dir_all(&root).unwrap();
1071 let lock =
1072 SearchWriterLock::acquire(&root, SearchCoordinationLimits::default(), &mut || Ok(()))
1073 .unwrap();
1074 let cancelled = AtomicBool::new(false);
1075 let limits = SearchCoordinationLimits {
1076 lock_timeout: Duration::from_secs(1),
1077 lock_poll_interval: Duration::from_millis(1),
1078 ..SearchCoordinationLimits::default()
1079 };
1080 let result = SearchWriterLock::acquire(&root, limits, &mut || {
1081 if cancelled.swap(true, Ordering::SeqCst) {
1082 Err(SearchArtifactError::Cancelled)
1083 } else {
1084 Ok(())
1085 }
1086 });
1087 drop(lock);
1088 assert!(matches!(result, Err(SearchArtifactError::Cancelled)));
1089 assert!(current_search_artifact(dir.path(), &key).unwrap().is_none());
1090 }
1091
1092 #[test]
1093 fn cancellation_after_build_does_not_publish_partial_output() {
1094 let dir = TempDir::new().unwrap();
1095 let key = key();
1096 let checkpoints = AtomicUsize::new(0);
1097 let error = coordinate_search_publication(
1098 dir.path(),
1099 plan(&key, SearchPublicationMode::Replace),
1100 SearchCoordinationLimits::default(),
1101 || Ok(snapshot(1)),
1102 |_| Ok(()),
1103 |path, _| {
1104 std::fs::write(path.join("index"), b"partial")
1105 .map_err(|error| SearchArtifactError::Build(error.to_string()))
1106 },
1107 || {
1108 if checkpoints.fetch_add(1, Ordering::SeqCst) >= 2 {
1109 Err(SearchArtifactError::Cancelled)
1110 } else {
1111 Ok(())
1112 }
1113 },
1114 )
1115 .unwrap_err();
1116 assert!(matches!(error, SearchArtifactError::Cancelled));
1117 assert!(current_search_artifact(dir.path(), &key).unwrap().is_none());
1118 }
1119
1120 #[test]
1121 fn cleanup_is_bounded_scoped_and_preserves_unknown_files() {
1122 let dir = TempDir::new().unwrap();
1123 let search = dir.path().join("indexes/search/text/key");
1124 let embeddings = dir.path().join("embeddings/space/key");
1125 let notes = dir.path().join("notes");
1126 for path in [&search, &embeddings, ¬es] {
1127 std::fs::create_dir_all(path).unwrap();
1128 }
1129 let stale = [
1130 search.join(".build-Abc123"),
1131 embeddings.join(".build-Xyz789"),
1132 ];
1133 for path in &stale {
1134 std::fs::create_dir_all(path).unwrap();
1135 std::fs::write(path.join("partial"), b"x").unwrap();
1136 }
1137 let pointer_temp = search.join("current.json.Qwe456.tmp");
1138 std::fs::write(&pointer_temp, b"x").unwrap();
1139 let preserved = [
1140 search.join(".build-bad-name"),
1141 embeddings.join("vectors.parquet"),
1142 notes.join(".build-Abc123"),
1143 ];
1144 for path in &preserved {
1145 if let Some(parent) = path.parent() {
1146 std::fs::create_dir_all(parent).unwrap();
1147 }
1148 std::fs::write(path, b"keep").unwrap();
1149 }
1150
1151 assert!(matches!(
1152 cleanup_abandoned_search_builds(dir.path(), 1),
1153 Err(SearchArtifactError::ResourceExhausted {
1154 resource: "cleanup_entries",
1155 ..
1156 })
1157 ));
1158 assert!(stale.iter().all(|path| path.exists()));
1159 assert_eq!(cleanup_abandoned_search_builds(dir.path(), 100).unwrap(), 3);
1160 assert!(stale.iter().all(|path| !path.exists()));
1161 assert!(!pointer_temp.exists());
1162 assert!(preserved.iter().all(|path| path.exists()));
1163 }
1164
1165 #[test]
1166 fn missing_published_text_manifest_is_rebuilt() {
1167 let dir = TempDir::new().unwrap();
1168 let key = key();
1169 let root = key.artifact_root(dir.path());
1170 std::fs::create_dir_all(&root).unwrap();
1171 std::fs::write(root.join(CURRENT_FILE), br#"{"version":"version-Abc123"}"#).unwrap();
1172 let builds = AtomicUsize::new(0);
1173 let outcome = coordinate_search_publication(
1174 dir.path(),
1175 plan(&key, SearchPublicationMode::ReuseFresh),
1176 SearchCoordinationLimits::default(),
1177 || Ok(snapshot(1)),
1178 |_| Ok(()),
1179 |path, _| {
1180 builds.fetch_add(1, Ordering::SeqCst);
1181 std::fs::write(path.join("index"), b"rebuilt")
1182 .map_err(|error| SearchArtifactError::Build(error.to_string()))
1183 },
1184 || Ok(()),
1185 )
1186 .unwrap();
1187 assert!(matches!(
1188 outcome,
1189 SearchPublicationOutcome::Published { attempts: 1, .. }
1190 ));
1191 assert_eq!(builds.load(Ordering::SeqCst), 1);
1192 }
1193
1194 #[test]
1195 fn corrupt_and_incompatible_text_manifests_are_rebuilt() {
1196 for incompatible in [false, true] {
1197 let dir = TempDir::new().unwrap();
1198 let key = key();
1199 let builds = AtomicUsize::new(0);
1200 let run = || {
1201 coordinate_search_publication(
1202 dir.path(),
1203 plan(&key, SearchPublicationMode::ReuseFresh),
1204 SearchCoordinationLimits::default(),
1205 || Ok(snapshot(1)),
1206 |_| Ok(()),
1207 |path, _| {
1208 builds.fetch_add(1, Ordering::SeqCst);
1209 std::fs::write(path.join("index"), b"complete")
1210 .map_err(|error| SearchArtifactError::Build(error.to_string()))
1211 },
1212 || Ok(()),
1213 )
1214 };
1215 run().unwrap();
1216 let current = current_search_artifact(dir.path(), &key).unwrap().unwrap();
1217 let manifest_path = current.path.join(MANIFEST_FILE);
1218 if incompatible {
1219 let mut manifest: serde_json::Value =
1220 serde_json::from_slice(&std::fs::read(&manifest_path).unwrap()).unwrap();
1221 manifest["manifest_version"] = serde_json::Value::from(99);
1222 std::fs::write(&manifest_path, serde_json::to_vec(&manifest).unwrap()).unwrap();
1223 } else {
1224 std::fs::write(&manifest_path, b"corrupt").unwrap();
1225 }
1226
1227 assert!(matches!(
1228 run().unwrap(),
1229 SearchPublicationOutcome::Published { attempts: 1, .. }
1230 ));
1231 assert_eq!(builds.load(Ordering::SeqCst), 2);
1232 }
1233 }
1234
1235 #[test]
1236 fn corrupt_vector_metadata_is_not_discarded() {
1237 let dir = TempDir::new().unwrap();
1238 let key = SearchArtifactKey::vector("Person", "semantic").unwrap();
1239 let root = key.artifact_root(dir.path());
1240 std::fs::create_dir_all(&root).unwrap();
1241 std::fs::write(root.join(CURRENT_FILE), b"corrupt").unwrap();
1242 let error = coordinate_search_publication(
1243 dir.path(),
1244 SearchPublicationPlan {
1245 key: &key,
1246 backend_version: "exact-cosine-v1",
1247 contract_version: "vector-v1",
1248 dimension: Some(3),
1249 mode: SearchPublicationMode::Replace,
1250 },
1251 SearchCoordinationLimits::default(),
1252 || Ok(snapshot(1)),
1253 |_| Ok(()),
1254 |_, _| Ok(()),
1255 || Ok(()),
1256 )
1257 .unwrap_err();
1258 assert!(matches!(
1259 error,
1260 SearchArtifactError::CorruptPrimaryVectors { .. }
1261 ));
1262 assert_eq!(std::fs::read(root.join(CURRENT_FILE)).unwrap(), b"corrupt");
1263 }
1264
1265 #[test]
1266 fn current_pointer_malformed_state_matrix_is_exact_and_non_mutating() {
1267 let key = key();
1268 let cases: Vec<Vec<u8>> = vec![
1269 b"not-json".to_vec(),
1270 br#"{}"#.to_vec(),
1271 br#"{"version":"../escape"}"#.to_vec(),
1272 br#"{"version":"version-bad/slash"}"#.to_vec(),
1273 vec![b'x'; MAX_CURRENT_BYTES + 1],
1274 ];
1275 for bytes in cases {
1276 let dir = TempDir::new().unwrap();
1277 let root = key.artifact_root(dir.path());
1278 std::fs::create_dir_all(&root).unwrap();
1279 let pointer = root.join(CURRENT_FILE);
1280 std::fs::write(&pointer, &bytes).unwrap();
1281 let result = current_search_artifact(dir.path(), &key);
1282 assert!(matches!(
1283 result,
1284 Err(SearchArtifactError::CorruptManifest { .. })
1285 | Err(SearchArtifactError::ResourceExhausted {
1286 resource: "current_pointer_bytes",
1287 ..
1288 })
1289 ));
1290 assert_eq!(std::fs::read(&pointer).unwrap(), bytes);
1291 assert!(!root.join(VERSIONS_DIR).exists());
1292 }
1293 }
1294
1295 #[test]
1296 fn rebuildability_primary_wrapping_and_owned_name_matrices_are_total() {
1297 let path = PathBuf::from("artifact");
1298 let rebuildable = [
1299 SearchArtifactError::Missing { path: path.clone() },
1300 SearchArtifactError::CorruptManifest {
1301 path: path.clone(),
1302 reason: "bad".into(),
1303 },
1304 SearchArtifactError::CorruptDerivedIndex {
1305 path: path.clone(),
1306 reason: "bad".into(),
1307 },
1308 SearchArtifactError::IncompatibleManifest {
1309 path: path.clone(),
1310 found: 2,
1311 supported: 1,
1312 },
1313 SearchArtifactError::Stale {
1314 reason: "old".into(),
1315 },
1316 SearchArtifactError::ResourceExhausted {
1317 resource: "manifest_bytes",
1318 limit: 1,
1319 },
1320 SearchArtifactError::ResourceExhausted {
1321 resource: "current_pointer_bytes",
1322 limit: 1,
1323 },
1324 ];
1325 for error in &rebuildable {
1326 assert!(rebuildable_metadata(error), "{error}");
1327 }
1328 for error in [
1329 SearchArtifactError::Cancelled,
1330 SearchArtifactError::ConcurrentMutation,
1331 SearchArtifactError::ResourceExhausted {
1332 resource: "other",
1333 limit: 1,
1334 },
1335 SearchArtifactError::Build("bad".into()),
1336 ] {
1337 assert!(!rebuildable_metadata(&error), "{error}");
1338 }
1339
1340 let primary = SearchArtifactError::CorruptPrimaryVectors {
1341 path: path.clone(),
1342 reason: "primary".into(),
1343 };
1344 assert!(matches!(
1345 primary_vector_error(path.clone(), primary),
1346 SearchArtifactError::CorruptPrimaryVectors { reason, .. } if reason == "primary"
1347 ));
1348 assert!(matches!(
1349 primary_vector_error(path.clone(), SearchArtifactError::Cancelled),
1350 SearchArtifactError::CorruptPrimaryVectors { path: actual, reason }
1351 if actual == path && reason.contains("cancelled")
1352 ));
1353
1354 for valid in ["build-A1", "version-z9"] {
1355 let prefix = if valid.starts_with("build") {
1356 "build-"
1357 } else {
1358 "version-"
1359 };
1360 assert!(valid_owned_name(valid, prefix));
1361 }
1362 for invalid in ["build-", "build-a/b", "build-a_b", "other-a"] {
1363 assert!(!valid_owned_name(invalid, "build-"));
1364 }
1365 for (name, expected) in [
1366 ("current.json.A1.tmp", true),
1367 ("current.json..tmp", false),
1368 ("current.json.a_b.tmp", false),
1369 ("current.json.a", false),
1370 ] {
1371 assert_eq!(valid_pointer_temp(name), expected);
1372 }
1373 }
1374
1375 #[cfg(unix)]
1376 #[test]
1377 fn syncing_build_with_symlink_fails_without_following_or_mutating_target() {
1378 use std::os::unix::fs::symlink;
1379
1380 let build = TempDir::new().unwrap();
1381 let external = TempDir::new().unwrap();
1382 let target = external.path().join("secret");
1383 std::fs::write(&target, b"caller bytes").unwrap();
1384 let link = build.path().join("linked");
1385 symlink(&target, &link).unwrap();
1386 assert!(matches!(
1387 sync_tree(build.path()),
1388 Err(SearchArtifactError::Build(_))
1389 ));
1390 assert_eq!(std::fs::read(&target).unwrap(), b"caller bytes");
1391 assert!(link.symlink_metadata().unwrap().file_type().is_symlink());
1392 }
1393
1394 #[test]
1395 fn wave10_writer_timeout_and_cleanup_bounds_are_fail_closed() {
1396 let project = TempDir::new().unwrap();
1397 let root = key().artifact_root(project.path());
1398 std::fs::create_dir_all(&root).unwrap();
1399 let first =
1400 SearchWriterLock::acquire(&root, SearchCoordinationLimits::default(), &mut || Ok(()))
1401 .unwrap();
1402 let zero_wait = SearchCoordinationLimits {
1403 lock_timeout: Duration::ZERO,
1404 lock_poll_interval: Duration::ZERO,
1405 ..SearchCoordinationLimits::default()
1406 };
1407 assert!(matches!(
1408 SearchWriterLock::acquire(&root, zero_wait, &mut || Ok(())),
1409 Err(SearchArtifactError::Lock { .. })
1410 ));
1411 drop(first);
1412
1413 let cleanup = project.path().join("indexes/search/owned");
1414 std::fs::create_dir_all(&cleanup).unwrap();
1415 std::fs::write(cleanup.join("caller"), b"preserve").unwrap();
1416 assert!(matches!(
1417 cleanup_abandoned_search_builds(project.path(), 0),
1418 Err(SearchArtifactError::ResourceExhausted {
1419 resource: "cleanup_entries",
1420 ..
1421 })
1422 ));
1423 assert_eq!(std::fs::read(cleanup.join("caller")).unwrap(), b"preserve");
1424 }
1425}