solana/src/ncp.rs

109 lines
2.9 KiB
Rust
Raw Normal View History

2018-06-07 15:06:32 -07:00
//! The `ncp` module implements the network control plane.
2018-06-06 16:36:54 -07:00
2018-06-27 11:33:56 -07:00
use crdt::Crdt;
2018-07-17 15:00:22 -07:00
use packet::BlobRecycler;
use result::Result;
use service::Service;
use std::net::UdpSocket;
2018-07-09 13:53:18 -07:00
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::channel;
use std::sync::{Arc, RwLock};
use std::thread::{self, JoinHandle};
use streamer;
2018-08-09 12:31:34 -07:00
use window::SharedWindow;
2018-06-07 15:06:32 -07:00
pub struct Ncp {
2018-07-09 13:53:18 -07:00
exit: Arc<AtomicBool>,
thread_hdls: Vec<JoinHandle<()>>,
}
2018-06-07 15:06:32 -07:00
impl Ncp {
pub fn new(
crdt: &Arc<RwLock<Crdt>>,
2018-08-09 12:31:34 -07:00
window: SharedWindow,
2018-08-06 12:35:38 -07:00
ledger_path: Option<&str>,
gossip_listen_socket: UdpSocket,
gossip_send_socket: UdpSocket,
exit: Arc<AtomicBool>,
2018-06-07 15:06:32 -07:00
) -> Result<Ncp> {
2018-06-27 11:33:56 -07:00
let blob_recycler = BlobRecycler::default();
let (request_sender, request_receiver) = channel();
trace!(
2018-06-07 15:06:32 -07:00
"Ncp: id: {:?}, listening on: {:?}",
2018-07-31 14:50:09 -07:00
&crdt.read().unwrap().me.as_ref()[..4],
gossip_listen_socket.local_addr().unwrap()
);
let t_receiver = streamer::blob_receiver(
exit.clone(),
blob_recycler.clone(),
gossip_listen_socket,
request_sender,
)?;
let (response_sender, response_receiver) = channel();
let t_responder = streamer::responder(
"ncp",
gossip_send_socket,
blob_recycler.clone(),
response_receiver,
);
2018-06-27 11:33:56 -07:00
let t_listen = Crdt::listen(
crdt.clone(),
window,
2018-08-06 12:35:38 -07:00
ledger_path,
blob_recycler.clone(),
request_receiver,
response_sender.clone(),
exit.clone(),
);
2018-07-09 13:53:18 -07:00
let t_gossip = Crdt::gossip(crdt.clone(), blob_recycler, response_sender, exit.clone());
let thread_hdls = vec![t_receiver, t_responder, t_listen, t_gossip];
2018-07-09 13:53:18 -07:00
Ok(Ncp { exit, thread_hdls })
}
pub fn close(self) -> thread::Result<()> {
self.exit.store(true, Ordering::Relaxed);
self.join()
}
}
impl Service for Ncp {
fn thread_hdls(self) -> Vec<JoinHandle<()>> {
self.thread_hdls
}
fn join(self) -> thread::Result<()> {
for thread_hdl in self.thread_hdls() {
thread_hdl.join()?;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use crdt::{Crdt, TestNode};
2018-06-07 15:06:32 -07:00
use ncp::Ncp;
2018-07-09 13:53:18 -07:00
use std::sync::atomic::AtomicBool;
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));
let tn = TestNode::new_localhost();
let crdt = Crdt::new(tn.data.clone()).expect("Crdt::new");
let c = Arc::new(RwLock::new(crdt));
let w = Arc::new(RwLock::new(vec![]));
2018-06-07 15:06:32 -07:00
let d = Ncp::new(
&c,
w,
2018-08-06 12:35:38 -07:00
None,
tn.sockets.gossip,
tn.sockets.gossip_send,
2018-05-30 09:50:28 -07:00
exit.clone(),
).unwrap();
2018-07-09 13:53:18 -07:00
d.close().expect("thread join");
}
}