Skip to main content

hermes_core/index/
metadata.rs

1//! Unified index metadata - segments list + vector index state
2//!
3//! This module manages all index-level metadata in a single `metadata.json` file:
4//! - List of committed segments
5//! - Vector index state per field (Flat/Built)
6//! - Trained centroids/codebooks paths
7//!
8//! The workflow is:
9//! 1. During initial accumulation, segments store flat vectors.
10//! 2. A manual build trains the first global codebook and ANN generation.
11//! 3. A manual retrain stages and atomically publishes a replacement generation.
12//! 4. On index open, metadata loads the currently published codebooks.
13
14use 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
22/// Metadata file name at index level
23pub const INDEX_META_FILENAME: &str = "metadata.json";
24/// Temp file for atomic writes (write here, then rename to INDEX_META_FILENAME)
25const INDEX_META_TMP_FILENAME: &str = "metadata.json.tmp";
26
27/// Current metadata.json format version written by this build.
28///
29/// `load` refuses metadata stamped with a newer version: serde_json silently
30/// drops fields it does not know about, so loading newer metadata would
31/// misread index state and the next save would destructively rewrite the
32/// unknown fields away.
33pub const INDEX_META_FORMAT_VERSION: u32 = 4;
34
35/// Index-level centroids/codebooks are deliberately bounded before they are
36/// read or decoded. Besides limiting ordinary corruption damage, the matching
37/// bincode limit prevents a tiny forged collection length from requesting an
38/// effectively unbounded allocation.
39pub(crate) const MAX_TRAINED_ARTIFACT_BYTES: usize = 512 * 1024 * 1024;
40
41/// State of vector index for a field
42#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
43pub enum VectorIndexState {
44    /// Accumulating vectors - using Flat (brute-force) search
45    #[default]
46    Flat,
47    /// Index structures built - using ANN search
48    Built {
49        /// Total vector count when training happened
50        vector_count: usize,
51        /// Number of clusters used
52        num_clusters: usize,
53    },
54}
55
56fn default_true() -> bool {
57    true
58}
59
60/// Per-segment metadata stored in index metadata
61/// This allows merge decisions without loading segment files
62#[derive(Debug, Clone, Serialize, Deserialize)]
63pub struct SegmentMetaInfo {
64    /// Number of documents in this segment
65    pub num_docs: u32,
66    /// Parent segment IDs that were merged to produce this segment (empty for fresh segments)
67    pub ancestors: Vec<String>,
68    /// Merge generation: 0 for fresh segments, max(parent generations) + 1 for merged segments
69    pub generation: u32,
70    /// Whether this segment has been reordered via Recursive Graph Bisection (BP).
71    /// Fresh segments and block-copy merges are not reordered. Only segments that have
72    /// been explicitly reordered (via background optimizer or reorder command) are marked true.
73    #[serde(default)]
74    pub reordered: bool,
75    /// Whether the last BP reorder pass ran to natural convergence. False when
76    /// a wall-clock BP budget ended the pass early — the segment is ordered
77    /// better than before, and a later warm-started pass can deepen it.
78    /// Old metadata (field absent) deserializes as converged.
79    #[serde(default = "default_true")]
80    pub bp_converged: bool,
81    /// Number of consecutive budget-exhausted BP rewrites in this segment's
82    /// current reordered lineage. Carried across replacement IDs so the
83    /// optimizer can impose a hard follow-up bound instead of rewriting forever.
84    #[serde(default)]
85    pub bp_unconverged_passes: u32,
86}
87
88/// Per-field vector index metadata
89#[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    /// Field ID
111    pub field_id: u32,
112    /// Configured index type (target type when built)
113    pub index_type: VectorFieldIndexType,
114    /// Current state
115    pub state: VectorIndexState,
116    /// Path to centroids file (relative to index dir)
117    #[serde(skip_serializing_if = "Option::is_none")]
118    pub centroids_file: Option<String>,
119    /// Path to the index-global codebook file, relative to the index directory.
120    #[serde(skip_serializing_if = "Option::is_none")]
121    pub codebook_file: Option<String>,
122}
123
124/// Unified index metadata - single source of truth for index state
125#[derive(Debug, Clone, Serialize, Deserialize)]
126pub struct IndexMetadata {
127    /// Version for compatibility
128    pub version: u32,
129    /// Index schema
130    pub schema: Schema,
131    /// Segment metadata: segment_id -> info (doc count, etc.)
132    /// Using HashMap allows O(1) lookup and stores doc counts for merge decisions
133    #[serde(default)]
134    pub segment_metas: HashMap<String, SegmentMetaInfo>,
135    /// Per-field vector index metadata
136    #[serde(default)]
137    pub vector_fields: HashMap<u32, FieldVectorMeta>,
138    /// Aggregate vector count recorded by all built vector fields.
139    ///
140    /// The per-field `VectorIndexState::Built::vector_count` values are the
141    /// source of truth. This cached aggregate is refreshed whenever a field is
142    /// marked built, rather than being overwritten with whichever field was
143    /// trained last.
144    #[serde(default)]
145    pub total_vectors: usize,
146}
147
148impl IndexMetadata {
149    /// Create new metadata with schema
150    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    /// Get segment IDs as a sorted Vec (deterministic ordering)
161    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    /// Add a fresh segment (gen=0, no ancestors, not reordered)
168    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    /// Add a merged segment with lineage info
183    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    /// Insert fully constructed lifecycle metadata. Merge/reorder code uses
206    /// this to carry bounded BP lineage; ordinary callers use the safer
207    /// constructors above, which start a fresh lineage.
208    pub(crate) fn add_segment_meta(&mut self, segment_id: String, info: SegmentMetaInfo) {
209        self.segment_metas.insert(segment_id, info);
210    }
211
212    /// Remove a segment
213    pub fn remove_segment(&mut self, segment_id: &str) {
214        self.segment_metas.remove(segment_id);
215    }
216
217    /// Check if segment exists
218    pub fn has_segment(&self, segment_id: &str) -> bool {
219        self.segment_metas.contains_key(segment_id)
220    }
221
222    /// Get segment doc count
223    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    /// Check if a field has been built
228    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    /// Get field metadata
236    pub fn get_field_meta(&self, field_id: u32) -> Option<&FieldVectorMeta> {
237        self.vector_fields.get(&field_id)
238    }
239
240    /// Initialize field metadata (called when field is first seen)
241    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    /// Mark field as built with trained structures
255    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    /// Refresh the cached aggregate from the authoritative per-field states.
275    ///
276    /// Saturation keeps this infallible metadata helper safe even if it is
277    /// called after loading externally modified metadata with impossible
278    /// counts.
279    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    /// Check if field should be built based on threshold
291    pub fn should_build_field(&self, field_id: u32, threshold: usize) -> bool {
292        // Don't build if already built
293        if self.is_field_built(field_id) {
294            return false;
295        }
296        // Build if we have enough vectors
297        self.total_vectors >= threshold
298    }
299
300    /// Load from directory
301    ///
302    /// If `metadata.json` is missing but `metadata.json.tmp` exists (crash
303    /// between write and rename), recovers from the temp file.
304    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                // Try recovering from temp file (crash between write and rename)
313                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    /// Deserialize only the current format. Vector artifacts are intentionally
325    /// rebuilt when the ANN format changes; silently accepting older metadata
326    /// would mix incompatible segment and global-codebook generations.
327    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    /// Save to directory (atomic: write temp file, then rename)
341    ///
342    /// Uses write-then-rename so a crash mid-write won't corrupt the
343    /// existing metadata file. On POSIX, rename is atomic.
344    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    /// Serialize metadata to bytes (cheap, no I/O).
350    /// Useful when you need to release a lock before doing disk I/O.
351    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    /// Write pre-serialized metadata bytes to directory (atomic rename + fsync).
356    ///
357    /// The fsync ensures durability: without it, a power failure after rename
358    /// could lose the metadata update on systems with volatile write caches.
359    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        // Metadata is tiny, but `DirectoryWriter::write` does not guarantee
366        // the file contents themselves are fsynced. Finish the streaming
367        // writer first (filesystem implementations call `File::sync_all`),
368        // then atomically publish the durable temp file by rename.
369        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        // Rename is the logical commit point: after it succeeds, readers can
373        // observe the new generation and callers must publish the matching
374        // in-memory/tracker state. Directory fsync only strengthens crash
375        // durability. It cannot safely turn an already-visible rename into a
376        // reported pre-commit failure, because cleanup could then delete files
377        // referenced by the metadata now on disk.
378        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    /// Fallible schema-aware loader used for lifecycle publication.
390    #[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    /// Load and validate the complete trained-artifact set described by a
400    /// `vector_fields` snapshot.
401    ///
402    /// This is intentionally all-or-nothing. A `Built` field is a durable
403    /// promise that every artifact required by its configured index exists and
404    /// is compatible with the schema. Returning a partial map would let some
405    /// segment builders publish ANN data while another field was silently
406    /// unusable, and would make the same index behave differently after a
407    /// restart.
408    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                        // IVF-TQ's leaf codec is derived, never trained.
536                        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        // Simulate a crash between write and rename: only the temp file exists.
844        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        // CoarseCentroids begins with num_clusters=1, dim=2, then the Vec
1038        // length. Bincode's standard varint marker 253 introduces a u64; this
1039        // tiny payload claims an impossible f32 vector and must hit the decode
1040        // limit before any large allocation is attempted.
1041        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        // Overwrites existing
1070        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        // Build in reverse field-id order to ensure the result is not tied to
1105        // HashMap or training iteration order.
1106        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        // Rebuilding a field replaces its contribution; it does not add a
1118        // duplicate training snapshot.
1119        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        // Below threshold
1129        meta.total_vectors = 500;
1130        assert!(!meta.should_build_field(0, 1000));
1131
1132        // Above threshold
1133        meta.total_vectors = 1500;
1134        assert!(meta.should_build_field(0, 1000));
1135
1136        // Already built - should not build again
1137        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        // Fresh segments: gen=0, no ancestors
1178        assert_eq!(meta.segment_metas["a"].generation, 0);
1179        assert!(meta.segment_metas["a"].ancestors.is_empty());
1180
1181        // Merge a+b → c
1182        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        // Merge c+d → e (gen should be 2)
1195        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}