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 centroid artifact paths
7//!
8//! The workflow is:
9//! 1. During initial accumulation, segments store flat vectors.
10//! 2. A manual build trains the first coarse-centroid ANN generation.
11//! 3. A manual retrain stages and atomically publishes a replacement generation.
12//! 4. On index open, metadata loads the currently published artifacts.
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    /// Legacy: path to a trained IVF-PQ codebook. Always `None` for current
120    /// formats; kept so pre-removal metadata deserializes into an actionable
121    /// error instead of dropping the field.
122    #[serde(skip_serializing_if = "Option::is_none")]
123    pub codebook_file: Option<String>,
124}
125
126/// Unified index metadata - single source of truth for index state
127#[derive(Debug, Clone, Serialize, Deserialize)]
128pub struct IndexMetadata {
129    /// Version for compatibility
130    pub version: u32,
131    /// Index schema
132    pub schema: Schema,
133    /// Segment metadata: segment_id -> info (doc count, etc.)
134    /// Using HashMap allows O(1) lookup and stores doc counts for merge decisions
135    #[serde(default)]
136    pub segment_metas: HashMap<String, SegmentMetaInfo>,
137    /// Per-field vector index metadata
138    #[serde(default)]
139    pub vector_fields: HashMap<u32, FieldVectorMeta>,
140    /// Aggregate vector count recorded by all built vector fields.
141    ///
142    /// The per-field `VectorIndexState::Built::vector_count` values are the
143    /// source of truth. This cached aggregate is refreshed whenever a field is
144    /// marked built, rather than being overwritten with whichever field was
145    /// trained last.
146    #[serde(default)]
147    pub total_vectors: usize,
148}
149
150impl IndexMetadata {
151    /// Create new metadata with schema
152    pub fn new(schema: Schema) -> Self {
153        Self {
154            version: INDEX_META_FORMAT_VERSION,
155            schema,
156            segment_metas: HashMap::new(),
157            vector_fields: HashMap::new(),
158            total_vectors: 0,
159        }
160    }
161
162    /// Get segment IDs as a sorted Vec (deterministic ordering)
163    pub fn segment_ids(&self) -> Vec<String> {
164        let mut ids: Vec<String> = self.segment_metas.keys().cloned().collect();
165        ids.sort();
166        ids
167    }
168
169    /// Add a fresh segment (gen=0, no ancestors, not reordered)
170    pub fn add_segment(&mut self, segment_id: String, num_docs: u32) {
171        self.segment_metas.insert(
172            segment_id,
173            SegmentMetaInfo {
174                num_docs,
175                ancestors: Vec::new(),
176                generation: 0,
177                reordered: false,
178                bp_converged: true,
179                bp_unconverged_passes: 0,
180            },
181        );
182    }
183
184    /// Add a merged segment with lineage info
185    pub fn add_merged_segment(
186        &mut self,
187        segment_id: String,
188        num_docs: u32,
189        ancestors: Vec<String>,
190        generation: u32,
191        reordered: bool,
192        bp_converged: bool,
193    ) {
194        self.add_segment_meta(
195            segment_id,
196            SegmentMetaInfo {
197                num_docs,
198                ancestors,
199                generation,
200                reordered,
201                bp_converged,
202                bp_unconverged_passes: 0,
203            },
204        );
205    }
206
207    /// Insert fully constructed lifecycle metadata. Merge/reorder code uses
208    /// this to carry bounded BP lineage; ordinary callers use the safer
209    /// constructors above, which start a fresh lineage.
210    pub(crate) fn add_segment_meta(&mut self, segment_id: String, info: SegmentMetaInfo) {
211        self.segment_metas.insert(segment_id, info);
212    }
213
214    /// Remove a segment
215    pub fn remove_segment(&mut self, segment_id: &str) {
216        self.segment_metas.remove(segment_id);
217    }
218
219    /// Check if segment exists
220    pub fn has_segment(&self, segment_id: &str) -> bool {
221        self.segment_metas.contains_key(segment_id)
222    }
223
224    /// Get segment doc count
225    pub fn segment_doc_count(&self, segment_id: &str) -> Option<u32> {
226        self.segment_metas.get(segment_id).map(|m| m.num_docs)
227    }
228
229    /// Check if a field has been built
230    pub fn is_field_built(&self, field_id: u32) -> bool {
231        self.vector_fields
232            .get(&field_id)
233            .map(|f| matches!(f.state, VectorIndexState::Built { .. }))
234            .unwrap_or(false)
235    }
236
237    /// Get field metadata
238    pub fn get_field_meta(&self, field_id: u32) -> Option<&FieldVectorMeta> {
239        self.vector_fields.get(&field_id)
240    }
241
242    /// Initialize field metadata (called when field is first seen)
243    pub fn init_field(&mut self, field_id: u32, index_type: impl Into<VectorFieldIndexType>) {
244        let index_type = index_type.into();
245        self.vector_fields
246            .entry(field_id)
247            .or_insert(FieldVectorMeta {
248                field_id,
249                index_type,
250                state: VectorIndexState::Flat,
251                centroids_file: None,
252                codebook_file: None,
253            });
254    }
255
256    /// Mark field as built with trained structures
257    pub fn mark_field_built(
258        &mut self,
259        field_id: u32,
260        vector_count: usize,
261        num_clusters: usize,
262        centroids_file: String,
263        codebook_file: Option<String>,
264    ) {
265        if let Some(field) = self.vector_fields.get_mut(&field_id) {
266            field.state = VectorIndexState::Built {
267                vector_count,
268                num_clusters,
269            };
270            field.centroids_file = Some(centroids_file);
271            field.codebook_file = codebook_file;
272            self.refresh_total_vectors();
273        }
274    }
275
276    /// Refresh the cached aggregate from the authoritative per-field states.
277    ///
278    /// Saturation keeps this infallible metadata helper safe even if it is
279    /// called after loading externally modified metadata with impossible
280    /// counts.
281    pub(crate) fn refresh_total_vectors(&mut self) {
282        self.total_vectors = self
283            .vector_fields
284            .values()
285            .filter_map(|field| match field.state {
286                VectorIndexState::Built { vector_count, .. } => Some(vector_count),
287                VectorIndexState::Flat => None,
288            })
289            .fold(0usize, usize::saturating_add);
290    }
291
292    /// Check if field should be built based on threshold
293    pub fn should_build_field(&self, field_id: u32, threshold: usize) -> bool {
294        // Don't build if already built
295        if self.is_field_built(field_id) {
296            return false;
297        }
298        // Build if we have enough vectors
299        self.total_vectors >= threshold
300    }
301
302    /// Load from directory
303    ///
304    /// If `metadata.json` is missing but `metadata.json.tmp` exists (crash
305    /// between write and rename), recovers from the temp file.
306    pub async fn load<D: crate::directories::Directory>(dir: &D) -> Result<Self> {
307        let path = Path::new(INDEX_META_FILENAME);
308        match dir.open_read(path).await {
309            Ok(slice) => {
310                let bytes = slice.read_bytes().await?;
311                Self::deserialize_versioned(bytes.as_slice())
312            }
313            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
314                // Try recovering from temp file (crash between write and rename)
315                let tmp_path = Path::new(INDEX_META_TMP_FILENAME);
316                let slice = dir.open_read(tmp_path).await?;
317                let bytes = slice.read_bytes().await?;
318                let meta = Self::deserialize_versioned(bytes.as_slice())?;
319                log::warn!("Recovered metadata from temp file (previous crash during save)");
320                Ok(meta)
321            }
322            Err(e) => Err(Error::Io(e)),
323        }
324    }
325
326    /// Deserialize only the current format. Vector artifacts are intentionally
327    /// rebuilt when the ANN format changes; silently accepting older metadata
328    /// would mix incompatible segment and global-codebook generations.
329    fn deserialize_versioned(bytes: &[u8]) -> Result<Self> {
330        let meta: Self =
331            serde_json::from_slice(bytes).map_err(|e| Error::Serialization(e.to_string()))?;
332        crate::dsl::reject_removed_vector_index_types(&meta.schema).map_err(Error::Schema)?;
333        if meta.version != INDEX_META_FORMAT_VERSION {
334            return Err(Error::Corruption(format!(
335                "metadata.json format version {} is incompatible with required version {}; \
336                 rebuild and republish the index with this Hermes version",
337                meta.version, INDEX_META_FORMAT_VERSION
338            )));
339        }
340        Ok(meta)
341    }
342
343    /// Save to directory (atomic: write temp file, then rename)
344    ///
345    /// Uses write-then-rename so a crash mid-write won't corrupt the
346    /// existing metadata file. On POSIX, rename is atomic.
347    pub async fn save<D: crate::directories::DirectoryWriter>(&self, dir: &D) -> Result<()> {
348        let bytes = self.serialize_to_bytes()?;
349        Self::save_bytes(dir, &bytes).await
350    }
351
352    /// Serialize metadata to bytes (cheap, no I/O).
353    /// Useful when you need to release a lock before doing disk I/O.
354    pub fn serialize_to_bytes(&self) -> Result<Vec<u8>> {
355        serde_json::to_vec_pretty(self).map_err(|e| Error::Serialization(e.to_string()))
356    }
357
358    /// Write pre-serialized metadata bytes to directory (atomic rename + fsync).
359    ///
360    /// The fsync ensures durability: without it, a power failure after rename
361    /// could lose the metadata update on systems with volatile write caches.
362    pub async fn save_bytes<D: crate::directories::DirectoryWriter>(
363        dir: &D,
364        bytes: &[u8],
365    ) -> Result<()> {
366        let tmp_path = Path::new(INDEX_META_TMP_FILENAME);
367        let final_path = Path::new(INDEX_META_FILENAME);
368        // Metadata is tiny, but `DirectoryWriter::write` does not guarantee
369        // the file contents themselves are fsynced. Finish the streaming
370        // writer first (filesystem implementations call `File::sync_all`),
371        // then atomically publish the durable temp file by rename.
372        let mut writer = dir.streaming_writer(tmp_path).await.map_err(Error::Io)?;
373        writer.write_all(bytes).map_err(Error::Io)?;
374        writer.finish().map_err(Error::Io)?;
375        // Rename is the logical commit point: after it succeeds, readers can
376        // observe the new generation and callers must publish the matching
377        // in-memory/tracker state. Directory fsync only strengthens crash
378        // durability. It cannot safely turn an already-visible rename into a
379        // reported pre-commit failure, because cleanup could then delete files
380        // referenced by the metadata now on disk.
381        dir.rename(tmp_path, final_path).await.map_err(Error::Io)?;
382        if let Err(error) = dir.sync().await {
383            log::error!(
384                "[metadata] directory fsync failed after committed rename: {}. \
385                 Continuing with the renamed generation; crash durability is not guaranteed",
386                error,
387            );
388        }
389        Ok(())
390    }
391
392    /// Fallible schema-aware loader used for lifecycle publication.
393    #[cfg_attr(not(feature = "native"), allow(dead_code))]
394    pub(crate) async fn try_load_trained_from_fields<D: crate::directories::Directory>(
395        vector_fields: &HashMap<u32, FieldVectorMeta>,
396        schema: &Schema,
397        dir: &D,
398    ) -> Result<Option<crate::segment::TrainedVectorStructures>> {
399        Self::load_trained_from_fields_impl(vector_fields, schema, dir).await
400    }
401
402    /// Load and validate the complete trained-artifact set described by a
403    /// `vector_fields` snapshot.
404    ///
405    /// This is intentionally all-or-nothing. A `Built` field is a durable
406    /// promise that every artifact required by its configured index exists and
407    /// is compatible with the schema. Returning a partial map would let some
408    /// segment builders publish ANN data while another field was silently
409    /// unusable, and would make the same index behave differently after a
410    /// restart.
411    async fn load_trained_from_fields_impl<D: crate::directories::Directory>(
412        vector_fields: &HashMap<u32, FieldVectorMeta>,
413        schema: &Schema,
414        dir: &D,
415    ) -> Result<Option<crate::segment::TrainedVectorStructures>> {
416        use std::sync::Arc;
417
418        let mut centroids = rustc_hash::FxHashMap::default();
419        let mut binary_quantizers = rustc_hash::FxHashMap::default();
420
421        let mut built_fields: Vec<_> = vector_fields
422            .iter()
423            .filter(|(_, meta)| matches!(meta.state, VectorIndexState::Built { .. }))
424            .collect();
425        built_fields.sort_unstable_by_key(|(field_id, _)| **field_id);
426
427        log::debug!(
428            "[trained] index={} loading trained structures, dense_vector_fields={:?}",
429            schema.index_label(),
430            vector_fields.keys().collect::<Vec<_>>()
431        );
432
433        for (field_id, field_meta) in built_fields {
434            log::debug!(
435                "[trained] index={} field {} state={:?} centroids_file={:?} codebook_file={:?}",
436                schema.index_label(),
437                field_id,
438                field_meta.state,
439                field_meta.centroids_file,
440                field_meta.codebook_file,
441            );
442            if field_meta.field_id != *field_id {
443                return Err(Error::Corruption(format!(
444                    "trained vector metadata key {field_id} contains field_id {}",
445                    field_meta.field_id
446                )));
447            }
448
449            let expected_clusters = match field_meta.state {
450                VectorIndexState::Built { num_clusters, .. } if num_clusters > 0 => num_clusters,
451                VectorIndexState::Built { .. } => {
452                    return Err(Error::Corruption(format!(
453                        "trained vector metadata field {field_id} has zero clusters"
454                    )));
455                }
456                VectorIndexState::Flat => unreachable!("built_fields contains only Built entries"),
457            };
458
459            let centroids_file = field_meta.centroids_file.as_deref().ok_or_else(|| {
460                Error::Corruption(format!(
461                    "trained vector metadata field {field_id} is Built but has no centroids_file"
462                ))
463            })?;
464            match field_meta.index_type {
465                VectorFieldIndexType::Float(VectorIndexType::IvfPq) => {
466                    return Err(Error::Corruption(format!(
467                        "field {field_id} was trained as IVF-PQ, which is no longer \
468                         supported; recreate the index with `ivf_tq` and reindex \
469                         (docs/turboquant-quantization.md)"
470                    )));
471                }
472                VectorFieldIndexType::Float(index_type @ VectorIndexType::IvfTq) => {
473                    let entry = schema
474                        .get_field_entry(crate::dsl::Field(*field_id))
475                        .ok_or_else(|| {
476                            Error::Corruption(format!(
477                                "trained vector metadata references missing field {field_id}"
478                            ))
479                        })?;
480                    let schema_config = entry
481                        .dense_vector_config
482                        .as_ref()
483                        .filter(|_| entry.field_type == crate::dsl::FieldType::DenseVector)
484                        .ok_or_else(|| {
485                            Error::Corruption(format!(
486                                "trained vector metadata field {field_id} is not a float dense field"
487                            ))
488                        })?;
489                    if schema_config.index_type != index_type {
490                        return Err(Error::Corruption(format!(
491                            "trained vector metadata field {field_id} uses {index_type:?}, schema requires {:?}",
492                            schema_config.index_type
493                        )));
494                    }
495                    let c: crate::structures::CoarseCentroids =
496                        load_trained_artifact(dir, *field_id, "centroids", centroids_file).await?;
497                    let expected_dim = schema_config.dim;
498                    let actual_clusters = c.num_clusters as usize;
499                    let expected_values =
500                        actual_clusters.checked_mul(expected_dim).ok_or_else(|| {
501                            Error::Corruption(format!(
502                                "trained centroid dimensions overflow for field {field_id}"
503                            ))
504                        })?;
505                    if actual_clusters == 0
506                        || actual_clusters > expected_clusters
507                        || c.dim == 0
508                        || c.dim != expected_dim
509                        || c.centroids.len() != expected_values
510                        || c.centroids.iter().any(|value| !value.is_finite())
511                    {
512                        return Err(Error::Corruption(format!(
513                            "trained centroids for field {field_id} do not match metadata/schema"
514                        )));
515                    }
516                    if !crate::structures::is_ivf_tq_cosine_generation(c.version) {
517                        return Err(Error::Corruption(format!(
518                            "trained IVF-TQ centroids for field {field_id} use an \
519                             unsupported legacy generation; rebuild the index"
520                        )));
521                    }
522                    c.validate_routing(schema_config.ivf_routing)
523                        .map_err(|error| {
524                            Error::Corruption(format!(
525                                "invalid trained centroid routing for field {field_id}: {error}"
526                            ))
527                        })?;
528                    // The TQ leaf codec is derived, never trained; ensure
529                    // `index_type` stays referenced for future variants.
530                    let _ = index_type;
531                    if field_meta.codebook_file.is_some() {
532                        return Err(Error::Corruption(format!(
533                            "trained IVF-TQ field {field_id} unexpectedly references a codebook file"
534                        )));
535                    }
536                    centroids.insert(*field_id, Arc::new(c));
537                }
538                VectorFieldIndexType::Binary(BinaryIndexType::Ivf) => {
539                    let entry = schema
540                        .get_field_entry(crate::dsl::Field(*field_id))
541                        .ok_or_else(|| {
542                            Error::Corruption(format!(
543                                "trained vector metadata references missing field {field_id}"
544                            ))
545                        })?;
546                    let schema_config = entry
547                        .binary_dense_vector_config
548                        .as_ref()
549                        .filter(|config| {
550                            entry.field_type == crate::dsl::FieldType::BinaryDenseVector
551                                && config.index_type == BinaryIndexType::Ivf
552                        })
553                        .ok_or_else(|| {
554                            Error::Corruption(format!(
555                                "trained vector metadata field {field_id} is not a binary IVF field"
556                            ))
557                        })?;
558                    let quantizer: crate::structures::BinaryCoarseQuantizer =
559                        load_trained_artifact(dir, *field_id, "binary centroids", centroids_file)
560                            .await?;
561                    quantizer.validate().map_err(|error| {
562                        Error::Corruption(format!(
563                            "invalid binary coarse quantizer for field {field_id}: {error}"
564                        ))
565                    })?;
566                    let actual_clusters = quantizer.num_clusters as usize;
567                    if actual_clusters > expected_clusters
568                        || schema_config.dim != quantizer.dim_bits
569                    {
570                        return Err(Error::Corruption(format!(
571                            "binary coarse quantizer for field {field_id} does not match metadata/schema"
572                        )));
573                    }
574                    quantizer
575                        .validate_routing(schema_config.ivf_routing)
576                        .map_err(|error| {
577                            Error::Corruption(format!(
578                                "invalid binary centroid routing for field {field_id}: {error}"
579                            ))
580                        })?;
581                    binary_quantizers.insert(*field_id, Arc::new(quantizer));
582                }
583                unsupported => {
584                    return Err(Error::Corruption(format!(
585                        "field {field_id} is Built for {unsupported:?}, which has no global IVF artifacts"
586                    )));
587                }
588            }
589        }
590
591        if centroids.is_empty() && binary_quantizers.is_empty() {
592            Ok(None)
593        } else {
594            let trained = crate::segment::TrainedVectorStructures {
595                #[cfg(feature = "native")]
596                _ann_pins: Default::default(),
597                centroids,
598                binary_quantizers,
599            };
600            #[cfg(feature = "native")]
601            let trained = {
602                let mut trained = trained;
603                trained.pin_ann_structures(crate::segment::pin::pin_policy());
604                trained
605            };
606            Ok(Some(trained))
607        }
608    }
609}
610
611fn validate_trained_artifact_path(field_id: u32, kind: &str, filename: &str) -> Result<()> {
612    use std::path::Component;
613
614    let path = Path::new(filename);
615    if filename.is_empty()
616        || path.is_absolute()
617        || path.components().any(|component| {
618            matches!(
619                component,
620                Component::ParentDir | Component::RootDir | Component::Prefix(_)
621            )
622        })
623    {
624        return Err(Error::Corruption(format!(
625            "trained {kind} path for field {field_id} is not a safe relative path: '{filename}'"
626        )));
627    }
628    Ok(())
629}
630
631async fn load_trained_artifact<T, D>(
632    dir: &D,
633    field_id: u32,
634    kind: &str,
635    filename: &str,
636) -> Result<T>
637where
638    T: serde::de::DeserializeOwned,
639    D: crate::directories::Directory,
640{
641    validate_trained_artifact_path(field_id, kind, filename)?;
642    let path = Path::new(filename);
643    let file_size = dir.file_size(path).await.map_err(|error| {
644        Error::Corruption(format!(
645            "failed to stat trained {kind} '{filename}' for field {field_id}: {error}"
646        ))
647    })?;
648    validate_trained_artifact_size(field_id, kind, filename, file_size)?;
649    let slice = dir.open_read(path).await.map_err(|error| {
650        Error::Corruption(format!(
651            "failed to open trained {kind} '{filename}' for field {field_id}: {error}"
652        ))
653    })?;
654    validate_trained_artifact_size(field_id, kind, filename, slice.len())?;
655    let bytes = slice.read_bytes().await.map_err(|error| {
656        Error::Corruption(format!(
657            "failed to read trained {kind} '{filename}' for field {field_id}: {error}"
658        ))
659    })?;
660    let (artifact, consumed) = bincode::serde::decode_from_slice::<T, _>(
661        bytes.as_slice(),
662        bincode::config::standard().with_limit::<MAX_TRAINED_ARTIFACT_BYTES>(),
663    )
664    .map_err(|error| {
665        Error::Corruption(format!(
666            "failed to deserialize trained {kind} '{filename}' for field {field_id}: {error}"
667        ))
668    })?;
669    if consumed != bytes.len() {
670        return Err(Error::Corruption(format!(
671            "trained {kind} '{filename}' for field {field_id} has {} trailing bytes",
672            bytes.len() - consumed
673        )));
674    }
675    Ok(artifact)
676}
677
678fn validate_trained_artifact_size(
679    field_id: u32,
680    kind: &str,
681    filename: &str,
682    file_size: u64,
683) -> Result<()> {
684    if file_size > MAX_TRAINED_ARTIFACT_BYTES as u64 {
685        return Err(Error::Corruption(format!(
686            "trained {kind} '{filename}' for field {field_id} is {file_size} bytes, \
687             exceeding the {MAX_TRAINED_ARTIFACT_BYTES}-byte safety limit"
688        )));
689    }
690    Ok(())
691}
692
693#[cfg(test)]
694mod tests {
695    use super::*;
696    use crate::directories::DirectoryWriter;
697
698    #[derive(Clone, Default)]
699    struct SyncFailDirectory(crate::directories::RamDirectory);
700
701    #[async_trait::async_trait]
702    impl crate::directories::Directory for SyncFailDirectory {
703        async fn exists(&self, path: &Path) -> std::io::Result<bool> {
704            self.0.exists(path).await
705        }
706
707        async fn file_size(&self, path: &Path) -> std::io::Result<u64> {
708            self.0.file_size(path).await
709        }
710
711        async fn open_read(&self, path: &Path) -> std::io::Result<crate::directories::FileHandle> {
712            self.0.open_read(path).await
713        }
714
715        async fn read_range(
716            &self,
717            path: &Path,
718            range: std::ops::Range<u64>,
719        ) -> std::io::Result<crate::directories::OwnedBytes> {
720            self.0.read_range(path, range).await
721        }
722
723        async fn list_files(&self, prefix: &Path) -> std::io::Result<Vec<std::path::PathBuf>> {
724            self.0.list_files(prefix).await
725        }
726
727        async fn open_lazy(&self, path: &Path) -> std::io::Result<crate::directories::FileHandle> {
728            self.0.open_lazy(path).await
729        }
730    }
731
732    #[async_trait::async_trait]
733    impl crate::directories::DirectoryWriter for SyncFailDirectory {
734        async fn write(&self, path: &Path, data: &[u8]) -> std::io::Result<()> {
735            self.0.write(path, data).await
736        }
737
738        async fn delete(&self, path: &Path) -> std::io::Result<()> {
739            self.0.delete(path).await
740        }
741
742        async fn rename(&self, from: &Path, to: &Path) -> std::io::Result<()> {
743            self.0.rename(from, to).await
744        }
745
746        async fn sync(&self) -> std::io::Result<()> {
747            Err(std::io::Error::other("injected directory fsync failure"))
748        }
749
750        async fn streaming_writer(
751            &self,
752            path: &Path,
753        ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
754            self.0.streaming_writer(path).await
755        }
756    }
757
758    fn test_schema() -> Schema {
759        Schema::default()
760    }
761
762    fn dense_schema(index_type: VectorIndexType) -> (Schema, crate::dsl::Field) {
763        let mut builder = crate::dsl::SchemaBuilder::default();
764        let config = match index_type {
765            VectorIndexType::IvfTq => crate::dsl::DenseVectorConfig::ivf_tq(2, Some(1), 1),
766            other => panic!("unsupported trained test index type: {other:?}"),
767        };
768        let field = builder.add_dense_vector_field_with_config("embedding", true, true, config);
769        (builder.build(), field)
770    }
771
772    fn test_centroids() -> crate::structures::CoarseCentroids {
773        crate::structures::CoarseCentroids {
774            num_clusters: 1,
775            dim: 2,
776            centroids: vec![0.25, 0.75],
777            version: crate::structures::mark_ivf_tq_cosine_generation(7),
778            soar_config: None,
779            routing_index: None,
780        }
781    }
782
783    async fn write_bincode(
784        directory: &crate::directories::RamDirectory,
785        filename: &str,
786        value: &impl serde::Serialize,
787    ) {
788        let bytes = bincode::serde::encode_to_vec(value, bincode::config::standard()).unwrap();
789        directory.write(Path::new(filename), &bytes).await.unwrap();
790    }
791
792    #[test]
793    fn test_metadata_init() {
794        let mut meta = IndexMetadata::new(test_schema());
795        assert_eq!(meta.total_vectors, 0);
796        assert!(meta.segment_metas.is_empty());
797        assert!(!meta.is_field_built(0));
798
799        meta.init_field(0, VectorIndexType::IvfTq);
800        assert!(!meta.is_field_built(0));
801        assert!(meta.vector_fields.contains_key(&0));
802    }
803
804    #[tokio::test]
805    async fn load_refuses_metadata_stamped_with_a_newer_format_version() {
806        let directory = crate::directories::RamDirectory::new();
807        let mut metadata = IndexMetadata::new(test_schema());
808        metadata.version = INDEX_META_FORMAT_VERSION + 1;
809        metadata.save(&directory).await.unwrap();
810
811        let error = IndexMetadata::load(&directory)
812            .await
813            .expect_err("metadata from a newer format version must be refused, not silently pruned")
814            .to_string();
815        assert!(error.contains("version 4"), "{error}");
816        assert!(error.contains("incompatible"), "{error}");
817    }
818
819    #[tokio::test]
820    async fn tmp_recovery_refuses_metadata_stamped_with_a_newer_format_version() {
821        let directory = crate::directories::RamDirectory::new();
822        let mut metadata = IndexMetadata::new(test_schema());
823        metadata.version = INDEX_META_FORMAT_VERSION + 1;
824        let bytes = metadata.serialize_to_bytes().unwrap();
825        // Simulate a crash between write and rename: only the temp file exists.
826        directory
827            .write(Path::new(INDEX_META_TMP_FILENAME), &bytes)
828            .await
829            .unwrap();
830
831        let error = IndexMetadata::load(&directory)
832            .await
833            .expect_err("temp-file recovery must apply the same version gate")
834            .to_string();
835        assert!(error.contains("version 4"), "{error}");
836    }
837
838    #[tokio::test]
839    async fn save_treats_post_rename_sync_failure_as_committed() {
840        let directory = SyncFailDirectory::default();
841        let mut metadata = IndexMetadata::new(test_schema());
842        metadata.add_segment("committed".to_string(), 7);
843
844        metadata.save(&directory).await.unwrap();
845
846        let loaded = IndexMetadata::load(&directory).await.unwrap();
847        assert_eq!(loaded.segment_doc_count("committed"), Some(7));
848    }
849
850    #[tokio::test]
851    async fn trained_artifacts_load_only_when_the_complete_built_set_is_valid() {
852        let mut builder = crate::dsl::SchemaBuilder::default();
853        let config = crate::dsl::DenseVectorConfig::ivf_tq(2, Some(1), 1);
854        let first = builder.add_dense_vector_field_with_config(
855            "first_embedding",
856            true,
857            true,
858            config.clone(),
859        );
860        let second =
861            builder.add_dense_vector_field_with_config("second_embedding", true, true, config);
862        let schema = builder.build();
863        let directory = crate::directories::RamDirectory::new();
864        let mut metadata = IndexMetadata::new(schema.clone());
865        metadata.init_field(first.0, VectorIndexType::IvfTq);
866        metadata.init_field(second.0, VectorIndexType::IvfTq);
867        metadata.mark_field_built(first.0, 10, 1, "field_0_centroids.bin".into(), None);
868        metadata.mark_field_built(second.0, 10, 1, "field_1_centroids.bin".into(), None);
869        write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
870
871        let error = IndexMetadata::try_load_trained_from_fields(
872            &metadata.vector_fields,
873            &schema,
874            &directory,
875        )
876        .await
877        .err()
878        .expect("missing artifact must fail the complete load")
879        .to_string();
880        assert!(error.contains("field_1_centroids.bin"), "{error}");
881        assert!(error.contains("field 1"), "{error}");
882    }
883
884    #[tokio::test]
885    async fn index_open_fails_closed_when_built_artifact_is_missing() {
886        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
887        let directory = crate::directories::RamDirectory::new();
888        let mut metadata = IndexMetadata::new(schema);
889        metadata.init_field(field.0, VectorIndexType::IvfTq);
890        metadata.mark_field_built(field.0, 10, 1, "missing_centroids.bin".into(), None);
891        metadata.save(&directory).await.unwrap();
892
893        let error = match crate::index::Index::open(directory, crate::index::IndexConfig::default())
894            .await
895        {
896            Ok(_) => panic!("Index::open accepted a Built field with no artifact"),
897            Err(error) => error.to_string(),
898        };
899        assert!(error.contains("missing_centroids.bin"), "{error}");
900    }
901
902    #[tokio::test]
903    async fn ivf_tq_built_state_rejects_a_codebook_file() {
904        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
905        let directory = crate::directories::RamDirectory::new();
906        let mut metadata = IndexMetadata::new(schema.clone());
907        metadata.init_field(field.0, VectorIndexType::IvfTq);
908        metadata.mark_field_built(
909            field.0,
910            10,
911            1,
912            "field_0_centroids.bin".into(),
913            Some("field_0_codebook.bin".into()),
914        );
915        write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
916
917        let error = IndexMetadata::try_load_trained_from_fields(
918            &metadata.vector_fields,
919            &schema,
920            &directory,
921        )
922        .await
923        .err()
924        .expect("IVF-TQ Built state with a codebook file must fail")
925        .to_string();
926        assert!(error.contains("codebook"), "{error}");
927    }
928
929    #[tokio::test]
930    async fn legacy_ivf_tq_centroid_generation_is_rejected_while_loading() {
931        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
932        let directory = crate::directories::RamDirectory::new();
933        let mut metadata = IndexMetadata::new(schema.clone());
934        metadata.init_field(field.0, VectorIndexType::IvfTq);
935        metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
936        let mut legacy = test_centroids();
937        legacy.version = 7;
938        write_bincode(&directory, "field_0_centroids.bin", &legacy).await;
939
940        let error = IndexMetadata::try_load_trained_from_fields(
941            &metadata.vector_fields,
942            &schema,
943            &directory,
944        )
945        .await
946        .err()
947        .expect("legacy IVF-TQ centroid state must fail while loading")
948        .to_string();
949        assert!(error.contains("unsupported legacy generation"), "{error}");
950        assert!(error.contains("rebuild the index"), "{error}");
951    }
952
953    #[tokio::test]
954    async fn legacy_ivf_pq_trained_field_fails_with_actionable_error() {
955        // Simulates metadata written by a pre-removal version: the schema
956        // gate rejects `ivf_pq` fields, so build the raw field-state map
957        // directly against a current-format schema.
958        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
959        let directory = crate::directories::RamDirectory::new();
960        let mut metadata = IndexMetadata::new(schema.clone());
961        metadata.init_field(field.0, VectorIndexType::IvfTq);
962        metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
963        // Overwrite the recorded type the way pre-removal metadata carries it
964        // (init_field never downgrades an existing entry).
965        metadata
966            .vector_fields
967            .get_mut(&field.0)
968            .expect("field initialized")
969            .index_type = VectorFieldIndexType::Float(VectorIndexType::IvfPq);
970        write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
971
972        let error = IndexMetadata::try_load_trained_from_fields(
973            &metadata.vector_fields,
974            &schema,
975            &directory,
976        )
977        .await
978        .err()
979        .expect("legacy IVF-PQ trained state must fail loudly")
980        .to_string();
981        assert!(error.contains("no longer"), "{error}");
982        assert!(error.contains("ivf_tq"), "{error}");
983    }
984
985    #[tokio::test]
986    async fn requested_cluster_count_accepts_a_quality_clamped_artifact() {
987        let mut builder = crate::dsl::SchemaBuilder::default();
988        let field = builder.add_dense_vector_field_with_config(
989            "embedding",
990            true,
991            true,
992            crate::dsl::DenseVectorConfig::ivf_tq(2, Some(4), 1),
993        );
994        let schema = builder.build();
995        let directory = crate::directories::RamDirectory::new();
996        let mut metadata = IndexMetadata::new(schema.clone());
997        metadata.init_field(field.0, VectorIndexType::IvfTq);
998        metadata.mark_field_built(field.0, 1, 4, "field_0_centroids.bin".into(), None);
999        write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
1000
1001        let trained = IndexMetadata::try_load_trained_from_fields(
1002            &metadata.vector_fields,
1003            &schema,
1004            &directory,
1005        )
1006        .await
1007        .unwrap()
1008        .unwrap();
1009        assert_eq!(trained.centroids[&field.0].num_clusters, 1);
1010    }
1011
1012    #[tokio::test]
1013    async fn trained_artifact_loader_rejects_trailing_data() {
1014        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1015        let directory = crate::directories::RamDirectory::new();
1016        let mut metadata = IndexMetadata::new(schema.clone());
1017        metadata.init_field(field.0, VectorIndexType::IvfTq);
1018        metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1019        let mut bytes =
1020            bincode::serde::encode_to_vec(test_centroids(), bincode::config::standard()).unwrap();
1021        bytes.extend_from_slice(&[0xaa, 0xbb]);
1022        directory
1023            .write(Path::new("field_0_centroids.bin"), &bytes)
1024            .await
1025            .unwrap();
1026
1027        let error = IndexMetadata::try_load_trained_from_fields(
1028            &metadata.vector_fields,
1029            &schema,
1030            &directory,
1031        )
1032        .await
1033        .err()
1034        .expect("trailing artifact bytes must fail validation")
1035        .to_string();
1036        assert!(error.contains("trailing bytes"), "{error}");
1037    }
1038
1039    #[test]
1040    fn trained_artifact_size_limit_rejects_before_reading() {
1041        let error = validate_trained_artifact_size(
1042            3,
1043            "centroids",
1044            "field_3_centroids.bin",
1045            MAX_TRAINED_ARTIFACT_BYTES as u64 + 1,
1046        )
1047        .unwrap_err()
1048        .to_string();
1049        assert!(error.contains("exceeding"), "{error}");
1050        assert!(error.contains("field 3"), "{error}");
1051    }
1052
1053    #[tokio::test]
1054    async fn trained_artifact_decode_limit_rejects_forged_collection_length() {
1055        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1056        let directory = crate::directories::RamDirectory::new();
1057        let mut metadata = IndexMetadata::new(schema.clone());
1058        metadata.init_field(field.0, VectorIndexType::IvfTq);
1059        metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1060
1061        // CoarseCentroids begins with num_clusters=1, dim=2, then the Vec
1062        // length. Bincode's standard varint marker 253 introduces a u64; this
1063        // tiny payload claims an impossible f32 vector and must hit the decode
1064        // limit before any large allocation is attempted.
1065        let mut bytes = vec![1, 2, 253];
1066        bytes.extend_from_slice(&u64::MAX.to_le_bytes());
1067        directory
1068            .write(Path::new("field_0_centroids.bin"), &bytes)
1069            .await
1070            .unwrap();
1071
1072        let error = IndexMetadata::try_load_trained_from_fields(
1073            &metadata.vector_fields,
1074            &schema,
1075            &directory,
1076        )
1077        .await
1078        .err()
1079        .expect("forged collection length must fail the bounded decoder")
1080        .to_string();
1081        assert!(error.contains("failed to deserialize"), "{error}");
1082    }
1083
1084    #[test]
1085    fn test_metadata_segments() {
1086        let mut meta = IndexMetadata::new(test_schema());
1087        meta.add_segment("abc123".to_string(), 50);
1088        meta.add_segment("def456".to_string(), 100);
1089        assert_eq!(meta.segment_metas.len(), 2);
1090        assert_eq!(meta.segment_doc_count("abc123"), Some(50));
1091        assert_eq!(meta.segment_doc_count("def456"), Some(100));
1092
1093        // Overwrites existing
1094        meta.add_segment("abc123".to_string(), 75);
1095        assert_eq!(meta.segment_metas.len(), 2);
1096        assert_eq!(meta.segment_doc_count("abc123"), Some(75));
1097
1098        meta.remove_segment("abc123");
1099        assert_eq!(meta.segment_metas.len(), 1);
1100        assert!(meta.has_segment("def456"));
1101        assert!(!meta.has_segment("abc123"));
1102    }
1103
1104    #[test]
1105    fn test_mark_field_built() {
1106        let mut meta = IndexMetadata::new(test_schema());
1107        meta.init_field(0, VectorIndexType::IvfTq);
1108        meta.total_vectors = 10000;
1109
1110        assert!(!meta.is_field_built(0));
1111
1112        meta.mark_field_built(0, 10000, 256, "field_0_centroids.bin".to_string(), None);
1113
1114        assert!(meta.is_field_built(0));
1115        let field = meta.get_field_meta(0).unwrap();
1116        assert_eq!(
1117            field.centroids_file.as_deref(),
1118            Some("field_0_centroids.bin")
1119        );
1120    }
1121
1122    #[test]
1123    fn total_vectors_is_aggregate_of_built_field_counts() {
1124        let mut meta = IndexMetadata::new(test_schema());
1125        meta.init_field(7, VectorIndexType::IvfTq);
1126        meta.init_field(3, VectorIndexType::IvfTq);
1127
1128        // Build in reverse field-id order to ensure the result is not tied to
1129        // HashMap or training iteration order.
1130        meta.mark_field_built(7, 400, 20, "field_7_centroids.bin".to_string(), None);
1131        assert_eq!(meta.total_vectors, 400);
1132        meta.mark_field_built(3, 250, 15, "field_3_centroids.bin".to_string(), None);
1133        assert_eq!(meta.total_vectors, 650);
1134
1135        // Rebuilding a field replaces its contribution; it does not add a
1136        // duplicate training snapshot.
1137        meta.mark_field_built(7, 425, 20, "field_7_centroids.bin".to_string(), None);
1138        assert_eq!(meta.total_vectors, 675);
1139    }
1140
1141    #[test]
1142    fn test_should_build_field() {
1143        let mut meta = IndexMetadata::new(test_schema());
1144        meta.init_field(0, VectorIndexType::IvfTq);
1145
1146        // Below threshold
1147        meta.total_vectors = 500;
1148        assert!(!meta.should_build_field(0, 1000));
1149
1150        // Above threshold
1151        meta.total_vectors = 1500;
1152        assert!(meta.should_build_field(0, 1000));
1153
1154        // Already built - should not build again
1155        meta.mark_field_built(0, 1500, 256, "centroids.bin".to_string(), None);
1156        assert!(!meta.should_build_field(0, 1000));
1157    }
1158
1159    #[test]
1160    fn test_serialization() {
1161        let mut meta = IndexMetadata::new(test_schema());
1162        meta.add_segment("seg1".to_string(), 100);
1163        meta.init_field(0, VectorIndexType::IvfTq);
1164        meta.total_vectors = 5000;
1165
1166        let json = serde_json::to_string_pretty(&meta).unwrap();
1167        let loaded: IndexMetadata = serde_json::from_str(&json).unwrap();
1168
1169        assert_eq!(loaded.segment_ids().len(), meta.segment_ids().len());
1170        assert_eq!(loaded.segment_doc_count("seg1"), Some(100));
1171        assert_eq!(loaded.total_vectors, meta.total_vectors);
1172        assert!(loaded.vector_fields.contains_key(&0));
1173    }
1174
1175    #[test]
1176    fn old_metadata_defaults_the_bp_retry_counter() {
1177        let mut meta = IndexMetadata::new(test_schema());
1178        meta.add_segment("legacy".to_string(), 10);
1179        let mut json = serde_json::to_value(&meta).unwrap();
1180        json["segment_metas"]["legacy"]
1181            .as_object_mut()
1182            .unwrap()
1183            .remove("bp_unconverged_passes");
1184
1185        let loaded: IndexMetadata = serde_json::from_value(json).unwrap();
1186        assert_eq!(loaded.segment_metas["legacy"].bp_unconverged_passes, 0);
1187    }
1188
1189    #[test]
1190    fn test_merged_segment_lineage() {
1191        let mut meta = IndexMetadata::new(test_schema());
1192        meta.add_segment("a".to_string(), 50);
1193        meta.add_segment("b".to_string(), 75);
1194
1195        // Fresh segments: gen=0, no ancestors
1196        assert_eq!(meta.segment_metas["a"].generation, 0);
1197        assert!(meta.segment_metas["a"].ancestors.is_empty());
1198
1199        // Merge a+b → c
1200        meta.add_merged_segment(
1201            "c".to_string(),
1202            125,
1203            vec!["a".to_string(), "b".to_string()],
1204            1,
1205            false,
1206            true,
1207        );
1208        assert_eq!(meta.segment_metas["c"].generation, 1);
1209        assert_eq!(meta.segment_metas["c"].ancestors, vec!["a", "b"]);
1210        assert_eq!(meta.segment_doc_count("c"), Some(125));
1211
1212        // Merge c+d → e (gen should be 2)
1213        meta.add_segment("d".to_string(), 30);
1214        meta.add_merged_segment(
1215            "e".to_string(),
1216            155,
1217            vec!["c".to_string(), "d".to_string()],
1218            2,
1219            false,
1220            true,
1221        );
1222        assert_eq!(meta.segment_metas["e"].generation, 2);
1223    }
1224}