nostralink 0.2.1

Linked data library for nostr
Documentation
//! API to send events between different stores

use super::manager::{RdfEventsStore, RdfStoreError};
use super::util::*;
use crate::querydb::nrq_get;
use nostr::{Filter, Timestamp};

use oxigraph::model::*;
use oxigraph::sparql::{QueryOptions, QueryResults};
use std::thread;

pub trait InterStoreOperations {
    fn leak(
        &self,
        oth_store: &RdfEventsStore,
        filter: Filter,
    ) -> Result<(), RdfStoreError>;
}

impl InterStoreOperations for RdfEventsStore {
    fn leak(
        &self,
        dst_store: &RdfEventsStore,
        filter: Filter,
    ) -> Result<(), RdfStoreError> {
        let mut subs = Vec::new();

        let ts_since = match filter.since {
            Some(ts) => ts,
            None => Timestamp::now() - (3600 * 12),
        };

        let ts_until = match filter.until {
            Some(ts) => ts,
            None => Timestamp::now(),
        };

        subs.push(
            subl("created_at_since", Literal::from(ts_since.as_u64()))
                .map_err(|_| RdfStoreError::SubstitutionError)?,
        );

        subs.push(
            subl("created_at_until", Literal::from(ts_until.as_u64()))
                .map_err(|_| RdfStoreError::SubstitutionError)?,
        );

        let q = self.prepare_query(
            &nrq_get("event_construct")
                .map_err(|_| RdfStoreError::QueryError)?,
        );

        thread::scope(|s| {
            s.spawn(move || {
                match self.store.query_opt_with_substituted_variables(
                    &q,
                    QueryOptions::default(),
                    subs,
                ) {
                    Ok(QueryResults::Graph(triples)) => {
                        for triple in triples {
                            let Ok(tri) = triple else {
                                continue;
                            };

                            if tri.subject.is_blank_node()
                                || tri.object.is_blank_node()
                            {
                                continue;
                            }
                            if let Err(err) =
                                dst_store.store.insert(QuadRef::new(
                                    &tri.subject,
                                    &tri.predicate,
                                    &tri.object,
                                    GraphNameRef::DefaultGraph,
                                ))
                            {
                                eprintln!("{err}");
                            }
                        }
                    }
                    Err(e) => println!("{e:?}"),
                    _ => {}
                }
            });
        });
        Ok(())
    }
}

#[tokio::test]
async fn test_interstore() -> Result<(), Box<dyn std::error::Error>> {
    use crate::prelude::*;

    let filter = Filter::new();

    let keys = Keys::generate();
    let mut note = NoteBuilder::default()
        .language(Language::FrFr)
        .headline(Some("test".to_string()))
        .body(Some("Hello".to_string()))
        .build()?;

    let event = note.event_builder()?.sign(&keys).await?;
    let ttl = ttlify_event(&event).await.unwrap();
    let store0 = RdfEventsStore::new(Store::new()?)?;

    let _ = store0.store_ttl(ttl);

    let store1 = RdfEventsStore::new(Store::new()?)?;

    let _ = store0.leak(&store1, filter)?;

    Ok(())
}