1#[cfg(test)]
4use std::cell::Cell;
5use std::{
6 collections::{BTreeMap, BTreeSet},
7 ops::Bound,
8 path::Path,
9 time::{Duration, Instant},
10};
11
12use hyphae_core::{
13 Q15Vector, VectorMetric, VectorSpaceDefinition, VectorSpaceName, VectorValueError,
14};
15use hyphae_query::{DocumentError, FieldPath, Value, decode_document};
16use hyphae_retrieval::{
17 ExactRetrievalError, LexicalError, LexicalField, LexicalIndexDefinition,
18 LexicalMaterializedCorpus, LexicalMaterializedDocument, tokenize_v1_checked,
19};
20use redb::{Database, Durability, ReadableDatabase, ReadableTable, TableDefinition};
21use thiserror::Error;
22
23use uuid::Uuid;
24
25use crate::{
26 CommitReceipt, Mutation, MutationError, RecoveredTransaction, RecoveryLimits, RecoveryReport,
27 StorageLimitError,
28 limits::OperationDeadline,
29 snapshot::{
30 SnapshotError, SnapshotInfo, SnapshotReadLimits, SnapshotRecordVisitor,
31 read_snapshot_records_with_policy,
32 },
33};
34
35const KV: TableDefinition<&[u8], &[u8]> = TableDefinition::new("hyphae_kv_v1");
36const METADATA: TableDefinition<&str, &[u8]> = TableDefinition::new("hyphae_metadata_v1");
37const IDEMPOTENCY: TableDefinition<&[u8], &[u8]> = TableDefinition::new("hyphae_idempotency_v1");
38const VECTOR_SPACES: TableDefinition<&str, &[u8]> = TableDefinition::new("hyphae_vector_spaces_v1");
39const VECTORS: TableDefinition<&[u8], &[u8]> = TableDefinition::new("hyphae_vectors_v1");
40const LEXICAL_INDEXES: TableDefinition<&str, &[u8]> =
41 TableDefinition::new("hyphae_lexical_indexes_v1");
42const LEXICAL_DOCUMENTS: TableDefinition<&[u8], &[u8]> =
43 TableDefinition::new("hyphae_lexical_documents_v1");
44const LEXICAL_POSTINGS: TableDefinition<&[u8], &[u8]> =
45 TableDefinition::new("hyphae_lexical_postings_v1");
46const LEXICAL_STATS: TableDefinition<&str, &[u8]> = TableDefinition::new("hyphae_lexical_stats_v1");
47const APPLIED_SEQUENCE: &str = "applied_sequence";
48const APPLIED_DIGEST: &str = "applied_digest";
49const RECEIPT_LENGTH: usize = 72;
50type RawKvEntry = (Vec<u8>, Vec<u8>);
51
52pub(crate) enum VectorScanError {
53 Index(MaterializedIndexError),
54 ExactRetrieval(ExactRetrievalError),
55}
56
57impl From<MaterializedIndexError> for VectorScanError {
58 fn from(source: MaterializedIndexError) -> Self {
59 Self::Index(source)
60 }
61}
62
63impl From<ExactRetrievalError> for VectorScanError {
64 fn from(source: ExactRetrievalError) -> Self {
65 Self::ExactRetrieval(source)
66 }
67}
68
69impl From<redb::TransactionError> for VectorScanError {
70 fn from(source: redb::TransactionError) -> Self {
71 Self::Index(MaterializedIndexError::from(source))
72 }
73}
74
75impl From<redb::TableError> for VectorScanError {
76 fn from(source: redb::TableError) -> Self {
77 Self::Index(MaterializedIndexError::from(source))
78 }
79}
80
81impl From<redb::StorageError> for VectorScanError {
82 fn from(source: redb::StorageError) -> Self {
83 Self::Index(MaterializedIndexError::from(source))
84 }
85}
86
87pub(crate) enum KvScanError {
88 Index(MaterializedIndexError),
89 ByteBudgetExceeded { maximum: u64 },
90}
91
92impl From<MaterializedIndexError> for KvScanError {
93 fn from(source: MaterializedIndexError) -> Self {
94 Self::Index(source)
95 }
96}
97
98impl From<redb::TransactionError> for KvScanError {
99 fn from(source: redb::TransactionError) -> Self {
100 Self::Index(MaterializedIndexError::from(source))
101 }
102}
103
104impl From<redb::TableError> for KvScanError {
105 fn from(source: redb::TableError) -> Self {
106 Self::Index(MaterializedIndexError::from(source))
107 }
108}
109
110impl From<redb::StorageError> for KvScanError {
111 fn from(source: redb::StorageError) -> Self {
112 Self::Index(MaterializedIndexError::from(source))
113 }
114}
115
116#[derive(Clone, Debug, Eq, PartialEq)]
117struct LexicalDocumentProjection {
118 field_lengths: Vec<u64>,
119 terms: BTreeMap<String, Vec<u64>>,
120}
121
122#[derive(Clone, Debug, Eq, PartialEq)]
123struct LexicalCorpusProjection {
124 document_count: u64,
125 token_count: u64,
126 total_field_lengths: Vec<u64>,
127}
128
129struct LexicalPreflightState {
130 definition: LexicalIndexDefinition,
131 corpus: LexicalCorpusProjection,
132 persisted: bool,
133 document_overrides: BTreeMap<Vec<u8>, Option<LexicalDocumentProjection>>,
134}
135
136#[derive(Clone, Debug, Eq, PartialEq)]
138pub struct VectorEntry {
139 pub key: Vec<u8>,
141 pub vector: Q15Vector,
143}
144
145#[derive(Debug, Error)]
147pub enum MaterializedIndexError {
148 #[error("failed to open materialized index: {0}")]
150 Database(#[from] redb::DatabaseError),
151
152 #[error("failed to begin materialized-index transaction: {0}")]
154 Transaction(#[from] redb::TransactionError),
155
156 #[error("failed to open materialized-index table: {0}")]
158 Table(#[from] redb::TableError),
159
160 #[error("materialized-index storage failure: {0}")]
162 Storage(#[from] redb::StorageError),
163
164 #[error("failed to commit materialized-index transaction: {0}")]
166 Commit(#[from] redb::CommitError),
167
168 #[error("failed to select materialized-index durability: {0}")]
170 Durability(#[from] redb::SetDurabilityError),
171
172 #[error("invalid committed mutation: {0}")]
174 Mutation(#[from] MutationError),
175
176 #[error("malformed materialized-index checkpoint")]
178 MalformedCheckpoint,
179
180 #[error("materialized index checkpoint at sequence {sequence} diverges from the log")]
182 Diverged {
183 sequence: u64,
185 },
186
187 #[error("materialized idempotency receipt for {transaction_id} diverges from the log")]
189 IdempotencyDiverged {
190 transaction_id: Uuid,
192 },
193
194 #[error("materialized idempotency key is malformed")]
196 MalformedIdempotencyKey,
197
198 #[error(transparent)]
200 Vector(#[from] VectorValueError),
201
202 #[error("vector space `{name}` is not defined")]
204 UnknownVectorSpace {
205 name: String,
207 },
208
209 #[error("vector space `{name}` already exists with a different definition")]
211 VectorSpaceConflict {
212 name: String,
214 },
215
216 #[error("malformed materialized vector index")]
218 MalformedVectorIndex,
219
220 #[error(transparent)]
222 Lexical(#[from] LexicalError),
223
224 #[error("lexical index `{name}` already exists with a different definition")]
226 LexicalIndexConflict {
227 name: String,
229 },
230
231 #[error("lexical index `{name}` is not defined")]
233 UnknownLexicalIndex {
234 name: String,
236 },
237
238 #[error("malformed materialized lexical index")]
240 MalformedLexicalIndex,
241
242 #[error("malformed materialized lexical projection")]
244 MalformedLexicalProjection,
245
246 #[error(transparent)]
248 Document(#[from] DocumentError),
249
250 #[error("vector candidate budget exceeded: {maximum}")]
252 VectorCandidateBudgetExceeded {
253 maximum: u64,
255 },
256
257 #[error("vector candidate byte budget exceeded: {maximum}")]
259 VectorByteBudgetExceeded {
260 maximum: u64,
262 },
263
264 #[cfg(test)]
266 #[error("injected materialized-index failure")]
267 InjectedFailure,
268}
269
270impl From<StorageLimitError> for MaterializedIndexError {
271 fn from(source: StorageLimitError) -> Self {
272 let source = match source {
273 StorageLimitError::TimedOut => LexicalError::TimedOut,
274 StorageLimitError::LexicalDocumentsExceeded { maximum } => {
275 LexicalError::DocumentBudgetExceeded { maximum }
276 }
277 StorageLimitError::LexicalTokensExceeded { maximum } => {
278 LexicalError::TokenBudgetExceeded { maximum }
279 }
280 source => {
281 debug_assert!(
282 false,
283 "non-lexical storage limit reached inside materialized index: {source}"
284 );
285 LexicalError::TimedOut
286 }
287 };
288 Self::Lexical(source)
289 }
290}
291
292#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
293pub(crate) struct IndexCheckpoint {
294 pub(crate) sequence: u64,
295 pub(crate) digest: Option<[u8; 32]>,
296}
297
298#[derive(Debug)]
299pub(crate) struct MaterializedIndex {
300 database: Database,
301 #[cfg(test)]
302 fail_next_apply: Cell<bool>,
303}
304
305impl MaterializedIndex {
306 pub(crate) fn open(path: impl AsRef<Path>) -> Result<Self, MaterializedIndexError> {
307 let database = Database::create(path)?;
308 let mut transaction = database.begin_write()?;
309 transaction.set_durability(Durability::Immediate)?;
310 {
311 let _table = transaction.open_table(KV)?;
312 }
313 {
314 let _table = transaction.open_table(METADATA)?;
315 }
316 {
317 let _table = transaction.open_table(IDEMPOTENCY)?;
318 }
319 {
320 let _table = transaction.open_table(VECTOR_SPACES)?;
321 }
322 {
323 let _table = transaction.open_table(VECTORS)?;
324 }
325 {
326 let _table = transaction.open_table(LEXICAL_INDEXES)?;
327 }
328 {
329 let _table = transaction.open_table(LEXICAL_DOCUMENTS)?;
330 }
331 {
332 let _table = transaction.open_table(LEXICAL_POSTINGS)?;
333 }
334 {
335 let _table = transaction.open_table(LEXICAL_STATS)?;
336 }
337 transaction.commit()?;
338 Ok(Self {
339 database,
340 #[cfg(test)]
341 fail_next_apply: Cell::new(false),
342 })
343 }
344
345 #[cfg(test)]
346 pub(crate) fn restore_from_snapshot(
347 index_path: &Path,
348 snapshot_path: &Path,
349 ) -> Result<SnapshotInfo, SnapshotError> {
350 let limits = RecoveryLimits::default();
351 let deadline = OperationDeadline::new(limits.timeout);
352 Self::restore_from_snapshot_with_limits(
353 index_path,
354 snapshot_path,
355 &limits.snapshot,
356 &limits,
357 &deadline,
358 )
359 }
360
361 pub(crate) fn restore_from_snapshot_with_limits(
362 index_path: &Path,
363 snapshot_path: &Path,
364 snapshot_limits: &SnapshotReadLimits,
365 recovery_limits: &RecoveryLimits,
366 deadline: &OperationDeadline,
367 ) -> Result<SnapshotInfo, SnapshotError> {
368 deadline.check().map_err(SnapshotError::from)?;
369 let index = Self::open(index_path)?;
370 let mut write = index
371 .database
372 .begin_write()
373 .map_err(MaterializedIndexError::from)?;
374 write
375 .set_durability(Durability::Immediate)
376 .map_err(MaterializedIndexError::from)?;
377 let snapshot = {
378 let mut visitor = IndexRestoreVisitor {
379 write: &mut write,
380 limits: recovery_limits,
381 deadline,
382 };
383 read_snapshot_records_with_policy(
384 snapshot_path,
385 &mut visitor,
386 snapshot_limits,
387 deadline,
388 )?
389 };
390 {
391 let mut metadata = write
392 .open_table(METADATA)
393 .map_err(MaterializedIndexError::from)?;
394 if snapshot.checkpoint_sequence > 0 {
395 let Some(digest) = snapshot.checkpoint_digest else {
396 return Err(SnapshotError::Invalid {
397 reason: "nonempty snapshot lacks a checkpoint digest",
398 });
399 };
400 metadata
401 .insert(
402 APPLIED_SEQUENCE,
403 snapshot.checkpoint_sequence.to_le_bytes().as_slice(),
404 )
405 .map_err(MaterializedIndexError::from)?;
406 metadata
407 .insert(APPLIED_DIGEST, digest.as_slice())
408 .map_err(MaterializedIndexError::from)?;
409 }
410 }
411 deadline.check().map_err(SnapshotError::from)?;
412 write.commit().map_err(MaterializedIndexError::from)?;
413 Ok(snapshot)
414 }
415
416 #[cfg(test)]
417 pub(crate) fn replay(&self, recovery: &RecoveryReport) -> Result<u64, MaterializedIndexError> {
418 let limits = RecoveryLimits::default();
419 let deadline = OperationDeadline::new(limits.timeout);
420 self.replay_with_limits(recovery, &limits, &deadline)
421 }
422
423 pub(crate) fn replay_with_limits(
424 &self,
425 recovery: &RecoveryReport,
426 limits: &RecoveryLimits,
427 deadline: &OperationDeadline,
428 ) -> Result<u64, MaterializedIndexError> {
429 deadline.check()?;
430 let checkpoint = self.checkpoint()?;
431 if checkpoint.sequence == 0 {
432 if checkpoint.digest.is_some() || recovery.base_sequence != 0 {
433 return Err(MaterializedIndexError::MalformedCheckpoint);
434 }
435 } else if checkpoint.sequence == recovery.base_sequence {
436 if checkpoint.digest != Some(recovery.base_digest) {
437 return Err(MaterializedIndexError::Diverged {
438 sequence: checkpoint.sequence,
439 });
440 }
441 } else {
442 let Some(transaction) = recovery
443 .transactions
444 .iter()
445 .find(|transaction| transaction.receipt.commit_sequence == checkpoint.sequence)
446 else {
447 return Err(MaterializedIndexError::Diverged {
448 sequence: checkpoint.sequence,
449 });
450 };
451 if checkpoint.digest != Some(transaction.receipt.commit_digest) {
452 return Err(MaterializedIndexError::Diverged {
453 sequence: checkpoint.sequence,
454 });
455 }
456 }
457
458 self.reconcile_idempotency(recovery, deadline)?;
459
460 let mut replayed = 0_u64;
461 for transaction in recovery
462 .transactions
463 .iter()
464 .filter(|transaction| transaction.receipt.commit_sequence > checkpoint.sequence)
465 {
466 deadline.check()?;
467 self.apply_with_limits(transaction, limits, deadline)?;
468 replayed = replayed.saturating_add(1);
469 }
470 self.rebuild_missing_lexical_projections(limits, deadline)?;
471 Ok(replayed)
472 }
473
474 fn rebuild_missing_lexical_projections(
475 &self,
476 limits: &RecoveryLimits,
477 deadline: &OperationDeadline,
478 ) -> Result<(), MaterializedIndexError> {
479 deadline.check()?;
480 let definitions = {
481 let read = self.database.begin_read()?;
482 let indexes = read.open_table(LEXICAL_INDEXES)?;
483 let stats = read.open_table(LEXICAL_STATS)?;
484 let mut definitions = Vec::new();
485 for entry in indexes.iter()? {
486 deadline.check()?;
487 let (name, value) = entry?;
488 if stats.get(name.value())?.is_some() {
489 continue;
490 }
491 let name = VectorSpaceName::new(name.value().to_owned())?;
492 definitions.push(decode_lexical_index_value(&name, value.value())?);
493 }
494 definitions
495 };
496 for definition in definitions {
497 deadline.check()?;
498 let mut write = self.database.begin_write()?;
499 write.set_durability(Durability::Immediate)?;
500 build_lexical_projection(&write, &definition, limits, deadline)?;
501 deadline.check()?;
502 write.commit()?;
503 }
504 deadline.check()?;
505 Ok(())
506 }
507
508 pub(crate) fn apply_with_limits(
509 &self,
510 transaction: &RecoveredTransaction,
511 limits: &RecoveryLimits,
512 deadline: &OperationDeadline,
513 ) -> Result<(), MaterializedIndexError> {
514 deadline.check()?;
515 #[cfg(test)]
516 if self.fail_next_apply.replace(false) {
517 return Err(MaterializedIndexError::InjectedFailure);
518 }
519 let mutations = transaction
520 .operations
521 .iter()
522 .map(|operation| Mutation::decode(operation))
523 .collect::<Result<Vec<_>, _>>()?;
524
525 let mut write = self.database.begin_write()?;
526 write.set_durability(redb::Durability::Immediate)?;
527 for mutation in mutations {
528 deadline.check()?;
529 if apply_lexical_state_mutation(&write, &mutation, limits, deadline)? {
530 continue;
531 }
532 match mutation {
533 Mutation::DefineVectorSpace { definition } => {
534 apply_vector_space_definition(&write, &definition)?;
535 }
536 Mutation::UpsertVector { space, key, vector } => {
537 let definition = require_vector_space(&write, &space)?;
538 definition.validate_vector(&vector)?;
539 let composite_key = encode_vector_key(&space, &key);
540 let encoded_vector = encode_vector_value(&vector);
541 let mut table = write.open_table(VECTORS)?;
542 table.insert(composite_key.as_slice(), encoded_vector.as_slice())?;
543 }
544 Mutation::DeleteVector { space, key } => {
545 let _definition = require_vector_space(&write, &space)?;
546 let composite_key = encode_vector_key(&space, &key);
547 let mut table = write.open_table(VECTORS)?;
548 table.remove(composite_key.as_slice())?;
549 }
550 Mutation::Put { .. }
551 | Mutation::Delete { .. }
552 | Mutation::DefineLexicalIndex { .. } => unreachable!(
553 "lexical state mutations are handled before vector materialization"
554 ),
555 }
556 }
557 deadline.check()?;
558 {
559 let mut metadata = write.open_table(METADATA)?;
560 metadata.insert(
561 APPLIED_SEQUENCE,
562 transaction.receipt.commit_sequence.to_le_bytes().as_slice(),
563 )?;
564 metadata.insert(APPLIED_DIGEST, transaction.receipt.commit_digest.as_slice())?;
565 }
566 {
567 let mut idempotency = write.open_table(IDEMPOTENCY)?;
568 idempotency.insert(
569 transaction.receipt.transaction_id.as_bytes().as_slice(),
570 encode_receipt(&transaction.receipt).as_slice(),
571 )?;
572 }
573 deadline.check()?;
574 write.commit()?;
575 Ok(())
576 }
577
578 pub(crate) fn preflight_lexical_mutations(
579 &self,
580 mutations: &[Mutation],
581 limits: &RecoveryLimits,
582 deadline: &OperationDeadline,
583 ) -> Result<(), MaterializedIndexError> {
584 deadline.check()?;
585 let has_lexical_state_mutation = mutations.iter().any(|mutation| {
586 matches!(
587 mutation,
588 Mutation::Put { .. }
589 | Mutation::Delete { .. }
590 | Mutation::DefineLexicalIndex { .. }
591 )
592 });
593 if !has_lexical_state_mutation {
594 return Ok(());
595 }
596
597 let read = self.database.begin_read()?;
598 preflight_lexical_mutations_from_read(&read, mutations, limits, deadline)
599 }
600
601 pub(crate) fn validate_mutations(
602 &self,
603 mutations: &[Mutation],
604 ) -> Result<(), MaterializedIndexError> {
605 let read = self.database.begin_read()?;
606 let table = read.open_table(VECTOR_SPACES)?;
607 let lexical_table = read.open_table(LEXICAL_INDEXES)?;
608 let mut pending: std::collections::BTreeMap<VectorSpaceName, VectorSpaceDefinition> =
609 std::collections::BTreeMap::new();
610 let mut pending_lexical: std::collections::BTreeMap<
611 VectorSpaceName,
612 LexicalIndexDefinition,
613 > = std::collections::BTreeMap::new();
614 for mutation in mutations {
615 match mutation {
616 Mutation::Put { .. } | Mutation::Delete { .. } => {}
617 Mutation::DefineVectorSpace { definition } => {
618 let existing = if let Some(existing) = pending.get(&definition.name) {
619 Some(existing.clone())
620 } else {
621 table
622 .get(definition.name.as_str())?
623 .map(|encoded| {
624 decode_vector_space_value(&definition.name, encoded.value())
625 })
626 .transpose()?
627 };
628 if existing
629 .as_ref()
630 .is_some_and(|existing| existing != definition)
631 {
632 return Err(MaterializedIndexError::VectorSpaceConflict {
633 name: definition.name.as_str().to_owned(),
634 });
635 }
636 pending.insert(definition.name.clone(), definition.clone());
637 }
638 Mutation::UpsertVector { space, vector, .. } => {
639 let definition = if let Some(definition) = pending.get(space) {
640 definition.clone()
641 } else {
642 table
643 .get(space.as_str())?
644 .map(|encoded| decode_vector_space_value(space, encoded.value()))
645 .transpose()?
646 .ok_or_else(|| MaterializedIndexError::UnknownVectorSpace {
647 name: space.as_str().to_owned(),
648 })?
649 };
650 definition.validate_vector(vector)?;
651 }
652 Mutation::DeleteVector { space, .. } => {
653 let exists =
654 pending.contains_key(space) || table.get(space.as_str())?.is_some();
655 if !exists {
656 return Err(MaterializedIndexError::UnknownVectorSpace {
657 name: space.as_str().to_owned(),
658 });
659 }
660 }
661 Mutation::DefineLexicalIndex { definition } => {
662 let existing = if let Some(existing) = pending_lexical.get(&definition.name) {
663 Some(existing.clone())
664 } else {
665 lexical_table
666 .get(definition.name.as_str())?
667 .map(|encoded| {
668 decode_lexical_index_value(&definition.name, encoded.value())
669 })
670 .transpose()?
671 };
672 if existing
673 .as_ref()
674 .is_some_and(|existing| existing != definition)
675 {
676 return Err(MaterializedIndexError::LexicalIndexConflict {
677 name: definition.name.as_str().to_owned(),
678 });
679 }
680 pending_lexical.insert(definition.name.clone(), definition.clone());
681 }
682 }
683 }
684 Ok(())
685 }
686
687 pub(crate) fn get(&self, key: &[u8]) -> Result<Option<Vec<u8>>, MaterializedIndexError> {
688 let read = self.database.begin_read()?;
689 let table = read.open_table(KV)?;
690 let value = table.get(key)?.map(|value| value.value().to_vec());
691 Ok(value)
692 }
693
694 pub(crate) fn vector_space(
695 &self,
696 name: &VectorSpaceName,
697 ) -> Result<Option<VectorSpaceDefinition>, MaterializedIndexError> {
698 let read = self.database.begin_read()?;
699 let table = read.open_table(VECTOR_SPACES)?;
700 table
701 .get(name.as_str())?
702 .map(|value| decode_vector_space_value(name, value.value()))
703 .transpose()
704 }
705
706 pub(crate) fn lexical_index(
707 &self,
708 name: &VectorSpaceName,
709 ) -> Result<Option<LexicalIndexDefinition>, MaterializedIndexError> {
710 let read = self.database.begin_read()?;
711 let table = read.open_table(LEXICAL_INDEXES)?;
712 table
713 .get(name.as_str())?
714 .map(|value| decode_lexical_index_value(name, value.value()))
715 .transpose()
716 }
717
718 pub(crate) fn lexical_corpus(
719 &self,
720 definition: &LexicalIndexDefinition,
721 query_tokens: &[String],
722 max_candidates: u64,
723 timeout: Duration,
724 ) -> Result<LexicalMaterializedCorpus, MaterializedIndexError> {
725 let started = Instant::now();
726 check_lexical_timeout(started, timeout)?;
727 let read = self.database.begin_read()?;
728 let stats = read.open_table(LEXICAL_STATS)?;
729 let encoded_stats = stats
730 .get(definition.name.as_str())?
731 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
732 let corpus = decode_lexical_corpus(encoded_stats.value(), definition.fields.len())?;
733 check_lexical_timeout(started, timeout)?;
734 let postings = read.open_table(LEXICAL_POSTINGS)?;
735 let documents = read.open_table(LEXICAL_DOCUMENTS)?;
736 let mut candidate_keys = BTreeSet::<Vec<u8>>::new();
737 for token in query_tokens {
738 check_lexical_timeout(started, timeout)?;
739 let prefix = encode_lexical_posting_prefix(&definition.name, token)?;
740 let upper = prefix_upper_bound(&prefix);
741 let bounds = (
742 Bound::Included(prefix.as_slice()),
743 upper.as_deref().map_or(Bound::Unbounded, Bound::Excluded),
744 );
745 for entry in postings.range::<&[u8]>(bounds)? {
746 check_lexical_timeout(started, timeout)?;
747 let (key, _value) = entry?;
748 let candidate = decode_lexical_posting_key(key.value(), &prefix)?;
749 candidate_keys.insert(candidate);
750 if u64::try_from(candidate_keys.len()).unwrap_or(u64::MAX) > max_candidates {
751 return Err(MaterializedIndexError::Lexical(
752 LexicalError::CandidateBudgetExceeded {
753 maximum: max_candidates,
754 },
755 ));
756 }
757 }
758 }
759 let mut materialized = Vec::with_capacity(candidate_keys.len());
760 for key in candidate_keys {
761 check_lexical_timeout(started, timeout)?;
762 let encoded_key = encode_lexical_document_key(&definition.name, &key)?;
763 let encoded = documents
764 .get(encoded_key.as_slice())?
765 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
766 let projection = decode_lexical_document_with_timeout(
767 encoded.value(),
768 definition.fields.len(),
769 started,
770 timeout,
771 )?;
772 let mut term_frequencies = BTreeMap::new();
773 for token in query_tokens {
774 check_lexical_timeout(started, timeout)?;
775 term_frequencies.insert(
776 token.clone(),
777 projection
778 .terms
779 .get(token)
780 .cloned()
781 .unwrap_or_else(|| vec![0; definition.fields.len()]),
782 );
783 }
784 materialized.push(LexicalMaterializedDocument {
785 key,
786 field_lengths: projection.field_lengths,
787 term_frequencies,
788 });
789 }
790 check_lexical_timeout(started, timeout)?;
791 Ok(LexicalMaterializedCorpus {
792 document_count: corpus.document_count,
793 token_count: corpus.token_count,
794 total_field_lengths: corpus.total_field_lengths,
795 documents: materialized,
796 })
797 }
798
799 pub(crate) fn scan_vectors(
800 &self,
801 space: &VectorSpaceName,
802 max_candidates: u64,
803 max_bytes: u64,
804 ) -> Result<Vec<VectorEntry>, MaterializedIndexError> {
805 match self.scan_vectors_inner(space, max_candidates, max_bytes, None) {
806 Ok(entries) => Ok(entries),
807 Err(VectorScanError::Index(source)) => Err(source),
808 Err(VectorScanError::ExactRetrieval(_)) => {
809 unreachable!("an unbounded vector scan cannot time out")
810 }
811 }
812 }
813
814 pub(crate) fn scan_vectors_with_timeout(
815 &self,
816 space: &VectorSpaceName,
817 max_candidates: u64,
818 max_bytes: u64,
819 timeout: Duration,
820 ) -> Result<Vec<VectorEntry>, VectorScanError> {
821 self.scan_vectors_inner(
822 space,
823 max_candidates,
824 max_bytes,
825 Some((Instant::now(), timeout)),
826 )
827 }
828
829 fn scan_vectors_inner(
830 &self,
831 space: &VectorSpaceName,
832 max_candidates: u64,
833 max_bytes: u64,
834 deadline: Option<(Instant, Duration)>,
835 ) -> Result<Vec<VectorEntry>, VectorScanError> {
836 check_exact_timeout(deadline)?;
837 let read = self.database.begin_read()?;
838 let table = read.open_table(VECTORS)?;
839 let mut entries = Vec::new();
840 let mut consumed_bytes = 0_u64;
841 for entry in table.iter()? {
842 check_exact_timeout(deadline)?;
843 let (raw_key, raw_vector) = entry?;
844 let Some(key) = decode_vector_key_for_space(raw_key.value(), space)? else {
845 continue;
846 };
847 if u64::try_from(entries.len()).unwrap_or(u64::MAX) >= max_candidates {
848 return Err(MaterializedIndexError::VectorCandidateBudgetExceeded {
849 maximum: max_candidates,
850 }
851 .into());
852 }
853 let vector_bytes = raw_vector
854 .value()
855 .len()
856 .checked_sub(2)
857 .ok_or(MaterializedIndexError::MalformedVectorIndex)?;
858 let record_bytes = u64::try_from(key.len())
859 .ok()
860 .and_then(|key_bytes| {
861 u64::try_from(vector_bytes)
862 .ok()
863 .and_then(|vector_bytes| key_bytes.checked_add(vector_bytes))
864 })
865 .ok_or(MaterializedIndexError::VectorByteBudgetExceeded { maximum: max_bytes })?;
866 consumed_bytes = consumed_bytes
867 .checked_add(record_bytes)
868 .ok_or(MaterializedIndexError::VectorByteBudgetExceeded { maximum: max_bytes })?;
869 if consumed_bytes > max_bytes {
870 return Err(MaterializedIndexError::VectorByteBudgetExceeded {
871 maximum: max_bytes,
872 }
873 .into());
874 }
875 check_exact_timeout(deadline)?;
876 entries.push(VectorEntry {
877 key,
878 vector: decode_vector_value(raw_vector.value())?,
879 });
880 }
881 check_exact_timeout(deadline)?;
882 Ok(entries)
883 }
884
885 pub(crate) fn scan_after_with_byte_limit(
886 &self,
887 after: Option<&[u8]>,
888 limit: usize,
889 max_bytes: u64,
890 ) -> Result<(Vec<RawKvEntry>, bool), KvScanError> {
891 let read = self.database.begin_read()?;
892 let table = read.open_table(KV)?;
893 let bounds = (
894 after.map_or(Bound::Unbounded, Bound::Excluded),
895 Bound::Unbounded,
896 );
897 let mut entries = Vec::with_capacity(limit);
898 let mut consumed_bytes = 0_u64;
899 let mut range = table.range::<&[u8]>(bounds)?;
900 for entry in range.by_ref().take(limit) {
901 let (key, value) = entry?;
902 let entry_bytes = u64::try_from(key.value().len())
903 .ok()
904 .and_then(|key_bytes| {
905 u64::try_from(value.value().len())
906 .ok()
907 .and_then(|value_bytes| key_bytes.checked_add(value_bytes))
908 })
909 .ok_or(KvScanError::ByteBudgetExceeded { maximum: max_bytes })?;
910 consumed_bytes = consumed_bytes
911 .checked_add(entry_bytes)
912 .ok_or(KvScanError::ByteBudgetExceeded { maximum: max_bytes })?;
913 if consumed_bytes > max_bytes {
914 return Err(KvScanError::ByteBudgetExceeded { maximum: max_bytes });
915 }
916 entries.push((key.value().to_vec(), value.value().to_vec()));
917 }
918 let has_more = range.next().transpose()?.is_some();
919 Ok((entries, has_more))
920 }
921
922 #[cfg(test)]
923 pub(crate) fn inject_apply_failure(&self) {
924 self.fail_next_apply.set(true);
925 }
926
927 pub(crate) fn for_each_entry(
928 &self,
929 mut visitor: impl FnMut(&[u8], &[u8]),
930 ) -> Result<(), MaterializedIndexError> {
931 let read = self.database.begin_read()?;
932 let table = read.open_table(KV)?;
933 for entry in table.iter()? {
934 let (key, value) = entry?;
935 visitor(key.value(), value.value());
936 }
937 Ok(())
938 }
939
940 pub(crate) fn for_each_vector_space(
941 &self,
942 mut visitor: impl FnMut(&VectorSpaceDefinition),
943 ) -> Result<(), MaterializedIndexError> {
944 let read = self.database.begin_read()?;
945 let table = read.open_table(VECTOR_SPACES)?;
946 for entry in table.iter()? {
947 let (name, value) = entry?;
948 let name = VectorSpaceName::new(name.value().to_owned())?;
949 let definition = decode_vector_space_value(&name, value.value())?;
950 visitor(&definition);
951 }
952 Ok(())
953 }
954
955 pub(crate) fn for_each_vector(
956 &self,
957 mut visitor: impl FnMut(&VectorSpaceName, &[u8], &Q15Vector),
958 ) -> Result<(), MaterializedIndexError> {
959 let read = self.database.begin_read()?;
960 let table = read.open_table(VECTORS)?;
961 for entry in table.iter()? {
962 let (raw_key, raw_vector) = entry?;
963 let (space, key) = decode_vector_key(raw_key.value())?;
964 let vector = decode_vector_value(raw_vector.value())?;
965 visitor(&space, &key, &vector);
966 }
967 Ok(())
968 }
969
970 pub(crate) fn for_each_lexical_index(
971 &self,
972 mut visitor: impl FnMut(&LexicalIndexDefinition),
973 ) -> Result<(), MaterializedIndexError> {
974 let read = self.database.begin_read()?;
975 let table = read.open_table(LEXICAL_INDEXES)?;
976 for entry in table.iter()? {
977 let (name, value) = entry?;
978 let name = VectorSpaceName::new(name.value().to_owned())?;
979 let definition = decode_lexical_index_value(&name, value.value())?;
980 visitor(&definition);
981 }
982 Ok(())
983 }
984
985 pub(crate) fn receipt(
986 &self,
987 transaction_id: Uuid,
988 ) -> Result<Option<CommitReceipt>, MaterializedIndexError> {
989 let read = self.database.begin_read()?;
990 let table = read.open_table(IDEMPOTENCY)?;
991 table
992 .get(transaction_id.as_bytes().as_slice())?
993 .map(|encoded| decode_receipt(transaction_id, encoded.value()))
994 .transpose()
995 }
996
997 pub(crate) fn for_each_receipt(
998 &self,
999 mut visitor: impl FnMut(&CommitReceipt),
1000 ) -> Result<(), MaterializedIndexError> {
1001 let read = self.database.begin_read()?;
1002 let table = read.open_table(IDEMPOTENCY)?;
1003 for entry in table.iter()? {
1004 let (key, value) = entry?;
1005 let transaction_id = Uuid::from_slice(key.value())
1006 .map_err(|_| MaterializedIndexError::MalformedIdempotencyKey)?;
1007 let receipt = decode_receipt(transaction_id, value.value())?;
1008 visitor(&receipt);
1009 }
1010 Ok(())
1011 }
1012
1013 pub(crate) fn checkpoint(&self) -> Result<IndexCheckpoint, MaterializedIndexError> {
1014 let read = self.database.begin_read()?;
1015 let metadata = read.open_table(METADATA)?;
1016 let sequence = metadata
1017 .get(APPLIED_SEQUENCE)?
1018 .map(|value| decode_sequence(value.value()))
1019 .transpose()?
1020 .unwrap_or(0);
1021 let digest = metadata
1022 .get(APPLIED_DIGEST)?
1023 .map(|value| decode_digest(value.value()))
1024 .transpose()?;
1025 Ok(IndexCheckpoint { sequence, digest })
1026 }
1027
1028 fn reconcile_idempotency(
1029 &self,
1030 recovery: &RecoveryReport,
1031 deadline: &OperationDeadline,
1032 ) -> Result<(), MaterializedIndexError> {
1033 deadline.check()?;
1034 let mut write = self.database.begin_write()?;
1035 write.set_durability(Durability::Immediate)?;
1036 {
1037 let mut table = write.open_table(IDEMPOTENCY)?;
1038 for transaction in &recovery.transactions {
1039 deadline.check()?;
1040 let receipt = transaction.receipt;
1041 if let Some(encoded) = table.get(receipt.transaction_id.as_bytes().as_slice())? {
1042 let existing = decode_receipt(receipt.transaction_id, encoded.value())?;
1043 if existing != receipt {
1044 return Err(MaterializedIndexError::IdempotencyDiverged {
1045 transaction_id: receipt.transaction_id,
1046 });
1047 }
1048 } else {
1049 table.insert(
1050 receipt.transaction_id.as_bytes().as_slice(),
1051 encode_receipt(&receipt).as_slice(),
1052 )?;
1053 }
1054 }
1055 }
1056 deadline.check()?;
1057 write.commit()?;
1058 Ok(())
1059 }
1060}
1061
1062fn check_lexical_timeout(
1063 started: Instant,
1064 timeout: Duration,
1065) -> Result<(), MaterializedIndexError> {
1066 if started.elapsed() >= timeout {
1067 Err(LexicalError::TimedOut.into())
1068 } else {
1069 Ok(())
1070 }
1071}
1072
1073fn check_exact_timeout(deadline: Option<(Instant, Duration)>) -> Result<(), ExactRetrievalError> {
1074 if deadline.is_some_and(|(started, timeout)| started.elapsed() >= timeout) {
1075 Err(ExactRetrievalError::TimedOut)
1076 } else {
1077 Ok(())
1078 }
1079}
1080
1081fn apply_vector_space_definition(
1082 write: &redb::WriteTransaction,
1083 definition: &VectorSpaceDefinition,
1084) -> Result<(), MaterializedIndexError> {
1085 let encoded = encode_vector_space_value(definition);
1086 let existing = {
1087 let table = write.open_table(VECTOR_SPACES)?;
1088 table
1089 .get(definition.name.as_str())?
1090 .map(|value| value.value().to_vec())
1091 };
1092 if let Some(existing) = existing {
1093 if existing == encoded {
1094 return Ok(());
1095 }
1096 return Err(MaterializedIndexError::VectorSpaceConflict {
1097 name: definition.name.as_str().to_owned(),
1098 });
1099 }
1100 let mut table = write.open_table(VECTOR_SPACES)?;
1101 table.insert(definition.name.as_str(), encoded.as_slice())?;
1102 Ok(())
1103}
1104
1105fn apply_lexical_index_definition(
1106 write: &redb::WriteTransaction,
1107 definition: &LexicalIndexDefinition,
1108) -> Result<(), MaterializedIndexError> {
1109 let encoded = encode_lexical_index_value(definition)?;
1110 let existing = {
1111 let table = write.open_table(LEXICAL_INDEXES)?;
1112 table
1113 .get(definition.name.as_str())?
1114 .map(|value| value.value().to_vec())
1115 };
1116 if let Some(existing) = existing {
1117 if existing == encoded {
1118 return Ok(());
1119 }
1120 return Err(MaterializedIndexError::LexicalIndexConflict {
1121 name: definition.name.as_str().to_owned(),
1122 });
1123 }
1124 let mut table = write.open_table(LEXICAL_INDEXES)?;
1125 table.insert(definition.name.as_str(), encoded.as_slice())?;
1126 Ok(())
1127}
1128
1129fn preflight_lexical_mutations_from_read(
1130 read: &redb::ReadTransaction,
1131 mutations: &[Mutation],
1132 limits: &RecoveryLimits,
1133 deadline: &OperationDeadline,
1134) -> Result<(), MaterializedIndexError> {
1135 deadline.check()?;
1136 let mut states = {
1137 let definitions = read.open_table(LEXICAL_INDEXES)?;
1138 let statistics = read.open_table(LEXICAL_STATS)?;
1139 let mut states = BTreeMap::new();
1140 for entry in definitions.iter()? {
1141 deadline.check()?;
1142 let (name, encoded_definition) = entry?;
1143 let name = VectorSpaceName::new(name.value().to_owned())?;
1144 let definition = decode_lexical_index_value(&name, encoded_definition.value())?;
1145 let encoded_corpus = statistics
1146 .get(name.as_str())?
1147 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1148 let corpus = decode_lexical_corpus(encoded_corpus.value(), definition.fields.len())?;
1149 validate_lexical_corpus_limits(&corpus, limits)?;
1150 states.insert(
1151 name,
1152 LexicalPreflightState {
1153 definition,
1154 corpus,
1155 persisted: true,
1156 document_overrides: BTreeMap::new(),
1157 },
1158 );
1159 }
1160 states
1161 };
1162 let mut kv_overlay = BTreeMap::<Vec<u8>, Option<&[u8]>>::new();
1163
1164 for mutation in mutations {
1165 deadline.check()?;
1166 match mutation {
1167 Mutation::Put { key, value } => {
1168 let decoded = if states.is_empty() {
1169 None
1170 } else {
1171 Some(decode_document(value)?)
1172 };
1173 preflight_transition_lexical_key(
1174 read,
1175 &mut states,
1176 &kv_overlay,
1177 key,
1178 decoded.as_ref(),
1179 limits,
1180 deadline,
1181 )?;
1182 kv_overlay.insert(key.clone(), Some(value.as_slice()));
1183 }
1184 Mutation::Delete { key } => {
1185 preflight_transition_lexical_key(
1186 read,
1187 &mut states,
1188 &kv_overlay,
1189 key,
1190 None,
1191 limits,
1192 deadline,
1193 )?;
1194 kv_overlay.insert(key.clone(), None);
1195 }
1196 Mutation::DefineLexicalIndex { definition } => {
1197 if states.contains_key(&definition.name) {
1198 continue;
1199 }
1200 let corpus = preflight_build_lexical_corpus(
1201 read,
1202 definition,
1203 &kv_overlay,
1204 limits,
1205 deadline,
1206 )?;
1207 states.insert(
1208 definition.name.clone(),
1209 LexicalPreflightState {
1210 definition: definition.clone(),
1211 corpus,
1212 persisted: false,
1213 document_overrides: BTreeMap::new(),
1214 },
1215 );
1216 }
1217 Mutation::DefineVectorSpace { .. }
1218 | Mutation::UpsertVector { .. }
1219 | Mutation::DeleteVector { .. } => {}
1220 }
1221 }
1222 deadline.check()?;
1223 Ok(())
1224}
1225
1226fn preflight_transition_lexical_key(
1227 read: &redb::ReadTransaction,
1228 states: &mut BTreeMap<VectorSpaceName, LexicalPreflightState>,
1229 kv_overlay: &BTreeMap<Vec<u8>, Option<&[u8]>>,
1230 key: &[u8],
1231 next_value: Option<&Value>,
1232 limits: &RecoveryLimits,
1233 deadline: &OperationDeadline,
1234) -> Result<(), MaterializedIndexError> {
1235 let names = states.keys().cloned().collect::<Vec<_>>();
1236 for name in names {
1237 deadline.check()?;
1238 let previous = {
1239 let state = states
1240 .get(&name)
1241 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1242 preflight_previous_lexical_projection(read, state, kv_overlay, key, deadline)?
1243 };
1244 let state = states
1245 .get_mut(&name)
1246 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1247 let (next_corpus, next_projection) = transition_lexical_projection(
1248 &state.corpus,
1249 previous.as_ref(),
1250 next_value,
1251 &state.definition,
1252 limits,
1253 deadline,
1254 )?;
1255 state.corpus = next_corpus;
1256 state
1257 .document_overrides
1258 .insert(key.to_vec(), next_projection);
1259 }
1260 Ok(())
1261}
1262
1263fn preflight_previous_lexical_projection(
1264 read: &redb::ReadTransaction,
1265 state: &LexicalPreflightState,
1266 kv_overlay: &BTreeMap<Vec<u8>, Option<&[u8]>>,
1267 key: &[u8],
1268 deadline: &OperationDeadline,
1269) -> Result<Option<LexicalDocumentProjection>, MaterializedIndexError> {
1270 deadline.check()?;
1271 if let Some(projection) = state.document_overrides.get(key) {
1272 return Ok(projection.clone());
1273 }
1274 if state.persisted {
1275 let encoded_key = encode_lexical_document_key(&state.definition.name, key)?;
1276 return read
1277 .open_table(LEXICAL_DOCUMENTS)?
1278 .get(encoded_key.as_slice())?
1279 .map(|encoded| {
1280 decode_lexical_document_with_deadline(
1281 encoded.value(),
1282 state.definition.fields.len(),
1283 deadline,
1284 )
1285 })
1286 .transpose();
1287 }
1288
1289 if let Some(overlaid) = kv_overlay.get(key) {
1290 return overlaid
1291 .as_deref()
1292 .map(|encoded| {
1293 project_encoded_lexical_document_unbounded(encoded, &state.definition, deadline)
1294 })
1295 .transpose();
1296 }
1297 read.open_table(KV)?
1298 .get(key)?
1299 .map(|encoded| {
1300 project_encoded_lexical_document_unbounded(encoded.value(), &state.definition, deadline)
1301 })
1302 .transpose()
1303}
1304
1305fn preflight_build_lexical_corpus(
1306 read: &redb::ReadTransaction,
1307 definition: &LexicalIndexDefinition,
1308 kv_overlay: &BTreeMap<Vec<u8>, Option<&[u8]>>,
1309 limits: &RecoveryLimits,
1310 deadline: &OperationDeadline,
1311) -> Result<LexicalCorpusProjection, MaterializedIndexError> {
1312 deadline.check()?;
1313 let table = read.open_table(KV)?;
1314 let mut corpus = empty_lexical_corpus(definition.fields.len());
1315 let mut shadowed_persisted_keys = BTreeSet::new();
1316 for entry in table.iter()? {
1317 deadline.check()?;
1318 let (key, value) = entry?;
1319 if let Some(overlaid) = kv_overlay.get(key.value()) {
1320 shadowed_persisted_keys.insert(key.value().to_vec());
1321 if let Some(encoded) = overlaid {
1322 preflight_add_encoded_lexical_document(
1323 &mut corpus,
1324 encoded,
1325 definition,
1326 limits,
1327 deadline,
1328 )?;
1329 }
1330 } else {
1331 preflight_add_encoded_lexical_document(
1332 &mut corpus,
1333 value.value(),
1334 definition,
1335 limits,
1336 deadline,
1337 )?;
1338 }
1339 }
1340 for (key, overlaid) in kv_overlay {
1341 deadline.check()?;
1342 if shadowed_persisted_keys.contains(key) {
1343 continue;
1344 }
1345 if let Some(encoded) = overlaid {
1346 preflight_add_encoded_lexical_document(
1347 &mut corpus,
1348 encoded,
1349 definition,
1350 limits,
1351 deadline,
1352 )?;
1353 }
1354 }
1355 validate_lexical_corpus_limits(&corpus, limits)?;
1356 Ok(corpus)
1357}
1358
1359fn preflight_add_encoded_lexical_document(
1360 corpus: &mut LexicalCorpusProjection,
1361 encoded: &[u8],
1362 definition: &LexicalIndexDefinition,
1363 limits: &RecoveryLimits,
1364 deadline: &OperationDeadline,
1365) -> Result<(), MaterializedIndexError> {
1366 deadline.check()?;
1367 let value = decode_document(encoded)?;
1368 let (next, _projection) =
1369 transition_lexical_projection(corpus, None, Some(&value), definition, limits, deadline)?;
1370 *corpus = next;
1371 Ok(())
1372}
1373
1374fn project_encoded_lexical_document_unbounded(
1375 encoded: &[u8],
1376 definition: &LexicalIndexDefinition,
1377 deadline: &OperationDeadline,
1378) -> Result<LexicalDocumentProjection, MaterializedIndexError> {
1379 let value = decode_document(encoded)?;
1380 let mut remaining_tokens = u64::MAX;
1381 project_lexical_document(
1382 &value,
1383 definition,
1384 &mut remaining_tokens,
1385 u64::MAX,
1386 deadline,
1387 )
1388}
1389
1390fn apply_lexical_state_mutation(
1391 write: &redb::WriteTransaction,
1392 mutation: &Mutation,
1393 limits: &RecoveryLimits,
1394 deadline: &OperationDeadline,
1395) -> Result<bool, MaterializedIndexError> {
1396 deadline.check()?;
1397 match mutation {
1398 Mutation::Put { key, value } => {
1399 update_lexical_projections(write, key, Some(value), limits, deadline)?;
1400 deadline.check()?;
1401 write
1402 .open_table(KV)?
1403 .insert(key.as_slice(), value.as_slice())?;
1404 Ok(true)
1405 }
1406 Mutation::Delete { key } => {
1407 update_lexical_projections(write, key, None, limits, deadline)?;
1408 deadline.check()?;
1409 write.open_table(KV)?.remove(key.as_slice())?;
1410 Ok(true)
1411 }
1412 Mutation::DefineLexicalIndex { definition } => {
1413 apply_lexical_index_definition(write, definition)?;
1414 build_lexical_projection(write, definition, limits, deadline)?;
1415 Ok(true)
1416 }
1417 Mutation::DefineVectorSpace { .. }
1418 | Mutation::UpsertVector { .. }
1419 | Mutation::DeleteVector { .. } => Ok(false),
1420 }
1421}
1422
1423fn build_lexical_projection(
1424 write: &redb::WriteTransaction,
1425 definition: &LexicalIndexDefinition,
1426 limits: &RecoveryLimits,
1427 deadline: &OperationDeadline,
1428) -> Result<(), MaterializedIndexError> {
1429 deadline.check()?;
1430 if write
1431 .open_table(LEXICAL_STATS)?
1432 .get(definition.name.as_str())?
1433 .is_some()
1434 {
1435 return Ok(());
1436 }
1437 let mut corpus = empty_lexical_corpus(definition.fields.len());
1438 let mut after = None;
1439 loop {
1440 deadline.check()?;
1441 let entry = {
1442 let table = write.open_table(KV)?;
1443 let bounds = (
1444 after.as_deref().map_or(Bound::Unbounded, Bound::Excluded),
1445 Bound::Unbounded,
1446 );
1447 table
1448 .range::<&[u8]>(bounds)?
1449 .next()
1450 .transpose()?
1451 .map(
1452 |(key, value)| -> Result<RawKvEntry, MaterializedIndexError> {
1453 deadline.check()?;
1454 Ok((key.value().to_vec(), value.value().to_vec()))
1455 },
1456 )
1457 .transpose()?
1458 };
1459 let Some((key, encoded)) = entry else {
1460 break;
1461 };
1462 deadline.check()?;
1463 let value = decode_document(&encoded)?;
1464 let (next_corpus, projection) = transition_lexical_projection(
1465 &corpus,
1466 None,
1467 Some(&value),
1468 definition,
1469 limits,
1470 deadline,
1471 )?;
1472 let projection = projection.ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1473 add_lexical_document(write, definition, &key, &projection, deadline)?;
1474 corpus = next_corpus;
1475 after = Some(key);
1476 }
1477 validate_lexical_corpus_limits(&corpus, limits)?;
1478 deadline.check()?;
1479 let encoded = encode_lexical_corpus(&corpus)?;
1480 write
1481 .open_table(LEXICAL_STATS)?
1482 .insert(definition.name.as_str(), encoded.as_slice())?;
1483 Ok(())
1484}
1485
1486fn update_lexical_projections(
1487 write: &redb::WriteTransaction,
1488 key: &[u8],
1489 encoded_document: Option<&[u8]>,
1490 limits: &RecoveryLimits,
1491 deadline: &OperationDeadline,
1492) -> Result<(), MaterializedIndexError> {
1493 deadline.check()?;
1494 let definitions = {
1495 let table = write.open_table(LEXICAL_INDEXES)?;
1496 let mut definitions = Vec::new();
1497 for entry in table.iter()? {
1498 deadline.check()?;
1499 let (name, value) = entry?;
1500 let name = VectorSpaceName::new(name.value().to_owned())?;
1501 definitions.push(decode_lexical_index_value(&name, value.value())?);
1502 }
1503 definitions
1504 };
1505 if definitions.is_empty() {
1506 return Ok(());
1507 }
1508 deadline.check()?;
1509 let decoded = encoded_document.map(decode_document).transpose()?;
1510 for definition in definitions {
1511 deadline.check()?;
1512 let encoded_key = encode_lexical_document_key(&definition.name, key)?;
1513 let existing = {
1514 let table = write.open_table(LEXICAL_DOCUMENTS)?;
1515 table
1516 .get(encoded_key.as_slice())?
1517 .map(|value| {
1518 decode_lexical_document_with_deadline(
1519 value.value(),
1520 definition.fields.len(),
1521 deadline,
1522 )
1523 })
1524 .transpose()?
1525 };
1526 let corpus = {
1527 let table = write.open_table(LEXICAL_STATS)?;
1528 let encoded = table
1529 .get(definition.name.as_str())?
1530 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1531 decode_lexical_corpus(encoded.value(), definition.fields.len())?
1532 };
1533 let (next_corpus, next_projection) = transition_lexical_projection(
1534 &corpus,
1535 existing.as_ref(),
1536 decoded.as_ref(),
1537 &definition,
1538 limits,
1539 deadline,
1540 )?;
1541 if let Some(existing) = &existing {
1542 remove_lexical_document(write, &definition, key, existing, deadline)?;
1543 }
1544 if let Some(projection) = &next_projection {
1545 add_lexical_document(write, &definition, key, projection, deadline)?;
1546 }
1547 deadline.check()?;
1548 let encoded = encode_lexical_corpus(&next_corpus)?;
1549 write
1550 .open_table(LEXICAL_STATS)?
1551 .insert(definition.name.as_str(), encoded.as_slice())?;
1552 }
1553 deadline.check()?;
1554 Ok(())
1555}
1556
1557fn transition_lexical_projection(
1558 corpus: &LexicalCorpusProjection,
1559 previous: Option<&LexicalDocumentProjection>,
1560 next_value: Option<&Value>,
1561 definition: &LexicalIndexDefinition,
1562 limits: &RecoveryLimits,
1563 deadline: &OperationDeadline,
1564) -> Result<(LexicalCorpusProjection, Option<LexicalDocumentProjection>), MaterializedIndexError> {
1565 deadline.check()?;
1566 let mut next_corpus = corpus.clone();
1567 if let Some(previous) = previous {
1568 remove_lexical_projection(&mut next_corpus, previous)?;
1569 }
1570 let next_projection = if let Some(value) = next_value {
1571 if next_corpus.document_count >= limits.max_lexical_documents {
1572 return Err(StorageLimitError::LexicalDocumentsExceeded {
1573 maximum: limits.max_lexical_documents,
1574 }
1575 .into());
1576 }
1577 let mut remaining_tokens = limits
1578 .max_lexical_tokens
1579 .checked_sub(next_corpus.token_count)
1580 .ok_or(StorageLimitError::LexicalTokensExceeded {
1581 maximum: limits.max_lexical_tokens,
1582 })?;
1583 let projection = project_lexical_document(
1584 value,
1585 definition,
1586 &mut remaining_tokens,
1587 limits.max_lexical_tokens,
1588 deadline,
1589 )?;
1590 add_lexical_projection(&mut next_corpus, &projection)?;
1591 Some(projection)
1592 } else {
1593 None
1594 };
1595 validate_lexical_corpus_limits(&next_corpus, limits)?;
1596 deadline.check()?;
1597 Ok((next_corpus, next_projection))
1598}
1599
1600fn add_lexical_projection(
1601 corpus: &mut LexicalCorpusProjection,
1602 projection: &LexicalDocumentProjection,
1603) -> Result<(), MaterializedIndexError> {
1604 corpus.document_count = corpus
1605 .document_count
1606 .checked_add(1)
1607 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1608 if corpus.total_field_lengths.len() != projection.field_lengths.len() {
1609 return Err(MaterializedIndexError::MalformedLexicalProjection);
1610 }
1611 for (total, length) in corpus
1612 .total_field_lengths
1613 .iter_mut()
1614 .zip(&projection.field_lengths)
1615 {
1616 *total = total
1617 .checked_add(*length)
1618 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1619 corpus.token_count = corpus
1620 .token_count
1621 .checked_add(*length)
1622 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1623 }
1624 Ok(())
1625}
1626
1627fn remove_lexical_projection(
1628 corpus: &mut LexicalCorpusProjection,
1629 projection: &LexicalDocumentProjection,
1630) -> Result<(), MaterializedIndexError> {
1631 corpus.document_count = corpus
1632 .document_count
1633 .checked_sub(1)
1634 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1635 if corpus.total_field_lengths.len() != projection.field_lengths.len() {
1636 return Err(MaterializedIndexError::MalformedLexicalProjection);
1637 }
1638 for (total, length) in corpus
1639 .total_field_lengths
1640 .iter_mut()
1641 .zip(&projection.field_lengths)
1642 {
1643 *total = total
1644 .checked_sub(*length)
1645 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1646 corpus.token_count = corpus
1647 .token_count
1648 .checked_sub(*length)
1649 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1650 }
1651 Ok(())
1652}
1653
1654fn project_lexical_document(
1655 value: &Value,
1656 definition: &LexicalIndexDefinition,
1657 remaining_tokens: &mut u64,
1658 maximum_tokens: u64,
1659 deadline: &OperationDeadline,
1660) -> Result<LexicalDocumentProjection, MaterializedIndexError> {
1661 deadline.check()?;
1662 let mut fields = Vec::with_capacity(definition.fields.len());
1663 for field in &definition.fields {
1664 deadline.check()?;
1665 let tokens = match field.path.resolve(value) {
1666 Some(Value::String(value)) => {
1667 tokenize_v1_with_policy(value, remaining_tokens, maximum_tokens, deadline)?
1668 }
1669 _ => Vec::new(),
1670 };
1671 fields.push(tokens);
1672 }
1673 let field_lengths = fields
1674 .iter()
1675 .map(|tokens| u64::try_from(tokens.len()).unwrap_or(u64::MAX))
1676 .collect();
1677 let mut terms = BTreeMap::<String, Vec<u64>>::new();
1678 for (field_index, tokens) in fields.iter().enumerate() {
1679 deadline.check()?;
1680 for token in tokens {
1681 deadline.check()?;
1682 let frequencies = terms
1683 .entry(token.clone())
1684 .or_insert_with(|| vec![0; definition.fields.len()]);
1685 frequencies[field_index] = frequencies[field_index].saturating_add(1);
1686 }
1687 }
1688 deadline.check()?;
1689 Ok(LexicalDocumentProjection {
1690 field_lengths,
1691 terms,
1692 })
1693}
1694
1695fn tokenize_v1_with_policy(
1696 input: &str,
1697 remaining_tokens: &mut u64,
1698 maximum_tokens: u64,
1699 deadline: &OperationDeadline,
1700) -> Result<Vec<String>, MaterializedIndexError> {
1701 tokenize_v1_checked(
1702 input,
1703 || -> Result<(), MaterializedIndexError> {
1704 deadline.check()?;
1705 Ok(())
1706 },
1707 || -> Result<(), MaterializedIndexError> {
1708 *remaining_tokens = remaining_tokens.checked_sub(1).ok_or(
1709 StorageLimitError::LexicalTokensExceeded {
1710 maximum: maximum_tokens,
1711 },
1712 )?;
1713 Ok(())
1714 },
1715 )
1716}
1717
1718fn add_lexical_document(
1719 write: &redb::WriteTransaction,
1720 definition: &LexicalIndexDefinition,
1721 key: &[u8],
1722 projection: &LexicalDocumentProjection,
1723 deadline: &OperationDeadline,
1724) -> Result<(), MaterializedIndexError> {
1725 deadline.check()?;
1726 {
1727 let mut postings = write.open_table(LEXICAL_POSTINGS)?;
1728 for term in projection.terms.keys() {
1729 deadline.check()?;
1730 let posting_key = encode_lexical_posting_key(&definition.name, term, key)?;
1731 postings.insert(posting_key.as_slice(), [1_u8].as_slice())?;
1732 }
1733 }
1734 deadline.check()?;
1735 let encoded_key = encode_lexical_document_key(&definition.name, key)?;
1736 let encoded_projection = encode_lexical_document_with_deadline(projection, deadline)?;
1737 write
1738 .open_table(LEXICAL_DOCUMENTS)?
1739 .insert(encoded_key.as_slice(), encoded_projection.as_slice())?;
1740 Ok(())
1741}
1742
1743fn remove_lexical_document(
1744 write: &redb::WriteTransaction,
1745 definition: &LexicalIndexDefinition,
1746 key: &[u8],
1747 projection: &LexicalDocumentProjection,
1748 deadline: &OperationDeadline,
1749) -> Result<(), MaterializedIndexError> {
1750 deadline.check()?;
1751 {
1752 let mut postings = write.open_table(LEXICAL_POSTINGS)?;
1753 for term in projection.terms.keys() {
1754 deadline.check()?;
1755 let posting_key = encode_lexical_posting_key(&definition.name, term, key)?;
1756 if postings.remove(posting_key.as_slice())?.is_none() {
1757 return Err(MaterializedIndexError::MalformedLexicalProjection);
1758 }
1759 }
1760 }
1761 let encoded_key = encode_lexical_document_key(&definition.name, key)?;
1762 if write
1763 .open_table(LEXICAL_DOCUMENTS)?
1764 .remove(encoded_key.as_slice())?
1765 .is_none()
1766 {
1767 return Err(MaterializedIndexError::MalformedLexicalProjection);
1768 }
1769 deadline.check()?;
1770 Ok(())
1771}
1772
1773fn empty_lexical_corpus(field_count: usize) -> LexicalCorpusProjection {
1774 LexicalCorpusProjection {
1775 document_count: 0,
1776 token_count: 0,
1777 total_field_lengths: vec![0; field_count],
1778 }
1779}
1780
1781fn validate_lexical_corpus_limits(
1782 corpus: &LexicalCorpusProjection,
1783 limits: &RecoveryLimits,
1784) -> Result<(), MaterializedIndexError> {
1785 if corpus.document_count > limits.max_lexical_documents {
1786 return Err(StorageLimitError::LexicalDocumentsExceeded {
1787 maximum: limits.max_lexical_documents,
1788 }
1789 .into());
1790 }
1791 if corpus.token_count > limits.max_lexical_tokens {
1792 return Err(StorageLimitError::LexicalTokensExceeded {
1793 maximum: limits.max_lexical_tokens,
1794 }
1795 .into());
1796 }
1797 Ok(())
1798}
1799
1800fn encode_lexical_document_with_deadline(
1801 projection: &LexicalDocumentProjection,
1802 deadline: &OperationDeadline,
1803) -> Result<Vec<u8>, MaterializedIndexError> {
1804 deadline.check()?;
1805 let field_count = u8::try_from(projection.field_lengths.len())
1806 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1807 let term_count = u32::try_from(projection.terms.len())
1808 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1809 let mut encoded = vec![1, field_count];
1810 for length in &projection.field_lengths {
1811 deadline.check()?;
1812 encoded.extend_from_slice(&length.to_le_bytes());
1813 }
1814 encoded.extend_from_slice(&term_count.to_le_bytes());
1815 for (term, frequencies) in &projection.terms {
1816 deadline.check()?;
1817 if frequencies.len() != projection.field_lengths.len() {
1818 return Err(MaterializedIndexError::MalformedLexicalProjection);
1819 }
1820 let length = u16::try_from(term.len())
1821 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1822 encoded.extend_from_slice(&length.to_le_bytes());
1823 encoded.extend_from_slice(term.as_bytes());
1824 for frequency in frequencies {
1825 deadline.check()?;
1826 encoded.extend_from_slice(&frequency.to_le_bytes());
1827 }
1828 }
1829 deadline.check()?;
1830 Ok(encoded)
1831}
1832
1833fn decode_lexical_document_with_deadline(
1834 encoded: &[u8],
1835 expected_fields: usize,
1836 deadline: &OperationDeadline,
1837) -> Result<LexicalDocumentProjection, MaterializedIndexError> {
1838 decode_lexical_document_inner(encoded, expected_fields, || {
1839 deadline.check().map_err(Into::into)
1840 })
1841}
1842
1843fn decode_lexical_document_with_timeout(
1844 encoded: &[u8],
1845 expected_fields: usize,
1846 started: Instant,
1847 timeout: Duration,
1848) -> Result<LexicalDocumentProjection, MaterializedIndexError> {
1849 decode_lexical_document_inner(encoded, expected_fields, || {
1850 check_lexical_timeout(started, timeout)
1851 })
1852}
1853
1854fn decode_lexical_document_inner(
1855 encoded: &[u8],
1856 expected_fields: usize,
1857 mut checkpoint: impl FnMut() -> Result<(), MaterializedIndexError>,
1858) -> Result<LexicalDocumentProjection, MaterializedIndexError> {
1859 checkpoint()?;
1860 if encoded.first() != Some(&1)
1861 || encoded.get(1).map(|value| usize::from(*value)) != Some(expected_fields)
1862 {
1863 return Err(MaterializedIndexError::MalformedLexicalProjection);
1864 }
1865 let mut cursor = 2_usize;
1866 let mut field_lengths = Vec::with_capacity(expected_fields);
1867 for _ in 0..expected_fields {
1868 checkpoint()?;
1869 field_lengths.push(read_u64(encoded, &mut cursor)?);
1870 }
1871 let term_count = usize::try_from(read_u32(encoded, &mut cursor)?)
1872 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1873 let mut terms = BTreeMap::new();
1874 let mut previous: Option<String> = None;
1875 for _ in 0..term_count {
1876 checkpoint()?;
1877 let term_length = usize::from(read_u16(encoded, &mut cursor)?);
1878 let end = cursor
1879 .checked_add(term_length)
1880 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1881 let term = std::str::from_utf8(
1882 encoded
1883 .get(cursor..end)
1884 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?,
1885 )
1886 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?
1887 .to_owned();
1888 cursor = end;
1889 if term.is_empty() || previous.as_ref().is_some_and(|previous| previous >= &term) {
1890 return Err(MaterializedIndexError::MalformedLexicalProjection);
1891 }
1892 let mut frequencies = Vec::with_capacity(expected_fields);
1893 for _ in 0..expected_fields {
1894 checkpoint()?;
1895 frequencies.push(read_u64(encoded, &mut cursor)?);
1896 }
1897 if frequencies.iter().all(|frequency| *frequency == 0) {
1898 return Err(MaterializedIndexError::MalformedLexicalProjection);
1899 }
1900 previous = Some(term.clone());
1901 terms.insert(term, frequencies);
1902 }
1903 if cursor != encoded.len() {
1904 return Err(MaterializedIndexError::MalformedLexicalProjection);
1905 }
1906 checkpoint()?;
1907 Ok(LexicalDocumentProjection {
1908 field_lengths,
1909 terms,
1910 })
1911}
1912
1913fn encode_lexical_corpus(
1914 corpus: &LexicalCorpusProjection,
1915) -> Result<Vec<u8>, MaterializedIndexError> {
1916 let field_count = u8::try_from(corpus.total_field_lengths.len())
1917 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1918 let mut encoded = vec![1, field_count];
1919 encoded.extend_from_slice(&corpus.document_count.to_le_bytes());
1920 encoded.extend_from_slice(&corpus.token_count.to_le_bytes());
1921 for length in &corpus.total_field_lengths {
1922 encoded.extend_from_slice(&length.to_le_bytes());
1923 }
1924 Ok(encoded)
1925}
1926
1927fn decode_lexical_corpus(
1928 encoded: &[u8],
1929 expected_fields: usize,
1930) -> Result<LexicalCorpusProjection, MaterializedIndexError> {
1931 if encoded.first() != Some(&1)
1932 || encoded.get(1).map(|value| usize::from(*value)) != Some(expected_fields)
1933 {
1934 return Err(MaterializedIndexError::MalformedLexicalProjection);
1935 }
1936 let mut cursor = 2_usize;
1937 let document_count = read_u64(encoded, &mut cursor)?;
1938 let token_count = read_u64(encoded, &mut cursor)?;
1939 let total_field_lengths = (0..expected_fields)
1940 .map(|_| read_u64(encoded, &mut cursor))
1941 .collect::<Result<Vec<_>, _>>()?;
1942 if cursor != encoded.len()
1943 || total_field_lengths
1944 .iter()
1945 .try_fold(0_u64, |sum, value| sum.checked_add(*value))
1946 != Some(token_count)
1947 {
1948 return Err(MaterializedIndexError::MalformedLexicalProjection);
1949 }
1950 Ok(LexicalCorpusProjection {
1951 document_count,
1952 token_count,
1953 total_field_lengths,
1954 })
1955}
1956
1957fn encode_lexical_document_key(
1958 name: &VectorSpaceName,
1959 key: &[u8],
1960) -> Result<Vec<u8>, MaterializedIndexError> {
1961 if key.is_empty() {
1962 return Err(MaterializedIndexError::MalformedLexicalProjection);
1963 }
1964 let name_length = u8::try_from(name.as_str().len())
1965 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1966 let mut encoded = Vec::with_capacity(1 + name.as_str().len() + key.len());
1967 encoded.push(name_length);
1968 encoded.extend_from_slice(name.as_str().as_bytes());
1969 encoded.extend_from_slice(key);
1970 Ok(encoded)
1971}
1972
1973fn encode_lexical_posting_prefix(
1974 name: &VectorSpaceName,
1975 term: &str,
1976) -> Result<Vec<u8>, MaterializedIndexError> {
1977 let name_length = u8::try_from(name.as_str().len())
1978 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1979 let term_length = u16::try_from(term.len())
1980 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1981 let mut encoded = Vec::with_capacity(3 + name.as_str().len() + term.len());
1982 encoded.push(name_length);
1983 encoded.extend_from_slice(name.as_str().as_bytes());
1984 encoded.extend_from_slice(&term_length.to_be_bytes());
1985 encoded.extend_from_slice(term.as_bytes());
1986 Ok(encoded)
1987}
1988
1989fn encode_lexical_posting_key(
1990 name: &VectorSpaceName,
1991 term: &str,
1992 key: &[u8],
1993) -> Result<Vec<u8>, MaterializedIndexError> {
1994 if key.is_empty() {
1995 return Err(MaterializedIndexError::MalformedLexicalProjection);
1996 }
1997 let mut encoded = encode_lexical_posting_prefix(name, term)?;
1998 encoded.extend_from_slice(key);
1999 Ok(encoded)
2000}
2001
2002fn decode_lexical_posting_key(
2003 encoded: &[u8],
2004 prefix: &[u8],
2005) -> Result<Vec<u8>, MaterializedIndexError> {
2006 let key = encoded
2007 .strip_prefix(prefix)
2008 .filter(|key| !key.is_empty())
2009 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
2010 Ok(key.to_vec())
2011}
2012
2013fn prefix_upper_bound(prefix: &[u8]) -> Option<Vec<u8>> {
2014 let mut upper = prefix.to_vec();
2015 for index in (0..upper.len()).rev() {
2016 if upper[index] != u8::MAX {
2017 upper[index] = upper[index].saturating_add(1);
2018 upper.truncate(index + 1);
2019 return Some(upper);
2020 }
2021 }
2022 None
2023}
2024
2025fn read_u16(encoded: &[u8], cursor: &mut usize) -> Result<u16, MaterializedIndexError> {
2026 read_array(encoded, cursor).map(u16::from_le_bytes)
2027}
2028
2029fn read_u32(encoded: &[u8], cursor: &mut usize) -> Result<u32, MaterializedIndexError> {
2030 read_array(encoded, cursor).map(u32::from_le_bytes)
2031}
2032
2033fn read_u64(encoded: &[u8], cursor: &mut usize) -> Result<u64, MaterializedIndexError> {
2034 read_array(encoded, cursor).map(u64::from_le_bytes)
2035}
2036
2037fn read_array<const N: usize>(
2038 encoded: &[u8],
2039 cursor: &mut usize,
2040) -> Result<[u8; N], MaterializedIndexError> {
2041 let end = cursor
2042 .checked_add(N)
2043 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
2044 let bytes = encoded
2045 .get(*cursor..end)
2046 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
2047 *cursor = end;
2048 Ok(copy_array(bytes))
2049}
2050
2051fn require_vector_space(
2052 write: &redb::WriteTransaction,
2053 name: &VectorSpaceName,
2054) -> Result<VectorSpaceDefinition, MaterializedIndexError> {
2055 let encoded = {
2056 let table = write.open_table(VECTOR_SPACES)?;
2057 table
2058 .get(name.as_str())?
2059 .map(|value| value.value().to_vec())
2060 };
2061 let Some(encoded) = encoded else {
2062 return Err(MaterializedIndexError::UnknownVectorSpace {
2063 name: name.as_str().to_owned(),
2064 });
2065 };
2066 decode_vector_space_value(name, &encoded)
2067}
2068
2069fn encode_vector_space_value(definition: &VectorSpaceDefinition) -> [u8; 4] {
2070 let dimension = definition.dimension.to_le_bytes();
2071 [dimension[0], dimension[1], definition.metric as u8, 1]
2072}
2073
2074fn decode_vector_space_value(
2075 name: &VectorSpaceName,
2076 encoded: &[u8],
2077) -> Result<VectorSpaceDefinition, MaterializedIndexError> {
2078 if encoded.len() != 4 || encoded[2] != VectorMetric::Cosine as u8 || encoded[3] != 1 {
2079 return Err(MaterializedIndexError::MalformedVectorIndex);
2080 }
2081 let dimension = u16::from_le_bytes(copy_array(&encoded[..2]));
2082 Ok(VectorSpaceDefinition::cosine(name.clone(), dimension)?)
2083}
2084
2085fn encode_lexical_index_value(
2086 definition: &LexicalIndexDefinition,
2087) -> Result<Vec<u8>, MaterializedIndexError> {
2088 let field_count = u8::try_from(definition.fields.len())
2089 .map_err(|_| MaterializedIndexError::MalformedLexicalIndex)?;
2090 let mut encoded = vec![1, field_count];
2091 for field in &definition.fields {
2092 let segment_count = u8::try_from(field.path.segments().len())
2093 .map_err(|_| MaterializedIndexError::MalformedLexicalIndex)?;
2094 encoded.push(segment_count);
2095 for segment in field.path.segments() {
2096 let length = u16::try_from(segment.len())
2097 .map_err(|_| MaterializedIndexError::MalformedLexicalIndex)?;
2098 encoded.extend_from_slice(&length.to_le_bytes());
2099 encoded.extend_from_slice(segment.as_bytes());
2100 }
2101 encoded.extend_from_slice(&field.weight_micros.to_le_bytes());
2102 }
2103 Ok(encoded)
2104}
2105
2106fn decode_lexical_index_value(
2107 name: &VectorSpaceName,
2108 encoded: &[u8],
2109) -> Result<LexicalIndexDefinition, MaterializedIndexError> {
2110 if encoded.first() != Some(&1) {
2111 return Err(MaterializedIndexError::MalformedLexicalIndex);
2112 }
2113 let field_count = usize::from(
2114 *encoded
2115 .get(1)
2116 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?,
2117 );
2118 let mut cursor = 2_usize;
2119 let mut fields = Vec::with_capacity(field_count);
2120 for _ in 0..field_count {
2121 let segment_count = usize::from(
2122 *encoded
2123 .get(cursor)
2124 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?,
2125 );
2126 cursor = cursor
2127 .checked_add(1)
2128 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?;
2129 let mut segments = Vec::with_capacity(segment_count);
2130 for _ in 0..segment_count {
2131 let length_end = cursor
2132 .checked_add(2)
2133 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?;
2134 let length = usize::from(u16::from_le_bytes(copy_array(
2135 encoded
2136 .get(cursor..length_end)
2137 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?,
2138 )));
2139 cursor = length_end;
2140 let segment_end = cursor
2141 .checked_add(length)
2142 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?;
2143 let segment = std::str::from_utf8(
2144 encoded
2145 .get(cursor..segment_end)
2146 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?,
2147 )
2148 .map_err(|_| MaterializedIndexError::MalformedLexicalIndex)?
2149 .to_owned();
2150 cursor = segment_end;
2151 segments.push(segment);
2152 }
2153 let weight_end = cursor
2154 .checked_add(4)
2155 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?;
2156 let weight_micros = u32::from_le_bytes(copy_array(
2157 encoded
2158 .get(cursor..weight_end)
2159 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?,
2160 ));
2161 cursor = weight_end;
2162 fields.push(LexicalField {
2163 path: FieldPath::new(segments),
2164 weight_micros,
2165 });
2166 }
2167 if cursor != encoded.len() {
2168 return Err(MaterializedIndexError::MalformedLexicalIndex);
2169 }
2170 LexicalIndexDefinition::new(name.clone(), fields).map_err(MaterializedIndexError::from)
2171}
2172
2173fn encode_vector_key(space: &VectorSpaceName, key: &[u8]) -> Vec<u8> {
2174 let mut encoded = Vec::with_capacity(space.as_str().len() + 1 + key.len());
2175 encoded.extend_from_slice(space.as_str().as_bytes());
2176 encoded.push(0);
2177 encoded.extend_from_slice(key);
2178 encoded
2179}
2180
2181fn decode_vector_key(encoded: &[u8]) -> Result<(VectorSpaceName, Vec<u8>), MaterializedIndexError> {
2182 let space_end = encoded
2183 .iter()
2184 .position(|byte| *byte == 0)
2185 .ok_or(MaterializedIndexError::MalformedVectorIndex)?;
2186 let space = encoded
2187 .get(..space_end)
2188 .filter(|space| !space.is_empty())
2189 .ok_or(MaterializedIndexError::MalformedVectorIndex)?;
2190 let key = encoded
2191 .get(space_end + 1..)
2192 .filter(|key| !key.is_empty())
2193 .ok_or(MaterializedIndexError::MalformedVectorIndex)?;
2194 let space =
2195 std::str::from_utf8(space).map_err(|_| MaterializedIndexError::MalformedVectorIndex)?;
2196 Ok((VectorSpaceName::new(space.to_owned())?, key.to_vec()))
2197}
2198
2199fn decode_vector_key_for_space(
2200 encoded: &[u8],
2201 expected: &VectorSpaceName,
2202) -> Result<Option<Vec<u8>>, MaterializedIndexError> {
2203 let (space, key) = decode_vector_key(encoded)?;
2204 Ok((space == *expected).then_some(key))
2205}
2206
2207fn encode_vector_value(vector: &Q15Vector) -> Vec<u8> {
2208 let mut encoded = Vec::with_capacity(2 + 2 * vector.as_slice().len());
2209 encoded.extend_from_slice(&vector.dimension().to_le_bytes());
2210 for value in vector.as_slice() {
2211 encoded.extend_from_slice(&value.to_le_bytes());
2212 }
2213 encoded
2214}
2215
2216fn decode_vector_value(encoded: &[u8]) -> Result<Q15Vector, MaterializedIndexError> {
2217 let dimension_bytes = encoded
2218 .get(..2)
2219 .ok_or(MaterializedIndexError::MalformedVectorIndex)?;
2220 let dimension = usize::from(u16::from_le_bytes(copy_array(dimension_bytes)));
2221 let expected_length = dimension
2222 .checked_mul(2)
2223 .and_then(|length| length.checked_add(2))
2224 .ok_or(MaterializedIndexError::MalformedVectorIndex)?;
2225 if encoded.len() != expected_length {
2226 return Err(MaterializedIndexError::MalformedVectorIndex);
2227 }
2228 let mut values = Vec::with_capacity(dimension);
2229 for chunk in encoded[2..].chunks_exact(2) {
2230 values.push(i16::from_le_bytes(copy_array(chunk)));
2231 }
2232 Ok(Q15Vector::new(values)?)
2233}
2234
2235struct IndexRestoreVisitor<'transaction, 'limits> {
2236 write: &'transaction mut redb::WriteTransaction,
2237 limits: &'limits RecoveryLimits,
2238 deadline: &'limits OperationDeadline,
2239}
2240
2241impl SnapshotRecordVisitor for IndexRestoreVisitor<'_, '_> {
2242 fn put(&mut self, key: &[u8], value: &[u8]) -> Result<(), SnapshotError> {
2243 update_lexical_projections(self.write, key, Some(value), self.limits, self.deadline)
2244 .map_err(SnapshotError::from)?;
2245 let mut table = self
2246 .write
2247 .open_table(KV)
2248 .map_err(MaterializedIndexError::from)?;
2249 table
2250 .insert(key, value)
2251 .map_err(MaterializedIndexError::from)?;
2252 Ok(())
2253 }
2254
2255 fn receipt(&mut self, receipt: &CommitReceipt) -> Result<(), SnapshotError> {
2256 let mut table = self
2257 .write
2258 .open_table(IDEMPOTENCY)
2259 .map_err(MaterializedIndexError::from)?;
2260 table
2261 .insert(
2262 receipt.transaction_id.as_bytes().as_slice(),
2263 encode_receipt(receipt).as_slice(),
2264 )
2265 .map_err(MaterializedIndexError::from)?;
2266 Ok(())
2267 }
2268
2269 fn vector_space(&mut self, definition: &VectorSpaceDefinition) -> Result<(), SnapshotError> {
2270 let mut table = self
2271 .write
2272 .open_table(VECTOR_SPACES)
2273 .map_err(MaterializedIndexError::from)?;
2274 let encoded = encode_vector_space_value(definition);
2275 table
2276 .insert(definition.name.as_str(), encoded.as_slice())
2277 .map_err(MaterializedIndexError::from)?;
2278 Ok(())
2279 }
2280
2281 fn lexical_index(&mut self, definition: &LexicalIndexDefinition) -> Result<(), SnapshotError> {
2282 apply_lexical_index_definition(self.write, definition).map_err(SnapshotError::from)?;
2283 build_lexical_projection(self.write, definition, self.limits, self.deadline)
2284 .map_err(SnapshotError::from)?;
2285 Ok(())
2286 }
2287
2288 fn vector(
2289 &mut self,
2290 space: &VectorSpaceName,
2291 key: &[u8],
2292 vector: &Q15Vector,
2293 ) -> Result<(), SnapshotError> {
2294 let mut table = self
2295 .write
2296 .open_table(VECTORS)
2297 .map_err(MaterializedIndexError::from)?;
2298 let encoded_key = encode_vector_key(space, key);
2299 let encoded_vector = encode_vector_value(vector);
2300 table
2301 .insert(encoded_key.as_slice(), encoded_vector.as_slice())
2302 .map_err(MaterializedIndexError::from)?;
2303 Ok(())
2304 }
2305}
2306
2307fn encode_receipt(receipt: &CommitReceipt) -> [u8; RECEIPT_LENGTH] {
2308 let mut encoded = [0_u8; RECEIPT_LENGTH];
2309 encoded[..8].copy_from_slice(&receipt.commit_sequence.to_le_bytes());
2310 encoded[8..40].copy_from_slice(&receipt.commit_digest);
2311 encoded[40..72].copy_from_slice(&receipt.transaction_digest);
2312 encoded
2313}
2314
2315fn decode_receipt(
2316 transaction_id: Uuid,
2317 encoded: &[u8],
2318) -> Result<CommitReceipt, MaterializedIndexError> {
2319 if encoded.len() != RECEIPT_LENGTH {
2320 return Err(MaterializedIndexError::IdempotencyDiverged { transaction_id });
2321 }
2322 Ok(CommitReceipt {
2323 transaction_id,
2324 commit_sequence: u64::from_le_bytes(copy_array(&encoded[..8])),
2325 commit_digest: copy_array(&encoded[8..40]),
2326 transaction_digest: copy_array(&encoded[40..72]),
2327 })
2328}
2329
2330fn decode_sequence(encoded: &[u8]) -> Result<u64, MaterializedIndexError> {
2331 if encoded.len() != 8 {
2332 return Err(MaterializedIndexError::MalformedCheckpoint);
2333 }
2334 Ok(u64::from_le_bytes(copy_array(encoded)))
2335}
2336
2337fn decode_digest(encoded: &[u8]) -> Result<[u8; 32], MaterializedIndexError> {
2338 if encoded.len() != 32 {
2339 return Err(MaterializedIndexError::MalformedCheckpoint);
2340 }
2341 Ok(copy_array(encoded))
2342}
2343
2344fn copy_array<const N: usize>(source: &[u8]) -> [u8; N] {
2345 let mut output = [0_u8; N];
2346 output.copy_from_slice(source);
2347 output
2348}
2349
2350#[cfg(test)]
2351mod tests {
2352 use std::{error::Error, time::Duration};
2353
2354 use hyphae_core::VectorSpaceName;
2355 use hyphae_retrieval::{ExactRetrievalError, LexicalError};
2356 use uuid::Uuid;
2357
2358 use super::{
2359 MaterializedIndex, MaterializedIndexError, VectorScanError, tokenize_v1_with_policy,
2360 };
2361 use crate::{DurableLog, Mutation, limits::OperationDeadline, test_support::TestDirectory};
2362
2363 fn recovery_with_operation(
2364 path: &std::path::Path,
2365 operation: Vec<u8>,
2366 ) -> Result<crate::RecoveryReport, Box<dyn Error>> {
2367 let (mut log, _) = DurableLog::open_file(path)?;
2368 log.append_transaction(Uuid::now_v7(), &[operation])?;
2369 drop(log);
2370 let (_, recovery) = DurableLog::open_file(path)?;
2371 Ok(recovery)
2372 }
2373
2374 #[test]
2375 fn checked_storage_tokenization_obeys_the_shared_deadline() {
2376 let deadline = OperationDeadline::new(Duration::ZERO);
2377 let mut remaining_tokens = u64::MAX;
2378 assert!(matches!(
2379 tokenize_v1_with_policy(
2380 &"a".repeat(2_048),
2381 &mut remaining_tokens,
2382 u64::MAX,
2383 &deadline
2384 ),
2385 Err(MaterializedIndexError::Lexical(LexicalError::TimedOut))
2386 ));
2387 }
2388
2389 #[test]
2390 fn vector_scan_checks_timeout_before_reading_candidates() -> Result<(), Box<dyn Error>> {
2391 let temporary = TestDirectory::new("index-vector-timeout")?;
2392 let index = MaterializedIndex::open(temporary.path().join("index.redb"))?;
2393 let space = VectorSpaceName::new("documents")?;
2394
2395 assert!(matches!(
2396 index.scan_vectors_with_timeout(&space, u64::MAX, u64::MAX, Duration::ZERO),
2397 Err(VectorScanError::ExactRetrieval(
2398 ExactRetrievalError::TimedOut
2399 ))
2400 ));
2401 Ok(())
2402 }
2403
2404 #[test]
2405 fn checkpoint_rejects_a_different_log_history() -> Result<(), Box<dyn Error>> {
2406 let temporary = TestDirectory::new("index-divergence")?;
2407 let first = recovery_with_operation(
2408 &temporary.path().join("first.hylog"),
2409 Mutation::put(b"key", b"first").encode()?,
2410 )?;
2411 let second = recovery_with_operation(
2412 &temporary.path().join("second.hylog"),
2413 Mutation::put(b"key", b"second").encode()?,
2414 )?;
2415 let index = MaterializedIndex::open(temporary.path().join("index.redb"))?;
2416 assert_eq!(index.replay(&first)?, 1);
2417
2418 let result = index.replay(&second);
2419 assert!(matches!(
2420 result,
2421 Err(MaterializedIndexError::Diverged { sequence: 3 })
2422 ));
2423 assert_eq!(index.get(b"key")?, Some(b"first".to_vec()));
2424 Ok(())
2425 }
2426
2427 #[test]
2428 fn invalid_committed_operation_never_advances_checkpoint() -> Result<(), Box<dyn Error>> {
2429 let temporary = TestDirectory::new("index-invalid-operation")?;
2430 let recovery = recovery_with_operation(
2431 &temporary.path().join("segment.hylog"),
2432 b"not-a-mutation".to_vec(),
2433 )?;
2434 let index = MaterializedIndex::open(temporary.path().join("index.redb"))?;
2435
2436 assert!(matches!(
2437 index.replay(&recovery),
2438 Err(MaterializedIndexError::Mutation(_))
2439 ));
2440 assert!(matches!(
2441 index.replay(&recovery),
2442 Err(MaterializedIndexError::Mutation(_))
2443 ));
2444 Ok(())
2445 }
2446
2447 #[test]
2448 fn replay_backfills_a_missing_idempotency_receipt() -> Result<(), Box<dyn Error>> {
2449 let temporary = TestDirectory::new("index-idempotency-backfill")?;
2450 let recovery = recovery_with_operation(
2451 &temporary.path().join("segment.hylog"),
2452 Mutation::put(b"key", b"value").encode()?,
2453 )?;
2454 let receipt = recovery.transactions[0].receipt;
2455 let index = MaterializedIndex::open(temporary.path().join("index.redb"))?;
2456 assert_eq!(index.replay(&recovery)?, 1);
2457 assert_eq!(index.receipt(receipt.transaction_id)?, Some(receipt));
2458
2459 let mut write = index.database.begin_write()?;
2460 write.set_durability(redb::Durability::Immediate)?;
2461 {
2462 let mut table = write.open_table(super::IDEMPOTENCY)?;
2463 table.remove(receipt.transaction_id.as_bytes().as_slice())?;
2464 }
2465 write.commit()?;
2466 assert_eq!(index.receipt(receipt.transaction_id)?, None);
2467
2468 assert_eq!(index.replay(&recovery)?, 0);
2469 assert_eq!(index.receipt(receipt.transaction_id)?, Some(receipt));
2470 Ok(())
2471 }
2472}