From 1e03f9094f4c2d6954917f1f7a6988fd8c05062b Mon Sep 17 00:00:00 2001 From: Mike Dilger Date: Thu, 15 Feb 2024 11:45:25 +1300 Subject: [PATCH] Handlers to listen to new_events channel --- src/main.rs | 36 +++++++++++++++++++++++++++++++++++- 1 file changed, 35 insertions(+), 1 deletion(-) diff --git a/src/main.rs b/src/main.rs index fd6ccee..330aba4 100644 --- a/src/main.rs +++ b/src/main.rs @@ -213,8 +213,10 @@ struct WebSocketService { impl WebSocketService { async fn handle_websocket_stream(&mut self) -> Result<(), Error> { + // Subscribe to the new_events broadcast channel + let mut new_events = GLOBALS.new_events.subscribe(); + loop { - // We will add more to this later tokio::select! { message_option = self.websocket.next() => { match message_option { @@ -225,6 +227,38 @@ impl WebSocketService { None => break, // websocket must be closed } }, + offset_result = new_events.recv() => { + let offset = offset_result?; + self.handle_new_event(offset).await?; + }, + } + } + + Ok(()) + } + + // If the event matches a subscription they have open, send them the event + async fn handle_new_event(&mut self, new_event_offset: usize) -> Result<(), Error> { + if self.subscriptions.is_empty() { + return Ok(()); + } + + if let Some(event) = GLOBALS + .store + .get() + .unwrap() + .get_event_by_offset(new_event_offset)? + { + 'subs: for (subid, filters) in self.subscriptions.iter() { + for filter in filters.iter() { + if filter.as_filter()?.event_matches(&event)? { + let message = NostrReply::Event(subid, event.clone()); + self.websocket + .send(Message::text(message.as_json())) + .await?; + continue 'subs; + } + } } }