use super::prelude::*;
use crate::niri::ToNamedNode;
use crate::querydb::nrq_get;
use nostr::{Event, Filter, Kind};
use nostr_database::prelude::*;
use nostr_database::{NostrDatabase, NostrEventsDatabase};
use std::str::FromStr;
use std::time::Duration;
use tokio::time::sleep;
impl RdfEventsStore {
fn sparqlify_filter(&self, filter: &Filter) -> (String, String) {
let mut ckey = String::new();
let mut q =
self.prepare_query(&nrq_get("events_matching_filter").unwrap());
if let Some(ref ids) = filter.ids {
let values = ids
.iter()
.map(|x| format!(r#"'{}'"#, x.to_hex()))
.collect::<Vec<_>>()
.join(",");
q = q.replace(
"@IDS@",
&format!("FILTER(?event_id IN ({}))", values),
);
ckey.push_str(&values);
} else {
q = q.replace("@IDS@", "");
}
if let Some(ref kinds) = filter.kinds {
let kvl = kinds
.iter()
.map(|x| x.to_string())
.collect::<Vec<_>>()
.join(",");
q = q.replace("@KINDS@", &format!("FILTER(?kind IN ({}))", kvl));
ckey.push_str(&kvl);
} else {
q = q.replace("@KINDS@", "");
}
if let Some(ref authors) = filter.authors {
let pk = authors
.iter()
.map(|x| format!(r#"'{}'"#, x.to_hex()))
.collect::<Vec<_>>()
.join(",");
q = q.replace("@PUBKS@", &format!("FILTER(?pubk IN ({}))", pk));
ckey.push_str(&pk);
} else {
q = q.replace("@PUBKS@", "");
}
if let Some(ref ts) = filter.since {
q = q.replace(
"@SINCE@",
&format!("FILTER(?created_at > {})", ts.as_u64()),
);
ckey.push_str(&format!("{}", ts.as_u64()));
} else {
q = q.replace("@SINCE@", "");
}
if let Some(ref ts) = filter.until {
q = q.replace(
"@UNTIL@",
&format!("FILTER(?created_at < {})", ts.as_u64()),
);
ckey.push_str(&format!("{}", ts.as_u64()));
} else {
q = q.replace("@UNTIL@", "");
}
(q, ckey)
}
}
pub fn event_chan_prio(event: &Event) -> i32 {
match event.kind {
Kind::TextNote | Kind::LongFormTextNote => 100,
Kind::Repost | Kind::GenericRepost => 50,
Kind::Metadata => 220,
Kind::RelayList | Kind::InboxRelays => 200,
Kind::ContactList => 200,
Kind::Reaction => 10,
Kind::Custom(39089) => 90,
_ => 0,
}
}
impl NostrDatabase for RdfEventsStore {
fn backend(&self) -> Backend {
Backend::Custom("RDF".to_string())
}
}
impl NostrEventsDatabase for RdfEventsStore {
fn save_event<'a>(
&'a self,
event: &'a Event,
) -> BoxedFuture<'a, Result<SaveEventStatus, DatabaseError>> {
Box::pin(async move {
match event.kind {
Kind::TextNote
| Kind::LongFormTextNote
| Kind::ContactList
| Kind::Repost
| Kind::Metadata
| Kind::Reaction
| Kind::ZapReceipt
| Kind::Custom(39089)
| Kind::Custom(7101)
| Kind::Custom(7102)
| Kind::Custom(7103) => {
self.process_event(
event.clone(),
Some(event_chan_prio(&event)),
);
Ok(SaveEventStatus::Success)
}
_ => Ok(SaveEventStatus::Rejected(RejectedReason::Other)),
}
})
}
fn check_id<'a>(
&'a self,
event_id: &'a EventId,
) -> BoxedFuture<'a, Result<DatabaseEventStatus, DatabaseError>> {
Box::pin(async move {
let enn = event_id.named_node().map_err(|_| {
DatabaseError::Backend(Box::from("Cannot parse event id"))
})?;
if self
.store
.quads_for_pattern(Some((&enn).into()), None, None, None)
.count()
> 0
{
return Ok(DatabaseEventStatus::Saved);
} else {
return Ok(DatabaseEventStatus::NotExistent);
}
})
}
fn has_coordinate_been_deleted<'a>(
&'a self,
_coordinate: &'a CoordinateBorrow<'a>,
_timestamp: &'a Timestamp,
) -> BoxedFuture<'a, Result<bool, DatabaseError>> {
Box::pin(async move { Ok(false) })
}
fn event_by_id<'a>(
&'a self,
event_id: &'a EventId,
) -> BoxedFuture<'a, Result<Option<Event>, DatabaseError>> {
Box::pin(async move {
Ok(self.query(Filter::new().id(*event_id)).await?.first_owned())
})
}
fn count(
&self,
_filter: Filter,
) -> BoxedFuture<Result<usize, DatabaseError>> {
Box::pin(async move { Ok(0) })
}
fn query(
&self,
filter: Filter,
) -> BoxedFuture<Result<Events, DatabaseError>> {
Box::pin(async move {
let mut events: Events = Events::new(&filter);
let (q, ckey) = self.sparqlify_filter(&filter);
if let Ok(set) = self.run_query(&q, [], Some(ckey)) {
for row in &set.rows {
events.insert(Event::new(
row.get(SPVars::EVENT_ID)
.unwrap()
.value
.to_event_id()
.map_err(|_| {
DatabaseError::Backend(Box::from(
"Invalid event ID",
))
})?,
row.get(SPVars::PUBK)
.unwrap()
.value
.to_public_key()
.map_err(|_| {
DatabaseError::Backend(Box::from(
"Invalid pubk",
))
})?,
row.get(SPVars::CREATED_AT)
.unwrap()
.try_into()
.map_err(|_| {
DatabaseError::Backend(Box::from("Invalid TS"))
})?,
Kind::from_str(
&row.get(SPVars::KIND).unwrap().to_string(),
)
.unwrap(),
vec![], row.get(SPVars::CONTENT).unwrap().to_string(),
Signature::from_str(
&row.get(SPVars::SIG).unwrap().to_string(),
)
.map_err(|_| {
DatabaseError::Backend(Box::from(
"Invalid signature",
))
})?,
));
sleep(Duration::from_millis(10)).await;
}
}
Ok(events)
})
}
fn negentropy_items(
&self,
filter: Filter,
) -> BoxedFuture<Result<Vec<(EventId, Timestamp)>, DatabaseError>> {
Box::pin(async move {
let (q, ckey) = self.sparqlify_filter(&filter);
if let Ok(results) = self.run_query(&q, [], Some(ckey)) {
Ok(results
.rows
.iter()
.filter_map(|r| {
let Ok(event_id) = r
.get(SPVars::EVENT_ID)
.unwrap()
.value
.to_event_id()
else {
return None;
};
let Ok(ts) =
r.get(SPVars::CREATED_AT).unwrap().try_into()
else {
return None;
};
Some((event_id, ts))
})
.collect())
} else {
Err(DatabaseError::Backend(Box::from(
"Error running SparQL query",
)))
}
})
}
fn delete(
&self,
_filter: Filter,
) -> BoxedFuture<Result<(), DatabaseError>> {
Box::pin(async move { Err(DatabaseError::NotSupported) })
}
}
impl NostrDatabaseWipe for RdfEventsStore {
#[inline]
fn wipe(&self) -> BoxedFuture<Result<(), DatabaseError>> {
Box::pin(async move { Err(DatabaseError::NotSupported) })
}
}