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