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
30const FORMAT_METADATA_KEY: &str = "inverted_index_format";
31const CLUSTERED_FORMAT_NAME: &str = "clustered-v1";
32const MIGRATION_PAGE_SIZE: usize = 1_024;
33
34#[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 terms.sort_unstable();
1168 store.put(
1169 &posting_document_key(&table, doc_id, &field)?,
1170 &encode_terms(&terms)?,
1171 )
1172}