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 use std::time::Duration;
924
925 use serde_json::json;
926
927 use super::*;
928 use crate::{
929 EmbeddingBatchRow, EmbeddingCompatibilityInput, EmbeddingDistance, EmbeddingNormalization,
930 EmbeddingProducerIdentity, EmbeddingValueType, validate_embedding_batch,
931 };
932
933 const NODE: [u8; 16] = [7; 16];
934
935 fn descriptor(model: &str, dimension: u32) -> EmbeddingCompatibilityDescriptor {
936 EmbeddingCompatibilityDescriptor::new(EmbeddingCompatibilityInput {
937 producer: EmbeddingProducerIdentity::Local {
938 implementation: "test-adapter".to_owned(),
939 model: model.to_owned(),
940 revision: "r1".to_owned(),
941 contract_version: "v1".to_owned(),
942 },
943 dimensions: dimension,
944 value_type: EmbeddingValueType::Float32,
945 normalization: EmbeddingNormalization::None,
946 distance: EmbeddingDistance::Cosine,
947 tokenizer: None,
948 chunking: None,
949 hyperparameters: BTreeMap::new(),
950 input_recipe: BTreeMap::from([("property".to_owned(), json!("body"))]),
951 source_projection_recipe: BTreeMap::from([("label".to_owned(), json!("Document"))]),
952 })
953 .unwrap()
954 }
955
956 fn source(generation: u64) -> EmbeddingSourceState {
957 EmbeddingSourceState::new(generation, [generation as u8; 32], [9; 32], 1)
958 }
959
960 fn batch(values: &[f32]) -> ValidatedEmbeddingBatch {
961 validate_embedding_batch(
962 vec![EmbeddingBatchRow {
963 node_uuid: NODE,
964 vector: values.to_vec(),
965 }],
966 &BTreeSet::from([NODE]),
967 values.len(),
968 EmbeddingNormalization::None,
969 VectorStoreLimits::default(),
970 || Ok(()),
971 )
972 .unwrap()
973 }
974
975 fn request<'a>(
976 descriptor: &'a EmbeddingCompatibilityDescriptor,
977 source: EmbeddingSourceState,
978 batch: &'a ValidatedEmbeddingBatch,
979 committed_at_micros: i64,
980 ) -> EmbeddingPublicationRequest<'a> {
981 EmbeddingPublicationRequest {
982 descriptor,
983 source,
984 batch,
985 generated_at_micros: 10,
986 committed_at_micros,
987 }
988 }
989
990 fn publish(
991 dir: &Path,
992 descriptor: &EmbeddingCompatibilityDescriptor,
993 source: EmbeddingSourceState,
994 batch: &ValidatedEmbeddingBatch,
995 committed_at_micros: i64,
996 ) -> EmbeddingPublicationOutcome {
997 publish_embedding_generation(
998 dir,
999 request(descriptor, source, batch, committed_at_micros),
1000 VectorStoreLimits::default(),
1001 SearchCoordinationLimits::default(),
1002 || Ok(()),
1003 )
1004 .unwrap()
1005 }
1006
1007 #[test]
1008 fn complete_publication_reopens_and_identical_content_is_reused() {
1009 let dir = tempfile::tempdir().unwrap();
1010 let descriptor = descriptor("model-a", 2);
1011 let batch = batch(&[1.0, 2.0]);
1012 let first = publish(dir.path(), &descriptor, source(1), &batch, 20);
1013 assert!(matches!(first, EmbeddingPublicationOutcome::Published(_)));
1014 let reopened = current_embedding_generation(
1015 dir.path(),
1016 &descriptor,
1017 VectorStoreLimits::default(),
1018 || Ok(()),
1019 )
1020 .unwrap()
1021 .unwrap();
1022 assert_eq!(reopened, first.publication().clone());
1023
1024 let second = publish(dir.path(), &descriptor, source(1), &batch, 99);
1025 assert!(matches!(second, EmbeddingPublicationOutcome::Reused(_)));
1026 assert_eq!(second.publication(), first.publication());
1027 assert_eq!(
1028 std::fs::read_dir(first.publication().path.parent().unwrap())
1029 .unwrap()
1030 .count(),
1031 1
1032 );
1033 }
1034
1035 #[test]
1036 fn failed_replacement_never_changes_the_active_generation() {
1037 let dir = tempfile::tempdir().unwrap();
1038 let descriptor = descriptor("model-a", 2);
1039 let first_batch = batch(&[1.0, 2.0]);
1040 let first = publish(dir.path(), &descriptor, source(1), &first_batch, 20);
1041 let second_batch = batch(&[3.0, 4.0]);
1042 let generations = first.publication().path.parent().unwrap().to_path_buf();
1043 let error = publish_embedding_generation(
1044 dir.path(),
1045 request(&descriptor, source(2), &second_batch, 30),
1046 VectorStoreLimits::default(),
1047 SearchCoordinationLimits::default(),
1048 || {
1049 if std::fs::read_dir(&generations)
1050 .map(|entries| entries.count() > 1)
1051 .unwrap_or(false)
1052 {
1053 Err(SearchArtifactError::Cancelled)
1054 } else {
1055 Ok(())
1056 }
1057 },
1058 )
1059 .unwrap_err();
1060 assert!(matches!(error, SearchArtifactError::Cancelled));
1061 let active = current_embedding_generation(
1062 dir.path(),
1063 &descriptor,
1064 VectorStoreLimits::default(),
1065 || Ok(()),
1066 )
1067 .unwrap()
1068 .unwrap();
1069 assert_eq!(active, first.publication().clone());
1070 }
1071
1072 #[test]
1073 fn compatibility_lineages_are_independent_for_the_same_uuid() {
1074 let dir = tempfile::tempdir().unwrap();
1075 let left_descriptor = descriptor("model-a", 2);
1076 let right_descriptor = descriptor("model-b", 3);
1077 let left_batch = batch(&[1.0, 2.0]);
1078 let right_batch = batch(&[1.0, 2.0, 3.0]);
1079 let left = publish(dir.path(), &left_descriptor, source(1), &left_batch, 20);
1080 let right = publish(dir.path(), &right_descriptor, source(1), &right_batch, 20);
1081 assert_ne!(left.publication().path, right.publication().path);
1082 assert!(left.publication().path.exists());
1083 assert!(right.publication().path.exists());
1084 }
1085
1086 #[test]
1087 fn pointer_descriptor_and_vector_corruption_fail_closed() {
1088 let cases = ["pointer", "descriptor", "vector"];
1089 for case in cases {
1090 let dir = tempfile::tempdir().unwrap();
1091 let descriptor = descriptor("model-a", 2);
1092 let batch = batch(&[1.0, 2.0]);
1093 let published = publish(dir.path(), &descriptor, source(1), &batch, 20);
1094 let compatibility = descriptor.compatibility_id().unwrap().to_hex();
1095 let root = dir.path().join("embeddings/spaces").join(compatibility);
1096 let path = match case {
1097 "pointer" => root.join(ACTIVE_FILE),
1098 "descriptor" => root.join(SPACE_FILE),
1099 "vector" => published.publication().path.join(VECTOR_DATA_FILE),
1100 _ => unreachable!(),
1101 };
1102 std::fs::write(path, b"corrupt").unwrap();
1103 assert!(matches!(
1104 current_embedding_generation(
1105 dir.path(),
1106 &descriptor,
1107 VectorStoreLimits::default(),
1108 || Ok(()),
1109 ),
1110 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1111 ));
1112 }
1113 }
1114
1115 #[test]
1116 fn canonical_manifest_identity_and_dimension_mismatches_fail_closed() {
1117 let descriptor = descriptor("model-a", 2);
1118 let compatibility = descriptor.compatibility_id().unwrap();
1119 let vectors = batch(&[1.0, 2.0]);
1120 for case in ["compatibility", "generation", "dimension"] {
1121 let dir = tempfile::tempdir().unwrap();
1122 let published = publish(dir.path(), &descriptor, source(1), &vectors, 20);
1123 let original = &published.publication().manifest;
1124 let manifest = EmbeddingGenerationManifest::new(EmbeddingGenerationManifestInput {
1125 compatibility_id: if case == "compatibility" {
1126 EmbeddingCompatibilityId::from_hex(&"11".repeat(32)).unwrap()
1127 } else {
1128 compatibility
1129 },
1130 source: if case == "generation" {
1131 source(2)
1132 } else {
1133 source(1)
1134 },
1135 content_digest: vectors.content_digest(),
1136 vector_count: 1,
1137 dimension: if case == "dimension" { 3 } else { 2 },
1138 generated_at_micros: 10,
1139 committed_at_micros: 20,
1140 publication_fingerprint: original.publication_fingerprint(),
1141 })
1142 .unwrap();
1143 std::fs::write(
1144 published.publication().path.join(MANIFEST_FILE),
1145 manifest.to_canonical_json().unwrap(),
1146 )
1147 .unwrap();
1148 let error = current_embedding_generation(
1149 dir.path(),
1150 &descriptor,
1151 VectorStoreLimits::default(),
1152 || Ok(()),
1153 )
1154 .unwrap_err();
1155 assert!(matches!(
1156 error,
1157 SearchArtifactError::CorruptPrimaryVectors { .. }
1158 ));
1159 }
1160 }
1161
1162 #[test]
1163 fn canonical_manifest_rejects_vector_row_count_and_content_drift() {
1164 let descriptor = descriptor("model-a", 2);
1165 let compatibility = descriptor.compatibility_id().unwrap();
1166 let original_vectors = batch(&[1.0, 2.0]);
1167 for case in ["row-count", "content"] {
1168 let dir = tempfile::tempdir().unwrap();
1169 let published = publish(dir.path(), &descriptor, source(1), &original_vectors, 20);
1170 let generation = &published.publication().path;
1171 let replacement = batch(&[3.0, 4.0]);
1172 let replacement_rows = replacement
1173 .rows()
1174 .iter()
1175 .map(|row| StoredVector {
1176 node_uuid: row.node_uuid,
1177 vector: row.vector.clone(),
1178 updated_at_micros: 10,
1179 })
1180 .collect::<Vec<_>>();
1181 let rows: &[StoredVector] = if case == "row-count" {
1182 &[]
1183 } else {
1184 &replacement_rows
1185 };
1186 let vector_path =
1187 write_vector_snapshot(generation, rows, 2, VectorStoreLimits::default(), || Ok(()))
1188 .unwrap();
1189 let fingerprint = hash_file(
1190 &vector_path,
1191 VectorStoreLimits::default().parquet_bytes,
1192 &mut || Ok(()),
1193 )
1194 .unwrap();
1195 let manifest = EmbeddingGenerationManifest::new(EmbeddingGenerationManifestInput {
1196 compatibility_id: compatibility,
1197 source: source(1),
1198 content_digest: original_vectors.content_digest(),
1199 vector_count: 1,
1200 dimension: 2,
1201 generated_at_micros: 10,
1202 committed_at_micros: 20,
1203 publication_fingerprint: fingerprint,
1204 })
1205 .unwrap();
1206 std::fs::write(
1207 generation.join(MANIFEST_FILE),
1208 manifest.to_canonical_json().unwrap(),
1209 )
1210 .unwrap();
1211 assert!(matches!(
1212 current_embedding_generation(
1213 dir.path(),
1214 &descriptor,
1215 VectorStoreLimits::default(),
1216 || Ok(()),
1217 ),
1218 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1219 ));
1220 }
1221 }
1222
1223 #[test]
1224 fn persisted_descriptor_must_match_the_exact_requested_contract() {
1225 let dir = tempfile::tempdir().unwrap();
1226 let expected_descriptor = descriptor("model-a", 2);
1227 let vectors = batch(&[1.0, 2.0]);
1228 publish(dir.path(), &expected_descriptor, source(1), &vectors, 20);
1229 let compatibility = expected_descriptor.compatibility_id().unwrap();
1230 let root = space_root(dir.path(), compatibility);
1231 let other = descriptor("model-b", 2);
1232 assert!(matches!(
1233 read_descriptor(&root, &other, compatibility),
1234 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1235 ));
1236 assert!(root.join(SPACE_FILE).is_file());
1237 }
1238
1239 #[test]
1240 fn private_and_traversal_like_trees_are_never_visible() {
1241 let dir = tempfile::tempdir().unwrap();
1242 let descriptor = descriptor("model-a", 2);
1243 let compatibility = descriptor.compatibility_id().unwrap().to_hex();
1244 let root = dir.path().join("embeddings/spaces").join(compatibility);
1245 std::fs::create_dir_all(root.join(".build-crashed")).unwrap();
1246 assert!(
1247 current_embedding_generation(
1248 dir.path(),
1249 &descriptor,
1250 VectorStoreLimits::default(),
1251 || Ok(()),
1252 )
1253 .unwrap()
1254 .is_none()
1255 );
1256 let batch = batch(&[1.0, 2.0]);
1257 publish(dir.path(), &descriptor, source(1), &batch, 20);
1258 std::fs::write(
1259 root.join(ACTIVE_FILE),
1260 br#"{"pointer_version":1,"compatibility_id":"..","generation_id":"..","checksum":"no"}"#,
1261 )
1262 .unwrap();
1263 assert!(matches!(
1264 current_embedding_generation(
1265 dir.path(),
1266 &descriptor,
1267 VectorStoreLimits::default(),
1268 || Ok(()),
1269 ),
1270 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1271 ));
1272 }
1273
1274 #[test]
1275 fn unexpected_generation_entries_fail_closed() {
1276 let dir = tempfile::tempdir().unwrap();
1277 let descriptor = descriptor("model-a", 2);
1278 let batch = batch(&[1.0, 2.0]);
1279 let published = publish(dir.path(), &descriptor, source(1), &batch, 20);
1280 std::fs::write(published.publication().path.join("unexpected"), b"data").unwrap();
1281 assert!(matches!(
1282 current_embedding_generation(
1283 dir.path(),
1284 &descriptor,
1285 VectorStoreLimits::default(),
1286 || Ok(()),
1287 ),
1288 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1289 ));
1290 }
1291
1292 #[test]
1293 fn active_pointer_malformed_state_matrix_fails_closed_without_repointing() {
1294 let descriptor = descriptor("model-a", 2);
1295 let compatibility = descriptor.compatibility_id().unwrap();
1296 let vectors = batch(&[1.0, 2.0]);
1297 let cases = [
1298 br#"{"pointer_version":2,"compatibility_id":"bad","generation_id":"bad","checksum":"bad"}"#.as_slice(),
1299 br#"{"pointer_version":1,"compatibility_id":"bad","generation_id":"bad","checksum":"bad"}"#,
1300 br#"{"pointer_version":1,"compatibility_id":"0000000000000000000000000000000000000000000000000000000000000000","generation_id":"bad","checksum":"bad"}"#,
1301 br#"{"pointer_version":1,"compatibility_id":"0000000000000000000000000000000000000000000000000000000000000000","generation_id":"0000000000000000000000000000000000000000000000000000000000000000","checksum":"bad"}"#,
1302 ];
1303 for bytes in cases {
1304 let dir = tempfile::tempdir().unwrap();
1305 let published = publish(dir.path(), &descriptor, source(1), &vectors, 20);
1306 let active = space_root(dir.path(), compatibility).join(ACTIVE_FILE);
1307 let generation = published.publication().path.clone();
1308 std::fs::write(&active, bytes).unwrap();
1309 let before = std::fs::read(&active).unwrap();
1310 assert!(matches!(
1311 current_embedding_generation(
1312 dir.path(),
1313 &descriptor,
1314 VectorStoreLimits::default(),
1315 || Ok(()),
1316 ),
1317 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1318 ));
1319 assert_eq!(std::fs::read(&active).unwrap(), before);
1320 assert!(generation.is_dir());
1321 }
1322 }
1323
1324 #[test]
1325 fn canonical_active_pointer_identity_and_encoding_mismatches_fail_closed() {
1326 let descriptor = descriptor("model-a", 2);
1327 let compatibility = descriptor.compatibility_id().unwrap();
1328 let vectors = batch(&[1.0, 2.0]);
1329
1330 let dir = tempfile::tempdir().unwrap();
1331 let published = publish(dir.path(), &descriptor, source(1), &vectors, 20);
1332 let active = space_root(dir.path(), compatibility).join(ACTIVE_FILE);
1333 let other = EmbeddingCompatibilityId::from_hex(&"11".repeat(32)).unwrap();
1334 std::fs::write(
1335 &active,
1336 active_pointer_bytes(other, published.publication().manifest.generation_id()).unwrap(),
1337 )
1338 .unwrap();
1339 assert!(matches!(
1340 current_embedding_generation(
1341 dir.path(),
1342 &descriptor,
1343 VectorStoreLimits::default(),
1344 || Ok(()),
1345 ),
1346 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1347 ));
1348
1349 let dir = tempfile::tempdir().unwrap();
1350 publish(dir.path(), &descriptor, source(1), &vectors, 20);
1351 let active = space_root(dir.path(), compatibility).join(ACTIVE_FILE);
1352 let mut noncanonical = vec![b' '];
1353 noncanonical.extend_from_slice(&std::fs::read(&active).unwrap());
1354 std::fs::write(&active, noncanonical).unwrap();
1355 assert!(matches!(
1356 current_embedding_generation(
1357 dir.path(),
1358 &descriptor,
1359 VectorStoreLimits::default(),
1360 || Ok(()),
1361 ),
1362 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1363 ));
1364 }
1365
1366 #[test]
1367 fn generation_layout_and_metadata_limits_fail_closed_without_mutation() {
1368 let descriptor = descriptor("model-a", 2);
1369 let vectors = batch(&[1.0, 2.0]);
1370 for case in [
1371 "missing-manifest",
1372 "manifest-dimension",
1373 "directory-entry",
1374 "oversized-active",
1375 ] {
1376 let dir = tempfile::tempdir().unwrap();
1377 let published = publish(dir.path(), &descriptor, source(1), &vectors, 20);
1378 let compatibility = descriptor.compatibility_id().unwrap();
1379 let root = space_root(dir.path(), compatibility);
1380 match case {
1381 "missing-manifest" => {
1382 std::fs::remove_file(published.publication().path.join(MANIFEST_FILE)).unwrap();
1383 }
1384 "manifest-dimension" => {
1385 let path = published.publication().path.join(MANIFEST_FILE);
1386 let mut manifest: serde_json::Value =
1387 serde_json::from_slice(&std::fs::read(&path).unwrap()).unwrap();
1388 manifest["dimension"] = serde_json::json!(3);
1389 std::fs::write(path, serde_json::to_vec(&manifest).unwrap()).unwrap();
1390 }
1391 "directory-entry" => {
1392 std::fs::create_dir(published.publication().path.join("nested")).unwrap();
1393 }
1394 "oversized-active" => {
1395 std::fs::write(
1396 root.join(ACTIVE_FILE),
1397 vec![b'x'; MAX_ACTIVE_BYTES as usize + 1],
1398 )
1399 .unwrap();
1400 }
1401 _ => unreachable!(),
1402 }
1403 let before = std::fs::read_dir(&root).unwrap().count();
1404 let error = current_embedding_generation(
1405 dir.path(),
1406 &descriptor,
1407 VectorStoreLimits::default(),
1408 || Ok(()),
1409 )
1410 .unwrap_err();
1411 assert!(matches!(
1412 error,
1413 SearchArtifactError::CorruptPrimaryVectors { .. }
1414 | SearchArtifactError::ResourceExhausted { .. }
1415 ));
1416 assert_eq!(std::fs::read_dir(&root).unwrap().count(), before);
1417 }
1418 }
1419
1420 #[test]
1421 fn publication_rejects_dimension_cancellation_and_deletion_marker_without_state() {
1422 let dir = tempfile::tempdir().unwrap();
1423 let descriptor = descriptor("model-a", 2);
1424 let wrong = batch(&[1.0, 2.0, 3.0]);
1425 let error = publish_embedding_generation(
1426 dir.path(),
1427 request(&descriptor, source(1), &wrong, 20),
1428 VectorStoreLimits::default(),
1429 SearchCoordinationLimits::default(),
1430 || Ok(()),
1431 )
1432 .unwrap_err();
1433 assert!(matches!(error, SearchArtifactError::InvalidSelector { .. }));
1434 assert!(!dir.path().join("embeddings").exists());
1435
1436 let vectors = batch(&[1.0, 2.0]);
1437 let error = publish_embedding_generation(
1438 dir.path(),
1439 request(&descriptor, source(1), &vectors, 20),
1440 VectorStoreLimits::default(),
1441 SearchCoordinationLimits::default(),
1442 || Err(SearchArtifactError::Cancelled),
1443 )
1444 .unwrap_err();
1445 assert!(matches!(error, SearchArtifactError::Cancelled));
1446
1447 let compatibility = descriptor.compatibility_id().unwrap();
1448 std::fs::create_dir_all(dir.path().join("embeddings")).unwrap();
1449 std::fs::write(deletion_marker(dir.path(), compatibility), b"deleting").unwrap();
1450 let error = publish_embedding_generation(
1451 dir.path(),
1452 request(&descriptor, source(1), &vectors, 20),
1453 VectorStoreLimits::default(),
1454 SearchCoordinationLimits::default(),
1455 || Ok(()),
1456 )
1457 .unwrap_err();
1458 assert!(matches!(error, SearchArtifactError::InvalidSelector { .. }));
1459 assert!(
1460 current_embedding_generation(
1461 dir.path(),
1462 &descriptor,
1463 VectorStoreLimits::default(),
1464 || Ok(()),
1465 )
1466 .unwrap()
1467 .is_none()
1468 );
1469 }
1470
1471 #[test]
1472 fn deletion_removes_generation_and_alias_and_is_idempotent_across_reopen() {
1473 let empty = tempfile::tempdir().unwrap();
1474 let absent = EmbeddingCompatibilityId::from_hex(&"11".repeat(32)).unwrap();
1475 assert!(
1476 !delete_embedding_space_lineage(
1477 empty.path(),
1478 absent,
1479 EmbeddingSpaceCatalogLimits::default(),
1480 SearchCoordinationLimits::default(),
1481 || Ok(()),
1482 )
1483 .unwrap()
1484 );
1485 assert!(!empty.path().join("embeddings").exists());
1486
1487 let project = tempfile::tempdir().unwrap();
1488 let descriptor = descriptor("delete-me", 2);
1489 let vectors = batch(&[1.0, 2.0]);
1490 let published = publish(project.path(), &descriptor, source(1), &vectors, 20);
1491 let compatibility = descriptor.compatibility_id().unwrap();
1492 crate::bind_existing_embedding_space_catalog_entry(
1493 project.path(),
1494 "semantic",
1495 compatibility,
1496 false,
1497 EmbeddingSpaceCatalogLimits::default(),
1498 || Ok(()),
1499 )
1500 .unwrap();
1501 assert!(
1502 delete_embedding_space_lineage(
1503 project.path(),
1504 compatibility,
1505 EmbeddingSpaceCatalogLimits::default(),
1506 SearchCoordinationLimits::default(),
1507 || Ok(()),
1508 )
1509 .unwrap()
1510 );
1511 assert!(!published.publication().path.exists());
1512 assert!(!deletion_marker(project.path(), compatibility).exists());
1513 assert!(
1514 crate::read_embedding_space_catalog(
1515 project.path(),
1516 EmbeddingSpaceCatalogLimits::default(),
1517 || Ok(()),
1518 )
1519 .unwrap()
1520 .is_empty()
1521 );
1522 assert!(
1523 current_embedding_generation(
1524 project.path(),
1525 &descriptor,
1526 VectorStoreLimits::default(),
1527 || Ok(()),
1528 )
1529 .unwrap()
1530 .is_none()
1531 );
1532 assert!(
1533 !delete_embedding_space_lineage(
1534 project.path(),
1535 compatibility,
1536 EmbeddingSpaceCatalogLimits::default(),
1537 SearchCoordinationLimits::default(),
1538 || Ok(()),
1539 )
1540 .unwrap()
1541 );
1542 }
1543
1544 #[test]
1545 fn interrupted_deletion_retains_marker_and_retry_completes_without_resurrection() {
1546 let project = tempfile::tempdir().unwrap();
1547 let descriptor = descriptor("interrupted-delete", 2);
1548 let vectors = batch(&[1.0, 2.0]);
1549 let published = publish(project.path(), &descriptor, source(1), &vectors, 20);
1550 let compatibility = descriptor.compatibility_id().unwrap();
1551 let marker = deletion_marker(project.path(), compatibility);
1552 let root = space_root(project.path(), compatibility);
1553 let error = delete_embedding_space_lineage(
1554 project.path(),
1555 compatibility,
1556 EmbeddingSpaceCatalogLimits::default(),
1557 SearchCoordinationLimits::default(),
1558 || {
1559 if marker.exists() && root.exists() {
1560 Err(SearchArtifactError::Cancelled)
1561 } else {
1562 Ok(())
1563 }
1564 },
1565 )
1566 .unwrap_err();
1567 assert!(matches!(error, SearchArtifactError::Cancelled));
1568 assert!(marker.is_file());
1569 assert!(published.publication().path.is_dir());
1570
1571 assert!(
1572 delete_embedding_space_lineage(
1573 project.path(),
1574 compatibility,
1575 EmbeddingSpaceCatalogLimits::default(),
1576 SearchCoordinationLimits::default(),
1577 || Ok(()),
1578 )
1579 .unwrap()
1580 );
1581 assert!(!marker.exists());
1582 assert!(!root.exists());
1583 }
1584
1585 #[test]
1586 fn publication_filesystem_and_error_boundaries_are_fail_closed() {
1587 let project = tempfile::tempdir().unwrap();
1588 let compatibility = EmbeddingCompatibilityId::from_hex(&"22".repeat(32)).unwrap();
1589 let first = EmbeddingWriterLock::acquire(
1590 project.path(),
1591 compatibility,
1592 SearchCoordinationLimits::default(),
1593 &mut || Ok(()),
1594 )
1595 .unwrap();
1596 let limits = SearchCoordinationLimits {
1597 lock_timeout: Duration::ZERO,
1598 lock_poll_interval: Duration::ZERO,
1599 ..SearchCoordinationLimits::default()
1600 };
1601 assert!(matches!(
1602 EmbeddingWriterLock::acquire(project.path(), compatibility, limits, &mut || Ok(())),
1603 Err(SearchArtifactError::Lock { .. })
1604 ));
1605 drop(first);
1606
1607 let tree = project.path().join("tree");
1608 std::fs::create_dir(&tree).unwrap();
1609 std::fs::create_dir(tree.join("nested")).unwrap();
1610 std::fs::write(tree.join("nested/data"), b"payload").unwrap();
1611 sync_tree(&tree).unwrap();
1612 assert!(ensure_owned_directory(&tree).is_ok());
1613 assert!(ensure_existing_directory(&tree).is_ok());
1614 assert!(ensure_regular_file(&tree.join("nested/data")).is_ok());
1615 assert!(path_exists(&tree.join("nested/data")).unwrap());
1616 assert!(!path_exists(&tree.join("absent")).unwrap());
1617
1618 assert!(matches!(
1619 read_bounded_file(&tree.join("nested/data"), 1),
1620 Err(SearchArtifactError::ResourceExhausted { .. })
1621 ));
1622 assert!(matches!(
1623 hash_file(&tree.join("nested/data"), 1, &mut || Ok(())),
1624 Err(SearchArtifactError::ResourceExhausted { .. })
1625 ));
1626 assert!(matches!(
1627 hash_file(&tree.join("nested/data"), 100, &mut || Err(
1628 SearchArtifactError::Cancelled
1629 )),
1630 Err(SearchArtifactError::Cancelled)
1631 ));
1632 assert!(matches!(
1633 ensure_existing_directory(&tree.join("nested/data")),
1634 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1635 ));
1636 assert!(matches!(
1637 ensure_regular_file(&tree),
1638 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1639 ));
1640
1641 let mapped = primary_from(
1642 &tree,
1643 SearchArtifactError::InvalidSelector {
1644 field: "test",
1645 reason: "invalid".into(),
1646 },
1647 );
1648 assert!(matches!(
1649 mapped,
1650 SearchArtifactError::CorruptPrimaryVectors { .. }
1651 ));
1652 for error in [
1653 SearchArtifactError::Cancelled,
1654 SearchArtifactError::ResourceExhausted {
1655 resource: "test",
1656 limit: 1,
1657 },
1658 ] {
1659 assert!(matches!(
1660 primary_from(&tree, error),
1661 SearchArtifactError::Cancelled | SearchArtifactError::ResourceExhausted { .. }
1662 ));
1663 }
1664 assert!(matches!(
1665 io("test operation", &tree, std::io::Error::other("failure")),
1666 SearchArtifactError::Io { .. }
1667 ));
1668 }
1669
1670 #[cfg(unix)]
1671 #[test]
1672 fn publication_tree_rejects_symlinks() {
1673 use std::os::unix::fs::symlink;
1674
1675 let project = tempfile::tempdir().unwrap();
1676 let tree = project.path().join("tree");
1677 std::fs::create_dir(&tree).unwrap();
1678 symlink(project.path(), tree.join("link")).unwrap();
1679 assert!(matches!(
1680 sync_tree(&tree),
1681 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1682 ));
1683 assert!(matches!(
1684 ensure_owned_directory(&tree.join("link")),
1685 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1686 ));
1687
1688 std::fs::remove_file(tree.join("link")).unwrap();
1689 std::fs::write(tree.join(MANIFEST_FILE), b"manifest").unwrap();
1690 std::fs::write(tree.join(VECTOR_DATA_FILE), b"vectors").unwrap();
1691 symlink(project.path(), tree.join("link")).unwrap();
1692 assert!(matches!(
1693 validate_generation_layout(&tree),
1694 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1695 ));
1696 }
1697
1698 #[cfg(unix)]
1699 #[test]
1700 fn symlinked_embedding_ancestor_fails_closed() {
1701 use std::os::unix::fs::symlink;
1702
1703 let project = tempfile::tempdir().unwrap();
1704 let external = tempfile::tempdir().unwrap();
1705 symlink(external.path(), project.path().join("embeddings")).unwrap();
1706 let descriptor = descriptor("model-a", 2);
1707 let compatibility = descriptor.compatibility_id().unwrap().to_hex();
1708 std::fs::create_dir_all(external.path().join("spaces").join(compatibility)).unwrap();
1709 assert!(matches!(
1710 current_embedding_generation(
1711 project.path(),
1712 &descriptor,
1713 VectorStoreLimits::default(),
1714 || Ok(()),
1715 ),
1716 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1717 ));
1718 }
1719
1720 #[test]
1721 fn publication_dimension_mismatch_fails_before_creating_storage() {
1722 let project = tempfile::tempdir().unwrap();
1723 let descriptor = descriptor("model-a", 3);
1724 let vectors = batch(&[1.0, 2.0]);
1725
1726 let error = publish_embedding_generation(
1727 project.path(),
1728 request(&descriptor, source(1), &vectors, 20),
1729 VectorStoreLimits::default(),
1730 SearchCoordinationLimits::default(),
1731 || Ok(()),
1732 )
1733 .unwrap_err();
1734
1735 assert!(matches!(
1736 error,
1737 SearchArtifactError::InvalidSelector {
1738 field: "embedding batch",
1739 ..
1740 }
1741 ));
1742 assert!(!project.path().join("embeddings").exists());
1743 }
1744
1745 #[test]
1746 fn public_reopen_distinguishes_absent_space_and_absent_active_pointer() {
1747 let dir = tempfile::tempdir().unwrap();
1748 let descriptor = descriptor("absent", 2);
1749 assert!(
1750 current_embedding_generation(
1751 dir.path(),
1752 &descriptor,
1753 VectorStoreLimits::default(),
1754 || Ok(())
1755 )
1756 .unwrap()
1757 .is_none()
1758 );
1759
1760 let compatibility_id = descriptor.compatibility_id().unwrap();
1761 let root = space_root(dir.path(), compatibility_id);
1762 std::fs::create_dir_all(&root).unwrap();
1763 assert!(
1764 current_embedding_generation(
1765 dir.path(),
1766 &descriptor,
1767 VectorStoreLimits::default(),
1768 || Ok(())
1769 )
1770 .unwrap()
1771 .is_none()
1772 );
1773 }
1774
1775 #[test]
1776 fn wave10_generation_layout_rejects_nonfiles_and_excess_inventory() {
1777 for kind in ["directory", "excess"] {
1778 let root = tempfile::tempdir().unwrap();
1779 std::fs::write(root.path().join(MANIFEST_FILE), b"manifest").unwrap();
1780 std::fs::write(root.path().join(VECTOR_DATA_FILE), b"vectors").unwrap();
1781 if kind == "directory" {
1782 std::fs::create_dir(root.path().join("unexpected")).unwrap();
1783 } else {
1784 std::fs::write(root.path().join("unexpected"), b"caller").unwrap();
1785 }
1786 assert!(matches!(
1787 validate_generation_layout(root.path()),
1788 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1789 ));
1790 }
1791 }
1792
1793 #[cfg(unix)]
1794 #[test]
1795 fn wave10_generation_sync_rejects_special_files() {
1796 use std::os::unix::net::UnixListener;
1797
1798 let root = tempfile::Builder::new()
1799 .prefix("gf")
1800 .tempdir_in("/tmp")
1801 .unwrap();
1802 let _socket = UnixListener::bind(root.path().join("socket")).unwrap();
1803 assert!(matches!(
1804 sync_tree(root.path()),
1805 Err(SearchArtifactError::CorruptPrimaryVectors { .. })
1806 ));
1807 }
1808
1809 #[test]
1810 fn wave13_embedding_path_helpers_fail_closed_on_regular_file_ancestors() {
1811 let root = tempfile::tempdir().unwrap();
1812 let ancestor = root.path().join("regular-file");
1813 std::fs::write(&ancestor, b"caller data").unwrap();
1814 let child = ancestor.join("child");
1815
1816 assert!(matches!(
1817 ensure_owned_directory(&child),
1818 Err(SearchArtifactError::Io { .. })
1819 ));
1820 assert!(matches!(
1821 path_exists(&child),
1822 Err(SearchArtifactError::Io { .. })
1823 ));
1824 assert!(matches!(
1825 ensure_existing_directory(&child),
1826 Err(SearchArtifactError::Io { .. })
1827 ));
1828 assert!(matches!(
1829 ensure_regular_file(&child),
1830 Err(SearchArtifactError::Io { .. })
1831 ));
1832 }
1833}