Skip to main content

uqa_storage/key_value/
inverted_index.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Clustered inverted-index adapter over an ordered key/value store.
8
9use super::codec::{
10    blob_to_positions, decode_u64_value, doc_length_doc_prefix, doc_length_key,
11    doc_length_key_prefix, field_stats_key, field_stats_key_prefix, key_with_tag, other_error,
12    posting_cluster_positions_key, posting_cluster_positions_key_prefix,
13    posting_cluster_score_field_prefix, posting_cluster_score_key,
14    posting_cluster_score_key_prefix, posting_cluster_score_term_prefix,
15    posting_document_doc_prefix, posting_document_key, posting_document_key_prefix, read_str,
16    read_u64, reverse_posting_key, single_str_key, string_value, u64_value, usize_to_u64,
17};
18use super::{
19    Analyzer, AnalyzerPhase, Arc, BTreeMap, BTreeSet, DocId, FieldName, IndexStats, InvertedIndex,
20    KeyValueBatch, KeyValueStore, Payload, PostingEntry, PostingList, StorageBackendResult,
21    TAG_METADATA, TAG_POSTING, TAG_REVERSE_POSTING,
22};
23use crate::clustered_postings::{
24    cluster_id, decode_all_scores, decode_cluster, decode_terms, encode_cluster, encode_terms,
25    score_count, ClusterPosting, ClusteredPostingCursor, EncodedScoreCluster,
26    MaterializedPostingCursor,
27};
28use crate::PostingCursor;
29
30mod migration;
31mod mutation;
32
33use migration::{migrate_legacy_forward_postings, migrate_legacy_reverse_postings};
34use mutation::{accumulate_field_changes, merge_cluster_changes};
35
36const FORMAT_METADATA_KEY: &str = "inverted_index_format";
37const CLUSTERED_FORMAT_NAME: &str = "clustered-v1";
38const MIGRATION_PAGE_SIZE: usize = 1_024;
39
40/// Inverted index implemented over [`KeyValueStore`].
41#[derive(Clone)]
42pub struct KeyValueInvertedIndex {
43    store: Arc<dyn KeyValueStore>,
44    table: String,
45    analyzer: Analyzer,
46    index_field_analyzers: BTreeMap<FieldName, Analyzer>,
47    search_field_analyzers: BTreeMap<FieldName, Analyzer>,
48}
49
50type KeyValueStagedPosting = (FieldName, String, Vec<u32>);
51type KeyValueAnalyzedFields = (BTreeMap<FieldName, u64>, Vec<KeyValueStagedPosting>);
52type ClusterKey = (FieldName, String, u64);
53type PostingChange = Option<(u64, Vec<u32>)>;
54type KeyValueStagedDocuments = BTreeMap<DocId, KeyValueAnalyzedFields>;
55type KeyValueClusterChanges = BTreeMap<ClusterKey, BTreeMap<DocId, PostingChange>>;
56type KeyValueFieldChanges = BTreeMap<FieldName, (u64, u64)>;
57type KeyValueMergedClusters = Vec<(ClusterKey, Vec<ClusterPosting>)>;
58
59impl KeyValueInvertedIndex {
60    pub fn new(
61        store: Arc<dyn KeyValueStore>,
62        table: impl Into<String>,
63        analyzer: Analyzer,
64    ) -> Self {
65        Self {
66            store,
67            table: table.into(),
68            analyzer,
69            index_field_analyzers: BTreeMap::new(),
70            search_field_analyzers: BTreeMap::new(),
71        }
72    }
73
74    pub(crate) fn migrate_legacy_storage(store: &dyn KeyValueStore) -> StorageBackendResult<()> {
75        let marker = single_str_key(TAG_METADATA, FORMAT_METADATA_KEY)?;
76        if let Some(format) = store.get(&marker)? {
77            if format == CLUSTERED_FORMAT_NAME.as_bytes() {
78                return Ok(());
79            }
80            return Err(other_error(format!(
81                "unsupported KeyValue inverted-index format `{}`",
82                String::from_utf8_lossy(&format)
83            )));
84        }
85        if store.in_transaction() {
86            return Err(other_error(
87                "cannot migrate KeyValue postings inside an active transaction",
88            ));
89        }
90
91        store.begin_transaction()?;
92        let migration = Self::migrate_legacy_storage_in_transaction(store, &marker);
93        match migration {
94            Ok(()) => store.commit_transaction(),
95            Err(error) => match store.rollback_transaction() {
96                Ok(()) => Err(error),
97                Err(rollback) => Err(other_error(format!(
98                    "{error}; KeyValue posting migration rollback also failed: {rollback}"
99                ))),
100            },
101        }
102    }
103
104    fn migrate_legacy_storage_in_transaction(
105        store: &dyn KeyValueStore,
106        marker: &[u8],
107    ) -> StorageBackendResult<()> {
108        let posting_count = migrate_legacy_forward_postings(store)?;
109        let reverse_count = migrate_legacy_reverse_postings(store)?;
110        if posting_count != reverse_count {
111            return Err(other_error(format!(
112                "cannot migrate inconsistent KeyValue postings: {posting_count} forward rows and {reverse_count} reverse rows"
113            )));
114        }
115        store.put(marker, &string_value(CLUSTERED_FORMAT_NAME))
116    }
117
118    fn old_doc_lengths(&self, doc_id: DocId) -> StorageBackendResult<BTreeMap<FieldName, u64>> {
119        let mut out = BTreeMap::new();
120        for (key, value) in self
121            .store
122            .scan_prefix(&doc_length_doc_prefix(&self.table, doc_id)?)?
123        {
124            let mut offset = 1;
125            let _table = read_str(&key, &mut offset)?;
126            let _doc_id = read_u64(&key, &mut offset)?;
127            let field = read_str(&key, &mut offset)?;
128            out.insert(field, decode_u64_value(&value)?);
129        }
130        Ok(out)
131    }
132
133    fn old_terms(&self, doc_id: DocId) -> StorageBackendResult<BTreeMap<FieldName, Vec<String>>> {
134        let mut out = BTreeMap::new();
135        for (key, value) in self
136            .store
137            .scan_prefix(&posting_document_doc_prefix(&self.table, doc_id)?)?
138        {
139            let mut offset = 1;
140            let _table = read_str(&key, &mut offset)?;
141            let _doc_id = read_u64(&key, &mut offset)?;
142            let field = read_str(&key, &mut offset)?;
143            out.insert(field, decode_terms(&value)?);
144        }
145        Ok(out)
146    }
147
148    fn analyze_fields(
149        &self,
150        fields: BTreeMap<FieldName, String>,
151    ) -> StorageBackendResult<KeyValueAnalyzedFields> {
152        let mut lengths = BTreeMap::new();
153        let mut postings = Vec::new();
154        for (field, text) in fields {
155            let analyzer = self
156                .index_field_analyzers
157                .get(&field)
158                .unwrap_or(&self.analyzer);
159            let tokens = analyzer.analyze(&text)?;
160            let token_count = usize_to_u64(tokens.len(), "document token count")?;
161            crate::inverted_index::validate_token_position_count(token_count)?;
162            lengths.insert(field.clone(), token_count);
163            let mut term_positions: BTreeMap<String, Vec<u32>> = BTreeMap::new();
164            for (position, token) in tokens.into_iter().enumerate() {
165                term_positions.entry(token).or_default().push(
166                    u32::try_from(position)
167                        .map_err(|_| other_error("token position exceeds u32 index format"))?,
168                );
169            }
170            for (term, mut positions) in term_positions {
171                positions.sort_unstable();
172                positions.dedup();
173                postings.push((field.clone(), term, positions));
174            }
175        }
176        Ok((lengths, postings))
177    }
178
179    fn set_total_length(
180        batch: &mut dyn KeyValueBatch,
181        table: &str,
182        field: &str,
183        value: u64,
184    ) -> StorageBackendResult<()> {
185        let key = field_stats_key(table, field)?;
186        if value == 0 {
187            batch.delete(&key)
188        } else {
189            batch.put(&key, &u64_value(value))
190        }
191    }
192
193    fn load_cluster(
194        &self,
195        field: &str,
196        term: &str,
197        posting_cluster: u64,
198    ) -> StorageBackendResult<Vec<ClusterPosting>> {
199        let score = self.store.get(&posting_cluster_score_key(
200            &self.table,
201            field,
202            term,
203            posting_cluster,
204        )?)?;
205        let positions = self.store.get(&posting_cluster_positions_key(
206            &self.table,
207            field,
208            term,
209            posting_cluster,
210        )?)?;
211        match (score, positions) {
212            (None, None) => Ok(Vec::new()),
213            (Some(score), Some(positions)) => decode_cluster(posting_cluster, &score, &positions),
214            _ => Err(other_error(
215                "clustered posting score and positions values disagree",
216            )),
217        }
218    }
219
220    fn stage_cluster_changes(
221        &self,
222        doc_id: DocId,
223        old_terms: &BTreeMap<FieldName, Vec<String>>,
224        lengths: &BTreeMap<FieldName, u64>,
225        postings: &[KeyValueStagedPosting],
226    ) -> StorageBackendResult<Vec<(ClusterKey, Vec<ClusterPosting>)>> {
227        let mut changes = BTreeMap::<(FieldName, String), PostingChange>::new();
228        for (field, terms) in old_terms {
229            for term in terms {
230                changes.insert((field.clone(), term.clone()), None);
231            }
232        }
233        for (field, term, positions) in postings {
234            changes.insert(
235                (field.clone(), term.clone()),
236                Some((lengths[field], positions.clone())),
237            );
238        }
239
240        let posting_cluster = cluster_id(doc_id);
241        let mut output = Vec::with_capacity(changes.len());
242        for ((field, term), replacement) in changes {
243            let mut entries = self.load_cluster(&field, &term, posting_cluster)?;
244            if let Ok(position) = entries.binary_search_by_key(&doc_id, |entry| entry.doc_id) {
245                entries.remove(position);
246            }
247            if let Some((doc_length, positions)) = replacement {
248                let position = entries.partition_point(|entry| entry.doc_id < doc_id);
249                entries.insert(
250                    position,
251                    ClusterPosting {
252                        doc_id,
253                        term_freq: positions.len() as u64,
254                        doc_length,
255                        positions,
256                    },
257                );
258            }
259            output.push(((field, term, posting_cluster), entries));
260        }
261        Ok(output)
262    }
263
264    fn apply_cluster_changes(
265        batch: &mut dyn KeyValueBatch,
266        table: &str,
267        changes: Vec<(ClusterKey, Vec<ClusterPosting>)>,
268    ) -> StorageBackendResult<()> {
269        for ((field, term, posting_cluster), entries) in changes {
270            let score_key = posting_cluster_score_key(table, &field, &term, posting_cluster)?;
271            let positions_key =
272                posting_cluster_positions_key(table, &field, &term, posting_cluster)?;
273            if entries.is_empty() {
274                batch.delete(&score_key)?;
275                batch.delete(&positions_key)?;
276            } else {
277                let (score, positions) = encode_cluster(&entries)?;
278                batch.put(&score_key, &score)?;
279                batch.put(&positions_key, &positions)?;
280            }
281        }
282        Ok(())
283    }
284
285    fn collect_batch_changes(
286        &self,
287        staged_documents: &KeyValueStagedDocuments,
288    ) -> StorageBackendResult<(KeyValueClusterChanges, KeyValueFieldChanges)> {
289        let mut cluster_changes = KeyValueClusterChanges::new();
290        let mut field_changes = KeyValueFieldChanges::new();
291        for (doc_id, (new_lengths, new_postings)) in staged_documents {
292            let old_lengths = self.old_doc_lengths(*doc_id)?;
293            let old_terms = self.old_terms(*doc_id)?;
294            let posting_cluster = cluster_id(*doc_id);
295            for (field, terms) in old_terms {
296                for term in terms {
297                    cluster_changes
298                        .entry((field.clone(), term, posting_cluster))
299                        .or_default()
300                        .insert(*doc_id, None);
301                }
302            }
303            for (field, term, positions) in new_postings {
304                cluster_changes
305                    .entry((field.clone(), term.clone(), posting_cluster))
306                    .or_default()
307                    .insert(*doc_id, Some((new_lengths[field], positions.clone())));
308            }
309            accumulate_field_changes(&mut field_changes, &old_lengths, new_lengths)?;
310        }
311        Ok((cluster_changes, field_changes))
312    }
313
314    fn plan_batch_totals(
315        &self,
316        field_changes: KeyValueFieldChanges,
317    ) -> StorageBackendResult<Vec<(FieldName, u64)>> {
318        let mut totals = Vec::with_capacity(field_changes.len());
319        for (field, (old_total, new_total)) in field_changes {
320            let base = self
321                .store
322                .get(&field_stats_key(&self.table, &field)?)?
323                .map(|value| decode_u64_value(&value))
324                .transpose()?
325                .unwrap_or(0);
326            let total = base
327                .checked_sub(old_total)
328                .ok_or_else(|| other_error("stored field length is smaller than batch length"))?
329                .checked_add(new_total)
330                .ok_or_else(|| other_error("total field length overflow"))?;
331            totals.push((field, total));
332        }
333        Ok(totals)
334    }
335
336    fn merge_batch_clusters(
337        &self,
338        cluster_changes: KeyValueClusterChanges,
339    ) -> StorageBackendResult<KeyValueMergedClusters> {
340        let mut merged = Vec::with_capacity(cluster_changes.len());
341        for ((field, term, posting_cluster), changes) in cluster_changes {
342            let entries = self.load_cluster(&field, &term, posting_cluster)?;
343            merged.push((
344                (field, term, posting_cluster),
345                merge_cluster_changes(entries, changes),
346            ));
347        }
348        Ok(merged)
349    }
350
351    fn write_batch_documents(
352        &self,
353        batch: &mut dyn KeyValueBatch,
354        staged_documents: KeyValueStagedDocuments,
355    ) -> StorageBackendResult<()> {
356        for (doc_id, (lengths, postings)) in staged_documents {
357            batch.delete_prefix(&posting_document_doc_prefix(&self.table, doc_id)?)?;
358            batch.delete_prefix(&doc_length_doc_prefix(&self.table, doc_id)?)?;
359            let mut terms_by_field = BTreeMap::<FieldName, Vec<String>>::new();
360            for (field, term, _) in postings {
361                terms_by_field.entry(field).or_default().push(term);
362            }
363            for (field, length) in lengths {
364                batch.put(
365                    &doc_length_key(&self.table, doc_id, &field)?,
366                    &u64_value(length),
367                )?;
368                batch.put(
369                    &posting_document_key(&self.table, doc_id, &field)?,
370                    &encode_terms(terms_by_field.get(&field).map_or(&[], Vec::as_slice))?,
371                )?;
372            }
373        }
374        Ok(())
375    }
376
377    fn add_documents(
378        &self,
379        documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
380    ) -> StorageBackendResult<()> {
381        let mut staged_documents = BTreeMap::new();
382        for (doc_id, fields) in documents {
383            staged_documents.insert(doc_id, self.analyze_fields(fields)?);
384        }
385        if staged_documents.is_empty() {
386            return Ok(());
387        }
388
389        let (cluster_changes, field_changes) = self.collect_batch_changes(&staged_documents)?;
390        let totals = self.plan_batch_totals(field_changes)?;
391        let merged_clusters = self.merge_batch_clusters(cluster_changes)?;
392        let mut batch = self.store.batch();
393        Self::apply_cluster_changes(batch.as_mut(), &self.table, merged_clusters)?;
394        for (field, total) in totals {
395            Self::set_total_length(batch.as_mut(), &self.table, &field, total)?;
396        }
397        self.write_batch_documents(batch.as_mut(), staged_documents)?;
398        batch.commit()
399    }
400
401    fn cursor_for_term(
402        &self,
403        field: &str,
404        term: &str,
405    ) -> StorageBackendResult<Box<dyn PostingCursor>> {
406        let mut clusters = Vec::new();
407        for (key, bytes) in self.store.scan_prefix(&posting_cluster_score_term_prefix(
408            &self.table,
409            field,
410            term,
411        )?)? {
412            let mut offset = 1;
413            let _table = read_str(&key, &mut offset)?;
414            let _field = read_str(&key, &mut offset)?;
415            let _term = read_str(&key, &mut offset)?;
416            let posting_cluster = read_u64(&key, &mut offset)?;
417            if offset != key.len() {
418                return Err(other_error("invalid clustered posting score key"));
419            }
420            clusters.push(EncodedScoreCluster {
421                cluster_id: posting_cluster,
422                bytes,
423            });
424        }
425        if clusters.is_empty() {
426            return Ok(Box::new(MaterializedPostingCursor::new(Vec::new())?));
427        }
428        Ok(Box::new(ClusteredPostingCursor::new(clusters)?))
429    }
430}
431
432impl InvertedIndex for KeyValueInvertedIndex {
433    fn analyzer(&self) -> &Analyzer {
434        &self.analyzer
435    }
436
437    fn add_document(
438        &mut self,
439        doc_id: DocId,
440        fields: BTreeMap<FieldName, String>,
441    ) -> StorageBackendResult<()> {
442        let old_lengths = self.old_doc_lengths(doc_id)?;
443        let old_terms = self.old_terms(doc_id)?;
444        let (new_lengths, new_postings) = self.analyze_fields(fields)?;
445        let cluster_changes =
446            self.stage_cluster_changes(doc_id, &old_terms, &new_lengths, &new_postings)?;
447
448        let mut fields_to_update = BTreeSet::new();
449        fields_to_update.extend(old_lengths.keys().cloned());
450        fields_to_update.extend(new_lengths.keys().cloned());
451        let mut totals = Vec::with_capacity(fields_to_update.len());
452        for field in fields_to_update {
453            let base = self
454                .store
455                .get(&field_stats_key(&self.table, &field)?)?
456                .map(|value| decode_u64_value(&value))
457                .transpose()?
458                .unwrap_or(0);
459            let old = old_lengths.get(&field).copied().unwrap_or(0);
460            let new = new_lengths.get(&field).copied().unwrap_or(0);
461            let total = base
462                .checked_sub(old)
463                .ok_or_else(|| other_error("stored field length is smaller than document length"))?
464                .checked_add(new)
465                .ok_or_else(|| other_error("total field length overflow"))?;
466            totals.push((field, total));
467        }
468
469        let mut terms_by_field = BTreeMap::<FieldName, Vec<String>>::new();
470        for (field, term, _) in &new_postings {
471            terms_by_field
472                .entry(field.clone())
473                .or_default()
474                .push(term.clone());
475        }
476        for field in new_lengths.keys() {
477            terms_by_field.entry(field.clone()).or_default();
478        }
479
480        let mut batch = self.store.batch();
481        Self::apply_cluster_changes(batch.as_mut(), &self.table, cluster_changes)?;
482        batch.delete_prefix(&posting_document_doc_prefix(&self.table, doc_id)?)?;
483        for field in old_lengths.keys() {
484            batch.delete(&doc_length_key(&self.table, doc_id, field)?)?;
485        }
486        for (field, total) in totals {
487            Self::set_total_length(batch.as_mut(), &self.table, &field, total)?;
488        }
489        for (field, length) in &new_lengths {
490            batch.put(
491                &doc_length_key(&self.table, doc_id, field)?,
492                &u64_value(*length),
493            )?;
494        }
495        for (field, terms) in terms_by_field {
496            batch.put(
497                &posting_document_key(&self.table, doc_id, &field)?,
498                &encode_terms(&terms)?,
499            )?;
500        }
501        batch.commit()
502    }
503
504    fn try_add_documents(
505        &mut self,
506        documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
507    ) -> StorageBackendResult<()> {
508        self.add_documents(documents)
509    }
510
511    fn remove_document(&mut self, doc_id: DocId) -> StorageBackendResult<()> {
512        let old_lengths = self.old_doc_lengths(doc_id)?;
513        let old_terms = self.old_terms(doc_id)?;
514        let cluster_changes =
515            self.stage_cluster_changes(doc_id, &old_terms, &BTreeMap::new(), &[])?;
516        let mut totals = Vec::with_capacity(old_lengths.len());
517        for (field, length) in &old_lengths {
518            let base = self
519                .store
520                .get(&field_stats_key(&self.table, field)?)?
521                .map(|value| decode_u64_value(&value))
522                .transpose()?
523                .unwrap_or(0);
524            totals.push((
525                field.clone(),
526                base.checked_sub(*length).ok_or_else(|| {
527                    other_error("stored field length is smaller than removed document length")
528                })?,
529            ));
530        }
531
532        let mut batch = self.store.batch();
533        Self::apply_cluster_changes(batch.as_mut(), &self.table, cluster_changes)?;
534        batch.delete_prefix(&posting_document_doc_prefix(&self.table, doc_id)?)?;
535        for (field, total) in totals {
536            Self::set_total_length(batch.as_mut(), &self.table, &field, total)?;
537            batch.delete(&doc_length_key(&self.table, doc_id, &field)?)?;
538        }
539        batch.commit()
540    }
541
542    fn try_rebuild_documents(
543        &mut self,
544        documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
545    ) -> StorageBackendResult<()> {
546        let mut staged = BTreeMap::new();
547        for (doc_id, fields) in documents {
548            if !fields.is_empty() {
549                staged.insert(doc_id, self.analyze_fields(fields)?);
550            }
551        }
552        let mut totals = BTreeMap::<FieldName, u64>::new();
553        let mut clusters = BTreeMap::<ClusterKey, Vec<ClusterPosting>>::new();
554        for (doc_id, (lengths, postings)) in &staged {
555            for (field, length) in lengths {
556                let total = totals.entry(field.clone()).or_default();
557                *total = total
558                    .checked_add(*length)
559                    .ok_or_else(|| other_error("total field length overflow"))?;
560            }
561            for (field, term, positions) in postings {
562                clusters
563                    .entry((field.clone(), term.clone(), cluster_id(*doc_id)))
564                    .or_default()
565                    .push(ClusterPosting {
566                        doc_id: *doc_id,
567                        term_freq: positions.len() as u64,
568                        doc_length: lengths[field],
569                        positions: positions.clone(),
570                    });
571            }
572        }
573
574        let mut batch = self.store.batch();
575        batch.delete_prefix(&posting_cluster_score_key_prefix(&self.table)?)?;
576        batch.delete_prefix(&posting_cluster_positions_key_prefix(&self.table)?)?;
577        batch.delete_prefix(&posting_document_key_prefix(&self.table)?)?;
578        batch.delete_prefix(&doc_length_key_prefix(&self.table)?)?;
579        batch.delete_prefix(&field_stats_key_prefix(&self.table)?)?;
580        for (field, total) in totals {
581            Self::set_total_length(batch.as_mut(), &self.table, &field, total)?;
582        }
583        for ((field, term, posting_cluster), entries) in clusters {
584            let (score, positions) = encode_cluster(&entries)?;
585            batch.put(
586                &posting_cluster_score_key(&self.table, &field, &term, posting_cluster)?,
587                &score,
588            )?;
589            batch.put(
590                &posting_cluster_positions_key(&self.table, &field, &term, posting_cluster)?,
591                &positions,
592            )?;
593        }
594        for (doc_id, (lengths, postings)) in staged {
595            for (field, length) in lengths {
596                batch.put(
597                    &doc_length_key(&self.table, doc_id, &field)?,
598                    &u64_value(length),
599                )?;
600                let terms = postings
601                    .iter()
602                    .filter(|(posting_field, _, _)| posting_field == &field)
603                    .map(|(_, term, _)| term.clone())
604                    .collect::<Vec<_>>();
605                batch.put(
606                    &posting_document_key(&self.table, doc_id, &field)?,
607                    &encode_terms(&terms)?,
608                )?;
609            }
610        }
611        batch.commit()
612    }
613
614    fn clear(&mut self) -> StorageBackendResult<()> {
615        let mut batch = self.store.batch();
616        batch.delete_prefix(&posting_cluster_score_key_prefix(&self.table)?)?;
617        batch.delete_prefix(&posting_cluster_positions_key_prefix(&self.table)?)?;
618        batch.delete_prefix(&posting_document_key_prefix(&self.table)?)?;
619        batch.delete_prefix(&doc_length_key_prefix(&self.table)?)?;
620        batch.delete_prefix(&field_stats_key_prefix(&self.table)?)?;
621        batch.commit()
622    }
623
624    fn get_posting_list(&self, field: &str, term: &str) -> StorageBackendResult<PostingList> {
625        let mut entries = Vec::new();
626        for (key, score) in self.store.scan_prefix(&posting_cluster_score_term_prefix(
627            &self.table,
628            field,
629            term,
630        )?)? {
631            let mut offset = 1;
632            let _table = read_str(&key, &mut offset)?;
633            let _field = read_str(&key, &mut offset)?;
634            let _term = read_str(&key, &mut offset)?;
635            let posting_cluster = read_u64(&key, &mut offset)?;
636            let positions = self
637                .store
638                .get(&posting_cluster_positions_key(
639                    &self.table,
640                    field,
641                    term,
642                    posting_cluster,
643                )?)?
644                .ok_or_else(|| other_error("clustered posting positions value is missing"))?;
645            entries.extend(
646                decode_cluster(posting_cluster, &score, &positions)?
647                    .into_iter()
648                    .map(|entry| {
649                        PostingEntry::new(
650                            entry.doc_id,
651                            Payload {
652                                positions: entry.positions,
653                                score: 0.0,
654                                fields: BTreeMap::new(),
655                            },
656                        )
657                    }),
658            );
659        }
660        Ok(PostingList::from_sorted_unchecked(entries))
661    }
662
663    fn posting_cursor(
664        &self,
665        field: &str,
666        term: &str,
667    ) -> StorageBackendResult<Box<dyn PostingCursor>> {
668        self.cursor_for_term(field, term)
669    }
670
671    fn for_each_term_freq(
672        &self,
673        field: &str,
674        term: &str,
675        visit: &mut dyn FnMut(DocId, u64),
676    ) -> StorageBackendResult<()> {
677        let mut cursor = self.cursor_for_term(field, term)?;
678        while let Some(entry) = cursor.current() {
679            visit(entry.doc_id, entry.term_freq);
680            cursor.advance()?;
681        }
682        Ok(())
683    }
684
685    fn doc_freq(&self, field: &str, term: &str) -> StorageBackendResult<u64> {
686        self.store
687            .scan_prefix(&posting_cluster_score_term_prefix(
688                &self.table,
689                field,
690                term,
691            )?)?
692            .into_iter()
693            .try_fold(0_u64, |total, (_, score)| {
694                total
695                    .checked_add(score_count(&score)?)
696                    .ok_or_else(|| other_error("document frequency overflow"))
697            })
698    }
699
700    fn get_doc_length(&self, doc_id: DocId, field: &str) -> StorageBackendResult<u64> {
701        Ok(self
702            .store
703            .get(&doc_length_key(&self.table, doc_id, field)?)?
704            .map(|value| decode_u64_value(&value))
705            .transpose()?
706            .unwrap_or(0))
707    }
708
709    fn get_scoring_inputs_bulk(
710        &self,
711        doc_ids: &[DocId],
712        field: &str,
713        terms: &[String],
714    ) -> StorageBackendResult<Vec<(u64, Vec<u64>)>> {
715        let mut output = doc_ids
716            .iter()
717            .map(|doc_id| Ok((self.get_doc_length(*doc_id, field)?, vec![0; terms.len()])))
718            .collect::<StorageBackendResult<Vec<_>>>()?;
719        let mut positions = BTreeMap::<DocId, Vec<usize>>::new();
720        for (position, doc_id) in doc_ids.iter().copied().enumerate() {
721            positions.entry(doc_id).or_default().push(position);
722        }
723        for (term_index, term) in terms.iter().enumerate() {
724            let mut cursor = self.cursor_for_term(field, term)?;
725            while let Some(entry) = cursor.current() {
726                if let Some(output_positions) = positions.get(&entry.doc_id) {
727                    for position in output_positions {
728                        output[*position].0 = entry.doc_length;
729                        output[*position].1[term_index] = entry.term_freq;
730                    }
731                }
732                cursor.advance()?;
733            }
734        }
735        Ok(output)
736    }
737
738    fn get_term_freq(&self, doc_id: DocId, field: &str, term: &str) -> StorageBackendResult<u64> {
739        let posting_cluster = cluster_id(doc_id);
740        self.store
741            .get(&posting_cluster_score_key(
742                &self.table,
743                field,
744                term,
745                posting_cluster,
746            )?)?
747            .map_or(Ok(0), |score| {
748                let entries = decode_all_scores(posting_cluster, &score)?;
749                Ok(entries
750                    .binary_search_by_key(&doc_id, |entry| entry.doc_id)
751                    .ok()
752                    .map_or(0, |position| entries[position].term_freq))
753            })
754    }
755
756    fn doc_count(&self) -> StorageBackendResult<u64> {
757        let mut doc_ids = BTreeSet::new();
758        for (key, _) in self
759            .store
760            .scan_prefix(&doc_length_key_prefix(&self.table)?)?
761        {
762            let mut offset = 1;
763            let _table = read_str(&key, &mut offset)?;
764            doc_ids.insert(read_u64(&key, &mut offset)?);
765        }
766        usize_to_u64(doc_ids.len(), "document count")
767    }
768
769    fn total_field_length(&self, field: &str) -> StorageBackendResult<u64> {
770        Ok(self
771            .store
772            .get(&field_stats_key(&self.table, field)?)?
773            .map(|value| decode_u64_value(&value))
774            .transpose()?
775            .unwrap_or(0))
776    }
777
778    fn vocabulary_terms(&self, field: &str) -> StorageBackendResult<Vec<String>> {
779        let mut terms = BTreeSet::new();
780        for (key, _) in self
781            .store
782            .scan_prefix(&posting_cluster_score_field_prefix(&self.table, field)?)?
783        {
784            let mut offset = 1;
785            let _table = read_str(&key, &mut offset)?;
786            let _field = read_str(&key, &mut offset)?;
787            terms.insert(read_str(&key, &mut offset)?);
788        }
789        Ok(terms.into_iter().collect())
790    }
791
792    fn stats(&self) -> StorageBackendResult<IndexStats> {
793        let doc_count = self.doc_count()?;
794        let mut stats = IndexStats::default();
795        stats.total_docs = doc_count;
796        if doc_count > 0 {
797            let mut total = 0_u64;
798            for (_, value) in self
799                .store
800                .scan_prefix(&field_stats_key_prefix(&self.table)?)?
801            {
802                total = total
803                    .checked_add(decode_u64_value(&value)?)
804                    .ok_or_else(|| other_error("index total field length overflow"))?;
805            }
806            stats.avg_doc_length = total as f64 / doc_count as f64;
807        }
808        let mut counts = BTreeMap::<(String, String), u64>::new();
809        for (key, value) in self
810            .store
811            .scan_prefix(&posting_cluster_score_key_prefix(&self.table)?)?
812        {
813            let mut offset = 1;
814            let _table = read_str(&key, &mut offset)?;
815            let field = read_str(&key, &mut offset)?;
816            let term = read_str(&key, &mut offset)?;
817            let count = counts.entry((field, term)).or_default();
818            *count = count
819                .checked_add(score_count(&value)?)
820                .ok_or_else(|| other_error("index document frequency overflow"))?;
821        }
822        for ((field, term), document_frequency) in counts {
823            stats.set_doc_freq(field, term, document_frequency);
824        }
825        Ok(stats)
826    }
827
828    fn posting_count(&self, field: Option<&str>) -> StorageBackendResult<u64> {
829        let prefix = match field {
830            Some(field) => posting_cluster_score_field_prefix(&self.table, field)?,
831            None => posting_cluster_score_key_prefix(&self.table)?,
832        };
833        self.store
834            .scan_prefix(&prefix)?
835            .into_iter()
836            .try_fold(0_u64, |total, (_, value)| {
837                total
838                    .checked_add(score_count(&value)?)
839                    .ok_or_else(|| other_error("posting count overflow"))
840            })
841    }
842
843    fn doc_length_count(&self, field: Option<&str>) -> StorageBackendResult<u64> {
844        let mut count = 0_u64;
845        for (key, _) in self
846            .store
847            .scan_prefix(&doc_length_key_prefix(&self.table)?)?
848        {
849            let mut offset = 1;
850            let _table = read_str(&key, &mut offset)?;
851            let _doc_id = read_u64(&key, &mut offset)?;
852            let indexed_field = read_str(&key, &mut offset)?;
853            if field.is_none_or(|target| target == indexed_field) {
854                count = count
855                    .checked_add(1)
856                    .ok_or_else(|| other_error("document-length row count overflow"))?;
857            }
858        }
859        Ok(count)
860    }
861
862    fn term_count(&self, field: Option<&str>) -> StorageBackendResult<u64> {
863        let prefix = match field {
864            Some(field) => posting_cluster_score_field_prefix(&self.table, field)?,
865            None => posting_cluster_score_key_prefix(&self.table)?,
866        };
867        let mut terms = BTreeSet::new();
868        for (key, _) in self.store.scan_prefix(&prefix)? {
869            let mut offset = 1;
870            let _table = read_str(&key, &mut offset)?;
871            let current_field = read_str(&key, &mut offset)?;
872            let term = read_str(&key, &mut offset)?;
873            terms.insert((current_field, term));
874        }
875        usize_to_u64(terms.len(), "term count")
876    }
877
878    fn snapshot(&self) -> StorageBackendResult<Arc<dyn InvertedIndex>> {
879        Ok(Arc::new(self.clone()))
880    }
881
882    fn field_names(&self) -> StorageBackendResult<Vec<FieldName>> {
883        let mut fields = Vec::new();
884        for (key, _) in self
885            .store
886            .scan_prefix(&field_stats_key_prefix(&self.table)?)?
887        {
888            let mut offset = 1;
889            let _table = read_str(&key, &mut offset)?;
890            fields.push(read_str(&key, &mut offset)?);
891        }
892        Ok(fields)
893    }
894
895    fn set_field_analyzer(
896        &mut self,
897        field: &str,
898        analyzer: Analyzer,
899        phase: AnalyzerPhase,
900    ) -> Result<(), String> {
901        match phase {
902            AnalyzerPhase::Index => {
903                self.index_field_analyzers
904                    .insert(field.to_string(), analyzer);
905            }
906            AnalyzerPhase::Search => {
907                self.search_field_analyzers
908                    .insert(field.to_string(), analyzer);
909            }
910            AnalyzerPhase::Both => {
911                self.index_field_analyzers
912                    .insert(field.to_string(), analyzer.clone());
913                self.search_field_analyzers
914                    .insert(field.to_string(), analyzer);
915            }
916        }
917        Ok(())
918    }
919
920    fn remove_field_analyzers(&mut self, field: &str) -> Result<(), String> {
921        self.index_field_analyzers.remove(field);
922        self.search_field_analyzers.remove(field);
923        Ok(())
924    }
925
926    fn get_field_analyzer(&self, field: &str) -> Analyzer {
927        self.index_field_analyzers
928            .get(field)
929            .cloned()
930            .unwrap_or_else(|| self.analyzer.clone())
931    }
932
933    fn get_search_analyzer(&self, field: &str) -> Analyzer {
934        if let Some(analyzer) = self.search_field_analyzers.get(field) {
935            return analyzer.clone();
936        }
937        if let Some(analyzer) = self.index_field_analyzers.get(field) {
938            return analyzer.clone();
939        }
940        self.analyzer.clone()
941    }
942}