use std::collections::BTreeSet;
use std::future::IntoFuture;
use std::time::Duration;
use futures::StreamExt;
use nostr::event::Event;
use nostr::filter::Filter;
use crate::error::Error;
use crate::future::BoxedFuture;
use crate::relay::constants::DEFAULT_FETCH_EVENTS_LIMIT;
use crate::relay::{Relay, ReqExitPolicy};
#[must_use = "Does nothing unless you await!"]
pub struct FetchEvents<'relay> {
relay: &'relay Relay,
filters: Vec<Filter>,
timeout: Option<Duration>,
policy: ReqExitPolicy,
max_events: usize,
}
impl<'relay> FetchEvents<'relay> {
pub(crate) fn new(relay: &'relay Relay, filters: Vec<Filter>) -> Self {
Self {
relay,
filters,
timeout: None,
policy: ReqExitPolicy::ExitOnEOSE,
max_events: DEFAULT_FETCH_EVENTS_LIMIT,
}
}
#[inline]
pub fn timeout(mut self, timeout: Duration) -> Self {
self.timeout = Some(timeout);
self
}
#[inline]
pub fn policy(mut self, policy: ReqExitPolicy) -> Self {
self.policy = policy;
self
}
#[inline]
pub fn max_events(mut self, max: usize) -> Self {
self.max_events = max;
self
}
}
impl<'relay> IntoFuture for FetchEvents<'relay> {
type Output = Result<BTreeSet<Event>, Error>;
type IntoFuture = BoxedFuture<'relay, Self::Output>;
fn into_future(self) -> Self::IntoFuture {
Box::pin(async move {
let mut events: BTreeSet<Event> = BTreeSet::new();
let mut stream = self
.relay
.stream_events(self.filters)
.maybe_timeout(self.timeout)
.policy(self.policy)
.await?;
while let Some(res) = stream.next().await {
let event: Event = res?;
if events.len() >= self.max_events && !events.contains(&event) {
return Err(Error::limit_exceeded("too many fetched events"));
}
events.insert(event);
}
Ok(events)
})
}
}
#[cfg(test)]
mod tests {
use nostr::event::{EventBuilder, FinalizeEvent, Kind};
use nostr::key::Keys;
use nostr::message::MachineReadablePrefix;
use nostr::nips::nip01::Metadata;
use super::*;
use crate::authenticator::SignerAuthenticator;
use crate::local_relay::*;
use crate::relay::{RelayOptions, RelayStatus};
use crate::test_utils::{
setup_nip42_read_local_relay, setup_relay, setup_relay_with_authenticator,
};
async fn setup_event_fetching_relay(num_events: usize) -> (Relay, MockRelay) {
let mock = MockRelay::run().await.unwrap();
let url = mock.url().await;
let relay = Relay::new(url);
relay.connect();
let keys = Keys::generate();
for i in 0..num_events {
let event = EventBuilder::new(Kind::TextNote, i.to_string())
.finalize(&keys)
.unwrap();
relay.send_event(&event).await.unwrap();
}
(relay, mock)
}
#[tokio::test]
async fn test_fetch_events_ban_relay() {
let opts = LocalRelayTestOptions {
unresponsive_connection: None,
send_random_events: true,
};
let mock = MockRelay::run_with_opts(opts).await.unwrap();
let url = mock.url().await;
let relay = Relay::builder(url)
.opts(
RelayOptions::default()
.verify_subscriptions(true)
.ban_relay_on_mismatch(true),
)
.build();
assert_eq!(relay.status(), RelayStatus::Initialized);
relay
.try_connect()
.timeout(Duration::from_secs(3))
.await
.unwrap();
assert_eq!(relay.status(), RelayStatus::Connected);
let filter = Filter::new().kind(Kind::Metadata);
relay
.fetch_events(filter)
.timeout(Duration::from_secs(3))
.await
.unwrap();
assert_eq!(relay.status(), RelayStatus::Banned);
assert!(!relay.inner.is_running());
}
#[tokio::test]
async fn test_fetch_events_dont_resubscribes_after_auth_required_closed_without_authenticator()
{
let local = setup_nip42_read_local_relay().await;
let keys = Keys::generate();
let event = EventBuilder::new(Kind::TextNote, "Test")
.finalize(&keys)
.unwrap();
local.add_event(event.clone()).await.unwrap();
let url = local.url().await;
let relay: Relay = setup_relay(url).await;
let filter = Filter::new().kind(Kind::TextNote).limit(3);
let err = relay
.fetch_events(filter.clone())
.timeout(Duration::from_secs(5))
.await
.unwrap_err();
assert_eq!(
MachineReadablePrefix::parse(&err.to_string()).unwrap(),
MachineReadablePrefix::AuthRequired
);
}
#[tokio::test]
async fn test_fetch_events_resubscribes_after_auth_required_closed() {
let local = setup_nip42_read_local_relay().await;
let keys = Keys::generate();
let expected = EventBuilder::new(Kind::TextNote, "Test")
.finalize(&keys)
.unwrap();
local.add_event(expected.clone()).await.unwrap();
let authenticator = SignerAuthenticator::new(keys);
let relay: Relay = setup_relay_with_authenticator(local.url().await, authenticator).await;
let filter = Filter::new().kind(Kind::TextNote).limit(1);
let events = relay
.fetch_events(filter)
.timeout(Duration::from_secs(5))
.await
.unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events.first().map(|event| event.id), Some(expected.id));
}
#[tokio::test]
async fn test_fetch_events_exit_on_eose() {
let (relay, _mock) = setup_event_fetching_relay(5).await;
let events = relay
.fetch_events(Filter::new().kind(Kind::TextNote))
.timeout(Duration::from_secs(5))
.await
.unwrap();
assert_eq!(events.len(), 5);
let events = relay
.fetch_events(Filter::new().kind(Kind::TextNote).limit(3))
.timeout(Duration::from_secs(5))
.policy(ReqExitPolicy::ExitOnEOSE)
.await
.unwrap();
assert_eq!(events.len(), 3);
}
#[tokio::test]
async fn test_fetch_events_enforces_buffer_limit() {
let (relay, _mock) = setup_event_fetching_relay(5).await;
let err = relay
.fetch_events(Filter::new().kind(Kind::TextNote))
.max_events(3)
.timeout(Duration::from_secs(5))
.await
.unwrap_err();
assert_eq!(err.kind(), crate::error::ErrorKind::LimitExceeded);
}
#[tokio::test]
async fn test_fetch_events_wait_for_events() {
let (relay, _mock) = setup_event_fetching_relay(5).await;
let events = relay
.fetch_events(Filter::new().kind(Kind::TextNote))
.timeout(Duration::from_secs(15))
.policy(ReqExitPolicy::WaitForEvents(2))
.await
.unwrap();
assert_eq!(events.len(), 2);
let r = relay.clone();
tokio::spawn(async move {
tokio::time::sleep(Duration::from_secs(2)).await;
let keys = Keys::generate();
let event = Metadata::new().name("Test").finalize(&keys).unwrap();
r.send_event(&event).await.unwrap();
});
let events = relay
.fetch_events(Filter::new().kind(Kind::Metadata))
.timeout(Duration::from_secs(5))
.policy(ReqExitPolicy::WaitForEvents(1))
.await
.unwrap();
assert_eq!(events.len(), 1);
}
#[tokio::test]
async fn test_fetch_events_wait_for_events_after_eose() {
let (relay, _mock) = setup_event_fetching_relay(10).await;
let r = relay.clone();
tokio::spawn(async move {
let keys = Keys::generate();
for _ in 0..2 {
tokio::time::sleep(Duration::from_secs(2)).await;
let event = EventBuilder::new(Kind::TextNote, "Additional")
.finalize(&keys)
.unwrap();
r.send_event(&event).await.unwrap();
}
});
let events = relay
.fetch_events(Filter::new().kind(Kind::TextNote).limit(3))
.timeout(Duration::from_secs(15))
.policy(ReqExitPolicy::WaitForEventsAfterEOSE(2))
.await
.unwrap();
assert_eq!(events.len(), 5); }
#[tokio::test]
async fn test_fetch_events_wait_for_duration_after_eose() {
let (relay, _mock) = setup_event_fetching_relay(5).await;
let r = relay.clone();
tokio::spawn(async move {
tokio::time::sleep(Duration::from_secs(2)).await;
let keys = Keys::generate();
for _ in 0..2 {
let event = EventBuilder::new(Kind::TextNote, "Additional")
.finalize(&keys)
.unwrap();
r.send_event(&event).await.unwrap();
tokio::time::sleep(Duration::from_secs(2)).await;
}
});
let events = relay
.fetch_events(Filter::new().kind(Kind::TextNote))
.timeout(Duration::from_secs(15))
.policy(ReqExitPolicy::WaitDurationAfterEOSE(Duration::from_secs(3)))
.await
.unwrap();
assert_eq!(events.len(), 6); }
}