node: Add eventing system to `Handler` and `Service`

Signed-off-by: xphoniex <dj.2dixx@gmail.com>
This commit is contained in:
xphoniex 2023-03-16 12:06:13 +00:00 committed by Alexis Sellier
parent ebfd445034
commit a49eec9892
No known key found for this signature in database
8 changed files with 85 additions and 30 deletions

View File

@ -3,6 +3,7 @@ mod handle;
use std::io::{BufRead, BufReader}; use std::io::{BufRead, BufReader};
use std::os::unix::net::UnixListener; use std::os::unix::net::UnixListener;
use std::path::PathBuf; use std::path::PathBuf;
use std::sync::{Arc, Mutex};
use std::{fs, io, net, thread, time}; use std::{fs, io, net, thread, time};
use crossbeam_channel as chan; use crossbeam_channel as chan;
@ -22,7 +23,7 @@ use crate::address;
use crate::control; use crate::control;
use crate::crypto::Signer; use crate::crypto::Signer;
use crate::node::{routing, NodeId}; use crate::node::{routing, NodeId};
use crate::service::tracking; use crate::service::{tracking, Event};
use crate::wire; use crate::wire;
use crate::wire::Wire; use crate::wire::Wire;
use crate::worker; use crate::worker;
@ -61,6 +62,39 @@ pub enum Error {
GitVersion(#[from] git::VersionError), GitVersion(#[from] git::VersionError),
} }
/// Publishes events to subscribers.
#[derive(Debug, Clone)]
pub struct Emitter<T> {
pub(crate) subscribers: Arc<Mutex<Vec<chan::Sender<T>>>>,
}
impl<T> Default for Emitter<T> {
fn default() -> Emitter<T> {
Emitter {
subscribers: Default::default(),
}
}
}
impl<T: Clone> Emitter<T> {
/// Emit event to subscribers and drop those who can't receive it.
pub(crate) fn emit(&self, event: T) {
self.subscribers
.lock()
.unwrap()
.retain(|s| s.try_send(event.clone()).is_ok());
}
/// Subscribe to events stream.
pub fn events(&mut self) -> chan::Receiver<T> {
let (sender, receiver) = chan::unbounded();
let mut subs = self.subscribers.lock().unwrap();
subs.push(sender);
receiver
}
}
/// Holds join handles to the client threads, as well as a client handle. /// Holds join handles to the client threads, as well as a client handle.
pub struct Runtime<G: Signer + Ecdh> { pub struct Runtime<G: Signer + Ecdh> {
pub id: NodeId, pub id: NodeId,
@ -113,6 +147,7 @@ impl<G: Signer + Ecdh + 'static> Runtime<G> {
log::info!(target: "node", "Default tracking policy set to '{}'", &config.policy); log::info!(target: "node", "Default tracking policy set to '{}'", &config.policy);
log::info!(target: "node", "Initializing service ({:?})..", network); log::info!(target: "node", "Initializing service ({:?})..", network);
let emitter: Emitter<Event> = Default::default();
let service = service::Service::new( let service = service::Service::new(
config, config,
clock, clock,
@ -122,6 +157,7 @@ impl<G: Signer + Ecdh + 'static> Runtime<G> {
tracking, tracking,
signer.clone(), signer.clone(),
rng, rng,
emitter.clone(),
); );
let (worker_send, worker_recv) = chan::unbounded::<worker::Task<G>>(); let (worker_send, worker_recv) = chan::unbounded::<worker::Task<G>>();
@ -138,7 +174,7 @@ impl<G: Signer + Ecdh + 'static> Runtime<G> {
log::info!(target: "node", "Listening on {local_addr}.."); log::info!(target: "node", "Listening on {local_addr}..");
} }
let reactor = Reactor::named(wire, popol::Poller::new(), id.to_human())?; let reactor = Reactor::named(wire, popol::Poller::new(), id.to_human())?;
let handle = Handle::new(home.clone(), reactor.controller()); let handle = Handle::new(home.clone(), reactor.controller(), emitter);
let atomic = git::version()? >= git::VERSION_REQUIRED; let atomic = git::version()? >= git::VERSION_REQUIRED;
if !atomic { if !atomic {

View File

@ -13,8 +13,10 @@ use crate::crypto::Signer;
use crate::identity::Id; use crate::identity::Id;
use crate::node::{Command, FetchResult}; use crate::node::{Command, FetchResult};
use crate::profile::Home; use crate::profile::Home;
use crate::runtime::Emitter;
use crate::service; use crate::service;
use crate::service::tracking; use crate::service::tracking;
use crate::service::Event;
use crate::service::{CommandError, QueryState}; use crate::service::{CommandError, QueryState};
use crate::service::{NodeId, Sessions}; use crate::service::{NodeId, Sessions};
use crate::wire; use crate::wire;
@ -64,6 +66,20 @@ pub struct Handle<G: Signer + Ecdh> {
/// Whether a shutdown was initiated or not. Prevents attempting to shutdown twice. /// Whether a shutdown was initiated or not. Prevents attempting to shutdown twice.
shutdown: Arc<AtomicBool>, shutdown: Arc<AtomicBool>,
/// Publishes events to subscribers.
emitter: Emitter<Event>,
}
impl<G: Signer + Ecdh> Handle<G> {
/// Subscribe to events stream.
pub fn events(&mut self) -> chan::Receiver<Event> {
let (sender, receiver) = chan::unbounded();
let mut subs = self.emitter.subscribers.lock().unwrap();
subs.push(sender);
receiver
}
} }
impl<G: Signer + Ecdh> fmt::Debug for Handle<G> { impl<G: Signer + Ecdh> fmt::Debug for Handle<G> {
@ -78,16 +94,22 @@ impl<G: Signer + Ecdh> Clone for Handle<G> {
home: self.home.clone(), home: self.home.clone(),
controller: self.controller.clone(), controller: self.controller.clone(),
shutdown: self.shutdown.clone(), shutdown: self.shutdown.clone(),
emitter: self.emitter.clone(),
} }
} }
} }
impl<G: Signer + Ecdh + 'static> Handle<G> { impl<G: Signer + Ecdh + 'static> Handle<G> {
pub fn new(home: Home, controller: reactor::Controller<wire::Control<G>>) -> Self { pub fn new(
home: Home,
controller: reactor::Controller<wire::Control<G>>,
emitter: Emitter<Event>,
) -> Self {
Self { Self {
home, home,
controller, controller,
shutdown: Arc::default(), shutdown: Arc::default(),
emitter,
} }
} }

View File

@ -29,6 +29,7 @@ use crate::node;
use crate::node::routing; use crate::node::routing;
use crate::node::{Address, Features, FetchResult, Seed, Seeds}; use crate::node::{Address, Features, FetchResult, Seed, Seeds};
use crate::prelude::*; use crate::prelude::*;
use crate::runtime::Emitter;
use crate::service::message::{Announcement, AnnouncementMessage, Ping}; use crate::service::message::{Announcement, AnnouncementMessage, Ping};
use crate::service::message::{NodeAnnouncement, RefsAnnouncement}; use crate::service::message::{NodeAnnouncement, RefsAnnouncement};
use crate::service::reactor::FetchDirection; use crate::service::reactor::FetchDirection;
@ -203,6 +204,8 @@ pub struct Service<R, A, S, G> {
last_announce: LocalTime, last_announce: LocalTime,
/// Time when the service was initialized. /// Time when the service was initialized.
start_time: LocalTime, start_time: LocalTime,
/// Publishes events to subscribers.
emitter: Emitter<Event>,
} }
impl<R, A, S, G> Service<R, A, S, G> impl<R, A, S, G> Service<R, A, S, G>
@ -236,6 +239,7 @@ where
tracking: tracking::Config, tracking: tracking::Config,
signer: G, signer: G,
rng: Rng, rng: Rng,
emitter: Emitter<Event>,
) -> Self { ) -> Self {
let sessions = Sessions::new(rng.clone()); let sessions = Sessions::new(rng.clone());
@ -261,6 +265,7 @@ where
last_prune: LocalTime::default(), last_prune: LocalTime::default(),
last_announce: LocalTime::default(), last_announce: LocalTime::default(),
start_time: LocalTime::default(), start_time: LocalTime::default(),
emitter,
} }
} }
@ -341,6 +346,11 @@ where
&self.signer &self.signer
} }
/// Subscriber to inner `Emitter` events.
pub fn events(&mut self) -> chan::Receiver<Event> {
self.emitter.events()
}
/// Get I/O reactor. /// Get I/O reactor.
pub fn reactor(&mut self) -> &mut Reactor { pub fn reactor(&mut self) -> &mut Reactor {
&mut self.reactor &mut self.reactor
@ -586,11 +596,12 @@ where
Ok(updated) => { Ok(updated) => {
log::debug!(target: "service", "Fetched {rid} from {remote}"); log::debug!(target: "service", "Fetched {rid} from {remote}");
self.reactor.event(Event::RefsFetched { self.emitter.emit(Event::RefsFetched {
remote, remote,
rid, rid,
updated: updated.clone(), updated: updated.clone(),
}); });
FetchResult::Success { updated } FetchResult::Success { updated }
} }
Err(err) => { Err(err) => {

View File

@ -22,8 +22,6 @@ pub enum Io {
Fetch(Fetch), Fetch(Fetch),
/// Ask for a wakeup in a specified amount of time. /// Ask for a wakeup in a specified amount of time.
Wakeup(LocalDuration), Wakeup(LocalDuration),
/// Emit an event.
Event(Event),
} }
/// Fetch job sent to worker thread. /// Fetch job sent to worker thread.
@ -81,11 +79,6 @@ pub struct Reactor {
} }
impl Reactor { impl Reactor {
/// Emit an event.
pub fn event(&mut self, event: Event) {
self.io.push_back(Io::Event(event));
}
/// Connect to a peer. /// Connect to a peer.
pub fn connect(&mut self, id: NodeId, addr: Address) { pub fn connect(&mut self, id: NodeId, addr: Address) {
self.io.push_back(Io::Connect(id, addr)); self.io.push_back(Io::Connect(id, addr));

View File

@ -3,6 +3,7 @@ use std::iter;
use std::net; use std::net;
use std::ops::{Deref, DerefMut}; use std::ops::{Deref, DerefMut};
use crossbeam_channel as chan;
use log::*; use log::*;
use crate::address; use crate::address;
@ -13,6 +14,7 @@ use crate::identity::Id;
use crate::node; use crate::node;
use crate::node::routing; use crate::node::routing;
use crate::prelude::*; use crate::prelude::*;
use crate::runtime::Emitter;
use crate::service; use crate::service;
use crate::service::message::*; use crate::service::message::*;
use crate::service::reactor::Io; use crate::service::reactor::Io;
@ -130,6 +132,8 @@ where
let tracking = tracking::Store::memory().unwrap(); let tracking = tracking::Store::memory().unwrap();
let tracking = tracking::Config::new(config.policy, config.scope, tracking); let tracking = tracking::Config::new(config.policy, config.scope, tracking);
let id = *config.signer.public_key(); let id = *config.signer.public_key();
let emitter: Emitter<Event> = Default::default();
let service = Service::new( let service = Service::new(
config.config, config.config,
config.local_time, config.local_time,
@ -139,6 +143,7 @@ where
tracking, tracking,
config.signer, config.signer,
config.rng.clone(), config.rng.clone(),
emitter,
); );
let ip = ip.into(); let ip = ip.into();
let local_addr = net::SocketAddr::new(ip, config.rng.u16(..)); let local_addr = net::SocketAddr::new(ip, config.rng.u16(..));
@ -320,10 +325,9 @@ where
msgs.into_iter() msgs.into_iter()
} }
/// Get a draining iterator over the peer's emitted events. /// Get a stream of the peer's emitted events.
pub fn events(&mut self) -> impl Iterator<Item = Event> + '_ { pub fn events(&mut self) -> chan::Receiver<Event> {
self.outbox() self.service.events()
.filter_map(|io| if let Io::Event(e) = io { Some(e) } else { None })
} }
/// Get a draining iterator over the peer's I/O outbox. /// Get a draining iterator over the peer's I/O outbox.

View File

@ -606,14 +606,6 @@ impl<S: WriteStorage + 'static, G: Signer> Simulation<S, G> {
); );
} }
} }
Io::Event(event) => {
let events = self.events.entry(node).or_insert_with(VecDeque::new);
if events.len() >= MAX_EVENTS {
warn!(target: "sim", "Dropping event: buffer is full");
} else {
events.push_back(event);
}
}
Io::Fetch(fetch) => { Io::Fetch(fetch) => {
let remote = fetch.remote; let remote = fetch.remote;

View File

@ -1069,6 +1069,8 @@ fn test_push_and_pull() {
) )
.initialize([&mut alice, &mut bob, &mut eve]); .initialize([&mut alice, &mut bob, &mut eve]);
let bob_events = bob.events();
// Here we expect Alice to connect to Eve. // Here we expect Alice to connect to Eve.
sim.run_while([&mut alice, &mut bob, &mut eve], |s| !s.is_settled()); sim.run_while([&mut alice, &mut bob, &mut eve], |s| !s.is_settled());
@ -1098,8 +1100,8 @@ fn test_push_and_pull() {
.unwrap() .unwrap()
.is_some()); .is_some());
assert_matches!( assert_matches!(
sim.events(&bob.id).next(), bob_events.recv(),
Some(service::Event::RefsFetched { remote, .. }) Ok(service::Event::RefsFetched { remote, .. })
if remote == eve.node_id(), if remote == eve.node_id(),
"Bob fetched from Eve" "Bob fetched from Eve"
); );

View File

@ -703,11 +703,6 @@ where
} }
self.actions.push_back(reactor::Action::Send(fd, data)); self.actions.push_back(reactor::Action::Send(fd, data));
} }
Io::Event(_e) => {
log::warn!(
target: "wire", "Events are not currently supported"
);
}
Io::Connect(node_id, addr) => { Io::Connect(node_id, addr) => {
if self.connected().any(|(_, id)| id == &node_id) { if self.connected().any(|(_, id)| id == &node_id) {
log::error!( log::error!(