1use 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#[derive(Clone, Copy, Debug)]
35pub struct EmbeddingPublicationRequest<'a> {
36 pub descriptor: &'a EmbeddingCompatibilityDescriptor,
38 pub source: EmbeddingSourceState,
40 pub batch: &'a ValidatedEmbeddingBatch,
42 pub generated_at_micros: i64,
44 pub committed_at_micros: i64,
46}
47
48#[derive(Clone, Debug, PartialEq)]
50pub struct EmbeddingGenerationPublication {
51 pub path: PathBuf,
53 pub descriptor: EmbeddingCompatibilityDescriptor,
55 pub manifest: EmbeddingGenerationManifest,
57}
58
59#[derive(Clone, Debug, PartialEq)]
61pub enum EmbeddingPublicationOutcome {
62 Reused(EmbeddingGenerationPublication),
64 Published(EmbeddingGenerationPublication),
66}
67
68impl EmbeddingPublicationOutcome {
69 #[must_use]
71 pub const fn publication(&self) -> &EmbeddingGenerationPublication {
72 match self {
73 Self::Reused(publication) | Self::Published(publication) => publication,
74 }
75 }
76}
77
78pub 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
192pub 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
315pub 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}