use rocksdb::WriteBatch;
use crate::lexer::itertools::UniqueBy;
use crate::lexer::preprocessor::{PreprocessorOutput, Token};
use crate::store::{StoreItemPart, StoreObjectOid};
use crate::util::hash::NoopU32HasherBuilder;
impl super::Executor {
pub fn push(
&self,
collection: StoreItemPart,
bucket: StoreItemPart,
oid: StoreObjectOid,
input: PreprocessorOutput,
assume_new: bool,
) -> Result<(), ()> {
let _kv_read_guard = self.kv_pool.lock_read_access();
let _fst_read_guard = self.fst_pool.lock_read_access();
let kv_store = self.kv_pool.acquire(true, collection, None, |_| {})?;
let fst_store = self.fst_pool.acquire(collection, bucket)?;
debug_assert!(kv_store.is_some());
let Some(kv_store) = kv_store else {
tracing::error!(
"collection store {collection:?} does not exist, but it should have been created"
);
return Err(());
};
let kv_action = kv_store.access_read_write(bucket);
let mut batch = WriteBatch::default();
let mut assign_new_iid = || {
tracing::trace!("must initialize push executor oid-to-iid and iid-to-oid");
let iid = (kv_action.get_new_iid(&mut batch))
.map_err(|error| tracing::error!("Error getting new IID: {error:?}"))?;
kv_action.set_oid_to_iid(&mut batch, oid, iid);
kv_action.set_iid_to_oid(&mut batch, iid, oid);
Ok(iid)
};
let iid = if assume_new {
assign_new_iid()?
} else {
match kv_action.get_oid_to_iid(oid) {
Ok(Some(iid)) => iid,
Ok(None) => assign_new_iid()?,
Err(error) => {
tracing::error!("Error getting OID-To-IID: {error:?}");
assign_new_iid()?
}
}
};
let mut tokens =
UniqueBy::new_with_hasher(input.tokens(), Token::hash, NoopU32HasherBuilder);
for token in &mut tokens {
let term = token.as_normalized();
let term_hash = token.hash();
if fst_store.push_word(&term, &self.app_conf.store.fst) {
tracing::trace!("push term committed to graph: {}", term);
}
kv_action.add_term_to_iids(&mut batch, term_hash, std::iter::once(iid));
}
kv_action.add_iid_to_terms(&mut batch, iid, tokens.seen().iter().copied());
executor_ensure_op!(kv_action.write(batch));
Ok(())
}
}