mirror of
https://github.com/sigp/lighthouse.git
synced 2026-03-09 19:51:47 +00:00
Integrate tracing (#6339)
Tracing Integration
- [reference](5bbf1859e9/projects/project-ideas.md (L297))
- [x] replace slog & log with tracing throughout the codebase
- [x] implement custom crit log
- [x] make relevant changes in the formatter
- [x] replace sloggers
- [x] re-write SSE logging components
cc: @macladson @eserilev
This commit is contained in:
@@ -4,10 +4,11 @@ use lighthouse_network::Enr;
|
||||
use lighthouse_network::EnrExt;
|
||||
use lighthouse_network::Multiaddr;
|
||||
use lighthouse_network::{NetworkConfig, NetworkEvent};
|
||||
use slog::{debug, error, o, Drain};
|
||||
use std::sync::Arc;
|
||||
use std::sync::Weak;
|
||||
use tokio::runtime::Runtime;
|
||||
use tracing::{debug, error, info_span, Instrument};
|
||||
use tracing_subscriber::EnvFilter;
|
||||
use types::{
|
||||
ChainSpec, EnrForkId, Epoch, EthSpec, FixedBytesExtended, ForkContext, ForkName, Hash256,
|
||||
MinimalEthSpec, Slot,
|
||||
@@ -67,15 +68,12 @@ impl std::ops::DerefMut for Libp2pInstance {
|
||||
}
|
||||
|
||||
#[allow(unused)]
|
||||
pub fn build_log(level: slog::Level, enabled: bool) -> slog::Logger {
|
||||
let decorator = slog_term::TermDecorator::new().build();
|
||||
let drain = slog_term::FullFormat::new(decorator).build().fuse();
|
||||
let drain = slog_async::Async::new(drain).build().fuse();
|
||||
|
||||
pub fn build_tracing_subscriber(level: &str, enabled: bool) {
|
||||
if enabled {
|
||||
slog::Logger::root(drain.filter_level(level).fuse(), o!())
|
||||
} else {
|
||||
slog::Logger::root(drain.filter(|_| false).fuse(), o!())
|
||||
tracing_subscriber::fmt()
|
||||
.with_env_filter(EnvFilter::try_new(level).unwrap())
|
||||
.try_init()
|
||||
.unwrap();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -101,16 +99,16 @@ pub fn build_config(mut boot_nodes: Vec<Enr>) -> Arc<NetworkConfig> {
|
||||
pub async fn build_libp2p_instance(
|
||||
rt: Weak<Runtime>,
|
||||
boot_nodes: Vec<Enr>,
|
||||
log: slog::Logger,
|
||||
fork_name: ForkName,
|
||||
chain_spec: Arc<ChainSpec>,
|
||||
service_name: String,
|
||||
) -> Libp2pInstance {
|
||||
let config = build_config(boot_nodes);
|
||||
// launch libp2p service
|
||||
|
||||
let (signal, exit) = async_channel::bounded(1);
|
||||
let (shutdown_tx, _) = futures::channel::mpsc::channel(1);
|
||||
let executor = task_executor::TaskExecutor::new(rt, exit, log.clone(), shutdown_tx);
|
||||
let executor = task_executor::TaskExecutor::new(rt, exit, shutdown_tx, service_name);
|
||||
let libp2p_context = lighthouse_network::Context {
|
||||
config,
|
||||
enr_fork_id: EnrForkId::default(),
|
||||
@@ -119,7 +117,7 @@ pub async fn build_libp2p_instance(
|
||||
libp2p_registry: None,
|
||||
};
|
||||
Libp2pInstance(
|
||||
LibP2PService::new(executor, libp2p_context, &log)
|
||||
LibP2PService::new(executor, libp2p_context)
|
||||
.await
|
||||
.expect("should build libp2p instance")
|
||||
.0,
|
||||
@@ -143,18 +141,20 @@ pub enum Protocol {
|
||||
#[allow(dead_code)]
|
||||
pub async fn build_node_pair(
|
||||
rt: Weak<Runtime>,
|
||||
log: &slog::Logger,
|
||||
fork_name: ForkName,
|
||||
spec: Arc<ChainSpec>,
|
||||
protocol: Protocol,
|
||||
) -> (Libp2pInstance, Libp2pInstance) {
|
||||
let sender_log = log.new(o!("who" => "sender"));
|
||||
let receiver_log = log.new(o!("who" => "receiver"));
|
||||
|
||||
let mut sender =
|
||||
build_libp2p_instance(rt.clone(), vec![], sender_log, fork_name, spec.clone()).await;
|
||||
let mut sender = build_libp2p_instance(
|
||||
rt.clone(),
|
||||
vec![],
|
||||
fork_name,
|
||||
spec.clone(),
|
||||
"sender".to_string(),
|
||||
)
|
||||
.await;
|
||||
let mut receiver =
|
||||
build_libp2p_instance(rt, vec![], receiver_log, fork_name, spec.clone()).await;
|
||||
build_libp2p_instance(rt, vec![], fork_name, spec.clone(), "receiver".to_string()).await;
|
||||
|
||||
// let the two nodes set up listeners
|
||||
let sender_fut = async {
|
||||
@@ -179,7 +179,8 @@ pub async fn build_node_pair(
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
.instrument(info_span!("Sender", who = "sender"));
|
||||
let receiver_fut = async {
|
||||
loop {
|
||||
if let NetworkEvent::NewListenAddr(addr) = receiver.next_event().await {
|
||||
@@ -201,7 +202,8 @@ pub async fn build_node_pair(
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
.instrument(info_span!("Receiver", who = "receiver"));
|
||||
|
||||
let joined = futures::future::join(sender_fut, receiver_fut);
|
||||
|
||||
@@ -209,9 +211,9 @@ pub async fn build_node_pair(
|
||||
|
||||
match sender.testing_dial(receiver_multiaddr.clone()) {
|
||||
Ok(()) => {
|
||||
debug!(log, "Sender dialed receiver"; "address" => format!("{:?}", receiver_multiaddr))
|
||||
debug!(address = ?receiver_multiaddr, "Sender dialed receiver")
|
||||
}
|
||||
Err(_) => error!(log, "Dialing failed"),
|
||||
Err(_) => error!("Dialing failed"),
|
||||
};
|
||||
(sender, receiver)
|
||||
}
|
||||
@@ -220,7 +222,6 @@ pub async fn build_node_pair(
|
||||
#[allow(dead_code)]
|
||||
pub async fn build_linear(
|
||||
rt: Weak<Runtime>,
|
||||
log: slog::Logger,
|
||||
n: usize,
|
||||
fork_name: ForkName,
|
||||
spec: Arc<ChainSpec>,
|
||||
@@ -228,7 +229,14 @@ pub async fn build_linear(
|
||||
let mut nodes = Vec::with_capacity(n);
|
||||
for _ in 0..n {
|
||||
nodes.push(
|
||||
build_libp2p_instance(rt.clone(), vec![], log.clone(), fork_name, spec.clone()).await,
|
||||
build_libp2p_instance(
|
||||
rt.clone(),
|
||||
vec![],
|
||||
fork_name,
|
||||
spec.clone(),
|
||||
"linear".to_string(),
|
||||
)
|
||||
.await,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -238,8 +246,8 @@ pub async fn build_linear(
|
||||
.collect();
|
||||
for i in 0..n - 1 {
|
||||
match nodes[i].testing_dial(multiaddrs[i + 1].clone()) {
|
||||
Ok(()) => debug!(log, "Connected"),
|
||||
Err(_) => error!(log, "Failed to connect"),
|
||||
Ok(()) => debug!("Connected"),
|
||||
Err(_) => error!("Failed to connect"),
|
||||
};
|
||||
}
|
||||
nodes
|
||||
|
||||
Reference in New Issue
Block a user