use super::nostrdb::DATABASE_AGENT_NODE;
use super::prelude::*;
use nostr_sdk::prelude::*;
use std::time::Duration;
pub trait EventStateOperations {
fn mark_seen(
&self,
agent_pubk: &PublicKey,
event_id: &EventId,
seen: bool,
) -> Result<bool, RdfStoreError>;
fn mark_row_seen<T: Clone>(
&self,
agent_pubk: &PublicKey,
row: &TRdfResultRow<T>,
seen: bool,
) -> Result<bool, RdfStoreError>;
fn mark_event_pinned_until(
&self,
event_id: &EventId,
until: Timestamp,
) -> Result<(), RdfStoreError>;
fn mark_event_unpin(&self, event_id: &EventId)
-> Result<(), RdfStoreError>;
}
pub const PRED_USER_SEEN_AT: NamedNodeRef<'static> =
NamedNodeRef::new_unchecked("https://w3id.org/nostr#user_seen_at");
pub const PRED_USER_SEEN_BY: NamedNodeRef<'static> =
NamedNodeRef::new_unchecked("https://w3id.org/nostr#user_seen_by");
pub const PRED_EVENT_PINNED_UNTIL: NamedNodeRef<'static> =
NamedNodeRef::new_unchecked("https://w3id.org/nostr#event_pinned_until");
pub const PRED_EVENT_BYEBYE: NamedNodeRef<'static> =
NamedNodeRef::new_unchecked("https://w3id.org/nostr#event_byebye");
impl EventStateOperations for RdfEventsStore {
fn mark_seen(
&self,
agent_pubk: &PublicKey,
event_id: &EventId,
_seen: bool,
) -> Result<bool, RdfStoreError> {
let agent_nn = agent_pubk.named_node()?;
let event_nn = event_id.named_node()?;
if self.store.contains(
QuadRef::new(
&event_nn,
PRED_USER_SEEN_BY,
&agent_nn,
&GraphName::DefaultGraph,
),
)? {
return Ok(false);
}
if !self
.store
.quads_for_pattern(
Some((&event_nn).into()),
Some(PRED_USER_SEEN_AT),
None,
None,
)
.collect::<Result<Vec<_>, _>>()?
.is_empty()
{
return Ok(false);
}
let Ok(_) = self.store.insert(QuadRef::new(
&event_nn,
PRED_USER_SEEN_AT,
&now_literal(),
&GraphName::DefaultGraph,
)) else {
return Err(RdfStoreError::QuadInsertError);
};
let Ok(_) = self.store.insert(QuadRef::new(
&event_nn,
PRED_USER_SEEN_BY,
&agent_nn,
&GraphName::DefaultGraph,
)) else {
return Err(RdfStoreError::QuadInsertError);
};
Ok(true)
}
fn mark_row_seen<T: Clone>(
&self,
agent_pubk: &PublicKey,
row: &TRdfResultRow<T>,
_seen: bool,
) -> Result<bool, RdfStoreError> {
let agent_nn = agent_pubk.named_node()?;
let Some(event) = row.get(SPVars::EVENT) else {
return Err(RdfStoreError::BadEventError);
};
let event_nn = event.value.to_named_node()?;
if self.store.contains(
QuadRef::new(
&event_nn,
PRED_USER_SEEN_BY,
&agent_nn,
&GraphName::DefaultGraph,
),
)? {
return Ok(false);
}
let Ok(_) = self.store.insert(QuadRef::new(
&event_nn,
PRED_USER_SEEN_AT,
&now_literal(),
&GraphName::DefaultGraph,
)) else {
return Err(RdfStoreError::QuadInsertError);
};
let Ok(_) = self.store.insert(QuadRef::new(
&event_nn,
PRED_USER_SEEN_BY,
&agent_nn,
&GraphName::DefaultGraph,
)) else {
return Err(RdfStoreError::QuadInsertError);
};
Ok(true)
}
fn mark_event_pinned_until(
&self,
event_id: &EventId,
until: Timestamp,
) -> Result<(), RdfStoreError> {
let event_nn = event_id.named_node()?;
self.store.insert(QuadRef::new(
&event_nn,
PRED_EVENT_PINNED_UNTIL,
&Literal::from(until.as_u64()),
&GraphName::DefaultGraph,
))?;
self.store.insert(QuadRef::new(
DATABASE_AGENT_NODE,
PRED_EVENT_BYEBYE,
&event_nn,
&GraphName::DefaultGraph,
))?;
Ok(())
}
fn mark_event_unpin(
&self,
event_id: &EventId,
) -> Result<(), RdfStoreError> {
let event_nn = event_id.named_node()?;
for quad in self
.store
.quads_for_pattern(
Some((&event_nn).into()),
Some(PRED_EVENT_PINNED_UNTIL),
None,
None,
)
.collect::<Result<Vec<_>, _>>()?
{
let _ = self.store.remove(QuadRef::new(
&event_nn,
PRED_EVENT_PINNED_UNTIL,
&quad.object,
&GraphName::DefaultGraph,
));
}
Ok(())
}
}
pub trait EventStoreOperations {
fn sweep(&self, duration: Duration) -> Result<(), RdfStoreError>;
fn delete_events_older_than(
&self,
ts: Timestamp,
ekinds: Option<Vec<Kind>>,
pubkeys: Option<Vec<PublicKey>>,
) -> Result<(), RdfStoreError>;
}
impl EventStoreOperations for RdfEventsStore {
fn delete_events_older_than(
&self,
ts: Timestamp,
ekinds: Option<Vec<Kind>>,
pubkeys: Option<Vec<PublicKey>>,
) -> Result<(), RdfStoreError> {
let now = Timestamp::now();
let mut kf = String::new();
let mut pubkf = String::new();
if let Some(kinds) = ekinds {
let values = kinds
.iter()
.map(|kind| format!("(?kind = {})", kind))
.collect::<Vec<_>>()
.join(" || ");
kf.push_str(&format!("FILTER({})", values));
}
if let Some(pubkl) = pubkeys {
let values = pubkl
.iter()
.map(|pk| format!("(?pubk = {})", pk.to_hex()))
.collect::<Vec<_>>()
.join(" || ");
pubkf.push_str(&format!("FILTER({})", values));
}
let q = self
.prepare_query(&nrq_get("delete_events_before")?)
.replace("@KIND_FILTER@", &kf)
.replace("@PUBK_FILTER@", &kf)
.replace("@NOW_TS@", &format!("{}", now.as_u64()))
.replace("@TS@", &format!("{}", ts.as_u64()));
match self.store.update(
Update::parse(&q, None).map_err(|_| RdfStoreError::QueryError)?,
) {
Ok(_r) => Ok(()),
Err(_e) => Err(RdfStoreError::QueryError),
}
}
fn sweep(&self, duration_before: Duration) -> Result<(), RdfStoreError> {
let kinds = vec![
Kind::TextNote,
Kind::LongFormTextNote,
Kind::Repost,
Kind::GenericRepost,
Kind::Reaction,
];
self.delete_events_older_than(
Timestamp::now() - duration_before,
Some(kinds),
None,
)
}
}