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 = 4;
34
35pub(crate) const MAX_TRAINED_ARTIFACT_BYTES: usize = 512 * 1024 * 1024;
40
41#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
43pub enum VectorIndexState {
44 #[default]
46 Flat,
47 Built {
49 vector_count: usize,
51 num_clusters: usize,
53 },
54}
55
56fn default_true() -> bool {
57 true
58}
59
60#[derive(Debug, Clone, Serialize, Deserialize)]
63pub struct SegmentMetaInfo {
64 pub num_docs: u32,
66 pub ancestors: Vec<String>,
68 pub generation: u32,
70 #[serde(default)]
74 pub reordered: bool,
75 #[serde(default = "default_true")]
80 pub bp_converged: bool,
81 #[serde(default)]
85 pub bp_unconverged_passes: u32,
86}
87
88#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
90#[serde(tag = "kind", content = "index", rename_all = "snake_case")]
91pub enum VectorFieldIndexType {
92 Float(VectorIndexType),
93 Binary(BinaryIndexType),
94}
95
96impl From<VectorIndexType> for VectorFieldIndexType {
97 fn from(value: VectorIndexType) -> Self {
98 Self::Float(value)
99 }
100}
101
102impl From<BinaryIndexType> for VectorFieldIndexType {
103 fn from(value: BinaryIndexType) -> Self {
104 Self::Binary(value)
105 }
106}
107
108#[derive(Debug, Clone, Serialize, Deserialize)]
109pub struct FieldVectorMeta {
110 pub field_id: u32,
112 pub index_type: VectorFieldIndexType,
114 pub state: VectorIndexState,
116 #[serde(skip_serializing_if = "Option::is_none")]
118 pub centroids_file: Option<String>,
119 #[serde(skip_serializing_if = "Option::is_none")]
123 pub codebook_file: Option<String>,
124}
125
126#[derive(Debug, Clone, Serialize, Deserialize)]
128pub struct IndexMetadata {
129 pub version: u32,
131 pub schema: Schema,
133 #[serde(default)]
136 pub segment_metas: HashMap<String, SegmentMetaInfo>,
137 #[serde(default)]
139 pub vector_fields: HashMap<u32, FieldVectorMeta>,
140 #[serde(default)]
147 pub total_vectors: usize,
148}
149
150impl IndexMetadata {
151 pub fn new(schema: Schema) -> Self {
153 Self {
154 version: INDEX_META_FORMAT_VERSION,
155 schema,
156 segment_metas: HashMap::new(),
157 vector_fields: HashMap::new(),
158 total_vectors: 0,
159 }
160 }
161
162 pub fn segment_ids(&self) -> Vec<String> {
164 let mut ids: Vec<String> = self.segment_metas.keys().cloned().collect();
165 ids.sort();
166 ids
167 }
168
169 pub fn add_segment(&mut self, segment_id: String, num_docs: u32) {
171 self.segment_metas.insert(
172 segment_id,
173 SegmentMetaInfo {
174 num_docs,
175 ancestors: Vec::new(),
176 generation: 0,
177 reordered: false,
178 bp_converged: true,
179 bp_unconverged_passes: 0,
180 },
181 );
182 }
183
184 pub fn add_merged_segment(
186 &mut self,
187 segment_id: String,
188 num_docs: u32,
189 ancestors: Vec<String>,
190 generation: u32,
191 reordered: bool,
192 bp_converged: bool,
193 ) {
194 self.add_segment_meta(
195 segment_id,
196 SegmentMetaInfo {
197 num_docs,
198 ancestors,
199 generation,
200 reordered,
201 bp_converged,
202 bp_unconverged_passes: 0,
203 },
204 );
205 }
206
207 pub(crate) fn add_segment_meta(&mut self, segment_id: String, info: SegmentMetaInfo) {
211 self.segment_metas.insert(segment_id, info);
212 }
213
214 pub fn remove_segment(&mut self, segment_id: &str) {
216 self.segment_metas.remove(segment_id);
217 }
218
219 pub fn has_segment(&self, segment_id: &str) -> bool {
221 self.segment_metas.contains_key(segment_id)
222 }
223
224 pub fn segment_doc_count(&self, segment_id: &str) -> Option<u32> {
226 self.segment_metas.get(segment_id).map(|m| m.num_docs)
227 }
228
229 pub fn is_field_built(&self, field_id: u32) -> bool {
231 self.vector_fields
232 .get(&field_id)
233 .map(|f| matches!(f.state, VectorIndexState::Built { .. }))
234 .unwrap_or(false)
235 }
236
237 pub fn get_field_meta(&self, field_id: u32) -> Option<&FieldVectorMeta> {
239 self.vector_fields.get(&field_id)
240 }
241
242 pub fn init_field(&mut self, field_id: u32, index_type: impl Into<VectorFieldIndexType>) {
244 let index_type = index_type.into();
245 self.vector_fields
246 .entry(field_id)
247 .or_insert(FieldVectorMeta {
248 field_id,
249 index_type,
250 state: VectorIndexState::Flat,
251 centroids_file: None,
252 codebook_file: None,
253 });
254 }
255
256 pub fn mark_field_built(
258 &mut self,
259 field_id: u32,
260 vector_count: usize,
261 num_clusters: usize,
262 centroids_file: String,
263 codebook_file: Option<String>,
264 ) {
265 if let Some(field) = self.vector_fields.get_mut(&field_id) {
266 field.state = VectorIndexState::Built {
267 vector_count,
268 num_clusters,
269 };
270 field.centroids_file = Some(centroids_file);
271 field.codebook_file = codebook_file;
272 self.refresh_total_vectors();
273 }
274 }
275
276 pub(crate) fn refresh_total_vectors(&mut self) {
282 self.total_vectors = self
283 .vector_fields
284 .values()
285 .filter_map(|field| match field.state {
286 VectorIndexState::Built { vector_count, .. } => Some(vector_count),
287 VectorIndexState::Flat => None,
288 })
289 .fold(0usize, usize::saturating_add);
290 }
291
292 pub fn should_build_field(&self, field_id: u32, threshold: usize) -> bool {
294 if self.is_field_built(field_id) {
296 return false;
297 }
298 self.total_vectors >= threshold
300 }
301
302 pub async fn load<D: crate::directories::Directory>(dir: &D) -> Result<Self> {
307 let path = Path::new(INDEX_META_FILENAME);
308 match dir.open_read(path).await {
309 Ok(slice) => {
310 let bytes = slice.read_bytes().await?;
311 Self::deserialize_versioned(bytes.as_slice())
312 }
313 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
314 let tmp_path = Path::new(INDEX_META_TMP_FILENAME);
316 let slice = dir.open_read(tmp_path).await?;
317 let bytes = slice.read_bytes().await?;
318 let meta = Self::deserialize_versioned(bytes.as_slice())?;
319 log::warn!("Recovered metadata from temp file (previous crash during save)");
320 Ok(meta)
321 }
322 Err(e) => Err(Error::Io(e)),
323 }
324 }
325
326 fn deserialize_versioned(bytes: &[u8]) -> Result<Self> {
330 let meta: Self =
331 serde_json::from_slice(bytes).map_err(|e| Error::Serialization(e.to_string()))?;
332 crate::dsl::reject_removed_vector_index_types(&meta.schema).map_err(Error::Schema)?;
333 if meta.version != INDEX_META_FORMAT_VERSION {
334 return Err(Error::Corruption(format!(
335 "metadata.json format version {} is incompatible with required version {}; \
336 rebuild and republish the index with this Hermes version",
337 meta.version, INDEX_META_FORMAT_VERSION
338 )));
339 }
340 Ok(meta)
341 }
342
343 pub async fn save<D: crate::directories::DirectoryWriter>(&self, dir: &D) -> Result<()> {
348 let bytes = self.serialize_to_bytes()?;
349 Self::save_bytes(dir, &bytes).await
350 }
351
352 pub fn serialize_to_bytes(&self) -> Result<Vec<u8>> {
355 serde_json::to_vec_pretty(self).map_err(|e| Error::Serialization(e.to_string()))
356 }
357
358 pub async fn save_bytes<D: crate::directories::DirectoryWriter>(
363 dir: &D,
364 bytes: &[u8],
365 ) -> Result<()> {
366 let tmp_path = Path::new(INDEX_META_TMP_FILENAME);
367 let final_path = Path::new(INDEX_META_FILENAME);
368 let mut writer = dir.streaming_writer(tmp_path).await.map_err(Error::Io)?;
373 writer.write_all(bytes).map_err(Error::Io)?;
374 writer.finish().map_err(Error::Io)?;
375 dir.rename(tmp_path, final_path).await.map_err(Error::Io)?;
382 if let Err(error) = dir.sync().await {
383 log::error!(
384 "[metadata] directory fsync failed after committed rename: {}. \
385 Continuing with the renamed generation; crash durability is not guaranteed",
386 error,
387 );
388 }
389 Ok(())
390 }
391
392 #[cfg_attr(not(feature = "native"), allow(dead_code))]
394 pub(crate) async fn try_load_trained_from_fields<D: crate::directories::Directory>(
395 vector_fields: &HashMap<u32, FieldVectorMeta>,
396 schema: &Schema,
397 dir: &D,
398 ) -> Result<Option<crate::segment::TrainedVectorStructures>> {
399 Self::load_trained_from_fields_impl(vector_fields, schema, dir).await
400 }
401
402 async fn load_trained_from_fields_impl<D: crate::directories::Directory>(
412 vector_fields: &HashMap<u32, FieldVectorMeta>,
413 schema: &Schema,
414 dir: &D,
415 ) -> Result<Option<crate::segment::TrainedVectorStructures>> {
416 use std::sync::Arc;
417
418 let mut centroids = rustc_hash::FxHashMap::default();
419 let mut binary_quantizers = rustc_hash::FxHashMap::default();
420
421 let mut built_fields: Vec<_> = vector_fields
422 .iter()
423 .filter(|(_, meta)| matches!(meta.state, VectorIndexState::Built { .. }))
424 .collect();
425 built_fields.sort_unstable_by_key(|(field_id, _)| **field_id);
426
427 log::debug!(
428 "[trained] index={} loading trained structures, dense_vector_fields={:?}",
429 schema.index_label(),
430 vector_fields.keys().collect::<Vec<_>>()
431 );
432
433 for (field_id, field_meta) in built_fields {
434 log::debug!(
435 "[trained] index={} field {} state={:?} centroids_file={:?} codebook_file={:?}",
436 schema.index_label(),
437 field_id,
438 field_meta.state,
439 field_meta.centroids_file,
440 field_meta.codebook_file,
441 );
442 if field_meta.field_id != *field_id {
443 return Err(Error::Corruption(format!(
444 "trained vector metadata key {field_id} contains field_id {}",
445 field_meta.field_id
446 )));
447 }
448
449 let expected_clusters = match field_meta.state {
450 VectorIndexState::Built { num_clusters, .. } if num_clusters > 0 => num_clusters,
451 VectorIndexState::Built { .. } => {
452 return Err(Error::Corruption(format!(
453 "trained vector metadata field {field_id} has zero clusters"
454 )));
455 }
456 VectorIndexState::Flat => unreachable!("built_fields contains only Built entries"),
457 };
458
459 let centroids_file = field_meta.centroids_file.as_deref().ok_or_else(|| {
460 Error::Corruption(format!(
461 "trained vector metadata field {field_id} is Built but has no centroids_file"
462 ))
463 })?;
464 match field_meta.index_type {
465 VectorFieldIndexType::Float(VectorIndexType::IvfPq) => {
466 return Err(Error::Corruption(format!(
467 "field {field_id} was trained as IVF-PQ, which is no longer \
468 supported; recreate the index with `ivf_tq` and reindex \
469 (docs/turboquant-quantization.md)"
470 )));
471 }
472 VectorFieldIndexType::Float(index_type @ VectorIndexType::IvfTq) => {
473 let entry = schema
474 .get_field_entry(crate::dsl::Field(*field_id))
475 .ok_or_else(|| {
476 Error::Corruption(format!(
477 "trained vector metadata references missing field {field_id}"
478 ))
479 })?;
480 let schema_config = entry
481 .dense_vector_config
482 .as_ref()
483 .filter(|_| entry.field_type == crate::dsl::FieldType::DenseVector)
484 .ok_or_else(|| {
485 Error::Corruption(format!(
486 "trained vector metadata field {field_id} is not a float dense field"
487 ))
488 })?;
489 if schema_config.index_type != index_type {
490 return Err(Error::Corruption(format!(
491 "trained vector metadata field {field_id} uses {index_type:?}, schema requires {:?}",
492 schema_config.index_type
493 )));
494 }
495 let c: crate::structures::CoarseCentroids =
496 load_trained_artifact(dir, *field_id, "centroids", centroids_file).await?;
497 let expected_dim = schema_config.dim;
498 let actual_clusters = c.num_clusters as usize;
499 let expected_values =
500 actual_clusters.checked_mul(expected_dim).ok_or_else(|| {
501 Error::Corruption(format!(
502 "trained centroid dimensions overflow for field {field_id}"
503 ))
504 })?;
505 if actual_clusters == 0
506 || actual_clusters > expected_clusters
507 || c.dim == 0
508 || c.dim != expected_dim
509 || c.centroids.len() != expected_values
510 || c.centroids.iter().any(|value| !value.is_finite())
511 {
512 return Err(Error::Corruption(format!(
513 "trained centroids for field {field_id} do not match metadata/schema"
514 )));
515 }
516 if !crate::structures::is_ivf_tq_cosine_generation(c.version) {
517 return Err(Error::Corruption(format!(
518 "trained IVF-TQ centroids for field {field_id} use an \
519 unsupported legacy generation; rebuild the index"
520 )));
521 }
522 c.validate_routing(schema_config.ivf_routing)
523 .map_err(|error| {
524 Error::Corruption(format!(
525 "invalid trained centroid routing for field {field_id}: {error}"
526 ))
527 })?;
528 let _ = index_type;
531 if field_meta.codebook_file.is_some() {
532 return Err(Error::Corruption(format!(
533 "trained IVF-TQ field {field_id} unexpectedly references a codebook file"
534 )));
535 }
536 centroids.insert(*field_id, Arc::new(c));
537 }
538 VectorFieldIndexType::Binary(BinaryIndexType::Ivf) => {
539 let entry = schema
540 .get_field_entry(crate::dsl::Field(*field_id))
541 .ok_or_else(|| {
542 Error::Corruption(format!(
543 "trained vector metadata references missing field {field_id}"
544 ))
545 })?;
546 let schema_config = entry
547 .binary_dense_vector_config
548 .as_ref()
549 .filter(|config| {
550 entry.field_type == crate::dsl::FieldType::BinaryDenseVector
551 && config.index_type == BinaryIndexType::Ivf
552 })
553 .ok_or_else(|| {
554 Error::Corruption(format!(
555 "trained vector metadata field {field_id} is not a binary IVF field"
556 ))
557 })?;
558 let quantizer: crate::structures::BinaryCoarseQuantizer =
559 load_trained_artifact(dir, *field_id, "binary centroids", centroids_file)
560 .await?;
561 quantizer.validate().map_err(|error| {
562 Error::Corruption(format!(
563 "invalid binary coarse quantizer for field {field_id}: {error}"
564 ))
565 })?;
566 let actual_clusters = quantizer.num_clusters as usize;
567 if actual_clusters > expected_clusters
568 || schema_config.dim != quantizer.dim_bits
569 {
570 return Err(Error::Corruption(format!(
571 "binary coarse quantizer for field {field_id} does not match metadata/schema"
572 )));
573 }
574 quantizer
575 .validate_routing(schema_config.ivf_routing)
576 .map_err(|error| {
577 Error::Corruption(format!(
578 "invalid binary centroid routing for field {field_id}: {error}"
579 ))
580 })?;
581 binary_quantizers.insert(*field_id, Arc::new(quantizer));
582 }
583 unsupported => {
584 return Err(Error::Corruption(format!(
585 "field {field_id} is Built for {unsupported:?}, which has no global IVF artifacts"
586 )));
587 }
588 }
589 }
590
591 if centroids.is_empty() && binary_quantizers.is_empty() {
592 Ok(None)
593 } else {
594 let trained = crate::segment::TrainedVectorStructures {
595 #[cfg(feature = "native")]
596 _ann_pins: Default::default(),
597 centroids,
598 binary_quantizers,
599 };
600 #[cfg(feature = "native")]
601 let trained = {
602 let mut trained = trained;
603 trained.pin_ann_structures(crate::segment::pin::pin_policy());
604 trained
605 };
606 Ok(Some(trained))
607 }
608 }
609}
610
611fn validate_trained_artifact_path(field_id: u32, kind: &str, filename: &str) -> Result<()> {
612 use std::path::Component;
613
614 let path = Path::new(filename);
615 if filename.is_empty()
616 || path.is_absolute()
617 || path.components().any(|component| {
618 matches!(
619 component,
620 Component::ParentDir | Component::RootDir | Component::Prefix(_)
621 )
622 })
623 {
624 return Err(Error::Corruption(format!(
625 "trained {kind} path for field {field_id} is not a safe relative path: '{filename}'"
626 )));
627 }
628 Ok(())
629}
630
631async fn load_trained_artifact<T, D>(
632 dir: &D,
633 field_id: u32,
634 kind: &str,
635 filename: &str,
636) -> Result<T>
637where
638 T: serde::de::DeserializeOwned,
639 D: crate::directories::Directory,
640{
641 validate_trained_artifact_path(field_id, kind, filename)?;
642 let path = Path::new(filename);
643 let file_size = dir.file_size(path).await.map_err(|error| {
644 Error::Corruption(format!(
645 "failed to stat trained {kind} '{filename}' for field {field_id}: {error}"
646 ))
647 })?;
648 validate_trained_artifact_size(field_id, kind, filename, file_size)?;
649 let slice = dir.open_read(path).await.map_err(|error| {
650 Error::Corruption(format!(
651 "failed to open trained {kind} '{filename}' for field {field_id}: {error}"
652 ))
653 })?;
654 validate_trained_artifact_size(field_id, kind, filename, slice.len())?;
655 let bytes = slice.read_bytes().await.map_err(|error| {
656 Error::Corruption(format!(
657 "failed to read trained {kind} '{filename}' for field {field_id}: {error}"
658 ))
659 })?;
660 let (artifact, consumed) = bincode::serde::decode_from_slice::<T, _>(
661 bytes.as_slice(),
662 bincode::config::standard().with_limit::<MAX_TRAINED_ARTIFACT_BYTES>(),
663 )
664 .map_err(|error| {
665 Error::Corruption(format!(
666 "failed to deserialize trained {kind} '{filename}' for field {field_id}: {error}"
667 ))
668 })?;
669 if consumed != bytes.len() {
670 return Err(Error::Corruption(format!(
671 "trained {kind} '{filename}' for field {field_id} has {} trailing bytes",
672 bytes.len() - consumed
673 )));
674 }
675 Ok(artifact)
676}
677
678fn validate_trained_artifact_size(
679 field_id: u32,
680 kind: &str,
681 filename: &str,
682 file_size: u64,
683) -> Result<()> {
684 if file_size > MAX_TRAINED_ARTIFACT_BYTES as u64 {
685 return Err(Error::Corruption(format!(
686 "trained {kind} '{filename}' for field {field_id} is {file_size} bytes, \
687 exceeding the {MAX_TRAINED_ARTIFACT_BYTES}-byte safety limit"
688 )));
689 }
690 Ok(())
691}
692
693#[cfg(test)]
694mod tests {
695 use super::*;
696 use crate::directories::DirectoryWriter;
697
698 #[derive(Clone, Default)]
699 struct SyncFailDirectory(crate::directories::RamDirectory);
700
701 #[async_trait::async_trait]
702 impl crate::directories::Directory for SyncFailDirectory {
703 async fn exists(&self, path: &Path) -> std::io::Result<bool> {
704 self.0.exists(path).await
705 }
706
707 async fn file_size(&self, path: &Path) -> std::io::Result<u64> {
708 self.0.file_size(path).await
709 }
710
711 async fn open_read(&self, path: &Path) -> std::io::Result<crate::directories::FileHandle> {
712 self.0.open_read(path).await
713 }
714
715 async fn read_range(
716 &self,
717 path: &Path,
718 range: std::ops::Range<u64>,
719 ) -> std::io::Result<crate::directories::OwnedBytes> {
720 self.0.read_range(path, range).await
721 }
722
723 async fn list_files(&self, prefix: &Path) -> std::io::Result<Vec<std::path::PathBuf>> {
724 self.0.list_files(prefix).await
725 }
726
727 async fn open_lazy(&self, path: &Path) -> std::io::Result<crate::directories::FileHandle> {
728 self.0.open_lazy(path).await
729 }
730 }
731
732 #[async_trait::async_trait]
733 impl crate::directories::DirectoryWriter for SyncFailDirectory {
734 async fn write(&self, path: &Path, data: &[u8]) -> std::io::Result<()> {
735 self.0.write(path, data).await
736 }
737
738 async fn delete(&self, path: &Path) -> std::io::Result<()> {
739 self.0.delete(path).await
740 }
741
742 async fn rename(&self, from: &Path, to: &Path) -> std::io::Result<()> {
743 self.0.rename(from, to).await
744 }
745
746 async fn sync(&self) -> std::io::Result<()> {
747 Err(std::io::Error::other("injected directory fsync failure"))
748 }
749
750 async fn streaming_writer(
751 &self,
752 path: &Path,
753 ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
754 self.0.streaming_writer(path).await
755 }
756 }
757
758 fn test_schema() -> Schema {
759 Schema::default()
760 }
761
762 fn dense_schema(index_type: VectorIndexType) -> (Schema, crate::dsl::Field) {
763 let mut builder = crate::dsl::SchemaBuilder::default();
764 let config = match index_type {
765 VectorIndexType::IvfTq => crate::dsl::DenseVectorConfig::ivf_tq(2, Some(1), 1),
766 other => panic!("unsupported trained test index type: {other:?}"),
767 };
768 let field = builder.add_dense_vector_field_with_config("embedding", true, true, config);
769 (builder.build(), field)
770 }
771
772 fn test_centroids() -> crate::structures::CoarseCentroids {
773 crate::structures::CoarseCentroids {
774 num_clusters: 1,
775 dim: 2,
776 centroids: vec![0.25, 0.75],
777 version: crate::structures::mark_ivf_tq_cosine_generation(7),
778 soar_config: None,
779 routing_index: None,
780 }
781 }
782
783 async fn write_bincode(
784 directory: &crate::directories::RamDirectory,
785 filename: &str,
786 value: &impl serde::Serialize,
787 ) {
788 let bytes = bincode::serde::encode_to_vec(value, bincode::config::standard()).unwrap();
789 directory.write(Path::new(filename), &bytes).await.unwrap();
790 }
791
792 #[test]
793 fn test_metadata_init() {
794 let mut meta = IndexMetadata::new(test_schema());
795 assert_eq!(meta.total_vectors, 0);
796 assert!(meta.segment_metas.is_empty());
797 assert!(!meta.is_field_built(0));
798
799 meta.init_field(0, VectorIndexType::IvfTq);
800 assert!(!meta.is_field_built(0));
801 assert!(meta.vector_fields.contains_key(&0));
802 }
803
804 #[tokio::test]
805 async fn load_refuses_metadata_stamped_with_a_newer_format_version() {
806 let directory = crate::directories::RamDirectory::new();
807 let mut metadata = IndexMetadata::new(test_schema());
808 metadata.version = INDEX_META_FORMAT_VERSION + 1;
809 metadata.save(&directory).await.unwrap();
810
811 let error = IndexMetadata::load(&directory)
812 .await
813 .expect_err("metadata from a newer format version must be refused, not silently pruned")
814 .to_string();
815 assert!(error.contains("version 4"), "{error}");
816 assert!(error.contains("incompatible"), "{error}");
817 }
818
819 #[tokio::test]
820 async fn tmp_recovery_refuses_metadata_stamped_with_a_newer_format_version() {
821 let directory = crate::directories::RamDirectory::new();
822 let mut metadata = IndexMetadata::new(test_schema());
823 metadata.version = INDEX_META_FORMAT_VERSION + 1;
824 let bytes = metadata.serialize_to_bytes().unwrap();
825 directory
827 .write(Path::new(INDEX_META_TMP_FILENAME), &bytes)
828 .await
829 .unwrap();
830
831 let error = IndexMetadata::load(&directory)
832 .await
833 .expect_err("temp-file recovery must apply the same version gate")
834 .to_string();
835 assert!(error.contains("version 4"), "{error}");
836 }
837
838 #[tokio::test]
839 async fn save_treats_post_rename_sync_failure_as_committed() {
840 let directory = SyncFailDirectory::default();
841 let mut metadata = IndexMetadata::new(test_schema());
842 metadata.add_segment("committed".to_string(), 7);
843
844 metadata.save(&directory).await.unwrap();
845
846 let loaded = IndexMetadata::load(&directory).await.unwrap();
847 assert_eq!(loaded.segment_doc_count("committed"), Some(7));
848 }
849
850 #[tokio::test]
851 async fn trained_artifacts_load_only_when_the_complete_built_set_is_valid() {
852 let mut builder = crate::dsl::SchemaBuilder::default();
853 let config = crate::dsl::DenseVectorConfig::ivf_tq(2, Some(1), 1);
854 let first = builder.add_dense_vector_field_with_config(
855 "first_embedding",
856 true,
857 true,
858 config.clone(),
859 );
860 let second =
861 builder.add_dense_vector_field_with_config("second_embedding", true, true, config);
862 let schema = builder.build();
863 let directory = crate::directories::RamDirectory::new();
864 let mut metadata = IndexMetadata::new(schema.clone());
865 metadata.init_field(first.0, VectorIndexType::IvfTq);
866 metadata.init_field(second.0, VectorIndexType::IvfTq);
867 metadata.mark_field_built(first.0, 10, 1, "field_0_centroids.bin".into(), None);
868 metadata.mark_field_built(second.0, 10, 1, "field_1_centroids.bin".into(), None);
869 write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
870
871 let error = IndexMetadata::try_load_trained_from_fields(
872 &metadata.vector_fields,
873 &schema,
874 &directory,
875 )
876 .await
877 .err()
878 .expect("missing artifact must fail the complete load")
879 .to_string();
880 assert!(error.contains("field_1_centroids.bin"), "{error}");
881 assert!(error.contains("field 1"), "{error}");
882 }
883
884 #[tokio::test]
885 async fn index_open_fails_closed_when_built_artifact_is_missing() {
886 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
887 let directory = crate::directories::RamDirectory::new();
888 let mut metadata = IndexMetadata::new(schema);
889 metadata.init_field(field.0, VectorIndexType::IvfTq);
890 metadata.mark_field_built(field.0, 10, 1, "missing_centroids.bin".into(), None);
891 metadata.save(&directory).await.unwrap();
892
893 let error = match crate::index::Index::open(directory, crate::index::IndexConfig::default())
894 .await
895 {
896 Ok(_) => panic!("Index::open accepted a Built field with no artifact"),
897 Err(error) => error.to_string(),
898 };
899 assert!(error.contains("missing_centroids.bin"), "{error}");
900 }
901
902 #[tokio::test]
903 async fn ivf_tq_built_state_rejects_a_codebook_file() {
904 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
905 let directory = crate::directories::RamDirectory::new();
906 let mut metadata = IndexMetadata::new(schema.clone());
907 metadata.init_field(field.0, VectorIndexType::IvfTq);
908 metadata.mark_field_built(
909 field.0,
910 10,
911 1,
912 "field_0_centroids.bin".into(),
913 Some("field_0_codebook.bin".into()),
914 );
915 write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
916
917 let error = IndexMetadata::try_load_trained_from_fields(
918 &metadata.vector_fields,
919 &schema,
920 &directory,
921 )
922 .await
923 .err()
924 .expect("IVF-TQ Built state with a codebook file must fail")
925 .to_string();
926 assert!(error.contains("codebook"), "{error}");
927 }
928
929 #[tokio::test]
930 async fn legacy_ivf_tq_centroid_generation_is_rejected_while_loading() {
931 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
932 let directory = crate::directories::RamDirectory::new();
933 let mut metadata = IndexMetadata::new(schema.clone());
934 metadata.init_field(field.0, VectorIndexType::IvfTq);
935 metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
936 let mut legacy = test_centroids();
937 legacy.version = 7;
938 write_bincode(&directory, "field_0_centroids.bin", &legacy).await;
939
940 let error = IndexMetadata::try_load_trained_from_fields(
941 &metadata.vector_fields,
942 &schema,
943 &directory,
944 )
945 .await
946 .err()
947 .expect("legacy IVF-TQ centroid state must fail while loading")
948 .to_string();
949 assert!(error.contains("unsupported legacy generation"), "{error}");
950 assert!(error.contains("rebuild the index"), "{error}");
951 }
952
953 #[tokio::test]
954 async fn legacy_ivf_pq_trained_field_fails_with_actionable_error() {
955 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
959 let directory = crate::directories::RamDirectory::new();
960 let mut metadata = IndexMetadata::new(schema.clone());
961 metadata.init_field(field.0, VectorIndexType::IvfTq);
962 metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
963 metadata
966 .vector_fields
967 .get_mut(&field.0)
968 .expect("field initialized")
969 .index_type = VectorFieldIndexType::Float(VectorIndexType::IvfPq);
970 write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
971
972 let error = IndexMetadata::try_load_trained_from_fields(
973 &metadata.vector_fields,
974 &schema,
975 &directory,
976 )
977 .await
978 .err()
979 .expect("legacy IVF-PQ trained state must fail loudly")
980 .to_string();
981 assert!(error.contains("no longer"), "{error}");
982 assert!(error.contains("ivf_tq"), "{error}");
983 }
984
985 #[tokio::test]
986 async fn requested_cluster_count_accepts_a_quality_clamped_artifact() {
987 let mut builder = crate::dsl::SchemaBuilder::default();
988 let field = builder.add_dense_vector_field_with_config(
989 "embedding",
990 true,
991 true,
992 crate::dsl::DenseVectorConfig::ivf_tq(2, Some(4), 1),
993 );
994 let schema = builder.build();
995 let directory = crate::directories::RamDirectory::new();
996 let mut metadata = IndexMetadata::new(schema.clone());
997 metadata.init_field(field.0, VectorIndexType::IvfTq);
998 metadata.mark_field_built(field.0, 1, 4, "field_0_centroids.bin".into(), None);
999 write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
1000
1001 let trained = IndexMetadata::try_load_trained_from_fields(
1002 &metadata.vector_fields,
1003 &schema,
1004 &directory,
1005 )
1006 .await
1007 .unwrap()
1008 .unwrap();
1009 assert_eq!(trained.centroids[&field.0].num_clusters, 1);
1010 }
1011
1012 #[tokio::test]
1013 async fn trained_artifact_loader_rejects_trailing_data() {
1014 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1015 let directory = crate::directories::RamDirectory::new();
1016 let mut metadata = IndexMetadata::new(schema.clone());
1017 metadata.init_field(field.0, VectorIndexType::IvfTq);
1018 metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1019 let mut bytes =
1020 bincode::serde::encode_to_vec(test_centroids(), bincode::config::standard()).unwrap();
1021 bytes.extend_from_slice(&[0xaa, 0xbb]);
1022 directory
1023 .write(Path::new("field_0_centroids.bin"), &bytes)
1024 .await
1025 .unwrap();
1026
1027 let error = IndexMetadata::try_load_trained_from_fields(
1028 &metadata.vector_fields,
1029 &schema,
1030 &directory,
1031 )
1032 .await
1033 .err()
1034 .expect("trailing artifact bytes must fail validation")
1035 .to_string();
1036 assert!(error.contains("trailing bytes"), "{error}");
1037 }
1038
1039 #[test]
1040 fn trained_artifact_size_limit_rejects_before_reading() {
1041 let error = validate_trained_artifact_size(
1042 3,
1043 "centroids",
1044 "field_3_centroids.bin",
1045 MAX_TRAINED_ARTIFACT_BYTES as u64 + 1,
1046 )
1047 .unwrap_err()
1048 .to_string();
1049 assert!(error.contains("exceeding"), "{error}");
1050 assert!(error.contains("field 3"), "{error}");
1051 }
1052
1053 #[tokio::test]
1054 async fn trained_artifact_decode_limit_rejects_forged_collection_length() {
1055 let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1056 let directory = crate::directories::RamDirectory::new();
1057 let mut metadata = IndexMetadata::new(schema.clone());
1058 metadata.init_field(field.0, VectorIndexType::IvfTq);
1059 metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1060
1061 let mut bytes = vec![1, 2, 253];
1066 bytes.extend_from_slice(&u64::MAX.to_le_bytes());
1067 directory
1068 .write(Path::new("field_0_centroids.bin"), &bytes)
1069 .await
1070 .unwrap();
1071
1072 let error = IndexMetadata::try_load_trained_from_fields(
1073 &metadata.vector_fields,
1074 &schema,
1075 &directory,
1076 )
1077 .await
1078 .err()
1079 .expect("forged collection length must fail the bounded decoder")
1080 .to_string();
1081 assert!(error.contains("failed to deserialize"), "{error}");
1082 }
1083
1084 #[test]
1085 fn test_metadata_segments() {
1086 let mut meta = IndexMetadata::new(test_schema());
1087 meta.add_segment("abc123".to_string(), 50);
1088 meta.add_segment("def456".to_string(), 100);
1089 assert_eq!(meta.segment_metas.len(), 2);
1090 assert_eq!(meta.segment_doc_count("abc123"), Some(50));
1091 assert_eq!(meta.segment_doc_count("def456"), Some(100));
1092
1093 meta.add_segment("abc123".to_string(), 75);
1095 assert_eq!(meta.segment_metas.len(), 2);
1096 assert_eq!(meta.segment_doc_count("abc123"), Some(75));
1097
1098 meta.remove_segment("abc123");
1099 assert_eq!(meta.segment_metas.len(), 1);
1100 assert!(meta.has_segment("def456"));
1101 assert!(!meta.has_segment("abc123"));
1102 }
1103
1104 #[test]
1105 fn test_mark_field_built() {
1106 let mut meta = IndexMetadata::new(test_schema());
1107 meta.init_field(0, VectorIndexType::IvfTq);
1108 meta.total_vectors = 10000;
1109
1110 assert!(!meta.is_field_built(0));
1111
1112 meta.mark_field_built(0, 10000, 256, "field_0_centroids.bin".to_string(), None);
1113
1114 assert!(meta.is_field_built(0));
1115 let field = meta.get_field_meta(0).unwrap();
1116 assert_eq!(
1117 field.centroids_file.as_deref(),
1118 Some("field_0_centroids.bin")
1119 );
1120 }
1121
1122 #[test]
1123 fn total_vectors_is_aggregate_of_built_field_counts() {
1124 let mut meta = IndexMetadata::new(test_schema());
1125 meta.init_field(7, VectorIndexType::IvfTq);
1126 meta.init_field(3, VectorIndexType::IvfTq);
1127
1128 meta.mark_field_built(7, 400, 20, "field_7_centroids.bin".to_string(), None);
1131 assert_eq!(meta.total_vectors, 400);
1132 meta.mark_field_built(3, 250, 15, "field_3_centroids.bin".to_string(), None);
1133 assert_eq!(meta.total_vectors, 650);
1134
1135 meta.mark_field_built(7, 425, 20, "field_7_centroids.bin".to_string(), None);
1138 assert_eq!(meta.total_vectors, 675);
1139 }
1140
1141 #[test]
1142 fn test_should_build_field() {
1143 let mut meta = IndexMetadata::new(test_schema());
1144 meta.init_field(0, VectorIndexType::IvfTq);
1145
1146 meta.total_vectors = 500;
1148 assert!(!meta.should_build_field(0, 1000));
1149
1150 meta.total_vectors = 1500;
1152 assert!(meta.should_build_field(0, 1000));
1153
1154 meta.mark_field_built(0, 1500, 256, "centroids.bin".to_string(), None);
1156 assert!(!meta.should_build_field(0, 1000));
1157 }
1158
1159 #[test]
1160 fn test_serialization() {
1161 let mut meta = IndexMetadata::new(test_schema());
1162 meta.add_segment("seg1".to_string(), 100);
1163 meta.init_field(0, VectorIndexType::IvfTq);
1164 meta.total_vectors = 5000;
1165
1166 let json = serde_json::to_string_pretty(&meta).unwrap();
1167 let loaded: IndexMetadata = serde_json::from_str(&json).unwrap();
1168
1169 assert_eq!(loaded.segment_ids().len(), meta.segment_ids().len());
1170 assert_eq!(loaded.segment_doc_count("seg1"), Some(100));
1171 assert_eq!(loaded.total_vectors, meta.total_vectors);
1172 assert!(loaded.vector_fields.contains_key(&0));
1173 }
1174
1175 #[test]
1176 fn old_metadata_defaults_the_bp_retry_counter() {
1177 let mut meta = IndexMetadata::new(test_schema());
1178 meta.add_segment("legacy".to_string(), 10);
1179 let mut json = serde_json::to_value(&meta).unwrap();
1180 json["segment_metas"]["legacy"]
1181 .as_object_mut()
1182 .unwrap()
1183 .remove("bp_unconverged_passes");
1184
1185 let loaded: IndexMetadata = serde_json::from_value(json).unwrap();
1186 assert_eq!(loaded.segment_metas["legacy"].bp_unconverged_passes, 0);
1187 }
1188
1189 #[test]
1190 fn test_merged_segment_lineage() {
1191 let mut meta = IndexMetadata::new(test_schema());
1192 meta.add_segment("a".to_string(), 50);
1193 meta.add_segment("b".to_string(), 75);
1194
1195 assert_eq!(meta.segment_metas["a"].generation, 0);
1197 assert!(meta.segment_metas["a"].ancestors.is_empty());
1198
1199 meta.add_merged_segment(
1201 "c".to_string(),
1202 125,
1203 vec!["a".to_string(), "b".to_string()],
1204 1,
1205 false,
1206 true,
1207 );
1208 assert_eq!(meta.segment_metas["c"].generation, 1);
1209 assert_eq!(meta.segment_metas["c"].ancestors, vec!["a", "b"]);
1210 assert_eq!(meta.segment_doc_count("c"), Some(125));
1211
1212 meta.add_segment("d".to_string(), 30);
1214 meta.add_merged_segment(
1215 "e".to_string(),
1216 155,
1217 vec!["c".to_string(), "d".to_string()],
1218 2,
1219 false,
1220 true,
1221 );
1222 assert_eq!(meta.segment_metas["e"].generation, 2);
1223 }
1224}