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