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