Skip to main content

summa_core/segment/
mod.rs

1//! Immutable index segments and their lifecycle building blocks.
2//!
3//! Builders write a complete segment, readers expose its text, stored-field,
4//! and vector data, and native mergers combine committed segments. Publication
5//! and deletion are coordinated at the index layer by
6//! [`crate::merge::SegmentManager`].
7
8pub(crate) mod ann_build;
9mod ann_disk;
10pub use ann_disk::AnnHealth;
11pub(crate) mod bmp_adaptive;
12pub(crate) mod bmp_forward;
13pub(crate) mod bmp_grid;
14#[cfg(any(feature = "native", feature = "wasm"))]
15mod builder;
16pub mod chunk_map;
17pub(crate) mod deletion;
18pub(crate) mod norms;
19pub use deletion::DeletionMeta;
20pub(crate) mod format;
21pub(crate) mod logical_address;
22#[cfg(feature = "native")]
23mod merger;
24pub(crate) mod reader;
25#[cfg(feature = "native")]
26pub(crate) mod reorder;
27#[cfg(feature = "native")]
28pub(crate) mod row_map;
29pub(crate) mod seismic;
30#[cfg(any(feature = "native", feature = "wasm"))]
31mod sparse_partitions;
32mod store;
33#[cfg(feature = "native")]
34pub(crate) mod text_reorder;
35#[cfg(feature = "native")]
36mod tracker;
37mod types;
38mod vector_data;
39mod vector_locations;
40
41#[cfg(test)]
42pub(crate) use builder::graph_bisection::{
43    build_forward_index_from_blocks, build_forward_index_from_bmps, graph_bisection,
44};
45#[cfg(any(feature = "native", feature = "wasm"))]
46pub(crate) use builder::validate_vector_value_counts;
47#[cfg(any(feature = "native", feature = "wasm"))]
48pub use builder::{
49    BpBudget, MemoryBreakdown, SegmentBuilder, SegmentBuilderConfig, SegmentBuilderStats,
50};
51#[cfg(feature = "native")]
52pub(crate) use merger::block_in_place_if_multithread;
53#[cfg(feature = "native")]
54pub use merger::{MergeStats, SegmentMerger, delete_segment};
55pub(crate) use reader::BmpIndex;
56pub(crate) use reader::bmp::BMP_SUPERBLOCK_SIZE;
57pub(crate) use reader::combine_ordinal_results;
58#[cfg(feature = "native")]
59pub mod pin;
60pub use reader::{
61    BmpDimStats, DensePlanCache, SegmentReader, SparseIndex, VectorIndex, VectorOrdinals,
62    VectorSearchResult,
63};
64pub use store::*;
65#[cfg(feature = "native")]
66pub use tracker::{PublishedIndexGeneration, SegmentSnapshot, SegmentTracker};
67pub use types::{
68    FieldStats, ScannTrainedArtifactBytes, SegmentFiles, SegmentId, SegmentMeta, SeismicStats,
69    TrainedVectorStructures,
70};
71pub use vector_data::{FlatVectorData, LazyFlatVectorData, dequantize_raw};
72
73/// Write adapter that tracks bytes written.
74///
75/// Concrete type so it works with generic `serialize<W: Write>` functions
76/// (unlike `dyn StreamingWriter` which isn't `Sized`).
77#[cfg(any(feature = "native", feature = "wasm"))]
78pub(crate) struct OffsetWriter {
79    inner: Box<dyn crate::directories::StreamingWriter>,
80    offset: u64,
81}
82
83#[cfg(any(feature = "native", feature = "wasm"))]
84impl OffsetWriter {
85    pub(crate) fn new(inner: Box<dyn crate::directories::StreamingWriter>) -> Self {
86        Self { inner, offset: 0 }
87    }
88
89    /// Current write position (total bytes written so far).
90    pub(crate) fn offset(&self) -> u64 {
91        self.offset
92    }
93
94    /// Finalize the underlying streaming writer.
95    pub(crate) fn finish(self) -> std::io::Result<()> {
96        self.inner.finish()
97    }
98
99    /// Copy a local source range at the current output position using the
100    /// writer backend's kernel-assisted path when available.
101    #[cfg(feature = "native")]
102    pub(crate) fn copy_from_file_range(
103        &mut self,
104        source: &std::fs::File,
105        source_offset: &mut u64,
106        len: usize,
107    ) -> std::io::Result<usize> {
108        let copied = self
109            .inner
110            .copy_from_file_range(source, source_offset, len)?;
111        self.offset = self
112            .offset
113            .checked_add(copied as u64)
114            .ok_or_else(|| std::io::Error::other("offset-writer byte count overflow"))?;
115        Ok(copied)
116    }
117}
118
119#[cfg(any(feature = "native", feature = "wasm"))]
120impl std::io::Write for OffsetWriter {
121    fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
122        let n = self.inner.write(buf)?;
123        self.offset += n as u64;
124        Ok(n)
125    }
126
127    fn flush(&mut self) -> std::io::Result<()> {
128        self.inner.flush()
129    }
130}
131
132#[cfg(test)]
133#[cfg(feature = "native")]
134mod tests {
135    use super::*;
136    use crate::directories::RamDirectory;
137    use crate::dsl::SchemaBuilder;
138    use std::sync::Arc;
139
140    #[tokio::test]
141    async fn test_async_segment_reader() {
142        let mut schema_builder = SchemaBuilder::default();
143        let title = schema_builder.add_text_field("title", true, true);
144        let schema = Arc::new(schema_builder.build());
145
146        let dir = RamDirectory::new();
147        let segment_id = SegmentId::new();
148
149        // Build segment using sync builder
150        let config = SegmentBuilderConfig::default();
151        let mut builder = SegmentBuilder::new(Arc::clone(&schema), config).unwrap();
152
153        let mut doc = crate::dsl::Document::new();
154        doc.add_text(title, "Hello World");
155        builder.add_document(doc).unwrap();
156
157        let mut doc = crate::dsl::Document::new();
158        doc.add_text(title, "Goodbye World");
159        builder.add_document(doc).unwrap();
160
161        builder.build(&dir, segment_id, None).await.unwrap();
162
163        // Open with async reader
164        let reader = SegmentReader::open(&dir, segment_id, schema.clone(), 16)
165            .await
166            .unwrap();
167
168        assert_eq!(reader.num_docs(), 2);
169
170        // Test postings lookup
171        let postings = reader.get_postings(title, b"hello").await.unwrap();
172        assert!(postings.is_some());
173        assert_eq!(postings.unwrap().doc_count(), 1);
174
175        let postings = reader.get_postings(title, b"world").await.unwrap();
176        assert!(postings.is_some());
177        assert_eq!(postings.unwrap().doc_count(), 2);
178
179        // Test document retrieval
180        let doc = reader.doc(0).await.unwrap().unwrap();
181        assert_eq!(doc.get_first(title).unwrap().as_text(), Some("Hello World"));
182    }
183
184    #[tokio::test]
185    async fn test_dense_vector_ordinal_tracking() {
186        use crate::query::MultiValueCombiner;
187
188        let mut schema_builder = SchemaBuilder::default();
189        // Use simple add method - defaults to Flat index
190        let embedding = schema_builder.add_dense_vector_field("embedding", 4, true, true);
191        let schema = Arc::new(schema_builder.build());
192
193        let dir = RamDirectory::new();
194        let segment_id = SegmentId::new();
195
196        let config = SegmentBuilderConfig::default();
197        let mut builder = SegmentBuilder::new(Arc::clone(&schema), config).unwrap();
198
199        // Doc 0: single vector
200        let mut doc = crate::dsl::Document::new();
201        doc.add_dense_vector(embedding, vec![1.0, 0.0, 0.0, 0.0]);
202        builder.add_document(doc).unwrap();
203
204        // Doc 1: multi-valued vectors (2 vectors)
205        let mut doc = crate::dsl::Document::new();
206        doc.add_dense_vector(embedding, vec![0.0, 1.0, 0.0, 0.0]);
207        doc.add_dense_vector(embedding, vec![0.0, 0.0, 1.0, 0.0]);
208        builder.add_document(doc).unwrap();
209
210        // Doc 2: single vector
211        let mut doc = crate::dsl::Document::new();
212        doc.add_dense_vector(embedding, vec![0.0, 0.0, 0.0, 1.0]);
213        builder.add_document(doc).unwrap();
214
215        builder.build(&dir, segment_id, None).await.unwrap();
216
217        let reader = SegmentReader::open(&dir, segment_id, schema.clone(), 16)
218            .await
219            .unwrap();
220
221        // Query close to doc 1's first vector
222        let query = vec![0.0, 0.9, 0.1, 0.0];
223        let results = reader
224            .search_dense_vector(embedding, &query, 10, 0, 1.0, MultiValueCombiner::Max)
225            .await
226            .unwrap();
227
228        // Doc 1 should be in results with ordinal tracking
229        let doc1_result = results.iter().find(|r| r.doc_id == 1);
230        assert!(doc1_result.is_some(), "Doc 1 should be in results");
231
232        let doc1 = doc1_result.unwrap();
233        // Should have 2 ordinals (0 and 1) for the two vectors
234        assert!(
235            doc1.ordinals.len() <= 2,
236            "Doc 1 should have at most 2 ordinals, got {}",
237            doc1.ordinals.len()
238        );
239
240        // Check ordinals are valid (0 or 1)
241        for (ordinal, _score) in &doc1.ordinals {
242            assert!(*ordinal <= 1, "Ordinal should be 0 or 1, got {}", ordinal);
243        }
244    }
245
246    #[tokio::test]
247    async fn test_sparse_vector_ordinal_tracking() {
248        use crate::query::MultiValueCombiner;
249
250        let mut schema_builder = SchemaBuilder::default();
251        let sparse = schema_builder.add_sparse_vector_field("sparse", true, true);
252        let schema = Arc::new(schema_builder.build());
253
254        let dir = RamDirectory::new();
255        let segment_id = SegmentId::new();
256
257        let config = SegmentBuilderConfig::default();
258        let mut builder = SegmentBuilder::new(Arc::clone(&schema), config).unwrap();
259
260        // Doc 0: single sparse vector
261        let mut doc = crate::dsl::Document::new();
262        doc.add_sparse_vector(sparse, vec![(0, 1.0), (1, 0.5)]);
263        builder.add_document(doc).unwrap();
264
265        // Doc 1: multi-valued sparse vectors (2 vectors)
266        let mut doc = crate::dsl::Document::new();
267        doc.add_sparse_vector(sparse, vec![(0, 0.8), (2, 0.3)]);
268        doc.add_sparse_vector(sparse, vec![(1, 0.9), (3, 0.4)]);
269        builder.add_document(doc).unwrap();
270
271        // Doc 2: single sparse vector
272        let mut doc = crate::dsl::Document::new();
273        doc.add_sparse_vector(sparse, vec![(2, 1.0), (3, 0.5)]);
274        builder.add_document(doc).unwrap();
275
276        builder.build(&dir, segment_id, None).await.unwrap();
277
278        let reader = SegmentReader::open(&dir, segment_id, schema.clone(), 16)
279            .await
280            .unwrap();
281
282        // Query matching dimension 0 via SparseVectorQuery
283        let query = crate::query::SparseVectorQuery::new(sparse, vec![(0, 1.0)])
284            .with_combiner(MultiValueCombiner::Sum);
285        let mut collector = crate::query::TopKCollector::new(10);
286        crate::query::collect_segment(&reader, &query, &mut collector)
287            .await
288            .unwrap();
289        let top_docs = collector.into_sorted_results();
290
291        // Both doc 0 and doc 1 have dimension 0
292        assert!(top_docs.len() >= 2, "Should have at least 2 results");
293    }
294}