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 b"{\"catalog_version\":1,\"default\":\"missing\",\"spaces\":[]}",
904 )
905 .unwrap();
906 assert!(matches!(
907 read_embedding_space_catalog(
908 project.path(),
909 EmbeddingSpaceCatalogLimits::default(),
910 || Ok(())
911 ),
912 Err(SearchArtifactError::CorruptManifest { .. })
913 ));
914 std::fs::write(
915 &path,
916 format!(
917 "{{\"catalog_version\":1,\"default\":null,\"spaces\":[{{\"display_name\":\"duplicate\",\"compatibility_id\":\"{}\"}},{{\"display_name\":\"duplicate\",\"compatibility_id\":\"{}\"}}]}}",
918 id(1).to_hex(),
919 id(2).to_hex()
920 ),
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 std::fs::write(
932 &path,
933 b"{ \"catalog_version\": 1, \"default\": null, \"spaces\": [] }",
934 )
935 .unwrap();
936 assert!(matches!(
937 read_embedding_space_catalog(
938 project.path(),
939 EmbeddingSpaceCatalogLimits::default(),
940 || Ok(())
941 ),
942 Err(SearchArtifactError::CorruptManifest { .. })
943 ));
944 assert!(matches!(
945 read_embedding_space_catalog(
946 project.path(),
947 EmbeddingSpaceCatalogLimits {
948 metadata_bytes: 8,
949 ..EmbeddingSpaceCatalogLimits::default()
950 },
951 || Ok(())
952 ),
953 Err(SearchArtifactError::ResourceExhausted { .. })
954 ));
955 }
956
957 #[test]
958 fn bind_existing_requires_owned_live_lineage_and_persists_exact_alias() {
959 let project = tempfile::tempdir().unwrap();
960 let compatibility_id = id(7);
961 assert!(matches!(
962 bind_existing_embedding_space_catalog_entry(
963 project.path(),
964 "semantic",
965 compatibility_id,
966 false,
967 EmbeddingSpaceCatalogLimits::default(),
968 || Ok(()),
969 ),
970 Err(SearchArtifactError::InvalidSelector {
971 field: "embedding compatibility identity",
972 ..
973 })
974 ));
975
976 let lineage = project
977 .path()
978 .join(EMBEDDINGS_DIR)
979 .join("spaces")
980 .join(compatibility_id.to_hex());
981 std::fs::create_dir_all(&lineage).unwrap();
982 std::fs::write(lineage.join("space.json"), b"descriptor-presence").unwrap();
983 let marker =
984 crate::embedding_publication::deletion_marker(project.path(), compatibility_id);
985 std::fs::write(&marker, b"deleting").unwrap();
986 assert!(
987 bind_existing_embedding_space_catalog_entry(
988 project.path(),
989 "semantic",
990 compatibility_id,
991 false,
992 EmbeddingSpaceCatalogLimits::default(),
993 || Ok(()),
994 )
995 .is_err()
996 );
997 assert!(!catalog_path(project.path()).exists());
998
999 std::fs::remove_file(marker).unwrap();
1000 let catalog = bind_existing_embedding_space_catalog_entry(
1001 project.path(),
1002 "semantic",
1003 compatibility_id,
1004 false,
1005 EmbeddingSpaceCatalogLimits::default(),
1006 || Ok(()),
1007 )
1008 .unwrap();
1009 assert_eq!(catalog.len(), 1);
1010 assert_eq!(
1011 catalog.get(&EmbeddingDisplayName::new("semantic").unwrap()),
1012 Some(compatibility_id)
1013 );
1014 assert_eq!(read(project.path()), catalog);
1015 }
1016
1017 #[test]
1018 fn catalog_limits_reject_before_creating_state() {
1019 for limits in [
1020 EmbeddingSpaceCatalogLimits {
1021 metadata_bytes: 0,
1022 ..EmbeddingSpaceCatalogLimits::default()
1023 },
1024 EmbeddingSpaceCatalogLimits {
1025 entries: 0,
1026 ..EmbeddingSpaceCatalogLimits::default()
1027 },
1028 ] {
1029 let project = tempfile::tempdir().unwrap();
1030 assert!(matches!(
1031 read_embedding_space_catalog(project.path(), limits, || Ok(())),
1032 Err(SearchArtifactError::InvalidSelector {
1033 field: "embedding space catalog limits",
1034 ..
1035 })
1036 ));
1037 assert!(!project.path().join(EMBEDDINGS_DIR).exists());
1038 }
1039 }
1040
1041 #[cfg(not(unix))]
1042 #[test]
1043 fn directory_sync_is_a_supported_noop() {
1044 let directory = tempfile::tempdir().unwrap();
1045 sync_directory(directory.path()).unwrap();
1046 }
1047
1048 #[test]
1049 fn wave10_catalog_lock_and_owned_path_corruption_are_structured() {
1050 use std::time::Duration;
1051
1052 let project = tempfile::tempdir().unwrap();
1053 let embeddings = project.path().join(EMBEDDINGS_DIR);
1054 std::fs::create_dir(&embeddings).unwrap();
1055 let first = CatalogWriterLock::acquire(
1056 &embeddings,
1057 SearchCoordinationLimits::default(),
1058 &mut || Ok(()),
1059 )
1060 .unwrap();
1061 let zero_wait = SearchCoordinationLimits {
1062 lock_timeout: Duration::ZERO,
1063 lock_poll_interval: Duration::ZERO,
1064 ..SearchCoordinationLimits::default()
1065 };
1066 assert!(matches!(
1067 CatalogWriterLock::acquire(&embeddings, zero_wait, &mut || Ok(())),
1068 Err(SearchArtifactError::Lock { .. })
1069 ));
1070 drop(first);
1071
1072 let owned = project.path().join("owned");
1073 std::fs::write(&owned, b"caller").unwrap();
1074 assert!(matches!(
1075 ensure_owned_directory(&owned),
1076 Err(SearchArtifactError::CorruptManifest { .. })
1077 ));
1078 assert!(matches!(
1079 ensure_regular_file(project.path()),
1080 Err(SearchArtifactError::CorruptManifest { .. })
1081 ));
1082 }
1083}