334 lines
9.7 KiB
Rust
Raw Normal View History

use std::fs::File;
use std::marker::PhantomData;
use crossbeam_channel::{IntoIter, Receiver, SendError, Sender};
2024-09-03 11:02:39 +02:00
use grenad::Merger;
2024-08-29 15:07:59 +02:00
use heed::types::Bytes;
use super::StdResult;
use crate::index::main_key::{DOCUMENTS_IDS_KEY, WORDS_FST_KEY};
use crate::update::new::KvReaderFieldId;
2024-09-03 11:02:39 +02:00
use crate::update::MergeDeladdCboRoaringBitmaps;
2024-09-04 12:17:13 +02:00
use crate::{CboRoaringBitmapCodec, DocumentId, Index};
2024-08-29 15:07:59 +02:00
/// The capacity of the channel is currently in number of messages.
2024-09-02 15:10:21 +02:00
pub fn merger_writer_channel(cap: usize) -> (MergerSender, WriterReceiver) {
2024-08-29 15:07:59 +02:00
let (sender, receiver) = crossbeam_channel::bounded(cap);
(MergerSender(sender), WriterReceiver(receiver))
2024-08-29 15:07:59 +02:00
}
/// The capacity of the channel is currently in number of messages.
pub fn extractors_merger_channels(cap: usize) -> (ExtractorSender, MergerReceiver) {
let (sender, receiver) = crossbeam_channel::bounded(cap);
(ExtractorSender(sender), MergerReceiver(receiver))
}
2024-08-29 15:07:59 +02:00
pub struct KeyValueEntry {
key_length: usize,
data: Box<[u8]>,
2024-08-29 15:07:59 +02:00
}
impl KeyValueEntry {
pub fn from_key_value(key: &[u8], value: &[u8]) -> Self {
let mut data = Vec::with_capacity(key.len() + value.len());
data.extend_from_slice(key);
data.extend_from_slice(value);
KeyValueEntry { key_length: key.len(), data: data.into_boxed_slice() }
}
pub fn key(&self) -> &[u8] {
&self.data.as_ref()[..self.key_length]
}
pub fn value(&self) -> &[u8] {
&self.data.as_ref()[self.key_length..]
2024-08-29 15:07:59 +02:00
}
}
pub struct KeyEntry {
data: Box<[u8]>,
}
impl KeyEntry {
pub fn from_key(key: &[u8]) -> Self {
KeyEntry { data: key.to_vec().into_boxed_slice() }
}
pub fn entry(&self) -> &[u8] {
self.data.as_ref()
}
}
pub enum EntryOperation {
Delete(KeyEntry),
Write(KeyValueEntry),
}
2024-08-29 15:07:59 +02:00
pub struct DocumentEntry {
docid: DocumentId,
content: Box<[u8]>,
}
impl DocumentEntry {
pub fn new_uncompressed(docid: DocumentId, content: Box<KvReaderFieldId>) -> Self {
DocumentEntry { docid, content: content.into() }
}
pub fn new_compressed(docid: DocumentId, content: Box<[u8]>) -> Self {
DocumentEntry { docid, content }
}
pub fn key(&self) -> [u8; 4] {
self.docid.to_be_bytes()
}
pub fn content(&self) -> &[u8] {
&self.content
2024-08-29 15:07:59 +02:00
}
}
pub struct DocumentDeletionEntry(DocumentId);
impl DocumentDeletionEntry {
pub fn key(&self) -> [u8; 4] {
self.0.to_be_bytes()
}
}
pub struct WriterOperation {
database: Database,
entry: EntryOperation,
}
pub enum Database {
WordDocids,
2024-09-04 12:17:13 +02:00
ExactWordDocids,
WordFidDocids,
2024-09-04 12:17:13 +02:00
WordPositionDocids,
Documents,
Main,
2024-08-29 15:07:59 +02:00
}
impl WriterOperation {
pub fn database(&self, index: &Index) -> heed::Database<Bytes, Bytes> {
match self.database {
Database::Main => index.main.remap_types(),
Database::Documents => index.documents.remap_types(),
Database::WordDocids => index.word_docids.remap_types(),
2024-09-04 12:17:13 +02:00
Database::ExactWordDocids => index.exact_word_docids.remap_types(),
Database::WordFidDocids => index.word_fid_docids.remap_types(),
2024-09-04 12:17:13 +02:00
Database::WordPositionDocids => index.word_position_docids.remap_types(),
2024-08-29 15:07:59 +02:00
}
}
pub fn entry(self) -> EntryOperation {
self.entry
}
2024-08-29 15:07:59 +02:00
}
pub struct WriterReceiver(Receiver<WriterOperation>);
impl IntoIterator for WriterReceiver {
type Item = WriterOperation;
type IntoIter = IntoIter<Self::Item>;
fn into_iter(self) -> Self::IntoIter {
self.0.into_iter()
2024-08-29 15:07:59 +02:00
}
}
pub struct MergerSender(Sender<WriterOperation>);
impl MergerSender {
pub fn main(&self) -> MainSender<'_> {
MainSender(&self.0)
}
2024-09-04 12:17:13 +02:00
pub fn docids<D: DatabaseType>(&self) -> DocidsSender<'_, D> {
DocidsSender { sender: &self.0, _marker: PhantomData }
}
pub fn documents(&self) -> DocumentsSender<'_> {
DocumentsSender(&self.0)
}
pub fn send_documents_ids(&self, bitmap: &[u8]) -> StdResult<(), SendError<()>> {
let entry = EntryOperation::Write(KeyValueEntry::from_key_value(
DOCUMENTS_IDS_KEY.as_bytes(),
bitmap,
));
match self.0.send(WriterOperation { database: Database::Main, entry }) {
Ok(()) => Ok(()),
Err(SendError(_)) => Err(SendError(())),
}
}
}
pub struct MainSender<'a>(&'a Sender<WriterOperation>);
impl MainSender<'_> {
pub fn write_words_fst(&self, value: &[u8]) -> StdResult<(), SendError<()>> {
let entry =
EntryOperation::Write(KeyValueEntry::from_key_value(WORDS_FST_KEY.as_bytes(), value));
match self.0.send(WriterOperation { database: Database::Main, entry }) {
Ok(()) => Ok(()),
Err(SendError(_)) => Err(SendError(())),
}
}
pub fn delete(&self, key: &[u8]) -> StdResult<(), SendError<()>> {
let entry = EntryOperation::Delete(KeyEntry::from_key(key));
match self.0.send(WriterOperation { database: Database::Main, entry }) {
Ok(()) => Ok(()),
Err(SendError(_)) => Err(SendError(())),
}
}
}
pub enum WordDocids {}
2024-09-04 12:17:13 +02:00
pub enum ExactWordDocids {}
pub enum WordFidDocids {}
2024-09-04 12:17:13 +02:00
pub enum WordPositionDocids {}
pub trait DatabaseType {
2024-09-04 12:17:13 +02:00
const DATABASE: Database;
fn new_merger_operation(merger: Merger<File, MergeDeladdCboRoaringBitmaps>) -> MergerOperation;
}
impl DatabaseType for WordDocids {
2024-09-04 12:17:13 +02:00
const DATABASE: Database = Database::WordDocids;
fn new_merger_operation(merger: Merger<File, MergeDeladdCboRoaringBitmaps>) -> MergerOperation {
MergerOperation::WordDocidsMerger(merger)
}
}
impl DatabaseType for ExactWordDocids {
const DATABASE: Database = Database::ExactWordDocids;
fn new_merger_operation(merger: Merger<File, MergeDeladdCboRoaringBitmaps>) -> MergerOperation {
MergerOperation::ExactWordDocidsMerger(merger)
}
}
impl DatabaseType for WordFidDocids {
2024-09-04 12:17:13 +02:00
const DATABASE: Database = Database::WordFidDocids;
fn new_merger_operation(merger: Merger<File, MergeDeladdCboRoaringBitmaps>) -> MergerOperation {
MergerOperation::WordFidDocidsMerger(merger)
}
}
impl DatabaseType for WordPositionDocids {
const DATABASE: Database = Database::WordPositionDocids;
fn new_merger_operation(merger: Merger<File, MergeDeladdCboRoaringBitmaps>) -> MergerOperation {
MergerOperation::WordPositionDocidsMerger(merger)
}
}
pub struct DocidsSender<'a, D> {
sender: &'a Sender<WriterOperation>,
_marker: PhantomData<D>,
}
impl<D: DatabaseType> DocidsSender<'_, D> {
pub fn write(&self, key: &[u8], value: &[u8]) -> StdResult<(), SendError<()>> {
let entry = EntryOperation::Write(KeyValueEntry::from_key_value(key, value));
2024-09-04 12:17:13 +02:00
match self.sender.send(WriterOperation { database: D::DATABASE, entry }) {
Ok(()) => Ok(()),
Err(SendError(_)) => Err(SendError(())),
}
}
pub fn delete(&self, key: &[u8]) -> StdResult<(), SendError<()>> {
let entry = EntryOperation::Delete(KeyEntry::from_key(key));
2024-09-04 12:17:13 +02:00
match self.sender.send(WriterOperation { database: D::DATABASE, entry }) {
Ok(()) => Ok(()),
Err(SendError(_)) => Err(SendError(())),
}
}
}
pub struct DocumentsSender<'a>(&'a Sender<WriterOperation>);
impl DocumentsSender<'_> {
/// TODO do that efficiently
pub fn uncompressed(
&self,
docid: DocumentId,
document: &KvReaderFieldId,
) -> StdResult<(), SendError<()>> {
let entry = EntryOperation::Write(KeyValueEntry::from_key_value(
&docid.to_be_bytes(),
document.as_bytes(),
));
match self.0.send(WriterOperation { database: Database::Documents, entry }) {
Ok(()) => Ok(()),
Err(SendError(_)) => Err(SendError(())),
}
}
2024-08-29 15:07:59 +02:00
pub fn delete(&self, docid: DocumentId) -> StdResult<(), SendError<()>> {
let entry = EntryOperation::Delete(KeyEntry::from_key(&docid.to_be_bytes()));
match self.0.send(WriterOperation { database: Database::Documents, entry }) {
2024-08-29 15:07:59 +02:00
Ok(()) => Ok(()),
Err(SendError(_)) => Err(SendError(())),
2024-08-29 15:07:59 +02:00
}
}
}
pub enum MergerOperation {
2024-09-03 11:02:39 +02:00
WordDocidsMerger(Merger<File, MergeDeladdCboRoaringBitmaps>),
2024-09-04 12:17:13 +02:00
ExactWordDocidsMerger(Merger<File, MergeDeladdCboRoaringBitmaps>),
WordFidDocidsMerger(Merger<File, MergeDeladdCboRoaringBitmaps>),
2024-09-04 12:17:13 +02:00
WordPositionDocidsMerger(Merger<File, MergeDeladdCboRoaringBitmaps>),
InsertDocument { docid: DocumentId, document: Box<KvReaderFieldId> },
DeleteDocument { docid: DocumentId },
}
pub struct MergerReceiver(Receiver<MergerOperation>);
impl IntoIterator for MergerReceiver {
type Item = MergerOperation;
type IntoIter = IntoIter<Self::Item>;
fn into_iter(self) -> Self::IntoIter {
self.0.into_iter()
}
}
pub struct ExtractorSender(Sender<MergerOperation>);
2024-09-03 11:02:39 +02:00
impl ExtractorSender {
pub fn document_insert(
2024-09-03 11:02:39 +02:00
&self,
docid: DocumentId,
document: Box<KvReaderFieldId>,
2024-09-03 11:02:39 +02:00
) -> StdResult<(), SendError<()>> {
match self.0.send(MergerOperation::InsertDocument { docid, document }) {
2024-09-03 11:02:39 +02:00
Ok(()) => Ok(()),
Err(SendError(_)) => Err(SendError(())),
}
}
pub fn document_delete(&self, docid: DocumentId) -> StdResult<(), SendError<()>> {
match self.0.send(MergerOperation::DeleteDocument { docid }) {
Ok(()) => Ok(()),
Err(SendError(_)) => Err(SendError(())),
}
}
2024-09-04 12:17:13 +02:00
pub fn send_searchable<D: DatabaseType>(
&self,
merger: Merger<File, MergeDeladdCboRoaringBitmaps>,
) -> StdResult<(), SendError<()>> {
2024-09-04 12:17:13 +02:00
match self.0.send(D::new_merger_operation(merger)) {
Ok(()) => Ok(()),
Err(SendError(_)) => Err(SendError(())),
}
}
}