Skip to main content

graphforge_storage/
embedding_publication.rs

1//! Crash-safe publication and reopen validation for complete embedding generations.
2
3use std::fs::{File, OpenOptions};
4use std::io::{Read, Write};
5use std::path::{Path, PathBuf};
6use std::time::Instant;
7
8use serde::Deserialize;
9use sha2::{Digest, Sha256};
10
11use crate::{
12    EmbeddingCompatibilityDescriptor, EmbeddingCompatibilityId, EmbeddingContentDigest,
13    EmbeddingGenerationId, EmbeddingGenerationManifest, EmbeddingGenerationManifestInput,
14    EmbeddingPublicationFingerprint, EmbeddingSourceState, EmbeddingSpaceCatalogLimits,
15    SearchArtifactError, SearchCoordinationLimits, StoredVector, VECTOR_DATA_FILE,
16    ValidatedEmbeddingBatch, VectorStoreLimits, read_vector_snapshot,
17    remove_embedding_space_catalog_identity, write_vector_snapshot,
18};
19
20const SPACE_FILE: &str = "space.json";
21const ACTIVE_FILE: &str = "active.json";
22const MANIFEST_FILE: &str = "manifest.json";
23const GENERATIONS_DIR: &str = "generations";
24const BUILD_PREFIX: &str = ".build-";
25const POINTER_VERSION: u32 = 1;
26const MAX_ACTIVE_BYTES: u64 = 4 * 1024;
27const MAX_DESCRIPTOR_BYTES: u64 = 64 * 1024;
28const MAX_MANIFEST_BYTES: u64 = 64 * 1024;
29const HASH_BUFFER_BYTES: usize = 64 * 1024;
30pub(crate) const EMBEDDING_DELETION_PREFIX: &str = ".deleting-";
31const EMBEDDING_WRITER_LOCK_PREFIX: &str = ".writer-";
32
33/// Complete inputs for one immutable embedding generation.
34#[derive(Clone, Copy, Debug)]
35pub struct EmbeddingPublicationRequest<'a> {
36    /// Exact versioned compatibility descriptor for this lineage.
37    pub descriptor: &'a EmbeddingCompatibilityDescriptor,
38    /// Exact committed graph inputs used by the producer.
39    pub source: EmbeddingSourceState,
40    /// Complete validated UUID/vector projection.
41    pub batch: &'a ValidatedEmbeddingBatch,
42    /// Producer completion time in UTC microseconds since Unix epoch.
43    pub generated_at_micros: i64,
44    /// Durable publication time in UTC microseconds since Unix epoch.
45    pub committed_at_micros: i64,
46}
47
48/// One fully validated active embedding generation.
49#[derive(Clone, Debug, PartialEq)]
50pub struct EmbeddingGenerationPublication {
51    /// Digest-addressed immutable generation directory.
52    pub path: PathBuf,
53    /// Exact compatibility descriptor reopened from `space.json`.
54    pub descriptor: EmbeddingCompatibilityDescriptor,
55    /// Exact completed generation manifest.
56    pub manifest: EmbeddingGenerationManifest,
57}
58
59/// Result of an atomic generation publication.
60#[derive(Clone, Debug, PartialEq)]
61pub enum EmbeddingPublicationOutcome {
62    /// The exact complete immutable generation was already present and verified.
63    Reused(EmbeddingGenerationPublication),
64    /// A new complete generation became active.
65    Published(EmbeddingGenerationPublication),
66}
67
68impl EmbeddingPublicationOutcome {
69    /// The verified active generation selected by this operation.
70    #[must_use]
71    pub const fn publication(&self) -> &EmbeddingGenerationPublication {
72        match self {
73            Self::Reused(publication) | Self::Published(publication) => publication,
74        }
75    }
76}
77
78/// Publish one complete generation without exposing private or partial data.
79///
80/// Publication is serialized per compatibility identity. The immutable vector
81/// tree and completed manifest are synchronized before the final atomic active
82/// pointer replacement. Repeating identical compatibility/source/content
83/// verifies and reuses the same generation directory.
84///
85/// # Errors
86/// Rejects incompatible dimensions or descriptors, corrupt primary data,
87/// configured resource exhaustion, cancellation, lock failure, and I/O errors.
88pub fn publish_embedding_generation<C>(
89    project_dir: &Path,
90    request: EmbeddingPublicationRequest<'_>,
91    vector_limits: VectorStoreLimits,
92    coordination: SearchCoordinationLimits,
93    mut checkpoint: C,
94) -> Result<EmbeddingPublicationOutcome, SearchArtifactError>
95where
96    C: FnMut() -> Result<(), SearchArtifactError>,
97{
98    checkpoint()?;
99    let dimension = usize::try_from(request.descriptor.dimensions()).map_err(|_| {
100        invalid(
101            "embedding dimension",
102            "cannot be represented on this platform",
103        )
104    })?;
105    if request.batch.dimension() != dimension {
106        return Err(invalid(
107            "embedding batch",
108            format!(
109                "dimension {} does not match compatibility dimension {dimension}",
110                request.batch.dimension()
111            ),
112        ));
113    }
114    let compatibility_id = request.descriptor.compatibility_id()?;
115    let root = prepare_space_root(project_dir, compatibility_id)?;
116    let _writer =
117        EmbeddingWriterLock::acquire(project_dir, compatibility_id, coordination, &mut checkpoint)?;
118    if path_exists(&deletion_marker(project_dir, compatibility_id))? {
119        return Err(invalid(
120            "embedding compatibility identity",
121            "deletion is in progress",
122        ));
123    }
124    persist_or_verify_descriptor(&root, request.descriptor, compatibility_id)?;
125
126    let provisional = EmbeddingGenerationManifest::new(EmbeddingGenerationManifestInput {
127        compatibility_id,
128        source: request.source,
129        content_digest: request.batch.content_digest(),
130        vector_count: u64::try_from(request.batch.rows().len()).map_err(|_| {
131            SearchArtifactError::ResourceExhausted {
132                resource: "embedding_rows",
133                limit: vector_limits.stored_vectors as u64,
134            }
135        })?,
136        dimension: request.descriptor.dimensions(),
137        generated_at_micros: request.generated_at_micros,
138        committed_at_micros: request.committed_at_micros,
139        publication_fingerprint: EmbeddingPublicationFingerprint::from_hex(&"0".repeat(64))?,
140    })?;
141    let generation_id = provisional.generation_id();
142    let generations = root.join(GENERATIONS_DIR);
143    ensure_owned_directory(&generations)?;
144    let generation_path = generations.join(generation_id.to_hex());
145
146    if path_exists(&generation_path)? {
147        let publication = validate_generation(
148            &root,
149            request.descriptor,
150            compatibility_id,
151            generation_id,
152            vector_limits,
153            &mut checkpoint,
154        )?;
155        checkpoint()?;
156        persist_active_pointer(&root, compatibility_id, generation_id)?;
157        return Ok(EmbeddingPublicationOutcome::Reused(publication));
158    }
159
160    let (private, manifest) = build_private_generation(
161        &root,
162        request,
163        vector_limits,
164        compatibility_id,
165        generation_id,
166        &mut checkpoint,
167    )?;
168
169    let private_path = private.keep();
170    if let Err(source) = std::fs::rename(&private_path, &generation_path) {
171        let _ = std::fs::remove_dir_all(&private_path);
172        return Err(io(
173            "publish immutable embedding generation",
174            &generation_path,
175            source,
176        ));
177    }
178    sync_directory(&generations)?;
179    checkpoint()?;
180    persist_active_pointer(&root, compatibility_id, generation_id)?;
181    sync_directory(&root)?;
182
183    Ok(EmbeddingPublicationOutcome::Published(
184        EmbeddingGenerationPublication {
185            path: generation_path,
186            descriptor: request.descriptor.clone(),
187            manifest,
188        },
189    ))
190}
191
192/// Delete one complete compatibility lineage and every catalog alias that targets it.
193///
194/// The per-lineage writer lock lives outside the deleted directory. A durable
195/// marker blocks alias binding if interruption occurs before catalog cleanup;
196/// because aliases are removed last, retry by the original display name can
197/// always complete an interrupted deletion.
198///
199/// # Errors
200/// Returns structured cancellation, lock, catalog, corruption, or I/O errors.
201pub fn delete_embedding_space_lineage<C>(
202    project_dir: &Path,
203    compatibility_id: EmbeddingCompatibilityId,
204    catalog_limits: EmbeddingSpaceCatalogLimits,
205    coordination: SearchCoordinationLimits,
206    mut checkpoint: C,
207) -> Result<bool, SearchArtifactError>
208where
209    C: FnMut() -> Result<(), SearchArtifactError>,
210{
211    checkpoint()?;
212    let embeddings = project_dir.join("embeddings");
213    if !path_exists(&embeddings)? {
214        return Ok(false);
215    }
216    ensure_existing_directory(&embeddings)?;
217    let spaces = embeddings.join("spaces");
218    let root = space_root(project_dir, compatibility_id);
219    if path_exists(&root)? {
220        ensure_space_ancestors(project_dir, &root)?;
221    }
222
223    let _writer =
224        EmbeddingWriterLock::acquire(project_dir, compatibility_id, coordination, &mut checkpoint)?;
225    let marker = deletion_marker(project_dir, compatibility_id);
226    if !path_exists(&marker)? {
227        write_synced_file(&marker, compatibility_id.to_hex().as_bytes())?;
228        sync_directory(&embeddings)?;
229    }
230    checkpoint()?;
231
232    let removed_root = if path_exists(&root)? {
233        ensure_space_ancestors(project_dir, &root)?;
234        std::fs::remove_dir_all(&root)
235            .map_err(|source| io("delete embedding space lineage", &root, source))?;
236        sync_directory(&spaces)?;
237        true
238    } else {
239        false
240    };
241    checkpoint()?;
242    std::fs::remove_file(&marker)
243        .map_err(|source| io("clear embedding deletion marker", &marker, source))?;
244    sync_directory(&embeddings)?;
245    checkpoint()?;
246    let removed_aliases = remove_embedding_space_catalog_identity(
247        project_dir,
248        compatibility_id,
249        catalog_limits,
250        &mut checkpoint,
251    )?;
252    Ok(removed_root || removed_aliases > 0)
253}
254
255fn build_private_generation<C>(
256    root: &Path,
257    request: EmbeddingPublicationRequest<'_>,
258    vector_limits: VectorStoreLimits,
259    compatibility_id: EmbeddingCompatibilityId,
260    generation_id: EmbeddingGenerationId,
261    checkpoint: &mut C,
262) -> Result<(tempfile::TempDir, EmbeddingGenerationManifest), SearchArtifactError>
263where
264    C: FnMut() -> Result<(), SearchArtifactError>,
265{
266    let private = tempfile::Builder::new()
267        .prefix(BUILD_PREFIX)
268        .tempdir_in(root)
269        .map_err(|source| io("create private embedding generation", root, source))?;
270    let rows = request
271        .batch
272        .rows()
273        .iter()
274        .map(|row| StoredVector {
275            node_uuid: row.node_uuid,
276            vector: row.vector.clone(),
277            updated_at_micros: request.generated_at_micros,
278        })
279        .collect::<Vec<_>>();
280    let vector_path = write_vector_snapshot(
281        private.path(),
282        &rows,
283        request.batch.dimension(),
284        vector_limits,
285        &mut *checkpoint,
286    )?;
287    let publication_fingerprint =
288        hash_file(&vector_path, vector_limits.parquet_bytes, &mut *checkpoint)?;
289    let manifest = EmbeddingGenerationManifest::new(EmbeddingGenerationManifestInput {
290        compatibility_id,
291        source: request.source,
292        content_digest: request.batch.content_digest(),
293        vector_count: u64::try_from(rows.len()).map_err(|_| {
294            SearchArtifactError::ResourceExhausted {
295                resource: "embedding_rows",
296                limit: vector_limits.stored_vectors as u64,
297            }
298        })?,
299        dimension: request.descriptor.dimensions(),
300        generated_at_micros: request.generated_at_micros,
301        committed_at_micros: request.committed_at_micros,
302        publication_fingerprint,
303    })?;
304    debug_assert_eq!(manifest.generation_id(), generation_id);
305    checkpoint()?;
306    write_synced_file(
307        &private.path().join(MANIFEST_FILE),
308        &manifest.to_canonical_json()?,
309    )?;
310    sync_tree(private.path())?;
311    checkpoint()?;
312    Ok((private, manifest))
313}
314
315/// Reopen and fully validate the active generation for one descriptor.
316///
317/// A missing space root or active pointer returns `Ok(None)`. Once an active
318/// pointer exists, descriptor, pointer, manifest, and vector corruption fail
319/// closed and are never represented as an absent generation.
320///
321/// # Errors
322/// Returns structured incompatibility, corruption, resource, cancellation, or
323/// filesystem errors.
324pub fn current_embedding_generation<C>(
325    project_dir: &Path,
326    descriptor: &EmbeddingCompatibilityDescriptor,
327    vector_limits: VectorStoreLimits,
328    mut checkpoint: C,
329) -> Result<Option<EmbeddingGenerationPublication>, SearchArtifactError>
330where
331    C: FnMut() -> Result<(), SearchArtifactError>,
332{
333    checkpoint()?;
334    let compatibility_id = descriptor.compatibility_id()?;
335    let root = space_root(project_dir, compatibility_id);
336    if !path_exists(&root)? {
337        return Ok(None);
338    }
339    ensure_space_ancestors(project_dir, &root)?;
340    let active_path = root.join(ACTIVE_FILE);
341    if !path_exists(&active_path)? {
342        return Ok(None);
343    }
344    let reopened_descriptor = read_descriptor(&root, descriptor, compatibility_id)?;
345    let pointer = read_active_pointer(&active_path)?;
346    if pointer.compatibility_id != compatibility_id {
347        return Err(corrupt_primary(
348            &active_path,
349            "active pointer compatibility identity does not match its space path",
350        ));
351    }
352    validate_generation(
353        &root,
354        &reopened_descriptor,
355        compatibility_id,
356        pointer.generation_id,
357        vector_limits,
358        &mut checkpoint,
359    )
360    .map(Some)
361}
362
363fn validate_generation<C>(
364    root: &Path,
365    descriptor: &EmbeddingCompatibilityDescriptor,
366    compatibility_id: EmbeddingCompatibilityId,
367    generation_id: EmbeddingGenerationId,
368    vector_limits: VectorStoreLimits,
369    checkpoint: &mut C,
370) -> Result<EmbeddingGenerationPublication, SearchArtifactError>
371where
372    C: FnMut() -> Result<(), SearchArtifactError>,
373{
374    checkpoint()?;
375    let path = root.join(GENERATIONS_DIR).join(generation_id.to_hex());
376    ensure_existing_directory(&path)?;
377    validate_generation_layout(&path)?;
378    let manifest_path = path.join(MANIFEST_FILE);
379    let manifest_bytes = read_bounded_file(&manifest_path, MAX_MANIFEST_BYTES)?;
380    let manifest = EmbeddingGenerationManifest::from_json(&manifest_path, &manifest_bytes)
381        .map_err(|error| primary_from(&path, error))?;
382    if manifest.compatibility_id() != compatibility_id {
383        return Err(corrupt_primary(
384            &manifest_path,
385            "manifest compatibility identity does not match its space path",
386        ));
387    }
388    if manifest.generation_id() != generation_id {
389        return Err(corrupt_primary(
390            &manifest_path,
391            "manifest generation identity does not match its generation path",
392        ));
393    }
394    if manifest.dimension() != descriptor.dimensions() {
395        return Err(corrupt_primary(
396            &manifest_path,
397            "manifest dimension does not match compatibility descriptor",
398        ));
399    }
400    let vector_path = path.join(VECTOR_DATA_FILE);
401    let fingerprint = hash_file(&vector_path, vector_limits.parquet_bytes, checkpoint)
402        .map_err(|error| primary_from(&path, error))?;
403    if fingerprint != manifest.publication_fingerprint() {
404        return Err(corrupt_primary(
405            &vector_path,
406            "vector file fingerprint does not match generation manifest",
407        ));
408    }
409    let dimension = usize::try_from(manifest.dimension())
410        .map_err(|_| corrupt_primary(&manifest_path, "manifest dimension cannot be represented"))?;
411    let rows = read_vector_snapshot(&path, dimension, vector_limits, &mut *checkpoint)
412        .map_err(|error| primary_from(&path, error))?;
413    if u64::try_from(rows.len()).ok() != Some(manifest.vector_count()) {
414        return Err(corrupt_primary(
415            &vector_path,
416            "vector row count does not match generation manifest",
417        ));
418    }
419    if content_digest(&rows, checkpoint)? != manifest.content_digest() {
420        return Err(corrupt_primary(
421            &vector_path,
422            "canonical UUID/vector content digest does not match generation manifest",
423        ));
424    }
425    Ok(EmbeddingGenerationPublication {
426        path,
427        descriptor: descriptor.clone(),
428        manifest,
429    })
430}
431
432fn prepare_space_root(
433    project_dir: &Path,
434    compatibility_id: EmbeddingCompatibilityId,
435) -> Result<PathBuf, SearchArtifactError> {
436    let embeddings = project_dir.join("embeddings");
437    ensure_owned_directory(&embeddings)?;
438    let spaces = embeddings.join("spaces");
439    ensure_owned_directory(&spaces)?;
440    let root = spaces.join(compatibility_id.to_hex());
441    ensure_owned_directory(&root)?;
442    Ok(root)
443}
444
445fn space_root(project_dir: &Path, compatibility_id: EmbeddingCompatibilityId) -> PathBuf {
446    project_dir
447        .join("embeddings")
448        .join("spaces")
449        .join(compatibility_id.to_hex())
450}
451
452fn ensure_space_ancestors(project_dir: &Path, root: &Path) -> Result<(), SearchArtifactError> {
453    ensure_existing_directory(&project_dir.join("embeddings"))?;
454    ensure_existing_directory(&project_dir.join("embeddings").join("spaces"))?;
455    ensure_existing_directory(root)
456}
457
458fn validate_generation_layout(path: &Path) -> Result<(), SearchArtifactError> {
459    let mut names = Vec::new();
460    for entry in std::fs::read_dir(path)
461        .map_err(|source| io("scan immutable embedding generation", path, source))?
462    {
463        let entry = entry
464            .map_err(|source| io("read immutable embedding generation entry", path, source))?;
465        let file_type = entry.file_type().map_err(|source| {
466            io(
467                "inspect immutable embedding generation entry",
468                &entry.path(),
469                source,
470            )
471        })?;
472        if file_type.is_symlink() || !file_type.is_file() {
473            return Err(corrupt_primary(
474                &entry.path(),
475                "immutable embedding generation entries must be regular files",
476            ));
477        }
478        let name = entry.file_name().into_string().map_err(|_| {
479            corrupt_primary(
480                &entry.path(),
481                "immutable embedding generation file name must be UTF-8",
482            )
483        })?;
484        names.push(name);
485        if names.len() > 2 {
486            return Err(corrupt_primary(
487                path,
488                "immutable embedding generation contains unexpected files",
489            ));
490        }
491    }
492    names.sort_unstable();
493    if names != [MANIFEST_FILE, VECTOR_DATA_FILE] {
494        return Err(corrupt_primary(
495            path,
496            "immutable embedding generation must contain manifest.json and vectors.parquet",
497        ));
498    }
499    Ok(())
500}
501
502fn persist_or_verify_descriptor(
503    root: &Path,
504    descriptor: &EmbeddingCompatibilityDescriptor,
505    compatibility_id: EmbeddingCompatibilityId,
506) -> Result<(), SearchArtifactError> {
507    let path = root.join(SPACE_FILE);
508    if path_exists(&path)? {
509        read_descriptor(root, descriptor, compatibility_id).map(|_| ())
510    } else {
511        let bytes = descriptor.to_canonical_json()?;
512        persist_synced_file(&path, ".space.json.", &bytes)?;
513        sync_directory(root)
514    }
515}
516
517fn read_descriptor(
518    root: &Path,
519    requested: &EmbeddingCompatibilityDescriptor,
520    compatibility_id: EmbeddingCompatibilityId,
521) -> Result<EmbeddingCompatibilityDescriptor, SearchArtifactError> {
522    let path = root.join(SPACE_FILE);
523    let bytes = read_bounded_file(&path, MAX_DESCRIPTOR_BYTES)?;
524    let descriptor = EmbeddingCompatibilityDescriptor::from_json(&path, &bytes)
525        .map_err(|error| primary_from(root, error))?;
526    let reopened_id = descriptor
527        .compatibility_id()
528        .map_err(|error| primary_from(root, error))?;
529    if reopened_id != compatibility_id || &descriptor != requested {
530        return Err(corrupt_primary(
531            &path,
532            "compatibility descriptor does not match requested identity",
533        ));
534    }
535    Ok(descriptor)
536}
537
538#[derive(Deserialize)]
539#[serde(deny_unknown_fields)]
540struct RawActivePointer {
541    pointer_version: u32,
542    compatibility_id: String,
543    generation_id: String,
544    checksum: String,
545}
546
547struct ActivePointer {
548    compatibility_id: EmbeddingCompatibilityId,
549    generation_id: EmbeddingGenerationId,
550}
551
552fn read_active_pointer(path: &Path) -> Result<ActivePointer, SearchArtifactError> {
553    let bytes = read_bounded_file(path, MAX_ACTIVE_BYTES)?;
554    let raw: RawActivePointer =
555        serde_json::from_slice(&bytes).map_err(|error| corrupt_primary(path, error.to_string()))?;
556    if raw.pointer_version != POINTER_VERSION {
557        return Err(corrupt_primary(path, "unsupported active pointer version"));
558    }
559    let compatibility_id = EmbeddingCompatibilityId::from_hex(&raw.compatibility_id)
560        .map_err(|error| corrupt_primary(path, error.to_string()))?;
561    let generation_id = EmbeddingGenerationId::from_hex(&raw.generation_id)
562        .map_err(|error| corrupt_primary(path, error.to_string()))?;
563    let checksum = active_checksum(compatibility_id, generation_id);
564    if raw.checksum != checksum {
565        return Err(corrupt_primary(path, "active pointer checksum mismatch"));
566    }
567    let canonical = active_pointer_bytes(compatibility_id, generation_id)?;
568    if canonical != bytes {
569        return Err(corrupt_primary(
570            path,
571            "active pointer bytes are not exact canonical JSON",
572        ));
573    }
574    Ok(ActivePointer {
575        compatibility_id,
576        generation_id,
577    })
578}
579
580fn persist_active_pointer(
581    root: &Path,
582    compatibility_id: EmbeddingCompatibilityId,
583    generation_id: EmbeddingGenerationId,
584) -> Result<(), SearchArtifactError> {
585    let bytes = active_pointer_bytes(compatibility_id, generation_id)?;
586    persist_synced_file(&root.join(ACTIVE_FILE), ".active.json.", &bytes)
587}
588
589fn active_pointer_bytes(
590    compatibility_id: EmbeddingCompatibilityId,
591    generation_id: EmbeddingGenerationId,
592) -> Result<Vec<u8>, SearchArtifactError> {
593    serde_json::to_vec(&serde_json::json!({
594        "checksum": active_checksum(compatibility_id, generation_id),
595        "compatibility_id": compatibility_id.to_hex(),
596        "generation_id": generation_id.to_hex(),
597        "pointer_version": POINTER_VERSION,
598    }))
599    .map_err(|error| SearchArtifactError::Build(error.to_string()))
600}
601
602fn active_checksum(
603    compatibility_id: EmbeddingCompatibilityId,
604    generation_id: EmbeddingGenerationId,
605) -> String {
606    let mut hasher = Sha256::new();
607    hasher.update(b"graphforge.embedding.active.v1\0");
608    hasher.update(compatibility_id.to_hex().as_bytes());
609    hasher.update(generation_id.to_hex().as_bytes());
610    format!("{:x}", hasher.finalize())
611}
612
613fn content_digest<C>(
614    rows: &[StoredVector],
615    checkpoint: &mut C,
616) -> Result<EmbeddingContentDigest, SearchArtifactError>
617where
618    C: FnMut() -> Result<(), SearchArtifactError>,
619{
620    let mut hasher = Sha256::new();
621    for row in rows {
622        checkpoint()?;
623        hasher.update(row.node_uuid);
624        for value in &row.vector {
625            hasher.update(value.to_le_bytes());
626        }
627    }
628    EmbeddingContentDigest::from_hex(&format!("{:x}", hasher.finalize()))
629}
630
631fn hash_file<C>(
632    path: &Path,
633    max_bytes: u64,
634    checkpoint: &mut C,
635) -> Result<EmbeddingPublicationFingerprint, SearchArtifactError>
636where
637    C: FnMut() -> Result<(), SearchArtifactError>,
638{
639    ensure_regular_file(path)?;
640    let metadata =
641        std::fs::metadata(path).map_err(|source| io("inspect embedding file", path, source))?;
642    if metadata.len() > max_bytes {
643        return Err(SearchArtifactError::ResourceExhausted {
644            resource: "vector_parquet_bytes",
645            limit: max_bytes,
646        });
647    }
648    let mut file = File::open(path).map_err(|source| io("open embedding file", path, source))?;
649    let mut hasher = Sha256::new();
650    let mut buffer = vec![0_u8; HASH_BUFFER_BYTES];
651    loop {
652        checkpoint()?;
653        let read = file
654            .read(&mut buffer)
655            .map_err(|source| io("read embedding file", path, source))?;
656        if read == 0 {
657            break;
658        }
659        hasher.update(&buffer[..read]);
660    }
661    EmbeddingPublicationFingerprint::from_hex(&format!("{:x}", hasher.finalize()))
662}
663
664struct EmbeddingWriterLock {
665    file: File,
666}
667
668impl EmbeddingWriterLock {
669    fn acquire<C>(
670        project_dir: &Path,
671        compatibility_id: EmbeddingCompatibilityId,
672        limits: SearchCoordinationLimits,
673        checkpoint: &mut C,
674    ) -> Result<Self, SearchArtifactError>
675    where
676        C: FnMut() -> Result<(), SearchArtifactError>,
677    {
678        let embeddings = project_dir.join("embeddings");
679        ensure_owned_directory(&embeddings)?;
680        let path = embeddings.join(format!(
681            "{EMBEDDING_WRITER_LOCK_PREFIX}{}.lock",
682            compatibility_id.to_hex()
683        ));
684        let file = OpenOptions::new()
685            .read(true)
686            .write(true)
687            .create(true)
688            .truncate(false)
689            .open(&path)
690            .map_err(|source| SearchArtifactError::Lock {
691                path: path.clone(),
692                reason: source.to_string(),
693            })?;
694        let started = Instant::now();
695        loop {
696            match file.try_lock() {
697                Ok(()) => return Ok(Self { file }),
698                Err(std::fs::TryLockError::WouldBlock) => {
699                    checkpoint()?;
700                    if started.elapsed() >= limits.lock_timeout {
701                        return Err(SearchArtifactError::Lock {
702                            path,
703                            reason: format!(
704                                "timed out after {} ms",
705                                limits.lock_timeout.as_millis()
706                            ),
707                        });
708                    }
709                    std::thread::sleep(limits.lock_poll_interval);
710                }
711                Err(std::fs::TryLockError::Error(source)) => {
712                    return Err(SearchArtifactError::Lock {
713                        path,
714                        reason: source.to_string(),
715                    });
716                }
717            }
718        }
719    }
720}
721
722pub(crate) fn deletion_marker(
723    project_dir: &Path,
724    compatibility_id: EmbeddingCompatibilityId,
725) -> PathBuf {
726    project_dir.join("embeddings").join(format!(
727        "{EMBEDDING_DELETION_PREFIX}{}",
728        compatibility_id.to_hex()
729    ))
730}
731
732impl Drop for EmbeddingWriterLock {
733    fn drop(&mut self) {
734        let _ = self.file.unlock();
735    }
736}
737
738fn persist_synced_file(path: &Path, prefix: &str, bytes: &[u8]) -> Result<(), SearchArtifactError> {
739    let parent = path
740        .parent()
741        .ok_or_else(|| SearchArtifactError::Build("publication file has no parent".to_owned()))?;
742    let mut temp = tempfile::Builder::new()
743        .prefix(prefix)
744        .suffix(".tmp")
745        .tempfile_in(parent)
746        .map_err(|source| io("create embedding metadata temp", path, source))?;
747    temp.write_all(bytes)
748        .map_err(|source| io("write embedding metadata temp", path, source))?;
749    temp.as_file()
750        .sync_all()
751        .map_err(|source| io("sync embedding metadata temp", path, source))?;
752    temp.persist(path)
753        .map_err(|error| io("publish embedding metadata", path, error.error))?;
754    sync_directory(parent)
755}
756
757fn write_synced_file(path: &Path, bytes: &[u8]) -> Result<(), SearchArtifactError> {
758    let mut file = OpenOptions::new()
759        .create_new(true)
760        .write(true)
761        .open(path)
762        .map_err(|source| io("create embedding generation file", path, source))?;
763    file.write_all(bytes)
764        .map_err(|source| io("write embedding generation file", path, source))?;
765    file.sync_all()
766        .map_err(|source| io("sync embedding generation file", path, source))
767}
768
769fn read_bounded_file(path: &Path, max_bytes: u64) -> Result<Vec<u8>, SearchArtifactError> {
770    ensure_regular_file(path)?;
771    let metadata =
772        std::fs::metadata(path).map_err(|source| io("inspect embedding metadata", path, source))?;
773    if metadata.len() > max_bytes {
774        return Err(SearchArtifactError::ResourceExhausted {
775            resource: "embedding_metadata_bytes",
776            limit: max_bytes,
777        });
778    }
779    std::fs::read(path).map_err(|source| io("read embedding metadata", path, source))
780}
781
782fn sync_tree(root: &Path) -> Result<(), SearchArtifactError> {
783    let mut directories = vec![root.to_path_buf()];
784    let mut files = Vec::new();
785    let mut cursor = 0;
786    while cursor < directories.len() {
787        let directory = directories[cursor].clone();
788        cursor += 1;
789        for entry in std::fs::read_dir(&directory)
790            .map_err(|source| io("scan embedding generation", &directory, source))?
791        {
792            let entry = entry
793                .map_err(|source| io("read embedding generation entry", &directory, source))?;
794            let file_type = entry.file_type().map_err(|source| {
795                io("inspect embedding generation entry", &entry.path(), source)
796            })?;
797            if file_type.is_symlink() {
798                return Err(corrupt_primary(
799                    &entry.path(),
800                    "embedding generation must not contain symlinks",
801                ));
802            }
803            if file_type.is_dir() {
804                directories.push(entry.path());
805            } else if file_type.is_file() {
806                files.push(entry.path());
807            } else {
808                return Err(corrupt_primary(
809                    &entry.path(),
810                    "embedding generation contains unsupported file type",
811                ));
812            }
813        }
814    }
815    files.sort_unstable();
816    for file in files {
817        OpenOptions::new()
818            .read(true)
819            .write(true)
820            .open(&file)
821            .and_then(|file| file.sync_all())
822            .map_err(|source| io("sync embedding generation file", &file, source))?;
823    }
824    directories.sort_unstable_by_key(|path| std::cmp::Reverse(path.components().count()));
825    for directory in directories {
826        sync_directory(&directory)?;
827    }
828    Ok(())
829}
830
831fn ensure_owned_directory(path: &Path) -> Result<(), SearchArtifactError> {
832    match std::fs::symlink_metadata(path) {
833        Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_dir() => Err(
834            corrupt_primary(path, "embedding path must be a real directory"),
835        ),
836        Ok(_) => Ok(()),
837        Err(source) if source.kind() == std::io::ErrorKind::NotFound => std::fs::create_dir(path)
838            .map_err(|source| io("create embedding directory", path, source)),
839        Err(source) => Err(io("inspect embedding directory", path, source)),
840    }
841}
842
843fn ensure_existing_directory(path: &Path) -> Result<(), SearchArtifactError> {
844    let metadata = std::fs::symlink_metadata(path)
845        .map_err(|source| io("inspect embedding directory", path, source))?;
846    if metadata.file_type().is_symlink() || !metadata.is_dir() {
847        return Err(corrupt_primary(
848            path,
849            "embedding path must be a real directory",
850        ));
851    }
852    Ok(())
853}
854
855fn ensure_regular_file(path: &Path) -> Result<(), SearchArtifactError> {
856    let metadata = std::fs::symlink_metadata(path)
857        .map_err(|source| io("inspect embedding file", path, source))?;
858    if metadata.file_type().is_symlink() || !metadata.is_file() {
859        return Err(corrupt_primary(
860            path,
861            "embedding path must be a regular file",
862        ));
863    }
864    Ok(())
865}
866
867fn path_exists(path: &Path) -> Result<bool, SearchArtifactError> {
868    match std::fs::symlink_metadata(path) {
869        Ok(_) => Ok(true),
870        Err(source) if source.kind() == std::io::ErrorKind::NotFound => Ok(false),
871        Err(source) => Err(io("inspect embedding path", path, source)),
872    }
873}
874
875#[cfg(unix)]
876fn sync_directory(path: &Path) -> Result<(), SearchArtifactError> {
877    File::open(path)
878        .and_then(|file| file.sync_all())
879        .map_err(|source| io("sync embedding directory", path, source))
880}
881
882#[cfg(not(unix))]
883fn sync_directory(_path: &Path) -> Result<(), SearchArtifactError> {
884    Ok(())
885}
886
887fn primary_from(path: &Path, error: SearchArtifactError) -> SearchArtifactError {
888    match error {
889        error @ (SearchArtifactError::Cancelled
890        | SearchArtifactError::ResourceExhausted { .. }
891        | SearchArtifactError::Lock { .. }
892        | SearchArtifactError::Io { .. }
893        | SearchArtifactError::CorruptPrimaryVectors { .. }) => error,
894        error => corrupt_primary(path, error.to_string()),
895    }
896}
897
898fn invalid(field: &'static str, reason: impl Into<String>) -> SearchArtifactError {
899    SearchArtifactError::InvalidSelector {
900        field,
901        reason: reason.into(),
902    }
903}
904
905fn corrupt_primary(path: &Path, reason: impl Into<String>) -> SearchArtifactError {
906    SearchArtifactError::CorruptPrimaryVectors {
907        path: path.to_path_buf(),
908        reason: reason.into(),
909    }
910}
911
912fn io(operation: &'static str, path: &Path, source: std::io::Error) -> SearchArtifactError {
913    SearchArtifactError::Io {
914        operation,
915        path: path.to_path_buf(),
916        source,
917    }
918}
919
920#[cfg(test)]
921mod tests {
922    use std::collections::{BTreeMap, BTreeSet};
923
924    use serde_json::json;
925
926    use super::*;
927    use crate::{
928        EmbeddingBatchRow, EmbeddingCompatibilityInput, EmbeddingDistance, EmbeddingNormalization,
929        EmbeddingProducerIdentity, EmbeddingValueType, validate_embedding_batch,
930    };
931
932    const NODE: [u8; 16] = [7; 16];
933
934    fn descriptor(model: &str, dimension: u32) -> EmbeddingCompatibilityDescriptor {
935        EmbeddingCompatibilityDescriptor::new(EmbeddingCompatibilityInput {
936            producer: EmbeddingProducerIdentity::Local {
937                implementation: "test-adapter".to_owned(),
938                model: model.to_owned(),
939                revision: "r1".to_owned(),
940                contract_version: "v1".to_owned(),
941            },
942            dimensions: dimension,
943            value_type: EmbeddingValueType::Float32,
944            normalization: EmbeddingNormalization::None,
945            distance: EmbeddingDistance::Cosine,
946            tokenizer: None,
947            chunking: None,
948            hyperparameters: BTreeMap::new(),
949            input_recipe: BTreeMap::from([("property".to_owned(), json!("body"))]),
950            source_projection_recipe: BTreeMap::from([("label".to_owned(), json!("Document"))]),
951        })
952        .unwrap()
953    }
954
955    fn source(generation: u64) -> EmbeddingSourceState {
956        EmbeddingSourceState::new(generation, [generation as u8; 32], [9; 32], 1)
957    }
958
959    fn batch(values: &[f32]) -> ValidatedEmbeddingBatch {
960        validate_embedding_batch(
961            vec![EmbeddingBatchRow {
962                node_uuid: NODE,
963                vector: values.to_vec(),
964            }],
965            &BTreeSet::from([NODE]),
966            values.len(),
967            EmbeddingNormalization::None,
968            VectorStoreLimits::default(),
969            || Ok(()),
970        )
971        .unwrap()
972    }
973
974    fn request<'a>(
975        descriptor: &'a EmbeddingCompatibilityDescriptor,
976        source: EmbeddingSourceState,
977        batch: &'a ValidatedEmbeddingBatch,
978        committed_at_micros: i64,
979    ) -> EmbeddingPublicationRequest<'a> {
980        EmbeddingPublicationRequest {
981            descriptor,
982            source,
983            batch,
984            generated_at_micros: 10,
985            committed_at_micros,
986        }
987    }
988
989    fn publish(
990        dir: &Path,
991        descriptor: &EmbeddingCompatibilityDescriptor,
992        source: EmbeddingSourceState,
993        batch: &ValidatedEmbeddingBatch,
994        committed_at_micros: i64,
995    ) -> EmbeddingPublicationOutcome {
996        publish_embedding_generation(
997            dir,
998            request(descriptor, source, batch, committed_at_micros),
999            VectorStoreLimits::default(),
1000            SearchCoordinationLimits::default(),
1001            || Ok(()),
1002        )
1003        .unwrap()
1004    }
1005
1006    #[test]
1007    fn complete_publication_reopens_and_identical_content_is_reused() {
1008        let dir = tempfile::tempdir().unwrap();
1009        let descriptor = descriptor("model-a", 2);
1010        let batch = batch(&[1.0, 2.0]);
1011        let first = publish(dir.path(), &descriptor, source(1), &batch, 20);
1012        assert!(matches!(first, EmbeddingPublicationOutcome::Published(_)));
1013        let reopened = current_embedding_generation(
1014            dir.path(),
1015            &descriptor,
1016            VectorStoreLimits::default(),
1017            || Ok(()),
1018        )
1019        .unwrap()
1020        .unwrap();
1021        assert_eq!(reopened, first.publication().clone());
1022
1023        let second = publish(dir.path(), &descriptor, source(1), &batch, 99);
1024        assert!(matches!(second, EmbeddingPublicationOutcome::Reused(_)));
1025        assert_eq!(second.publication(), first.publication());
1026        assert_eq!(
1027            std::fs::read_dir(first.publication().path.parent().unwrap())
1028                .unwrap()
1029                .count(),
1030            1
1031        );
1032    }
1033
1034    #[test]
1035    fn failed_replacement_never_changes_the_active_generation() {
1036        let dir = tempfile::tempdir().unwrap();
1037        let descriptor = descriptor("model-a", 2);
1038        let first_batch = batch(&[1.0, 2.0]);
1039        let first = publish(dir.path(), &descriptor, source(1), &first_batch, 20);
1040        let second_batch = batch(&[3.0, 4.0]);
1041        let generations = first.publication().path.parent().unwrap().to_path_buf();
1042        let error = publish_embedding_generation(
1043            dir.path(),
1044            request(&descriptor, source(2), &second_batch, 30),
1045            VectorStoreLimits::default(),
1046            SearchCoordinationLimits::default(),
1047            || {
1048                if std::fs::read_dir(&generations)
1049                    .map(|entries| entries.count() > 1)
1050                    .unwrap_or(false)
1051                {
1052                    Err(SearchArtifactError::Cancelled)
1053                } else {
1054                    Ok(())
1055                }
1056            },
1057        )
1058        .unwrap_err();
1059        assert!(matches!(error, SearchArtifactError::Cancelled));
1060        let active = current_embedding_generation(
1061            dir.path(),
1062            &descriptor,
1063            VectorStoreLimits::default(),
1064            || Ok(()),
1065        )
1066        .unwrap()
1067        .unwrap();
1068        assert_eq!(active, first.publication().clone());
1069    }
1070
1071    #[test]
1072    fn compatibility_lineages_are_independent_for_the_same_uuid() {
1073        let dir = tempfile::tempdir().unwrap();
1074        let left_descriptor = descriptor("model-a", 2);
1075        let right_descriptor = descriptor("model-b", 3);
1076        let left_batch = batch(&[1.0, 2.0]);
1077        let right_batch = batch(&[1.0, 2.0, 3.0]);
1078        let left = publish(dir.path(), &left_descriptor, source(1), &left_batch, 20);
1079        let right = publish(dir.path(), &right_descriptor, source(1), &right_batch, 20);
1080        assert_ne!(left.publication().path, right.publication().path);
1081        assert!(left.publication().path.exists());
1082        assert!(right.publication().path.exists());
1083    }
1084
1085    #[test]
1086    fn pointer_descriptor_and_vector_corruption_fail_closed() {
1087        let cases = ["pointer", "descriptor", "vector"];
1088        for case in cases {
1089            let dir = tempfile::tempdir().unwrap();
1090            let descriptor = descriptor("model-a", 2);
1091            let batch = batch(&[1.0, 2.0]);
1092            let published = publish(dir.path(), &descriptor, source(1), &batch, 20);
1093            let compatibility = descriptor.compatibility_id().unwrap().to_hex();
1094            let root = dir.path().join("embeddings/spaces").join(compatibility);
1095            let path = match case {
1096                "pointer" => root.join(ACTIVE_FILE),
1097                "descriptor" => root.join(SPACE_FILE),
1098                "vector" => published.publication().path.join(VECTOR_DATA_FILE),
1099                _ => unreachable!(),
1100            };
1101            std::fs::write(path, b"corrupt").unwrap();
1102            assert!(matches!(
1103                current_embedding_generation(
1104                    dir.path(),
1105                    &descriptor,
1106                    VectorStoreLimits::default(),
1107                    || Ok(()),
1108                ),
1109                Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1110            ));
1111        }
1112    }
1113
1114    #[test]
1115    fn private_and_traversal_like_trees_are_never_visible() {
1116        let dir = tempfile::tempdir().unwrap();
1117        let descriptor = descriptor("model-a", 2);
1118        let compatibility = descriptor.compatibility_id().unwrap().to_hex();
1119        let root = dir.path().join("embeddings/spaces").join(compatibility);
1120        std::fs::create_dir_all(root.join(".build-crashed")).unwrap();
1121        assert!(
1122            current_embedding_generation(
1123                dir.path(),
1124                &descriptor,
1125                VectorStoreLimits::default(),
1126                || Ok(()),
1127            )
1128            .unwrap()
1129            .is_none()
1130        );
1131        let batch = batch(&[1.0, 2.0]);
1132        publish(dir.path(), &descriptor, source(1), &batch, 20);
1133        std::fs::write(
1134            root.join(ACTIVE_FILE),
1135            br#"{"pointer_version":1,"compatibility_id":"..","generation_id":"..","checksum":"no"}"#,
1136        )
1137        .unwrap();
1138        assert!(matches!(
1139            current_embedding_generation(
1140                dir.path(),
1141                &descriptor,
1142                VectorStoreLimits::default(),
1143                || Ok(()),
1144            ),
1145            Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1146        ));
1147    }
1148
1149    #[test]
1150    fn unexpected_generation_entries_fail_closed() {
1151        let dir = tempfile::tempdir().unwrap();
1152        let descriptor = descriptor("model-a", 2);
1153        let batch = batch(&[1.0, 2.0]);
1154        let published = publish(dir.path(), &descriptor, source(1), &batch, 20);
1155        std::fs::write(published.publication().path.join("unexpected"), b"data").unwrap();
1156        assert!(matches!(
1157            current_embedding_generation(
1158                dir.path(),
1159                &descriptor,
1160                VectorStoreLimits::default(),
1161                || Ok(()),
1162            ),
1163            Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1164        ));
1165    }
1166
1167    #[cfg(unix)]
1168    #[test]
1169    fn symlinked_embedding_ancestor_fails_closed() {
1170        use std::os::unix::fs::symlink;
1171
1172        let project = tempfile::tempdir().unwrap();
1173        let external = tempfile::tempdir().unwrap();
1174        symlink(external.path(), project.path().join("embeddings")).unwrap();
1175        let descriptor = descriptor("model-a", 2);
1176        let compatibility = descriptor.compatibility_id().unwrap().to_hex();
1177        std::fs::create_dir_all(external.path().join("spaces").join(compatibility)).unwrap();
1178        assert!(matches!(
1179            current_embedding_generation(
1180                project.path(),
1181                &descriptor,
1182                VectorStoreLimits::default(),
1183                || Ok(()),
1184            ),
1185            Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1186        ));
1187    }
1188}