1use 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#[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}