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
30const FORMAT_METADATA_KEY: &str = "inverted_index_format";
31const CLUSTERED_FORMAT_NAME: &str = "clustered-v1";
32const MIGRATION_PAGE_SIZE: usize = 1_024;
33
34/// Inverted index implemented over [`KeyValueStore`].
35#[derive(Clone)]
36pub struct KeyValueInvertedIndex {
37    store: Arc<dyn KeyValueStore>,
38    table: String,
39    analyzer: Analyzer,
40    index_field_analyzers: BTreeMap<FieldName, Analyzer>,
41    search_field_analyzers: BTreeMap<FieldName, Analyzer>,
42}
43
44type KeyValueStagedPosting = (FieldName, String, Vec<u32>);
45type KeyValueAnalyzedFields = (BTreeMap<FieldName, u64>, Vec<KeyValueStagedPosting>);
46type ClusterKey = (FieldName, String, u64);
47type PostingChange = Option<(u64, Vec<u32>)>;
48type KeyValueStagedDocuments = BTreeMap<DocId, KeyValueAnalyzedFields>;
49type KeyValueClusterChanges = BTreeMap<ClusterKey, BTreeMap<DocId, PostingChange>>;
50type KeyValueFieldChanges = BTreeMap<FieldName, (u64, u64)>;
51type KeyValueMergedClusters = Vec<(ClusterKey, Vec<ClusterPosting>)>;
52
53fn merge_cluster_changes(
54    entries: Vec<ClusterPosting>,
55    changes: BTreeMap<DocId, PostingChange>,
56) -> Vec<ClusterPosting> {
57    fn push_replacement(
58        entries: &mut Vec<ClusterPosting>,
59        doc_id: DocId,
60        replacement: PostingChange,
61    ) {
62        if let Some((doc_length, positions)) = replacement {
63            entries.push(ClusterPosting {
64                doc_id,
65                term_freq: positions.len() as u64,
66                doc_length,
67                positions,
68            });
69        }
70    }
71
72    let mut merged = Vec::with_capacity(entries.len().saturating_add(changes.len()));
73    let mut changes = changes.into_iter().peekable();
74    for entry in entries {
75        while changes
76            .peek()
77            .is_some_and(|(doc_id, _)| *doc_id < entry.doc_id)
78        {
79            let (doc_id, replacement) = changes.next().expect("peeked posting change exists");
80            push_replacement(&mut merged, doc_id, replacement);
81        }
82        if changes
83            .peek()
84            .is_some_and(|(doc_id, _)| *doc_id == entry.doc_id)
85        {
86            let (doc_id, replacement) = changes.next().expect("peeked posting change exists");
87            push_replacement(&mut merged, doc_id, replacement);
88        } else {
89            merged.push(entry);
90        }
91    }
92    for (doc_id, replacement) in changes {
93        push_replacement(&mut merged, doc_id, replacement);
94    }
95    merged
96}
97
98fn accumulate_field_changes(
99    field_changes: &mut KeyValueFieldChanges,
100    old_lengths: &BTreeMap<FieldName, u64>,
101    new_lengths: &BTreeMap<FieldName, u64>,
102) -> StorageBackendResult<()> {
103    let mut affected_fields = BTreeSet::new();
104    affected_fields.extend(old_lengths.keys().cloned());
105    affected_fields.extend(new_lengths.keys().cloned());
106    for field in affected_fields {
107        let (old_total, new_total) = field_changes.entry(field.clone()).or_default();
108        if let Some(length) = old_lengths.get(&field) {
109            *old_total = old_total
110                .checked_add(*length)
111                .ok_or_else(|| other_error("old field length overflow"))?;
112        }
113        if let Some(length) = new_lengths.get(&field) {
114            *new_total = new_total
115                .checked_add(*length)
116                .ok_or_else(|| other_error("new field length overflow"))?;
117        }
118    }
119    Ok(())
120}
121
122impl KeyValueInvertedIndex {
123    pub fn new(
124        store: Arc<dyn KeyValueStore>,
125        table: impl Into<String>,
126        analyzer: Analyzer,
127    ) -> Self {
128        Self {
129            store,
130            table: table.into(),
131            analyzer,
132            index_field_analyzers: BTreeMap::new(),
133            search_field_analyzers: BTreeMap::new(),
134        }
135    }
136
137    pub(crate) fn migrate_legacy_storage(store: &dyn KeyValueStore) -> StorageBackendResult<()> {
138        let marker = single_str_key(TAG_METADATA, FORMAT_METADATA_KEY)?;
139        if let Some(format) = store.get(&marker)? {
140            if format == CLUSTERED_FORMAT_NAME.as_bytes() {
141                return Ok(());
142            }
143            return Err(other_error(format!(
144                "unsupported KeyValue inverted-index format `{}`",
145                String::from_utf8_lossy(&format)
146            )));
147        }
148        if store.in_transaction() {
149            return Err(other_error(
150                "cannot migrate KeyValue postings inside an active transaction",
151            ));
152        }
153
154        store.begin_transaction()?;
155        let migration = Self::migrate_legacy_storage_in_transaction(store, &marker);
156        match migration {
157            Ok(()) => store.commit_transaction(),
158            Err(error) => match store.rollback_transaction() {
159                Ok(()) => Err(error),
160                Err(rollback) => Err(other_error(format!(
161                    "{error}; KeyValue posting migration rollback also failed: {rollback}"
162                ))),
163            },
164        }
165    }
166
167    fn migrate_legacy_storage_in_transaction(
168        store: &dyn KeyValueStore,
169        marker: &[u8],
170    ) -> StorageBackendResult<()> {
171        let posting_count = migrate_legacy_forward_postings(store)?;
172        let reverse_count = migrate_legacy_reverse_postings(store)?;
173        if posting_count != reverse_count {
174            return Err(other_error(format!(
175                "cannot migrate inconsistent KeyValue postings: {posting_count} forward rows and {reverse_count} reverse rows"
176            )));
177        }
178        store.put(marker, &string_value(CLUSTERED_FORMAT_NAME))
179    }
180
181    fn old_doc_lengths(&self, doc_id: DocId) -> StorageBackendResult<BTreeMap<FieldName, u64>> {
182        let mut out = BTreeMap::new();
183        for (key, value) in self
184            .store
185            .scan_prefix(&doc_length_doc_prefix(&self.table, doc_id)?)?
186        {
187            let mut offset = 1;
188            let _table = read_str(&key, &mut offset)?;
189            let _doc_id = read_u64(&key, &mut offset)?;
190            let field = read_str(&key, &mut offset)?;
191            out.insert(field, decode_u64_value(&value)?);
192        }
193        Ok(out)
194    }
195
196    fn old_terms(&self, doc_id: DocId) -> StorageBackendResult<BTreeMap<FieldName, Vec<String>>> {
197        let mut out = BTreeMap::new();
198        for (key, value) in self
199            .store
200            .scan_prefix(&posting_document_doc_prefix(&self.table, doc_id)?)?
201        {
202            let mut offset = 1;
203            let _table = read_str(&key, &mut offset)?;
204            let _doc_id = read_u64(&key, &mut offset)?;
205            let field = read_str(&key, &mut offset)?;
206            out.insert(field, decode_terms(&value)?);
207        }
208        Ok(out)
209    }
210
211    fn analyze_fields(
212        &self,
213        fields: BTreeMap<FieldName, String>,
214    ) -> StorageBackendResult<KeyValueAnalyzedFields> {
215        let mut lengths = BTreeMap::new();
216        let mut postings = Vec::new();
217        for (field, text) in fields {
218            let analyzer = self
219                .index_field_analyzers
220                .get(&field)
221                .unwrap_or(&self.analyzer);
222            let tokens = analyzer.analyze(&text)?;
223            let token_count = usize_to_u64(tokens.len(), "document token count")?;
224            crate::inverted_index::validate_token_position_count(token_count)?;
225            lengths.insert(field.clone(), token_count);
226            let mut term_positions: BTreeMap<String, Vec<u32>> = BTreeMap::new();
227            for (position, token) in tokens.into_iter().enumerate() {
228                term_positions.entry(token).or_default().push(
229                    u32::try_from(position)
230                        .map_err(|_| other_error("token position exceeds u32 index format"))?,
231                );
232            }
233            for (term, mut positions) in term_positions {
234                positions.sort_unstable();
235                positions.dedup();
236                postings.push((field.clone(), term, positions));
237            }
238        }
239        Ok((lengths, postings))
240    }
241
242    fn set_total_length(
243        batch: &mut dyn KeyValueBatch,
244        table: &str,
245        field: &str,
246        value: u64,
247    ) -> StorageBackendResult<()> {
248        let key = field_stats_key(table, field)?;
249        if value == 0 {
250            batch.delete(&key)
251        } else {
252            batch.put(&key, &u64_value(value))
253        }
254    }
255
256    fn load_cluster(
257        &self,
258        field: &str,
259        term: &str,
260        posting_cluster: u64,
261    ) -> StorageBackendResult<Vec<ClusterPosting>> {
262        let score = self.store.get(&posting_cluster_score_key(
263            &self.table,
264            field,
265            term,
266            posting_cluster,
267        )?)?;
268        let positions = self.store.get(&posting_cluster_positions_key(
269            &self.table,
270            field,
271            term,
272            posting_cluster,
273        )?)?;
274        match (score, positions) {
275            (None, None) => Ok(Vec::new()),
276            (Some(score), Some(positions)) => decode_cluster(posting_cluster, &score, &positions),
277            _ => Err(other_error(
278                "clustered posting score and positions values disagree",
279            )),
280        }
281    }
282
283    fn stage_cluster_changes(
284        &self,
285        doc_id: DocId,
286        old_terms: &BTreeMap<FieldName, Vec<String>>,
287        lengths: &BTreeMap<FieldName, u64>,
288        postings: &[KeyValueStagedPosting],
289    ) -> StorageBackendResult<Vec<(ClusterKey, Vec<ClusterPosting>)>> {
290        let mut changes = BTreeMap::<(FieldName, String), PostingChange>::new();
291        for (field, terms) in old_terms {
292            for term in terms {
293                changes.insert((field.clone(), term.clone()), None);
294            }
295        }
296        for (field, term, positions) in postings {
297            changes.insert(
298                (field.clone(), term.clone()),
299                Some((lengths[field], positions.clone())),
300            );
301        }
302
303        let posting_cluster = cluster_id(doc_id);
304        let mut output = Vec::with_capacity(changes.len());
305        for ((field, term), replacement) in changes {
306            let mut entries = self.load_cluster(&field, &term, posting_cluster)?;
307            if let Ok(position) = entries.binary_search_by_key(&doc_id, |entry| entry.doc_id) {
308                entries.remove(position);
309            }
310            if let Some((doc_length, positions)) = replacement {
311                let position = entries.partition_point(|entry| entry.doc_id < doc_id);
312                entries.insert(
313                    position,
314                    ClusterPosting {
315                        doc_id,
316                        term_freq: positions.len() as u64,
317                        doc_length,
318                        positions,
319                    },
320                );
321            }
322            output.push(((field, term, posting_cluster), entries));
323        }
324        Ok(output)
325    }
326
327    fn apply_cluster_changes(
328        batch: &mut dyn KeyValueBatch,
329        table: &str,
330        changes: Vec<(ClusterKey, Vec<ClusterPosting>)>,
331    ) -> StorageBackendResult<()> {
332        for ((field, term, posting_cluster), entries) in changes {
333            let score_key = posting_cluster_score_key(table, &field, &term, posting_cluster)?;
334            let positions_key =
335                posting_cluster_positions_key(table, &field, &term, posting_cluster)?;
336            if entries.is_empty() {
337                batch.delete(&score_key)?;
338                batch.delete(&positions_key)?;
339            } else {
340                let (score, positions) = encode_cluster(&entries)?;
341                batch.put(&score_key, &score)?;
342                batch.put(&positions_key, &positions)?;
343            }
344        }
345        Ok(())
346    }
347
348    fn collect_batch_changes(
349        &self,
350        staged_documents: &KeyValueStagedDocuments,
351    ) -> StorageBackendResult<(KeyValueClusterChanges, KeyValueFieldChanges)> {
352        let mut cluster_changes = KeyValueClusterChanges::new();
353        let mut field_changes = KeyValueFieldChanges::new();
354        for (doc_id, (new_lengths, new_postings)) in staged_documents {
355            let old_lengths = self.old_doc_lengths(*doc_id)?;
356            let old_terms = self.old_terms(*doc_id)?;
357            let posting_cluster = cluster_id(*doc_id);
358            for (field, terms) in old_terms {
359                for term in terms {
360                    cluster_changes
361                        .entry((field.clone(), term, posting_cluster))
362                        .or_default()
363                        .insert(*doc_id, None);
364                }
365            }
366            for (field, term, positions) in new_postings {
367                cluster_changes
368                    .entry((field.clone(), term.clone(), posting_cluster))
369                    .or_default()
370                    .insert(*doc_id, Some((new_lengths[field], positions.clone())));
371            }
372            accumulate_field_changes(&mut field_changes, &old_lengths, new_lengths)?;
373        }
374        Ok((cluster_changes, field_changes))
375    }
376
377    fn plan_batch_totals(
378        &self,
379        field_changes: KeyValueFieldChanges,
380    ) -> StorageBackendResult<Vec<(FieldName, u64)>> {
381        let mut totals = Vec::with_capacity(field_changes.len());
382        for (field, (old_total, new_total)) in field_changes {
383            let base = self
384                .store
385                .get(&field_stats_key(&self.table, &field)?)?
386                .map(|value| decode_u64_value(&value))
387                .transpose()?
388                .unwrap_or(0);
389            let total = base
390                .checked_sub(old_total)
391                .ok_or_else(|| other_error("stored field length is smaller than batch length"))?
392                .checked_add(new_total)
393                .ok_or_else(|| other_error("total field length overflow"))?;
394            totals.push((field, total));
395        }
396        Ok(totals)
397    }
398
399    fn merge_batch_clusters(
400        &self,
401        cluster_changes: KeyValueClusterChanges,
402    ) -> StorageBackendResult<KeyValueMergedClusters> {
403        let mut merged = Vec::with_capacity(cluster_changes.len());
404        for ((field, term, posting_cluster), changes) in cluster_changes {
405            let entries = self.load_cluster(&field, &term, posting_cluster)?;
406            merged.push((
407                (field, term, posting_cluster),
408                merge_cluster_changes(entries, changes),
409            ));
410        }
411        Ok(merged)
412    }
413
414    fn write_batch_documents(
415        &self,
416        batch: &mut dyn KeyValueBatch,
417        staged_documents: KeyValueStagedDocuments,
418    ) -> StorageBackendResult<()> {
419        for (doc_id, (lengths, postings)) in staged_documents {
420            batch.delete_prefix(&posting_document_doc_prefix(&self.table, doc_id)?)?;
421            batch.delete_prefix(&doc_length_doc_prefix(&self.table, doc_id)?)?;
422            let mut terms_by_field = BTreeMap::<FieldName, Vec<String>>::new();
423            for (field, term, _) in postings {
424                terms_by_field.entry(field).or_default().push(term);
425            }
426            for (field, length) in lengths {
427                batch.put(
428                    &doc_length_key(&self.table, doc_id, &field)?,
429                    &u64_value(length),
430                )?;
431                batch.put(
432                    &posting_document_key(&self.table, doc_id, &field)?,
433                    &encode_terms(terms_by_field.get(&field).map_or(&[], Vec::as_slice))?,
434                )?;
435            }
436        }
437        Ok(())
438    }
439
440    fn add_documents(
441        &self,
442        documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
443    ) -> StorageBackendResult<()> {
444        let mut staged_documents = BTreeMap::new();
445        for (doc_id, fields) in documents {
446            staged_documents.insert(doc_id, self.analyze_fields(fields)?);
447        }
448        if staged_documents.is_empty() {
449            return Ok(());
450        }
451
452        let (cluster_changes, field_changes) = self.collect_batch_changes(&staged_documents)?;
453        let totals = self.plan_batch_totals(field_changes)?;
454        let merged_clusters = self.merge_batch_clusters(cluster_changes)?;
455        let mut batch = self.store.batch();
456        Self::apply_cluster_changes(batch.as_mut(), &self.table, merged_clusters)?;
457        for (field, total) in totals {
458            Self::set_total_length(batch.as_mut(), &self.table, &field, total)?;
459        }
460        self.write_batch_documents(batch.as_mut(), staged_documents)?;
461        batch.commit()
462    }
463
464    fn cursor_for_term(
465        &self,
466        field: &str,
467        term: &str,
468    ) -> StorageBackendResult<Box<dyn PostingCursor>> {
469        let mut clusters = Vec::new();
470        for (key, bytes) in self.store.scan_prefix(&posting_cluster_score_term_prefix(
471            &self.table,
472            field,
473            term,
474        )?)? {
475            let mut offset = 1;
476            let _table = read_str(&key, &mut offset)?;
477            let _field = read_str(&key, &mut offset)?;
478            let _term = read_str(&key, &mut offset)?;
479            let posting_cluster = read_u64(&key, &mut offset)?;
480            if offset != key.len() {
481                return Err(other_error("invalid clustered posting score key"));
482            }
483            clusters.push(EncodedScoreCluster {
484                cluster_id: posting_cluster,
485                bytes,
486            });
487        }
488        if clusters.is_empty() {
489            return Ok(Box::new(MaterializedPostingCursor::new(Vec::new())?));
490        }
491        Ok(Box::new(ClusteredPostingCursor::new(clusters)?))
492    }
493}
494
495fn migrate_legacy_forward_postings(store: &dyn KeyValueStore) -> StorageBackendResult<u64> {
496    let posting_prefix = key_with_tag(TAG_POSTING);
497    let mut after = None::<Vec<u8>>;
498    let mut group = None::<(String, String, String, u64, Vec<ClusterPosting>)>;
499    let mut posting_count = 0_u64;
500    loop {
501        let page =
502            store.scan_prefix_after(&posting_prefix, after.as_deref(), MIGRATION_PAGE_SIZE)?;
503        if page.is_empty() {
504            break;
505        }
506        for (key, value) in page {
507            after = Some(key.clone());
508            let (table, field, term, doc_id) = decode_legacy_posting_key(&key)?;
509            let positions = blob_to_positions(&value)?;
510            let doc_length = store
511                .get(&doc_length_key(&table, doc_id, &field)?)?
512                .map(|value| decode_u64_value(&value))
513                .transpose()?
514                .ok_or_else(|| {
515                    other_error(format!(
516                        "cannot migrate posting `{table}.{field}.{term}` for document {doc_id}: missing document length"
517                    ))
518                })?;
519            if !store.contains_key(&reverse_posting_key(&table, doc_id, &field, &term)?)? {
520                return Err(other_error(format!(
521                    "cannot migrate posting `{table}.{field}.{term}` for document {doc_id}: missing reverse posting"
522                )));
523            }
524            let posting_cluster = cluster_id(doc_id);
525            let same_group = group.as_ref().is_some_and(
526                |(group_table, group_field, group_term, group_cluster, _)| {
527                    group_table == &table
528                        && group_field == &field
529                        && group_term == &term
530                        && *group_cluster == posting_cluster
531                },
532            );
533            if !same_group {
534                if let Some(cluster) = group.take() {
535                    put_migrated_cluster(store, cluster)?;
536                }
537                group = Some((table, field, term, posting_cluster, Vec::new()));
538            }
539            group
540                .as_mut()
541                .expect("posting migration group exists")
542                .4
543                .push(ClusterPosting {
544                    doc_id,
545                    term_freq: positions.len() as u64,
546                    doc_length,
547                    positions,
548                });
549            posting_count = posting_count
550                .checked_add(1)
551                .ok_or_else(|| other_error("legacy posting count overflow"))?;
552            store.delete(&key)?;
553        }
554    }
555    if let Some(cluster) = group {
556        put_migrated_cluster(store, cluster)?;
557    }
558    Ok(posting_count)
559}
560
561fn migrate_legacy_reverse_postings(store: &dyn KeyValueStore) -> StorageBackendResult<u64> {
562    let reverse_prefix = key_with_tag(TAG_REVERSE_POSTING);
563    let mut after = None::<Vec<u8>>;
564    let mut group = None::<(String, DocId, FieldName, Vec<String>)>;
565    let mut reverse_count = 0_u64;
566    loop {
567        let page =
568            store.scan_prefix_after(&reverse_prefix, after.as_deref(), MIGRATION_PAGE_SIZE)?;
569        if page.is_empty() {
570            break;
571        }
572        for (key, _) in page {
573            after = Some(key.clone());
574            let (table, doc_id, field, term) = decode_legacy_reverse_key(&key)?;
575            let same_group =
576                group
577                    .as_ref()
578                    .is_some_and(|(group_table, group_doc_id, group_field, _)| {
579                        group_table == &table && *group_doc_id == doc_id && group_field == &field
580                    });
581            if !same_group {
582                if let Some(document) = group.take() {
583                    put_migrated_document(store, document)?;
584                }
585                group = Some((table, doc_id, field, Vec::new()));
586            }
587            group
588                .as_mut()
589                .expect("reverse posting migration group exists")
590                .3
591                .push(term);
592            reverse_count = reverse_count
593                .checked_add(1)
594                .ok_or_else(|| other_error("legacy reverse posting count overflow"))?;
595            store.delete(&key)?;
596        }
597    }
598    if let Some(document) = group {
599        put_migrated_document(store, document)?;
600    }
601    Ok(reverse_count)
602}
603
604impl InvertedIndex for KeyValueInvertedIndex {
605    fn analyzer(&self) -> &Analyzer {
606        &self.analyzer
607    }
608
609    fn add_document(
610        &mut self,
611        doc_id: DocId,
612        fields: BTreeMap<FieldName, String>,
613    ) -> StorageBackendResult<()> {
614        let old_lengths = self.old_doc_lengths(doc_id)?;
615        let old_terms = self.old_terms(doc_id)?;
616        let (new_lengths, new_postings) = self.analyze_fields(fields)?;
617        let cluster_changes =
618            self.stage_cluster_changes(doc_id, &old_terms, &new_lengths, &new_postings)?;
619
620        let mut fields_to_update = BTreeSet::new();
621        fields_to_update.extend(old_lengths.keys().cloned());
622        fields_to_update.extend(new_lengths.keys().cloned());
623        let mut totals = Vec::with_capacity(fields_to_update.len());
624        for field in fields_to_update {
625            let base = self
626                .store
627                .get(&field_stats_key(&self.table, &field)?)?
628                .map(|value| decode_u64_value(&value))
629                .transpose()?
630                .unwrap_or(0);
631            let old = old_lengths.get(&field).copied().unwrap_or(0);
632            let new = new_lengths.get(&field).copied().unwrap_or(0);
633            let total = base
634                .checked_sub(old)
635                .ok_or_else(|| other_error("stored field length is smaller than document length"))?
636                .checked_add(new)
637                .ok_or_else(|| other_error("total field length overflow"))?;
638            totals.push((field, total));
639        }
640
641        let mut terms_by_field = BTreeMap::<FieldName, Vec<String>>::new();
642        for (field, term, _) in &new_postings {
643            terms_by_field
644                .entry(field.clone())
645                .or_default()
646                .push(term.clone());
647        }
648        for field in new_lengths.keys() {
649            terms_by_field.entry(field.clone()).or_default();
650        }
651
652        let mut batch = self.store.batch();
653        Self::apply_cluster_changes(batch.as_mut(), &self.table, cluster_changes)?;
654        batch.delete_prefix(&posting_document_doc_prefix(&self.table, doc_id)?)?;
655        for field in old_lengths.keys() {
656            batch.delete(&doc_length_key(&self.table, doc_id, field)?)?;
657        }
658        for (field, total) in totals {
659            Self::set_total_length(batch.as_mut(), &self.table, &field, total)?;
660        }
661        for (field, length) in &new_lengths {
662            batch.put(
663                &doc_length_key(&self.table, doc_id, field)?,
664                &u64_value(*length),
665            )?;
666        }
667        for (field, terms) in terms_by_field {
668            batch.put(
669                &posting_document_key(&self.table, doc_id, &field)?,
670                &encode_terms(&terms)?,
671            )?;
672        }
673        batch.commit()
674    }
675
676    fn try_add_documents(
677        &mut self,
678        documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
679    ) -> StorageBackendResult<()> {
680        self.add_documents(documents)
681    }
682
683    fn remove_document(&mut self, doc_id: DocId) -> StorageBackendResult<()> {
684        let old_lengths = self.old_doc_lengths(doc_id)?;
685        let old_terms = self.old_terms(doc_id)?;
686        let cluster_changes =
687            self.stage_cluster_changes(doc_id, &old_terms, &BTreeMap::new(), &[])?;
688        let mut totals = Vec::with_capacity(old_lengths.len());
689        for (field, length) in &old_lengths {
690            let base = self
691                .store
692                .get(&field_stats_key(&self.table, field)?)?
693                .map(|value| decode_u64_value(&value))
694                .transpose()?
695                .unwrap_or(0);
696            totals.push((
697                field.clone(),
698                base.checked_sub(*length).ok_or_else(|| {
699                    other_error("stored field length is smaller than removed document length")
700                })?,
701            ));
702        }
703
704        let mut batch = self.store.batch();
705        Self::apply_cluster_changes(batch.as_mut(), &self.table, cluster_changes)?;
706        batch.delete_prefix(&posting_document_doc_prefix(&self.table, doc_id)?)?;
707        for (field, total) in totals {
708            Self::set_total_length(batch.as_mut(), &self.table, &field, total)?;
709            batch.delete(&doc_length_key(&self.table, doc_id, &field)?)?;
710        }
711        batch.commit()
712    }
713
714    fn try_rebuild_documents(
715        &mut self,
716        documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
717    ) -> StorageBackendResult<()> {
718        let mut staged = BTreeMap::new();
719        for (doc_id, fields) in documents {
720            if !fields.is_empty() {
721                staged.insert(doc_id, self.analyze_fields(fields)?);
722            }
723        }
724        let mut totals = BTreeMap::<FieldName, u64>::new();
725        let mut clusters = BTreeMap::<ClusterKey, Vec<ClusterPosting>>::new();
726        for (doc_id, (lengths, postings)) in &staged {
727            for (field, length) in lengths {
728                let total = totals.entry(field.clone()).or_default();
729                *total = total
730                    .checked_add(*length)
731                    .ok_or_else(|| other_error("total field length overflow"))?;
732            }
733            for (field, term, positions) in postings {
734                clusters
735                    .entry((field.clone(), term.clone(), cluster_id(*doc_id)))
736                    .or_default()
737                    .push(ClusterPosting {
738                        doc_id: *doc_id,
739                        term_freq: positions.len() as u64,
740                        doc_length: lengths[field],
741                        positions: positions.clone(),
742                    });
743            }
744        }
745
746        let mut batch = self.store.batch();
747        batch.delete_prefix(&posting_cluster_score_key_prefix(&self.table)?)?;
748        batch.delete_prefix(&posting_cluster_positions_key_prefix(&self.table)?)?;
749        batch.delete_prefix(&posting_document_key_prefix(&self.table)?)?;
750        batch.delete_prefix(&doc_length_key_prefix(&self.table)?)?;
751        batch.delete_prefix(&field_stats_key_prefix(&self.table)?)?;
752        for (field, total) in totals {
753            Self::set_total_length(batch.as_mut(), &self.table, &field, total)?;
754        }
755        for ((field, term, posting_cluster), entries) in clusters {
756            let (score, positions) = encode_cluster(&entries)?;
757            batch.put(
758                &posting_cluster_score_key(&self.table, &field, &term, posting_cluster)?,
759                &score,
760            )?;
761            batch.put(
762                &posting_cluster_positions_key(&self.table, &field, &term, posting_cluster)?,
763                &positions,
764            )?;
765        }
766        for (doc_id, (lengths, postings)) in staged {
767            for (field, length) in lengths {
768                batch.put(
769                    &doc_length_key(&self.table, doc_id, &field)?,
770                    &u64_value(length),
771                )?;
772                let terms = postings
773                    .iter()
774                    .filter(|(posting_field, _, _)| posting_field == &field)
775                    .map(|(_, term, _)| term.clone())
776                    .collect::<Vec<_>>();
777                batch.put(
778                    &posting_document_key(&self.table, doc_id, &field)?,
779                    &encode_terms(&terms)?,
780                )?;
781            }
782        }
783        batch.commit()
784    }
785
786    fn clear(&mut self) -> StorageBackendResult<()> {
787        let mut batch = self.store.batch();
788        batch.delete_prefix(&posting_cluster_score_key_prefix(&self.table)?)?;
789        batch.delete_prefix(&posting_cluster_positions_key_prefix(&self.table)?)?;
790        batch.delete_prefix(&posting_document_key_prefix(&self.table)?)?;
791        batch.delete_prefix(&doc_length_key_prefix(&self.table)?)?;
792        batch.delete_prefix(&field_stats_key_prefix(&self.table)?)?;
793        batch.commit()
794    }
795
796    fn get_posting_list(&self, field: &str, term: &str) -> StorageBackendResult<PostingList> {
797        let mut entries = Vec::new();
798        for (key, score) in self.store.scan_prefix(&posting_cluster_score_term_prefix(
799            &self.table,
800            field,
801            term,
802        )?)? {
803            let mut offset = 1;
804            let _table = read_str(&key, &mut offset)?;
805            let _field = read_str(&key, &mut offset)?;
806            let _term = read_str(&key, &mut offset)?;
807            let posting_cluster = read_u64(&key, &mut offset)?;
808            let positions = self
809                .store
810                .get(&posting_cluster_positions_key(
811                    &self.table,
812                    field,
813                    term,
814                    posting_cluster,
815                )?)?
816                .ok_or_else(|| other_error("clustered posting positions value is missing"))?;
817            entries.extend(
818                decode_cluster(posting_cluster, &score, &positions)?
819                    .into_iter()
820                    .map(|entry| {
821                        PostingEntry::new(
822                            entry.doc_id,
823                            Payload {
824                                positions: entry.positions,
825                                score: 0.0,
826                                fields: BTreeMap::new(),
827                            },
828                        )
829                    }),
830            );
831        }
832        Ok(PostingList::from_sorted_unchecked(entries))
833    }
834
835    fn posting_cursor(
836        &self,
837        field: &str,
838        term: &str,
839    ) -> StorageBackendResult<Box<dyn PostingCursor>> {
840        self.cursor_for_term(field, term)
841    }
842
843    fn for_each_term_freq(
844        &self,
845        field: &str,
846        term: &str,
847        visit: &mut dyn FnMut(DocId, u64),
848    ) -> StorageBackendResult<()> {
849        let mut cursor = self.cursor_for_term(field, term)?;
850        while let Some(entry) = cursor.current() {
851            visit(entry.doc_id, entry.term_freq);
852            cursor.advance()?;
853        }
854        Ok(())
855    }
856
857    fn doc_freq(&self, field: &str, term: &str) -> StorageBackendResult<u64> {
858        self.store
859            .scan_prefix(&posting_cluster_score_term_prefix(
860                &self.table,
861                field,
862                term,
863            )?)?
864            .into_iter()
865            .try_fold(0_u64, |total, (_, score)| {
866                total
867                    .checked_add(score_count(&score)?)
868                    .ok_or_else(|| other_error("document frequency overflow"))
869            })
870    }
871
872    fn get_doc_length(&self, doc_id: DocId, field: &str) -> StorageBackendResult<u64> {
873        Ok(self
874            .store
875            .get(&doc_length_key(&self.table, doc_id, field)?)?
876            .map(|value| decode_u64_value(&value))
877            .transpose()?
878            .unwrap_or(0))
879    }
880
881    fn get_scoring_inputs_bulk(
882        &self,
883        doc_ids: &[DocId],
884        field: &str,
885        terms: &[String],
886    ) -> StorageBackendResult<Vec<(u64, Vec<u64>)>> {
887        let mut output = doc_ids
888            .iter()
889            .map(|doc_id| Ok((self.get_doc_length(*doc_id, field)?, vec![0; terms.len()])))
890            .collect::<StorageBackendResult<Vec<_>>>()?;
891        let mut positions = BTreeMap::<DocId, Vec<usize>>::new();
892        for (position, doc_id) in doc_ids.iter().copied().enumerate() {
893            positions.entry(doc_id).or_default().push(position);
894        }
895        for (term_index, term) in terms.iter().enumerate() {
896            let mut cursor = self.cursor_for_term(field, term)?;
897            while let Some(entry) = cursor.current() {
898                if let Some(output_positions) = positions.get(&entry.doc_id) {
899                    for position in output_positions {
900                        output[*position].0 = entry.doc_length;
901                        output[*position].1[term_index] = entry.term_freq;
902                    }
903                }
904                cursor.advance()?;
905            }
906        }
907        Ok(output)
908    }
909
910    fn get_term_freq(&self, doc_id: DocId, field: &str, term: &str) -> StorageBackendResult<u64> {
911        let posting_cluster = cluster_id(doc_id);
912        self.store
913            .get(&posting_cluster_score_key(
914                &self.table,
915                field,
916                term,
917                posting_cluster,
918            )?)?
919            .map_or(Ok(0), |score| {
920                let entries = decode_all_scores(posting_cluster, &score)?;
921                Ok(entries
922                    .binary_search_by_key(&doc_id, |entry| entry.doc_id)
923                    .ok()
924                    .map_or(0, |position| entries[position].term_freq))
925            })
926    }
927
928    fn doc_count(&self) -> StorageBackendResult<u64> {
929        let mut doc_ids = BTreeSet::new();
930        for (key, _) in self
931            .store
932            .scan_prefix(&doc_length_key_prefix(&self.table)?)?
933        {
934            let mut offset = 1;
935            let _table = read_str(&key, &mut offset)?;
936            doc_ids.insert(read_u64(&key, &mut offset)?);
937        }
938        usize_to_u64(doc_ids.len(), "document count")
939    }
940
941    fn total_field_length(&self, field: &str) -> StorageBackendResult<u64> {
942        Ok(self
943            .store
944            .get(&field_stats_key(&self.table, field)?)?
945            .map(|value| decode_u64_value(&value))
946            .transpose()?
947            .unwrap_or(0))
948    }
949
950    fn vocabulary_terms(&self, field: &str) -> StorageBackendResult<Vec<String>> {
951        let mut terms = BTreeSet::new();
952        for (key, _) in self
953            .store
954            .scan_prefix(&posting_cluster_score_field_prefix(&self.table, field)?)?
955        {
956            let mut offset = 1;
957            let _table = read_str(&key, &mut offset)?;
958            let _field = read_str(&key, &mut offset)?;
959            terms.insert(read_str(&key, &mut offset)?);
960        }
961        Ok(terms.into_iter().collect())
962    }
963
964    fn stats(&self) -> StorageBackendResult<IndexStats> {
965        let doc_count = self.doc_count()?;
966        let mut stats = IndexStats::default();
967        stats.total_docs = doc_count;
968        if doc_count > 0 {
969            let mut total = 0_u64;
970            for (_, value) in self
971                .store
972                .scan_prefix(&field_stats_key_prefix(&self.table)?)?
973            {
974                total = total
975                    .checked_add(decode_u64_value(&value)?)
976                    .ok_or_else(|| other_error("index total field length overflow"))?;
977            }
978            stats.avg_doc_length = total as f64 / doc_count as f64;
979        }
980        let mut counts = BTreeMap::<(String, String), u64>::new();
981        for (key, value) in self
982            .store
983            .scan_prefix(&posting_cluster_score_key_prefix(&self.table)?)?
984        {
985            let mut offset = 1;
986            let _table = read_str(&key, &mut offset)?;
987            let field = read_str(&key, &mut offset)?;
988            let term = read_str(&key, &mut offset)?;
989            let count = counts.entry((field, term)).or_default();
990            *count = count
991                .checked_add(score_count(&value)?)
992                .ok_or_else(|| other_error("index document frequency overflow"))?;
993        }
994        for ((field, term), document_frequency) in counts {
995            stats.set_doc_freq(field, term, document_frequency);
996        }
997        Ok(stats)
998    }
999
1000    fn posting_count(&self, field: Option<&str>) -> StorageBackendResult<u64> {
1001        let prefix = match field {
1002            Some(field) => posting_cluster_score_field_prefix(&self.table, field)?,
1003            None => posting_cluster_score_key_prefix(&self.table)?,
1004        };
1005        self.store
1006            .scan_prefix(&prefix)?
1007            .into_iter()
1008            .try_fold(0_u64, |total, (_, value)| {
1009                total
1010                    .checked_add(score_count(&value)?)
1011                    .ok_or_else(|| other_error("posting count overflow"))
1012            })
1013    }
1014
1015    fn doc_length_count(&self, field: Option<&str>) -> StorageBackendResult<u64> {
1016        let mut count = 0_u64;
1017        for (key, _) in self
1018            .store
1019            .scan_prefix(&doc_length_key_prefix(&self.table)?)?
1020        {
1021            let mut offset = 1;
1022            let _table = read_str(&key, &mut offset)?;
1023            let _doc_id = read_u64(&key, &mut offset)?;
1024            let indexed_field = read_str(&key, &mut offset)?;
1025            if field.is_none_or(|target| target == indexed_field) {
1026                count = count
1027                    .checked_add(1)
1028                    .ok_or_else(|| other_error("document-length row count overflow"))?;
1029            }
1030        }
1031        Ok(count)
1032    }
1033
1034    fn term_count(&self, field: Option<&str>) -> StorageBackendResult<u64> {
1035        let prefix = match field {
1036            Some(field) => posting_cluster_score_field_prefix(&self.table, field)?,
1037            None => posting_cluster_score_key_prefix(&self.table)?,
1038        };
1039        let mut terms = BTreeSet::new();
1040        for (key, _) in self.store.scan_prefix(&prefix)? {
1041            let mut offset = 1;
1042            let _table = read_str(&key, &mut offset)?;
1043            let current_field = read_str(&key, &mut offset)?;
1044            let term = read_str(&key, &mut offset)?;
1045            terms.insert((current_field, term));
1046        }
1047        usize_to_u64(terms.len(), "term count")
1048    }
1049
1050    fn snapshot(&self) -> StorageBackendResult<Arc<dyn InvertedIndex>> {
1051        Ok(Arc::new(self.clone()))
1052    }
1053
1054    fn field_names(&self) -> StorageBackendResult<Vec<FieldName>> {
1055        let mut fields = Vec::new();
1056        for (key, _) in self
1057            .store
1058            .scan_prefix(&field_stats_key_prefix(&self.table)?)?
1059        {
1060            let mut offset = 1;
1061            let _table = read_str(&key, &mut offset)?;
1062            fields.push(read_str(&key, &mut offset)?);
1063        }
1064        Ok(fields)
1065    }
1066
1067    fn set_field_analyzer(
1068        &mut self,
1069        field: &str,
1070        analyzer: Analyzer,
1071        phase: AnalyzerPhase,
1072    ) -> Result<(), String> {
1073        match phase {
1074            AnalyzerPhase::Index => {
1075                self.index_field_analyzers
1076                    .insert(field.to_string(), analyzer);
1077            }
1078            AnalyzerPhase::Search => {
1079                self.search_field_analyzers
1080                    .insert(field.to_string(), analyzer);
1081            }
1082            AnalyzerPhase::Both => {
1083                self.index_field_analyzers
1084                    .insert(field.to_string(), analyzer.clone());
1085                self.search_field_analyzers
1086                    .insert(field.to_string(), analyzer);
1087            }
1088        }
1089        Ok(())
1090    }
1091
1092    fn remove_field_analyzers(&mut self, field: &str) -> Result<(), String> {
1093        self.index_field_analyzers.remove(field);
1094        self.search_field_analyzers.remove(field);
1095        Ok(())
1096    }
1097
1098    fn get_field_analyzer(&self, field: &str) -> Analyzer {
1099        self.index_field_analyzers
1100            .get(field)
1101            .cloned()
1102            .unwrap_or_else(|| self.analyzer.clone())
1103    }
1104
1105    fn get_search_analyzer(&self, field: &str) -> Analyzer {
1106        if let Some(analyzer) = self.search_field_analyzers.get(field) {
1107            return analyzer.clone();
1108        }
1109        if let Some(analyzer) = self.index_field_analyzers.get(field) {
1110            return analyzer.clone();
1111        }
1112        self.analyzer.clone()
1113    }
1114}
1115
1116fn decode_legacy_posting_key(
1117    key: &[u8],
1118) -> StorageBackendResult<(String, FieldName, String, DocId)> {
1119    let mut offset = 1;
1120    let table = read_str(key, &mut offset)?;
1121    let field = read_str(key, &mut offset)?;
1122    let term = read_str(key, &mut offset)?;
1123    let doc_id = read_u64(key, &mut offset)?;
1124    if offset != key.len() {
1125        return Err(other_error("invalid legacy posting key"));
1126    }
1127    Ok((table, field, term, doc_id))
1128}
1129
1130fn decode_legacy_reverse_key(
1131    key: &[u8],
1132) -> StorageBackendResult<(String, DocId, FieldName, String)> {
1133    let mut offset = 1;
1134    let table = read_str(key, &mut offset)?;
1135    let doc_id = read_u64(key, &mut offset)?;
1136    let field = read_str(key, &mut offset)?;
1137    let term = read_str(key, &mut offset)?;
1138    if offset != key.len() {
1139        return Err(other_error("invalid legacy reverse posting key"));
1140    }
1141    Ok((table, doc_id, field, term))
1142}
1143
1144fn put_migrated_cluster(
1145    store: &dyn KeyValueStore,
1146    cluster: (String, FieldName, String, u64, Vec<ClusterPosting>),
1147) -> StorageBackendResult<()> {
1148    let (table, field, term, posting_cluster, entries) = cluster;
1149    let (score_blob, positions_blob) = encode_cluster(&entries)?;
1150    store.put(
1151        &posting_cluster_score_key(&table, &field, &term, posting_cluster)?,
1152        &score_blob,
1153    )?;
1154    store.put(
1155        &posting_cluster_positions_key(&table, &field, &term, posting_cluster)?,
1156        &positions_blob,
1157    )
1158}
1159
1160fn put_migrated_document(
1161    store: &dyn KeyValueStore,
1162    document: (String, DocId, FieldName, Vec<String>),
1163) -> StorageBackendResult<()> {
1164    let (table, doc_id, field, mut terms) = document;
1165    // Length-prefixed key segments sort by encoded length before text bytes,
1166    // while the shared terms codec requires ordinary lexical ordering.
1167    terms.sort_unstable();
1168    store.put(
1169        &posting_document_key(&table, doc_id, &field)?,
1170        &encode_terms(&terms)?,
1171    )
1172}