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, 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(index_type @ VectorIndexType::IvfPq) => {
461                    let entry = schema
462                        .get_field_entry(crate::dsl::Field(*field_id))
463                        .ok_or_else(|| {
464                            Error::Corruption(format!(
465                                "trained vector metadata references missing field {field_id}"
466                            ))
467                        })?;
468                    let schema_config = entry
469                        .dense_vector_config
470                        .as_ref()
471                        .filter(|_| entry.field_type == crate::dsl::FieldType::DenseVector)
472                        .ok_or_else(|| {
473                            Error::Corruption(format!(
474                                "trained vector metadata field {field_id} is not a float dense field"
475                            ))
476                        })?;
477                    if schema_config.index_type != index_type {
478                        return Err(Error::Corruption(format!(
479                            "trained vector metadata field {field_id} uses {index_type:?}, schema requires {:?}",
480                            schema_config.index_type
481                        )));
482                    }
483                    let c: crate::structures::CoarseCentroids =
484                        load_trained_artifact(dir, *field_id, "centroids", centroids_file).await?;
485                    let expected_dim = schema_config.dim;
486                    let actual_clusters = c.num_clusters as usize;
487                    let expected_values =
488                        actual_clusters.checked_mul(expected_dim).ok_or_else(|| {
489                            Error::Corruption(format!(
490                                "trained centroid dimensions overflow for field {field_id}"
491                            ))
492                        })?;
493                    if actual_clusters == 0
494                        || actual_clusters > expected_clusters
495                        || c.dim == 0
496                        || c.dim != expected_dim
497                        || c.centroids.len() != expected_values
498                        || c.centroids.iter().any(|value| !value.is_finite())
499                    {
500                        return Err(Error::Corruption(format!(
501                            "trained centroids for field {field_id} do not match metadata/schema"
502                        )));
503                    }
504                    c.validate_routing(schema_config.ivf_routing)
505                        .map_err(|error| {
506                            Error::Corruption(format!(
507                                "invalid trained centroid routing for field {field_id}: {error}"
508                            ))
509                        })?;
510                    let codebook_file = field_meta.codebook_file.as_deref().ok_or_else(|| {
511                        Error::Corruption(format!(
512                            "trained float IVF-PQ field {field_id} has no codebook_file"
513                        ))
514                    })?;
515                    let codebook: crate::structures::PQCodebook =
516                        load_trained_artifact(dir, *field_id, "codebook", codebook_file).await?;
517                    codebook.validate().map_err(|error| {
518                        Error::Corruption(format!(
519                            "invalid trained codebook for field {field_id}: {error}"
520                        ))
521                    })?;
522                    if codebook.config.dim != expected_dim {
523                        return Err(Error::Corruption(format!(
524                            "trained codebook for field {field_id} has dimension {}, expected {expected_dim}",
525                            codebook.config.dim
526                        )));
527                    }
528                    codebooks.insert(*field_id, Arc::new(codebook));
529                    centroids.insert(*field_id, Arc::new(c));
530                }
531                VectorFieldIndexType::Binary(BinaryIndexType::Ivf) => {
532                    let entry = schema
533                        .get_field_entry(crate::dsl::Field(*field_id))
534                        .ok_or_else(|| {
535                            Error::Corruption(format!(
536                                "trained vector metadata references missing field {field_id}"
537                            ))
538                        })?;
539                    let schema_config = entry
540                        .binary_dense_vector_config
541                        .as_ref()
542                        .filter(|config| {
543                            entry.field_type == crate::dsl::FieldType::BinaryDenseVector
544                                && config.index_type == BinaryIndexType::Ivf
545                        })
546                        .ok_or_else(|| {
547                            Error::Corruption(format!(
548                                "trained vector metadata field {field_id} is not a binary IVF field"
549                            ))
550                        })?;
551                    let quantizer: crate::structures::BinaryCoarseQuantizer =
552                        load_trained_artifact(dir, *field_id, "binary centroids", centroids_file)
553                            .await?;
554                    quantizer.validate().map_err(|error| {
555                        Error::Corruption(format!(
556                            "invalid binary coarse quantizer for field {field_id}: {error}"
557                        ))
558                    })?;
559                    let actual_clusters = quantizer.num_clusters as usize;
560                    if actual_clusters > expected_clusters
561                        || schema_config.dim != quantizer.dim_bits
562                    {
563                        return Err(Error::Corruption(format!(
564                            "binary coarse quantizer for field {field_id} does not match metadata/schema"
565                        )));
566                    }
567                    quantizer
568                        .validate_routing(schema_config.ivf_routing)
569                        .map_err(|error| {
570                            Error::Corruption(format!(
571                                "invalid binary centroid routing for field {field_id}: {error}"
572                            ))
573                        })?;
574                    binary_quantizers.insert(*field_id, Arc::new(quantizer));
575                }
576                unsupported => {
577                    return Err(Error::Corruption(format!(
578                        "field {field_id} is Built for {unsupported:?}, which has no global IVF artifacts"
579                    )));
580                }
581            }
582        }
583
584        if centroids.is_empty() && binary_quantizers.is_empty() {
585            Ok(None)
586        } else {
587            let trained = crate::segment::TrainedVectorStructures {
588                #[cfg(feature = "native")]
589                _ann_pins: Default::default(),
590                centroids,
591                binary_quantizers,
592                codebooks,
593            };
594            #[cfg(feature = "native")]
595            let trained = {
596                let mut trained = trained;
597                trained.pin_ann_structures(crate::segment::pin::pin_policy());
598                trained
599            };
600            Ok(Some(trained))
601        }
602    }
603}
604
605fn validate_trained_artifact_path(field_id: u32, kind: &str, filename: &str) -> Result<()> {
606    use std::path::Component;
607
608    let path = Path::new(filename);
609    if filename.is_empty()
610        || path.is_absolute()
611        || path.components().any(|component| {
612            matches!(
613                component,
614                Component::ParentDir | Component::RootDir | Component::Prefix(_)
615            )
616        })
617    {
618        return Err(Error::Corruption(format!(
619            "trained {kind} path for field {field_id} is not a safe relative path: '{filename}'"
620        )));
621    }
622    Ok(())
623}
624
625async fn load_trained_artifact<T, D>(
626    dir: &D,
627    field_id: u32,
628    kind: &str,
629    filename: &str,
630) -> Result<T>
631where
632    T: serde::de::DeserializeOwned,
633    D: crate::directories::Directory,
634{
635    validate_trained_artifact_path(field_id, kind, filename)?;
636    let path = Path::new(filename);
637    let file_size = dir.file_size(path).await.map_err(|error| {
638        Error::Corruption(format!(
639            "failed to stat trained {kind} '{filename}' for field {field_id}: {error}"
640        ))
641    })?;
642    validate_trained_artifact_size(field_id, kind, filename, file_size)?;
643    let slice = dir.open_read(path).await.map_err(|error| {
644        Error::Corruption(format!(
645            "failed to open trained {kind} '{filename}' for field {field_id}: {error}"
646        ))
647    })?;
648    validate_trained_artifact_size(field_id, kind, filename, slice.len())?;
649    let bytes = slice.read_bytes().await.map_err(|error| {
650        Error::Corruption(format!(
651            "failed to read trained {kind} '{filename}' for field {field_id}: {error}"
652        ))
653    })?;
654    let (artifact, consumed) = bincode::serde::decode_from_slice::<T, _>(
655        bytes.as_slice(),
656        bincode::config::standard().with_limit::<MAX_TRAINED_ARTIFACT_BYTES>(),
657    )
658    .map_err(|error| {
659        Error::Corruption(format!(
660            "failed to deserialize trained {kind} '{filename}' for field {field_id}: {error}"
661        ))
662    })?;
663    if consumed != bytes.len() {
664        return Err(Error::Corruption(format!(
665            "trained {kind} '{filename}' for field {field_id} has {} trailing bytes",
666            bytes.len() - consumed
667        )));
668    }
669    Ok(artifact)
670}
671
672fn validate_trained_artifact_size(
673    field_id: u32,
674    kind: &str,
675    filename: &str,
676    file_size: u64,
677) -> Result<()> {
678    if file_size > MAX_TRAINED_ARTIFACT_BYTES as u64 {
679        return Err(Error::Corruption(format!(
680            "trained {kind} '{filename}' for field {field_id} is {file_size} bytes, \
681             exceeding the {MAX_TRAINED_ARTIFACT_BYTES}-byte safety limit"
682        )));
683    }
684    Ok(())
685}
686
687#[cfg(test)]
688mod tests {
689    use super::*;
690    use crate::directories::DirectoryWriter;
691
692    #[derive(Clone, Default)]
693    struct SyncFailDirectory(crate::directories::RamDirectory);
694
695    #[async_trait::async_trait]
696    impl crate::directories::Directory for SyncFailDirectory {
697        async fn exists(&self, path: &Path) -> std::io::Result<bool> {
698            self.0.exists(path).await
699        }
700
701        async fn file_size(&self, path: &Path) -> std::io::Result<u64> {
702            self.0.file_size(path).await
703        }
704
705        async fn open_read(&self, path: &Path) -> std::io::Result<crate::directories::FileHandle> {
706            self.0.open_read(path).await
707        }
708
709        async fn read_range(
710            &self,
711            path: &Path,
712            range: std::ops::Range<u64>,
713        ) -> std::io::Result<crate::directories::OwnedBytes> {
714            self.0.read_range(path, range).await
715        }
716
717        async fn list_files(&self, prefix: &Path) -> std::io::Result<Vec<std::path::PathBuf>> {
718            self.0.list_files(prefix).await
719        }
720
721        async fn open_lazy(&self, path: &Path) -> std::io::Result<crate::directories::FileHandle> {
722            self.0.open_lazy(path).await
723        }
724    }
725
726    #[async_trait::async_trait]
727    impl crate::directories::DirectoryWriter for SyncFailDirectory {
728        async fn write(&self, path: &Path, data: &[u8]) -> std::io::Result<()> {
729            self.0.write(path, data).await
730        }
731
732        async fn delete(&self, path: &Path) -> std::io::Result<()> {
733            self.0.delete(path).await
734        }
735
736        async fn rename(&self, from: &Path, to: &Path) -> std::io::Result<()> {
737            self.0.rename(from, to).await
738        }
739
740        async fn sync(&self) -> std::io::Result<()> {
741            Err(std::io::Error::other("injected directory fsync failure"))
742        }
743
744        async fn streaming_writer(
745            &self,
746            path: &Path,
747        ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
748            self.0.streaming_writer(path).await
749        }
750    }
751
752    fn test_schema() -> Schema {
753        Schema::default()
754    }
755
756    fn dense_schema(index_type: VectorIndexType) -> (Schema, crate::dsl::Field) {
757        let mut builder = crate::dsl::SchemaBuilder::default();
758        let config = match index_type {
759            VectorIndexType::IvfPq => crate::dsl::DenseVectorConfig::with_ivf_pq(2, Some(1), 1),
760            other => panic!("unsupported trained test index type: {other:?}"),
761        };
762        let field = builder.add_dense_vector_field_with_config("embedding", true, true, config);
763        (builder.build(), field)
764    }
765
766    fn test_centroids() -> crate::structures::CoarseCentroids {
767        crate::structures::CoarseCentroids {
768            num_clusters: 1,
769            dim: 2,
770            centroids: vec![0.25, 0.75],
771            version: 7,
772            soar_config: None,
773            routing_index: None,
774        }
775    }
776
777    fn test_codebook() -> crate::structures::PQCodebook {
778        let config = crate::structures::PQConfig::new(2).with_opq(false, 0);
779        crate::structures::PQCodebook {
780            centroids: vec![
781                0.0;
782                config.num_subspaces * config.num_centroids * config.dims_per_block
783            ],
784            rotation_matrix: None,
785            version: 11,
786            config,
787        }
788    }
789
790    async fn write_bincode(
791        directory: &crate::directories::RamDirectory,
792        filename: &str,
793        value: &impl serde::Serialize,
794    ) {
795        let bytes = bincode::serde::encode_to_vec(value, bincode::config::standard()).unwrap();
796        directory.write(Path::new(filename), &bytes).await.unwrap();
797    }
798
799    #[test]
800    fn test_metadata_init() {
801        let mut meta = IndexMetadata::new(test_schema());
802        assert_eq!(meta.total_vectors, 0);
803        assert!(meta.segment_metas.is_empty());
804        assert!(!meta.is_field_built(0));
805
806        meta.init_field(0, VectorIndexType::IvfPq);
807        assert!(!meta.is_field_built(0));
808        assert!(meta.vector_fields.contains_key(&0));
809    }
810
811    #[tokio::test]
812    async fn load_refuses_metadata_stamped_with_a_newer_format_version() {
813        let directory = crate::directories::RamDirectory::new();
814        let mut metadata = IndexMetadata::new(test_schema());
815        metadata.version = INDEX_META_FORMAT_VERSION + 1;
816        metadata.save(&directory).await.unwrap();
817
818        let error = IndexMetadata::load(&directory)
819            .await
820            .expect_err("metadata from a newer format version must be refused, not silently pruned")
821            .to_string();
822        assert!(error.contains("version 4"), "{error}");
823        assert!(error.contains("incompatible"), "{error}");
824    }
825
826    #[tokio::test]
827    async fn tmp_recovery_refuses_metadata_stamped_with_a_newer_format_version() {
828        let directory = crate::directories::RamDirectory::new();
829        let mut metadata = IndexMetadata::new(test_schema());
830        metadata.version = INDEX_META_FORMAT_VERSION + 1;
831        let bytes = metadata.serialize_to_bytes().unwrap();
832        // Simulate a crash between write and rename: only the temp file exists.
833        directory
834            .write(Path::new(INDEX_META_TMP_FILENAME), &bytes)
835            .await
836            .unwrap();
837
838        let error = IndexMetadata::load(&directory)
839            .await
840            .expect_err("temp-file recovery must apply the same version gate")
841            .to_string();
842        assert!(error.contains("version 4"), "{error}");
843    }
844
845    #[tokio::test]
846    async fn save_treats_post_rename_sync_failure_as_committed() {
847        let directory = SyncFailDirectory::default();
848        let mut metadata = IndexMetadata::new(test_schema());
849        metadata.add_segment("committed".to_string(), 7);
850
851        metadata.save(&directory).await.unwrap();
852
853        let loaded = IndexMetadata::load(&directory).await.unwrap();
854        assert_eq!(loaded.segment_doc_count("committed"), Some(7));
855    }
856
857    #[tokio::test]
858    async fn trained_artifacts_load_only_when_the_complete_built_set_is_valid() {
859        let mut builder = crate::dsl::SchemaBuilder::default();
860        let config = crate::dsl::DenseVectorConfig::with_ivf_pq(2, Some(1), 1);
861        let first = builder.add_dense_vector_field_with_config(
862            "first_embedding",
863            true,
864            true,
865            config.clone(),
866        );
867        let second =
868            builder.add_dense_vector_field_with_config("second_embedding", true, true, config);
869        let schema = builder.build();
870        let directory = crate::directories::RamDirectory::new();
871        let mut metadata = IndexMetadata::new(schema.clone());
872        metadata.init_field(first.0, VectorIndexType::IvfPq);
873        metadata.init_field(second.0, VectorIndexType::IvfPq);
874        metadata.mark_field_built(
875            first.0,
876            10,
877            1,
878            "field_0_centroids.bin".into(),
879            Some("field_0_codebook.bin".into()),
880        );
881        metadata.mark_field_built(
882            second.0,
883            10,
884            1,
885            "field_1_centroids.bin".into(),
886            Some("field_1_codebook.bin".into()),
887        );
888        write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
889        write_bincode(&directory, "field_0_codebook.bin", &test_codebook()).await;
890
891        let error = IndexMetadata::try_load_trained_from_fields(
892            &metadata.vector_fields,
893            &schema,
894            &directory,
895        )
896        .await
897        .err()
898        .expect("missing artifact must fail the complete load")
899        .to_string();
900        assert!(error.contains("field_1_centroids.bin"), "{error}");
901        assert!(error.contains("field 1"), "{error}");
902    }
903
904    #[tokio::test]
905    async fn index_open_fails_closed_when_built_artifact_is_missing() {
906        let (schema, field) = dense_schema(VectorIndexType::IvfPq);
907        let directory = crate::directories::RamDirectory::new();
908        let mut metadata = IndexMetadata::new(schema);
909        metadata.init_field(field.0, VectorIndexType::IvfPq);
910        metadata.mark_field_built(field.0, 10, 1, "missing_centroids.bin".into(), None);
911        metadata.save(&directory).await.unwrap();
912
913        let error = match crate::index::Index::open(directory, crate::index::IndexConfig::default())
914            .await
915        {
916            Ok(_) => panic!("Index::open accepted a Built field with no artifact"),
917            Err(error) => error.to_string(),
918        };
919        assert!(error.contains("missing_centroids.bin"), "{error}");
920    }
921
922    #[tokio::test]
923    async fn ivf_pq_built_state_requires_a_codebook() {
924        let (schema, field) = dense_schema(VectorIndexType::IvfPq);
925        let directory = crate::directories::RamDirectory::new();
926        let mut metadata = IndexMetadata::new(schema.clone());
927        metadata.init_field(field.0, VectorIndexType::IvfPq);
928        metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
929        write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
930
931        let error = IndexMetadata::try_load_trained_from_fields(
932            &metadata.vector_fields,
933            &schema,
934            &directory,
935        )
936        .await
937        .err()
938        .expect("IVF-PQ Built state without a codebook must fail")
939        .to_string();
940        assert!(error.contains("has no codebook_file"), "{error}");
941    }
942
943    #[tokio::test]
944    async fn requested_cluster_count_accepts_a_quality_clamped_artifact() {
945        let mut builder = crate::dsl::SchemaBuilder::default();
946        let field = builder.add_dense_vector_field_with_config(
947            "embedding",
948            true,
949            true,
950            crate::dsl::DenseVectorConfig::with_ivf_pq(2, Some(4), 1),
951        );
952        let schema = builder.build();
953        let directory = crate::directories::RamDirectory::new();
954        let mut metadata = IndexMetadata::new(schema.clone());
955        metadata.init_field(field.0, VectorIndexType::IvfPq);
956        metadata.mark_field_built(
957            field.0,
958            1,
959            4,
960            "field_0_centroids.bin".into(),
961            Some("field_0_codebook.bin".into()),
962        );
963        write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
964        write_bincode(&directory, "field_0_codebook.bin", &test_codebook()).await;
965
966        let trained = IndexMetadata::try_load_trained_from_fields(
967            &metadata.vector_fields,
968            &schema,
969            &directory,
970        )
971        .await
972        .unwrap()
973        .unwrap();
974        assert_eq!(trained.centroids[&field.0].num_clusters, 1);
975    }
976
977    #[tokio::test]
978    async fn trained_artifact_loader_rejects_trailing_data() {
979        let (schema, field) = dense_schema(VectorIndexType::IvfPq);
980        let directory = crate::directories::RamDirectory::new();
981        let mut metadata = IndexMetadata::new(schema.clone());
982        metadata.init_field(field.0, VectorIndexType::IvfPq);
983        metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
984        let mut bytes =
985            bincode::serde::encode_to_vec(test_centroids(), bincode::config::standard()).unwrap();
986        bytes.extend_from_slice(&[0xaa, 0xbb]);
987        directory
988            .write(Path::new("field_0_centroids.bin"), &bytes)
989            .await
990            .unwrap();
991
992        let error = IndexMetadata::try_load_trained_from_fields(
993            &metadata.vector_fields,
994            &schema,
995            &directory,
996        )
997        .await
998        .err()
999        .expect("trailing artifact bytes must fail validation")
1000        .to_string();
1001        assert!(error.contains("trailing bytes"), "{error}");
1002    }
1003
1004    #[test]
1005    fn trained_artifact_size_limit_rejects_before_reading() {
1006        let error = validate_trained_artifact_size(
1007            3,
1008            "centroids",
1009            "field_3_centroids.bin",
1010            MAX_TRAINED_ARTIFACT_BYTES as u64 + 1,
1011        )
1012        .unwrap_err()
1013        .to_string();
1014        assert!(error.contains("exceeding"), "{error}");
1015        assert!(error.contains("field 3"), "{error}");
1016    }
1017
1018    #[tokio::test]
1019    async fn trained_artifact_decode_limit_rejects_forged_collection_length() {
1020        let (schema, field) = dense_schema(VectorIndexType::IvfPq);
1021        let directory = crate::directories::RamDirectory::new();
1022        let mut metadata = IndexMetadata::new(schema.clone());
1023        metadata.init_field(field.0, VectorIndexType::IvfPq);
1024        metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1025
1026        // CoarseCentroids begins with num_clusters=1, dim=2, then the Vec
1027        // length. Bincode's standard varint marker 253 introduces a u64; this
1028        // tiny payload claims an impossible f32 vector and must hit the decode
1029        // limit before any large allocation is attempted.
1030        let mut bytes = vec![1, 2, 253];
1031        bytes.extend_from_slice(&u64::MAX.to_le_bytes());
1032        directory
1033            .write(Path::new("field_0_centroids.bin"), &bytes)
1034            .await
1035            .unwrap();
1036
1037        let error = IndexMetadata::try_load_trained_from_fields(
1038            &metadata.vector_fields,
1039            &schema,
1040            &directory,
1041        )
1042        .await
1043        .err()
1044        .expect("forged collection length must fail the bounded decoder")
1045        .to_string();
1046        assert!(error.contains("failed to deserialize"), "{error}");
1047    }
1048
1049    #[test]
1050    fn test_metadata_segments() {
1051        let mut meta = IndexMetadata::new(test_schema());
1052        meta.add_segment("abc123".to_string(), 50);
1053        meta.add_segment("def456".to_string(), 100);
1054        assert_eq!(meta.segment_metas.len(), 2);
1055        assert_eq!(meta.segment_doc_count("abc123"), Some(50));
1056        assert_eq!(meta.segment_doc_count("def456"), Some(100));
1057
1058        // Overwrites existing
1059        meta.add_segment("abc123".to_string(), 75);
1060        assert_eq!(meta.segment_metas.len(), 2);
1061        assert_eq!(meta.segment_doc_count("abc123"), Some(75));
1062
1063        meta.remove_segment("abc123");
1064        assert_eq!(meta.segment_metas.len(), 1);
1065        assert!(meta.has_segment("def456"));
1066        assert!(!meta.has_segment("abc123"));
1067    }
1068
1069    #[test]
1070    fn test_mark_field_built() {
1071        let mut meta = IndexMetadata::new(test_schema());
1072        meta.init_field(0, VectorIndexType::IvfPq);
1073        meta.total_vectors = 10000;
1074
1075        assert!(!meta.is_field_built(0));
1076
1077        meta.mark_field_built(0, 10000, 256, "field_0_centroids.bin".to_string(), None);
1078
1079        assert!(meta.is_field_built(0));
1080        let field = meta.get_field_meta(0).unwrap();
1081        assert_eq!(
1082            field.centroids_file.as_deref(),
1083            Some("field_0_centroids.bin")
1084        );
1085    }
1086
1087    #[test]
1088    fn total_vectors_is_aggregate_of_built_field_counts() {
1089        let mut meta = IndexMetadata::new(test_schema());
1090        meta.init_field(7, VectorIndexType::IvfPq);
1091        meta.init_field(3, VectorIndexType::IvfPq);
1092
1093        // Build in reverse field-id order to ensure the result is not tied to
1094        // HashMap or training iteration order.
1095        meta.mark_field_built(7, 400, 20, "field_7_centroids.bin".to_string(), None);
1096        assert_eq!(meta.total_vectors, 400);
1097        meta.mark_field_built(
1098            3,
1099            250,
1100            15,
1101            "field_3_centroids.bin".to_string(),
1102            Some("field_3_codebook.bin".to_string()),
1103        );
1104        assert_eq!(meta.total_vectors, 650);
1105
1106        // Rebuilding a field replaces its contribution; it does not add a
1107        // duplicate training snapshot.
1108        meta.mark_field_built(7, 425, 20, "field_7_centroids.bin".to_string(), None);
1109        assert_eq!(meta.total_vectors, 675);
1110    }
1111
1112    #[test]
1113    fn test_should_build_field() {
1114        let mut meta = IndexMetadata::new(test_schema());
1115        meta.init_field(0, VectorIndexType::IvfPq);
1116
1117        // Below threshold
1118        meta.total_vectors = 500;
1119        assert!(!meta.should_build_field(0, 1000));
1120
1121        // Above threshold
1122        meta.total_vectors = 1500;
1123        assert!(meta.should_build_field(0, 1000));
1124
1125        // Already built - should not build again
1126        meta.mark_field_built(0, 1500, 256, "centroids.bin".to_string(), None);
1127        assert!(!meta.should_build_field(0, 1000));
1128    }
1129
1130    #[test]
1131    fn test_serialization() {
1132        let mut meta = IndexMetadata::new(test_schema());
1133        meta.add_segment("seg1".to_string(), 100);
1134        meta.init_field(0, VectorIndexType::IvfPq);
1135        meta.total_vectors = 5000;
1136
1137        let json = serde_json::to_string_pretty(&meta).unwrap();
1138        let loaded: IndexMetadata = serde_json::from_str(&json).unwrap();
1139
1140        assert_eq!(loaded.segment_ids().len(), meta.segment_ids().len());
1141        assert_eq!(loaded.segment_doc_count("seg1"), Some(100));
1142        assert_eq!(loaded.total_vectors, meta.total_vectors);
1143        assert!(loaded.vector_fields.contains_key(&0));
1144    }
1145
1146    #[test]
1147    fn old_metadata_defaults_the_bp_retry_counter() {
1148        let mut meta = IndexMetadata::new(test_schema());
1149        meta.add_segment("legacy".to_string(), 10);
1150        let mut json = serde_json::to_value(&meta).unwrap();
1151        json["segment_metas"]["legacy"]
1152            .as_object_mut()
1153            .unwrap()
1154            .remove("bp_unconverged_passes");
1155
1156        let loaded: IndexMetadata = serde_json::from_value(json).unwrap();
1157        assert_eq!(loaded.segment_metas["legacy"].bp_unconverged_passes, 0);
1158    }
1159
1160    #[test]
1161    fn test_merged_segment_lineage() {
1162        let mut meta = IndexMetadata::new(test_schema());
1163        meta.add_segment("a".to_string(), 50);
1164        meta.add_segment("b".to_string(), 75);
1165
1166        // Fresh segments: gen=0, no ancestors
1167        assert_eq!(meta.segment_metas["a"].generation, 0);
1168        assert!(meta.segment_metas["a"].ancestors.is_empty());
1169
1170        // Merge a+b → c
1171        meta.add_merged_segment(
1172            "c".to_string(),
1173            125,
1174            vec!["a".to_string(), "b".to_string()],
1175            1,
1176            false,
1177            true,
1178        );
1179        assert_eq!(meta.segment_metas["c"].generation, 1);
1180        assert_eq!(meta.segment_metas["c"].ancestors, vec!["a", "b"]);
1181        assert_eq!(meta.segment_doc_count("c"), Some(125));
1182
1183        // Merge c+d → e (gen should be 2)
1184        meta.add_segment("d".to_string(), 30);
1185        meta.add_merged_segment(
1186            "e".to_string(),
1187            155,
1188            vec!["c".to_string(), "d".to_string()],
1189            2,
1190            false,
1191            true,
1192        );
1193        assert_eq!(meta.segment_metas["e"].generation, 2);
1194    }
1195}