use super::prelude::*;
use crate::rdfify::event_quads;
use fs4::fs_std::FileExt;
use nostr::Kind;
use nostr_sdk::client::Error as NostrClientError;
use serde_json::Error as SerdeJsonError;
use std::fmt::Debug;
use std::fs;
use std::sync::Mutex;
use std::sync::RwLock;
#[derive(Debug)]
pub struct RRCache<T: Clone + Send + Sync + 'static> {
pub cache: Cache<String, Arc<TRdfResultSet<T>>>,
}
impl<T: Clone + Send + Sync + 'static> RRCache<T> {
pub fn new(ttl_secs: Option<u64>, idle_secs: Option<u64>) -> Self {
let cache = Cache::builder()
.max_capacity(250)
.time_to_live(Duration::from_secs(ttl_secs.unwrap_or(10)))
.time_to_idle(Duration::from_secs(idle_secs.unwrap_or(15)))
.build();
Self { cache }
}
}
impl<T: Clone + Send + Sync + 'static> std::default::Default for RRCache<T> {
fn default() -> Self {
let cache = Cache::builder()
.max_capacity(250)
.time_to_live(Duration::from_secs(10))
.time_to_idle(Duration::from_secs(15))
.build();
Self { cache }
}
}
impl Debug for RdfEventsStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{:?}", self.dump_path)
}
}
pub enum DatabaseEventsSaveMode {
Direct,
Queue,
}
pub struct RdfEventsStore {
pub database_save_mode: DatabaseEventsSaveMode,
curie_mappings: PrefixMapping,
pub cache: Cache<String, Arc<RdfResultSet>>,
pub store: Store,
pub dump_path: Option<PathBuf>,
pub events_queue: Arc<Mutex<VecDeque<Event>>>,
pub event_tx: Sender<Event, i32>,
pub event_rx: Receiver<Event, i32>,
ev_callbacks: RwLock<
HashMap<u16, Vec<Box<(dyn Fn(&Event, &Vec<Triple>) + Send + Sync)>>>,
>,
ev_quads_callbacks: RwLock<
HashMap<u16, Vec<Box<(dyn Fn(&Event, &Vec<Quad>) + Send + Sync)>>>,
>,
}
#[derive(Debug)]
pub enum RdfStoreError {
IriParseError,
FileLoadError,
DumpFileLockError,
DumpStoreError,
AvailableSpaceError,
EventStoreError(StorageError),
QueryError,
QuadError,
QuadInsertError,
URNError,
NamedNodeError,
CurieError,
SubstitutionError,
RdfCellError,
LDErr(LDError),
ResultsSendError,
QueryChannelLockedError,
QueryChannelEmptyError,
QueryChannelReceiveError,
QueryInProgress,
BadEventError,
EventBuilderError(nostr::event::builder::Error),
SerdeError(SerdeJsonError),
QueryExplanationError,
NostrClientError(NostrClientError),
CallbacksLockError,
DatabaseError(nostr_database::DatabaseError),
}
impl From<LDError> for RdfStoreError {
fn from(e: LDError) -> Self {
RdfStoreError::LDErr(e)
}
}
impl From<IriParseError> for RdfStoreError {
fn from(_e: IriParseError) -> Self {
RdfStoreError::IriParseError
}
}
impl From<StorageError> for RdfStoreError {
fn from(e: StorageError) -> Self {
RdfStoreError::EventStoreError(e)
}
}
impl From<nostr::event::builder::Error> for RdfStoreError {
fn from(e: nostr::event::builder::Error) -> Self {
RdfStoreError::EventBuilderError(e)
}
}
impl From<nostr_sdk::client::Error> for RdfStoreError {
fn from(e: nostr_sdk::client::Error) -> Self {
RdfStoreError::NostrClientError(e)
}
}
impl From<SerdeJsonError> for RdfStoreError {
fn from(e: SerdeJsonError) -> Self {
RdfStoreError::SerdeError(e)
}
}
#[cfg(feature = "nostrdb")]
impl From<nostr_database::DatabaseError> for RdfStoreError {
fn from(e: nostr_database::DatabaseError) -> Self {
RdfStoreError::DatabaseError(e)
}
}
impl Display for RdfStoreError {
fn fmt(&self, f: &mut Formatter) -> std::fmt::Result {
match self {
RdfStoreError::CurieError => write!(f, "Prefix mapping error"),
_ => write!(f, "default"),
}
}
}
impl std::error::Error for RdfStoreError {}
impl RdfEventsStore {
pub fn new(store: Store) -> Result<RdfEventsStore, RdfStoreError> {
let (e_tx, e_rx) = unbounded();
let cache = Cache::builder()
.max_capacity(10_000)
.time_to_live(Duration::from_secs(10))
.time_to_idle(Duration::from_secs(30))
.build();
let mut curie_mappings = PrefixMapping::default();
curie_mappings
.add_prefix("w3nostr", "https://w3id.org/nostr#")
.map_err(|_| RdfStoreError::CurieError)?;
curie_mappings
.add_prefix("schema", "http://schema.org/")
.map_err(|_| RdfStoreError::CurieError)?;
curie_mappings
.add_prefix("nostralink", "http://nostralink.org/")
.map_err(|_| RdfStoreError::CurieError)?;
Ok(RdfEventsStore {
database_save_mode: DatabaseEventsSaveMode::Direct,
cache,
ev_callbacks: RwLock::new(HashMap::new()),
ev_quads_callbacks: RwLock::new(HashMap::new()),
store,
dump_path: None,
curie_mappings,
events_queue: Arc::new(Mutex::new(VecDeque::new())),
event_tx: e_tx,
event_rx: e_rx,
})
}
#[deprecated]
pub fn register_event_triples_callback(
&self,
kind: Kind,
callback: impl Fn(&Event, &Vec<Triple>) + Send + Sync + 'static,
) -> Result<(), RdfStoreError> {
match self.ev_callbacks.write() {
Ok(mut callbacks) => {
let kind_callbacks =
callbacks.entry(kind.as_u16()).or_default();
kind_callbacks.push(Box::new(callback));
Ok(())
}
Err(_) => Err(RdfStoreError::CallbacksLockError),
}
}
pub fn register_event_callback(
&self,
kind: Kind,
callback: impl Fn(&Event, &Vec<Quad>) + Send + Sync + 'static,
) -> Result<(), RdfStoreError> {
match self.ev_quads_callbacks.write() {
Ok(mut callbacks) => {
let kind_callbacks =
callbacks.entry(kind.as_u16()).or_default();
kind_callbacks.push(Box::new(callback));
Ok(())
}
Err(_) => Err(RdfStoreError::CallbacksLockError),
}
}
pub fn open(path: PathBuf) -> Result<Self, RdfStoreError> {
Ok(Self::new(Store::open(path)?)?)
}
pub fn new_inmem(
dump_path: Option<PathBuf>,
) -> Result<RdfEventsStore, RdfStoreError> {
let store = Self::new(Store::new()?)?.with_dump_path(dump_path);
let _ = store.load_from_dump();
Ok(store)
}
pub fn prefix_mappings(&self) -> String {
let mut pmappings = String::new();
for (pfx, uri) in self.curie_mappings.mappings() {
pmappings.push_str(&format!("PREFIX {pfx}: <{uri}>\n"));
}
pmappings
}
pub fn prepare_query(&self, q: &String) -> String {
let mut query = String::new();
query.push_str(&self.prefix_mappings());
query.push_str(&q);
query
}
pub fn with_dump_path(mut self, path: Option<PathBuf>) -> Self {
self.dump_path = path;
self
}
pub fn with_database_save_mode(
mut self,
mode: DatabaseEventsSaveMode,
) -> Self {
self.database_save_mode = mode;
self
}
pub fn load_from_dump(&self) -> Result<(), RdfStoreError> {
match &self.dump_path {
Some(path) => self.bulk_load_file(&path, None),
None => Ok(()),
}
}
pub fn dump(&self, _force: bool) -> Result<(), RdfStoreError> {
match &self.dump_path {
Some(path) => self.dump_graph_to_file(path.clone(), None),
None => Ok(()),
}
}
pub fn bulk_load_file(
&self,
file_path: &PathBuf,
rdf_format: Option<RdfFormat>,
) -> Result<(), RdfStoreError> {
let format = rdf_format.unwrap_or(RdfFormat::NQuads);
match File::open(file_path) {
Ok(file) => self
.store
.bulk_loader()
.load_from_reader(format, file)
.map_err(|_| RdfStoreError::FileLoadError),
Err(_e) => Err(RdfStoreError::FileLoadError),
}
}
pub fn dump_graph_to_file(
&self,
path: PathBuf,
format: Option<RdfFormat>,
) -> Result<(), RdfStoreError> {
let mut tmp_path = path.clone();
tmp_path.set_extension("tmp");
let file = File::create(&tmp_path)
.map_err(|_| RdfStoreError::DumpStoreError)?;
file.lock_exclusive()
.map_err(|_| RdfStoreError::DumpFileLockError)?;
self.store
.dump_graph_to_writer(
GraphNameRef::DefaultGraph,
format.unwrap_or(RdfFormat::NQuads),
file,
)
.map_err(|_| RdfStoreError::DumpStoreError)?;
fs::rename(tmp_path, path)
.map_err(|_| RdfStoreError::DumpFileLockError)?;
Ok(())
}
pub fn event_already_in(
&self,
event: &Event,
) -> Result<bool, Box<dyn std::error::Error>> {
let Ok(ev_iri) = event.id.named_node() else {
return Err(Box::from("Invalid event ID"));
};
let results = self
.store
.quads_for_pattern(Some((&ev_iri).into()), None, None, None)
.collect::<Result<Vec<_>, _>>()?;
if results.len() > 0 {
return Ok(true);
}
Ok(false)
}
pub fn feed_nobject(
&self,
obj: Box<(dyn NostraObject + Send)>,
) -> Result<bool, LDError> {
match thread::spawn(move || match ttlify_nobject(obj) {
Ok(ttl) => Some(ttl),
Err(_e) => None,
})
.join()
{
Ok(Some(ttl)) => self
.store_ttl(ttl)
.map_err(|_| LDError::TTLSerializationError),
Ok(None) => Err(LDError::TTLSerializationError),
Err(_e) => Err(LDError::TTLSerializationError),
}
}
pub fn feed_event(&self, event: Box<Event>) -> Result<bool, LDError> {
match thread::spawn(move || match ttlify_event_sync(&*event) {
Ok(ttl) => Some(ttl),
Err(_e) => None,
})
.join()
{
Ok(Some(ttl)) => self
.store_ttl(ttl)
.map_err(|_| LDError::TTLSerializationError),
Ok(None) => Err(LDError::TTLSerializationError),
Err(_e) => Err(LDError::TTLSerializationError),
}
}
pub async fn feed_event_async(
&self,
event: &Event,
) -> Result<(), RdfStoreError> {
match ntify_event(event).await {
Ok(nt) => {
let _ = self.store_event_nt(event, nt);
Ok(())
}
Err(_e) => Err(RdfStoreError::LDErr(LDError::NTSerializationError)),
}
}
pub(crate) fn store_ttl(&self, ttl: String) -> Result<bool, RdfStoreError> {
for tri in TurtleParser::new().for_slice(ttl.as_bytes()) {
let Ok(triple) = tri else {
continue;
};
let _ = self.store.insert(QuadRef::new(
&triple.subject,
&triple.predicate,
&triple.object,
&GraphName::DefaultGraph,
));
}
Ok(true)
}
pub(crate) fn store_event_nt(
&self,
event: &Event,
nt: String,
) -> Result<(), RdfStoreError> {
let triples: Vec<_> = NTriplesParser::new()
.for_slice(nt.as_bytes())
.filter_map(|tri| tri.ok())
.collect();
if let Ok(cbk) = self.ev_callbacks.read() {
if let Some(kcallbacks) = cbk.get(&event.kind.as_u16()) {
for callback in kcallbacks {
(callback)(&event, &triples)
}
}
}
Ok(self.store.transaction(|mut transaction| {
for triple in &triples {
let _ = transaction.insert(QuadRef::new(
&triple.subject,
&triple.predicate,
&triple.object,
&GraphName::DefaultGraph,
));
}
Result::<_, StorageError>::Ok(())
})?)
}
#[allow(dead_code)]
pub(crate) fn store_nt(&self, nt: String) -> Result<(), RdfStoreError> {
Ok(self.store.transaction(|mut transaction| {
for tri in NTriplesParser::new().for_slice(nt.as_bytes()) {
let Ok(triple) = tri else {
continue;
};
let _ = transaction.insert(QuadRef::new(
&triple.subject,
&triple.predicate,
&triple.object,
&GraphName::DefaultGraph,
));
}
Result::<_, StorageError>::Ok(())
})?)
}
pub(crate) fn insert_event(
&self,
event: &Event,
) -> Result<(), LDError> {
let quads = event_quads(event)?;
if let Ok(cbk) = self.ev_quads_callbacks.read() {
if let Some(kcallbacks) = cbk.get(&event.kind.as_u16()) {
for callback in kcallbacks {
(callback)(&event, &quads)
}
}
}
for ev_quad in &quads {
if let Err(e) = self.store.insert(ev_quad) {
eprintln!("Failed to store quad: {e:?}")
}
}
Ok(())
}
}