use super::channel::{QueryChannel, RX};
use super::manager::{RRCache, RdfEventsStore, RdfStoreError};
use super::prelude::*;
use crate::querydb::nrq_get;
use nostr::PublicKey;
use oxigraph::sparql::Update;
use std::{
sync::{atomic::Ordering, Arc},
thread::sleep,
time::Duration,
};
use thread_priority::*;
impl RdfEventsStore {
pub fn run_query_nosubs(
&self,
query: &String,
) -> Result<Arc<RdfResultSet>, RdfStoreError> {
self.run_query(query, [], None)
}
pub fn run_query(
&self,
query: &String,
substitutions: impl IntoIterator<Item = (Variable, Term)>,
cache_key: Option<String>,
) -> Result<Arc<RdfResultSet>, RdfStoreError> {
if cache_key.is_some() {
let cached = self.cache.get(&cache_key.clone().unwrap());
if cached.is_some() {
return Ok(Arc::clone(&cached.unwrap()));
}
}
let mut q = String::new();
q.push_str(&self.prefix_mappings());
q.push_str(query);
let mut column_headings: HashSet<String> = HashSet::new();
let mut result_set_rows: Vec<HashMap<String, RdfCell>> = Vec::new();
if let QueryResults::Solutions(solutions) = self
.store
.query_opt_with_substituted_variables(
&q,
QueryOptions::default(),
substitutions,
)
.map_err(|_| RdfStoreError::QueryError)?
{
for solution in solutions {
let Ok(row) = solution else {
continue;
};
let mut result_set_row: HashMap<String, RdfCell> =
HashMap::new();
for (variable, term) in row.iter() {
let Ok(cell_value) =
RdfCell::new_cell_from_value_term(variable, term)
else {
continue;
};
column_headings.insert(cell_value.name.clone());
result_set_row.insert(cell_value.name.clone(), cell_value);
}
result_set_rows.insert(result_set_rows.len(), result_set_row);
}
}
let rdf_results = RdfResultSet {
when: SystemTime::now(),
column_headings: column_headings.into_iter().collect(),
rows: result_set_rows,
};
if cache_key.is_some() {
self.cache
.insert(cache_key.unwrap(), Arc::new(rdf_results.clone()));
}
Ok(Arc::new(rdf_results))
}
pub fn run_query_typed<T: Clone + Send + Sync + 'static>(
&self,
query: &String,
s_cache: Option<&RRCache<T>>,
substitutions: impl IntoIterator<Item = (Variable, Term)>,
cache_key: Option<String>,
) -> Result<Arc<TRdfResultSet<T>>, RdfStoreError> {
if s_cache.is_some() && cache_key.is_some() {
if let Some(cached) =
s_cache.unwrap().cache.get(&cache_key.clone().unwrap())
{
return Ok(Arc::clone(&cached));
}
}
let mut q = String::new();
q.push_str(&self.prefix_mappings());
q.push_str(query);
let mut column_headings: HashSet<String> = HashSet::new();
let mut result_set_rows: Vec<TRdfResultRow<T>> = Vec::new();
if let QueryResults::Solutions(solutions) = self
.store
.query_opt_with_substituted_variables(
&q,
QueryOptions::default(),
substitutions,
)
.map_err(|_| RdfStoreError::QueryError)?
{
for solution in solutions {
let Ok(row) = solution else {
continue;
};
let mut result_set_row: TRdfResultRow<T> =
TRdfResultRow::<T>(HashMap::new(), PhantomData);
for (variable, term) in row.iter() {
let Ok(cell_value) =
RdfCell::new_cell_from_value_term(variable, term)
else {
continue;
};
column_headings.insert(cell_value.name.clone());
result_set_row
.0
.insert(cell_value.name.clone(), cell_value);
}
result_set_rows.insert(result_set_rows.len(), result_set_row);
}
}
let rdf_results = TRdfResultSet::<T> {
when: SystemTime::now(),
column_headings: column_headings.into_iter().collect(),
rows: result_set_rows,
};
if s_cache.is_some() && cache_key.is_some() {
s_cache
.unwrap()
.cache
.insert(cache_key.unwrap(), Arc::new(rdf_results.clone()));
}
Ok(Arc::new(rdf_results.into()))
}
pub fn send<T: Clone + Send + Sync + 'static>(
&self,
qchannel: &QueryChannel<T>,
query: &String,
substitutions: impl IntoIterator<Item = (Variable, Term)> + Send,
s_cache: Option<&RRCache<T>>,
cache_key: String,
) -> Result<(), RdfStoreError> {
let mut column_headings: HashSet<String> = HashSet::new();
let mut result_set_rows: Vec<TRdfResultRow<T>> = Vec::new();
let q = self.prepare_query(&query);
let cstatus = Arc::clone(&qchannel.status);
std::thread::scope(|s| {
if let Err(e) = set_current_thread_priority(ThreadPriority::Min) {
eprintln!("Error setting thread priority: {e}");
}
s.spawn(move || {
cstatus.store(1, Ordering::Release);
if let Ok(QueryResults::Solutions(solutions)) = self
.store
.query_opt_with_substituted_variables(
&q,
QueryOptions::default(),
substitutions,
)
.map_err(|_| RdfStoreError::SubstitutionError)
{
for solution in solutions {
let Ok(row) =
solution.map_err(|_| RdfStoreError::RdfCellError)
else {
continue;
};
let mut result_set_row: TRdfResultRow<T> =
TRdfResultRow::<T>(HashMap::new(), PhantomData);
for (variable, term) in row.iter() {
let Ok(cell_value) =
RdfCell::new_cell_from_value_term(
variable, term,
)
.map_err(|_| RdfStoreError::RdfCellError)
else {
continue;
};
column_headings.insert(cell_value.name.clone());
result_set_row
.0
.insert(cell_value.name.clone(), cell_value);
}
result_set_rows
.insert(result_set_rows.len(), result_set_row);
sleep(Duration::from_millis(1));
}
let results = Arc::new(TRdfResultSet::<T> {
when: SystemTime::now(),
column_headings: column_headings.into_iter().collect(),
rows: result_set_rows,
});
let c_results = Arc::clone(&results);
if let Err(err) = qchannel.send(&cache_key, results) {
eprintln!(
"Error sending results on {cache_key}: {err:?}"
);
}
if let Some(results_cache) = s_cache {
results_cache
.cache
.insert(cache_key.clone(), c_results);
}
}
cstatus.store(0, Ordering::Release);
});
});
Ok(())
}
pub fn channel_query<'a, T: Clone + Send + Sync>(
self: Arc<RdfEventsStore>,
qchannel: Arc<QueryChannel<T>>,
query: String,
substitutions: impl IntoIterator<Item = (Variable, Term)> + Send + 'static,
cache_key: String,
) -> Result<Arc<TRdfResultSet<T>>, RdfStoreError> {
if let Some(cached) = qchannel.rcache.cache.get(&cache_key.clone()) {
return Ok(cached.clone());
}
let store = Arc::clone(&self);
let qchannel_t = Arc::clone(&qchannel);
let qchannel_r = Arc::clone(&qchannel);
let ckey = cache_key.clone();
let ckey2 = cache_key.clone();
if qchannel_t.get_thread_status(&ckey) != 1 {
std::thread::spawn(move || {
qchannel_t.set_thread_status(&ckey, 1);
sleep(Duration::from_millis(50));
let _ = store.send(
&qchannel_t,
&query,
substitutions,
Some(&qchannel.rcache),
cache_key,
);
qchannel_t.set_thread_status(&ckey, 0);
});
}
match qchannel_r.recv(&ckey2) {
Ok(set) => Ok(set),
Err(_) => Err(RdfStoreError::QueryChannelReceiveError),
}
}
pub fn channel_continuous_query<'a, T: Clone + Send + Sync>(
self: Arc<RdfEventsStore>,
qchannel: Arc<QueryChannel<T>>,
query: String,
substitutions: impl IntoIterator<Item = (Variable, Term)>
+ Send
+ Clone
+ 'static,
cache_key: String,
) -> Result<RX<T>, RdfStoreError> {
let store = Arc::clone(&self);
let qchannel_t = Arc::clone(&qchannel);
let ckey = cache_key.clone();
let reader = qchannel.clone().reader_for_key(&ckey.clone());
if qchannel.get_thread_status(&ckey) != 1 {
std::thread::spawn(move || {
qchannel_t.set_thread_status(&ckey, 1);
std::thread::scope(|s| {
if let Err(e) =
set_current_thread_priority(ThreadPriority::Min)
{
eprintln!("Error setting thread priority: {e}");
}
s.spawn(move || loop {
if let Ok(set) = store.run_query_typed(
&query,
None::<&RRCache<T>>,
substitutions.clone(),
Some(cache_key.clone()),
) {
qchannel
.clone()
.rcache
.cache
.insert(cache_key.clone(), set);
}
sleep(Duration::from_secs(10));
});
});
qchannel_t.set_thread_status(&ckey, 0);
});
}
Ok(reader)
}
pub fn explain_with_subs(
&self,
query: &String,
substitutions: impl IntoIterator<Item = (Variable, Term)>,
) -> Result<Value, RdfStoreError> {
let mut buf = Vec::new();
let q = self.prepare_query(query);
if let (Ok(QueryResults::Solutions(_solutions)), explanation) = self
.store
.explain_query_opt_with_substituted_variables(
&q,
QueryOptions::default(),
true,
substitutions,
)
.map_err(|_| RdfStoreError::QueryExplanationError)?
{
explanation
.write_in_json(&mut buf)
.map_err(|_| RdfStoreError::QueryExplanationError)?;
}
let s: &str = std::str::from_utf8(&buf).unwrap();
let value: Value = serde_json::from_str(s)?;
Ok(value)
}
pub fn delete_metadata_events_for(
&self,
pubk: PublicKey,
) -> Result<(), RdfStoreError> {
let q = self
.prepare_query(&nrq_get("delete_metadata")?)
.replace("@PUBK@", &pubk.to_hex());
let update =
Update::parse(&q, None).map_err(|_| RdfStoreError::QueryError)?;
match self.store.update(update) {
Ok(_r) => {
let _ = self.store.flush();
Ok(())
}
Err(e) => {
eprintln!("SparQL UPDATE error: {e:?}");
Err(RdfStoreError::QueryError)
}
}
}
pub fn delete_previous_events(
&self,
event: &Event,
) -> Result<(), RdfStoreError> {
let q = self
.prepare_query(&nrq_get("delete_previous_events")?)
.replace("@BEFORE_TS@", &event.created_at.to_string())
.replace("@KIND@", &event.kind.to_string())
.replace("@PUBK@", &event.pubkey.to_hex());
let update =
Update::parse(&q, None).map_err(|_| RdfStoreError::QueryError)?;
match self.store.update(update) {
Ok(_r) => {
let _ = self.store.flush();
Ok(())
}
Err(_e) => Err(RdfStoreError::QueryError),
}
}
}