use super::manager::{RdfEventsStore, RdfStoreError};
use super::util::*;
use crate::err::LDError;
use crate::niri::ToNamedNode;
use crate::querydb::nrq_get;
use nostr::PublicKey;
use nostr_sdk::prelude::*;
use oxigraph::model::*;
impl ToNamedNode for RelayUrl {
fn named_node(&self) -> Result<NamedNode, LDError> {
Ok(NamedNode::new_unchecked(self.to_string()))
}
}
pub trait RelaysPrefsManager {
fn set_flags_for_relay(
&self,
pubk: &PublicKey,
relay_node: NamedNode,
flags: AtomicRelayServiceFlags,
) -> Result<(), RdfStoreError>;
fn relay_flags_for_pubk(
&self,
pubk: &PublicKey,
relay_node: &NamedNode,
) -> Result<AtomicRelayServiceFlags, RdfStoreError>;
fn read_relays_for_pubk(
&self,
pubk: &PublicKey,
) -> Result<Vec<RelayUrl>, RdfStoreError>;
fn write_relays_for_pubk(
&self,
pubk: &PublicKey,
) -> Result<Vec<RelayUrl>, RdfStoreError>;
fn relay_set_last_connection_ts(
&self,
relay_node: NamedNodeRef,
ts: Timestamp,
) -> Result<(), RdfStoreError>;
fn relay_exists(
&self,
relay_node: NamedNodeRef,
) -> Result<bool, RdfStoreError>;
}
fn pred_write() -> NamedNode {
NamedNode::new_unchecked("http://nostralink.org/NostrRelay#writes")
}
fn pred_read() -> NamedNode {
NamedNode::new_unchecked("http://nostralink.org/NostrRelay#reads")
}
fn pred_discovery() -> NamedNode {
NamedNode::new_unchecked("http://nostralink.org/NostrRelay#discovers")
}
pub const PRED_RELAY_LAST_CONNECTED_AT: NamedNodeRef<'static> =
NamedNodeRef::new_unchecked(
"http://nostralink.org/NostrRelay#last_connected_at",
);
impl RelaysPrefsManager for RdfEventsStore {
fn set_flags_for_relay(
&self,
pubk: &PublicKey,
relay_node: NamedNode,
flags: AtomicRelayServiceFlags,
) -> Result<(), RdfStoreError> {
let user_nn = pubk.named_node()?;
let qr = QuadRef::new(
&user_nn,
NamedNodeRef::new_unchecked(
"http://nostralink.org/NostrRelay#writes",
),
&relay_node,
&GraphName::DefaultGraph,
);
if flags.has_write() {
let _ = self.store.insert(qr);
} else {
let _ = self.store.remove(qr);
}
let qr = QuadRef::new(
&user_nn,
NamedNodeRef::new_unchecked(
"http://nostralink.org/NostrRelay#reads",
),
&relay_node,
&GraphName::DefaultGraph,
);
if flags.has_read() {
let _ = self.store.insert(qr);
} else {
let _ = self.store.remove(qr);
}
let qr = QuadRef::new(
&user_nn,
NamedNodeRef::new_unchecked(
"http://nostralink.org/NostrRelay#discovers",
),
&relay_node,
&GraphName::DefaultGraph,
);
if flags.has_discovery() {
let _ = self.store.insert(qr);
} else {
let _ = self.store.remove(qr);
}
Ok(())
}
fn relay_flags_for_pubk(
&self,
pubk: &PublicKey,
relay_node: &NamedNode,
) -> Result<AtomicRelayServiceFlags, RdfStoreError> {
let user_nn = pubk.named_node()?;
let flags = AtomicRelayServiceFlags::new(RelayServiceFlags::NONE);
let results = self.store.quads_for_pattern(
Some((&user_nn).into()),
None,
Some((relay_node).into()),
None,
);
for quad_item in results {
let quad = quad_item.map_err(|_| RdfStoreError::QuadError)?;
if quad.predicate == pred_write() {
flags.add(RelayServiceFlags::WRITE);
}
if quad.predicate == pred_read() {
flags.add(RelayServiceFlags::READ);
}
if quad.predicate == pred_discovery() {
flags.add(RelayServiceFlags::DISCOVERY);
}
}
Ok(flags)
}
fn read_relays_for_pubk(
&self,
pubk: &PublicKey,
) -> Result<Vec<RelayUrl>, RdfStoreError> {
let subs = [
subl("query_relay_readers", Literal::from(1))
.map_err(|_| RdfStoreError::SubstitutionError)?,
subl("query_relay_writers", Literal::from(0))
.map_err(|_| RdfStoreError::SubstitutionError)?,
subl("query_relay_discoverers", Literal::from(0))
.map_err(|_| RdfStoreError::SubstitutionError)?,
subnn("nip21", pubk.named_node()?)
.map_err(|_| RdfStoreError::SubstitutionError)?,
];
let results = self
.run_query(
&nrq_get("user_relays")
.map_err(|_| RdfStoreError::QueryError)?,
subs,
None,
)
.map_err(|_| RdfStoreError::QueryError)?;
Ok(results
.rows
.iter()
.filter_map(|row| match row.get("url") {
Some(url_s) => RelayUrl::parse(&url_s.to_string()).ok(),
None => None,
})
.collect())
}
fn write_relays_for_pubk(
&self,
pubk: &PublicKey,
) -> Result<Vec<RelayUrl>, RdfStoreError> {
let subs = [
subl("query_relay_readers", Literal::from(0))
.map_err(|_| RdfStoreError::SubstitutionError)?,
subl("query_relay_writers", Literal::from(1))
.map_err(|_| RdfStoreError::SubstitutionError)?,
subl("query_relay_discoverers", Literal::from(0))
.map_err(|_| RdfStoreError::SubstitutionError)?,
subnn("nip21", pubk.named_node()?)
.map_err(|_| RdfStoreError::SubstitutionError)?,
];
let results = self
.run_query(
&nrq_get("user_relays")
.map_err(|_| RdfStoreError::QueryError)?,
subs,
None,
)
.map_err(|_| RdfStoreError::QueryError)?;
Ok(results
.rows
.iter()
.filter_map(|row| match row.get("url") {
Some(url_s) => RelayUrl::parse(&url_s.to_string()).ok(),
None => None,
})
.collect())
}
fn relay_set_last_connection_ts(
&self,
relay_node: NamedNodeRef,
ts: Timestamp,
) -> Result<(), RdfStoreError> {
for quad in self.store.quads_for_pattern(
Some(relay_node.into()),
Some(PRED_RELAY_LAST_CONNECTED_AT),
None,
None,
) {
self.store.remove(&quad?)?;
}
let _ = self.store.insert(QuadRef::new(
relay_node,
PRED_RELAY_LAST_CONNECTED_AT,
&Literal::from(ts.as_secs()),
&GraphName::DefaultGraph,
));
Ok(())
}
fn relay_exists(
&self,
relay_node: NamedNodeRef,
) -> Result<bool, RdfStoreError> {
Ok(self
.store
.quads_for_pattern(Some(relay_node.into()), None, None, None)
.count()
> 0)
}
}