2024-08-29 17:51:42 +02:00
|
|
|
use std::fs::File;
|
|
|
|
|
2024-08-29 18:27:02 +02:00
|
|
|
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;
|
2024-09-04 09:59:19 +02:00
|
|
|
use crate::index::main_key::{DOCUMENTS_IDS_KEY, WORDS_FST_KEY};
|
2024-09-02 10:42:19 +02:00
|
|
|
use crate::update::new::KvReaderFieldId;
|
2024-09-03 11:02:39 +02:00
|
|
|
use crate::update::MergeDeladdCboRoaringBitmaps;
|
2024-08-29 15:07:59 +02:00
|
|
|
use crate::{DocumentId, Index};
|
|
|
|
|
|
|
|
/// 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);
|
2024-08-29 18:27:02 +02:00
|
|
|
(MergerSender(sender), WriterReceiver(receiver))
|
2024-08-29 15:07:59 +02:00
|
|
|
}
|
|
|
|
|
2024-08-29 17:51:42 +02:00
|
|
|
/// The capacity of the channel is currently in number of messages.
|
|
|
|
pub fn extractors_merger_channels(cap: usize) -> ExtractorsMergerChannels {
|
|
|
|
let (sender, receiver) = crossbeam_channel::bounded(cap);
|
|
|
|
|
|
|
|
ExtractorsMergerChannels {
|
|
|
|
merger_receiver: MergerReceiver(receiver),
|
|
|
|
deladd_cbo_roaring_bitmap_sender: DeladdCboRoaringBitmapSender(sender.clone()),
|
2024-09-04 09:59:19 +02:00
|
|
|
extracted_documents_sender: ExtractedDocumentsSender(sender.clone()),
|
2024-08-29 17:51:42 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
pub struct ExtractorsMergerChannels {
|
|
|
|
pub merger_receiver: MergerReceiver,
|
|
|
|
pub deladd_cbo_roaring_bitmap_sender: DeladdCboRoaringBitmapSender,
|
2024-09-04 09:59:19 +02:00
|
|
|
pub extracted_documents_sender: ExtractedDocumentsSender,
|
2024-08-29 17:51:42 +02:00
|
|
|
}
|
|
|
|
|
2024-08-29 15:07:59 +02:00
|
|
|
pub struct KeyValueEntry {
|
2024-08-29 17:51:42 +02:00
|
|
|
key_length: usize,
|
|
|
|
data: Box<[u8]>,
|
2024-08-29 15:07:59 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
impl KeyValueEntry {
|
2024-08-29 17:51:42 +02:00
|
|
|
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() }
|
|
|
|
}
|
|
|
|
|
2024-08-29 18:27:02 +02:00
|
|
|
pub fn key(&self) -> &[u8] {
|
2024-08-29 19:20:10 +02:00
|
|
|
&self.data.as_ref()[..self.key_length]
|
2024-08-29 18:27:02 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
pub fn value(&self) -> &[u8] {
|
2024-08-29 19:20:10 +02:00
|
|
|
&self.data.as_ref()[self.key_length..]
|
2024-08-29 15:07:59 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2024-08-29 17:51:42 +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()
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2024-08-29 18:27:02 +02:00
|
|
|
pub enum EntryOperation {
|
2024-08-29 17:51:42 +02:00
|
|
|
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 }
|
|
|
|
}
|
|
|
|
|
2024-08-29 18:27:02 +02:00
|
|
|
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
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2024-09-04 09:59:19 +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,
|
|
|
|
Documents,
|
|
|
|
Main,
|
2024-08-29 15:07:59 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
impl WriterOperation {
|
|
|
|
pub fn database(&self, index: &Index) -> heed::Database<Bytes, Bytes> {
|
2024-09-04 09:59:19 +02:00
|
|
|
match self.database {
|
|
|
|
Database::Main => index.main.remap_types(),
|
|
|
|
Database::Documents => index.documents.remap_types(),
|
|
|
|
Database::WordDocids => index.word_docids.remap_types(),
|
2024-08-29 15:07:59 +02:00
|
|
|
}
|
|
|
|
}
|
2024-09-04 09:59:19 +02:00
|
|
|
|
|
|
|
pub fn entry(self) -> EntryOperation {
|
|
|
|
self.entry
|
|
|
|
}
|
2024-08-29 15:07:59 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
pub struct WriterReceiver(Receiver<WriterOperation>);
|
|
|
|
|
2024-08-29 18:27:02 +02:00
|
|
|
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>);
|
|
|
|
|
2024-08-29 17:51:42 +02:00
|
|
|
impl MergerSender {
|
2024-09-04 09:59:19 +02:00
|
|
|
pub fn main(&self) -> MainSender<'_> {
|
|
|
|
MainSender(&self.0)
|
|
|
|
}
|
|
|
|
|
2024-08-29 17:51:42 +02:00
|
|
|
pub fn word_docids(&self) -> WordDocidsSender<'_> {
|
|
|
|
WordDocidsSender(&self.0)
|
|
|
|
}
|
2024-09-04 09:59:19 +02:00
|
|
|
|
|
|
|
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(())),
|
|
|
|
}
|
|
|
|
}
|
2024-08-29 17:51:42 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
pub struct WordDocidsSender<'a>(&'a Sender<WriterOperation>);
|
|
|
|
|
|
|
|
impl WordDocidsSender<'_> {
|
|
|
|
pub fn write(&self, key: &[u8], value: &[u8]) -> StdResult<(), SendError<()>> {
|
2024-09-04 09:59:19 +02:00
|
|
|
let entry = EntryOperation::Write(KeyValueEntry::from_key_value(key, value));
|
|
|
|
match self.0.send(WriterOperation { database: Database::WordDocids, entry }) {
|
2024-08-29 17:51:42 +02:00
|
|
|
Ok(()) => Ok(()),
|
|
|
|
Err(SendError(_)) => Err(SendError(())),
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
pub fn delete(&self, key: &[u8]) -> StdResult<(), SendError<()>> {
|
2024-09-04 09:59:19 +02:00
|
|
|
let entry = EntryOperation::Delete(KeyEntry::from_key(key));
|
|
|
|
match self.0.send(WriterOperation { database: Database::WordDocids, entry }) {
|
2024-08-29 17:51:42 +02:00
|
|
|
Ok(()) => Ok(()),
|
|
|
|
Err(SendError(_)) => Err(SendError(())),
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2024-09-04 09:59:19 +02:00
|
|
|
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
|
|
|
|
2024-09-04 09:59:19 +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(()),
|
2024-08-29 17:51:42 +02:00
|
|
|
Err(SendError(_)) => Err(SendError(())),
|
2024-08-29 15:07:59 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2024-08-29 17:51:42 +02:00
|
|
|
|
|
|
|
pub enum MergerOperation {
|
2024-09-03 11:02:39 +02:00
|
|
|
WordDocidsMerger(Merger<File, MergeDeladdCboRoaringBitmaps>),
|
2024-09-04 09:59:19 +02:00
|
|
|
InsertDocument { docid: DocumentId, document: Box<KvReaderFieldId> },
|
|
|
|
DeleteDocument { docid: DocumentId },
|
2024-08-29 17:51:42 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
pub struct MergerReceiver(Receiver<MergerOperation>);
|
|
|
|
|
|
|
|
impl IntoIterator for MergerReceiver {
|
|
|
|
type Item = MergerOperation;
|
2024-08-29 18:27:02 +02:00
|
|
|
type IntoIter = IntoIter<Self::Item>;
|
2024-08-29 17:51:42 +02:00
|
|
|
|
|
|
|
fn into_iter(self) -> Self::IntoIter {
|
|
|
|
self.0.into_iter()
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
#[derive(Clone)]
|
|
|
|
pub struct DeladdCboRoaringBitmapSender(Sender<MergerOperation>);
|
2024-09-03 11:02:39 +02:00
|
|
|
|
|
|
|
impl DeladdCboRoaringBitmapSender {
|
|
|
|
pub fn word_docids(
|
|
|
|
&self,
|
|
|
|
merger: Merger<File, MergeDeladdCboRoaringBitmaps>,
|
|
|
|
) -> StdResult<(), SendError<()>> {
|
|
|
|
let operation = MergerOperation::WordDocidsMerger(merger);
|
|
|
|
match self.0.send(operation) {
|
|
|
|
Ok(()) => Ok(()),
|
|
|
|
Err(SendError(_)) => Err(SendError(())),
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2024-09-04 09:59:19 +02:00
|
|
|
|
|
|
|
#[derive(Clone)]
|
|
|
|
pub struct ExtractedDocumentsSender(Sender<MergerOperation>);
|
|
|
|
|
|
|
|
impl ExtractedDocumentsSender {
|
|
|
|
pub fn insert(
|
|
|
|
&self,
|
|
|
|
docid: DocumentId,
|
|
|
|
document: Box<KvReaderFieldId>,
|
|
|
|
) -> StdResult<(), SendError<()>> {
|
|
|
|
match self.0.send(MergerOperation::InsertDocument { docid, document }) {
|
|
|
|
Ok(()) => Ok(()),
|
|
|
|
Err(SendError(_)) => Err(SendError(())),
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
pub fn delete(&self, docid: DocumentId) -> StdResult<(), SendError<()>> {
|
|
|
|
match self.0.send(MergerOperation::DeleteDocument { docid }) {
|
|
|
|
Ok(()) => Ok(()),
|
|
|
|
Err(SendError(_)) => Err(SendError(())),
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|