1use serde::{Deserialize, Serialize};
15use std::collections::HashMap;
16use std::io::Write;
17use std::path::Path;
18
19use crate::dsl::{BinaryIndexType, Schema, VectorIndexType};
20use crate::error::{Error, Result};
21
22pub const INDEX_META_FILENAME: &str = "metadata.json";
24const INDEX_META_TMP_FILENAME: &str = "metadata.json.tmp";
26
27pub const INDEX_META_FORMAT_VERSION: u32 = 9;
34
35pub const OLDEST_MIGRATABLE_FORMAT_VERSION: u32 = 6;
47
48pub(crate) const MAX_TRAINED_ARTIFACT_BYTES: usize = 512 * 1024 * 1024;
53
54#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
56pub enum VectorIndexState {
57 #[default]
59 Flat,
60 Built {
62 vector_count: usize,
64 num_clusters: usize,
66 },
67}
68
69fn is_zero(value: &u32) -> bool {
70 *value == 0
71}
72
73fn default_true() -> bool {
74 true
75}
76
77#[derive(Debug, Clone, Serialize, Deserialize)]
80pub struct SegmentMetaInfo {
81 #[serde(default, skip_serializing_if = "Option::is_none")]
83 pub deletions: Option<crate::segment::DeletionMeta>,
84 pub num_docs: u32,
86 pub ancestors: Vec<String>,
88 pub generation: u32,
90 #[serde(default)]
94 pub reordered: bool,
95 #[serde(default = "default_true")]
100 pub bp_converged: bool,
101 #[serde(default)]
105 pub bp_unconverged_passes: u32,
106 #[serde(default)]
108 pub seismic_pending_terms: u32,
109 #[serde(default)]
112 pub seismic_maintenance_passes: u32,
113 #[serde(default, skip_serializing_if = "is_zero")]
116 pub seismic_no_progress_passes: u32,
117 #[serde(default)]
119 pub ann_fragmented: bool,
120}
121
122impl SegmentMetaInfo {
123 pub fn num_deleted_docs(&self) -> u32 {
124 self.deletions.as_ref().map_or(0, |meta| meta.num_deleted)
125 }
126
127 pub fn num_live_docs(&self) -> u32 {
128 self.num_docs - self.num_deleted_docs()
129 }
130
131 pub fn deleted_ratio(&self) -> f64 {
132 if self.num_docs == 0 {
133 0.0
134 } else {
135 f64::from(self.num_deleted_docs()) / f64::from(self.num_docs)
136 }
137 }
138}
139
140#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
142#[serde(tag = "kind", content = "index", rename_all = "snake_case")]
143pub enum VectorFieldIndexType {
144 Float(VectorIndexType),
145 Binary(BinaryIndexType),
146}
147
148impl From<VectorIndexType> for VectorFieldIndexType {
149 fn from(value: VectorIndexType) -> Self {
150 Self::Float(value)
151 }
152}
153
154impl From<BinaryIndexType> for VectorFieldIndexType {
155 fn from(value: BinaryIndexType) -> Self {
156 Self::Binary(value)
157 }
158}
159
160#[derive(Debug, Clone, Serialize, Deserialize)]
161pub struct FieldVectorMeta {
162 pub field_id: u32,
164 pub index_type: VectorFieldIndexType,
166 pub state: VectorIndexState,
168 #[serde(skip_serializing_if = "Option::is_none")]
170 pub centroids_file: Option<String>,
171 #[serde(skip_serializing_if = "Option::is_none")]
175 pub codebook_file: Option<String>,
176 #[serde(default, skip_serializing_if = "Option::is_none")]
179 pub artifact_generation: Option<u64>,
180 #[serde(default, skip_serializing_if = "Option::is_none")]
183 pub artifact_id: Option<u64>,
184}
185
186#[derive(Debug, Clone, Serialize, Deserialize)]
188pub struct IndexMetadata {
189 pub version: u32,
191 #[serde(default)]
195 pub publication_generation: u64,
196 pub schema: Schema,
198 #[serde(default)]
201 pub segment_metas: HashMap<String, SegmentMetaInfo>,
202 #[serde(default)]
204 pub vector_fields: HashMap<u32, FieldVectorMeta>,
205 #[serde(default)]
212 pub total_vectors: usize,
213}
214
215impl IndexMetadata {
216 #[cfg(feature = "native")]
218 pub(crate) fn owned_ids(&self) -> Vec<String> {
219 self.segment_metas
220 .keys()
221 .cloned()
222 .chain(
223 self.segment_metas
224 .values()
225 .filter_map(|info| info.deletions.as_ref().map(|d| d.id.clone())),
226 )
227 .collect()
228 }
229
230 #[cfg(feature = "native")]
231 pub(crate) fn owns_id(&self, id: &str) -> bool {
232 self.has_segment(id)
233 || self
234 .segment_metas
235 .values()
236 .any(|info| info.deletions.as_ref().is_some_and(|d| d.id == id))
237 }
238 pub fn new(schema: Schema) -> Self {
240 Self {
241 version: INDEX_META_FORMAT_VERSION,
242 publication_generation: 0,
243 schema,
244 segment_metas: HashMap::new(),
245 vector_fields: HashMap::new(),
246 total_vectors: 0,
247 }
248 }
249
250 pub fn segment_ids(&self) -> Vec<String> {
252 let mut ids: Vec<String> = self.segment_metas.keys().cloned().collect();
253 ids.sort();
254 ids
255 }
256
257 pub fn add_segment(&mut self, segment_id: String, num_docs: u32) {
259 self.segment_metas.insert(
260 segment_id,
261 SegmentMetaInfo {
262 deletions: None,
263 num_docs,
264 ancestors: Vec::new(),
265 generation: 0,
266 reordered: false,
267 bp_converged: true,
268 bp_unconverged_passes: 0,
269 seismic_pending_terms: 0,
270 seismic_maintenance_passes: 0,
271 seismic_no_progress_passes: 0,
272 ann_fragmented: false,
273 },
274 );
275 }
276
277 pub fn add_merged_segment(
279 &mut self,
280 segment_id: String,
281 num_docs: u32,
282 ancestors: Vec<String>,
283 generation: u32,
284 reordered: bool,
285 bp_converged: bool,
286 ) {
287 self.add_segment_meta(
288 segment_id,
289 SegmentMetaInfo {
290 deletions: None,
291 num_docs,
292 ancestors,
293 generation,
294 reordered,
295 bp_converged,
296 bp_unconverged_passes: 0,
297 seismic_pending_terms: 0,
298 seismic_maintenance_passes: 0,
299 seismic_no_progress_passes: 0,
300 ann_fragmented: false,
301 },
302 );
303 }
304
305 pub(crate) fn add_segment_meta(&mut self, segment_id: String, info: SegmentMetaInfo) {
309 self.segment_metas.insert(segment_id, info);
310 }
311
312 pub fn remove_segment(&mut self, segment_id: &str) {
314 self.segment_metas.remove(segment_id);
315 }
316
317 pub fn has_segment(&self, segment_id: &str) -> bool {
319 self.segment_metas.contains_key(segment_id)
320 }
321
322 pub fn segment_doc_count(&self, segment_id: &str) -> Option<u32> {
324 self.segment_metas.get(segment_id).map(|m| m.num_docs)
325 }
326
327 pub fn is_field_built(&self, field_id: u32) -> bool {
329 self.vector_fields
330 .get(&field_id)
331 .map(|f| matches!(f.state, VectorIndexState::Built { .. }))
332 .unwrap_or(false)
333 }
334
335 pub fn get_field_meta(&self, field_id: u32) -> Option<&FieldVectorMeta> {
337 self.vector_fields.get(&field_id)
338 }
339
340 pub fn init_field(&mut self, field_id: u32, index_type: impl Into<VectorFieldIndexType>) {
342 let index_type = index_type.into();
343 self.vector_fields
344 .entry(field_id)
345 .or_insert(FieldVectorMeta {
346 field_id,
347 index_type,
348 state: VectorIndexState::Flat,
349 centroids_file: None,
350 codebook_file: None,
351 artifact_generation: None,
352 artifact_id: None,
353 });
354 }
355
356 pub fn mark_field_built(
358 &mut self,
359 field_id: u32,
360 vector_count: usize,
361 num_clusters: usize,
362 centroids_file: String,
363 codebook_file: Option<String>,
364 ) {
365 if let Some(field) = self.vector_fields.get_mut(&field_id) {
366 field.state = VectorIndexState::Built {
367 vector_count,
368 num_clusters,
369 };
370 field.centroids_file = Some(centroids_file);
371 field.codebook_file = codebook_file;
372 field.artifact_generation = None;
373 field.artifact_id = None;
374 self.refresh_total_vectors();
375 }
376 }
377
378 pub fn mark_scann_field_built(
382 &mut self,
383 field_id: u32,
384 vector_count: usize,
385 num_leaves: usize,
386 artifact_file: String,
387 artifact_generation: u64,
388 artifact_id: u64,
389 ) -> Result<()> {
390 if artifact_generation == 0 || artifact_id == 0 {
391 return Err(Error::Corruption(format!(
392 "ScaNN field {field_id} cannot publish a zero generation or artifact fingerprint"
393 )));
394 }
395 let field = self.vector_fields.get_mut(&field_id).ok_or_else(|| {
396 Error::Corruption(format!(
397 "ScaNN field {field_id} must be initialized before it is marked built"
398 ))
399 })?;
400 if !matches!(
401 field.index_type,
402 VectorFieldIndexType::Float(VectorIndexType::Scann)
403 | VectorFieldIndexType::Binary(BinaryIndexType::Scann)
404 ) {
405 return Err(Error::Corruption(format!(
406 "field {field_id} is not configured as ScaNN"
407 )));
408 }
409 field.state = VectorIndexState::Built {
410 vector_count,
411 num_clusters: num_leaves,
412 };
413 field.centroids_file = Some(artifact_file);
414 field.codebook_file = None;
415 field.artifact_generation = Some(artifact_generation);
416 field.artifact_id = Some(artifact_id);
417 self.refresh_total_vectors();
418 Ok(())
419 }
420
421 pub(crate) fn refresh_total_vectors(&mut self) {
427 self.total_vectors = self
428 .vector_fields
429 .values()
430 .filter_map(|field| match field.state {
431 VectorIndexState::Built { vector_count, .. } => Some(vector_count),
432 VectorIndexState::Flat => None,
433 })
434 .fold(0usize, usize::saturating_add);
435 }
436
437 pub fn should_build_field(&self, field_id: u32, threshold: usize) -> bool {
439 if self.is_field_built(field_id) {
441 return false;
442 }
443 self.total_vectors >= threshold
445 }
446
447 pub async fn load<D: crate::directories::Directory>(dir: &D) -> Result<Self> {
457 let (meta, migrated_from) = Self::load_reporting_migration(dir).await?;
458 if let Some(from) = migrated_from {
459 log::warn!(
460 "[metadata_migration] metadata.json format version {from} upgraded in memory \
461 to {INDEX_META_FORMAT_VERSION} ({} segment(s)); the upgrade is persisted \
462 when a writer opens this index. Segments from builds before 1.8.125 fail at \
463 segment open and need a rebuild; segments without .rowstats cannot be \
464 compacted until they are merged",
465 meta.segment_metas.len()
466 );
467 }
468 Ok(meta)
469 }
470
471 pub async fn load_reporting_migration<D: crate::directories::Directory>(
477 dir: &D,
478 ) -> Result<(Self, Option<u32>)> {
479 let path = Path::new(INDEX_META_FILENAME);
480 match dir.open_read(path).await {
481 Ok(slice) => {
482 let bytes = slice.read_bytes().await?;
483 Self::deserialize_versioned(bytes.as_slice())
484 }
485 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
486 let tmp_path = Path::new(INDEX_META_TMP_FILENAME);
488 let slice = dir.open_read(tmp_path).await?;
489 let bytes = slice.read_bytes().await?;
490 let recovered = Self::deserialize_versioned(bytes.as_slice())?;
491 log::warn!("Recovered metadata from temp file (previous crash during save)");
492 Ok(recovered)
493 }
494 Err(e) => Err(Error::Io(e)),
495 }
496 }
497
498 fn deserialize_versioned(bytes: &[u8]) -> Result<(Self, Option<u32>)> {
504 let mut meta: Self =
505 serde_json::from_slice(bytes).map_err(|e| Error::Serialization(e.to_string()))?;
506 meta.schema.validate()?;
507 crate::dsl::reject_removed_vector_index_types(&meta.schema).map_err(Error::Schema)?;
508 let migrated_from = if meta.version == INDEX_META_FORMAT_VERSION {
509 None
510 } else if (OLDEST_MIGRATABLE_FORMAT_VERSION..INDEX_META_FORMAT_VERSION)
511 .contains(&meta.version)
512 {
513 let from = meta.version;
514 meta.version = INDEX_META_FORMAT_VERSION;
515 Some(from)
516 } else {
517 return Err(Error::Corruption(format!(
518 "metadata.json format version {} is incompatible with required version {} \
519 (formats {}..={} are upgraded on open); \
520 rebuild and republish the index with this Summa version",
521 meta.version,
522 INDEX_META_FORMAT_VERSION,
523 OLDEST_MIGRATABLE_FORMAT_VERSION,
524 INDEX_META_FORMAT_VERSION - 1
525 )));
526 };
527 let mut deletion_ids = std::collections::HashSet::new();
528 for info in meta.segment_metas.values() {
529 if let Some(deletion) = &info.deletions
530 && (deletion.num_deleted == 0
531 || deletion.num_deleted > info.num_docs
532 || crate::segment::SegmentId::from_hex(&deletion.id).is_none()
533 || meta.segment_metas.contains_key(&deletion.id)
534 || !deletion_ids.insert(&deletion.id))
535 {
536 return Err(Error::Corruption(
537 "invalid or aliased deletion metadata".into(),
538 ));
539 }
540 }
541 Ok((meta, migrated_from))
542 }
543
544 #[cfg(any(feature = "native", feature = "wasm"))]
548 pub(crate) async fn load_persisting_migration<D: crate::directories::DirectoryWriter>(
549 dir: &D,
550 ) -> Result<Self> {
551 let (meta, migrated_from) = Self::load_reporting_migration(dir).await?;
552 if let Some(from) = migrated_from {
553 meta.save(dir).await?;
554 log::warn!(
555 "[metadata_migration] metadata.json format version {from} upgraded to \
556 {INDEX_META_FORMAT_VERSION} and persisted ({} segment(s)); builds that do not \
557 support format {INDEX_META_FORMAT_VERSION} can no longer open this index. Segments from builds before 1.8.125 \
558 fail at segment open and need a rebuild; segments without .rowstats cannot \
559 be compacted until they are merged",
560 meta.segment_metas.len()
561 );
562 }
563 Ok(meta)
564 }
565
566 pub async fn save<D: crate::directories::DirectoryWriter>(&self, dir: &D) -> Result<()> {
571 let bytes = self.serialize_to_bytes()?;
572 Self::save_bytes(dir, &bytes).await
573 }
574
575 pub fn serialize_to_bytes(&self) -> Result<Vec<u8>> {
578 serde_json::to_vec_pretty(self).map_err(|e| Error::Serialization(e.to_string()))
579 }
580
581 pub async fn save_bytes<D: crate::directories::DirectoryWriter>(
586 dir: &D,
587 bytes: &[u8],
588 ) -> Result<()> {
589 let tmp_path = Path::new(INDEX_META_TMP_FILENAME);
590 let final_path = Path::new(INDEX_META_FILENAME);
591 let mut writer = dir.streaming_writer(tmp_path).await.map_err(Error::Io)?;
596 writer.write_all(bytes).map_err(Error::Io)?;
597 writer.finish().map_err(Error::Io)?;
598 dir.rename(tmp_path, final_path).await.map_err(Error::Io)?;
605 if let Err(error) = dir.sync().await {
606 log::error!(
607 "[metadata] directory fsync failed after committed rename: {}. \
608 Continuing with the renamed generation; crash durability is not guaranteed",
609 error,
610 );
611 }
612 Ok(())
613 }
614
615 #[cfg_attr(not(feature = "native"), allow(dead_code))]
617 pub(crate) async fn try_load_trained_from_fields<D: crate::directories::Directory>(
618 vector_fields: &HashMap<u32, FieldVectorMeta>,
619 schema: &Schema,
620 dir: &D,
621 ) -> Result<Option<crate::segment::TrainedVectorStructures>> {
622 Self::load_trained_from_fields_impl(vector_fields, schema, dir).await
623 }
624
625 async fn load_trained_from_fields_impl<D: crate::directories::Directory>(
635 vector_fields: &HashMap<u32, FieldVectorMeta>,
636 schema: &Schema,
637 dir: &D,
638 ) -> Result<Option<crate::segment::TrainedVectorStructures>> {
639 use std::sync::Arc;
640
641 let mut centroids = rustc_hash::FxHashMap::default();
642 let mut binary_quantizers = rustc_hash::FxHashMap::default();
643 let mut scann_artifacts = rustc_hash::FxHashMap::default();
644
645 let mut built_fields: Vec<_> = vector_fields
646 .iter()
647 .filter(|(_, meta)| matches!(meta.state, VectorIndexState::Built { .. }))
648 .collect();
649 built_fields.sort_unstable_by_key(|(field_id, _)| **field_id);
650
651 log::debug!(
652 "[trained] index={} loading trained structures, dense_vector_fields={:?}",
653 schema.index_label(),
654 vector_fields.keys().collect::<Vec<_>>()
655 );
656
657 for (field_id, field_meta) in built_fields {
658 log::debug!(
659 "[trained] index={} field {} state={:?} centroids_file={:?} codebook_file={:?}",
660 schema.index_label(),
661 field_id,
662 field_meta.state,
663 field_meta.centroids_file,
664 field_meta.codebook_file,
665 );
666 if field_meta.field_id != *field_id {
667 return Err(Error::Corruption(format!(
668 "trained vector metadata key {field_id} contains field_id {}",
669 field_meta.field_id
670 )));
671 }
672
673 let expected_clusters = match field_meta.state {
674 VectorIndexState::Built { num_clusters, .. } if num_clusters > 0 => num_clusters,
675 VectorIndexState::Built { .. } => {
676 return Err(Error::Corruption(format!(
677 "trained vector metadata field {field_id} has zero clusters"
678 )));
679 }
680 VectorIndexState::Flat => unreachable!("built_fields contains only Built entries"),
681 };
682
683 let centroids_file = field_meta.centroids_file.as_deref().ok_or_else(|| {
684 Error::Corruption(format!(
685 "trained vector metadata field {field_id} is Built but has no centroids_file"
686 ))
687 })?;
688 match field_meta.index_type {
689 VectorFieldIndexType::Float(VectorIndexType::IvfPq) => {
690 return Err(Error::Corruption(format!(
691 "field {field_id} was trained as IVF-PQ, which is no longer \
692 supported; recreate the index with `ivf_tq` and reindex \
693 (docs/turboquant-quantization.md)"
694 )));
695 }
696 VectorFieldIndexType::Float(index_type @ VectorIndexType::IvfTq) => {
697 let entry = schema
698 .get_field_entry(crate::dsl::Field(*field_id))
699 .ok_or_else(|| {
700 Error::Corruption(format!(
701 "trained vector metadata references missing field {field_id}"
702 ))
703 })?;
704 let schema_config = entry
705 .dense_vector_config
706 .as_ref()
707 .filter(|_| entry.field_type == crate::dsl::FieldType::DenseVector)
708 .ok_or_else(|| {
709 Error::Corruption(format!(
710 "trained vector metadata field {field_id} is not a float dense field"
711 ))
712 })?;
713 if schema_config.index_type != index_type {
714 return Err(Error::Corruption(format!(
715 "trained vector metadata field {field_id} uses {index_type:?}, schema requires {:?}",
716 schema_config.index_type
717 )));
718 }
719 let c: crate::structures::CoarseCentroids =
720 load_trained_artifact(dir, *field_id, "centroids", centroids_file).await?;
721 let expected_dim = schema_config.dim;
722 let actual_clusters = c.num_clusters as usize;
723 let expected_values =
724 actual_clusters.checked_mul(expected_dim).ok_or_else(|| {
725 Error::Corruption(format!(
726 "trained centroid dimensions overflow for field {field_id}"
727 ))
728 })?;
729 if actual_clusters == 0
730 || actual_clusters > expected_clusters
731 || c.dim == 0
732 || c.dim != expected_dim
733 || c.centroids.len() != expected_values
734 || c.centroids.iter().any(|value| !value.is_finite())
735 {
736 return Err(Error::Corruption(format!(
737 "trained centroids for field {field_id} do not match metadata/schema"
738 )));
739 }
740 if !crate::structures::is_ivf_tq_cosine_generation(c.version) {
741 return Err(Error::Corruption(format!(
742 "trained IVF-TQ centroids for field {field_id} use an \
743 unsupported legacy generation; rebuild the index"
744 )));
745 }
746 c.validate_routing(schema_config.ivf_routing)
747 .map_err(|error| {
748 Error::Corruption(format!(
749 "invalid trained centroid routing for field {field_id}: {error}"
750 ))
751 })?;
752 let _ = index_type;
755 if field_meta.codebook_file.is_some() {
756 return Err(Error::Corruption(format!(
757 "trained IVF-TQ field {field_id} unexpectedly references a codebook file"
758 )));
759 }
760 centroids.insert(*field_id, Arc::new(c));
761 }
762 VectorFieldIndexType::Binary(BinaryIndexType::Ivf) => {
763 let entry = schema
764 .get_field_entry(crate::dsl::Field(*field_id))
765 .ok_or_else(|| {
766 Error::Corruption(format!(
767 "trained vector metadata references missing field {field_id}"
768 ))
769 })?;
770 let schema_config = entry
771 .binary_dense_vector_config
772 .as_ref()
773 .filter(|config| {
774 entry.field_type == crate::dsl::FieldType::BinaryDenseVector
775 && config.index_type == BinaryIndexType::Ivf
776 })
777 .ok_or_else(|| {
778 Error::Corruption(format!(
779 "trained vector metadata field {field_id} is not a binary IVF field"
780 ))
781 })?;
782 let quantizer: crate::structures::BinaryCoarseQuantizer =
783 load_trained_artifact(dir, *field_id, "binary centroids", centroids_file)
784 .await?;
785 quantizer.validate().map_err(|error| {
786 Error::Corruption(format!(
787 "invalid binary coarse quantizer for field {field_id}: {error}"
788 ))
789 })?;
790 let actual_clusters = quantizer.num_clusters as usize;
791 if actual_clusters > expected_clusters
792 || schema_config.dim != quantizer.dim_bits
793 {
794 return Err(Error::Corruption(format!(
795 "binary coarse quantizer for field {field_id} does not match metadata/schema"
796 )));
797 }
798 quantizer
799 .validate_routing(schema_config.ivf_routing)
800 .map_err(|error| {
801 Error::Corruption(format!(
802 "invalid binary centroid routing for field {field_id}: {error}"
803 ))
804 })?;
805 binary_quantizers.insert(*field_id, Arc::new(quantizer));
806 }
807 VectorFieldIndexType::Float(VectorIndexType::Scann)
808 | VectorFieldIndexType::Binary(BinaryIndexType::Scann) => {
809 if field_meta.codebook_file.is_some() {
810 return Err(Error::Corruption(format!(
811 "trained ScaNN field {field_id} unexpectedly references a separate codebook file"
812 )));
813 }
814 let expected_generation = field_meta.artifact_generation.ok_or_else(|| {
815 Error::Corruption(format!(
816 "trained ScaNN field {field_id} has no artifact generation"
817 ))
818 })?;
819 let expected_artifact_id = field_meta.artifact_id.ok_or_else(|| {
820 Error::Corruption(format!(
821 "trained ScaNN field {field_id} has no artifact fingerprint"
822 ))
823 })?;
824 validate_trained_artifact_path(
825 field_id.to_owned(),
826 "ScaNN artifact",
827 centroids_file,
828 )?;
829 let path = Path::new(centroids_file);
830 let slice = dir.open_read(path).await.map_err(|error| {
831 Error::Corruption(format!(
832 "failed to open trained ScaNN artifact '{centroids_file}' for field {field_id}: {error}"
833 ))
834 })?;
835 let raw = slice.read_bytes().await.map_err(|error| {
836 Error::Corruption(format!(
837 "failed to map trained ScaNN artifact '{centroids_file}' for field {field_id}: {error}"
838 ))
839 })?;
840 let artifact =
841 crate::segment::ScannTrainedArtifactBytes::open(raw).map_err(|error| {
842 Error::Corruption(format!(
843 "invalid trained ScaNN artifact for field {field_id}: {error}"
844 ))
845 })?;
846 if artifact.generation() != expected_generation
847 || artifact.artifact_id() != expected_artifact_id
848 || artifact.config().num_leaves as usize != expected_clusters
849 {
850 return Err(Error::Corruption(format!(
851 "trained ScaNN artifact for field {field_id} does not match metadata"
852 )));
853 }
854 let entry = schema
855 .get_field_entry(crate::dsl::Field(*field_id))
856 .ok_or_else(|| {
857 Error::Corruption(format!(
858 "trained vector metadata references missing field {field_id}"
859 ))
860 })?;
861 let schema_matches = match field_meta.index_type {
862 VectorFieldIndexType::Float(VectorIndexType::Scann) => entry
863 .dense_vector_config
864 .as_ref()
865 .filter(|_| entry.field_type == crate::dsl::FieldType::DenseVector)
866 .is_some_and(|config| {
867 config.index_type == VectorIndexType::Scann
868 && config.dim == artifact.config().dimension as usize
869 && scann_explicit_geometry_matches(
870 config.num_clusters,
871 config.tree_levels,
872 artifact.config().num_leaves as usize,
873 artifact.config().tree_levels,
874 )
875 && matches!(
876 artifact.config().encoding,
877 crate::structures::vector::scann::ScannEncoding::AsymmetricHash { .. }
878 )
879 }),
880 VectorFieldIndexType::Binary(BinaryIndexType::Scann) => entry
881 .binary_dense_vector_config
882 .as_ref()
883 .filter(|_| {
884 entry.field_type == crate::dsl::FieldType::BinaryDenseVector
885 })
886 .is_some_and(|config| {
887 config.index_type == BinaryIndexType::Scann
888 && config.dim == artifact.config().dimension as usize
889 && scann_explicit_geometry_matches(
890 config.num_clusters,
891 config.tree_levels,
892 artifact.config().num_leaves as usize,
893 artifact.config().tree_levels,
894 )
895 && artifact.config().encoding
896 == crate::structures::vector::scann::ScannEncoding::BinaryHamming
897 }),
898 _ => false,
899 };
900 if !schema_matches {
901 return Err(Error::Corruption(format!(
902 "trained ScaNN artifact for field {field_id} does not match schema geometry/encoding"
903 )));
904 }
905 scann_artifacts.insert(*field_id, Arc::new(artifact));
906 }
907 unsupported => {
908 return Err(Error::Corruption(format!(
909 "field {field_id} is Built for {unsupported:?}, which has no global IVF artifacts"
910 )));
911 }
912 }
913 }
914
915 if centroids.is_empty() && binary_quantizers.is_empty() && scann_artifacts.is_empty() {
916 Ok(None)
917 } else {
918 let trained = crate::segment::TrainedVectorStructures {
919 #[cfg(feature = "native")]
920 _ann_pins: Default::default(),
921 centroids,
922 binary_quantizers,
923 scann_artifacts,
924 };
925 #[cfg(feature = "native")]
926 let trained = {
927 let mut trained = trained;
928 trained.pin_ann_structures(crate::segment::pin::pin_policy());
929 trained
930 };
931 Ok(Some(trained))
932 }
933 }
934}
935
936fn scann_explicit_geometry_matches(
940 configured_leaves: Option<usize>,
941 configured_levels: Option<u8>,
942 resolved_leaves: usize,
943 resolved_levels: u8,
944) -> bool {
945 configured_leaves.is_none_or(|leaves| leaves == resolved_leaves)
946 && configured_levels.is_none_or(|levels| levels == resolved_levels)
947}
948
949fn validate_trained_artifact_path(field_id: u32, kind: &str, filename: &str) -> Result<()> {
950 use std::path::Component;
951
952 let path = Path::new(filename);
953 if filename.is_empty()
954 || path.is_absolute()
955 || path.components().any(|component| {
956 matches!(
957 component,
958 Component::ParentDir | Component::RootDir | Component::Prefix(_)
959 )
960 })
961 {
962 return Err(Error::Corruption(format!(
963 "trained {kind} path for field {field_id} is not a safe relative path: '{filename}'"
964 )));
965 }
966 Ok(())
967}
968
969async fn load_trained_artifact<T, D>(
970 dir: &D,
971 field_id: u32,
972 kind: &str,
973 filename: &str,
974) -> Result<T>
975where
976 T: serde::de::DeserializeOwned,
977 D: crate::directories::Directory,
978{
979 validate_trained_artifact_path(field_id, kind, filename)?;
980 let path = Path::new(filename);
981 let file_size = dir.file_size(path).await.map_err(|error| {
982 Error::Corruption(format!(
983 "failed to stat trained {kind} '{filename}' for field {field_id}: {error}"
984 ))
985 })?;
986 validate_trained_artifact_size(field_id, kind, filename, file_size)?;
987 let slice = dir.open_read(path).await.map_err(|error| {
988 Error::Corruption(format!(
989 "failed to open trained {kind} '{filename}' for field {field_id}: {error}"
990 ))
991 })?;
992 validate_trained_artifact_size(field_id, kind, filename, slice.len())?;
993 let bytes = slice.read_bytes().await.map_err(|error| {
994 Error::Corruption(format!(
995 "failed to read trained {kind} '{filename}' for field {field_id}: {error}"
996 ))
997 })?;
998 let (artifact, consumed) = bincode::serde::decode_from_slice::<T, _>(
999 bytes.as_slice(),
1000 bincode::config::standard().with_limit::<MAX_TRAINED_ARTIFACT_BYTES>(),
1001 )
1002 .map_err(|error| {
1003 Error::Corruption(format!(
1004 "failed to deserialize trained {kind} '{filename}' for field {field_id}: {error}"
1005 ))
1006 })?;
1007 if consumed != bytes.len() {
1008 return Err(Error::Corruption(format!(
1009 "trained {kind} '{filename}' for field {field_id} has {} trailing bytes",
1010 bytes.len() - consumed
1011 )));
1012 }
1013 Ok(artifact)
1014}
1015
1016fn validate_trained_artifact_size(
1017 field_id: u32,
1018 kind: &str,
1019 filename: &str,
1020 file_size: u64,
1021) -> Result<()> {
1022 if file_size > MAX_TRAINED_ARTIFACT_BYTES as u64 {
1023 return Err(Error::Corruption(format!(
1024 "trained {kind} '{filename}' for field {field_id} is {file_size} bytes, \
1025 exceeding the {MAX_TRAINED_ARTIFACT_BYTES}-byte safety limit"
1026 )));
1027 }
1028 Ok(())
1029}
1030
1031#[cfg(test)]
1032mod tests {
1033 use super::*;
1034 use crate::directories::DirectoryWriter;
1035
1036 #[derive(Clone, Default)]
1037 struct SyncFailDirectory(crate::directories::RamDirectory);
1038
1039 #[async_trait::async_trait]
1040 impl crate::directories::Directory for SyncFailDirectory {
1041 async fn exists(&self, path: &Path) -> std::io::Result<bool> {
1042 self.0.exists(path).await
1043 }
1044
1045 async fn file_size(&self, path: &Path) -> std::io::Result<u64> {
1046 self.0.file_size(path).await
1047 }
1048
1049 async fn open_read(&self, path: &Path) -> std::io::Result<crate::directories::FileHandle> {
1050 self.0.open_read(path).await
1051 }
1052
1053 async fn read_range(
1054 &self,
1055 path: &Path,
1056 range: std::ops::Range<u64>,
1057 ) -> std::io::Result<crate::directories::OwnedBytes> {
1058 self.0.read_range(path, range).await
1059 }
1060
1061 async fn list_files(&self, prefix: &Path) -> std::io::Result<Vec<std::path::PathBuf>> {
1062 self.0.list_files(prefix).await
1063 }
1064
1065 async fn open_lazy(&self, path: &Path) -> std::io::Result<crate::directories::FileHandle> {
1066 self.0.open_lazy(path).await
1067 }
1068 }
1069
1070 #[async_trait::async_trait]
1071 impl crate::directories::DirectoryWriter for SyncFailDirectory {
1072 async fn write(&self, path: &Path, data: &[u8]) -> std::io::Result<()> {
1073 self.0.write(path, data).await
1074 }
1075
1076 async fn delete(&self, path: &Path) -> std::io::Result<()> {
1077 self.0.delete(path).await
1078 }
1079
1080 async fn rename(&self, from: &Path, to: &Path) -> std::io::Result<()> {
1081 self.0.rename(from, to).await
1082 }
1083
1084 async fn sync(&self) -> std::io::Result<()> {
1085 Err(std::io::Error::other("injected directory fsync failure"))
1086 }
1087
1088 async fn streaming_writer(
1089 &self,
1090 path: &Path,
1091 ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
1092 self.0.streaming_writer(path).await
1093 }
1094 }
1095
1096 fn test_schema() -> Schema {
1097 Schema::default()
1098 }
1099
1100 fn dense_schema(index_type: VectorIndexType) -> (Schema, crate::dsl::Field) {
1101 let mut builder = crate::dsl::SchemaBuilder::default();
1102 let config = match index_type {
1103 VectorIndexType::IvfTq => crate::dsl::DenseVectorConfig::ivf_tq(2, Some(1), 1),
1104 other => panic!("unsupported trained test index type: {other:?}"),
1105 };
1106 let field = builder.add_dense_vector_field_with_config("embedding", true, true, config);
1107 (builder.build(), field)
1108 }
1109
1110 fn test_centroids() -> crate::structures::CoarseCentroids {
1111 crate::structures::CoarseCentroids {
1112 num_clusters: 1,
1113 dim: 2,
1114 centroids: vec![0.25, 0.75],
1115 version: crate::structures::mark_ivf_tq_cosine_generation(7),
1116 soar_config: None,
1117 routing_index: None,
1118 }
1119 }
1120
1121 async fn write_bincode(
1122 directory: &crate::directories::RamDirectory,
1123 filename: &str,
1124 value: &impl serde::Serialize,
1125 ) {
1126 let bytes = bincode::serde::encode_to_vec(value, bincode::config::standard()).unwrap();
1127 directory.write(Path::new(filename), &bytes).await.unwrap();
1128 }
1129
1130 #[test]
1131 fn test_metadata_init() {
1132 let mut meta = IndexMetadata::new(test_schema());
1133 assert_eq!(meta.total_vectors, 0);
1134 assert!(meta.segment_metas.is_empty());
1135 assert!(!meta.is_field_built(0));
1136
1137 meta.init_field(0, VectorIndexType::IvfTq);
1138 assert!(!meta.is_field_built(0));
1139 assert!(meta.vector_fields.contains_key(&0));
1140 }
1141
1142 #[tokio::test]
1143 async fn phrase_limits_load_from_legacy_and_configured_metadata_and_reject_corruption() {
1144 let directory = crate::directories::RamDirectory::new();
1145 let legacy = IndexMetadata::new(test_schema())
1146 .serialize_to_bytes()
1147 .unwrap();
1148 assert!(!String::from_utf8_lossy(&legacy).contains("max_l1_phrase_terms"));
1149 directory
1150 .write(Path::new(INDEX_META_FILENAME), &legacy)
1151 .await
1152 .unwrap();
1153 let loaded = IndexMetadata::load(&directory).await.unwrap();
1154 assert_eq!(loaded.schema.max_l1_phrase_terms(), 64);
1155 assert_eq!(loaded.serialize_to_bytes().unwrap(), legacy);
1156
1157 let mut configured: serde_json::Value = serde_json::from_slice(&legacy).unwrap();
1158 configured["schema"]["max_l1_phrase_terms"] = 300.into();
1159 directory
1160 .write(
1161 Path::new(INDEX_META_FILENAME),
1162 &serde_json::to_vec(&configured).unwrap(),
1163 )
1164 .await
1165 .unwrap();
1166 let loaded = IndexMetadata::load(&directory).await.unwrap();
1167 assert_eq!(loaded.schema.max_l1_phrase_terms(), 300);
1168 loaded.save(&directory).await.unwrap();
1169 assert_eq!(
1170 IndexMetadata::load(&directory)
1171 .await
1172 .unwrap()
1173 .schema
1174 .max_l1_phrase_terms(),
1175 300
1176 );
1177
1178 for invalid in [
1179 serde_json::json!(0),
1180 serde_json::json!(-1),
1181 serde_json::json!(1.5),
1182 serde_json::json!(4294967296u64),
1183 ] {
1184 configured["schema"]["max_l1_phrase_terms"] = invalid;
1185 directory
1186 .write(
1187 Path::new(INDEX_META_FILENAME),
1188 &serde_json::to_vec(&configured).unwrap(),
1189 )
1190 .await
1191 .unwrap();
1192 assert!(
1193 IndexMetadata::load(&directory).await.is_err(),
1194 "must not replace a corrupt cap with the default"
1195 );
1196 }
1197 }
1198
1199 #[tokio::test]
1200 async fn load_refuses_metadata_stamped_with_a_newer_format_version() {
1201 let directory = crate::directories::RamDirectory::new();
1202 let mut metadata = IndexMetadata::new(test_schema());
1203 metadata.version = INDEX_META_FORMAT_VERSION + 1;
1204 metadata.save(&directory).await.unwrap();
1205
1206 let error = IndexMetadata::load(&directory)
1207 .await
1208 .expect_err("metadata from a newer format version must be refused, not silently pruned")
1209 .to_string();
1210 assert!(
1211 error.contains(&format!("version {}", INDEX_META_FORMAT_VERSION + 1)),
1212 "{error}"
1213 );
1214 assert!(error.contains("incompatible"), "{error}");
1215 }
1216
1217 fn stamped_previous_format_bytes(metadata: &IndexMetadata) -> Vec<u8> {
1218 let mut raw: serde_json::Value =
1219 serde_json::from_slice(&metadata.serialize_to_bytes().unwrap()).unwrap();
1220 raw["version"] = serde_json::Value::from(INDEX_META_FORMAT_VERSION - 1);
1221 serde_json::to_vec(&raw).unwrap()
1222 }
1223
1224 #[tokio::test]
1225 async fn read_only_migration_preserves_format_6_and_7_metadata_bytes() {
1226 use crate::directories::Directory;
1227 for version in [6, 7] {
1228 let directory = crate::directories::RamDirectory::new();
1229 let mut metadata = IndexMetadata::new(test_schema());
1230 metadata.version = version;
1231 metadata.add_segment("kept".into(), 7);
1232 metadata.save(&directory).await.unwrap();
1233 let before = directory
1234 .open_read(Path::new(INDEX_META_FILENAME))
1235 .await
1236 .unwrap()
1237 .read_bytes()
1238 .await
1239 .unwrap();
1240 let (loaded, migrated) = IndexMetadata::load_reporting_migration(&directory)
1241 .await
1242 .unwrap();
1243 assert_eq!(migrated, Some(version));
1244 assert_eq!(loaded.version, INDEX_META_FORMAT_VERSION);
1245 assert_eq!(loaded.segment_metas["kept"].num_docs, 7);
1246 let after = directory
1247 .open_read(Path::new(INDEX_META_FILENAME))
1248 .await
1249 .unwrap()
1250 .read_bytes()
1251 .await
1252 .unwrap();
1253 assert_eq!(after.as_slice(), before.as_slice());
1254 loaded.save(&directory).await.unwrap();
1255 assert_eq!(
1256 IndexMetadata::load_reporting_migration(&directory)
1257 .await
1258 .unwrap()
1259 .1,
1260 None
1261 );
1262 }
1263 }
1264
1265 #[tokio::test]
1266 async fn load_migrates_previous_metadata_to_the_current_format() {
1267 let directory = crate::directories::RamDirectory::new();
1268 let mut metadata = IndexMetadata::new(test_schema());
1269 metadata.add_segment("kept".to_string(), 7);
1270 directory
1271 .write(
1272 Path::new(INDEX_META_FILENAME),
1273 &stamped_previous_format_bytes(&metadata),
1274 )
1275 .await
1276 .unwrap();
1277
1278 let (loaded, migrated_from) = IndexMetadata::load_reporting_migration(&directory)
1279 .await
1280 .expect("the previous format remains readable");
1281 assert_eq!(migrated_from, Some(INDEX_META_FORMAT_VERSION - 1));
1282 assert_eq!(loaded.version, INDEX_META_FORMAT_VERSION);
1283 assert_eq!(loaded.segment_metas["kept"].num_docs, 7);
1284 assert!(loaded.segment_metas["kept"].deletions.is_none());
1285
1286 let (_, current) = IndexMetadata::load_reporting_migration(&directory)
1287 .await
1288 .unwrap();
1289 assert_eq!(
1290 current,
1291 Some(INDEX_META_FORMAT_VERSION - 1),
1292 "load alone never persists"
1293 );
1294 metadata.save(&directory).await.unwrap();
1295 let (_, current) = IndexMetadata::load_reporting_migration(&directory)
1296 .await
1297 .unwrap();
1298 assert_eq!(
1299 current, None,
1300 "current-format metadata reports no migration"
1301 );
1302 }
1303
1304 #[tokio::test]
1305 async fn tmp_recovery_migrates_previous_metadata() {
1306 let directory = crate::directories::RamDirectory::new();
1307 let mut metadata = IndexMetadata::new(test_schema());
1308 metadata.add_segment("kept".to_string(), 3);
1309 directory
1310 .write(
1311 Path::new(INDEX_META_TMP_FILENAME),
1312 &stamped_previous_format_bytes(&metadata),
1313 )
1314 .await
1315 .unwrap();
1316
1317 let loaded = IndexMetadata::load(&directory)
1318 .await
1319 .expect("temp-file recovery applies the same migration");
1320 assert_eq!(loaded.version, INDEX_META_FORMAT_VERSION);
1321 assert_eq!(loaded.segment_metas["kept"].num_docs, 3);
1322 }
1323
1324 #[tokio::test]
1325 async fn load_refuses_metadata_older_than_format_6() {
1326 let directory = crate::directories::RamDirectory::new();
1329 let mut metadata = IndexMetadata::new(test_schema());
1330 metadata.version = OLDEST_MIGRATABLE_FORMAT_VERSION - 1;
1331 metadata.save(&directory).await.unwrap();
1332
1333 let error = IndexMetadata::load(&directory)
1334 .await
1335 .expect_err("metadata predating compatible position streams must be refused")
1336 .to_string();
1337 assert!(
1338 error.contains(&format!("version {}", OLDEST_MIGRATABLE_FORMAT_VERSION - 1)),
1339 "{error}"
1340 );
1341 assert!(error.contains("incompatible"), "{error}");
1342 }
1343
1344 #[tokio::test]
1345 async fn tmp_recovery_refuses_metadata_stamped_with_a_newer_format_version() {
1346 let directory = crate::directories::RamDirectory::new();
1347 let mut metadata = IndexMetadata::new(test_schema());
1348 metadata.version = INDEX_META_FORMAT_VERSION + 1;
1349 let bytes = metadata.serialize_to_bytes().unwrap();
1350 directory
1352 .write(Path::new(INDEX_META_TMP_FILENAME), &bytes)
1353 .await
1354 .unwrap();
1355
1356 let error = IndexMetadata::load(&directory)
1357 .await
1358 .expect_err("temp-file recovery must apply the same version gate")
1359 .to_string();
1360 assert!(
1361 error.contains(&format!("version {}", INDEX_META_FORMAT_VERSION + 1)),
1362 "{error}"
1363 );
1364 }
1365
1366 #[tokio::test]
1367 async fn save_treats_post_rename_sync_failure_as_committed() {
1368 let directory = SyncFailDirectory::default();
1369 let mut metadata = IndexMetadata::new(test_schema());
1370 metadata.add_segment("committed".to_string(), 7);
1371
1372 metadata.save(&directory).await.unwrap();
1373
1374 let loaded = IndexMetadata::load(&directory).await.unwrap();
1375 assert_eq!(loaded.segment_doc_count("committed"), Some(7));
1376 }
1377
1378 #[tokio::test]
1379 async fn trained_artifacts_load_only_when_the_complete_built_set_is_valid() {
1380 let mut builder = crate::dsl::SchemaBuilder::default();
1381 let config = crate::dsl::DenseVectorConfig::ivf_tq(2, Some(1), 1);
1382 let first = builder.add_dense_vector_field_with_config(
1383 "first_embedding",
1384 true,
1385 true,
1386 config.clone(),
1387 );
1388 let second =
1389 builder.add_dense_vector_field_with_config("second_embedding", true, true, config);
1390 let schema = builder.build();
1391 let directory = crate::directories::RamDirectory::new();
1392 let mut metadata = IndexMetadata::new(schema.clone());
1393 metadata.init_field(first.0, VectorIndexType::IvfTq);
1394 metadata.init_field(second.0, VectorIndexType::IvfTq);
1395 metadata.mark_field_built(first.0, 10, 1, "field_0_centroids.bin".into(), None);
1396 metadata.mark_field_built(second.0, 10, 1, "field_1_centroids.bin".into(), None);
1397 write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
1398
1399 let error = IndexMetadata::try_load_trained_from_fields(
1400 &metadata.vector_fields,
1401 &schema,
1402 &directory,
1403 )
1404 .await
1405 .err()
1406 .expect("missing artifact must fail the complete load")
1407 .to_string();
1408 assert!(error.contains("field_1_centroids.bin"), "{error}");
1409 assert!(error.contains("field 1"), "{error}");
1410 }
1411
1412 #[tokio::test]
1413 async fn index_open_fails_closed_when_built_artifact_is_missing() {
1414 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1415 let directory = crate::directories::RamDirectory::new();
1416 let mut metadata = IndexMetadata::new(schema);
1417 metadata.init_field(field.0, VectorIndexType::IvfTq);
1418 metadata.mark_field_built(field.0, 10, 1, "missing_centroids.bin".into(), None);
1419 metadata.save(&directory).await.unwrap();
1420
1421 let error = match crate::index::Index::open(directory, crate::index::IndexConfig::default())
1422 .await
1423 {
1424 Ok(_) => panic!("Index::open accepted a Built field with no artifact"),
1425 Err(error) => error.to_string(),
1426 };
1427 assert!(error.contains("missing_centroids.bin"), "{error}");
1428 }
1429
1430 #[tokio::test]
1431 async fn ivf_tq_built_state_rejects_a_codebook_file() {
1432 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1433 let directory = crate::directories::RamDirectory::new();
1434 let mut metadata = IndexMetadata::new(schema.clone());
1435 metadata.init_field(field.0, VectorIndexType::IvfTq);
1436 metadata.mark_field_built(
1437 field.0,
1438 10,
1439 1,
1440 "field_0_centroids.bin".into(),
1441 Some("field_0_codebook.bin".into()),
1442 );
1443 write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
1444
1445 let error = IndexMetadata::try_load_trained_from_fields(
1446 &metadata.vector_fields,
1447 &schema,
1448 &directory,
1449 )
1450 .await
1451 .err()
1452 .expect("IVF-TQ Built state with a codebook file must fail")
1453 .to_string();
1454 assert!(error.contains("codebook"), "{error}");
1455 }
1456
1457 #[tokio::test]
1458 async fn legacy_ivf_tq_centroid_generation_is_rejected_while_loading() {
1459 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1460 let directory = crate::directories::RamDirectory::new();
1461 let mut metadata = IndexMetadata::new(schema.clone());
1462 metadata.init_field(field.0, VectorIndexType::IvfTq);
1463 metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1464 let mut legacy = test_centroids();
1465 legacy.version = 7;
1466 write_bincode(&directory, "field_0_centroids.bin", &legacy).await;
1467
1468 let error = IndexMetadata::try_load_trained_from_fields(
1469 &metadata.vector_fields,
1470 &schema,
1471 &directory,
1472 )
1473 .await
1474 .err()
1475 .expect("legacy IVF-TQ centroid state must fail while loading")
1476 .to_string();
1477 assert!(error.contains("unsupported legacy generation"), "{error}");
1478 assert!(error.contains("rebuild the index"), "{error}");
1479 }
1480
1481 #[tokio::test]
1482 async fn legacy_ivf_pq_trained_field_fails_with_actionable_error() {
1483 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1487 let directory = crate::directories::RamDirectory::new();
1488 let mut metadata = IndexMetadata::new(schema.clone());
1489 metadata.init_field(field.0, VectorIndexType::IvfTq);
1490 metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1491 metadata
1494 .vector_fields
1495 .get_mut(&field.0)
1496 .expect("field initialized")
1497 .index_type = VectorFieldIndexType::Float(VectorIndexType::IvfPq);
1498 write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
1499
1500 let error = IndexMetadata::try_load_trained_from_fields(
1501 &metadata.vector_fields,
1502 &schema,
1503 &directory,
1504 )
1505 .await
1506 .err()
1507 .expect("legacy IVF-PQ trained state must fail loudly")
1508 .to_string();
1509 assert!(error.contains("no longer"), "{error}");
1510 assert!(error.contains("ivf_tq"), "{error}");
1511 }
1512
1513 #[tokio::test]
1514 async fn requested_cluster_count_accepts_a_quality_clamped_artifact() {
1515 let mut builder = crate::dsl::SchemaBuilder::default();
1516 let field = builder.add_dense_vector_field_with_config(
1517 "embedding",
1518 true,
1519 true,
1520 crate::dsl::DenseVectorConfig::ivf_tq(2, Some(4), 1),
1521 );
1522 let schema = builder.build();
1523 let directory = crate::directories::RamDirectory::new();
1524 let mut metadata = IndexMetadata::new(schema.clone());
1525 metadata.init_field(field.0, VectorIndexType::IvfTq);
1526 metadata.mark_field_built(field.0, 1, 4, "field_0_centroids.bin".into(), None);
1527 write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
1528
1529 let trained = IndexMetadata::try_load_trained_from_fields(
1530 &metadata.vector_fields,
1531 &schema,
1532 &directory,
1533 )
1534 .await
1535 .unwrap()
1536 .unwrap();
1537 assert_eq!(trained.centroids[&field.0].num_clusters, 1);
1538 }
1539
1540 #[tokio::test]
1541 async fn trained_artifact_loader_rejects_trailing_data() {
1542 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1543 let directory = crate::directories::RamDirectory::new();
1544 let mut metadata = IndexMetadata::new(schema.clone());
1545 metadata.init_field(field.0, VectorIndexType::IvfTq);
1546 metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1547 let mut bytes =
1548 bincode::serde::encode_to_vec(test_centroids(), bincode::config::standard()).unwrap();
1549 bytes.extend_from_slice(&[0xaa, 0xbb]);
1550 directory
1551 .write(Path::new("field_0_centroids.bin"), &bytes)
1552 .await
1553 .unwrap();
1554
1555 let error = IndexMetadata::try_load_trained_from_fields(
1556 &metadata.vector_fields,
1557 &schema,
1558 &directory,
1559 )
1560 .await
1561 .err()
1562 .expect("trailing artifact bytes must fail validation")
1563 .to_string();
1564 assert!(error.contains("trailing bytes"), "{error}");
1565 }
1566
1567 #[test]
1568 fn trained_artifact_size_limit_rejects_before_reading() {
1569 let error = validate_trained_artifact_size(
1570 3,
1571 "centroids",
1572 "field_3_centroids.bin",
1573 MAX_TRAINED_ARTIFACT_BYTES as u64 + 1,
1574 )
1575 .unwrap_err()
1576 .to_string();
1577 assert!(error.contains("exceeding"), "{error}");
1578 assert!(error.contains("field 3"), "{error}");
1579 }
1580
1581 #[tokio::test]
1582 async fn trained_artifact_decode_limit_rejects_forged_collection_length() {
1583 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1584 let directory = crate::directories::RamDirectory::new();
1585 let mut metadata = IndexMetadata::new(schema.clone());
1586 metadata.init_field(field.0, VectorIndexType::IvfTq);
1587 metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1588
1589 let mut bytes = vec![1, 2, 253];
1594 bytes.extend_from_slice(&u64::MAX.to_le_bytes());
1595 directory
1596 .write(Path::new("field_0_centroids.bin"), &bytes)
1597 .await
1598 .unwrap();
1599
1600 let error = IndexMetadata::try_load_trained_from_fields(
1601 &metadata.vector_fields,
1602 &schema,
1603 &directory,
1604 )
1605 .await
1606 .err()
1607 .expect("forged collection length must fail the bounded decoder")
1608 .to_string();
1609 assert!(error.contains("failed to deserialize"), "{error}");
1610 }
1611
1612 #[test]
1613 fn test_metadata_segments() {
1614 let mut meta = IndexMetadata::new(test_schema());
1615 meta.add_segment("abc123".to_string(), 50);
1616 meta.add_segment("def456".to_string(), 100);
1617 assert_eq!(meta.segment_metas.len(), 2);
1618 assert_eq!(meta.segment_doc_count("abc123"), Some(50));
1619 assert_eq!(meta.segment_doc_count("def456"), Some(100));
1620
1621 meta.add_segment("abc123".to_string(), 75);
1623 assert_eq!(meta.segment_metas.len(), 2);
1624 assert_eq!(meta.segment_doc_count("abc123"), Some(75));
1625
1626 meta.remove_segment("abc123");
1627 assert_eq!(meta.segment_metas.len(), 1);
1628 assert!(meta.has_segment("def456"));
1629 assert!(!meta.has_segment("abc123"));
1630 }
1631
1632 #[test]
1633 fn test_mark_field_built() {
1634 let mut meta = IndexMetadata::new(test_schema());
1635 meta.init_field(0, VectorIndexType::IvfTq);
1636 meta.total_vectors = 10000;
1637
1638 assert!(!meta.is_field_built(0));
1639
1640 meta.mark_field_built(0, 10000, 256, "field_0_centroids.bin".to_string(), None);
1641
1642 assert!(meta.is_field_built(0));
1643 let field = meta.get_field_meta(0).unwrap();
1644 assert_eq!(
1645 field.centroids_file.as_deref(),
1646 Some("field_0_centroids.bin")
1647 );
1648 }
1649
1650 #[test]
1651 fn scann_metadata_persists_generation_and_fingerprint_and_defaults_old_json() {
1652 let mut meta = IndexMetadata::new(test_schema());
1653 meta.init_field(3, VectorIndexType::Scann);
1654 meta.mark_scann_field_built(
1655 3,
1656 100_000,
1657 1_000,
1658 "field_3_scann_17.bin".to_string(),
1659 17,
1660 0xdecafbad,
1661 )
1662 .unwrap();
1663
1664 let bytes = meta.serialize_to_bytes().unwrap();
1665 let decoded: IndexMetadata = serde_json::from_slice(&bytes).unwrap();
1666 let field = decoded.get_field_meta(3).unwrap();
1667 assert_eq!(field.artifact_generation, Some(17));
1668 assert_eq!(field.artifact_id, Some(0xdecafbad));
1669
1670 let mut legacy_json = serde_json::to_value(&decoded).unwrap();
1671 legacy_json["vector_fields"]["3"]
1672 .as_object_mut()
1673 .unwrap()
1674 .remove("artifact_generation");
1675 legacy_json["vector_fields"]["3"]
1676 .as_object_mut()
1677 .unwrap()
1678 .remove("artifact_id");
1679 let legacy: IndexMetadata = serde_json::from_value(legacy_json).unwrap();
1680 let legacy_field = legacy.get_field_meta(3).unwrap();
1681 assert_eq!(legacy_field.artifact_generation, None);
1682 assert_eq!(legacy_field.artifact_id, None);
1683 }
1684
1685 #[test]
1686 fn scann_metadata_refuses_zero_or_non_scann_generation() {
1687 let mut meta = IndexMetadata::new(test_schema());
1688 meta.init_field(0, VectorIndexType::Scann);
1689 assert!(
1690 meta.mark_scann_field_built(0, 100_000, 1_000, "artifact.bin".into(), 0, 1)
1691 .is_err()
1692 );
1693 meta.init_field(1, VectorIndexType::IvfTq);
1694 assert!(
1695 meta.mark_scann_field_built(1, 100_000, 1_000, "artifact.bin".into(), 1, 2)
1696 .is_err()
1697 );
1698 }
1699
1700 #[test]
1701 fn scann_autopilot_accepts_resolved_billion_scale_three_level_geometry() {
1702 assert!(scann_explicit_geometry_matches(None, None, 10_000_000, 3));
1703 assert!(!scann_explicit_geometry_matches(
1704 Some(1_000_000),
1705 None,
1706 10_000_000,
1707 3
1708 ));
1709 assert!(!scann_explicit_geometry_matches(
1710 None,
1711 Some(1),
1712 10_000_000,
1713 3
1714 ));
1715 }
1716
1717 #[test]
1718 fn total_vectors_is_aggregate_of_built_field_counts() {
1719 let mut meta = IndexMetadata::new(test_schema());
1720 meta.init_field(7, VectorIndexType::IvfTq);
1721 meta.init_field(3, VectorIndexType::IvfTq);
1722
1723 meta.mark_field_built(7, 400, 20, "field_7_centroids.bin".to_string(), None);
1726 assert_eq!(meta.total_vectors, 400);
1727 meta.mark_field_built(3, 250, 15, "field_3_centroids.bin".to_string(), None);
1728 assert_eq!(meta.total_vectors, 650);
1729
1730 meta.mark_field_built(7, 425, 20, "field_7_centroids.bin".to_string(), None);
1733 assert_eq!(meta.total_vectors, 675);
1734 }
1735
1736 #[test]
1737 fn test_should_build_field() {
1738 let mut meta = IndexMetadata::new(test_schema());
1739 meta.init_field(0, VectorIndexType::IvfTq);
1740
1741 meta.total_vectors = 500;
1743 assert!(!meta.should_build_field(0, 1000));
1744
1745 meta.total_vectors = 1500;
1747 assert!(meta.should_build_field(0, 1000));
1748
1749 meta.mark_field_built(0, 1500, 256, "centroids.bin".to_string(), None);
1751 assert!(!meta.should_build_field(0, 1000));
1752 }
1753
1754 #[test]
1755 fn test_serialization() {
1756 let mut meta = IndexMetadata::new(test_schema());
1757 meta.add_segment("seg1".to_string(), 100);
1758 meta.init_field(0, VectorIndexType::IvfTq);
1759 meta.total_vectors = 5000;
1760
1761 let json = serde_json::to_string_pretty(&meta).unwrap();
1762 let loaded: IndexMetadata = serde_json::from_str(&json).unwrap();
1763
1764 assert_eq!(loaded.segment_ids().len(), meta.segment_ids().len());
1765 assert_eq!(loaded.segment_doc_count("seg1"), Some(100));
1766 assert_eq!(loaded.total_vectors, meta.total_vectors);
1767 assert!(loaded.vector_fields.contains_key(&0));
1768 }
1769
1770 #[test]
1771 fn old_metadata_defaults_the_bp_retry_counter() {
1772 let mut meta = IndexMetadata::new(test_schema());
1773 meta.add_segment("legacy".to_string(), 10);
1774 let mut json = serde_json::to_value(&meta).unwrap();
1775 json["segment_metas"]["legacy"]
1776 .as_object_mut()
1777 .unwrap()
1778 .remove("bp_unconverged_passes");
1779
1780 let loaded: IndexMetadata = serde_json::from_value(json).unwrap();
1781 assert_eq!(loaded.segment_metas["legacy"].bp_unconverged_passes, 0);
1782 }
1783
1784 #[test]
1785 fn test_merged_segment_lineage() {
1786 let mut meta = IndexMetadata::new(test_schema());
1787 meta.add_segment("a".to_string(), 50);
1788 meta.add_segment("b".to_string(), 75);
1789
1790 assert_eq!(meta.segment_metas["a"].generation, 0);
1792 assert!(meta.segment_metas["a"].ancestors.is_empty());
1793
1794 meta.add_merged_segment(
1796 "c".to_string(),
1797 125,
1798 vec!["a".to_string(), "b".to_string()],
1799 1,
1800 false,
1801 true,
1802 );
1803 assert_eq!(meta.segment_metas["c"].generation, 1);
1804 assert_eq!(meta.segment_metas["c"].ancestors, vec!["a", "b"]);
1805 assert_eq!(meta.segment_doc_count("c"), Some(125));
1806
1807 meta.add_segment("d".to_string(), 30);
1809 meta.add_merged_segment(
1810 "e".to_string(),
1811 155,
1812 vec!["c".to_string(), "d".to_string()],
1813 2,
1814 false,
1815 true,
1816 );
1817 assert_eq!(meta.segment_metas["e"].generation, 2);
1818 }
1819}