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