2018-12-06 12:52:47 -08:00
|
|
|
//! The `gossip_service` module implements the network control plane.
|
2018-06-06 16:36:54 -07:00
|
|
|
|
2018-10-08 19:55:54 -07:00
|
|
|
use cluster_info::ClusterInfo;
|
2018-07-03 21:14:08 -07:00
|
|
|
use service::Service;
|
2018-05-27 18:21:39 -07:00
|
|
|
use std::net::UdpSocket;
|
2018-07-09 13:53:18 -07:00
|
|
|
use std::sync::atomic::{AtomicBool, Ordering};
|
2018-05-27 18:21:39 -07:00
|
|
|
use std::sync::mpsc::channel;
|
|
|
|
use std::sync::{Arc, RwLock};
|
2018-07-03 21:14:08 -07:00
|
|
|
use std::thread::{self, JoinHandle};
|
2018-05-27 18:21:39 -07:00
|
|
|
use streamer;
|
2018-08-09 12:31:34 -07:00
|
|
|
use window::SharedWindow;
|
2018-05-27 18:21:39 -07:00
|
|
|
|
2018-12-06 12:52:47 -08:00
|
|
|
pub struct GossipService {
|
2018-07-09 13:53:18 -07:00
|
|
|
exit: Arc<AtomicBool>,
|
2018-07-03 21:14:08 -07:00
|
|
|
thread_hdls: Vec<JoinHandle<()>>,
|
2018-05-27 18:21:39 -07:00
|
|
|
}
|
|
|
|
|
2018-12-06 12:52:47 -08:00
|
|
|
impl GossipService {
|
2018-05-27 18:21:39 -07:00
|
|
|
pub fn new(
|
2018-10-08 19:55:54 -07:00
|
|
|
cluster_info: &Arc<RwLock<ClusterInfo>>,
|
2018-08-09 12:31:34 -07:00
|
|
|
window: SharedWindow,
|
2018-08-06 12:35:38 -07:00
|
|
|
ledger_path: Option<&str>,
|
2018-08-28 16:32:40 -07:00
|
|
|
gossip_socket: UdpSocket,
|
2018-05-27 18:21:39 -07:00
|
|
|
exit: Arc<AtomicBool>,
|
2018-09-03 02:23:43 -07:00
|
|
|
) -> Self {
|
2018-05-27 18:21:39 -07:00
|
|
|
let (request_sender, request_receiver) = channel();
|
2018-08-28 16:32:40 -07:00
|
|
|
let gossip_socket = Arc::new(gossip_socket);
|
2018-05-27 18:21:39 -07:00
|
|
|
trace!(
|
2018-12-06 12:52:47 -08:00
|
|
|
"GossipService: id: {}, listening on: {:?}",
|
2018-11-15 13:23:26 -08:00
|
|
|
&cluster_info.read().unwrap().my_data().id,
|
2018-08-28 16:32:40 -07:00
|
|
|
gossip_socket.local_addr().unwrap()
|
2018-05-27 18:21:39 -07:00
|
|
|
);
|
2018-09-18 08:02:57 -07:00
|
|
|
let t_receiver =
|
|
|
|
streamer::blob_receiver(gossip_socket.clone(), exit.clone(), request_sender);
|
2018-05-27 18:21:39 -07:00
|
|
|
let (response_sender, response_receiver) = channel();
|
2018-12-06 12:52:47 -08:00
|
|
|
let t_responder = streamer::responder("gossip", gossip_socket, response_receiver);
|
2018-10-08 19:55:54 -07:00
|
|
|
let t_listen = ClusterInfo::listen(
|
|
|
|
cluster_info.clone(),
|
2018-05-27 18:21:39 -07:00
|
|
|
window,
|
2018-08-06 12:35:38 -07:00
|
|
|
ledger_path,
|
2018-05-27 18:21:39 -07:00
|
|
|
request_receiver,
|
|
|
|
response_sender.clone(),
|
|
|
|
exit.clone(),
|
|
|
|
);
|
2018-10-08 19:55:54 -07:00
|
|
|
let t_gossip = ClusterInfo::gossip(cluster_info.clone(), response_sender, exit.clone());
|
2018-05-27 18:21:39 -07:00
|
|
|
let thread_hdls = vec![t_receiver, t_responder, t_listen, t_gossip];
|
2018-12-06 12:52:47 -08:00
|
|
|
Self { exit, thread_hdls }
|
2018-07-09 13:53:18 -07:00
|
|
|
}
|
|
|
|
|
|
|
|
pub fn close(self) -> thread::Result<()> {
|
|
|
|
self.exit.store(true, Ordering::Relaxed);
|
|
|
|
self.join()
|
2018-05-27 18:21:39 -07:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2018-12-06 12:52:47 -08:00
|
|
|
impl Service for GossipService {
|
2018-09-13 14:00:17 -07:00
|
|
|
type JoinReturnType = ();
|
2018-07-03 21:14:08 -07:00
|
|
|
|
|
|
|
fn join(self) -> thread::Result<()> {
|
2018-09-13 14:00:17 -07:00
|
|
|
for thread_hdl in self.thread_hdls {
|
2018-07-03 21:14:08 -07:00
|
|
|
thread_hdl.join()?;
|
|
|
|
}
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2018-05-27 18:21:39 -07:00
|
|
|
#[cfg(test)]
|
|
|
|
mod tests {
|
2018-12-06 12:52:47 -08:00
|
|
|
use super::*;
|
2018-10-08 19:55:54 -07:00
|
|
|
use cluster_info::{ClusterInfo, Node};
|
2018-07-09 13:53:18 -07:00
|
|
|
use std::sync::atomic::AtomicBool;
|
2018-05-27 18:21:39 -07:00
|
|
|
use std::sync::{Arc, RwLock};
|
|
|
|
|
2018-05-30 09:50:28 -07:00
|
|
|
#[test]
|
2018-07-09 17:35:23 -07:00
|
|
|
#[ignore]
|
2018-05-30 09:50:28 -07:00
|
|
|
// test that stage will exit when flag is set
|
|
|
|
fn test_exit() {
|
|
|
|
let exit = Arc::new(AtomicBool::new(false));
|
2018-08-28 16:32:40 -07:00
|
|
|
let tn = Node::new_localhost();
|
2018-11-19 11:25:14 -08:00
|
|
|
let cluster_info = ClusterInfo::new(tn.info.clone());
|
2018-10-08 19:55:54 -07:00
|
|
|
let c = Arc::new(RwLock::new(cluster_info));
|
2018-05-27 18:21:39 -07:00
|
|
|
let w = Arc::new(RwLock::new(vec![]));
|
2018-12-06 12:52:47 -08:00
|
|
|
let d = GossipService::new(&c, w, None, tn.sockets.gossip, exit.clone());
|
2018-07-09 13:53:18 -07:00
|
|
|
d.close().expect("thread join");
|
2018-05-27 18:21:39 -07:00
|
|
|
}
|
|
|
|
}
|