1use std::collections::{BTreeMap, BTreeSet};
4use std::fs::{File, OpenOptions};
5use std::io::Write;
6use std::path::{Path, PathBuf};
7use std::time::Instant;
8
9use serde::{Deserialize, Serialize};
10
11use crate::{
12 EmbeddingCompatibilityId, EmbeddingDisplayName, SearchArtifactError, SearchCoordinationLimits,
13};
14
15pub const EMBEDDING_SPACE_CATALOG_VERSION: u32 = 1;
17pub const MAX_EMBEDDING_SPACE_CATALOG_BYTES: usize = 64 * 1024;
19pub const MAX_EMBEDDING_SPACE_CATALOG_ENTRIES: usize = 1_024;
21
22const EMBEDDINGS_DIR: &str = "embeddings";
23const CATALOG_FILE: &str = "catalog.json";
24const CATALOG_LOCK: &str = ".catalog.lock";
25
26#[derive(Clone, Copy, Debug)]
28pub struct EmbeddingSpaceCatalogLimits {
29 pub metadata_bytes: usize,
31 pub entries: usize,
33 pub coordination: SearchCoordinationLimits,
35}
36
37impl Default for EmbeddingSpaceCatalogLimits {
38 fn default() -> Self {
39 Self {
40 metadata_bytes: MAX_EMBEDDING_SPACE_CATALOG_BYTES,
41 entries: MAX_EMBEDDING_SPACE_CATALOG_ENTRIES,
42 coordination: SearchCoordinationLimits::default(),
43 }
44 }
45}
46
47#[derive(Clone, Debug, PartialEq, Eq)]
49pub struct EmbeddingSpaceCatalogEntry {
50 display_name: EmbeddingDisplayName,
51 compatibility_id: EmbeddingCompatibilityId,
52}
53
54impl EmbeddingSpaceCatalogEntry {
55 #[must_use]
57 pub const fn display_name(&self) -> &EmbeddingDisplayName {
58 &self.display_name
59 }
60
61 #[must_use]
63 pub const fn compatibility_id(&self) -> EmbeddingCompatibilityId {
64 self.compatibility_id
65 }
66}
67
68#[derive(Clone, Debug, Default, PartialEq, Eq)]
70pub struct EmbeddingSpaceCatalog {
71 spaces: BTreeMap<EmbeddingDisplayName, EmbeddingCompatibilityId>,
72 default: Option<EmbeddingDisplayName>,
73}
74
75impl EmbeddingSpaceCatalog {
76 #[must_use]
78 pub fn len(&self) -> usize {
79 self.spaces.len()
80 }
81
82 #[must_use]
84 pub fn is_empty(&self) -> bool {
85 self.spaces.is_empty()
86 }
87
88 #[must_use]
90 pub fn get(&self, display_name: &EmbeddingDisplayName) -> Option<EmbeddingCompatibilityId> {
91 self.spaces.get(display_name).copied()
92 }
93
94 #[must_use]
96 pub fn selected_default(&self) -> Option<EmbeddingSpaceCatalogEntry> {
97 self.default.as_ref().map(|display_name| {
98 let compatibility_id = self.spaces[display_name];
99 EmbeddingSpaceCatalogEntry {
100 display_name: display_name.clone(),
101 compatibility_id,
102 }
103 })
104 }
105
106 #[must_use]
108 pub fn entries(&self) -> Vec<EmbeddingSpaceCatalogEntry> {
109 self.spaces
110 .iter()
111 .map(
112 |(display_name, compatibility_id)| EmbeddingSpaceCatalogEntry {
113 display_name: display_name.clone(),
114 compatibility_id: *compatibility_id,
115 },
116 )
117 .collect()
118 }
119
120 fn to_canonical_json(
121 &self,
122 limits: EmbeddingSpaceCatalogLimits,
123 ) -> Result<Vec<u8>, SearchArtifactError> {
124 validate_limits(limits)?;
125 if self.spaces.len() > limits.entries {
126 return Err(exhausted("embedding_space_catalog_entries", limits.entries));
127 }
128 let spaces = self
129 .spaces
130 .iter()
131 .map(|(display_name, compatibility_id)| WireEntry {
132 display_name: display_name.as_str(),
133 compatibility_id: compatibility_id.to_hex(),
134 })
135 .collect::<Vec<_>>();
136 let bytes = serde_json::to_vec(&WireCatalog {
137 catalog_version: EMBEDDING_SPACE_CATALOG_VERSION,
138 default: self.default.as_ref().map(EmbeddingDisplayName::as_str),
139 spaces,
140 })
141 .map_err(|error| invalid("embedding space catalog", error.to_string()))?;
142 if bytes.len() > limits.metadata_bytes {
143 return Err(exhausted(
144 "embedding_space_catalog_bytes",
145 limits.metadata_bytes,
146 ));
147 }
148 Ok(bytes)
149 }
150}
151
152#[derive(Clone, Copy, Debug)]
154pub enum EmbeddingSpaceCatalogUpdate<'a> {
155 Bind {
157 display_name: &'a str,
159 compatibility_id: EmbeddingCompatibilityId,
161 replace: bool,
163 },
164 Remove {
166 display_name: &'a str,
168 },
169 SetDefault {
171 display_name: Option<&'a str>,
173 },
174}
175
176pub fn read_embedding_space_catalog<C>(
183 project_dir: &Path,
184 limits: EmbeddingSpaceCatalogLimits,
185 mut checkpoint: C,
186) -> Result<EmbeddingSpaceCatalog, SearchArtifactError>
187where
188 C: FnMut() -> Result<(), SearchArtifactError>,
189{
190 validate_limits(limits)?;
191 checkpoint()?;
192 read_catalog_file(&catalog_path(project_dir), limits)
193}
194
195pub fn update_embedding_space_catalog<C>(
203 project_dir: &Path,
204 update: EmbeddingSpaceCatalogUpdate<'_>,
205 limits: EmbeddingSpaceCatalogLimits,
206 checkpoint: C,
207) -> Result<EmbeddingSpaceCatalog, SearchArtifactError>
208where
209 C: FnMut() -> Result<(), SearchArtifactError>,
210{
211 mutate_embedding_space_catalog(project_dir, update, limits, checkpoint)
212 .map(|(catalog, _changed)| catalog)
213}
214
215pub fn remove_embedding_space_catalog_entry<C>(
224 project_dir: &Path,
225 display_name: &str,
226 limits: EmbeddingSpaceCatalogLimits,
227 checkpoint: C,
228) -> Result<bool, SearchArtifactError>
229where
230 C: FnMut() -> Result<(), SearchArtifactError>,
231{
232 mutate_embedding_space_catalog(
233 project_dir,
234 EmbeddingSpaceCatalogUpdate::Remove { display_name },
235 limits,
236 checkpoint,
237 )
238 .map(|(_catalog, changed)| changed)
239}
240
241pub fn bind_existing_embedding_space_catalog_entry<C>(
250 project_dir: &Path,
251 display_name: &str,
252 compatibility_id: EmbeddingCompatibilityId,
253 replace: bool,
254 limits: EmbeddingSpaceCatalogLimits,
255 mut checkpoint: C,
256) -> Result<EmbeddingSpaceCatalog, SearchArtifactError>
257where
258 C: FnMut() -> Result<(), SearchArtifactError>,
259{
260 validate_limits(limits)?;
261 checkpoint()?;
262 let embeddings = project_dir.join(EMBEDDINGS_DIR);
263 ensure_owned_directory(&embeddings)?;
264 let _lock = CatalogWriterLock::acquire(&embeddings, limits.coordination, &mut checkpoint)?;
265 checkpoint()?;
266 ensure_bindable_lineage(project_dir, compatibility_id)?;
267
268 let path = embeddings.join(CATALOG_FILE);
269 let mut catalog = read_catalog_file(&path, limits)?;
270 let changed = apply_update(
271 &mut catalog,
272 EmbeddingSpaceCatalogUpdate::Bind {
273 display_name,
274 compatibility_id,
275 replace,
276 },
277 limits,
278 )?;
279 if changed {
280 checkpoint()?;
281 persist_synced_file(&path, &catalog.to_canonical_json(limits)?)?;
282 }
283 Ok(catalog)
284}
285
286pub fn remove_embedding_space_catalog_identity<C>(
293 project_dir: &Path,
294 compatibility_id: EmbeddingCompatibilityId,
295 limits: EmbeddingSpaceCatalogLimits,
296 mut checkpoint: C,
297) -> Result<usize, SearchArtifactError>
298where
299 C: FnMut() -> Result<(), SearchArtifactError>,
300{
301 validate_limits(limits)?;
302 checkpoint()?;
303 let embeddings = project_dir.join(EMBEDDINGS_DIR);
304 ensure_owned_directory(&embeddings)?;
305 let _lock = CatalogWriterLock::acquire(&embeddings, limits.coordination, &mut checkpoint)?;
306 checkpoint()?;
307
308 let path = embeddings.join(CATALOG_FILE);
309 let mut catalog = read_catalog_file(&path, limits)?;
310 let removed_names = catalog
311 .spaces
312 .iter()
313 .filter_map(|(name, current)| (*current == compatibility_id).then_some(name.clone()))
314 .collect::<Vec<_>>();
315 if removed_names.is_empty() {
316 return Ok(0);
317 }
318 for name in &removed_names {
319 catalog.spaces.remove(name);
320 }
321 if catalog
322 .default
323 .as_ref()
324 .is_some_and(|name| removed_names.contains(name))
325 {
326 catalog.default = None;
327 }
328 checkpoint()?;
329 persist_synced_file(&path, &catalog.to_canonical_json(limits)?)?;
330 Ok(removed_names.len())
331}
332
333fn mutate_embedding_space_catalog<C>(
334 project_dir: &Path,
335 update: EmbeddingSpaceCatalogUpdate<'_>,
336 limits: EmbeddingSpaceCatalogLimits,
337 mut checkpoint: C,
338) -> Result<(EmbeddingSpaceCatalog, bool), SearchArtifactError>
339where
340 C: FnMut() -> Result<(), SearchArtifactError>,
341{
342 validate_limits(limits)?;
343 checkpoint()?;
344 let embeddings = project_dir.join(EMBEDDINGS_DIR);
345 ensure_owned_directory(&embeddings)?;
346 let _lock = CatalogWriterLock::acquire(&embeddings, limits.coordination, &mut checkpoint)?;
347 checkpoint()?;
348
349 let path = embeddings.join(CATALOG_FILE);
350 let mut catalog = read_catalog_file(&path, limits)?;
351 let changed = apply_update(&mut catalog, update, limits)?;
352 if changed {
353 checkpoint()?;
354 let bytes = catalog.to_canonical_json(limits)?;
355 checkpoint()?;
356 persist_synced_file(&path, &bytes)?;
357 }
358 Ok((catalog, changed))
359}
360
361fn apply_update(
362 catalog: &mut EmbeddingSpaceCatalog,
363 update: EmbeddingSpaceCatalogUpdate<'_>,
364 limits: EmbeddingSpaceCatalogLimits,
365) -> Result<bool, SearchArtifactError> {
366 match update {
367 EmbeddingSpaceCatalogUpdate::Bind {
368 display_name,
369 compatibility_id,
370 replace,
371 } => {
372 let display_name = EmbeddingDisplayName::new(display_name)?;
373 match catalog.spaces.get(&display_name).copied() {
374 Some(current) if current == compatibility_id => Ok(false),
375 Some(_) if !replace => Err(invalid(
376 "embedding display name",
377 "is already bound to a different compatibility identity",
378 )),
379 _ => {
380 if !catalog.spaces.contains_key(&display_name)
381 && catalog.spaces.len() >= limits.entries
382 {
383 return Err(exhausted("embedding_space_catalog_entries", limits.entries));
384 }
385 catalog.spaces.insert(display_name, compatibility_id);
386 Ok(true)
387 }
388 }
389 }
390 EmbeddingSpaceCatalogUpdate::Remove { display_name } => {
391 let display_name = EmbeddingDisplayName::new(display_name)?;
392 let removed = catalog.spaces.remove(&display_name).is_some();
393 if catalog.default.as_ref() == Some(&display_name) {
394 catalog.default = None;
395 }
396 Ok(removed)
397 }
398 EmbeddingSpaceCatalogUpdate::SetDefault { display_name } => {
399 let display_name = display_name.map(EmbeddingDisplayName::new).transpose()?;
400 if display_name
401 .as_ref()
402 .is_some_and(|name| !catalog.spaces.contains_key(name))
403 {
404 return Err(invalid(
405 "embedding default space",
406 "must name an existing catalog entry",
407 ));
408 }
409 if catalog.default == display_name {
410 Ok(false)
411 } else {
412 catalog.default = display_name;
413 Ok(true)
414 }
415 }
416 }
417}
418
419fn read_catalog_file(
420 path: &Path,
421 limits: EmbeddingSpaceCatalogLimits,
422) -> Result<EmbeddingSpaceCatalog, SearchArtifactError> {
423 if !path_exists(path)? {
424 return Ok(EmbeddingSpaceCatalog::default());
425 }
426 ensure_regular_file(path)?;
427 let metadata = std::fs::metadata(path)
428 .map_err(|source| io("inspect embedding space catalog", path, source))?;
429 if metadata.len() > limits.metadata_bytes as u64 {
430 return Err(exhausted(
431 "embedding_space_catalog_bytes",
432 limits.metadata_bytes,
433 ));
434 }
435 let bytes =
436 std::fs::read(path).map_err(|source| io("read embedding space catalog", path, source))?;
437 let raw: RawCatalog =
438 serde_json::from_slice(&bytes).map_err(|error| corrupt(path, error.to_string()))?;
439 if raw.catalog_version != u64::from(EMBEDDING_SPACE_CATALOG_VERSION) {
440 return Err(SearchArtifactError::IncompatibleManifest {
441 path: path.to_path_buf(),
442 found: raw.catalog_version,
443 supported: EMBEDDING_SPACE_CATALOG_VERSION,
444 });
445 }
446 if raw.spaces.len() > limits.entries {
447 return Err(exhausted("embedding_space_catalog_entries", limits.entries));
448 }
449
450 let mut spaces = BTreeMap::new();
451 let mut seen = BTreeSet::new();
452 for entry in raw.spaces {
453 let display_name = EmbeddingDisplayName::new(&entry.display_name)
454 .map_err(|error| corrupt(path, error.to_string()))?;
455 if !seen.insert(display_name.clone()) {
456 return Err(corrupt(path, "duplicate embedding display name"));
457 }
458 let compatibility_id = EmbeddingCompatibilityId::from_hex(&entry.compatibility_id)
459 .map_err(|error| corrupt(path, error.to_string()))?;
460 spaces.insert(display_name, compatibility_id);
461 }
462 let default = raw
463 .default
464 .map(|value| EmbeddingDisplayName::new(&value))
465 .transpose()
466 .map_err(|error| corrupt(path, error.to_string()))?;
467 if default
468 .as_ref()
469 .is_some_and(|display_name| !spaces.contains_key(display_name))
470 {
471 return Err(corrupt(path, "default does not name a catalog entry"));
472 }
473 let catalog = EmbeddingSpaceCatalog { spaces, default };
474 let canonical = catalog
475 .to_canonical_json(limits)
476 .map_err(|error| corrupt(path, error.to_string()))?;
477 if canonical != bytes {
478 return Err(corrupt(path, "catalog bytes are not exact canonical JSON"));
479 }
480 Ok(catalog)
481}
482
483#[derive(Serialize)]
484struct WireCatalog<'a> {
485 catalog_version: u32,
486 default: Option<&'a str>,
487 spaces: Vec<WireEntry<'a>>,
488}
489
490#[derive(Serialize)]
491struct WireEntry<'a> {
492 display_name: &'a str,
493 compatibility_id: String,
494}
495
496#[derive(Deserialize)]
497#[serde(deny_unknown_fields)]
498struct RawCatalog {
499 catalog_version: u64,
500 default: Option<String>,
501 spaces: Vec<RawEntry>,
502}
503
504#[derive(Deserialize)]
505#[serde(deny_unknown_fields)]
506struct RawEntry {
507 display_name: String,
508 compatibility_id: String,
509}
510
511struct CatalogWriterLock {
512 file: File,
513}
514
515impl CatalogWriterLock {
516 fn acquire<C>(
517 embeddings: &Path,
518 limits: SearchCoordinationLimits,
519 checkpoint: &mut C,
520 ) -> Result<Self, SearchArtifactError>
521 where
522 C: FnMut() -> Result<(), SearchArtifactError>,
523 {
524 let path = embeddings.join(CATALOG_LOCK);
525 let file = OpenOptions::new()
526 .read(true)
527 .write(true)
528 .create(true)
529 .truncate(false)
530 .open(&path)
531 .map_err(|source| SearchArtifactError::Lock {
532 path: path.clone(),
533 reason: source.to_string(),
534 })?;
535 let started = Instant::now();
536 loop {
537 match file.try_lock() {
538 Ok(()) => return Ok(Self { file }),
539 Err(std::fs::TryLockError::WouldBlock) => {
540 checkpoint()?;
541 if started.elapsed() >= limits.lock_timeout {
542 return Err(SearchArtifactError::Lock {
543 path,
544 reason: format!(
545 "timed out after {} ms",
546 limits.lock_timeout.as_millis()
547 ),
548 });
549 }
550 std::thread::sleep(limits.lock_poll_interval);
551 }
552 Err(std::fs::TryLockError::Error(source)) => {
553 return Err(SearchArtifactError::Lock {
554 path,
555 reason: source.to_string(),
556 });
557 }
558 }
559 }
560 }
561}
562
563impl Drop for CatalogWriterLock {
564 fn drop(&mut self) {
565 let _ = self.file.unlock();
566 }
567}
568
569fn validate_limits(limits: EmbeddingSpaceCatalogLimits) -> Result<(), SearchArtifactError> {
570 if limits.metadata_bytes == 0 || limits.entries == 0 {
571 Err(invalid(
572 "embedding space catalog limits",
573 "must be non-zero",
574 ))
575 } else {
576 Ok(())
577 }
578}
579
580fn catalog_path(project_dir: &Path) -> PathBuf {
581 project_dir.join(EMBEDDINGS_DIR).join(CATALOG_FILE)
582}
583
584fn ensure_bindable_lineage(
585 project_dir: &Path,
586 compatibility_id: EmbeddingCompatibilityId,
587) -> Result<(), SearchArtifactError> {
588 let marker = crate::embedding_publication::deletion_marker(project_dir, compatibility_id);
589 if path_exists(&marker)? {
590 return Err(invalid(
591 "embedding compatibility identity",
592 "deletion is in progress",
593 ));
594 }
595 let root = project_dir
596 .join(EMBEDDINGS_DIR)
597 .join("spaces")
598 .join(compatibility_id.to_hex());
599 let metadata = std::fs::symlink_metadata(&root).map_err(|source| {
600 if source.kind() == std::io::ErrorKind::NotFound {
601 invalid(
602 "embedding compatibility identity",
603 "is not a published lineage",
604 )
605 } else {
606 io("inspect embedding lineage", &root, source)
607 }
608 })?;
609 if metadata.file_type().is_symlink() || !metadata.is_dir() {
610 return Err(corrupt(
611 &root,
612 "expected an owned embedding lineage directory",
613 ));
614 }
615 ensure_regular_file(&root.join("space.json"))
616}
617
618fn ensure_owned_directory(path: &Path) -> Result<(), SearchArtifactError> {
619 if path_exists(path)? {
620 let metadata = std::fs::symlink_metadata(path)
621 .map_err(|source| io("inspect embedding catalog directory", path, source))?;
622 if metadata.file_type().is_symlink() || !metadata.is_dir() {
623 return Err(corrupt(path, "expected an owned directory"));
624 }
625 return Ok(());
626 }
627 std::fs::create_dir_all(path)
628 .map_err(|source| io("create embedding catalog directory", path, source))?;
629 sync_directory(path.parent().unwrap_or(path))
630}
631
632fn ensure_regular_file(path: &Path) -> Result<(), SearchArtifactError> {
633 let metadata = std::fs::symlink_metadata(path)
634 .map_err(|source| io("inspect embedding space catalog", path, source))?;
635 if metadata.file_type().is_symlink() || !metadata.is_file() {
636 return Err(corrupt(path, "expected a regular file"));
637 }
638 Ok(())
639}
640
641fn persist_synced_file(path: &Path, bytes: &[u8]) -> Result<(), SearchArtifactError> {
642 let parent = path
643 .parent()
644 .ok_or_else(|| invalid("embedding space catalog", "path has no parent"))?;
645 let mut temp = tempfile::Builder::new()
646 .prefix(".catalog.json.")
647 .suffix(".tmp")
648 .tempfile_in(parent)
649 .map_err(|source| io("create embedding catalog temp", path, source))?;
650 temp.write_all(bytes)
651 .map_err(|source| io("write embedding catalog temp", path, source))?;
652 temp.as_file()
653 .sync_all()
654 .map_err(|source| io("sync embedding catalog temp", path, source))?;
655 temp.persist(path)
656 .map_err(|error| io("publish embedding space catalog", path, error.error))?;
657 sync_directory(parent)
658}
659
660#[cfg(unix)]
661fn sync_directory(path: &Path) -> Result<(), SearchArtifactError> {
662 File::open(path)
663 .and_then(|directory| directory.sync_all())
664 .map_err(|source| io("sync embedding catalog directory", path, source))
665}
666
667#[cfg(not(unix))]
668fn sync_directory(_path: &Path) -> Result<(), SearchArtifactError> {
669 Ok(())
670}
671
672fn path_exists(path: &Path) -> Result<bool, SearchArtifactError> {
673 match std::fs::symlink_metadata(path) {
674 Ok(_) => Ok(true),
675 Err(source) if source.kind() == std::io::ErrorKind::NotFound => Ok(false),
676 Err(source) => Err(io("inspect embedding catalog path", path, source)),
677 }
678}
679
680fn invalid(field: &'static str, reason: impl Into<String>) -> SearchArtifactError {
681 SearchArtifactError::InvalidSelector {
682 field,
683 reason: reason.into(),
684 }
685}
686
687fn corrupt(path: &Path, reason: impl Into<String>) -> SearchArtifactError {
688 SearchArtifactError::CorruptManifest {
689 path: path.to_path_buf(),
690 reason: reason.into(),
691 }
692}
693
694fn exhausted(resource: &'static str, limit: usize) -> SearchArtifactError {
695 SearchArtifactError::ResourceExhausted {
696 resource,
697 limit: limit as u64,
698 }
699}
700
701fn io(operation: &'static str, path: &Path, source: std::io::Error) -> SearchArtifactError {
702 SearchArtifactError::Io {
703 operation,
704 path: path.to_path_buf(),
705 source,
706 }
707}
708
709#[cfg(test)]
710mod tests {
711 use std::cell::Cell;
712
713 use super::*;
714
715 fn id(value: u8) -> EmbeddingCompatibilityId {
716 EmbeddingCompatibilityId::from_hex(&format!("{value:02x}").repeat(32)).unwrap()
717 }
718
719 fn read(project: &Path) -> EmbeddingSpaceCatalog {
720 read_embedding_space_catalog(project, EmbeddingSpaceCatalogLimits::default(), || Ok(()))
721 .unwrap()
722 }
723
724 fn update(project: &Path, mutation: EmbeddingSpaceCatalogUpdate<'_>) -> EmbeddingSpaceCatalog {
725 update_embedding_space_catalog(
726 project,
727 mutation,
728 EmbeddingSpaceCatalogLimits::default(),
729 || Ok(()),
730 )
731 .unwrap()
732 }
733
734 #[test]
735 fn missing_catalog_is_empty_and_does_not_create_files() {
736 let project = tempfile::tempdir().unwrap();
737 assert!(read(project.path()).is_empty());
738 assert!(!project.path().join(EMBEDDINGS_DIR).exists());
739 }
740
741 #[test]
742 fn bindings_replacement_default_and_removal_are_durable_and_ordered() {
743 let project = tempfile::tempdir().unwrap();
744 update(
745 project.path(),
746 EmbeddingSpaceCatalogUpdate::Bind {
747 display_name: "zeta",
748 compatibility_id: id(1),
749 replace: false,
750 },
751 );
752 let path = catalog_path(project.path());
753 let first_bytes = std::fs::read(&path).unwrap();
754 update(
755 project.path(),
756 EmbeddingSpaceCatalogUpdate::Bind {
757 display_name: "zeta",
758 compatibility_id: id(1),
759 replace: false,
760 },
761 );
762 assert_eq!(std::fs::read(&path).unwrap(), first_bytes);
763
764 assert!(matches!(
765 update_embedding_space_catalog(
766 project.path(),
767 EmbeddingSpaceCatalogUpdate::Bind {
768 display_name: "zeta",
769 compatibility_id: id(2),
770 replace: false,
771 },
772 EmbeddingSpaceCatalogLimits::default(),
773 || Ok(())
774 ),
775 Err(SearchArtifactError::InvalidSelector { .. })
776 ));
777 update(
778 project.path(),
779 EmbeddingSpaceCatalogUpdate::Bind {
780 display_name: "zeta",
781 compatibility_id: id(2),
782 replace: true,
783 },
784 );
785 update(
786 project.path(),
787 EmbeddingSpaceCatalogUpdate::Bind {
788 display_name: "alpha",
789 compatibility_id: id(3),
790 replace: false,
791 },
792 );
793 let selected = update(
794 project.path(),
795 EmbeddingSpaceCatalogUpdate::SetDefault {
796 display_name: Some("zeta"),
797 },
798 );
799 assert_eq!(
800 selected
801 .entries()
802 .iter()
803 .map(|entry| entry.display_name().as_str())
804 .collect::<Vec<_>>(),
805 ["alpha", "zeta"]
806 );
807 assert_eq!(
808 selected.selected_default().unwrap().compatibility_id(),
809 id(2)
810 );
811 assert_eq!(read(project.path()), selected);
812
813 let removed = update(
814 project.path(),
815 EmbeddingSpaceCatalogUpdate::Remove {
816 display_name: "zeta",
817 },
818 );
819 assert!(removed.selected_default().is_none());
820 let bytes = std::fs::read_to_string(&path).unwrap();
821 assert!(!bytes.contains("zeta"));
822 assert!(!bytes.contains("credential"));
823 assert!(!bytes.contains("vector"));
824 assert!(!bytes.contains("payload"));
825 }
826
827 #[test]
828 fn invalid_default_corruption_version_limits_and_cancellation_fail_closed() {
829 let project = tempfile::tempdir().unwrap();
830 assert!(matches!(
831 update_embedding_space_catalog(
832 project.path(),
833 EmbeddingSpaceCatalogUpdate::SetDefault {
834 display_name: Some("missing"),
835 },
836 EmbeddingSpaceCatalogLimits::default(),
837 || Ok(())
838 ),
839 Err(SearchArtifactError::InvalidSelector { .. })
840 ));
841 update(
842 project.path(),
843 EmbeddingSpaceCatalogUpdate::Bind {
844 display_name: "stable",
845 compatibility_id: id(1),
846 replace: false,
847 },
848 );
849 let path = catalog_path(project.path());
850 let stable = std::fs::read(&path).unwrap();
851 let calls = Cell::new(0_u8);
852 assert!(matches!(
853 update_embedding_space_catalog(
854 project.path(),
855 EmbeddingSpaceCatalogUpdate::Bind {
856 display_name: "cancelled",
857 compatibility_id: id(2),
858 replace: false,
859 },
860 EmbeddingSpaceCatalogLimits::default(),
861 || {
862 let next = calls.get() + 1;
863 calls.set(next);
864 if next >= 4 {
865 Err(SearchArtifactError::Cancelled)
866 } else {
867 Ok(())
868 }
869 }
870 ),
871 Err(SearchArtifactError::Cancelled)
872 ));
873 assert_eq!(std::fs::read(&path).unwrap(), stable);
874
875 std::fs::write(
876 &path,
877 b"{\"catalog_version\":2,\"default\":null,\"spaces\":[]}",
878 )
879 .unwrap();
880 assert!(matches!(
881 read_embedding_space_catalog(
882 project.path(),
883 EmbeddingSpaceCatalogLimits::default(),
884 || Ok(())
885 ),
886 Err(SearchArtifactError::IncompatibleManifest { .. })
887 ));
888 std::fs::write(
889 &path,
890 b"{\"catalog_version\":1,\"default\":null,\"spaces\":[],\"extra\":1}",
891 )
892 .unwrap();
893 assert!(matches!(
894 read_embedding_space_catalog(
895 project.path(),
896 EmbeddingSpaceCatalogLimits::default(),
897 || Ok(())
898 ),
899 Err(SearchArtifactError::CorruptManifest { .. })
900 ));
901 std::fs::write(
902 &path,
903 format!(
904 "{{\"catalog_version\":1,\"default\":null,\"spaces\":[{{\"display_name\":\"duplicate\",\"compatibility_id\":\"{}\"}},{{\"display_name\":\"duplicate\",\"compatibility_id\":\"{}\"}}]}}",
905 id(1).to_hex(),
906 id(2).to_hex()
907 ),
908 )
909 .unwrap();
910 assert!(matches!(
911 read_embedding_space_catalog(
912 project.path(),
913 EmbeddingSpaceCatalogLimits::default(),
914 || Ok(())
915 ),
916 Err(SearchArtifactError::CorruptManifest { .. })
917 ));
918 std::fs::write(
919 &path,
920 b"{ \"catalog_version\": 1, \"default\": null, \"spaces\": [] }",
921 )
922 .unwrap();
923 assert!(matches!(
924 read_embedding_space_catalog(
925 project.path(),
926 EmbeddingSpaceCatalogLimits::default(),
927 || Ok(())
928 ),
929 Err(SearchArtifactError::CorruptManifest { .. })
930 ));
931 assert!(matches!(
932 read_embedding_space_catalog(
933 project.path(),
934 EmbeddingSpaceCatalogLimits {
935 metadata_bytes: 8,
936 ..EmbeddingSpaceCatalogLimits::default()
937 },
938 || Ok(())
939 ),
940 Err(SearchArtifactError::ResourceExhausted { .. })
941 ));
942 }
943
944 #[cfg(not(unix))]
945 #[test]
946 fn directory_sync_is_a_supported_noop() {
947 let directory = tempfile::tempdir().unwrap();
948 sync_directory(directory.path()).unwrap();
949 }
950}