Skip to main content

uqa_storage/key_value/inverted_index/
trait_impl.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Key/value index provider operations and retained revision installation.
8
9use super::super::codec::usize_to_u64;
10use super::queries::require_score_version;
11use super::{
12    cluster_id, decode_all_scores, keys, other_error, score_count, Analyzer, AnalyzerPhase, Arc,
13    BTreeMap, BTreeSet, DocId, FieldName, IndexStats, IndexedFieldMetadata, InvertedIndex,
14    KeyValueInvertedIndex, OccurrencePosting, Payload, PostingCursor, PostingEntry, PostingList,
15    StorageBackendResult, TokenTermKey,
16};
17
18impl InvertedIndex for KeyValueInvertedIndex {
19    fn visit_score_clusters(
20        &self,
21        field: &str,
22        term: &TokenTermKey,
23        after: Option<u64>,
24        limit: usize,
25        control: &crate::read_control::StorageReadControl,
26        visit: &mut crate::clustered_postings::ScoreClusterVisitor<'_>,
27    ) -> StorageBackendResult<()> {
28        self.visit_clusters_budgeted(field, term, after, limit, control, visit)
29    }
30
31    fn get_occurrences_budgeted(
32        &self,
33        doc_id: DocId,
34        field: &str,
35        term: &TokenTermKey,
36        control: &crate::read_control::StorageReadControl,
37    ) -> StorageBackendResult<uqa_core::memory::Budgeted<Vec<uqa_core::TokenOccurrence>>> {
38        self.occurrences_budgeted(doc_id, field, term, control)
39    }
40
41    fn field_stats_scalar_budgeted(
42        &self,
43        field: &str,
44        control: &crate::read_control::StorageReadControl,
45    ) -> StorageBackendResult<IndexStats> {
46        self.scalar_stats_budgeted(field, control)
47    }
48
49    fn analyzer(&self) -> &Analyzer {
50        self.bindings.default_configuration()
51    }
52
53    fn source_rebuild_required(&self) -> StorageBackendResult<bool> {
54        self.needs_source_rebuild()
55    }
56
57    fn add_document(
58        &mut self,
59        doc_id: DocId,
60        fields: BTreeMap<FieldName, String>,
61    ) -> StorageBackendResult<()> {
62        self.add_documents(vec![(doc_id, fields)])
63    }
64
65    fn try_add_documents(
66        &mut self,
67        documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
68    ) -> StorageBackendResult<()> {
69        self.add_documents(documents)
70    }
71
72    fn remove_document(&mut self, doc_id: DocId) -> StorageBackendResult<()> {
73        self.add_documents(vec![(doc_id, BTreeMap::new())])
74    }
75
76    fn try_rebuild_documents(
77        &mut self,
78        documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
79    ) -> StorageBackendResult<()> {
80        self.rebuild_documents(documents)
81    }
82
83    fn try_rebuild_documents_cancellable(
84        &mut self,
85        documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
86        cancellation: &uqa_core::CancellationToken,
87    ) -> StorageBackendResult<()> {
88        self.rebuild_documents_inner(documents, Some(cancellation))
89    }
90
91    fn clear(&mut self) -> StorageBackendResult<()> {
92        let mut batch = self.store.batch();
93        self.clear_index_batch(batch.as_mut())?;
94        batch.commit()
95    }
96
97    fn get_posting_list(&self, field: &str, term: &str) -> StorageBackendResult<PostingList> {
98        self.get_posting_list_key(field, &TokenTermKey::from_text(term))
99    }
100
101    fn get_posting_list_key(
102        &self,
103        field: &str,
104        term: &TokenTermKey,
105    ) -> StorageBackendResult<PostingList> {
106        let entries = self
107            .occurrence_postings(field, term)?
108            .into_iter()
109            .map(|entry| {
110                PostingEntry::new(
111                    entry.doc_id,
112                    Payload {
113                        positions: entry.positions(),
114                        score: 0.0,
115                        fields: BTreeMap::new(),
116                    },
117                )
118            })
119            .collect();
120        Ok(PostingList::from_sorted_unchecked(entries))
121    }
122
123    fn posting_cursor(
124        &self,
125        field: &str,
126        term: &str,
127    ) -> StorageBackendResult<Box<dyn PostingCursor>> {
128        self.cursor_for_term(field, &TokenTermKey::from_text(term))
129    }
130
131    fn posting_cursor_key(
132        &self,
133        field: &str,
134        term: &TokenTermKey,
135    ) -> StorageBackendResult<Box<dyn PostingCursor>> {
136        self.cursor_for_term(field, term)
137    }
138
139    fn get_occurrence_postings(
140        &self,
141        field: &str,
142        term: &TokenTermKey,
143    ) -> StorageBackendResult<Vec<OccurrencePosting>> {
144        self.occurrence_postings(field, term)
145    }
146
147    fn get_occurrences(
148        &self,
149        doc_id: DocId,
150        field: &str,
151        term: &TokenTermKey,
152    ) -> StorageBackendResult<Vec<uqa_core::TokenOccurrence>> {
153        self.require_graph_format()?;
154        let entries = self.load_cluster(field, term, cluster_id(doc_id))?;
155        let Some(posting) = entries.into_iter().find(|entry| entry.doc_id == doc_id) else {
156            return Ok(Vec::new());
157        };
158        self.validate_posting_metadata(field, &posting)?;
159        Ok(posting.occurrences)
160    }
161
162    fn indexed_field_metadata(
163        &self,
164        doc_id: DocId,
165        field: &str,
166    ) -> StorageBackendResult<Option<IndexedFieldMetadata>> {
167        self.require_graph_format()?;
168        self.read_field_metadata(doc_id, field)
169    }
170
171    fn for_each_term_freq(
172        &self,
173        field: &str,
174        term: &str,
175        visit: &mut dyn FnMut(DocId, u64),
176    ) -> StorageBackendResult<()> {
177        let mut cursor = self.posting_cursor(field, term)?;
178        while let Some(entry) = cursor.current() {
179            visit(entry.doc_id, entry.term_freq);
180            cursor.advance()?;
181        }
182        Ok(())
183    }
184
185    fn doc_freq(&self, field: &str, term: &str) -> StorageBackendResult<u64> {
186        self.doc_freq_key(field, &TokenTermKey::from_text(term))
187    }
188
189    fn doc_freq_key(&self, field: &str, term: &TokenTermKey) -> StorageBackendResult<u64> {
190        self.require_graph_format()?;
191        self.store
192            .scan_prefix(&keys::term_prefix(&self.table, keys::SCORE, field, term)?)?
193            .into_iter()
194            .try_fold(0_u64, |total, (key, score)| {
195                keys::read_cluster(&key, keys::SCORE)?;
196                require_score_version(&score)?;
197                total
198                    .checked_add(score_count(&score)?)
199                    .ok_or_else(|| other_error("document frequency overflow"))
200            })
201    }
202
203    fn get_doc_length(&self, doc_id: DocId, field: &str) -> StorageBackendResult<u64> {
204        self.document_length(doc_id, field)
205    }
206
207    fn get_scoring_inputs_bulk(
208        &self,
209        doc_ids: &[DocId],
210        field: &str,
211        terms: &[String],
212    ) -> StorageBackendResult<Vec<(u64, Vec<u64>)>> {
213        let keys = terms
214            .iter()
215            .map(|term| TokenTermKey::from_text(term))
216            .collect::<Vec<_>>();
217        self.get_scoring_inputs_keys_bulk(doc_ids, field, &keys)
218    }
219
220    fn get_scoring_inputs_keys_bulk(
221        &self,
222        doc_ids: &[DocId],
223        field: &str,
224        terms: &[TokenTermKey],
225    ) -> StorageBackendResult<Vec<(u64, Vec<u64>)>> {
226        let mut output = doc_ids
227            .iter()
228            .map(|id| Ok((self.document_length(*id, field)?, vec![0; terms.len()])))
229            .collect::<StorageBackendResult<Vec<_>>>()?;
230        let mut positions = BTreeMap::<DocId, Vec<usize>>::new();
231        for (position, doc_id) in doc_ids.iter().copied().enumerate() {
232            positions.entry(doc_id).or_default().push(position);
233        }
234        for (term_index, term) in terms.iter().enumerate() {
235            let mut cursor = self.posting_cursor_key(field, term)?;
236            while let Some(entry) = cursor.current() {
237                if let Some(output_positions) = positions.get(&entry.doc_id) {
238                    for position in output_positions {
239                        output[*position].0 = entry.doc_length;
240                        output[*position].1[term_index] = entry.term_freq;
241                    }
242                }
243                cursor.advance()?;
244            }
245        }
246        Ok(output)
247    }
248
249    fn get_term_freq(&self, doc_id: DocId, field: &str, term: &str) -> StorageBackendResult<u64> {
250        self.get_term_freq_key(doc_id, field, &TokenTermKey::from_text(term))
251    }
252
253    fn get_term_freq_key(
254        &self,
255        doc_id: DocId,
256        field: &str,
257        term: &TokenTermKey,
258    ) -> StorageBackendResult<u64> {
259        self.require_graph_format()?;
260        let cluster = cluster_id(doc_id);
261        self.store
262            .get(&keys::cluster_key(
263                &self.table,
264                keys::SCORE,
265                field,
266                term,
267                cluster,
268            )?)?
269            .map_or(Ok(0), |score| {
270                require_score_version(&score)?;
271                let entries = decode_all_scores(cluster, &score)?;
272                Ok(entries
273                    .binary_search_by_key(&doc_id, |entry| entry.doc_id)
274                    .ok()
275                    .map_or(0, |position| entries[position].term_freq))
276            })
277    }
278
279    fn doc_count(&self) -> StorageBackendResult<u64> {
280        self.require_graph_format()?;
281        let mut doc_ids = BTreeSet::new();
282        for (key, _) in self
283            .store
284            .scan_prefix(&keys::kind_prefix(&self.table, keys::LENGTH)?)?
285        {
286            doc_ids.insert(keys::read_document(&key, keys::LENGTH)?.0);
287        }
288        usize_to_u64(doc_ids.len(), "document count")
289    }
290
291    fn total_field_length(&self, field: &str) -> StorageBackendResult<u64> {
292        self.require_graph_format()?;
293        Ok(self
294            .stored_field_stats(field)?
295            .map_or(0, |stats| stats.total_length))
296    }
297
298    fn field_doc_count(&self, field: &str) -> StorageBackendResult<u64> {
299        self.require_graph_format()?;
300        Ok(self
301            .stored_field_stats(field)?
302            .map_or(0, |stats| stats.doc_count))
303    }
304
305    fn vocabulary_terms(&self, field: &str) -> StorageBackendResult<Vec<String>> {
306        self.vocabulary_keys(field)?
307            .into_iter()
308            .map(|key| Ok(key.to_term().into_string()?))
309            .collect()
310    }
311
312    fn vocabulary_keys(&self, field: &str) -> StorageBackendResult<Vec<TokenTermKey>> {
313        self.indexed_terms(Some(field))
314    }
315
316    fn stats(&self) -> StorageBackendResult<IndexStats> {
317        self.index_statistics()
318    }
319
320    fn posting_count(&self, field: Option<&str>) -> StorageBackendResult<u64> {
321        self.store
322            .scan_prefix(&self.score_prefix(field)?)?
323            .into_iter()
324            .try_fold(0_u64, |total, (key, value)| {
325                keys::read_cluster(&key, keys::SCORE)?;
326                require_score_version(&value)?;
327                total
328                    .checked_add(score_count(&value)?)
329                    .ok_or_else(|| other_error("posting count overflow"))
330            })
331    }
332
333    fn doc_length_count(&self, field: Option<&str>) -> StorageBackendResult<u64> {
334        self.require_graph_format()?;
335        if let Some(field) = field {
336            return self.field_doc_count(field);
337        }
338        let mut count = 0_u64;
339        for (key, _) in self
340            .store
341            .scan_prefix(&keys::kind_prefix(&self.table, keys::LENGTH)?)?
342        {
343            keys::read_document(&key, keys::LENGTH)?;
344            count = count
345                .checked_add(1)
346                .ok_or_else(|| other_error("document-length row count overflow"))?;
347        }
348        Ok(count)
349    }
350
351    fn term_count(&self, field: Option<&str>) -> StorageBackendResult<u64> {
352        usize_to_u64(self.indexed_terms(field)?.len(), "term count")
353    }
354
355    fn snapshot(&self) -> StorageBackendResult<Arc<dyn InvertedIndex>> {
356        Ok(Arc::new(self.clone()))
357    }
358
359    fn field_names(&self) -> StorageBackendResult<Vec<FieldName>> {
360        self.require_graph_format()?;
361        let mut fields = Vec::new();
362        for (key, value) in self
363            .store
364            .scan_prefix(&keys::kind_prefix(&self.table, keys::FIELD)?)?
365        {
366            super::FieldStats::from_bytes(&value)?;
367            fields.push(keys::read_field(&key)?);
368        }
369        fields.sort();
370        Ok(fields)
371    }
372
373    fn set_field_analyzer(
374        &mut self,
375        field: &str,
376        analyzer: Analyzer,
377        phase: AnalyzerPhase,
378    ) -> Result<(), String> {
379        let mut candidate = self.bindings.clone();
380        candidate
381            .bind(field, &analyzer, phase)
382            .map_err(|error| error.to_string())?;
383        self.validate_index_revision_change(field, &candidate)
384            .map_err(|error| error.to_string())?;
385        self.bindings = candidate;
386        Ok(())
387    }
388
389    fn remove_field_analyzers(&mut self, field: &str) -> Result<(), String> {
390        let mut candidate = self.bindings.clone();
391        candidate.remove(field);
392        self.validate_index_revision_change(field, &candidate)
393            .map_err(|error| error.to_string())?;
394        self.bindings = candidate;
395        Ok(())
396    }
397
398    fn get_field_analyzer(&self, field: &str) -> Analyzer {
399        self.bindings.index_configuration(field).clone()
400    }
401    fn get_search_analyzer(&self, field: &str) -> Analyzer {
402        self.bindings.search_configuration(field).clone()
403    }
404    fn index_analyzer_revision(
405        &self,
406        field: &str,
407    ) -> StorageBackendResult<Arc<uqa_analysis::CompiledAnalyzer>> {
408        Ok(self.bindings.index_revision(field)?)
409    }
410    fn search_analyzer_revision(
411        &self,
412        field: &str,
413    ) -> StorageBackendResult<Arc<uqa_analysis::CompiledAnalyzer>> {
414        Ok(self.bindings.search_revision(field)?)
415    }
416
417    fn set_field_analyzer_revision(
418        &mut self,
419        field: &str,
420        revision: Arc<uqa_analysis::CompiledAnalyzer>,
421        phase: AnalyzerPhase,
422    ) -> Result<(), String> {
423        let mut candidate = self.bindings.clone();
424        candidate
425            .bind_revision(field, revision, phase)
426            .map_err(|error| error.to_string())?;
427        self.validate_index_revision_change(field, &candidate)
428            .map_err(|error| error.to_string())?;
429        self.bindings = candidate;
430        Ok(())
431    }
432
433    fn set_field_analyzer_revisions(
434        &mut self,
435        field: &str,
436        index: Arc<uqa_analysis::CompiledAnalyzer>,
437        search: Arc<uqa_analysis::CompiledAnalyzer>,
438    ) -> Result<(), String> {
439        let mut candidate = self.bindings.clone();
440        candidate
441            .bind_revisions(field, index, search)
442            .map_err(|error| error.to_string())?;
443        self.validate_index_revision_change(field, &candidate)
444            .map_err(|error| error.to_string())?;
445        self.bindings = candidate;
446        Ok(())
447    }
448
449    fn rebuild_with_analyzer_revision(
450        &mut self,
451        field: &str,
452        revision: Arc<uqa_analysis::CompiledAnalyzer>,
453        phase: AnalyzerPhase,
454        documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
455    ) -> StorageBackendResult<()> {
456        let mut replacement = self.clone();
457        replacement.bindings.bind_revision(field, revision, phase)?;
458        replacement.rebuild_documents(documents)?;
459        *self = replacement;
460        Ok(())
461    }
462
463    fn rebuild_with_analyzer_revision_cancellable(
464        &mut self,
465        field: &str,
466        revision: Arc<uqa_analysis::CompiledAnalyzer>,
467        phase: AnalyzerPhase,
468        documents: Vec<(DocId, BTreeMap<FieldName, String>)>,
469        cancellation: &uqa_core::CancellationToken,
470    ) -> StorageBackendResult<()> {
471        cancellation.check()?;
472        let mut replacement = self.clone();
473        replacement.bindings.bind_revision(field, revision, phase)?;
474        replacement.rebuild_documents_inner(documents, Some(cancellation))?;
475        *self = replacement;
476        Ok(())
477    }
478}