Clean up gossip service exit flag handling
This commit is contained in:
parent
eb90d8d463
commit
e7cde846cb
|
@ -16,7 +16,6 @@ use std::thread::{self, JoinHandle};
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
||||||
pub struct GossipService {
|
pub struct GossipService {
|
||||||
exit: Arc<AtomicBool>,
|
|
||||||
thread_hdls: Vec<JoinHandle<()>>,
|
thread_hdls: Vec<JoinHandle<()>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -53,15 +52,7 @@ impl GossipService {
|
||||||
exit.clone(),
|
exit.clone(),
|
||||||
);
|
);
|
||||||
let thread_hdls = vec![t_receiver, t_responder, t_listen, t_gossip];
|
let thread_hdls = vec![t_receiver, t_responder, t_listen, t_gossip];
|
||||||
Self {
|
Self { thread_hdls }
|
||||||
exit: exit.clone(),
|
|
||||||
thread_hdls,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
pub fn close(self) -> thread::Result<()> {
|
|
||||||
self.exit.store(true, Ordering::Relaxed);
|
|
||||||
self.join()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -99,8 +90,10 @@ pub fn discover(entry_point_info: &NodeInfo, num_nodes: usize) -> Vec<NodeInfo>
|
||||||
//TODO: deprecate this in favor of discover
|
//TODO: deprecate this in favor of discover
|
||||||
pub fn converge(node: &NodeInfo, num_nodes: usize) -> Vec<NodeInfo> {
|
pub fn converge(node: &NodeInfo, num_nodes: usize) -> Vec<NodeInfo> {
|
||||||
info!("Wait for convergence with {} nodes", num_nodes);
|
info!("Wait for convergence with {} nodes", num_nodes);
|
||||||
|
|
||||||
|
let exit = Arc::new(AtomicBool::new(false));
|
||||||
// Let's spy on the network
|
// Let's spy on the network
|
||||||
let (gossip_service, spy_ref, id) = make_spy_node(node);
|
let (gossip_service, spy_ref, id) = make_spy_node(node, &exit);
|
||||||
trace!(
|
trace!(
|
||||||
"converge spy_node {} looking for at least {} nodes",
|
"converge spy_node {} looking for at least {} nodes",
|
||||||
id,
|
id,
|
||||||
|
@ -117,7 +110,8 @@ pub fn converge(node: &NodeInfo, num_nodes: usize) -> Vec<NodeInfo> {
|
||||||
num_nodes,
|
num_nodes,
|
||||||
rpc_peers
|
rpc_peers
|
||||||
);
|
);
|
||||||
gossip_service.close().unwrap();
|
exit.store(true, Ordering::Relaxed);
|
||||||
|
gossip_service.join().unwrap();
|
||||||
return rpc_peers;
|
return rpc_peers;
|
||||||
}
|
}
|
||||||
debug!(
|
debug!(
|
||||||
|
@ -132,9 +126,11 @@ pub fn converge(node: &NodeInfo, num_nodes: usize) -> Vec<NodeInfo> {
|
||||||
panic!("Failed to converge");
|
panic!("Failed to converge");
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn make_spy_node(leader: &NodeInfo) -> (GossipService, Arc<RwLock<ClusterInfo>>, Pubkey) {
|
fn make_spy_node(
|
||||||
|
leader: &NodeInfo,
|
||||||
|
exit: &Arc<AtomicBool>,
|
||||||
|
) -> (GossipService, Arc<RwLock<ClusterInfo>>, Pubkey) {
|
||||||
let keypair = Keypair::new();
|
let keypair = Keypair::new();
|
||||||
let exit = Arc::new(AtomicBool::new(false));
|
|
||||||
let mut spy = Node::new_localhost_with_pubkey(keypair.pubkey());
|
let mut spy = Node::new_localhost_with_pubkey(keypair.pubkey());
|
||||||
let id = spy.info.id;
|
let id = spy.info.id;
|
||||||
let daddr = "0.0.0.0:0".parse().unwrap();
|
let daddr = "0.0.0.0:0".parse().unwrap();
|
||||||
|
@ -145,7 +141,7 @@ pub fn make_spy_node(leader: &NodeInfo) -> (GossipService, Arc<RwLock<ClusterInf
|
||||||
spy_cluster_info.set_leader(leader.id);
|
spy_cluster_info.set_leader(leader.id);
|
||||||
let spy_cluster_info_ref = Arc::new(RwLock::new(spy_cluster_info));
|
let spy_cluster_info_ref = Arc::new(RwLock::new(spy_cluster_info));
|
||||||
let gossip_service =
|
let gossip_service =
|
||||||
GossipService::new(&spy_cluster_info_ref, None, None, spy.sockets.gossip, &exit);
|
GossipService::new(&spy_cluster_info_ref, None, None, spy.sockets.gossip, exit);
|
||||||
|
|
||||||
(gossip_service, spy_cluster_info_ref, id)
|
(gossip_service, spy_cluster_info_ref, id)
|
||||||
}
|
}
|
||||||
|
@ -177,6 +173,7 @@ mod tests {
|
||||||
let cluster_info = ClusterInfo::new(tn.info.clone());
|
let cluster_info = ClusterInfo::new(tn.info.clone());
|
||||||
let c = Arc::new(RwLock::new(cluster_info));
|
let c = Arc::new(RwLock::new(cluster_info));
|
||||||
let d = GossipService::new(&c, None, None, tn.sockets.gossip, &exit);
|
let d = GossipService::new(&c, None, None, tn.sockets.gossip, &exit);
|
||||||
d.close().expect("thread join");
|
exit.store(true, Ordering::Relaxed);
|
||||||
|
d.join().unwrap();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
@ -9,6 +9,7 @@ use crate::gossip_service::GossipService;
|
||||||
use crate::packet::PACKET_DATA_SIZE;
|
use crate::packet::PACKET_DATA_SIZE;
|
||||||
use crate::result::{Error, Result};
|
use crate::result::{Error, Result};
|
||||||
use crate::rpc_request::{RpcClient, RpcRequest, RpcRequestHandler};
|
use crate::rpc_request::{RpcClient, RpcRequest, RpcRequestHandler};
|
||||||
|
use crate::service::Service;
|
||||||
use bincode::serialize_into;
|
use bincode::serialize_into;
|
||||||
use bs58;
|
use bs58;
|
||||||
use serde_json;
|
use serde_json;
|
||||||
|
@ -24,7 +25,7 @@ use solana_sdk::transaction::Transaction;
|
||||||
use std;
|
use std;
|
||||||
use std::io;
|
use std::io;
|
||||||
use std::net::{SocketAddr, UdpSocket};
|
use std::net::{SocketAddr, UdpSocket};
|
||||||
use std::sync::atomic::AtomicBool;
|
use std::sync::atomic::{AtomicBool, Ordering};
|
||||||
use std::sync::{Arc, RwLock};
|
use std::sync::{Arc, RwLock};
|
||||||
use std::thread::sleep;
|
use std::thread::sleep;
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
@ -432,7 +433,8 @@ pub fn poll_gossip_for_leader(leader_gossip: SocketAddr, timeout: Option<u64>) -
|
||||||
sleep(Duration::from_millis(100));
|
sleep(Duration::from_millis(100));
|
||||||
}
|
}
|
||||||
|
|
||||||
gossip_service.close()?;
|
exit.store(true, Ordering::Relaxed);
|
||||||
|
gossip_service.join()?;
|
||||||
|
|
||||||
if log_enabled!(log::Level::Trace) {
|
if log_enabled!(log::Level::Trace) {
|
||||||
trace!("{}", cluster_info.read().unwrap().node_info_trace());
|
trace!("{}", cluster_info.read().unwrap().node_info_trace());
|
||||||
|
|
Loading…
Reference in New Issue