use hashbrown::HashSet;
use rocksdb::WriteBatch;
use std::sync::Arc;
use crate::lexer::TokenLexer;
use crate::store::StoreItem;
use crate::store::fst::StoreFSTActionBuilder;
use crate::store::kv::{StoreKVAcquireMode, StoreKVActionBuilder};
use crate::util::hash::NoopU32HasherBuilder;
impl super::Executor {
pub fn push(&self, item: StoreItem, lexer: TokenLexer, assume_new: bool) -> Result<(), ()> {
let StoreItem(collection, Some(bucket), Some(object)) = item else {
return Err(());
};
let _kv_read_guard = self.kv_pool.lock_read_access();
let _fst_read_guard = self.fst_pool.lock_read_access();
let (Ok(kv_store), Ok(fst_store)) = (
self.kv_pool
.acquire(StoreKVAcquireMode::Any, collection, None, |_| {}),
self.fst_pool.acquire(collection, bucket),
) else {
return Err(());
};
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, fst_action) = (
StoreKVActionBuilder::access_read_write(bucket, Arc::clone(&kv_store)),
StoreFSTActionBuilder::access(fst_store),
);
let mut batch = WriteBatch::default();
let oid = object.as_str();
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);
kv_action.set_oid_to_iid(&mut batch, oid, iid);
kv_action.set_iid_to_oid(&mut batch, iid, oid);
iid
};
let iid = if assume_new {
assign_new_iid()
} else {
(kv_action.get_oid_to_iid(oid))
.unwrap_or_else(|()| {
tracing::error!("Error getting OID-To-IID");
None
})
.unwrap_or_else(assign_new_iid)
};
let mut tokens = HashSet::with_capacity_and_hasher(128, NoopU32HasherBuilder);
for (token, term_hashed, _) in lexer {
let term = token.as_str();
tokens.insert(term_hashed);
if fst_action.push_word(&term, &self.app_conf.store.fst) {
tracing::trace!("push term committed to graph: {}", term);
}
}
for &term_hashed in tokens.iter() {
kv_action.add_term_to_iids(&mut batch, term_hashed, std::iter::once(iid));
}
kv_action.add_iid_to_terms(&mut batch, iid, tokens.into_iter());
executor_ensure_op!(kv_action.write(batch));
Ok(())
}
}