nostralink 0.2.4

Linked data library for nostr
Documentation
//! Event state and event store operations

use super::nostrdb::DATABASE_AGENT_NODE;
use super::prelude::*;
use nostr_sdk::prelude::*;
use std::time::Duration;

pub trait EventStateOperations {
    /// Tag an event as seen
    fn mark_seen(
        &self,
        agent_pubk: &PublicKey,
        event_id: &EventId,
        seen: bool,
    ) -> Result<bool, RdfStoreError>;

    fn mark_read(
        &self,
        agent_pubk: &PublicKey,
        event_id: &EventId,
    ) -> Result<bool, RdfStoreError>;

    fn mark_row_seen<T: Clone>(
        &self,
        agent_pubk: &PublicKey,
        row: &TRdfResultRow<T>,
        seen: bool,
    ) -> Result<bool, RdfStoreError>;

    fn mark_event_notify_reply(
        &self,
        for_pubk: &PublicKey,
        event_id: &EventId,
    ) -> Result<bool, RdfStoreError>;

    /// Tag an event as pinned until a certain time
    /// The event cannot be deleted (garbaged-collected) until it is unpinned
    fn mark_event_pinned_until(
        &self,
        event_id: &EventId,
        until: Timestamp,
    ) -> Result<(), RdfStoreError>;

    /// Unpin an event
    fn mark_event_unpin(&self, event_id: &EventId)
        -> Result<(), RdfStoreError>;
}

/// User (pubk) has seen an event (used mostly on notes) at a given time
pub const PRED_USER_SEEN_AT: NamedNodeRef<'static> =
    NamedNodeRef::new_unchecked("https://w3id.org/nostr#user_seen_at");

pub const PRED_USER_READ_AT: NamedNodeRef<'static> =
    NamedNodeRef::new_unchecked("https://w3id.org/nostr#user_read_at");

/// User (pubk) has seen an event
pub const PRED_USER_SEEN_BY: NamedNodeRef<'static> =
    NamedNodeRef::new_unchecked("https://w3id.org/nostr#user_seen_by");

pub const PRED_USER_READ_BY: NamedNodeRef<'static> =
    NamedNodeRef::new_unchecked("https://w3id.org/nostr#user_read_by");

/// Predicate to mark an event (and anything related to it) as pinned until
/// a certain time (won't be garbage-collected by sweeping calls)
pub const PRED_EVENT_PINNED_UNTIL: NamedNodeRef<'static> =
    NamedNodeRef::new_unchecked("https://w3id.org/nostr#event_pinned_until");

pub const PRED_EVENT_INBOX_NOTIFICATION: NamedNodeRef<'static> =
    NamedNodeRef::new_unchecked(
        "https://w3id.org/nostr#event_inbox_notification",
    );

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(
            // Already added
            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_read(
        &self,
        agent_pubk: &PublicKey,
        event_id: &EventId,
    ) -> 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_READ_BY,
            &agent_nn,
            &GraphName::DefaultGraph,
        ))? {
            return Ok(false);
        }

        if !self
            .store
            .quads_for_pattern(
                Some((&event_nn).into()),
                Some(PRED_USER_READ_AT),
                None,
                None,
            )
            .collect::<Result<Vec<_>, _>>()?
            .is_empty()
        {
            return Ok(false);
        }

        let Ok(_) = self.store.insert(QuadRef::new(
            &event_nn,
            PRED_USER_READ_AT,
            &now_literal(),
            &GraphName::DefaultGraph,
        )) else {
            return Err(RdfStoreError::QuadInsertError);
        };

        let Ok(_) = self.store.insert(QuadRef::new(
            &event_nn,
            PRED_USER_READ_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(
            // Already added
            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_notify_reply(
        &self,
        notified_pubk: &PublicKey,
        event_id: &EventId,
    ) -> Result<bool, RdfStoreError> {
        let notified_nn = notified_pubk.named_node()?;
        let event_nn = event_id.named_node()?;

        self.store
            .insert(QuadRef::new(
                &event_nn,
                PRED_EVENT_INBOX_NOTIFICATION,
                &notified_nn,
                &GraphName::DefaultGraph,
            ))
            .map_err(|_| RdfStoreError::QuadError)
    }

    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_secs()),
            &GraphName::DefaultGraph,
        ))?;

        // The database agent will refuse to store the event again
        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 {
    /// Delete events older than the given timestamp
    /// Only pinned events can survive this function ^_^
    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 {
            // Build a FILTER for these requested event kinds
            let values = kinds
                .iter()
                .map(|kind| format!("(?kind = {})", kind))
                .collect::<Vec<_>>()
                .join(" || ");

            kf.push_str(&format!("FILTER({})", values));
        }

        if let Some(pubkl) = pubkeys {
            // Build a FILTER for these 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_secs()))
            .replace("@TS@", &format!("{}", ts.as_secs()));

        match self.store.update(
            Update::parse(&q, None).map_err(|_| RdfStoreError::QueryError)?,
        ) {
            Ok(_r) => Ok(()),
            Err(_e) => Err(RdfStoreError::QueryError),
        }
    }

    /// Sweep the store, purging old events to maintain an optimal store size
    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,
        )
    }
}