use super::manager::{RdfEventsStore, RdfStoreError};
use super::util::*;
use crate::querydb::nrq_get;
use nostr::{Filter, Timestamp};
use oxigraph::model::*;
use oxigraph::sparql::{QueryOptions, QueryResults};
use std::thread;
pub trait InterStoreOperations {
fn leak(
&self,
oth_store: &RdfEventsStore,
filter: Filter,
) -> Result<(), RdfStoreError>;
}
impl InterStoreOperations for RdfEventsStore {
fn leak(
&self,
dst_store: &RdfEventsStore,
filter: Filter,
) -> Result<(), RdfStoreError> {
let mut subs = Vec::new();
let ts_since = match filter.since {
Some(ts) => ts,
None => Timestamp::now() - (3600 * 12),
};
let ts_until = match filter.until {
Some(ts) => ts,
None => Timestamp::now(),
};
subs.push(
subl("created_at_since", Literal::from(ts_since.as_u64()))
.map_err(|_| RdfStoreError::SubstitutionError)?,
);
subs.push(
subl("created_at_until", Literal::from(ts_until.as_u64()))
.map_err(|_| RdfStoreError::SubstitutionError)?,
);
let q = self.prepare_query(
&nrq_get("event_construct")
.map_err(|_| RdfStoreError::QueryError)?,
);
thread::scope(|s| {
s.spawn(move || {
match self.store.query_opt_with_substituted_variables(
&q,
QueryOptions::default(),
subs,
) {
Ok(QueryResults::Graph(triples)) => {
for triple in triples {
let Ok(tri) = triple else {
continue;
};
if tri.subject.is_blank_node()
|| tri.object.is_blank_node()
{
continue;
}
if let Err(err) =
dst_store.store.insert(QuadRef::new(
&tri.subject,
&tri.predicate,
&tri.object,
GraphNameRef::DefaultGraph,
))
{
eprintln!("{err}");
}
}
}
Err(e) => println!("{e:?}"),
_ => {}
}
});
});
Ok(())
}
}
#[tokio::test]
async fn test_interstore() -> Result<(), Box<dyn std::error::Error>> {
use crate::prelude::*;
let filter = Filter::new();
let keys = Keys::generate();
let mut note = NoteBuilder::default()
.language(Language::FrFr)
.headline(Some("test".to_string()))
.body(Some("Hello".to_string()))
.build()?;
let event = note.event_builder()?.sign(&keys).await?;
let ttl = ttlify_event(&event).await.unwrap();
let store0 = RdfEventsStore::new(Store::new()?)?;
let _ = store0.store_ttl(ttl);
let store1 = RdfEventsStore::new(Store::new()?)?;
let _ = store0.leak(&store1, filter)?;
Ok(())
}