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 LexicalError, LexicalField, LexicalIndexDefinition, LexicalMaterializedCorpus,
18 LexicalMaterializedDocument, tokenize_v1,
19};
20use redb::{Database, Durability, ReadableDatabase, ReadableTable, TableDefinition};
21use thiserror::Error;
22
23use uuid::Uuid;
24
25use crate::{
26 CommitReceipt, Mutation, MutationError, RecoveredTransaction, RecoveryReport,
27 snapshot::{SnapshotError, SnapshotInfo, SnapshotRecordVisitor, read_snapshot_records},
28};
29
30const KV: TableDefinition<&[u8], &[u8]> = TableDefinition::new("hyphae_kv_v1");
31const METADATA: TableDefinition<&str, &[u8]> = TableDefinition::new("hyphae_metadata_v1");
32const IDEMPOTENCY: TableDefinition<&[u8], &[u8]> = TableDefinition::new("hyphae_idempotency_v1");
33const VECTOR_SPACES: TableDefinition<&str, &[u8]> = TableDefinition::new("hyphae_vector_spaces_v1");
34const VECTORS: TableDefinition<&[u8], &[u8]> = TableDefinition::new("hyphae_vectors_v1");
35const LEXICAL_INDEXES: TableDefinition<&str, &[u8]> =
36 TableDefinition::new("hyphae_lexical_indexes_v1");
37const LEXICAL_DOCUMENTS: TableDefinition<&[u8], &[u8]> =
38 TableDefinition::new("hyphae_lexical_documents_v1");
39const LEXICAL_POSTINGS: TableDefinition<&[u8], &[u8]> =
40 TableDefinition::new("hyphae_lexical_postings_v1");
41const LEXICAL_STATS: TableDefinition<&str, &[u8]> = TableDefinition::new("hyphae_lexical_stats_v1");
42const APPLIED_SEQUENCE: &str = "applied_sequence";
43const APPLIED_DIGEST: &str = "applied_digest";
44const RECEIPT_LENGTH: usize = 72;
45type RawKvEntry = (Vec<u8>, Vec<u8>);
46
47#[derive(Clone, Debug, Eq, PartialEq)]
48struct LexicalDocumentProjection {
49 field_lengths: Vec<u64>,
50 terms: BTreeMap<String, Vec<u64>>,
51}
52
53#[derive(Clone, Debug, Eq, PartialEq)]
54struct LexicalCorpusProjection {
55 document_count: u64,
56 token_count: u64,
57 total_field_lengths: Vec<u64>,
58}
59
60#[derive(Clone, Debug, Eq, PartialEq)]
62pub struct VectorEntry {
63 pub key: Vec<u8>,
65 pub vector: Q15Vector,
67}
68
69#[derive(Debug, Error)]
71pub enum MaterializedIndexError {
72 #[error("failed to open materialized index: {0}")]
74 Database(#[from] redb::DatabaseError),
75
76 #[error("failed to begin materialized-index transaction: {0}")]
78 Transaction(#[from] redb::TransactionError),
79
80 #[error("failed to open materialized-index table: {0}")]
82 Table(#[from] redb::TableError),
83
84 #[error("materialized-index storage failure: {0}")]
86 Storage(#[from] redb::StorageError),
87
88 #[error("failed to commit materialized-index transaction: {0}")]
90 Commit(#[from] redb::CommitError),
91
92 #[error("failed to select materialized-index durability: {0}")]
94 Durability(#[from] redb::SetDurabilityError),
95
96 #[error("invalid committed mutation: {0}")]
98 Mutation(#[from] MutationError),
99
100 #[error("malformed materialized-index checkpoint")]
102 MalformedCheckpoint,
103
104 #[error("materialized index checkpoint at sequence {sequence} diverges from the log")]
106 Diverged {
107 sequence: u64,
109 },
110
111 #[error("materialized idempotency receipt for {transaction_id} diverges from the log")]
113 IdempotencyDiverged {
114 transaction_id: Uuid,
116 },
117
118 #[error("materialized idempotency key is malformed")]
120 MalformedIdempotencyKey,
121
122 #[error(transparent)]
124 Vector(#[from] VectorValueError),
125
126 #[error("vector space `{name}` is not defined")]
128 UnknownVectorSpace {
129 name: String,
131 },
132
133 #[error("vector space `{name}` already exists with a different definition")]
135 VectorSpaceConflict {
136 name: String,
138 },
139
140 #[error("malformed materialized vector index")]
142 MalformedVectorIndex,
143
144 #[error(transparent)]
146 Lexical(#[from] LexicalError),
147
148 #[error("lexical index `{name}` already exists with a different definition")]
150 LexicalIndexConflict {
151 name: String,
153 },
154
155 #[error("lexical index `{name}` is not defined")]
157 UnknownLexicalIndex {
158 name: String,
160 },
161
162 #[error("malformed materialized lexical index")]
164 MalformedLexicalIndex,
165
166 #[error("malformed materialized lexical projection")]
168 MalformedLexicalProjection,
169
170 #[error(transparent)]
172 Document(#[from] DocumentError),
173
174 #[error("vector candidate budget exceeded: {maximum}")]
176 VectorCandidateBudgetExceeded {
177 maximum: u64,
179 },
180
181 #[error("vector candidate byte budget exceeded: {maximum}")]
183 VectorByteBudgetExceeded {
184 maximum: u64,
186 },
187
188 #[cfg(test)]
190 #[error("injected materialized-index failure")]
191 InjectedFailure,
192}
193
194#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
195pub(crate) struct IndexCheckpoint {
196 pub(crate) sequence: u64,
197 pub(crate) digest: Option<[u8; 32]>,
198}
199
200#[derive(Debug)]
201pub(crate) struct MaterializedIndex {
202 database: Database,
203 #[cfg(test)]
204 fail_next_apply: Cell<bool>,
205}
206
207impl MaterializedIndex {
208 pub(crate) fn open(path: impl AsRef<Path>) -> Result<Self, MaterializedIndexError> {
209 let database = Database::create(path)?;
210 let mut transaction = database.begin_write()?;
211 transaction.set_durability(Durability::Immediate)?;
212 {
213 let _table = transaction.open_table(KV)?;
214 }
215 {
216 let _table = transaction.open_table(METADATA)?;
217 }
218 {
219 let _table = transaction.open_table(IDEMPOTENCY)?;
220 }
221 {
222 let _table = transaction.open_table(VECTOR_SPACES)?;
223 }
224 {
225 let _table = transaction.open_table(VECTORS)?;
226 }
227 {
228 let _table = transaction.open_table(LEXICAL_INDEXES)?;
229 }
230 {
231 let _table = transaction.open_table(LEXICAL_DOCUMENTS)?;
232 }
233 {
234 let _table = transaction.open_table(LEXICAL_POSTINGS)?;
235 }
236 {
237 let _table = transaction.open_table(LEXICAL_STATS)?;
238 }
239 transaction.commit()?;
240 Ok(Self {
241 database,
242 #[cfg(test)]
243 fail_next_apply: Cell::new(false),
244 })
245 }
246
247 pub(crate) fn restore_from_snapshot(
248 index_path: &Path,
249 snapshot_path: &Path,
250 ) -> Result<SnapshotInfo, SnapshotError> {
251 let index = Self::open(index_path)?;
252 let mut write = index
253 .database
254 .begin_write()
255 .map_err(MaterializedIndexError::from)?;
256 write
257 .set_durability(Durability::Immediate)
258 .map_err(MaterializedIndexError::from)?;
259 let snapshot = {
260 let mut visitor = IndexRestoreVisitor { write: &mut write };
261 read_snapshot_records(snapshot_path, &mut visitor)?
262 };
263 {
264 let mut metadata = write
265 .open_table(METADATA)
266 .map_err(MaterializedIndexError::from)?;
267 if snapshot.checkpoint_sequence > 0 {
268 let Some(digest) = snapshot.checkpoint_digest else {
269 return Err(SnapshotError::Invalid {
270 reason: "nonempty snapshot lacks a checkpoint digest",
271 });
272 };
273 metadata
274 .insert(
275 APPLIED_SEQUENCE,
276 snapshot.checkpoint_sequence.to_le_bytes().as_slice(),
277 )
278 .map_err(MaterializedIndexError::from)?;
279 metadata
280 .insert(APPLIED_DIGEST, digest.as_slice())
281 .map_err(MaterializedIndexError::from)?;
282 }
283 }
284 write.commit().map_err(MaterializedIndexError::from)?;
285 Ok(snapshot)
286 }
287
288 pub(crate) fn replay(&self, recovery: &RecoveryReport) -> Result<u64, MaterializedIndexError> {
289 let checkpoint = self.checkpoint()?;
290 if checkpoint.sequence == 0 {
291 if checkpoint.digest.is_some() || recovery.base_sequence != 0 {
292 return Err(MaterializedIndexError::MalformedCheckpoint);
293 }
294 } else if checkpoint.sequence == recovery.base_sequence {
295 if checkpoint.digest != Some(recovery.base_digest) {
296 return Err(MaterializedIndexError::Diverged {
297 sequence: checkpoint.sequence,
298 });
299 }
300 } else {
301 let Some(transaction) = recovery
302 .transactions
303 .iter()
304 .find(|transaction| transaction.receipt.commit_sequence == checkpoint.sequence)
305 else {
306 return Err(MaterializedIndexError::Diverged {
307 sequence: checkpoint.sequence,
308 });
309 };
310 if checkpoint.digest != Some(transaction.receipt.commit_digest) {
311 return Err(MaterializedIndexError::Diverged {
312 sequence: checkpoint.sequence,
313 });
314 }
315 }
316
317 self.reconcile_idempotency(recovery)?;
318
319 let mut replayed = 0_u64;
320 for transaction in recovery
321 .transactions
322 .iter()
323 .filter(|transaction| transaction.receipt.commit_sequence > checkpoint.sequence)
324 {
325 self.apply(transaction)?;
326 replayed = replayed.saturating_add(1);
327 }
328 self.rebuild_missing_lexical_projections()?;
329 Ok(replayed)
330 }
331
332 fn rebuild_missing_lexical_projections(&self) -> Result<(), MaterializedIndexError> {
333 let definitions = {
334 let read = self.database.begin_read()?;
335 let indexes = read.open_table(LEXICAL_INDEXES)?;
336 let stats = read.open_table(LEXICAL_STATS)?;
337 indexes
338 .iter()?
339 .filter_map(|entry| {
340 let (name, value) = match entry {
341 Ok(entry) => entry,
342 Err(source) => return Some(Err(MaterializedIndexError::from(source))),
343 };
344 match stats.get(name.value()) {
345 Ok(Some(_)) => None,
346 Ok(None) => {
347 let name = match VectorSpaceName::new(name.value().to_owned()) {
348 Ok(name) => name,
349 Err(source) => {
350 return Some(Err(MaterializedIndexError::from(source)));
351 }
352 };
353 Some(decode_lexical_index_value(&name, value.value()))
354 }
355 Err(source) => Some(Err(MaterializedIndexError::from(source))),
356 }
357 })
358 .collect::<Result<Vec<_>, _>>()?
359 };
360 for definition in definitions {
361 let mut write = self.database.begin_write()?;
362 write.set_durability(Durability::Immediate)?;
363 build_lexical_projection(&write, &definition)?;
364 write.commit()?;
365 }
366 Ok(())
367 }
368
369 pub(crate) fn apply(
370 &self,
371 transaction: &RecoveredTransaction,
372 ) -> Result<(), MaterializedIndexError> {
373 #[cfg(test)]
374 if self.fail_next_apply.replace(false) {
375 return Err(MaterializedIndexError::InjectedFailure);
376 }
377 let mutations = transaction
378 .operations
379 .iter()
380 .map(|operation| Mutation::decode(operation))
381 .collect::<Result<Vec<_>, _>>()?;
382
383 let mut write = self.database.begin_write()?;
384 write.set_durability(redb::Durability::Immediate)?;
385 for mutation in mutations {
386 match mutation {
387 Mutation::Put { key, value } => {
388 update_lexical_projections(&write, &key, Some(&value))?;
389 let mut table = write.open_table(KV)?;
390 table.insert(key.as_slice(), value.as_slice())?;
391 }
392 Mutation::Delete { key } => {
393 update_lexical_projections(&write, &key, None)?;
394 let mut table = write.open_table(KV)?;
395 table.remove(key.as_slice())?;
396 }
397 Mutation::DefineVectorSpace { definition } => {
398 apply_vector_space_definition(&write, &definition)?;
399 }
400 Mutation::UpsertVector { space, key, vector } => {
401 let definition = require_vector_space(&write, &space)?;
402 definition.validate_vector(&vector)?;
403 let composite_key = encode_vector_key(&space, &key);
404 let encoded_vector = encode_vector_value(&vector);
405 let mut table = write.open_table(VECTORS)?;
406 table.insert(composite_key.as_slice(), encoded_vector.as_slice())?;
407 }
408 Mutation::DeleteVector { space, key } => {
409 let _definition = require_vector_space(&write, &space)?;
410 let composite_key = encode_vector_key(&space, &key);
411 let mut table = write.open_table(VECTORS)?;
412 table.remove(composite_key.as_slice())?;
413 }
414 Mutation::DefineLexicalIndex { definition } => {
415 apply_lexical_index_definition(&write, &definition)?;
416 build_lexical_projection(&write, &definition)?;
417 }
418 }
419 }
420 {
421 let mut metadata = write.open_table(METADATA)?;
422 metadata.insert(
423 APPLIED_SEQUENCE,
424 transaction.receipt.commit_sequence.to_le_bytes().as_slice(),
425 )?;
426 metadata.insert(APPLIED_DIGEST, transaction.receipt.commit_digest.as_slice())?;
427 }
428 {
429 let mut idempotency = write.open_table(IDEMPOTENCY)?;
430 idempotency.insert(
431 transaction.receipt.transaction_id.as_bytes().as_slice(),
432 encode_receipt(&transaction.receipt).as_slice(),
433 )?;
434 }
435 write.commit()?;
436 Ok(())
437 }
438
439 pub(crate) fn validate_mutations(
440 &self,
441 mutations: &[Mutation],
442 ) -> Result<(), MaterializedIndexError> {
443 let read = self.database.begin_read()?;
444 let table = read.open_table(VECTOR_SPACES)?;
445 let lexical_table = read.open_table(LEXICAL_INDEXES)?;
446 let mut pending: std::collections::BTreeMap<VectorSpaceName, VectorSpaceDefinition> =
447 std::collections::BTreeMap::new();
448 let mut pending_lexical: std::collections::BTreeMap<
449 VectorSpaceName,
450 LexicalIndexDefinition,
451 > = std::collections::BTreeMap::new();
452 for mutation in mutations {
453 match mutation {
454 Mutation::Put { .. } | Mutation::Delete { .. } => {}
455 Mutation::DefineVectorSpace { definition } => {
456 let existing = if let Some(existing) = pending.get(&definition.name) {
457 Some(existing.clone())
458 } else {
459 table
460 .get(definition.name.as_str())?
461 .map(|encoded| {
462 decode_vector_space_value(&definition.name, encoded.value())
463 })
464 .transpose()?
465 };
466 if existing
467 .as_ref()
468 .is_some_and(|existing| existing != definition)
469 {
470 return Err(MaterializedIndexError::VectorSpaceConflict {
471 name: definition.name.as_str().to_owned(),
472 });
473 }
474 pending.insert(definition.name.clone(), definition.clone());
475 }
476 Mutation::UpsertVector { space, vector, .. } => {
477 let definition = if let Some(definition) = pending.get(space) {
478 definition.clone()
479 } else {
480 table
481 .get(space.as_str())?
482 .map(|encoded| decode_vector_space_value(space, encoded.value()))
483 .transpose()?
484 .ok_or_else(|| MaterializedIndexError::UnknownVectorSpace {
485 name: space.as_str().to_owned(),
486 })?
487 };
488 definition.validate_vector(vector)?;
489 }
490 Mutation::DeleteVector { space, .. } => {
491 let exists =
492 pending.contains_key(space) || table.get(space.as_str())?.is_some();
493 if !exists {
494 return Err(MaterializedIndexError::UnknownVectorSpace {
495 name: space.as_str().to_owned(),
496 });
497 }
498 }
499 Mutation::DefineLexicalIndex { definition } => {
500 let existing = if let Some(existing) = pending_lexical.get(&definition.name) {
501 Some(existing.clone())
502 } else {
503 lexical_table
504 .get(definition.name.as_str())?
505 .map(|encoded| {
506 decode_lexical_index_value(&definition.name, encoded.value())
507 })
508 .transpose()?
509 };
510 if existing
511 .as_ref()
512 .is_some_and(|existing| existing != definition)
513 {
514 return Err(MaterializedIndexError::LexicalIndexConflict {
515 name: definition.name.as_str().to_owned(),
516 });
517 }
518 pending_lexical.insert(definition.name.clone(), definition.clone());
519 }
520 }
521 }
522 Ok(())
523 }
524
525 pub(crate) fn get(&self, key: &[u8]) -> Result<Option<Vec<u8>>, MaterializedIndexError> {
526 let read = self.database.begin_read()?;
527 let table = read.open_table(KV)?;
528 let value = table.get(key)?.map(|value| value.value().to_vec());
529 Ok(value)
530 }
531
532 pub(crate) fn vector_space(
533 &self,
534 name: &VectorSpaceName,
535 ) -> Result<Option<VectorSpaceDefinition>, MaterializedIndexError> {
536 let read = self.database.begin_read()?;
537 let table = read.open_table(VECTOR_SPACES)?;
538 table
539 .get(name.as_str())?
540 .map(|value| decode_vector_space_value(name, value.value()))
541 .transpose()
542 }
543
544 pub(crate) fn lexical_index(
545 &self,
546 name: &VectorSpaceName,
547 ) -> Result<Option<LexicalIndexDefinition>, MaterializedIndexError> {
548 let read = self.database.begin_read()?;
549 let table = read.open_table(LEXICAL_INDEXES)?;
550 table
551 .get(name.as_str())?
552 .map(|value| decode_lexical_index_value(name, value.value()))
553 .transpose()
554 }
555
556 pub(crate) fn lexical_corpus(
557 &self,
558 definition: &LexicalIndexDefinition,
559 query_tokens: &[String],
560 max_candidates: u64,
561 timeout: Duration,
562 ) -> Result<LexicalMaterializedCorpus, MaterializedIndexError> {
563 let started = Instant::now();
564 let read = self.database.begin_read()?;
565 let stats = read.open_table(LEXICAL_STATS)?;
566 let encoded_stats = stats
567 .get(definition.name.as_str())?
568 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
569 let corpus = decode_lexical_corpus(encoded_stats.value(), definition.fields.len())?;
570 let postings = read.open_table(LEXICAL_POSTINGS)?;
571 let documents = read.open_table(LEXICAL_DOCUMENTS)?;
572 let mut candidate_keys = BTreeSet::<Vec<u8>>::new();
573 for token in query_tokens {
574 check_lexical_timeout(started, timeout)?;
575 let prefix = encode_lexical_posting_prefix(&definition.name, token)?;
576 let upper = prefix_upper_bound(&prefix);
577 let bounds = (
578 Bound::Included(prefix.as_slice()),
579 upper.as_deref().map_or(Bound::Unbounded, Bound::Excluded),
580 );
581 for entry in postings.range::<&[u8]>(bounds)? {
582 check_lexical_timeout(started, timeout)?;
583 let (key, _value) = entry?;
584 let candidate = decode_lexical_posting_key(key.value(), &prefix)?;
585 candidate_keys.insert(candidate);
586 if u64::try_from(candidate_keys.len()).unwrap_or(u64::MAX) > max_candidates {
587 return Err(MaterializedIndexError::Lexical(
588 LexicalError::CandidateBudgetExceeded {
589 maximum: max_candidates,
590 },
591 ));
592 }
593 }
594 }
595 let mut materialized = Vec::with_capacity(candidate_keys.len());
596 for key in candidate_keys {
597 check_lexical_timeout(started, timeout)?;
598 let encoded_key = encode_lexical_document_key(&definition.name, &key)?;
599 let encoded = documents
600 .get(encoded_key.as_slice())?
601 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
602 let projection = decode_lexical_document(encoded.value(), definition.fields.len())?;
603 materialized.push(LexicalMaterializedDocument {
604 key,
605 field_lengths: projection.field_lengths,
606 term_frequencies: query_tokens
607 .iter()
608 .map(|token| {
609 (
610 token.clone(),
611 projection
612 .terms
613 .get(token)
614 .cloned()
615 .unwrap_or_else(|| vec![0; definition.fields.len()]),
616 )
617 })
618 .collect(),
619 });
620 }
621 Ok(LexicalMaterializedCorpus {
622 document_count: corpus.document_count,
623 token_count: corpus.token_count,
624 total_field_lengths: corpus.total_field_lengths,
625 documents: materialized,
626 })
627 }
628
629 pub(crate) fn scan_vectors(
630 &self,
631 space: &VectorSpaceName,
632 max_candidates: u64,
633 max_bytes: u64,
634 ) -> Result<Vec<VectorEntry>, MaterializedIndexError> {
635 let read = self.database.begin_read()?;
636 let table = read.open_table(VECTORS)?;
637 let mut entries = Vec::new();
638 let mut consumed_bytes = 0_u64;
639 for entry in table.iter()? {
640 let (raw_key, raw_vector) = entry?;
641 let Some(key) = decode_vector_key_for_space(raw_key.value(), space)? else {
642 continue;
643 };
644 if u64::try_from(entries.len()).unwrap_or(u64::MAX) >= max_candidates {
645 return Err(MaterializedIndexError::VectorCandidateBudgetExceeded {
646 maximum: max_candidates,
647 });
648 }
649 let record_bytes = u64::try_from(key.len())
650 .ok()
651 .and_then(|key_bytes| {
652 u64::try_from(raw_vector.value().len())
653 .ok()
654 .and_then(|vector_bytes| key_bytes.checked_add(vector_bytes))
655 })
656 .ok_or(MaterializedIndexError::VectorByteBudgetExceeded { maximum: max_bytes })?;
657 consumed_bytes = consumed_bytes
658 .checked_add(record_bytes)
659 .ok_or(MaterializedIndexError::VectorByteBudgetExceeded { maximum: max_bytes })?;
660 if consumed_bytes > max_bytes {
661 return Err(MaterializedIndexError::VectorByteBudgetExceeded {
662 maximum: max_bytes,
663 });
664 }
665 entries.push(VectorEntry {
666 key,
667 vector: decode_vector_value(raw_vector.value())?,
668 });
669 }
670 Ok(entries)
671 }
672
673 pub(crate) fn scan_after(
674 &self,
675 after: Option<&[u8]>,
676 limit: usize,
677 ) -> Result<Vec<RawKvEntry>, MaterializedIndexError> {
678 let read = self.database.begin_read()?;
679 let table = read.open_table(KV)?;
680 let bounds = (
681 after.map_or(Bound::Unbounded, Bound::Excluded),
682 Bound::Unbounded,
683 );
684 let mut entries = Vec::with_capacity(limit);
685 for entry in table.range::<&[u8]>(bounds)?.take(limit) {
686 let (key, value) = entry?;
687 entries.push((key.value().to_vec(), value.value().to_vec()));
688 }
689 Ok(entries)
690 }
691
692 #[cfg(test)]
693 pub(crate) fn inject_apply_failure(&self) {
694 self.fail_next_apply.set(true);
695 }
696
697 pub(crate) fn for_each_entry(
698 &self,
699 mut visitor: impl FnMut(&[u8], &[u8]),
700 ) -> Result<(), MaterializedIndexError> {
701 let read = self.database.begin_read()?;
702 let table = read.open_table(KV)?;
703 for entry in table.iter()? {
704 let (key, value) = entry?;
705 visitor(key.value(), value.value());
706 }
707 Ok(())
708 }
709
710 pub(crate) fn for_each_vector_space(
711 &self,
712 mut visitor: impl FnMut(&VectorSpaceDefinition),
713 ) -> Result<(), MaterializedIndexError> {
714 let read = self.database.begin_read()?;
715 let table = read.open_table(VECTOR_SPACES)?;
716 for entry in table.iter()? {
717 let (name, value) = entry?;
718 let name = VectorSpaceName::new(name.value().to_owned())?;
719 let definition = decode_vector_space_value(&name, value.value())?;
720 visitor(&definition);
721 }
722 Ok(())
723 }
724
725 pub(crate) fn for_each_vector(
726 &self,
727 mut visitor: impl FnMut(&VectorSpaceName, &[u8], &Q15Vector),
728 ) -> Result<(), MaterializedIndexError> {
729 let read = self.database.begin_read()?;
730 let table = read.open_table(VECTORS)?;
731 for entry in table.iter()? {
732 let (raw_key, raw_vector) = entry?;
733 let (space, key) = decode_vector_key(raw_key.value())?;
734 let vector = decode_vector_value(raw_vector.value())?;
735 visitor(&space, &key, &vector);
736 }
737 Ok(())
738 }
739
740 pub(crate) fn for_each_lexical_index(
741 &self,
742 mut visitor: impl FnMut(&LexicalIndexDefinition),
743 ) -> Result<(), MaterializedIndexError> {
744 let read = self.database.begin_read()?;
745 let table = read.open_table(LEXICAL_INDEXES)?;
746 for entry in table.iter()? {
747 let (name, value) = entry?;
748 let name = VectorSpaceName::new(name.value().to_owned())?;
749 let definition = decode_lexical_index_value(&name, value.value())?;
750 visitor(&definition);
751 }
752 Ok(())
753 }
754
755 pub(crate) fn receipt(
756 &self,
757 transaction_id: Uuid,
758 ) -> Result<Option<CommitReceipt>, MaterializedIndexError> {
759 let read = self.database.begin_read()?;
760 let table = read.open_table(IDEMPOTENCY)?;
761 table
762 .get(transaction_id.as_bytes().as_slice())?
763 .map(|encoded| decode_receipt(transaction_id, encoded.value()))
764 .transpose()
765 }
766
767 pub(crate) fn for_each_receipt(
768 &self,
769 mut visitor: impl FnMut(&CommitReceipt),
770 ) -> Result<(), MaterializedIndexError> {
771 let read = self.database.begin_read()?;
772 let table = read.open_table(IDEMPOTENCY)?;
773 for entry in table.iter()? {
774 let (key, value) = entry?;
775 let transaction_id = Uuid::from_slice(key.value())
776 .map_err(|_| MaterializedIndexError::MalformedIdempotencyKey)?;
777 let receipt = decode_receipt(transaction_id, value.value())?;
778 visitor(&receipt);
779 }
780 Ok(())
781 }
782
783 pub(crate) fn checkpoint(&self) -> Result<IndexCheckpoint, MaterializedIndexError> {
784 let read = self.database.begin_read()?;
785 let metadata = read.open_table(METADATA)?;
786 let sequence = metadata
787 .get(APPLIED_SEQUENCE)?
788 .map(|value| decode_sequence(value.value()))
789 .transpose()?
790 .unwrap_or(0);
791 let digest = metadata
792 .get(APPLIED_DIGEST)?
793 .map(|value| decode_digest(value.value()))
794 .transpose()?;
795 Ok(IndexCheckpoint { sequence, digest })
796 }
797
798 fn reconcile_idempotency(
799 &self,
800 recovery: &RecoveryReport,
801 ) -> Result<(), MaterializedIndexError> {
802 let mut write = self.database.begin_write()?;
803 write.set_durability(Durability::Immediate)?;
804 {
805 let mut table = write.open_table(IDEMPOTENCY)?;
806 for transaction in &recovery.transactions {
807 let receipt = transaction.receipt;
808 if let Some(encoded) = table.get(receipt.transaction_id.as_bytes().as_slice())? {
809 let existing = decode_receipt(receipt.transaction_id, encoded.value())?;
810 if existing != receipt {
811 return Err(MaterializedIndexError::IdempotencyDiverged {
812 transaction_id: receipt.transaction_id,
813 });
814 }
815 } else {
816 table.insert(
817 receipt.transaction_id.as_bytes().as_slice(),
818 encode_receipt(&receipt).as_slice(),
819 )?;
820 }
821 }
822 }
823 write.commit()?;
824 Ok(())
825 }
826}
827
828fn check_lexical_timeout(
829 started: Instant,
830 timeout: Duration,
831) -> Result<(), MaterializedIndexError> {
832 if started.elapsed() >= timeout {
833 Err(LexicalError::TimedOut.into())
834 } else {
835 Ok(())
836 }
837}
838
839fn apply_vector_space_definition(
840 write: &redb::WriteTransaction,
841 definition: &VectorSpaceDefinition,
842) -> Result<(), MaterializedIndexError> {
843 let encoded = encode_vector_space_value(definition);
844 let existing = {
845 let table = write.open_table(VECTOR_SPACES)?;
846 table
847 .get(definition.name.as_str())?
848 .map(|value| value.value().to_vec())
849 };
850 if let Some(existing) = existing {
851 if existing == encoded {
852 return Ok(());
853 }
854 return Err(MaterializedIndexError::VectorSpaceConflict {
855 name: definition.name.as_str().to_owned(),
856 });
857 }
858 let mut table = write.open_table(VECTOR_SPACES)?;
859 table.insert(definition.name.as_str(), encoded.as_slice())?;
860 Ok(())
861}
862
863fn apply_lexical_index_definition(
864 write: &redb::WriteTransaction,
865 definition: &LexicalIndexDefinition,
866) -> Result<(), MaterializedIndexError> {
867 let encoded = encode_lexical_index_value(definition)?;
868 let existing = {
869 let table = write.open_table(LEXICAL_INDEXES)?;
870 table
871 .get(definition.name.as_str())?
872 .map(|value| value.value().to_vec())
873 };
874 if let Some(existing) = existing {
875 if existing == encoded {
876 return Ok(());
877 }
878 return Err(MaterializedIndexError::LexicalIndexConflict {
879 name: definition.name.as_str().to_owned(),
880 });
881 }
882 let mut table = write.open_table(LEXICAL_INDEXES)?;
883 table.insert(definition.name.as_str(), encoded.as_slice())?;
884 Ok(())
885}
886
887fn build_lexical_projection(
888 write: &redb::WriteTransaction,
889 definition: &LexicalIndexDefinition,
890) -> Result<(), MaterializedIndexError> {
891 if write
892 .open_table(LEXICAL_STATS)?
893 .get(definition.name.as_str())?
894 .is_some()
895 {
896 return Ok(());
897 }
898 let entries = {
899 let table = write.open_table(KV)?;
900 table
901 .iter()?
902 .map(|entry| {
903 let (key, value) = entry?;
904 Ok((key.value().to_vec(), value.value().to_vec()))
905 })
906 .collect::<Result<Vec<_>, redb::StorageError>>()?
907 };
908 let mut corpus = empty_lexical_corpus(definition.fields.len());
909 for (key, encoded) in entries {
910 let value = decode_document(&encoded)?;
911 let projection = project_lexical_document(&value, definition);
912 add_lexical_document(write, definition, &key, &projection, &mut corpus)?;
913 }
914 let encoded = encode_lexical_corpus(&corpus)?;
915 write
916 .open_table(LEXICAL_STATS)?
917 .insert(definition.name.as_str(), encoded.as_slice())?;
918 Ok(())
919}
920
921fn update_lexical_projections(
922 write: &redb::WriteTransaction,
923 key: &[u8],
924 encoded_document: Option<&[u8]>,
925) -> Result<(), MaterializedIndexError> {
926 let definitions = {
927 let table = write.open_table(LEXICAL_INDEXES)?;
928 table
929 .iter()?
930 .map(|entry| {
931 let (name, value) = entry?;
932 let name = VectorSpaceName::new(name.value().to_owned())?;
933 decode_lexical_index_value(&name, value.value())
934 })
935 .collect::<Result<Vec<_>, MaterializedIndexError>>()?
936 };
937 if definitions.is_empty() {
938 return Ok(());
939 }
940 let decoded = encoded_document.map(decode_document).transpose()?;
941 for definition in definitions {
942 let encoded_key = encode_lexical_document_key(&definition.name, key)?;
943 let existing = {
944 let table = write.open_table(LEXICAL_DOCUMENTS)?;
945 table
946 .get(encoded_key.as_slice())?
947 .map(|value| decode_lexical_document(value.value(), definition.fields.len()))
948 .transpose()?
949 };
950 let mut corpus = {
951 let table = write.open_table(LEXICAL_STATS)?;
952 let encoded = table
953 .get(definition.name.as_str())?
954 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
955 decode_lexical_corpus(encoded.value(), definition.fields.len())?
956 };
957 if let Some(existing) = existing {
958 remove_lexical_document(write, &definition, key, &existing, &mut corpus)?;
959 }
960 if let Some(value) = &decoded {
961 let projection = project_lexical_document(value, &definition);
962 add_lexical_document(write, &definition, key, &projection, &mut corpus)?;
963 }
964 let encoded = encode_lexical_corpus(&corpus)?;
965 write
966 .open_table(LEXICAL_STATS)?
967 .insert(definition.name.as_str(), encoded.as_slice())?;
968 }
969 Ok(())
970}
971
972fn project_lexical_document(
973 value: &Value,
974 definition: &LexicalIndexDefinition,
975) -> LexicalDocumentProjection {
976 let fields = definition
977 .fields
978 .iter()
979 .map(|field| match field.path.resolve(value) {
980 Some(Value::String(value)) => tokenize_v1(value),
981 _ => Vec::new(),
982 })
983 .collect::<Vec<_>>();
984 let field_lengths = fields
985 .iter()
986 .map(|tokens| u64::try_from(tokens.len()).unwrap_or(u64::MAX))
987 .collect();
988 let mut terms = BTreeMap::<String, Vec<u64>>::new();
989 for (field_index, tokens) in fields.iter().enumerate() {
990 for token in tokens {
991 let frequencies = terms
992 .entry(token.clone())
993 .or_insert_with(|| vec![0; definition.fields.len()]);
994 frequencies[field_index] = frequencies[field_index].saturating_add(1);
995 }
996 }
997 LexicalDocumentProjection {
998 field_lengths,
999 terms,
1000 }
1001}
1002
1003fn add_lexical_document(
1004 write: &redb::WriteTransaction,
1005 definition: &LexicalIndexDefinition,
1006 key: &[u8],
1007 projection: &LexicalDocumentProjection,
1008 corpus: &mut LexicalCorpusProjection,
1009) -> Result<(), MaterializedIndexError> {
1010 corpus.document_count = corpus
1011 .document_count
1012 .checked_add(1)
1013 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1014 for (total, length) in corpus
1015 .total_field_lengths
1016 .iter_mut()
1017 .zip(&projection.field_lengths)
1018 {
1019 *total = total
1020 .checked_add(*length)
1021 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1022 corpus.token_count = corpus
1023 .token_count
1024 .checked_add(*length)
1025 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1026 }
1027 {
1028 let mut postings = write.open_table(LEXICAL_POSTINGS)?;
1029 for term in projection.terms.keys() {
1030 let posting_key = encode_lexical_posting_key(&definition.name, term, key)?;
1031 postings.insert(posting_key.as_slice(), [1_u8].as_slice())?;
1032 }
1033 }
1034 let encoded_key = encode_lexical_document_key(&definition.name, key)?;
1035 let encoded_projection = encode_lexical_document(projection)?;
1036 write
1037 .open_table(LEXICAL_DOCUMENTS)?
1038 .insert(encoded_key.as_slice(), encoded_projection.as_slice())?;
1039 Ok(())
1040}
1041
1042fn remove_lexical_document(
1043 write: &redb::WriteTransaction,
1044 definition: &LexicalIndexDefinition,
1045 key: &[u8],
1046 projection: &LexicalDocumentProjection,
1047 corpus: &mut LexicalCorpusProjection,
1048) -> Result<(), MaterializedIndexError> {
1049 corpus.document_count = corpus
1050 .document_count
1051 .checked_sub(1)
1052 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1053 for (total, length) in corpus
1054 .total_field_lengths
1055 .iter_mut()
1056 .zip(&projection.field_lengths)
1057 {
1058 *total = total
1059 .checked_sub(*length)
1060 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1061 corpus.token_count = corpus
1062 .token_count
1063 .checked_sub(*length)
1064 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1065 }
1066 {
1067 let mut postings = write.open_table(LEXICAL_POSTINGS)?;
1068 for term in projection.terms.keys() {
1069 let posting_key = encode_lexical_posting_key(&definition.name, term, key)?;
1070 if postings.remove(posting_key.as_slice())?.is_none() {
1071 return Err(MaterializedIndexError::MalformedLexicalProjection);
1072 }
1073 }
1074 }
1075 let encoded_key = encode_lexical_document_key(&definition.name, key)?;
1076 if write
1077 .open_table(LEXICAL_DOCUMENTS)?
1078 .remove(encoded_key.as_slice())?
1079 .is_none()
1080 {
1081 return Err(MaterializedIndexError::MalformedLexicalProjection);
1082 }
1083 Ok(())
1084}
1085
1086fn empty_lexical_corpus(field_count: usize) -> LexicalCorpusProjection {
1087 LexicalCorpusProjection {
1088 document_count: 0,
1089 token_count: 0,
1090 total_field_lengths: vec![0; field_count],
1091 }
1092}
1093
1094fn encode_lexical_document(
1095 projection: &LexicalDocumentProjection,
1096) -> Result<Vec<u8>, MaterializedIndexError> {
1097 let field_count = u8::try_from(projection.field_lengths.len())
1098 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1099 let term_count = u32::try_from(projection.terms.len())
1100 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1101 let mut encoded = vec![1, field_count];
1102 for length in &projection.field_lengths {
1103 encoded.extend_from_slice(&length.to_le_bytes());
1104 }
1105 encoded.extend_from_slice(&term_count.to_le_bytes());
1106 for (term, frequencies) in &projection.terms {
1107 if frequencies.len() != projection.field_lengths.len() {
1108 return Err(MaterializedIndexError::MalformedLexicalProjection);
1109 }
1110 let length = u16::try_from(term.len())
1111 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1112 encoded.extend_from_slice(&length.to_le_bytes());
1113 encoded.extend_from_slice(term.as_bytes());
1114 for frequency in frequencies {
1115 encoded.extend_from_slice(&frequency.to_le_bytes());
1116 }
1117 }
1118 Ok(encoded)
1119}
1120
1121fn decode_lexical_document(
1122 encoded: &[u8],
1123 expected_fields: usize,
1124) -> Result<LexicalDocumentProjection, MaterializedIndexError> {
1125 if encoded.first() != Some(&1)
1126 || encoded.get(1).map(|value| usize::from(*value)) != Some(expected_fields)
1127 {
1128 return Err(MaterializedIndexError::MalformedLexicalProjection);
1129 }
1130 let mut cursor = 2_usize;
1131 let mut field_lengths = Vec::with_capacity(expected_fields);
1132 for _ in 0..expected_fields {
1133 field_lengths.push(read_u64(encoded, &mut cursor)?);
1134 }
1135 let term_count = usize::try_from(read_u32(encoded, &mut cursor)?)
1136 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1137 let mut terms = BTreeMap::new();
1138 let mut previous: Option<String> = None;
1139 for _ in 0..term_count {
1140 let term_length = usize::from(read_u16(encoded, &mut cursor)?);
1141 let end = cursor
1142 .checked_add(term_length)
1143 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1144 let term = std::str::from_utf8(
1145 encoded
1146 .get(cursor..end)
1147 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?,
1148 )
1149 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?
1150 .to_owned();
1151 cursor = end;
1152 if term.is_empty() || previous.as_ref().is_some_and(|previous| previous >= &term) {
1153 return Err(MaterializedIndexError::MalformedLexicalProjection);
1154 }
1155 let frequencies = (0..expected_fields)
1156 .map(|_| read_u64(encoded, &mut cursor))
1157 .collect::<Result<Vec<_>, _>>()?;
1158 if frequencies.iter().all(|frequency| *frequency == 0) {
1159 return Err(MaterializedIndexError::MalformedLexicalProjection);
1160 }
1161 previous = Some(term.clone());
1162 terms.insert(term, frequencies);
1163 }
1164 if cursor != encoded.len() {
1165 return Err(MaterializedIndexError::MalformedLexicalProjection);
1166 }
1167 Ok(LexicalDocumentProjection {
1168 field_lengths,
1169 terms,
1170 })
1171}
1172
1173fn encode_lexical_corpus(
1174 corpus: &LexicalCorpusProjection,
1175) -> Result<Vec<u8>, MaterializedIndexError> {
1176 let field_count = u8::try_from(corpus.total_field_lengths.len())
1177 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1178 let mut encoded = vec![1, field_count];
1179 encoded.extend_from_slice(&corpus.document_count.to_le_bytes());
1180 encoded.extend_from_slice(&corpus.token_count.to_le_bytes());
1181 for length in &corpus.total_field_lengths {
1182 encoded.extend_from_slice(&length.to_le_bytes());
1183 }
1184 Ok(encoded)
1185}
1186
1187fn decode_lexical_corpus(
1188 encoded: &[u8],
1189 expected_fields: usize,
1190) -> Result<LexicalCorpusProjection, MaterializedIndexError> {
1191 if encoded.first() != Some(&1)
1192 || encoded.get(1).map(|value| usize::from(*value)) != Some(expected_fields)
1193 {
1194 return Err(MaterializedIndexError::MalformedLexicalProjection);
1195 }
1196 let mut cursor = 2_usize;
1197 let document_count = read_u64(encoded, &mut cursor)?;
1198 let token_count = read_u64(encoded, &mut cursor)?;
1199 let total_field_lengths = (0..expected_fields)
1200 .map(|_| read_u64(encoded, &mut cursor))
1201 .collect::<Result<Vec<_>, _>>()?;
1202 if cursor != encoded.len()
1203 || total_field_lengths
1204 .iter()
1205 .try_fold(0_u64, |sum, value| sum.checked_add(*value))
1206 != Some(token_count)
1207 {
1208 return Err(MaterializedIndexError::MalformedLexicalProjection);
1209 }
1210 Ok(LexicalCorpusProjection {
1211 document_count,
1212 token_count,
1213 total_field_lengths,
1214 })
1215}
1216
1217fn encode_lexical_document_key(
1218 name: &VectorSpaceName,
1219 key: &[u8],
1220) -> Result<Vec<u8>, MaterializedIndexError> {
1221 if key.is_empty() {
1222 return Err(MaterializedIndexError::MalformedLexicalProjection);
1223 }
1224 let name_length = u8::try_from(name.as_str().len())
1225 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1226 let mut encoded = Vec::with_capacity(1 + name.as_str().len() + key.len());
1227 encoded.push(name_length);
1228 encoded.extend_from_slice(name.as_str().as_bytes());
1229 encoded.extend_from_slice(key);
1230 Ok(encoded)
1231}
1232
1233fn encode_lexical_posting_prefix(
1234 name: &VectorSpaceName,
1235 term: &str,
1236) -> Result<Vec<u8>, MaterializedIndexError> {
1237 let name_length = u8::try_from(name.as_str().len())
1238 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1239 let term_length = u16::try_from(term.len())
1240 .map_err(|_| MaterializedIndexError::MalformedLexicalProjection)?;
1241 let mut encoded = Vec::with_capacity(3 + name.as_str().len() + term.len());
1242 encoded.push(name_length);
1243 encoded.extend_from_slice(name.as_str().as_bytes());
1244 encoded.extend_from_slice(&term_length.to_be_bytes());
1245 encoded.extend_from_slice(term.as_bytes());
1246 Ok(encoded)
1247}
1248
1249fn encode_lexical_posting_key(
1250 name: &VectorSpaceName,
1251 term: &str,
1252 key: &[u8],
1253) -> Result<Vec<u8>, MaterializedIndexError> {
1254 if key.is_empty() {
1255 return Err(MaterializedIndexError::MalformedLexicalProjection);
1256 }
1257 let mut encoded = encode_lexical_posting_prefix(name, term)?;
1258 encoded.extend_from_slice(key);
1259 Ok(encoded)
1260}
1261
1262fn decode_lexical_posting_key(
1263 encoded: &[u8],
1264 prefix: &[u8],
1265) -> Result<Vec<u8>, MaterializedIndexError> {
1266 let key = encoded
1267 .strip_prefix(prefix)
1268 .filter(|key| !key.is_empty())
1269 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1270 Ok(key.to_vec())
1271}
1272
1273fn prefix_upper_bound(prefix: &[u8]) -> Option<Vec<u8>> {
1274 let mut upper = prefix.to_vec();
1275 for index in (0..upper.len()).rev() {
1276 if upper[index] != u8::MAX {
1277 upper[index] = upper[index].saturating_add(1);
1278 upper.truncate(index + 1);
1279 return Some(upper);
1280 }
1281 }
1282 None
1283}
1284
1285fn read_u16(encoded: &[u8], cursor: &mut usize) -> Result<u16, MaterializedIndexError> {
1286 read_array(encoded, cursor).map(u16::from_le_bytes)
1287}
1288
1289fn read_u32(encoded: &[u8], cursor: &mut usize) -> Result<u32, MaterializedIndexError> {
1290 read_array(encoded, cursor).map(u32::from_le_bytes)
1291}
1292
1293fn read_u64(encoded: &[u8], cursor: &mut usize) -> Result<u64, MaterializedIndexError> {
1294 read_array(encoded, cursor).map(u64::from_le_bytes)
1295}
1296
1297fn read_array<const N: usize>(
1298 encoded: &[u8],
1299 cursor: &mut usize,
1300) -> Result<[u8; N], MaterializedIndexError> {
1301 let end = cursor
1302 .checked_add(N)
1303 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1304 let bytes = encoded
1305 .get(*cursor..end)
1306 .ok_or(MaterializedIndexError::MalformedLexicalProjection)?;
1307 *cursor = end;
1308 Ok(copy_array(bytes))
1309}
1310
1311fn require_vector_space(
1312 write: &redb::WriteTransaction,
1313 name: &VectorSpaceName,
1314) -> Result<VectorSpaceDefinition, MaterializedIndexError> {
1315 let encoded = {
1316 let table = write.open_table(VECTOR_SPACES)?;
1317 table
1318 .get(name.as_str())?
1319 .map(|value| value.value().to_vec())
1320 };
1321 let Some(encoded) = encoded else {
1322 return Err(MaterializedIndexError::UnknownVectorSpace {
1323 name: name.as_str().to_owned(),
1324 });
1325 };
1326 decode_vector_space_value(name, &encoded)
1327}
1328
1329fn encode_vector_space_value(definition: &VectorSpaceDefinition) -> [u8; 4] {
1330 let dimension = definition.dimension.to_le_bytes();
1331 [dimension[0], dimension[1], definition.metric as u8, 1]
1332}
1333
1334fn decode_vector_space_value(
1335 name: &VectorSpaceName,
1336 encoded: &[u8],
1337) -> Result<VectorSpaceDefinition, MaterializedIndexError> {
1338 if encoded.len() != 4 || encoded[2] != VectorMetric::Cosine as u8 || encoded[3] != 1 {
1339 return Err(MaterializedIndexError::MalformedVectorIndex);
1340 }
1341 let dimension = u16::from_le_bytes(copy_array(&encoded[..2]));
1342 Ok(VectorSpaceDefinition::cosine(name.clone(), dimension)?)
1343}
1344
1345fn encode_lexical_index_value(
1346 definition: &LexicalIndexDefinition,
1347) -> Result<Vec<u8>, MaterializedIndexError> {
1348 let field_count = u8::try_from(definition.fields.len())
1349 .map_err(|_| MaterializedIndexError::MalformedLexicalIndex)?;
1350 let mut encoded = vec![1, field_count];
1351 for field in &definition.fields {
1352 let segment_count = u8::try_from(field.path.segments().len())
1353 .map_err(|_| MaterializedIndexError::MalformedLexicalIndex)?;
1354 encoded.push(segment_count);
1355 for segment in field.path.segments() {
1356 let length = u16::try_from(segment.len())
1357 .map_err(|_| MaterializedIndexError::MalformedLexicalIndex)?;
1358 encoded.extend_from_slice(&length.to_le_bytes());
1359 encoded.extend_from_slice(segment.as_bytes());
1360 }
1361 encoded.extend_from_slice(&field.weight_micros.to_le_bytes());
1362 }
1363 Ok(encoded)
1364}
1365
1366fn decode_lexical_index_value(
1367 name: &VectorSpaceName,
1368 encoded: &[u8],
1369) -> Result<LexicalIndexDefinition, MaterializedIndexError> {
1370 if encoded.first() != Some(&1) {
1371 return Err(MaterializedIndexError::MalformedLexicalIndex);
1372 }
1373 let field_count = usize::from(
1374 *encoded
1375 .get(1)
1376 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?,
1377 );
1378 let mut cursor = 2_usize;
1379 let mut fields = Vec::with_capacity(field_count);
1380 for _ in 0..field_count {
1381 let segment_count = usize::from(
1382 *encoded
1383 .get(cursor)
1384 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?,
1385 );
1386 cursor = cursor
1387 .checked_add(1)
1388 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?;
1389 let mut segments = Vec::with_capacity(segment_count);
1390 for _ in 0..segment_count {
1391 let length_end = cursor
1392 .checked_add(2)
1393 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?;
1394 let length = usize::from(u16::from_le_bytes(copy_array(
1395 encoded
1396 .get(cursor..length_end)
1397 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?,
1398 )));
1399 cursor = length_end;
1400 let segment_end = cursor
1401 .checked_add(length)
1402 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?;
1403 let segment = std::str::from_utf8(
1404 encoded
1405 .get(cursor..segment_end)
1406 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?,
1407 )
1408 .map_err(|_| MaterializedIndexError::MalformedLexicalIndex)?
1409 .to_owned();
1410 cursor = segment_end;
1411 segments.push(segment);
1412 }
1413 let weight_end = cursor
1414 .checked_add(4)
1415 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?;
1416 let weight_micros = u32::from_le_bytes(copy_array(
1417 encoded
1418 .get(cursor..weight_end)
1419 .ok_or(MaterializedIndexError::MalformedLexicalIndex)?,
1420 ));
1421 cursor = weight_end;
1422 fields.push(LexicalField {
1423 path: FieldPath::new(segments),
1424 weight_micros,
1425 });
1426 }
1427 if cursor != encoded.len() {
1428 return Err(MaterializedIndexError::MalformedLexicalIndex);
1429 }
1430 LexicalIndexDefinition::new(name.clone(), fields).map_err(MaterializedIndexError::from)
1431}
1432
1433fn encode_vector_key(space: &VectorSpaceName, key: &[u8]) -> Vec<u8> {
1434 let mut encoded = Vec::with_capacity(space.as_str().len() + 1 + key.len());
1435 encoded.extend_from_slice(space.as_str().as_bytes());
1436 encoded.push(0);
1437 encoded.extend_from_slice(key);
1438 encoded
1439}
1440
1441fn decode_vector_key(encoded: &[u8]) -> Result<(VectorSpaceName, Vec<u8>), MaterializedIndexError> {
1442 let space_end = encoded
1443 .iter()
1444 .position(|byte| *byte == 0)
1445 .ok_or(MaterializedIndexError::MalformedVectorIndex)?;
1446 let space = encoded
1447 .get(..space_end)
1448 .filter(|space| !space.is_empty())
1449 .ok_or(MaterializedIndexError::MalformedVectorIndex)?;
1450 let key = encoded
1451 .get(space_end + 1..)
1452 .filter(|key| !key.is_empty())
1453 .ok_or(MaterializedIndexError::MalformedVectorIndex)?;
1454 let space =
1455 std::str::from_utf8(space).map_err(|_| MaterializedIndexError::MalformedVectorIndex)?;
1456 Ok((VectorSpaceName::new(space.to_owned())?, key.to_vec()))
1457}
1458
1459fn decode_vector_key_for_space(
1460 encoded: &[u8],
1461 expected: &VectorSpaceName,
1462) -> Result<Option<Vec<u8>>, MaterializedIndexError> {
1463 let (space, key) = decode_vector_key(encoded)?;
1464 Ok((space == *expected).then_some(key))
1465}
1466
1467fn encode_vector_value(vector: &Q15Vector) -> Vec<u8> {
1468 let mut encoded = Vec::with_capacity(2 + 2 * vector.as_slice().len());
1469 encoded.extend_from_slice(&vector.dimension().to_le_bytes());
1470 for value in vector.as_slice() {
1471 encoded.extend_from_slice(&value.to_le_bytes());
1472 }
1473 encoded
1474}
1475
1476fn decode_vector_value(encoded: &[u8]) -> Result<Q15Vector, MaterializedIndexError> {
1477 let dimension_bytes = encoded
1478 .get(..2)
1479 .ok_or(MaterializedIndexError::MalformedVectorIndex)?;
1480 let dimension = usize::from(u16::from_le_bytes(copy_array(dimension_bytes)));
1481 let expected_length = dimension
1482 .checked_mul(2)
1483 .and_then(|length| length.checked_add(2))
1484 .ok_or(MaterializedIndexError::MalformedVectorIndex)?;
1485 if encoded.len() != expected_length {
1486 return Err(MaterializedIndexError::MalformedVectorIndex);
1487 }
1488 let mut values = Vec::with_capacity(dimension);
1489 for chunk in encoded[2..].chunks_exact(2) {
1490 values.push(i16::from_le_bytes(copy_array(chunk)));
1491 }
1492 Ok(Q15Vector::new(values)?)
1493}
1494
1495struct IndexRestoreVisitor<'transaction> {
1496 write: &'transaction mut redb::WriteTransaction,
1497}
1498
1499impl SnapshotRecordVisitor for IndexRestoreVisitor<'_> {
1500 fn put(&mut self, key: &[u8], value: &[u8]) -> Result<(), SnapshotError> {
1501 update_lexical_projections(self.write, key, Some(value)).map_err(SnapshotError::from)?;
1502 let mut table = self
1503 .write
1504 .open_table(KV)
1505 .map_err(MaterializedIndexError::from)?;
1506 table
1507 .insert(key, value)
1508 .map_err(MaterializedIndexError::from)?;
1509 Ok(())
1510 }
1511
1512 fn receipt(&mut self, receipt: &CommitReceipt) -> Result<(), SnapshotError> {
1513 let mut table = self
1514 .write
1515 .open_table(IDEMPOTENCY)
1516 .map_err(MaterializedIndexError::from)?;
1517 table
1518 .insert(
1519 receipt.transaction_id.as_bytes().as_slice(),
1520 encode_receipt(receipt).as_slice(),
1521 )
1522 .map_err(MaterializedIndexError::from)?;
1523 Ok(())
1524 }
1525
1526 fn vector_space(&mut self, definition: &VectorSpaceDefinition) -> Result<(), SnapshotError> {
1527 let mut table = self
1528 .write
1529 .open_table(VECTOR_SPACES)
1530 .map_err(MaterializedIndexError::from)?;
1531 let encoded = encode_vector_space_value(definition);
1532 table
1533 .insert(definition.name.as_str(), encoded.as_slice())
1534 .map_err(MaterializedIndexError::from)?;
1535 Ok(())
1536 }
1537
1538 fn lexical_index(&mut self, definition: &LexicalIndexDefinition) -> Result<(), SnapshotError> {
1539 apply_lexical_index_definition(self.write, definition).map_err(SnapshotError::from)?;
1540 build_lexical_projection(self.write, definition).map_err(SnapshotError::from)?;
1541 Ok(())
1542 }
1543
1544 fn vector(
1545 &mut self,
1546 space: &VectorSpaceName,
1547 key: &[u8],
1548 vector: &Q15Vector,
1549 ) -> Result<(), SnapshotError> {
1550 let mut table = self
1551 .write
1552 .open_table(VECTORS)
1553 .map_err(MaterializedIndexError::from)?;
1554 let encoded_key = encode_vector_key(space, key);
1555 let encoded_vector = encode_vector_value(vector);
1556 table
1557 .insert(encoded_key.as_slice(), encoded_vector.as_slice())
1558 .map_err(MaterializedIndexError::from)?;
1559 Ok(())
1560 }
1561}
1562
1563fn encode_receipt(receipt: &CommitReceipt) -> [u8; RECEIPT_LENGTH] {
1564 let mut encoded = [0_u8; RECEIPT_LENGTH];
1565 encoded[..8].copy_from_slice(&receipt.commit_sequence.to_le_bytes());
1566 encoded[8..40].copy_from_slice(&receipt.commit_digest);
1567 encoded[40..72].copy_from_slice(&receipt.transaction_digest);
1568 encoded
1569}
1570
1571fn decode_receipt(
1572 transaction_id: Uuid,
1573 encoded: &[u8],
1574) -> Result<CommitReceipt, MaterializedIndexError> {
1575 if encoded.len() != RECEIPT_LENGTH {
1576 return Err(MaterializedIndexError::IdempotencyDiverged { transaction_id });
1577 }
1578 Ok(CommitReceipt {
1579 transaction_id,
1580 commit_sequence: u64::from_le_bytes(copy_array(&encoded[..8])),
1581 commit_digest: copy_array(&encoded[8..40]),
1582 transaction_digest: copy_array(&encoded[40..72]),
1583 })
1584}
1585
1586fn decode_sequence(encoded: &[u8]) -> Result<u64, MaterializedIndexError> {
1587 if encoded.len() != 8 {
1588 return Err(MaterializedIndexError::MalformedCheckpoint);
1589 }
1590 Ok(u64::from_le_bytes(copy_array(encoded)))
1591}
1592
1593fn decode_digest(encoded: &[u8]) -> Result<[u8; 32], MaterializedIndexError> {
1594 if encoded.len() != 32 {
1595 return Err(MaterializedIndexError::MalformedCheckpoint);
1596 }
1597 Ok(copy_array(encoded))
1598}
1599
1600fn copy_array<const N: usize>(source: &[u8]) -> [u8; N] {
1601 let mut output = [0_u8; N];
1602 output.copy_from_slice(source);
1603 output
1604}
1605
1606#[cfg(test)]
1607mod tests {
1608 use std::error::Error;
1609
1610 use uuid::Uuid;
1611
1612 use super::{MaterializedIndex, MaterializedIndexError};
1613 use crate::{DurableLog, Mutation, test_support::TestDirectory};
1614
1615 fn recovery_with_operation(
1616 path: &std::path::Path,
1617 operation: Vec<u8>,
1618 ) -> Result<crate::RecoveryReport, Box<dyn Error>> {
1619 let (mut log, _) = DurableLog::open_file(path)?;
1620 log.append_transaction(Uuid::now_v7(), &[operation])?;
1621 drop(log);
1622 let (_, recovery) = DurableLog::open_file(path)?;
1623 Ok(recovery)
1624 }
1625
1626 #[test]
1627 fn checkpoint_rejects_a_different_log_history() -> Result<(), Box<dyn Error>> {
1628 let temporary = TestDirectory::new("index-divergence")?;
1629 let first = recovery_with_operation(
1630 &temporary.path().join("first.hylog"),
1631 Mutation::put(b"key", b"first").encode()?,
1632 )?;
1633 let second = recovery_with_operation(
1634 &temporary.path().join("second.hylog"),
1635 Mutation::put(b"key", b"second").encode()?,
1636 )?;
1637 let index = MaterializedIndex::open(temporary.path().join("index.redb"))?;
1638 assert_eq!(index.replay(&first)?, 1);
1639
1640 let result = index.replay(&second);
1641 assert!(matches!(
1642 result,
1643 Err(MaterializedIndexError::Diverged { sequence: 3 })
1644 ));
1645 assert_eq!(index.get(b"key")?, Some(b"first".to_vec()));
1646 Ok(())
1647 }
1648
1649 #[test]
1650 fn invalid_committed_operation_never_advances_checkpoint() -> Result<(), Box<dyn Error>> {
1651 let temporary = TestDirectory::new("index-invalid-operation")?;
1652 let recovery = recovery_with_operation(
1653 &temporary.path().join("segment.hylog"),
1654 b"not-a-mutation".to_vec(),
1655 )?;
1656 let index = MaterializedIndex::open(temporary.path().join("index.redb"))?;
1657
1658 assert!(matches!(
1659 index.replay(&recovery),
1660 Err(MaterializedIndexError::Mutation(_))
1661 ));
1662 assert!(matches!(
1663 index.replay(&recovery),
1664 Err(MaterializedIndexError::Mutation(_))
1665 ));
1666 Ok(())
1667 }
1668
1669 #[test]
1670 fn replay_backfills_a_missing_idempotency_receipt() -> Result<(), Box<dyn Error>> {
1671 let temporary = TestDirectory::new("index-idempotency-backfill")?;
1672 let recovery = recovery_with_operation(
1673 &temporary.path().join("segment.hylog"),
1674 Mutation::put(b"key", b"value").encode()?,
1675 )?;
1676 let receipt = recovery.transactions[0].receipt;
1677 let index = MaterializedIndex::open(temporary.path().join("index.redb"))?;
1678 assert_eq!(index.replay(&recovery)?, 1);
1679 assert_eq!(index.receipt(receipt.transaction_id)?, Some(receipt));
1680
1681 let mut write = index.database.begin_write()?;
1682 write.set_durability(redb::Durability::Immediate)?;
1683 {
1684 let mut table = write.open_table(super::IDEMPOTENCY)?;
1685 table.remove(receipt.transaction_id.as_bytes().as_slice())?;
1686 }
1687 write.commit()?;
1688 assert_eq!(index.receipt(receipt.transaction_id)?, None);
1689
1690 assert_eq!(index.replay(&recovery)?, 0);
1691 assert_eq!(index.receipt(receipt.transaction_id)?, Some(receipt));
1692 Ok(())
1693 }
1694}