Skip to main content

summa_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` requires this version or one it can upgrade in memory (see
30/// [`OLDEST_MIGRATABLE_FORMAT_VERSION`]). Anything else is a clean rebuild
31/// boundary; serde_json would otherwise silently drop fields it does not know
32/// and a later save could destructively rewrite index state.
33pub const INDEX_META_FORMAT_VERSION: u32 = 9;
34
35/// Oldest metadata.json format `load` upgrades in place.
36///
37/// Format 9 adds optional SIMD blocks, compact posting/position directories
38/// and byte norms; older encodings retain their scoring semantics. Format 8
39/// protects the optional content-hash field marker from older writers. Format 7 added the optional per-segment `deletions` entry,
40/// so format 6 metadata (1.8.121..=1.8.133) is also readable. The
41/// upgrade is loud and one-way: the writer persists the new stamp on open so
42/// older builds cannot reopen the index and misread position codec tags or
43/// drop deletion generations. No encoded blocks are rewritten by migration.
44/// Segment-level breaks inside the format 6 window (BMP blob
45/// magic BMP9 -> BMPA in 1.8.125) are still refused at segment open.
46pub const OLDEST_MIGRATABLE_FORMAT_VERSION: u32 = 6;
47
48/// Index-level centroids/codebooks are deliberately bounded before they are
49/// read or decoded. Besides limiting ordinary corruption damage, the matching
50/// bincode limit prevents a tiny forged collection length from requesting an
51/// effectively unbounded allocation.
52pub(crate) const MAX_TRAINED_ARTIFACT_BYTES: usize = 512 * 1024 * 1024;
53
54/// State of vector index for a field
55#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
56pub enum VectorIndexState {
57    /// Accumulating vectors - using Flat (brute-force) search
58    #[default]
59    Flat,
60    /// Index structures built - using ANN search
61    Built {
62        /// Total vector count when training happened
63        vector_count: usize,
64        /// Number of clusters used
65        num_clusters: usize,
66    },
67}
68
69fn is_zero(value: &u32) -> bool {
70    *value == 0
71}
72
73fn default_true() -> bool {
74    true
75}
76
77/// Per-segment metadata stored in index metadata
78/// This allows merge decisions without loading segment files
79#[derive(Debug, Clone, Serialize, Deserialize)]
80pub struct SegmentMetaInfo {
81    /// Exact immutable row visibility generation; absence means all rows live.
82    #[serde(default, skip_serializing_if = "Option::is_none")]
83    pub deletions: Option<crate::segment::DeletionMeta>,
84    /// Number of documents in this segment
85    pub num_docs: u32,
86    /// Parent segment IDs that were merged to produce this segment (empty for fresh segments)
87    pub ancestors: Vec<String>,
88    /// Merge generation: 0 for fresh segments, max(parent generations) + 1 for merged segments
89    pub generation: u32,
90    /// Whether this segment has been reordered via Recursive Graph Bisection (BP).
91    /// Fresh segments and block-copy merges are not reordered. Only segments that have
92    /// been explicitly reordered (via background optimizer or reorder command) are marked true.
93    #[serde(default)]
94    pub reordered: bool,
95    /// Whether the last BP reorder pass ran to natural convergence. False when
96    /// a wall-clock BP budget ended the pass early — the segment is ordered
97    /// better than before, and a later warm-started pass can deepen it.
98    /// Old metadata (field absent) deserializes as converged.
99    #[serde(default = "default_true")]
100    pub bp_converged: bool,
101    /// Number of consecutive budget-exhausted BP rewrites in this segment's
102    /// current reordered lineage. Carried across replacement IDs so the
103    /// optimizer can impose a hard follow-up bound instead of rewriting forever.
104    #[serde(default)]
105    pub bp_unconverged_passes: u32,
106    /// Terms whose copied Seismic nominations still span multiple runs.
107    #[serde(default)]
108    pub seismic_pending_terms: u32,
109    /// Published partial-maintenance passes in this replacement lineage. Used
110    /// to distinguish initial work from cooldown-paced follow-ups, not as a cap.
111    #[serde(default)]
112    pub seismic_maintenance_passes: u32,
113    /// Consecutive published maintenance passes that did not reduce Seismic
114    /// nomination debt. Progress resets this count; the optimizer bounds stalls.
115    #[serde(default, skip_serializing_if = "is_zero")]
116    pub seismic_no_progress_passes: u32,
117    /// Binary ANN leaves with multiple copied runs require lossless coalescing.
118    #[serde(default)]
119    pub ann_fragmented: bool,
120}
121
122impl SegmentMetaInfo {
123    pub fn num_deleted_docs(&self) -> u32 {
124        self.deletions.as_ref().map_or(0, |meta| meta.num_deleted)
125    }
126
127    pub fn num_live_docs(&self) -> u32 {
128        self.num_docs - self.num_deleted_docs()
129    }
130
131    pub fn deleted_ratio(&self) -> f64 {
132        if self.num_docs == 0 {
133            0.0
134        } else {
135            f64::from(self.num_deleted_docs()) / f64::from(self.num_docs)
136        }
137    }
138}
139
140/// Per-field vector index metadata
141#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
142#[serde(tag = "kind", content = "index", rename_all = "snake_case")]
143pub enum VectorFieldIndexType {
144    Float(VectorIndexType),
145    Binary(BinaryIndexType),
146}
147
148impl From<VectorIndexType> for VectorFieldIndexType {
149    fn from(value: VectorIndexType) -> Self {
150        Self::Float(value)
151    }
152}
153
154impl From<BinaryIndexType> for VectorFieldIndexType {
155    fn from(value: BinaryIndexType) -> Self {
156        Self::Binary(value)
157    }
158}
159
160#[derive(Debug, Clone, Serialize, Deserialize)]
161pub struct FieldVectorMeta {
162    /// Field ID
163    pub field_id: u32,
164    /// Configured index type (target type when built)
165    pub index_type: VectorFieldIndexType,
166    /// Current state
167    pub state: VectorIndexState,
168    /// Path to centroids file (relative to index dir)
169    #[serde(skip_serializing_if = "Option::is_none")]
170    pub centroids_file: Option<String>,
171    /// Legacy: path to a trained IVF-PQ codebook. Always `None` for current
172    /// formats; kept so pre-removal metadata deserializes into an actionable
173    /// error instead of dropping the field.
174    #[serde(skip_serializing_if = "Option::is_none")]
175    pub codebook_file: Option<String>,
176    /// ScaNN global model generation encoded into every segment payload.
177    /// Absent for legacy IVF/TQ fields and older metadata.
178    #[serde(default, skip_serializing_if = "Option::is_none")]
179    pub artifact_generation: Option<u64>,
180    /// Content fingerprint of the exact ScaNN model encoded into every
181    /// segment payload. This prevents same-generation accidental mixing.
182    #[serde(default, skip_serializing_if = "Option::is_none")]
183    pub artifact_id: Option<u64>,
184}
185
186/// Unified index metadata - single source of truth for index state
187#[derive(Debug, Clone, Serialize, Deserialize)]
188pub struct IndexMetadata {
189    /// Version for compatibility
190    pub version: u32,
191    /// Monotonic publication identity for schema/vector generations. Segment
192    /// commits are already detected by their IDs; this also makes query-only
193    /// ALTERs (for example `nprobe`) visible to cached readers.
194    #[serde(default)]
195    pub publication_generation: u64,
196    /// Index schema
197    pub schema: Schema,
198    /// Segment metadata: segment_id -> info (doc count, etc.)
199    /// Using HashMap allows O(1) lookup and stores doc counts for merge decisions
200    #[serde(default)]
201    pub segment_metas: HashMap<String, SegmentMetaInfo>,
202    /// Per-field vector index metadata
203    #[serde(default)]
204    pub vector_fields: HashMap<u32, FieldVectorMeta>,
205    /// Aggregate vector count recorded by all built vector fields.
206    ///
207    /// The per-field `VectorIndexState::Built::vector_count` values are the
208    /// source of truth. This cached aggregate is refreshed whenever a field is
209    /// marked built, rather than being overwritten with whichever field was
210    /// trained last.
211    #[serde(default)]
212    pub total_vectors: usize,
213}
214
215impl IndexMetadata {
216    /// Every file identity protected by this metadata generation.
217    #[cfg(feature = "native")]
218    pub(crate) fn owned_ids(&self) -> Vec<String> {
219        self.segment_metas
220            .keys()
221            .cloned()
222            .chain(
223                self.segment_metas
224                    .values()
225                    .filter_map(|info| info.deletions.as_ref().map(|d| d.id.clone())),
226            )
227            .collect()
228    }
229
230    #[cfg(feature = "native")]
231    pub(crate) fn owns_id(&self, id: &str) -> bool {
232        self.has_segment(id)
233            || self
234                .segment_metas
235                .values()
236                .any(|info| info.deletions.as_ref().is_some_and(|d| d.id == id))
237    }
238    /// Create new metadata with schema
239    pub fn new(schema: Schema) -> Self {
240        Self {
241            version: INDEX_META_FORMAT_VERSION,
242            publication_generation: 0,
243            schema,
244            segment_metas: HashMap::new(),
245            vector_fields: HashMap::new(),
246            total_vectors: 0,
247        }
248    }
249
250    /// Get segment IDs as a sorted Vec (deterministic ordering)
251    pub fn segment_ids(&self) -> Vec<String> {
252        let mut ids: Vec<String> = self.segment_metas.keys().cloned().collect();
253        ids.sort();
254        ids
255    }
256
257    /// Add a fresh segment (gen=0, no ancestors, not reordered)
258    pub fn add_segment(&mut self, segment_id: String, num_docs: u32) {
259        self.segment_metas.insert(
260            segment_id,
261            SegmentMetaInfo {
262                deletions: None,
263                num_docs,
264                ancestors: Vec::new(),
265                generation: 0,
266                reordered: false,
267                bp_converged: true,
268                bp_unconverged_passes: 0,
269                seismic_pending_terms: 0,
270                seismic_maintenance_passes: 0,
271                seismic_no_progress_passes: 0,
272                ann_fragmented: false,
273            },
274        );
275    }
276
277    /// Add a merged segment with lineage info
278    pub fn add_merged_segment(
279        &mut self,
280        segment_id: String,
281        num_docs: u32,
282        ancestors: Vec<String>,
283        generation: u32,
284        reordered: bool,
285        bp_converged: bool,
286    ) {
287        self.add_segment_meta(
288            segment_id,
289            SegmentMetaInfo {
290                deletions: None,
291                num_docs,
292                ancestors,
293                generation,
294                reordered,
295                bp_converged,
296                bp_unconverged_passes: 0,
297                seismic_pending_terms: 0,
298                seismic_maintenance_passes: 0,
299                seismic_no_progress_passes: 0,
300                ann_fragmented: false,
301            },
302        );
303    }
304
305    /// Insert fully constructed lifecycle metadata. Merge/reorder code uses
306    /// this to carry bounded BP lineage; ordinary callers use the safer
307    /// constructors above, which start a fresh lineage.
308    pub(crate) fn add_segment_meta(&mut self, segment_id: String, info: SegmentMetaInfo) {
309        self.segment_metas.insert(segment_id, info);
310    }
311
312    /// Remove a segment
313    pub fn remove_segment(&mut self, segment_id: &str) {
314        self.segment_metas.remove(segment_id);
315    }
316
317    /// Check if segment exists
318    pub fn has_segment(&self, segment_id: &str) -> bool {
319        self.segment_metas.contains_key(segment_id)
320    }
321
322    /// Get segment doc count
323    pub fn segment_doc_count(&self, segment_id: &str) -> Option<u32> {
324        self.segment_metas.get(segment_id).map(|m| m.num_docs)
325    }
326
327    /// Check if a field has been built
328    pub fn is_field_built(&self, field_id: u32) -> bool {
329        self.vector_fields
330            .get(&field_id)
331            .map(|f| matches!(f.state, VectorIndexState::Built { .. }))
332            .unwrap_or(false)
333    }
334
335    /// Get field metadata
336    pub fn get_field_meta(&self, field_id: u32) -> Option<&FieldVectorMeta> {
337        self.vector_fields.get(&field_id)
338    }
339
340    /// Initialize field metadata (called when field is first seen)
341    pub fn init_field(&mut self, field_id: u32, index_type: impl Into<VectorFieldIndexType>) {
342        let index_type = index_type.into();
343        self.vector_fields
344            .entry(field_id)
345            .or_insert(FieldVectorMeta {
346                field_id,
347                index_type,
348                state: VectorIndexState::Flat,
349                centroids_file: None,
350                codebook_file: None,
351                artifact_generation: None,
352                artifact_id: None,
353            });
354    }
355
356    /// Mark field as built with trained structures
357    pub fn mark_field_built(
358        &mut self,
359        field_id: u32,
360        vector_count: usize,
361        num_clusters: usize,
362        centroids_file: String,
363        codebook_file: Option<String>,
364    ) {
365        if let Some(field) = self.vector_fields.get_mut(&field_id) {
366            field.state = VectorIndexState::Built {
367                vector_count,
368                num_clusters,
369            };
370            field.centroids_file = Some(centroids_file);
371            field.codebook_file = codebook_file;
372            field.artifact_generation = None;
373            field.artifact_id = None;
374            self.refresh_total_vectors();
375        }
376    }
377
378    /// Mark a ScaNN field built against one immutable global model. Reuses
379    /// `centroids_file` as the trained-artifact path for wire compatibility;
380    /// the explicit generation/fingerprint fields make its semantics loud.
381    pub fn mark_scann_field_built(
382        &mut self,
383        field_id: u32,
384        vector_count: usize,
385        num_leaves: usize,
386        artifact_file: String,
387        artifact_generation: u64,
388        artifact_id: u64,
389    ) -> Result<()> {
390        if artifact_generation == 0 || artifact_id == 0 {
391            return Err(Error::Corruption(format!(
392                "ScaNN field {field_id} cannot publish a zero generation or artifact fingerprint"
393            )));
394        }
395        let field = self.vector_fields.get_mut(&field_id).ok_or_else(|| {
396            Error::Corruption(format!(
397                "ScaNN field {field_id} must be initialized before it is marked built"
398            ))
399        })?;
400        if !matches!(
401            field.index_type,
402            VectorFieldIndexType::Float(VectorIndexType::Scann)
403                | VectorFieldIndexType::Binary(BinaryIndexType::Scann)
404        ) {
405            return Err(Error::Corruption(format!(
406                "field {field_id} is not configured as ScaNN"
407            )));
408        }
409        field.state = VectorIndexState::Built {
410            vector_count,
411            num_clusters: num_leaves,
412        };
413        field.centroids_file = Some(artifact_file);
414        field.codebook_file = None;
415        field.artifact_generation = Some(artifact_generation);
416        field.artifact_id = Some(artifact_id);
417        self.refresh_total_vectors();
418        Ok(())
419    }
420
421    /// Refresh the cached aggregate from the authoritative per-field states.
422    ///
423    /// Saturation keeps this infallible metadata helper safe even if it is
424    /// called after loading externally modified metadata with impossible
425    /// counts.
426    pub(crate) fn refresh_total_vectors(&mut self) {
427        self.total_vectors = self
428            .vector_fields
429            .values()
430            .filter_map(|field| match field.state {
431                VectorIndexState::Built { vector_count, .. } => Some(vector_count),
432                VectorIndexState::Flat => None,
433            })
434            .fold(0usize, usize::saturating_add);
435    }
436
437    /// Check if field should be built based on threshold
438    pub fn should_build_field(&self, field_id: u32, threshold: usize) -> bool {
439        // Don't build if already built
440        if self.is_field_built(field_id) {
441            return false;
442        }
443        // Build if we have enough vectors
444        self.total_vectors >= threshold
445    }
446
447    /// Load from directory
448    ///
449    /// If `metadata.json` is missing but `metadata.json.tmp` exists (crash
450    /// between write and rename), recovers from the temp file.
451    ///
452    /// Metadata stamped with a migratable older format is upgraded in memory
453    /// and reported with `log::warn!`; nothing is written here. Writers call
454    /// [`load_reporting_migration`](Self::load_reporting_migration) and persist
455    /// the upgrade themselves.
456    pub async fn load<D: crate::directories::Directory>(dir: &D) -> Result<Self> {
457        let (meta, migrated_from) = Self::load_reporting_migration(dir).await?;
458        if let Some(from) = migrated_from {
459            log::warn!(
460                "[metadata_migration] metadata.json format version {from} upgraded in memory \
461                 to {INDEX_META_FORMAT_VERSION} ({} segment(s)); the upgrade is persisted \
462                 when a writer opens this index. Segments from builds before 1.8.125 fail at \
463                 segment open and need a rebuild; segments without .rowstats cannot be \
464                 compacted until they are merged",
465                meta.segment_metas.len()
466            );
467        }
468        Ok(meta)
469    }
470
471    /// [`load`](Self::load) that also returns the on-disk format version when
472    /// it differed from `INDEX_META_FORMAT_VERSION` and was upgraded.
473    ///
474    /// The caller owns the log line and any persistence, so a writer can
475    /// report "persisted" instead of the reader's "in memory only" warning.
476    pub async fn load_reporting_migration<D: crate::directories::Directory>(
477        dir: &D,
478    ) -> Result<(Self, Option<u32>)> {
479        let path = Path::new(INDEX_META_FILENAME);
480        match dir.open_read(path).await {
481            Ok(slice) => {
482                let bytes = slice.read_bytes().await?;
483                Self::deserialize_versioned(bytes.as_slice())
484            }
485            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
486                // Try recovering from temp file (crash between write and rename)
487                let tmp_path = Path::new(INDEX_META_TMP_FILENAME);
488                let slice = dir.open_read(tmp_path).await?;
489                let bytes = slice.read_bytes().await?;
490                let recovered = Self::deserialize_versioned(bytes.as_slice())?;
491                log::warn!("Recovered metadata from temp file (previous crash during save)");
492                Ok(recovered)
493            }
494            Err(e) => Err(Error::Io(e)),
495        }
496    }
497
498    /// Deserialize the current format or upgrade a migratable older one,
499    /// returning the original stamp when an upgrade happened. Vector artifacts
500    /// are intentionally rebuilt when the ANN format changes; accepting
501    /// arbitrary older metadata would mix incompatible segment and
502    /// global-codebook generations.
503    fn deserialize_versioned(bytes: &[u8]) -> Result<(Self, Option<u32>)> {
504        let mut meta: Self =
505            serde_json::from_slice(bytes).map_err(|e| Error::Serialization(e.to_string()))?;
506        meta.schema.validate()?;
507        crate::dsl::reject_removed_vector_index_types(&meta.schema).map_err(Error::Schema)?;
508        let migrated_from = if meta.version == INDEX_META_FORMAT_VERSION {
509            None
510        } else if (OLDEST_MIGRATABLE_FORMAT_VERSION..INDEX_META_FORMAT_VERSION)
511            .contains(&meta.version)
512        {
513            let from = meta.version;
514            meta.version = INDEX_META_FORMAT_VERSION;
515            Some(from)
516        } else {
517            return Err(Error::Corruption(format!(
518                "metadata.json format version {} is incompatible with required version {} \
519                 (formats {}..={} are upgraded on open); \
520                 rebuild and republish the index with this Summa version",
521                meta.version,
522                INDEX_META_FORMAT_VERSION,
523                OLDEST_MIGRATABLE_FORMAT_VERSION,
524                INDEX_META_FORMAT_VERSION - 1
525            )));
526        };
527        let mut deletion_ids = std::collections::HashSet::new();
528        for info in meta.segment_metas.values() {
529            if let Some(deletion) = &info.deletions
530                && (deletion.num_deleted == 0
531                    || deletion.num_deleted > info.num_docs
532                    || crate::segment::SegmentId::from_hex(&deletion.id).is_none()
533                    || meta.segment_metas.contains_key(&deletion.id)
534                    || !deletion_ids.insert(&deletion.id))
535            {
536                return Err(Error::Corruption(
537                    "invalid or aliased deletion metadata".into(),
538                ));
539            }
540        }
541        Ok((meta, migrated_from))
542    }
543
544    /// Load for a writer: upgrade a migratable older format and persist the
545    /// new stamp immediately, so the index is never left readable by a build
546    /// that would silently drop the fields this format added.
547    #[cfg(any(feature = "native", feature = "wasm"))]
548    pub(crate) async fn load_persisting_migration<D: crate::directories::DirectoryWriter>(
549        dir: &D,
550    ) -> Result<Self> {
551        let (meta, migrated_from) = Self::load_reporting_migration(dir).await?;
552        if let Some(from) = migrated_from {
553            meta.save(dir).await?;
554            log::warn!(
555                "[metadata_migration] metadata.json format version {from} upgraded to \
556                 {INDEX_META_FORMAT_VERSION} and persisted ({} segment(s)); builds that do not \
557                 support format {INDEX_META_FORMAT_VERSION} can no longer open this index. Segments from builds before 1.8.125 \
558                 fail at segment open and need a rebuild; segments without .rowstats cannot \
559                 be compacted until they are merged",
560                meta.segment_metas.len()
561            );
562        }
563        Ok(meta)
564    }
565
566    /// Save to directory (atomic: write temp file, then rename)
567    ///
568    /// Uses write-then-rename so a crash mid-write won't corrupt the
569    /// existing metadata file. On POSIX, rename is atomic.
570    pub async fn save<D: crate::directories::DirectoryWriter>(&self, dir: &D) -> Result<()> {
571        let bytes = self.serialize_to_bytes()?;
572        Self::save_bytes(dir, &bytes).await
573    }
574
575    /// Serialize metadata to bytes (cheap, no I/O).
576    /// Useful when you need to release a lock before doing disk I/O.
577    pub fn serialize_to_bytes(&self) -> Result<Vec<u8>> {
578        serde_json::to_vec_pretty(self).map_err(|e| Error::Serialization(e.to_string()))
579    }
580
581    /// Write pre-serialized metadata bytes to directory (atomic rename + fsync).
582    ///
583    /// The fsync ensures durability: without it, a power failure after rename
584    /// could lose the metadata update on systems with volatile write caches.
585    pub async fn save_bytes<D: crate::directories::DirectoryWriter>(
586        dir: &D,
587        bytes: &[u8],
588    ) -> Result<()> {
589        let tmp_path = Path::new(INDEX_META_TMP_FILENAME);
590        let final_path = Path::new(INDEX_META_FILENAME);
591        // Metadata is tiny, but `DirectoryWriter::write` does not guarantee
592        // the file contents themselves are fsynced. Finish the streaming
593        // writer first (filesystem implementations call `File::sync_all`),
594        // then atomically publish the durable temp file by rename.
595        let mut writer = dir.streaming_writer(tmp_path).await.map_err(Error::Io)?;
596        writer.write_all(bytes).map_err(Error::Io)?;
597        writer.finish().map_err(Error::Io)?;
598        // Rename is the logical commit point: after it succeeds, readers can
599        // observe the new generation and callers must publish the matching
600        // in-memory/tracker state. Directory fsync only strengthens crash
601        // durability. It cannot safely turn an already-visible rename into a
602        // reported pre-commit failure, because cleanup could then delete files
603        // referenced by the metadata now on disk.
604        dir.rename(tmp_path, final_path).await.map_err(Error::Io)?;
605        if let Err(error) = dir.sync().await {
606            log::error!(
607                "[metadata] directory fsync failed after committed rename: {}. \
608                 Continuing with the renamed generation; crash durability is not guaranteed",
609                error,
610            );
611        }
612        Ok(())
613    }
614
615    /// Fallible schema-aware loader used for lifecycle publication.
616    #[cfg_attr(not(feature = "native"), allow(dead_code))]
617    pub(crate) async fn try_load_trained_from_fields<D: crate::directories::Directory>(
618        vector_fields: &HashMap<u32, FieldVectorMeta>,
619        schema: &Schema,
620        dir: &D,
621    ) -> Result<Option<crate::segment::TrainedVectorStructures>> {
622        Self::load_trained_from_fields_impl(vector_fields, schema, dir).await
623    }
624
625    /// Load and validate the complete trained-artifact set described by a
626    /// `vector_fields` snapshot.
627    ///
628    /// This is intentionally all-or-nothing. A `Built` field is a durable
629    /// promise that every artifact required by its configured index exists and
630    /// is compatible with the schema. Returning a partial map would let some
631    /// segment builders publish ANN data while another field was silently
632    /// unusable, and would make the same index behave differently after a
633    /// restart.
634    async fn load_trained_from_fields_impl<D: crate::directories::Directory>(
635        vector_fields: &HashMap<u32, FieldVectorMeta>,
636        schema: &Schema,
637        dir: &D,
638    ) -> Result<Option<crate::segment::TrainedVectorStructures>> {
639        use std::sync::Arc;
640
641        let mut centroids = rustc_hash::FxHashMap::default();
642        let mut binary_quantizers = rustc_hash::FxHashMap::default();
643        let mut scann_artifacts = rustc_hash::FxHashMap::default();
644
645        let mut built_fields: Vec<_> = vector_fields
646            .iter()
647            .filter(|(_, meta)| matches!(meta.state, VectorIndexState::Built { .. }))
648            .collect();
649        built_fields.sort_unstable_by_key(|(field_id, _)| **field_id);
650
651        log::debug!(
652            "[trained] index={} loading trained structures, dense_vector_fields={:?}",
653            schema.index_label(),
654            vector_fields.keys().collect::<Vec<_>>()
655        );
656
657        for (field_id, field_meta) in built_fields {
658            log::debug!(
659                "[trained] index={} field {} state={:?} centroids_file={:?} codebook_file={:?}",
660                schema.index_label(),
661                field_id,
662                field_meta.state,
663                field_meta.centroids_file,
664                field_meta.codebook_file,
665            );
666            if field_meta.field_id != *field_id {
667                return Err(Error::Corruption(format!(
668                    "trained vector metadata key {field_id} contains field_id {}",
669                    field_meta.field_id
670                )));
671            }
672
673            let expected_clusters = match field_meta.state {
674                VectorIndexState::Built { num_clusters, .. } if num_clusters > 0 => num_clusters,
675                VectorIndexState::Built { .. } => {
676                    return Err(Error::Corruption(format!(
677                        "trained vector metadata field {field_id} has zero clusters"
678                    )));
679                }
680                VectorIndexState::Flat => unreachable!("built_fields contains only Built entries"),
681            };
682
683            let centroids_file = field_meta.centroids_file.as_deref().ok_or_else(|| {
684                Error::Corruption(format!(
685                    "trained vector metadata field {field_id} is Built but has no centroids_file"
686                ))
687            })?;
688            match field_meta.index_type {
689                VectorFieldIndexType::Float(VectorIndexType::IvfPq) => {
690                    return Err(Error::Corruption(format!(
691                        "field {field_id} was trained as IVF-PQ, which is no longer \
692                         supported; recreate the index with `ivf_tq` and reindex \
693                         (docs/turboquant-quantization.md)"
694                    )));
695                }
696                VectorFieldIndexType::Float(index_type @ VectorIndexType::IvfTq) => {
697                    let entry = schema
698                        .get_field_entry(crate::dsl::Field(*field_id))
699                        .ok_or_else(|| {
700                            Error::Corruption(format!(
701                                "trained vector metadata references missing field {field_id}"
702                            ))
703                        })?;
704                    let schema_config = entry
705                        .dense_vector_config
706                        .as_ref()
707                        .filter(|_| entry.field_type == crate::dsl::FieldType::DenseVector)
708                        .ok_or_else(|| {
709                            Error::Corruption(format!(
710                                "trained vector metadata field {field_id} is not a float dense field"
711                            ))
712                        })?;
713                    if schema_config.index_type != index_type {
714                        return Err(Error::Corruption(format!(
715                            "trained vector metadata field {field_id} uses {index_type:?}, schema requires {:?}",
716                            schema_config.index_type
717                        )));
718                    }
719                    let c: crate::structures::CoarseCentroids =
720                        load_trained_artifact(dir, *field_id, "centroids", centroids_file).await?;
721                    let expected_dim = schema_config.dim;
722                    let actual_clusters = c.num_clusters as usize;
723                    let expected_values =
724                        actual_clusters.checked_mul(expected_dim).ok_or_else(|| {
725                            Error::Corruption(format!(
726                                "trained centroid dimensions overflow for field {field_id}"
727                            ))
728                        })?;
729                    if actual_clusters == 0
730                        || actual_clusters > expected_clusters
731                        || c.dim == 0
732                        || c.dim != expected_dim
733                        || c.centroids.len() != expected_values
734                        || c.centroids.iter().any(|value| !value.is_finite())
735                    {
736                        return Err(Error::Corruption(format!(
737                            "trained centroids for field {field_id} do not match metadata/schema"
738                        )));
739                    }
740                    if !crate::structures::is_ivf_tq_cosine_generation(c.version) {
741                        return Err(Error::Corruption(format!(
742                            "trained IVF-TQ centroids for field {field_id} use an \
743                             unsupported legacy generation; rebuild the index"
744                        )));
745                    }
746                    c.validate_routing(schema_config.ivf_routing)
747                        .map_err(|error| {
748                            Error::Corruption(format!(
749                                "invalid trained centroid routing for field {field_id}: {error}"
750                            ))
751                        })?;
752                    // The TQ leaf codec is derived, never trained; ensure
753                    // `index_type` stays referenced for future variants.
754                    let _ = index_type;
755                    if field_meta.codebook_file.is_some() {
756                        return Err(Error::Corruption(format!(
757                            "trained IVF-TQ field {field_id} unexpectedly references a codebook file"
758                        )));
759                    }
760                    centroids.insert(*field_id, Arc::new(c));
761                }
762                VectorFieldIndexType::Binary(BinaryIndexType::Ivf) => {
763                    let entry = schema
764                        .get_field_entry(crate::dsl::Field(*field_id))
765                        .ok_or_else(|| {
766                            Error::Corruption(format!(
767                                "trained vector metadata references missing field {field_id}"
768                            ))
769                        })?;
770                    let schema_config = entry
771                        .binary_dense_vector_config
772                        .as_ref()
773                        .filter(|config| {
774                            entry.field_type == crate::dsl::FieldType::BinaryDenseVector
775                                && config.index_type == BinaryIndexType::Ivf
776                        })
777                        .ok_or_else(|| {
778                            Error::Corruption(format!(
779                                "trained vector metadata field {field_id} is not a binary IVF field"
780                            ))
781                        })?;
782                    let quantizer: crate::structures::BinaryCoarseQuantizer =
783                        load_trained_artifact(dir, *field_id, "binary centroids", centroids_file)
784                            .await?;
785                    quantizer.validate().map_err(|error| {
786                        Error::Corruption(format!(
787                            "invalid binary coarse quantizer for field {field_id}: {error}"
788                        ))
789                    })?;
790                    let actual_clusters = quantizer.num_clusters as usize;
791                    if actual_clusters > expected_clusters
792                        || schema_config.dim != quantizer.dim_bits
793                    {
794                        return Err(Error::Corruption(format!(
795                            "binary coarse quantizer for field {field_id} does not match metadata/schema"
796                        )));
797                    }
798                    quantizer
799                        .validate_routing(schema_config.ivf_routing)
800                        .map_err(|error| {
801                            Error::Corruption(format!(
802                                "invalid binary centroid routing for field {field_id}: {error}"
803                            ))
804                        })?;
805                    binary_quantizers.insert(*field_id, Arc::new(quantizer));
806                }
807                VectorFieldIndexType::Float(VectorIndexType::Scann)
808                | VectorFieldIndexType::Binary(BinaryIndexType::Scann) => {
809                    if field_meta.codebook_file.is_some() {
810                        return Err(Error::Corruption(format!(
811                            "trained ScaNN field {field_id} unexpectedly references a separate codebook file"
812                        )));
813                    }
814                    let expected_generation = field_meta.artifact_generation.ok_or_else(|| {
815                        Error::Corruption(format!(
816                            "trained ScaNN field {field_id} has no artifact generation"
817                        ))
818                    })?;
819                    let expected_artifact_id = field_meta.artifact_id.ok_or_else(|| {
820                        Error::Corruption(format!(
821                            "trained ScaNN field {field_id} has no artifact fingerprint"
822                        ))
823                    })?;
824                    validate_trained_artifact_path(
825                        field_id.to_owned(),
826                        "ScaNN artifact",
827                        centroids_file,
828                    )?;
829                    let path = Path::new(centroids_file);
830                    let slice = dir.open_read(path).await.map_err(|error| {
831                        Error::Corruption(format!(
832                            "failed to open trained ScaNN artifact '{centroids_file}' for field {field_id}: {error}"
833                        ))
834                    })?;
835                    let raw = slice.read_bytes().await.map_err(|error| {
836                        Error::Corruption(format!(
837                            "failed to map trained ScaNN artifact '{centroids_file}' for field {field_id}: {error}"
838                        ))
839                    })?;
840                    let artifact =
841                        crate::segment::ScannTrainedArtifactBytes::open(raw).map_err(|error| {
842                            Error::Corruption(format!(
843                                "invalid trained ScaNN artifact for field {field_id}: {error}"
844                            ))
845                        })?;
846                    if artifact.generation() != expected_generation
847                        || artifact.artifact_id() != expected_artifact_id
848                        || artifact.config().num_leaves as usize != expected_clusters
849                    {
850                        return Err(Error::Corruption(format!(
851                            "trained ScaNN artifact for field {field_id} does not match metadata"
852                        )));
853                    }
854                    let entry = schema
855                        .get_field_entry(crate::dsl::Field(*field_id))
856                        .ok_or_else(|| {
857                            Error::Corruption(format!(
858                                "trained vector metadata references missing field {field_id}"
859                            ))
860                        })?;
861                    let schema_matches = match field_meta.index_type {
862                        VectorFieldIndexType::Float(VectorIndexType::Scann) => entry
863                            .dense_vector_config
864                            .as_ref()
865                            .filter(|_| entry.field_type == crate::dsl::FieldType::DenseVector)
866                            .is_some_and(|config| {
867                                config.index_type == VectorIndexType::Scann
868                                    && config.dim == artifact.config().dimension as usize
869                                    && scann_explicit_geometry_matches(
870                                        config.num_clusters,
871                                        config.tree_levels,
872                                        artifact.config().num_leaves as usize,
873                                        artifact.config().tree_levels,
874                                    )
875                                    && matches!(
876                                        artifact.config().encoding,
877                                        crate::structures::vector::scann::ScannEncoding::AsymmetricHash { .. }
878                                    )
879                            }),
880                        VectorFieldIndexType::Binary(BinaryIndexType::Scann) => entry
881                            .binary_dense_vector_config
882                            .as_ref()
883                            .filter(|_| {
884                                entry.field_type == crate::dsl::FieldType::BinaryDenseVector
885                            })
886                            .is_some_and(|config| {
887                                config.index_type == BinaryIndexType::Scann
888                                    && config.dim == artifact.config().dimension as usize
889                                    && scann_explicit_geometry_matches(
890                                        config.num_clusters,
891                                        config.tree_levels,
892                                        artifact.config().num_leaves as usize,
893                                        artifact.config().tree_levels,
894                                    )
895                                    && artifact.config().encoding
896                                        == crate::structures::vector::scann::ScannEncoding::BinaryHamming
897                            }),
898                        _ => false,
899                    };
900                    if !schema_matches {
901                        return Err(Error::Corruption(format!(
902                            "trained ScaNN artifact for field {field_id} does not match schema geometry/encoding"
903                        )));
904                    }
905                    scann_artifacts.insert(*field_id, Arc::new(artifact));
906                }
907                unsupported => {
908                    return Err(Error::Corruption(format!(
909                        "field {field_id} is Built for {unsupported:?}, which has no global IVF artifacts"
910                    )));
911                }
912            }
913        }
914
915        if centroids.is_empty() && binary_quantizers.is_empty() && scann_artifacts.is_empty() {
916            Ok(None)
917        } else {
918            let trained = crate::segment::TrainedVectorStructures {
919                #[cfg(feature = "native")]
920                _ann_pins: Default::default(),
921                centroids,
922                binary_quantizers,
923                scann_artifacts,
924            };
925            #[cfg(feature = "native")]
926            let trained = {
927                let mut trained = trained;
928                trained.pin_ann_structures(crate::segment::pin::pin_policy());
929                trained
930            };
931            Ok(Some(trained))
932        }
933    }
934}
935
936/// Omitted ScaNN geometry is autopilot: the trained artifact's resolved
937/// values are authoritative and are also pinned by `VectorIndexState::Built`.
938/// Explicit schema values remain strict compatibility requirements.
939fn scann_explicit_geometry_matches(
940    configured_leaves: Option<usize>,
941    configured_levels: Option<u8>,
942    resolved_leaves: usize,
943    resolved_levels: u8,
944) -> bool {
945    configured_leaves.is_none_or(|leaves| leaves == resolved_leaves)
946        && configured_levels.is_none_or(|levels| levels == resolved_levels)
947}
948
949fn validate_trained_artifact_path(field_id: u32, kind: &str, filename: &str) -> Result<()> {
950    use std::path::Component;
951
952    let path = Path::new(filename);
953    if filename.is_empty()
954        || path.is_absolute()
955        || path.components().any(|component| {
956            matches!(
957                component,
958                Component::ParentDir | Component::RootDir | Component::Prefix(_)
959            )
960        })
961    {
962        return Err(Error::Corruption(format!(
963            "trained {kind} path for field {field_id} is not a safe relative path: '{filename}'"
964        )));
965    }
966    Ok(())
967}
968
969async fn load_trained_artifact<T, D>(
970    dir: &D,
971    field_id: u32,
972    kind: &str,
973    filename: &str,
974) -> Result<T>
975where
976    T: serde::de::DeserializeOwned,
977    D: crate::directories::Directory,
978{
979    validate_trained_artifact_path(field_id, kind, filename)?;
980    let path = Path::new(filename);
981    let file_size = dir.file_size(path).await.map_err(|error| {
982        Error::Corruption(format!(
983            "failed to stat trained {kind} '{filename}' for field {field_id}: {error}"
984        ))
985    })?;
986    validate_trained_artifact_size(field_id, kind, filename, file_size)?;
987    let slice = dir.open_read(path).await.map_err(|error| {
988        Error::Corruption(format!(
989            "failed to open trained {kind} '{filename}' for field {field_id}: {error}"
990        ))
991    })?;
992    validate_trained_artifact_size(field_id, kind, filename, slice.len())?;
993    let bytes = slice.read_bytes().await.map_err(|error| {
994        Error::Corruption(format!(
995            "failed to read trained {kind} '{filename}' for field {field_id}: {error}"
996        ))
997    })?;
998    let (artifact, consumed) = bincode::serde::decode_from_slice::<T, _>(
999        bytes.as_slice(),
1000        bincode::config::standard().with_limit::<MAX_TRAINED_ARTIFACT_BYTES>(),
1001    )
1002    .map_err(|error| {
1003        Error::Corruption(format!(
1004            "failed to deserialize trained {kind} '{filename}' for field {field_id}: {error}"
1005        ))
1006    })?;
1007    if consumed != bytes.len() {
1008        return Err(Error::Corruption(format!(
1009            "trained {kind} '{filename}' for field {field_id} has {} trailing bytes",
1010            bytes.len() - consumed
1011        )));
1012    }
1013    Ok(artifact)
1014}
1015
1016fn validate_trained_artifact_size(
1017    field_id: u32,
1018    kind: &str,
1019    filename: &str,
1020    file_size: u64,
1021) -> Result<()> {
1022    if file_size > MAX_TRAINED_ARTIFACT_BYTES as u64 {
1023        return Err(Error::Corruption(format!(
1024            "trained {kind} '{filename}' for field {field_id} is {file_size} bytes, \
1025             exceeding the {MAX_TRAINED_ARTIFACT_BYTES}-byte safety limit"
1026        )));
1027    }
1028    Ok(())
1029}
1030
1031#[cfg(test)]
1032mod tests {
1033    use super::*;
1034    use crate::directories::DirectoryWriter;
1035
1036    #[derive(Clone, Default)]
1037    struct SyncFailDirectory(crate::directories::RamDirectory);
1038
1039    #[async_trait::async_trait]
1040    impl crate::directories::Directory for SyncFailDirectory {
1041        async fn exists(&self, path: &Path) -> std::io::Result<bool> {
1042            self.0.exists(path).await
1043        }
1044
1045        async fn file_size(&self, path: &Path) -> std::io::Result<u64> {
1046            self.0.file_size(path).await
1047        }
1048
1049        async fn open_read(&self, path: &Path) -> std::io::Result<crate::directories::FileHandle> {
1050            self.0.open_read(path).await
1051        }
1052
1053        async fn read_range(
1054            &self,
1055            path: &Path,
1056            range: std::ops::Range<u64>,
1057        ) -> std::io::Result<crate::directories::OwnedBytes> {
1058            self.0.read_range(path, range).await
1059        }
1060
1061        async fn list_files(&self, prefix: &Path) -> std::io::Result<Vec<std::path::PathBuf>> {
1062            self.0.list_files(prefix).await
1063        }
1064
1065        async fn open_lazy(&self, path: &Path) -> std::io::Result<crate::directories::FileHandle> {
1066            self.0.open_lazy(path).await
1067        }
1068    }
1069
1070    #[async_trait::async_trait]
1071    impl crate::directories::DirectoryWriter for SyncFailDirectory {
1072        async fn write(&self, path: &Path, data: &[u8]) -> std::io::Result<()> {
1073            self.0.write(path, data).await
1074        }
1075
1076        async fn delete(&self, path: &Path) -> std::io::Result<()> {
1077            self.0.delete(path).await
1078        }
1079
1080        async fn rename(&self, from: &Path, to: &Path) -> std::io::Result<()> {
1081            self.0.rename(from, to).await
1082        }
1083
1084        async fn sync(&self) -> std::io::Result<()> {
1085            Err(std::io::Error::other("injected directory fsync failure"))
1086        }
1087
1088        async fn streaming_writer(
1089            &self,
1090            path: &Path,
1091        ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
1092            self.0.streaming_writer(path).await
1093        }
1094    }
1095
1096    fn test_schema() -> Schema {
1097        Schema::default()
1098    }
1099
1100    fn dense_schema(index_type: VectorIndexType) -> (Schema, crate::dsl::Field) {
1101        let mut builder = crate::dsl::SchemaBuilder::default();
1102        let config = match index_type {
1103            VectorIndexType::IvfTq => crate::dsl::DenseVectorConfig::ivf_tq(2, Some(1), 1),
1104            other => panic!("unsupported trained test index type: {other:?}"),
1105        };
1106        let field = builder.add_dense_vector_field_with_config("embedding", true, true, config);
1107        (builder.build(), field)
1108    }
1109
1110    fn test_centroids() -> crate::structures::CoarseCentroids {
1111        crate::structures::CoarseCentroids {
1112            num_clusters: 1,
1113            dim: 2,
1114            centroids: vec![0.25, 0.75],
1115            version: crate::structures::mark_ivf_tq_cosine_generation(7),
1116            soar_config: None,
1117            routing_index: None,
1118        }
1119    }
1120
1121    async fn write_bincode(
1122        directory: &crate::directories::RamDirectory,
1123        filename: &str,
1124        value: &impl serde::Serialize,
1125    ) {
1126        let bytes = bincode::serde::encode_to_vec(value, bincode::config::standard()).unwrap();
1127        directory.write(Path::new(filename), &bytes).await.unwrap();
1128    }
1129
1130    #[test]
1131    fn test_metadata_init() {
1132        let mut meta = IndexMetadata::new(test_schema());
1133        assert_eq!(meta.total_vectors, 0);
1134        assert!(meta.segment_metas.is_empty());
1135        assert!(!meta.is_field_built(0));
1136
1137        meta.init_field(0, VectorIndexType::IvfTq);
1138        assert!(!meta.is_field_built(0));
1139        assert!(meta.vector_fields.contains_key(&0));
1140    }
1141
1142    #[tokio::test]
1143    async fn phrase_limits_load_from_legacy_and_configured_metadata_and_reject_corruption() {
1144        let directory = crate::directories::RamDirectory::new();
1145        let legacy = IndexMetadata::new(test_schema())
1146            .serialize_to_bytes()
1147            .unwrap();
1148        assert!(!String::from_utf8_lossy(&legacy).contains("max_l1_phrase_terms"));
1149        directory
1150            .write(Path::new(INDEX_META_FILENAME), &legacy)
1151            .await
1152            .unwrap();
1153        let loaded = IndexMetadata::load(&directory).await.unwrap();
1154        assert_eq!(loaded.schema.max_l1_phrase_terms(), 64);
1155        assert_eq!(loaded.serialize_to_bytes().unwrap(), legacy);
1156
1157        let mut configured: serde_json::Value = serde_json::from_slice(&legacy).unwrap();
1158        configured["schema"]["max_l1_phrase_terms"] = 300.into();
1159        directory
1160            .write(
1161                Path::new(INDEX_META_FILENAME),
1162                &serde_json::to_vec(&configured).unwrap(),
1163            )
1164            .await
1165            .unwrap();
1166        let loaded = IndexMetadata::load(&directory).await.unwrap();
1167        assert_eq!(loaded.schema.max_l1_phrase_terms(), 300);
1168        loaded.save(&directory).await.unwrap();
1169        assert_eq!(
1170            IndexMetadata::load(&directory)
1171                .await
1172                .unwrap()
1173                .schema
1174                .max_l1_phrase_terms(),
1175            300
1176        );
1177
1178        for invalid in [
1179            serde_json::json!(0),
1180            serde_json::json!(-1),
1181            serde_json::json!(1.5),
1182            serde_json::json!(4294967296u64),
1183        ] {
1184            configured["schema"]["max_l1_phrase_terms"] = invalid;
1185            directory
1186                .write(
1187                    Path::new(INDEX_META_FILENAME),
1188                    &serde_json::to_vec(&configured).unwrap(),
1189                )
1190                .await
1191                .unwrap();
1192            assert!(
1193                IndexMetadata::load(&directory).await.is_err(),
1194                "must not replace a corrupt cap with the default"
1195            );
1196        }
1197    }
1198
1199    #[tokio::test]
1200    async fn load_refuses_metadata_stamped_with_a_newer_format_version() {
1201        let directory = crate::directories::RamDirectory::new();
1202        let mut metadata = IndexMetadata::new(test_schema());
1203        metadata.version = INDEX_META_FORMAT_VERSION + 1;
1204        metadata.save(&directory).await.unwrap();
1205
1206        let error = IndexMetadata::load(&directory)
1207            .await
1208            .expect_err("metadata from a newer format version must be refused, not silently pruned")
1209            .to_string();
1210        assert!(
1211            error.contains(&format!("version {}", INDEX_META_FORMAT_VERSION + 1)),
1212            "{error}"
1213        );
1214        assert!(error.contains("incompatible"), "{error}");
1215    }
1216
1217    fn stamped_previous_format_bytes(metadata: &IndexMetadata) -> Vec<u8> {
1218        let mut raw: serde_json::Value =
1219            serde_json::from_slice(&metadata.serialize_to_bytes().unwrap()).unwrap();
1220        raw["version"] = serde_json::Value::from(INDEX_META_FORMAT_VERSION - 1);
1221        serde_json::to_vec(&raw).unwrap()
1222    }
1223
1224    #[tokio::test]
1225    async fn read_only_migration_preserves_format_6_and_7_metadata_bytes() {
1226        use crate::directories::Directory;
1227        for version in [6, 7] {
1228            let directory = crate::directories::RamDirectory::new();
1229            let mut metadata = IndexMetadata::new(test_schema());
1230            metadata.version = version;
1231            metadata.add_segment("kept".into(), 7);
1232            metadata.save(&directory).await.unwrap();
1233            let before = directory
1234                .open_read(Path::new(INDEX_META_FILENAME))
1235                .await
1236                .unwrap()
1237                .read_bytes()
1238                .await
1239                .unwrap();
1240            let (loaded, migrated) = IndexMetadata::load_reporting_migration(&directory)
1241                .await
1242                .unwrap();
1243            assert_eq!(migrated, Some(version));
1244            assert_eq!(loaded.version, INDEX_META_FORMAT_VERSION);
1245            assert_eq!(loaded.segment_metas["kept"].num_docs, 7);
1246            let after = directory
1247                .open_read(Path::new(INDEX_META_FILENAME))
1248                .await
1249                .unwrap()
1250                .read_bytes()
1251                .await
1252                .unwrap();
1253            assert_eq!(after.as_slice(), before.as_slice());
1254            loaded.save(&directory).await.unwrap();
1255            assert_eq!(
1256                IndexMetadata::load_reporting_migration(&directory)
1257                    .await
1258                    .unwrap()
1259                    .1,
1260                None
1261            );
1262        }
1263    }
1264
1265    #[tokio::test]
1266    async fn load_migrates_previous_metadata_to_the_current_format() {
1267        let directory = crate::directories::RamDirectory::new();
1268        let mut metadata = IndexMetadata::new(test_schema());
1269        metadata.add_segment("kept".to_string(), 7);
1270        directory
1271            .write(
1272                Path::new(INDEX_META_FILENAME),
1273                &stamped_previous_format_bytes(&metadata),
1274            )
1275            .await
1276            .unwrap();
1277
1278        let (loaded, migrated_from) = IndexMetadata::load_reporting_migration(&directory)
1279            .await
1280            .expect("the previous format remains readable");
1281        assert_eq!(migrated_from, Some(INDEX_META_FORMAT_VERSION - 1));
1282        assert_eq!(loaded.version, INDEX_META_FORMAT_VERSION);
1283        assert_eq!(loaded.segment_metas["kept"].num_docs, 7);
1284        assert!(loaded.segment_metas["kept"].deletions.is_none());
1285
1286        let (_, current) = IndexMetadata::load_reporting_migration(&directory)
1287            .await
1288            .unwrap();
1289        assert_eq!(
1290            current,
1291            Some(INDEX_META_FORMAT_VERSION - 1),
1292            "load alone never persists"
1293        );
1294        metadata.save(&directory).await.unwrap();
1295        let (_, current) = IndexMetadata::load_reporting_migration(&directory)
1296            .await
1297            .unwrap();
1298        assert_eq!(
1299            current, None,
1300            "current-format metadata reports no migration"
1301        );
1302    }
1303
1304    #[tokio::test]
1305    async fn tmp_recovery_migrates_previous_metadata() {
1306        let directory = crate::directories::RamDirectory::new();
1307        let mut metadata = IndexMetadata::new(test_schema());
1308        metadata.add_segment("kept".to_string(), 3);
1309        directory
1310            .write(
1311                Path::new(INDEX_META_TMP_FILENAME),
1312                &stamped_previous_format_bytes(&metadata),
1313            )
1314            .await
1315            .unwrap();
1316
1317        let loaded = IndexMetadata::load(&directory)
1318            .await
1319            .expect("temp-file recovery applies the same migration");
1320        assert_eq!(loaded.version, INDEX_META_FORMAT_VERSION);
1321        assert_eq!(loaded.segment_metas["kept"].num_docs, 3);
1322    }
1323
1324    #[tokio::test]
1325    async fn load_refuses_metadata_older_than_format_6() {
1326        // Format 5 predates the position-list v3 layout; its segments are
1327        // unreadable, so the stamp must stay a rebuild boundary.
1328        let directory = crate::directories::RamDirectory::new();
1329        let mut metadata = IndexMetadata::new(test_schema());
1330        metadata.version = OLDEST_MIGRATABLE_FORMAT_VERSION - 1;
1331        metadata.save(&directory).await.unwrap();
1332
1333        let error = IndexMetadata::load(&directory)
1334            .await
1335            .expect_err("metadata predating compatible position streams must be refused")
1336            .to_string();
1337        assert!(
1338            error.contains(&format!("version {}", OLDEST_MIGRATABLE_FORMAT_VERSION - 1)),
1339            "{error}"
1340        );
1341        assert!(error.contains("incompatible"), "{error}");
1342    }
1343
1344    #[tokio::test]
1345    async fn tmp_recovery_refuses_metadata_stamped_with_a_newer_format_version() {
1346        let directory = crate::directories::RamDirectory::new();
1347        let mut metadata = IndexMetadata::new(test_schema());
1348        metadata.version = INDEX_META_FORMAT_VERSION + 1;
1349        let bytes = metadata.serialize_to_bytes().unwrap();
1350        // Simulate a crash between write and rename: only the temp file exists.
1351        directory
1352            .write(Path::new(INDEX_META_TMP_FILENAME), &bytes)
1353            .await
1354            .unwrap();
1355
1356        let error = IndexMetadata::load(&directory)
1357            .await
1358            .expect_err("temp-file recovery must apply the same version gate")
1359            .to_string();
1360        assert!(
1361            error.contains(&format!("version {}", INDEX_META_FORMAT_VERSION + 1)),
1362            "{error}"
1363        );
1364    }
1365
1366    #[tokio::test]
1367    async fn save_treats_post_rename_sync_failure_as_committed() {
1368        let directory = SyncFailDirectory::default();
1369        let mut metadata = IndexMetadata::new(test_schema());
1370        metadata.add_segment("committed".to_string(), 7);
1371
1372        metadata.save(&directory).await.unwrap();
1373
1374        let loaded = IndexMetadata::load(&directory).await.unwrap();
1375        assert_eq!(loaded.segment_doc_count("committed"), Some(7));
1376    }
1377
1378    #[tokio::test]
1379    async fn trained_artifacts_load_only_when_the_complete_built_set_is_valid() {
1380        let mut builder = crate::dsl::SchemaBuilder::default();
1381        let config = crate::dsl::DenseVectorConfig::ivf_tq(2, Some(1), 1);
1382        let first = builder.add_dense_vector_field_with_config(
1383            "first_embedding",
1384            true,
1385            true,
1386            config.clone(),
1387        );
1388        let second =
1389            builder.add_dense_vector_field_with_config("second_embedding", true, true, config);
1390        let schema = builder.build();
1391        let directory = crate::directories::RamDirectory::new();
1392        let mut metadata = IndexMetadata::new(schema.clone());
1393        metadata.init_field(first.0, VectorIndexType::IvfTq);
1394        metadata.init_field(second.0, VectorIndexType::IvfTq);
1395        metadata.mark_field_built(first.0, 10, 1, "field_0_centroids.bin".into(), None);
1396        metadata.mark_field_built(second.0, 10, 1, "field_1_centroids.bin".into(), None);
1397        write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
1398
1399        let error = IndexMetadata::try_load_trained_from_fields(
1400            &metadata.vector_fields,
1401            &schema,
1402            &directory,
1403        )
1404        .await
1405        .err()
1406        .expect("missing artifact must fail the complete load")
1407        .to_string();
1408        assert!(error.contains("field_1_centroids.bin"), "{error}");
1409        assert!(error.contains("field 1"), "{error}");
1410    }
1411
1412    #[tokio::test]
1413    async fn index_open_fails_closed_when_built_artifact_is_missing() {
1414        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1415        let directory = crate::directories::RamDirectory::new();
1416        let mut metadata = IndexMetadata::new(schema);
1417        metadata.init_field(field.0, VectorIndexType::IvfTq);
1418        metadata.mark_field_built(field.0, 10, 1, "missing_centroids.bin".into(), None);
1419        metadata.save(&directory).await.unwrap();
1420
1421        let error = match crate::index::Index::open(directory, crate::index::IndexConfig::default())
1422            .await
1423        {
1424            Ok(_) => panic!("Index::open accepted a Built field with no artifact"),
1425            Err(error) => error.to_string(),
1426        };
1427        assert!(error.contains("missing_centroids.bin"), "{error}");
1428    }
1429
1430    #[tokio::test]
1431    async fn ivf_tq_built_state_rejects_a_codebook_file() {
1432        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1433        let directory = crate::directories::RamDirectory::new();
1434        let mut metadata = IndexMetadata::new(schema.clone());
1435        metadata.init_field(field.0, VectorIndexType::IvfTq);
1436        metadata.mark_field_built(
1437            field.0,
1438            10,
1439            1,
1440            "field_0_centroids.bin".into(),
1441            Some("field_0_codebook.bin".into()),
1442        );
1443        write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
1444
1445        let error = IndexMetadata::try_load_trained_from_fields(
1446            &metadata.vector_fields,
1447            &schema,
1448            &directory,
1449        )
1450        .await
1451        .err()
1452        .expect("IVF-TQ Built state with a codebook file must fail")
1453        .to_string();
1454        assert!(error.contains("codebook"), "{error}");
1455    }
1456
1457    #[tokio::test]
1458    async fn legacy_ivf_tq_centroid_generation_is_rejected_while_loading() {
1459        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1460        let directory = crate::directories::RamDirectory::new();
1461        let mut metadata = IndexMetadata::new(schema.clone());
1462        metadata.init_field(field.0, VectorIndexType::IvfTq);
1463        metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1464        let mut legacy = test_centroids();
1465        legacy.version = 7;
1466        write_bincode(&directory, "field_0_centroids.bin", &legacy).await;
1467
1468        let error = IndexMetadata::try_load_trained_from_fields(
1469            &metadata.vector_fields,
1470            &schema,
1471            &directory,
1472        )
1473        .await
1474        .err()
1475        .expect("legacy IVF-TQ centroid state must fail while loading")
1476        .to_string();
1477        assert!(error.contains("unsupported legacy generation"), "{error}");
1478        assert!(error.contains("rebuild the index"), "{error}");
1479    }
1480
1481    #[tokio::test]
1482    async fn legacy_ivf_pq_trained_field_fails_with_actionable_error() {
1483        // Simulates metadata written by a pre-removal version: the schema
1484        // gate rejects `ivf_pq` fields, so build the raw field-state map
1485        // directly against a current-format schema.
1486        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1487        let directory = crate::directories::RamDirectory::new();
1488        let mut metadata = IndexMetadata::new(schema.clone());
1489        metadata.init_field(field.0, VectorIndexType::IvfTq);
1490        metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1491        // Overwrite the recorded type the way pre-removal metadata carries it
1492        // (init_field never downgrades an existing entry).
1493        metadata
1494            .vector_fields
1495            .get_mut(&field.0)
1496            .expect("field initialized")
1497            .index_type = VectorFieldIndexType::Float(VectorIndexType::IvfPq);
1498        write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
1499
1500        let error = IndexMetadata::try_load_trained_from_fields(
1501            &metadata.vector_fields,
1502            &schema,
1503            &directory,
1504        )
1505        .await
1506        .err()
1507        .expect("legacy IVF-PQ trained state must fail loudly")
1508        .to_string();
1509        assert!(error.contains("no longer"), "{error}");
1510        assert!(error.contains("ivf_tq"), "{error}");
1511    }
1512
1513    #[tokio::test]
1514    async fn requested_cluster_count_accepts_a_quality_clamped_artifact() {
1515        let mut builder = crate::dsl::SchemaBuilder::default();
1516        let field = builder.add_dense_vector_field_with_config(
1517            "embedding",
1518            true,
1519            true,
1520            crate::dsl::DenseVectorConfig::ivf_tq(2, Some(4), 1),
1521        );
1522        let schema = builder.build();
1523        let directory = crate::directories::RamDirectory::new();
1524        let mut metadata = IndexMetadata::new(schema.clone());
1525        metadata.init_field(field.0, VectorIndexType::IvfTq);
1526        metadata.mark_field_built(field.0, 1, 4, "field_0_centroids.bin".into(), None);
1527        write_bincode(&directory, "field_0_centroids.bin", &test_centroids()).await;
1528
1529        let trained = IndexMetadata::try_load_trained_from_fields(
1530            &metadata.vector_fields,
1531            &schema,
1532            &directory,
1533        )
1534        .await
1535        .unwrap()
1536        .unwrap();
1537        assert_eq!(trained.centroids[&field.0].num_clusters, 1);
1538    }
1539
1540    #[tokio::test]
1541    async fn trained_artifact_loader_rejects_trailing_data() {
1542        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1543        let directory = crate::directories::RamDirectory::new();
1544        let mut metadata = IndexMetadata::new(schema.clone());
1545        metadata.init_field(field.0, VectorIndexType::IvfTq);
1546        metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1547        let mut bytes =
1548            bincode::serde::encode_to_vec(test_centroids(), bincode::config::standard()).unwrap();
1549        bytes.extend_from_slice(&[0xaa, 0xbb]);
1550        directory
1551            .write(Path::new("field_0_centroids.bin"), &bytes)
1552            .await
1553            .unwrap();
1554
1555        let error = IndexMetadata::try_load_trained_from_fields(
1556            &metadata.vector_fields,
1557            &schema,
1558            &directory,
1559        )
1560        .await
1561        .err()
1562        .expect("trailing artifact bytes must fail validation")
1563        .to_string();
1564        assert!(error.contains("trailing bytes"), "{error}");
1565    }
1566
1567    #[test]
1568    fn trained_artifact_size_limit_rejects_before_reading() {
1569        let error = validate_trained_artifact_size(
1570            3,
1571            "centroids",
1572            "field_3_centroids.bin",
1573            MAX_TRAINED_ARTIFACT_BYTES as u64 + 1,
1574        )
1575        .unwrap_err()
1576        .to_string();
1577        assert!(error.contains("exceeding"), "{error}");
1578        assert!(error.contains("field 3"), "{error}");
1579    }
1580
1581    #[tokio::test]
1582    async fn trained_artifact_decode_limit_rejects_forged_collection_length() {
1583        let (schema, field) = dense_schema(VectorIndexType::IvfTq);
1584        let directory = crate::directories::RamDirectory::new();
1585        let mut metadata = IndexMetadata::new(schema.clone());
1586        metadata.init_field(field.0, VectorIndexType::IvfTq);
1587        metadata.mark_field_built(field.0, 10, 1, "field_0_centroids.bin".into(), None);
1588
1589        // CoarseCentroids begins with num_clusters=1, dim=2, then the Vec
1590        // length. Bincode's standard varint marker 253 introduces a u64; this
1591        // tiny payload claims an impossible f32 vector and must hit the decode
1592        // limit before any large allocation is attempted.
1593        let mut bytes = vec![1, 2, 253];
1594        bytes.extend_from_slice(&u64::MAX.to_le_bytes());
1595        directory
1596            .write(Path::new("field_0_centroids.bin"), &bytes)
1597            .await
1598            .unwrap();
1599
1600        let error = IndexMetadata::try_load_trained_from_fields(
1601            &metadata.vector_fields,
1602            &schema,
1603            &directory,
1604        )
1605        .await
1606        .err()
1607        .expect("forged collection length must fail the bounded decoder")
1608        .to_string();
1609        assert!(error.contains("failed to deserialize"), "{error}");
1610    }
1611
1612    #[test]
1613    fn test_metadata_segments() {
1614        let mut meta = IndexMetadata::new(test_schema());
1615        meta.add_segment("abc123".to_string(), 50);
1616        meta.add_segment("def456".to_string(), 100);
1617        assert_eq!(meta.segment_metas.len(), 2);
1618        assert_eq!(meta.segment_doc_count("abc123"), Some(50));
1619        assert_eq!(meta.segment_doc_count("def456"), Some(100));
1620
1621        // Overwrites existing
1622        meta.add_segment("abc123".to_string(), 75);
1623        assert_eq!(meta.segment_metas.len(), 2);
1624        assert_eq!(meta.segment_doc_count("abc123"), Some(75));
1625
1626        meta.remove_segment("abc123");
1627        assert_eq!(meta.segment_metas.len(), 1);
1628        assert!(meta.has_segment("def456"));
1629        assert!(!meta.has_segment("abc123"));
1630    }
1631
1632    #[test]
1633    fn test_mark_field_built() {
1634        let mut meta = IndexMetadata::new(test_schema());
1635        meta.init_field(0, VectorIndexType::IvfTq);
1636        meta.total_vectors = 10000;
1637
1638        assert!(!meta.is_field_built(0));
1639
1640        meta.mark_field_built(0, 10000, 256, "field_0_centroids.bin".to_string(), None);
1641
1642        assert!(meta.is_field_built(0));
1643        let field = meta.get_field_meta(0).unwrap();
1644        assert_eq!(
1645            field.centroids_file.as_deref(),
1646            Some("field_0_centroids.bin")
1647        );
1648    }
1649
1650    #[test]
1651    fn scann_metadata_persists_generation_and_fingerprint_and_defaults_old_json() {
1652        let mut meta = IndexMetadata::new(test_schema());
1653        meta.init_field(3, VectorIndexType::Scann);
1654        meta.mark_scann_field_built(
1655            3,
1656            100_000,
1657            1_000,
1658            "field_3_scann_17.bin".to_string(),
1659            17,
1660            0xdecafbad,
1661        )
1662        .unwrap();
1663
1664        let bytes = meta.serialize_to_bytes().unwrap();
1665        let decoded: IndexMetadata = serde_json::from_slice(&bytes).unwrap();
1666        let field = decoded.get_field_meta(3).unwrap();
1667        assert_eq!(field.artifact_generation, Some(17));
1668        assert_eq!(field.artifact_id, Some(0xdecafbad));
1669
1670        let mut legacy_json = serde_json::to_value(&decoded).unwrap();
1671        legacy_json["vector_fields"]["3"]
1672            .as_object_mut()
1673            .unwrap()
1674            .remove("artifact_generation");
1675        legacy_json["vector_fields"]["3"]
1676            .as_object_mut()
1677            .unwrap()
1678            .remove("artifact_id");
1679        let legacy: IndexMetadata = serde_json::from_value(legacy_json).unwrap();
1680        let legacy_field = legacy.get_field_meta(3).unwrap();
1681        assert_eq!(legacy_field.artifact_generation, None);
1682        assert_eq!(legacy_field.artifact_id, None);
1683    }
1684
1685    #[test]
1686    fn scann_metadata_refuses_zero_or_non_scann_generation() {
1687        let mut meta = IndexMetadata::new(test_schema());
1688        meta.init_field(0, VectorIndexType::Scann);
1689        assert!(
1690            meta.mark_scann_field_built(0, 100_000, 1_000, "artifact.bin".into(), 0, 1)
1691                .is_err()
1692        );
1693        meta.init_field(1, VectorIndexType::IvfTq);
1694        assert!(
1695            meta.mark_scann_field_built(1, 100_000, 1_000, "artifact.bin".into(), 1, 2)
1696                .is_err()
1697        );
1698    }
1699
1700    #[test]
1701    fn scann_autopilot_accepts_resolved_billion_scale_three_level_geometry() {
1702        assert!(scann_explicit_geometry_matches(None, None, 10_000_000, 3));
1703        assert!(!scann_explicit_geometry_matches(
1704            Some(1_000_000),
1705            None,
1706            10_000_000,
1707            3
1708        ));
1709        assert!(!scann_explicit_geometry_matches(
1710            None,
1711            Some(1),
1712            10_000_000,
1713            3
1714        ));
1715    }
1716
1717    #[test]
1718    fn total_vectors_is_aggregate_of_built_field_counts() {
1719        let mut meta = IndexMetadata::new(test_schema());
1720        meta.init_field(7, VectorIndexType::IvfTq);
1721        meta.init_field(3, VectorIndexType::IvfTq);
1722
1723        // Build in reverse field-id order to ensure the result is not tied to
1724        // HashMap or training iteration order.
1725        meta.mark_field_built(7, 400, 20, "field_7_centroids.bin".to_string(), None);
1726        assert_eq!(meta.total_vectors, 400);
1727        meta.mark_field_built(3, 250, 15, "field_3_centroids.bin".to_string(), None);
1728        assert_eq!(meta.total_vectors, 650);
1729
1730        // Rebuilding a field replaces its contribution; it does not add a
1731        // duplicate training snapshot.
1732        meta.mark_field_built(7, 425, 20, "field_7_centroids.bin".to_string(), None);
1733        assert_eq!(meta.total_vectors, 675);
1734    }
1735
1736    #[test]
1737    fn test_should_build_field() {
1738        let mut meta = IndexMetadata::new(test_schema());
1739        meta.init_field(0, VectorIndexType::IvfTq);
1740
1741        // Below threshold
1742        meta.total_vectors = 500;
1743        assert!(!meta.should_build_field(0, 1000));
1744
1745        // Above threshold
1746        meta.total_vectors = 1500;
1747        assert!(meta.should_build_field(0, 1000));
1748
1749        // Already built - should not build again
1750        meta.mark_field_built(0, 1500, 256, "centroids.bin".to_string(), None);
1751        assert!(!meta.should_build_field(0, 1000));
1752    }
1753
1754    #[test]
1755    fn test_serialization() {
1756        let mut meta = IndexMetadata::new(test_schema());
1757        meta.add_segment("seg1".to_string(), 100);
1758        meta.init_field(0, VectorIndexType::IvfTq);
1759        meta.total_vectors = 5000;
1760
1761        let json = serde_json::to_string_pretty(&meta).unwrap();
1762        let loaded: IndexMetadata = serde_json::from_str(&json).unwrap();
1763
1764        assert_eq!(loaded.segment_ids().len(), meta.segment_ids().len());
1765        assert_eq!(loaded.segment_doc_count("seg1"), Some(100));
1766        assert_eq!(loaded.total_vectors, meta.total_vectors);
1767        assert!(loaded.vector_fields.contains_key(&0));
1768    }
1769
1770    #[test]
1771    fn old_metadata_defaults_the_bp_retry_counter() {
1772        let mut meta = IndexMetadata::new(test_schema());
1773        meta.add_segment("legacy".to_string(), 10);
1774        let mut json = serde_json::to_value(&meta).unwrap();
1775        json["segment_metas"]["legacy"]
1776            .as_object_mut()
1777            .unwrap()
1778            .remove("bp_unconverged_passes");
1779
1780        let loaded: IndexMetadata = serde_json::from_value(json).unwrap();
1781        assert_eq!(loaded.segment_metas["legacy"].bp_unconverged_passes, 0);
1782    }
1783
1784    #[test]
1785    fn test_merged_segment_lineage() {
1786        let mut meta = IndexMetadata::new(test_schema());
1787        meta.add_segment("a".to_string(), 50);
1788        meta.add_segment("b".to_string(), 75);
1789
1790        // Fresh segments: gen=0, no ancestors
1791        assert_eq!(meta.segment_metas["a"].generation, 0);
1792        assert!(meta.segment_metas["a"].ancestors.is_empty());
1793
1794        // Merge a+b → c
1795        meta.add_merged_segment(
1796            "c".to_string(),
1797            125,
1798            vec!["a".to_string(), "b".to_string()],
1799            1,
1800            false,
1801            true,
1802        );
1803        assert_eq!(meta.segment_metas["c"].generation, 1);
1804        assert_eq!(meta.segment_metas["c"].ancestors, vec!["a", "b"]);
1805        assert_eq!(meta.segment_doc_count("c"), Some(125));
1806
1807        // Merge c+d → e (gen should be 2)
1808        meta.add_segment("d".to_string(), 30);
1809        meta.add_merged_segment(
1810            "e".to_string(),
1811            155,
1812            vec!["c".to_string(), "d".to_string()],
1813            2,
1814            false,
1815            true,
1816        );
1817        assert_eq!(meta.segment_metas["e"].generation, 2);
1818    }
1819}