Skip to main content

graphforge_storage/
embedding_catalog.rs

1//! Durable caller-facing names for compatibility-addressed embedding spaces.
2
3use 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
15/// Embedding-space catalog schema implemented by this release.
16pub const EMBEDDING_SPACE_CATALOG_VERSION: u32 = 1;
17/// Maximum accepted catalog bytes by default.
18pub const MAX_EMBEDDING_SPACE_CATALOG_BYTES: usize = 64 * 1024;
19/// Maximum accepted named spaces by default.
20pub 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/// Resource and coordination bounds for embedding-space catalog access.
27#[derive(Clone, Copy, Debug)]
28pub struct EmbeddingSpaceCatalogLimits {
29    /// Maximum canonical JSON bytes.
30    pub metadata_bytes: usize,
31    /// Maximum distinct display names.
32    pub entries: usize,
33    /// Catalog-writer lock bounds.
34    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/// One deterministic caller-facing name to compatibility-identity binding.
48#[derive(Clone, Debug, PartialEq, Eq)]
49pub struct EmbeddingSpaceCatalogEntry {
50    display_name: EmbeddingDisplayName,
51    compatibility_id: EmbeddingCompatibilityId,
52}
53
54impl EmbeddingSpaceCatalogEntry {
55    /// Normalized caller-facing name, never a path component.
56    #[must_use]
57    pub const fn display_name(&self) -> &EmbeddingDisplayName {
58        &self.display_name
59    }
60
61    /// Exact compatibility lineage selected by this name.
62    #[must_use]
63    pub const fn compatibility_id(&self) -> EmbeddingCompatibilityId {
64        self.compatibility_id
65    }
66}
67
68/// Fully validated catalog returned in deterministic display-name order.
69#[derive(Clone, Debug, Default, PartialEq, Eq)]
70pub struct EmbeddingSpaceCatalog {
71    spaces: BTreeMap<EmbeddingDisplayName, EmbeddingCompatibilityId>,
72    default: Option<EmbeddingDisplayName>,
73}
74
75impl EmbeddingSpaceCatalog {
76    /// Number of distinct named spaces.
77    #[must_use]
78    pub fn len(&self) -> usize {
79        self.spaces.len()
80    }
81
82    /// Whether no display names are configured.
83    #[must_use]
84    pub fn is_empty(&self) -> bool {
85        self.spaces.is_empty()
86    }
87
88    /// Exact compatibility identity bound to one normalized name.
89    #[must_use]
90    pub fn get(&self, display_name: &EmbeddingDisplayName) -> Option<EmbeddingCompatibilityId> {
91        self.spaces.get(display_name).copied()
92    }
93
94    /// Optional configured default and its exact compatibility identity.
95    #[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    /// All bindings in deterministic normalized display-name order.
107    #[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/// One explicit mutation of the durable display-name catalog.
153#[derive(Clone, Copy, Debug)]
154pub enum EmbeddingSpaceCatalogUpdate<'a> {
155    /// Bind a normalized display name to one exact compatibility identity.
156    Bind {
157        /// Caller-facing name.
158        display_name: &'a str,
159        /// Exact compatibility lineage.
160        compatibility_id: EmbeddingCompatibilityId,
161        /// Permit replacing a different existing identity when `true`.
162        replace: bool,
163    },
164    /// Remove one name. Missing names are idempotent.
165    Remove {
166        /// Caller-facing name.
167        display_name: &'a str,
168    },
169    /// Select an existing name as default, or clear the default with `None`.
170    SetDefault {
171        /// Existing caller-facing name, or `None`.
172        display_name: Option<&'a str>,
173    },
174}
175
176/// Read and validate the durable embedding-space catalog.
177///
178/// A missing file returns an empty catalog without creating directories.
179///
180/// # Errors
181/// Returns structured cancellation, limit, corruption, incompatibility, or I/O errors.
182pub 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
195/// Apply one catalog mutation under a bounded cross-process writer lock.
196///
197/// Exact-idempotent binds and missing removals do not rewrite durable bytes.
198/// Publication uses a synchronized temporary file and atomic replacement.
199///
200/// # Errors
201/// Returns structured validation, cancellation, lock, limit, corruption, or I/O errors.
202pub 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
215/// Remove one display name and report the mutation outcome from inside the
216/// catalog writer lock.
217///
218/// Missing names are idempotent and return `false`. If the removed name was
219/// the selected default, the default is cleared in the same atomic update.
220///
221/// # Errors
222/// Returns structured validation, cancellation, lock, limit, corruption, or I/O errors.
223pub 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
241/// Bind one display name only while its compatibility lineage is present and
242/// not being deleted.
243///
244/// The existence check occurs inside the catalog writer lock so a concurrent
245/// deletion cannot leave a newly dangling alias.
246///
247/// # Errors
248/// Returns structured validation, cancellation, lock, corruption, or I/O errors.
249pub 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
286/// Remove every display name that targets one compatibility identity.
287///
288/// The configured default is cleared when it names any removed alias.
289///
290/// # Errors
291/// Returns structured cancellation, lock, limit, corruption, or I/O errors.
292pub 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}