From a957bf4aa2b08fcfecc99b72afe58d458775f5dd Mon Sep 17 00:00:00 2001 From: Benjamin Saunders Date: Fri, 25 Oct 2019 17:05:33 -0700 Subject: [PATCH] Replace slog with tracing --- interop/Cargo.toml | 4 +- interop/src/main.rs | 40 ++-- quinn-h3/Cargo.toml | 5 +- quinn-h3/examples/h3.rs | 30 +-- quinn-h3/examples/shared/mod.rs | 16 +- quinn-h3/examples/simple_client.rs | 12 +- quinn-h3/examples/simple_server.rs | 14 +- quinn-proto/Cargo.toml | 4 +- quinn-proto/src/connection.rs | 366 +++++++++++------------------ quinn-proto/src/endpoint.rs | 79 ++----- quinn-proto/src/frame.rs | 11 - quinn-proto/src/lib.rs | 36 --- quinn-proto/src/packet.rs | 34 --- quinn-proto/src/shared.rs | 24 +- quinn-proto/src/streams.rs | 9 +- quinn-proto/src/tests/mod.rs | 107 +++++---- quinn-proto/src/tests/util.rs | 110 +++------ quinn-proto/src/timer.rs | 11 - quinn-proto/src/transport_error.rs | 23 -- quinn/Cargo.toml | 5 +- quinn/benches/bench.rs | 64 +++-- quinn/examples/client.rs | 26 +- quinn/examples/server.rs | 66 +++--- quinn/src/builders.rs | 20 +- quinn/src/connection.rs | 23 +- quinn/src/endpoint.rs | 19 +- quinn/src/lib.rs | 3 - quinn/src/tests.rs | 107 +++++---- quinn/tests/many_connections.rs | 83 +------ 29 files changed, 475 insertions(+), 876 deletions(-) diff --git a/interop/Cargo.toml b/interop/Cargo.toml index ec941b090..8716118f4 100644 --- a/interop/Cargo.toml +++ b/interop/Cargo.toml @@ -10,16 +10,16 @@ quinn = { path = "../quinn" } quinn-h3 = { path = "../quinn-h3" } quinn-proto = { path = "../quinn-proto" } http = { git = "https://github.com/hyperium/http/", rev = "912534f1ef27d8a9050a4bd40d5ea0ee35136ea7" } -slog-term = "2" bytes = "0.4.7" structopt = "0.3.0" tokio = "0.2.0-alpha.5" tokio-net = "0.2.0-alpha.5" # tokio doesn't reexport everything we use rustls = { version = "0.16", features = ["dangerous_configuration"] } failure = "0.1" -slog = "2.2" futures = { package = "futures-preview", version = "0.3.0-alpha.18" } webpki = "0.21" +tracing = "0.1.10" +tracing-subscriber = "0.1.5" [[bin]] name = "main" diff --git a/interop/src/main.rs b/interop/src/main.rs index 8523ae343..5d4e90e10 100644 --- a/interop/src/main.rs +++ b/interop/src/main.rs @@ -1,14 +1,11 @@ use std::net::{SocketAddr, ToSocketAddrs}; use std::sync::{Arc, Mutex}; +use failure::{format_err, Error}; use futures::TryFutureExt; -// use quinn_h3::qpack; use structopt::StructOpt; use tokio::runtime::current_thread::Runtime; - -// use bytes::{Bytes, BytesMut}; -use failure::{format_err, Error}; -use slog::{info, o, warn, Drain, Logger}; +use tracing::{info, warn}; type Result = std::result::Result; @@ -29,14 +26,13 @@ struct Opt { fn main() { let opt = Opt::from_args(); let code = { - let decorator = slog_term::TermDecorator::new().stderr().build(); - let drain = slog_term::FullFormat::new(decorator) - .use_original_order() - .build() - .fuse(); - // We use a mutex-protected drain for simplicity; this tool is single-threaded anyway. - let drain = std::sync::Mutex::new(drain).fuse(); - if let Err(e) = run(Logger::root(drain, o!()), opt) { + tracing::subscriber::set_global_default( + tracing_subscriber::FmtSubscriber::builder() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .finish(), + ) + .unwrap(); + if let Err(e) = run(opt) { eprintln!("ERROR: {}", e); 1 } else { @@ -51,7 +47,6 @@ struct State { client_config: quinn::ClientConfig, remote: SocketAddr, host: String, - log: Logger, options: Opt, results: Arc>, } @@ -104,7 +99,7 @@ impl State { new_conn.connection } Err(conn) => { - info!(self.log, "0-RTT unsupported"); + info!("0-RTT unsupported"); let new_conn = conn .await .map_err(|e| format_err!("failed to connect: {}", e))?; @@ -168,9 +163,8 @@ impl State { } async fn rebind(self: Arc) -> Result<()> { - let mut builder = quinn::Endpoint::builder(); - builder.logger(self.log.clone()); - let (endpoint_driver, endpoint, _) = builder.bind(&"[::]:0".parse().unwrap())?; + let (endpoint_driver, endpoint, _) = + quinn::Endpoint::builder().bind(&"[::]:0".parse().unwrap())?; tokio::runtime::current_thread::spawn( endpoint_driver.unwrap_or_else(|e| eprintln!("IO error: {}", e)), ); @@ -228,7 +222,7 @@ struct Results { h3: bool, } -fn run(log: Logger, options: Opt) -> Result<()> { +fn run(options: Opt) -> Result<()> { let remote = format!("{}:{}", options.host, options.port) .to_socket_addrs()? .next() @@ -236,7 +230,7 @@ fn run(log: Logger, options: Opt) -> Result<()> { let host = if webpki::DNSNameRef::try_from_ascii_str(&options.host).is_ok() { &options.host } else { - warn!(log, "invalid hostname, using \"example.com\""); + warn!("invalid hostname, using \"example.com\""); "example.com" }; @@ -245,7 +239,6 @@ fn run(log: Logger, options: Opt) -> Result<()> { let results = Arc::new(Mutex::new(Results::default())); let protocols = vec![b"hq-23"[..].into(), quinn_h3::ALPN.into()]; - let mut builder = quinn::Endpoint::builder(); let mut tls_config = rustls::ClientConfig::new(); tls_config.versions = vec![rustls::ProtocolVersion::TLSv1_3]; tls_config.enable_early_data = true; @@ -265,8 +258,8 @@ fn run(log: Logger, options: Opt) -> Result<()> { ..Default::default() }; - builder.logger(log.clone()); - let (endpoint_driver, endpoint, _) = builder.bind(&"[::]:0".parse().unwrap())?; + let (endpoint_driver, endpoint, _) = + quinn::Endpoint::builder().bind(&"[::]:0".parse().unwrap())?; runtime.spawn(endpoint_driver.unwrap_or_else(|e| eprintln!("IO error: {}", e))); let state = Arc::new(State { @@ -274,7 +267,6 @@ fn run(log: Logger, options: Opt) -> Result<()> { client_config, remote, host: host.into(), - log, options, results, }); diff --git a/quinn-h3/Cargo.toml b/quinn-h3/Cargo.toml index 7e45a9a4d..1cf7fa8e6 100644 --- a/quinn-h3/Cargo.toml +++ b/quinn-h3/Cargo.toml @@ -31,8 +31,6 @@ http = { git = "https://github.com/hyperium/http/", rev = "912534f1ef27d8a9050a4 lazy_static = "1" quinn-proto = { path = "../quinn-proto", version = "0.4.0" } quinn = { path = "../quinn", version = "0.4.0" } -# slog = { version = "2.1", features = ["max_level_trace", "release_max_level_warn"]} -slog = { version = "2.1" } string = "0.2" tokio = "0.2.0-alpha.5" tokio-io = "0.2.0-alpha.5" @@ -45,9 +43,10 @@ failure = "0.1" proptest = "0.9.1" rand = "0.7.0" rcgen = "0.7" -slog-term = "2" structopt = "0.3.0" url = "2" +tracing = "0.1.10" +tracing-subscriber = "0.1.5" [[example]] name = "h3" diff --git a/quinn-h3/examples/h3.rs b/quinn-h3/examples/h3.rs index 15d88b2ab..7dcbbc1f3 100644 --- a/quinn-h3/examples/h3.rs +++ b/quinn-h3/examples/h3.rs @@ -7,7 +7,6 @@ use std::time::Instant; use failure::{format_err, Error}; use futures::{StreamExt, TryFutureExt}; use http::{header::HeaderValue, method::Method, HeaderMap, Request, Response, StatusCode}; -use slog::{o, Logger}; use structopt::{self, StructOpt}; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use url::Url; @@ -24,7 +23,7 @@ use quinn_h3::{ use quinn_proto::crypto::rustls::{Certificate, CertificateChain, PrivateKey}; mod shared; -use shared::{build_certs, logger}; +use shared::build_certs; #[derive(StructOpt, Debug, Clone)] #[structopt(name = "h3")] @@ -44,19 +43,19 @@ struct Opt { #[tokio::main] async fn main() -> Result<(), Box> { - let opt = Opt::from_args(); - let log = logger("h3".into()); - let certs = build_certs(log.clone(), &opt.key, &opt.cert).expect("failed to build certs"); - - let server = server( - log.new(o!("server" => "")), - opt.clone(), - (certs.0.clone(), certs.2.clone()), + tracing::subscriber::set_global_default( + tracing_subscriber::FmtSubscriber::builder() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .finish(), ) - .expect("init server failed"); + .unwrap(); + let opt = Opt::from_args(); + let certs = build_certs(&opt.key, &opt.cert).expect("failed to build certs"); - let (client, client_driver) = - build_client(log.new(o!("client" => "")), certs.1).expect("build client failed"); + let server = + server(opt.clone(), (certs.0.clone(), certs.2.clone())).expect("init server failed"); + + let (client, client_driver) = build_client(certs.1).expect("build client failed"); tokio::spawn(async move { println!("server running"); @@ -87,7 +86,6 @@ async fn main() -> Result<(), Box> { } fn server( - log: Logger, options: Opt, certs: (CertificateChain, PrivateKey), ) -> Result { @@ -103,7 +101,6 @@ fn server( server_config.certificate(certs.0, certs.1)?; let mut endpoint = quinn::Endpoint::builder(); - endpoint.logger(log.clone()); endpoint.listen(server_config.build()); let server = ServerBuilder::new(endpoint); @@ -195,11 +192,10 @@ async fn handle_request(request: Request, sender: Sender) -> Result<() Ok(()) } -fn build_client(log: Logger, cert: Certificate) -> Result<(Client, quinn::EndpointDriver), Error> { +fn build_client(cert: Certificate) -> Result<(Client, quinn::EndpointDriver), Error> { let mut endpoint = quinn::Endpoint::builder(); let mut client_config = quinn::ClientConfigBuilder::default(); client_config.protocols(&[quinn_h3::ALPN]); - endpoint.logger(log.clone()); client_config.add_certificate_authority(cert)?; endpoint.default_client_config(client_config.build()); diff --git a/quinn-h3/examples/shared/mod.rs b/quinn-h3/examples/shared/mod.rs index 90795e39a..31e51dcfe 100644 --- a/quinn-h3/examples/shared/mod.rs +++ b/quinn-h3/examples/shared/mod.rs @@ -1,9 +1,9 @@ use std::fmt; +use std::{fs, io, path::PathBuf}; use failure::{bail, Error, Fail, ResultExt}; use quinn_proto::crypto::rustls::{Certificate, CertificateChain, PrivateKey}; -use slog::{info, o, Drain, Logger}; -use std::{fs, io, path::PathBuf}; +use tracing::info; pub type Result = std::result::Result; @@ -32,7 +32,6 @@ impl ErrorExt for Error { } pub fn build_certs( - log: Logger, key: &Option, cert: &Option, ) -> Result<(CertificateChain, Certificate, PrivateKey)> { @@ -51,7 +50,7 @@ pub fn build_certs( let (cert, key) = match fs::read(&cert_path).and_then(|x| Ok((x, fs::read(&key_path)?))) { Ok(x) => x, Err(ref e) if e.kind() == io::ErrorKind::NotFound => { - info!(log, "generating self-signed certificate"); + info!("generating self-signed certificate"); let cert = rcgen::generate_simple_self_signed(vec!["localhost".into()]).unwrap(); let key = cert.serialize_private_key_der(); let cert = cert.serialize_der().unwrap(); @@ -73,12 +72,3 @@ pub fn build_certs( )) } } - -pub fn logger(label: String) -> Logger { - let decorator = slog_term::PlainSyncDecorator::new(std::io::stderr()); - let drain = slog_term::FullFormat::new(decorator) - .use_original_order() - .build() - .fuse(); - Logger::root(drain, o!("h3" => label)) -} diff --git a/quinn-h3/examples/simple_client.rs b/quinn-h3/examples/simple_client.rs index b751879ae..75c94d14b 100644 --- a/quinn-h3/examples/simple_client.rs +++ b/quinn-h3/examples/simple_client.rs @@ -9,7 +9,7 @@ use url::Url; use quinn_h3::{self, client::Builder as ClientBuilder, client::Client}; mod shared; -use shared::{build_certs, logger}; +use shared::build_certs; #[derive(StructOpt, Debug, Clone)] #[structopt(name = "h3_client")] @@ -29,9 +29,14 @@ const MAX_LEN: usize = 256 * 1024; #[tokio::main] async fn main() -> Result<(), Box> { + tracing::subscriber::set_global_default( + tracing_subscriber::FmtSubscriber::builder() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .finish(), + ) + .unwrap(); let opt = Opt::from_args(); - let log = logger("h3".into()); - let certs = build_certs(log.clone(), &opt.key, &opt.cert).expect("failed to build certs"); + let certs = build_certs(&opt.key, &opt.cert).expect("failed to build certs"); let remote = (opt.url.host_str().unwrap(), opt.url.port().unwrap_or(4433)) .to_socket_addrs() @@ -41,7 +46,6 @@ async fn main() -> Result<(), Box> { let mut endpoint = quinn::Endpoint::builder(); let mut client_config = quinn::ClientConfigBuilder::default(); - endpoint.logger(log.clone()); client_config.protocols(&[quinn_h3::ALPN]); client_config diff --git a/quinn-h3/examples/simple_server.rs b/quinn-h3/examples/simple_server.rs index a7308566d..2d94e9a17 100644 --- a/quinn-h3/examples/simple_server.rs +++ b/quinn-h3/examples/simple_server.rs @@ -1,10 +1,10 @@ use std::net::SocketAddr; use std::path::PathBuf; +use std::sync::Arc; use failure::{format_err, Error}; use futures::StreamExt; use http::{Request, Response, StatusCode}; -use std::sync::Arc; use structopt::{self, StructOpt}; use quinn::ConnectionDriver as QuicDriver; @@ -16,7 +16,7 @@ use quinn_h3::{ }; mod shared; -use shared::{build_certs, logger}; +use shared::build_certs; #[derive(StructOpt, Debug, Clone)] #[structopt(name = "h3_server")] @@ -35,9 +35,14 @@ struct Opt { #[tokio::main] async fn main() -> Result<(), Box> { + tracing::subscriber::set_global_default( + tracing_subscriber::FmtSubscriber::builder() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .finish(), + ) + .unwrap(); let opt = Opt::from_args(); - let log = logger("h3".into()); - let certs = build_certs(log.clone(), &opt.key, &opt.cert).expect("failed to build certs"); + let certs = build_certs(&opt.key, &opt.cert).expect("failed to build certs"); let server_config = quinn::ServerConfig { transport: Arc::new(quinn::TransportConfig { @@ -53,7 +58,6 @@ async fn main() -> Result<(), Box> { .expect("failed to add cert"); let mut endpoint = quinn::Endpoint::builder(); - endpoint.logger(log.clone()); endpoint.listen(server_config.build()); let server = ServerBuilder::new(endpoint); diff --git a/quinn-proto/Cargo.toml b/quinn-proto/Cargo.toml index 8c21fc647..072853810 100644 --- a/quinn-proto/Cargo.toml +++ b/quinn-proto/Cargo.toml @@ -29,11 +29,11 @@ rand = "0.7" ring = { version = "0.16.7", optional = true } rustls = { version = "0.16", features = ["quic"], optional = true } slab = "0.4" -slog = "2.2" +tracing = "0.1.10" webpki = { version = "0.21", optional = true } [dev-dependencies] assert_matches = "1.1" hex-literal = "0.2.0" rcgen = "0.7" -#slog = { version = "2.2", features = ["max_level_trace"] } # For debugging +tracing-subscriber = "0.1.5" diff --git a/quinn-proto/src/connection.rs b/quinn-proto/src/connection.rs index 2d78181bd..23f60eae2 100644 --- a/quinn-proto/src/connection.rs +++ b/quinn-proto/src/connection.rs @@ -8,7 +8,7 @@ use bytes::{Bytes, BytesMut}; use err_derive::Error; use fnv::FnvHashSet; use rand::{rngs::StdRng, Rng, SeedableRng}; -use slog::Logger; +use tracing::{debug, error, info, trace, trace_span, warn}; use crate::coding::BufMutExt; use crate::crypto::{self, HeaderKeys, Keys}; @@ -42,7 +42,6 @@ pub struct Connection where S: crypto::Session, { - log: Logger, endpoint_config: Arc, server_config: Option>>, config: Arc, @@ -163,7 +162,6 @@ where S: crypto::Session, { pub(crate) fn new( - log: Logger, endpoint_config: Arc, server_config: Option>>, config: Arc, @@ -193,7 +191,6 @@ where .as_ref() .map_or(false, |c| c.use_stateless_retry); let mut this = Self { - log, endpoint_config, server_config, tls, @@ -272,11 +269,6 @@ where this } - /// Replace the diagnostic logger - pub fn set_logger(&mut self, log: Logger) { - self.log = log; - } - /// Returns timer updates /// /// Connections should be polled for timer updates after: @@ -352,7 +344,6 @@ where space: SpaceId, ack: frame::Ack, ) -> Result<(), TransportError> { - trace!(self.log, "handling ack"; "ranges" => ?ack.iter().collect::>()); if ack.largest >= self.space(space).next_packet_number { return Err(TransportError::PROTOCOL_VIOLATION("unsent packet acked")); } @@ -434,7 +425,7 @@ where } } else { // We always start out sending ECN, so any ack that doesn't acknowledge it disables it. - debug!(self.log, "ECN not acknowledged by peer"); + debug!("ECN not acknowledged by peer"); self.path.sending_ecn = false; } } @@ -459,11 +450,7 @@ where ) { match self.space_mut(space).detect_ecn(newly_acked, ecn) { Err(e) => { - debug!( - self.log, - "halting ECN due to verification failure: {error}", - error = e - ); + debug!("halting ECN due to verification failure: {}", e); self.path.sending_ecn = false; // Wipe out the existing value because it might be garbage and could interfere with // future attempts to use ECN on new paths. @@ -501,7 +488,7 @@ where let ss = match self.streams.send_mut(id) { Some(ss) => ss, None => { - info!(self.log, "no send stream found for acked reset: {:?}", id); + info!("no send stream found for acked reset: {:?}", id); continue; } }; @@ -549,7 +536,7 @@ where self.endpoint_events.push_back(EndpointEventInner::Drained); } TimerKind::KeepAlive => { - trace!(self.log, "sending keep-alive"); + trace!("sending keep-alive"); self.ping(); } TimerKind::LossDetection => { @@ -560,7 +547,7 @@ where self.prev_crypto = None; } TimerKind::PathValidation => { - debug!(self.log, "path validation failed"); + debug!("path validation failed"); self.path_challenge = None; self.path_challenge_pending = false; if let Some(prev) = self.prev_path.take() { @@ -594,7 +581,11 @@ where return; } - trace!(self.log, "PTO fired"; "in flight" => self.in_flight.bytes, "count" => self.pto_count); + trace!( + in_flight = self.in_flight.bytes, + count = self.pto_count, + "PTO fired" + ); // Send two probes to improve odds of getting through under lossy conditions self.loss_probes = self.loss_probes.saturating_add(2); self.pto_count = self.pto_count.saturating_add(1); @@ -671,7 +662,7 @@ where let old_bytes_in_flight = self.in_flight.bytes; let largest_lost_sent = self.space(pn_space).sent_packets[&largest_lost].time_sent; self.lost_packets += lost_packets.len() as u64; - trace!(self.log, "packets lost: {:?}", lost_packets); + trace!("packets lost: {:?}", lost_packets); for packet in &lost_packets { let info = self .space_mut(pn_space) @@ -787,12 +778,7 @@ where Some(x) => x, None => return, }; - trace!( - self.log, - "{space:?} packet {packet} authenticated", - space = space_id, - packet = packet - ); + trace!("authenticated"); if self.side.is_server() { if self.spaces[SpaceId::Initial as usize].crypto.is_some() && space_id == SpaceId::Handshake @@ -928,6 +914,8 @@ where packet: Packet, remaining: Option, ) -> Result<(), TransportError> { + let span = trace_span!("recv initial"); + let _guard = span.enter(); debug_assert!(self.side.is_server()); let len = packet.header_data.len() + packet.payload.len(); self.total_recvd = len as u64; @@ -942,7 +930,7 @@ where ); self.process_early_payload(now, packet)?; if self.space(SpaceId::Initial).crypto_stream.offset() == 0 { - debug!(self.log, "dropping Initial with no CRYPTO data"); + debug!("dropping Initial with no CRYPTO data"); return Ok(()); } if self.state.is_closed() { @@ -994,15 +982,12 @@ where self.set_params(params); } Err(e) => { - error!( - self.log, - "session ticket has malformed transport parameters: {}", e - ); + error!("session ticket has malformed transport parameters: {}", e); return; } } } - trace!(self.log, "0-RTT enabled"); + trace!("0-RTT enabled"); self.zero_rtt_enabled = true; self.zero_rtt_crypto = Some(CryptoSpace { header: packet.header_keys(), @@ -1021,10 +1006,8 @@ where let end = crypto.offset + crypto.data.len() as u64; if space < expected && end > self.space(space).crypto_stream.offset() { warn!( - self.log, - "received new {actual:?} CRYPTO data when expecting {expected:?}", - actual = space, - expected = expected + "received new {:?} CRYPTO data when expecting {:?}", + space, expected ); return Err(TransportError::PROTOCOL_VIOLATION( "new data at unexpected encryption level", @@ -1045,7 +1028,7 @@ where if n == 0 { return Ok(()); } - trace!(self.log, "read {} TLS bytes", n); + trace!("read {} TLS bytes", n); self.tls.read_handshake(&buf[..n])?; } } @@ -1076,12 +1059,7 @@ where } } self.space_mut(space).crypto_offset += outgoing.len() as u64; - trace!( - self.log, - "wrote {} {space:?} TLS bytes", - outgoing.len(), - space = space - ); + trace!("wrote {} {:?} TLS bytes", outgoing.len(), space); self.space_mut(space) .pending .crypto @@ -1099,7 +1077,7 @@ where "already reached packet space {:?}", space ); - trace!(self.log, "{space:?} keys ready", space = space); + trace!("{:?} keys ready", space); self.spaces[space as usize].crypto = Some(CryptoSpace::new(crypto)); debug_assert!(space as usize > self.highest_space as usize); self.highest_space = space; @@ -1110,7 +1088,7 @@ where } fn discard_space(&mut self, space: SpaceId) { - trace!(self.log, "discarding {space:?} keys", space = space); + trace!("discarding {:?} keys", space); self.space_mut(space).crypto = None; let sent_packets = mem::replace(&mut self.space_mut(space).sent_packets, BTreeMap::new()); for (_, packet) in sent_packets.into_iter() { @@ -1139,11 +1117,7 @@ where if remote != self.path.remote && self.server_config.as_ref().map_or(true, |x| !x.migration) { - trace!( - self.log, - "discarding packet from unrecognized peer {address}", - address = format!("{}", remote) - ); + trace!("discarding packet from unrecognized peer {}", remote); return; } @@ -1179,7 +1153,7 @@ where self.handle_decode(now, remote, ecn, partial_decode); } Err(e) => { - trace!(self.log, "malformed header"; "reason" => %e); + trace!("malformed header: {}", e); return; } } @@ -1197,7 +1171,7 @@ where if let Some(ref crypto) = self.zero_rtt_crypto { Some(&crypto.header) } else { - debug!(self.log, "dropping unexpected 0-RTT packet"); + debug!("dropping unexpected 0-RTT packet"); return; } } else if let Some(space) = partial_decode.space() { @@ -1205,10 +1179,9 @@ where Some(&crypto.header) } else { debug!( - self.log, - "discarding unexpected {space:?} packet ({len} bytes)", - space = space, - len = partial_decode.len(), + "discarding unexpected {:?} packet ({} bytes)", + space, + partial_decode.len(), ); return; } @@ -1220,7 +1193,7 @@ where match partial_decode.finish(header_crypto) { Ok(packet) => self.handle_packet(now, remote, ecn, packet), Err(e) => { - trace!(self.log, "unable to complete packet decoding"; "reason" => %e); + trace!("unable to complete packet decoding: {}", e); } } } @@ -1233,12 +1206,11 @@ where mut packet: Packet, ) { trace!( - self.log, - "got {space:?} packet ({len} bytes) from {remote} using id {connection}", - space = packet.header.space(), - len = packet.payload.len() + packet.header_data.len(), - remote = remote, - connection = packet.header.dst_cid(), + "got {:?} packet ({} bytes) from {} using id {}", + packet.header.space(), + packet.payload.len() + packet.header_data.len(), + remote, + packet.header.dst_cid(), ); let was_closed = self.state.is_closed(); let was_drained = self.state.is_drained(); @@ -1249,19 +1221,24 @@ where let result = match self.decrypt_packet(now, &mut packet) { Err(Some(e)) => { - warn!(self.log, "got illegal packet"; "reason" => %e); + warn!("illegal packet: {}", e); Err(e.into()) } Err(None) => { if stateless_reset { - debug!(self.log, "got stateless reset"); + debug!("got stateless reset"); Err(ConnectionError::Reset) } else { - debug!(self.log, "failed to authenticate packet"); + debug!("failed to authenticate packet"); return; } } Ok(number) => { + let span = match number { + Some(pn) => trace_span!("recv", space = ?packet.header.space(), pn), + None => trace_span!("recv", space = ?packet.header.space()), + }; + let _guard = span.enter(); let duplicate = number.and_then(|n| { if self.space_mut(packet.header.space()).dedup.insert(n) { Some(n) @@ -1274,16 +1251,12 @@ where if stateless_reset { Err(ConnectionError::Reset) } else { - warn!( - self.log, - "discarding possible duplicate packet {packet}", - packet = number - ); + warn!("discarding possible duplicate packet {}", number); return; } } else if self.state.is_handshake() && packet.header.is_short() { // TODO: SHOULD buffer these to improve reordering tolerance. - trace!(self.log, "dropping short packet during handshake"); + trace!("dropping short packet during handshake"); return; } else { if !self.state.is_closed() { @@ -1316,11 +1289,7 @@ where unreachable!("timeouts aren't generated by packet processing"); } ConnectionError::TransportError(err) => { - debug!( - self.log, - "closing connection due to transport error: {error}", - error = &err - ); + debug!("closing connection due to transport error: {}", err); State::closed(err) } ConnectionError::VersionMismatch => State::Draining, @@ -1376,7 +1345,7 @@ where // match the Destination Connection ID from its Initial packet. return Ok(()); } - trace!(self.log, "retrying with CID {rem_cid}", rem_cid = rem_cid); + trace!("retrying with CID {}", rem_cid); let client_hello = state.client_hello.take().unwrap(); self.orig_rem_cid = Some(self.rem_cid); self.rem_cid = rem_cid; @@ -1425,7 +1394,10 @@ where .. } => { if rem_cid != self.rem_handshake_cid { - debug!(self.log, "discarding packet with mismatched remote CID: {expected} != {actual}", expected = self.rem_handshake_cid, actual = rem_cid); + debug!( + "discarding packet with mismatched remote CID: {} != {}", + self.rem_handshake_cid, rem_cid + ); return Ok(()); } self.remote_validated = true; @@ -1437,7 +1409,7 @@ where } if self.tls.is_handshaking() { - trace!(self.log, "handshake ongoing"); + trace!("handshake ongoing"); self.state = State::Handshake(state::Handshake { token: None, ..state @@ -1497,25 +1469,24 @@ where self.events.push_back(Event::Connected); self.state = State::Established; - trace!(self.log, "established"); + trace!("established"); Ok(()) } Header::Initial { src_cid: rem_cid, .. } => { if !state.rem_cid_set { - trace!( - self.log, - "switching remote CID to {rem_cid}", - rem_cid = rem_cid - ); + trace!("switching remote CID to {}", rem_cid); let mut state = state.clone(); self.rem_cid = rem_cid; self.rem_handshake_cid = rem_cid; state.rem_cid_set = true; self.state = State::Handshake(state); } else if rem_cid != self.rem_handshake_cid { - debug!(self.log, "discarding packet with mismatched remote CID: {expected} != {actual}", expected = self.rem_handshake_cid, actual = rem_cid); + debug!( + "discarding packet with mismatched remote CID: {} != {}", + self.rem_handshake_cid, rem_cid + ); return Ok(()); } self.process_early_payload(now, packet)?; @@ -1529,7 +1500,7 @@ where Ok(()) } Header::VersionNegotiate { .. } => { - debug!(self.log, "remote doesn't support our version"); + debug!("remote doesn't support our version"); Err(ConnectionError::VersionMismatch) } Header::Short { .. } => unreachable!( @@ -1550,7 +1521,7 @@ where for frame in frame::Iter::new(packet.payload.into()) { match frame { Frame::Close(_) => { - trace!(self.log, "draining"); + trace!("draining"); self.state = State::Draining; return Ok(()); } @@ -1571,12 +1542,11 @@ where ) -> Result<(), TransportError> { debug_assert_ne!(packet.header.space(), SpaceId::Data); for frame in frame::Iter::new(packet.payload.into()) { - match frame { - Frame::Padding => {} - _ => { - trace!(self.log, "got {type}", type=frame.ty()); - } - } + let span = match frame { + Frame::Padding => None, + _ => Some(trace_span!("frame", ty = %frame.ty())), + }; + let _guard = span.as_ref().map(|x| x.enter()); match frame { Frame::Ack(_) | Frame::Padding => {} _ => { @@ -1623,12 +1593,11 @@ where let is_0rtt = self.space(SpaceId::Data).crypto.is_none(); let mut is_probing_packet = true; for frame in frame::Iter::new(payload) { - match frame { - Frame::Padding => {} - _ => { - trace!(self.log, "got {type}", type=frame.ty()); - } - } + let span = match frame { + Frame::Padding => None, + _ => Some(trace_span!("frame", ty = %frame.ty())), + }; + let _guard = span.as_ref().map(|x| x.enter()); if is_0rtt { match frame { Frame::Crypto(_) | Frame::Close(_) => { @@ -1664,26 +1633,25 @@ where self.read_tls(SpaceId::Data, &frame)?; } Frame::Stream(frame) => { - trace!(self.log, "got stream"; "id" => frame.id.0, "offset" => frame.offset, "len" => frame.data.len(), "fin" => frame.fin); + trace!(id = %frame.id, offset = frame.offset, len = frame.data.len(), fin = frame.fin, "got stream"); let stream = frame.id; let rs = match self.streams.recv_stream(self.side, stream) { Err(e) => { - debug!(self.log, "received illegal stream frame"; "stream" => stream.0); + debug!("received illegal stream frame"); return Err(e); } Ok(None) => { - trace!(self.log, "dropping frame for closed stream"); + trace!("dropping frame for closed stream"); continue; } Ok(Some(rs)) => rs, }; if rs.is_finished() { - trace!(self.log, "dropping frame for finished stream"); + trace!("dropping frame for finished stream"); continue; } self.data_recvd += rs.ingest( - &self.log, frame, self.data_recvd, self.local_max_data, @@ -1722,7 +1690,7 @@ where if self.path_challenge != Some(token) || remote != self.path.remote { continue; } - trace!(self.log, "path validated"); + trace!("path validated"); self.io.timer_stop(TimerKind::PathValidation); self.path_challenge = None; } @@ -1737,30 +1705,21 @@ where } Frame::MaxStreamData { id, offset } => { if id.initiator() != self.side && id.dir() == Dir::Uni { - debug!( - self.log, - "got MAX_STREAM_DATA on recv-only {stream}", - stream = id - ); + debug!("got MAX_STREAM_DATA on recv-only {}", id); return Err(TransportError::STREAM_STATE_ERROR( "MAX_STREAM_DATA on recv-only stream", )); } if let Some(ss) = self.streams.send_mut(id) { if offset > ss.max_data { - trace!(self.log, "stream limit increased"; "stream" => id.0, - "old" => ss.max_data, "new" => offset, "current offset" => ss.offset); + trace!(stream = %id, old = ss.max_data, new = offset, current_offset = ss.offset, "stream limit increased"); if ss.offset == ss.max_data { self.events.push_back(Event::StreamWritable { stream: id }); } ss.max_data = offset; } } else if id.initiator() == self.side() && self.streams.is_local_unopened(id) { - debug!( - self.log, - "got MAX_STREAM_DATA on unopened {stream}", - stream = id - ); + debug!("got MAX_STREAM_DATA on unopened {}", id); return Err(TransportError::STREAM_STATE_ERROR( "MAX_STREAM_DATA on unopened stream", )); @@ -1786,11 +1745,11 @@ where }) => { let rs = match self.streams.recv_stream(self.side, id) { Err(e) => { - debug!(self.log, "received illegal RST_STREAM"); + debug!("received illegal RST_STREAM"); return Err(e); } Ok(None) => { - trace!(self.log, "received RST_STREAM on closed stream"); + trace!("received RST_STREAM on closed stream"); continue; } Ok(Some(stream)) => stream, @@ -1823,37 +1782,30 @@ where self.on_stream_frame(true, id); } Frame::DataBlocked { offset } => { - debug!(self.log, "peer claims to be blocked at connection level"; "offset" => offset); + debug!(offset, "peer claims to be blocked at connection level"); } Frame::StreamDataBlocked { id, offset } => { if id.initiator() == self.side && id.dir() == Dir::Uni { - debug!( - self.log, - "got STREAM_DATA_BLOCKED on send-only {stream}", - stream = id - ); + debug!("got STREAM_DATA_BLOCKED on send-only {}", id); return Err(TransportError::STREAM_STATE_ERROR( "STREAM_DATA_BLOCKED on send-only stream", )); } - debug!(self.log, "peer claims to be blocked at stream level"; "stream" => id, "offset" => offset); + debug!( + stream = %id, + offset, "peer claims to be blocked at stream level" + ); } Frame::StreamsBlocked { dir, limit } => { debug!( - self.log, - "peer claims to be blocked opening more than {limit} {dir} streams", - limit = limit, - dir = dir + "peer claims to be blocked opening more than {} {} streams", + limit, dir ); } Frame::StopSending(frame::StopSending { id, error_code }) => { if id.initiator() != self.side { if id.dir() == Dir::Uni { - debug!( - self.log, - "got STOP_SENDING on recv-only {stream}", - stream = id - ); + debug!("got STOP_SENDING on recv-only {}", id); return Err(TransportError::STREAM_STATE_ERROR( "STOP_SENDING on recv-only stream", )); @@ -1880,9 +1832,8 @@ where } if sequence > self.cids_issued { debug!( - self.log, - "got RETIRE_CONNECTION_ID for unissued cid sequence number {sequence}", - sequence = sequence, + sequence, + "got RETIRE_CONNECTION_ID for unissued sequence number" ); return Err(TransportError::PROTOCOL_VIOLATION( "RETIRE_CONNECTION_ID for unissued sequence number", @@ -1893,10 +1844,8 @@ where } Frame::NewConnectionId(frame) => { trace!( - self.log, - "NEW_CONNECTION_ID {sequence} = {id}", sequence = frame.sequence, - id = frame.id, + id = %frame.id, ); if self.rem_cid.is_empty() { return Err(TransportError::PROTOCOL_VIOLATION( @@ -1944,7 +1893,7 @@ where if self.side.is_server() { return Err(TransportError::PROTOCOL_VIOLATION("client sent NEW_TOKEN")); } - trace!(self.log, "got new token"); + trace!("got new token"); // TODO: Cache, or perhaps forward to user? } Frame::Datagram(datagram) => { @@ -1963,7 +1912,7 @@ where self.events.push_back(Event::DatagramReceived); } while datagram.data.len() + self.datagrams.recv_buffered > window { - debug!(self.log, "dropping stale datagram"); + debug!("dropping stale datagram"); self.recv_datagram(); } self.datagrams.recv_buffered += datagram.data.len(); @@ -2017,11 +1966,7 @@ where } fn migrate(&mut self, now: Instant, remote: SocketAddr) { - trace!( - self.log, - "migration initiated from {remote}", - remote = remote - ); + trace!(%remote, "migration initiated"); // Reset rtt/congestion state for new path unless it looks like a NAT rebinding. let maybe_rebinding = remote.is_ipv4() && remote.ip() == self.path.remote.ip(); // Note that the congestion window will not grow until validation terminates. Helps mitigate @@ -2065,12 +2010,7 @@ where } fn update_rem_cid(&mut self, new: IssuedCid) { - trace!( - self.log, - "switching to remote CID {sequence}: {connection_id}", - sequence = new.sequence, - connection_id = new.id - ); + trace!("switching to remote CID {}: {}", new.sequence, new.id); let retired = self.rem_cid_seq; self.space_mut(SpaceId::Data) .pending @@ -2097,7 +2037,7 @@ where && self.side.is_server() && self.total_recvd * 3 < self.total_sent + u64::from(self.mtu) { - trace!(self.log, "blocked by anti-amplification"); + trace!("blocked by anti-amplification"); return None; } @@ -2171,12 +2111,8 @@ where let space = &mut self.spaces[space_id as usize]; let exact_number = space.get_tx_number(); - trace!( - self.log, - "sending {space:?} packet {number}", - space = space_id, - number = exact_number - ); + let span = trace_span!("send", space = ?space_id, pn = exact_number); + let _guard = span.enter(); let number = PacketNumber::new(exact_number, space.largest_acked_packet.unwrap_or(0)); let header = match space_id { SpaceId::Data if space.crypto.is_some() => Header::Short { @@ -2217,7 +2153,7 @@ where coalesce = coalesce && !header.is_short(); let sent = if close { - trace!(self.log, "sending CONNECTION_CLOSE"); + trace!("sending CONNECTION_CLOSE"); let max_len = buf.capacity() - partial_encode.start - partial_encode.header_len @@ -2274,7 +2210,7 @@ where { let padding = padding_minus_one + 1; padded = true; - trace!(self.log, "PADDING * {count}", count = padding); + trace!("PADDING * {}", padding); buf.resize(buf.len() + padding, 0); } @@ -2320,7 +2256,7 @@ where return None; } - trace!(self.log, "{len} bytes", len = buf.len()); + trace!("sending {} byte datagram", buf.len()); self.total_sent = self.total_sent.wrapping_add(buf.len() as u64); Some(Transmit { @@ -2357,7 +2293,7 @@ where // PING if mem::replace(&mut self.ping_pending, false) { - trace!(self.log, "PING"); + trace!("PING"); buf.write(frame::Type::PING); } @@ -2365,7 +2301,7 @@ where // 0-RTT packets must never carry acks (which would have to be of handshake packets) let acks = if !space.pending_acks.is_empty() { debug_assert!(space.crypto.is_some(), "tried to send ACK in 0-RTT"); - trace!(self.log, "ACK"; "ranges" => ?space.pending_acks.iter().collect::>()); + trace!("ACK"); let ecn = if self.receiving_ecn { Some(&self.ecn_counters) } else { @@ -2383,7 +2319,7 @@ where if let Some(token) = self.path_challenge { // But only send a packet solely for that purpose at most once self.path_challenge_pending = false; - trace!(self.log, "PATH_CHALLENGE {token:08x}", token = token); + trace!("PATH_CHALLENGE {:08x}", token); buf.write(frame::Type::PATH_CHALLENGE); buf.write(token); } @@ -2392,11 +2328,7 @@ where // PATH_RESPONSE if buf.len() + 9 < max_size && space_id == SpaceId::Data { if let Some(response) = self.path_response.take() { - trace!( - self.log, - "PATH_RESPONSE {token:08x}", - token = response.token - ); + trace!("PATH_RESPONSE {:08x}", response.token); buf.write(frame::Type::PATH_RESPONSE); buf.write(response.token); } @@ -2418,10 +2350,9 @@ where data, }; trace!( - self.log, - "CRYPTO: off {offset} len {length}", - offset = truncated.offset, - length = truncated.data.len() + "CRYPTO: off {} len {}", + truncated.offset, + truncated.data.len() ); truncated.encode(buf); sent.crypto.push_back(truncated); @@ -2441,7 +2372,7 @@ where Some(x) => x, None => continue, }; - trace!(self.log, "RESET_STREAM"; "stream" => id.0); + trace!(stream = %id, "RESET_STREAM"); sent.rst_stream.push((id, error_code)); frame::ResetStream { id, @@ -2464,14 +2395,14 @@ where if stream.is_finished() { continue; } - trace!(self.log, "STOP_SENDING"; "stream" => frame.id); + trace!(stream = %frame.id, "STOP_SENDING"); frame.encode(buf); sent.stop_sending.push(frame); } // MAX_DATA if space.pending.max_data && buf.len() + 9 < max_size { - trace!(self.log, "MAX_DATA"; "value" => self.local_max_data); + trace!(value = self.local_max_data, "MAX_DATA"); space.pending.max_data = false; sent.max_data = true; buf.write(frame::Type::MAX_DATA); @@ -2494,12 +2425,7 @@ where } sent.max_stream_data.insert(id); let max = rs.bytes_read + self.config.stream_receive_window; - trace!( - self.log, - "MAX_STREAM_DATA: {stream} = {max}", - stream = id, - max = max - ); + trace!(stream = %id, max = max, "MAX_STREAM_DATA"); buf.write(frame::Type::MAX_STREAM_DATA); buf.write(id); buf.write_var(max); @@ -2509,7 +2435,10 @@ where if space.pending.max_uni_stream_id && buf.len() + 9 < max_size { space.pending.max_uni_stream_id = false; sent.max_uni_stream_id = true; - trace!(self.log, "MAX_STREAMS (unidirectional)"; "value" => self.streams.max_remote[Dir::Uni as usize]); + trace!( + value = self.streams.max_remote[Dir::Uni as usize], + "MAX_STREAMS (unidirectional)" + ); buf.write(frame::Type::MAX_STREAMS_UNI); buf.write_var(self.streams.max_remote[Dir::Uni as usize]); } @@ -2518,7 +2447,10 @@ where if space.pending.max_bi_stream_id && buf.len() + 9 < max_size { space.pending.max_bi_stream_id = false; sent.max_bi_stream_id = true; - trace!(self.log, "MAX_STREAMS (bidirectional)"; "value" => self.streams.max_remote[Dir::Bi as usize] - 1); + trace!( + value = self.streams.max_remote[Dir::Bi as usize], + "MAX_STREAMS (bidirectional)" + ); buf.write(frame::Type::MAX_STREAMS_BIDI); buf.write_var(self.streams.max_remote[Dir::Bi as usize]); } @@ -2530,10 +2462,9 @@ where None => break, }; trace!( - self.log, - "NEW_CONNECTION_ID {sequence} = {id}", sequence = issued.sequence, - id = issued.id, + id = %issued.id, + "NEW_CONNECTION_ID" ); frame::NewConnectionId { sequence: issued.sequence, @@ -2551,7 +2482,7 @@ where Some(x) => x, None => break, }; - trace!(self.log, "RETIRE_CONNECTION_ID {sequence}", sequence = seq); + trace!(sequence = seq, "RETIRE_CONNECTION_ID"); buf.write(frame::Type::RETIRE_CONNECTION_ID); buf.write_var(seq); sent.retire_cids.push(seq); @@ -2596,7 +2527,7 @@ where ); let data = stream.data.split_to(len); let fin = stream.fin && stream.data.is_empty(); - trace!(self.log, "STREAM"; "id" => stream.id.0, "off" => stream.offset, "len" => len, "fin" => fin); + trace!(id = %stream.id, off = stream.offset, len, fin, "STREAM"); let frame = frame::Stream { id: stream.id, offset: stream.offset, @@ -2631,7 +2562,7 @@ where } fn close_common(&mut self) { - trace!(self.log, "connection closed"); + trace!("connection closed"); for (_, timer) in &mut self.io.timers { *timer = Some(TimerSetting::Stop); } @@ -2645,10 +2576,8 @@ where fn validate_params(&mut self, params: &TransportParameters) -> Result<(), TransportError> { if self.side.is_client() && self.orig_rem_cid != params.original_connection_id { debug!( - self.log, - "original connection ID mismatch: expected {expected:x?}, actual {actual:x?}", - expected = self.orig_rem_cid, - actual = params.original_connection_id + "original connection ID mismatch: expected {:x?}, actual {:x?}", + self.orig_rem_cid, params.original_connection_id ); return Err(TransportError::TRANSPORT_PARAMETER_ERROR( "original CID mismatch", @@ -2856,11 +2785,7 @@ where crypto .decrypt(number, &packet.header_data, &mut packet.payload) .map_err(|()| { - trace!( - self.log, - "decryption failed with packet number {packet}", - packet = number - ); + trace!("decryption failed with packet number {}", number); None })?; @@ -2894,7 +2819,7 @@ where "illegal key update", ))); } - trace!(self.log, "key update authenticated"); + trace!("key update authenticated"); self.update_keys(crypto, Some((number, now)), true); self.set_key_discard_timer(now); } @@ -2915,23 +2840,15 @@ where pub fn write(&mut self, stream: StreamId, data: &[u8]) -> Result { assert!(stream.dir() == Dir::Bi || stream.initiator() == self.side); if self.state.is_closed() { - trace!(self.log, "write blocked; connection draining"; "stream" => stream.0); + trace!(%stream, "write blocked; connection draining"); return Err(WriteError::Blocked); } if self.blocked() { if self.congestion_blocked() { - trace!( - self.log, - "write on {stream} blocked by congestion", - stream = stream - ); + trace!(%stream, "write blocked by congestion"); } else { - trace!( - self.log, - "write on {stream} blocked by connection-level flow control", - stream = stream - ); + trace!(%stream, "write blocked by connection-level flow control"); } self.blocked_streams.insert(stream); return Err(WriteError::Blocked); @@ -2950,11 +2867,7 @@ where return Err(e); } Err(e @ WriteError::Blocked) => { - trace!( - self.log, - "write on {stream} blocked by flow control", - stream = stream - ); + trace!(%stream, "write blocked by flow control"); return Err(e); } Err(WriteError::UnknownStream) => unreachable!("not returned here"), @@ -2966,12 +2879,7 @@ where ); let n = conn_budget.min(stream_budget).min(data.len() as u64) as usize; self.queue_stream_data(stream, (&data[0..n]).into())?; - trace!( - self.log, - "wrote {len} bytes to {stream}", - len = n, - stream = stream - ); + trace!(%stream, "wrote {} bytes", n); Ok(n) } @@ -3149,7 +3057,7 @@ where /// Reset state to account for 0-RTT being ignored by the server fn reject_0rtt(&mut self) { debug_assert!(self.side.is_client()); - debug!(self.log, "0-RTT rejected"); + debug!("0-RTT rejected"); self.accepted_0rtt = false; self.streams.zero_rtt_rejected(self.side); // Discard already-queued frames diff --git a/quinn-proto/src/endpoint.rs b/quinn-proto/src/endpoint.rs index 7c263fd06..700aeb15c 100644 --- a/quinn-proto/src/endpoint.rs +++ b/quinn-proto/src/endpoint.rs @@ -11,7 +11,7 @@ use err_derive::Error; use fnv::FnvHashMap; use rand::{rngs::StdRng, Rng, RngCore, SeedableRng}; use slab::Slab; -use slog::{self, Logger}; +use tracing::{debug, trace, warn}; use crate::coding::BufMutExt; use crate::connection::{initial_close, Connection}; @@ -38,7 +38,6 @@ pub struct Endpoint where S: crypto::Session, { - log: Logger, rng: StdRng, transmits: VecDeque, connection_ids_initial: FnvHashMap, @@ -70,13 +69,11 @@ where /// /// Returns `Err` if the configuration is invalid. pub fn new( - log: Logger, config: Arc, server_config: Option>>, ) -> Result { config.validate()?; Ok(Self { - log, rng: StdRng::from_entropy(), transmits: VecDeque::new(), connection_ids_initial: FnvHashMap::default(), @@ -128,17 +125,12 @@ where self.connection_reset_tokens.remove(&old).unwrap(); } if self.connection_reset_tokens.insert(token, ch).is_some() { - warn!(self.log, "duplicate reset token"); + warn!("duplicate reset token"); } } RetireConnectionId(seq) => { if let Some(cid) = self.connections[ch].loc_cids.remove(&seq) { - trace!( - self.log, - "peer retired CID {sequence}: {cid}", - sequence = seq, - cid = cid, - ); + trace!("peer retired CID {}: {}", seq, cid); self.connection_ids.remove(&cid); return Some(self.send_new_identifiers(ch, 1)); } @@ -176,10 +168,10 @@ where destination, }) => { if !self.is_server() { - debug!(self.log, "dropping packet with unsupported version"); + debug!("dropping packet with unsupported version"); return None; } - trace!(self.log, "sending version negotiation"); + trace!("sending version negotiation"); // Negotiate versions let mut buf = Vec::::new(); Header::VersionNegotiate { @@ -198,7 +190,7 @@ where return None; } Err(e) => { - trace!(self.log, "malformed header"; "reason" => %e); + trace!("malformed header: {}", e); return None; } }; @@ -251,11 +243,7 @@ where // if !self.is_server() { - debug!( - self.log, - "got unexpected packet on unrecognized connection {connection}", - connection = dst_cid - ); + debug!("packet for unrecognized connection {}", dst_cid); self.stateless_reset(datagram_len, remote, &dst_cid); return None; } @@ -263,11 +251,7 @@ where if first_decode.has_long_header() { return if first_decode.is_initial() { if datagram_len < MIN_INITIAL_SIZE { - debug!( - self.log, - "ignoring short initial on {connection}", - connection = dst_cid - ); + debug!("ignoring short initial for connection {}", dst_cid); return None; } @@ -286,15 +270,14 @@ where ) .map(|(ch, conn)| (ch, DatagramEvent::NewConnection(conn))), Err(e) => { - trace!(self.log, "unable to decode packet"; "reason" => %e); + trace!("unable to decode packet: {}", e); None } } } else { debug!( - self.log, - "ignoring non-initial packet for unknown connection {connection}", - connection = dst_cid + "ignoring non-initial packet for unknown connection {}", + dst_cid ); None }; @@ -308,7 +291,7 @@ where if !dst_cid.is_empty() { self.stateless_reset(datagram_len, remote, &dst_cid); } else { - trace!(self.log, "dropping unrecognized short packet without ID"); + trace!("dropping unrecognized short packet without ID"); } None } @@ -327,17 +310,12 @@ where let max_padding_len = match inciting_dgram_len.checked_sub(RESET_TOKEN_SIZE) { Some(headroom) if headroom > MIN_PADDING_LEN => headroom - 1, _ => { - debug!(self.log, "ignoring unexpected {len} byte packet: not larger than minimum stateless reset size", len=inciting_dgram_len); + debug!("ignoring unexpected {} byte packet: not larger than minimum stateless reset size", inciting_dgram_len); return; } }; - debug!( - self.log, - "sending stateless reset for {connection} to {remote}", - connection = dst_cid, - remote = remote, - ); + debug!("sending stateless reset for {} to {}", dst_cid, remote); let mut buf = Vec::::new(); // Resets with at least this much padding can't possibly be distinguished from real packets const IDEAL_MIN_PADDING_LEN: usize = MIN_PADDING_LEN + MAX_CID_SIZE; @@ -368,9 +346,9 @@ where remote: SocketAddr, server_name: &str, ) -> Result<(ConnectionHandle, Connection), ConnectError> { - config.transport.validate(&self.log)?; + config.transport.validate()?; let remote_id = ConnectionId::random(&mut self.rng, MAX_CID_SIZE); - trace!(self.log, "initial dcid"; "value" => %remote_id); + trace!(initial_dcid = %remote_id); let (ch, conn) = self.add_connection( remote_id, remote_id, @@ -421,7 +399,7 @@ where now: Instant, ) -> Result<(ConnectionHandle, Connection), ConnectError> { let loc_cid = self.new_cid(); - let (server_config, tls, transport_config, log) = match opts { + let (server_config, tls, transport_config) = match opts { ConnectionOpts::Client { config, server_name, @@ -431,9 +409,6 @@ where None, config.crypto.start_session(&server_name, ¶ms)?, config.transport, - config - .log - .unwrap_or_else(|| self.log.new(o!("connection" => loc_cid))), ) } ConnectionOpts::Server { orig_dst_cid } => { @@ -448,13 +423,11 @@ where Some(config.clone()), config.crypto.start_session(&server_params), config.transport.clone(), - self.log.new(o!("connection" => loc_cid)), ) } }; let conn = Connection::new( - log, Arc::clone(&self.config), server_config, transport_config, @@ -511,7 +484,7 @@ where ) .is_err() { - debug!(self.log, "failed to authenticate initial packet"; "pn" => packet_number); + debug!(packet_number, "failed to authenticate initial packet"); return None; }; @@ -522,7 +495,7 @@ where if self.incoming_handshakes == server_config.accept_buffer as usize || self.reject_new_connections { - debug!(self.log, "rejecting connection due to full accept buffer"); + debug!("rejecting connection due to full accept buffer"); self.transmits.push_back(Transmit { destination: remote, ecn: None, @@ -542,9 +515,8 @@ where && (!server_config.use_stateless_retry || dst_cid.len() != self.config.local_cid_len) { debug!( - self.log, - "rejecting connection due to invalid DCID length {len}", - len = dst_cid.len() + "rejecting connection due to invalid DCID length {}", + dst_cid.len() ); self.transmits.push_back(Transmit { destination: remote, @@ -573,10 +545,10 @@ where if expires > SystemTime::now() { retry_cid = Some(token_dst_cid); } else { - trace!(self.log, "sending stateless retry due to expired token"); + trace!("sending stateless retry due to expired token"); } } else { - trace!(self.log, "sending stateless retry due to invalid token"); + trace!("sending stateless retry due to invalid token"); } if retry_cid.is_none() { let token = token::generate( @@ -624,12 +596,12 @@ where } match conn.handle_initial(now, remote, ecn, packet_number as u64, packet, rest) { Ok(()) => { - trace!(self.log, "connection incoming; ICID {icid}", icid = dst_cid); + trace!(id = ch.0, icid = %dst_cid, "connection incoming"); self.incoming_handshakes += 1; Some((ch, conn)) } Err(e) => { - debug!(self.log, "handshake failed"; "reason" => %e); + debug!("handshake failed: {}", e); self.handle_event(ch, EndpointEvent(EndpointEventInner::Drained)); self.transmits.push_back(Transmit { destination: remote, @@ -679,7 +651,6 @@ where { fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result { fmt.debug_struct("Endpoint") - .field("log", &self.log) .field("rng", &self.rng) .field("transmits", &self.transmits) .field("connection_ids_initial", &self.connection_ids_initial) diff --git a/quinn-proto/src/frame.rs b/quinn-proto/src/frame.rs index 0b84f01d0..0cc3f800a 100644 --- a/quinn-proto/src/frame.rs +++ b/quinn-proto/src/frame.rs @@ -40,17 +40,6 @@ impl coding::Codec for Type { } } -impl slog::Value for Type { - fn serialize( - &self, - _: &slog::Record<'_>, - key: slog::Key, - serializer: &mut dyn slog::Serializer, - ) -> slog::Result { - serializer.emit_arguments(key, &format_args!("{:?}", self)) - } -} - pub trait FrameStruct { /// Smallest number of bytes this type of frame is guaranteed to fit within. const SIZE_BOUND: usize; diff --git a/quinn-proto/src/lib.rs b/quinn-proto/src/lib.rs index 6107e49d7..20888025e 100644 --- a/quinn-proto/src/lib.rs +++ b/quinn-proto/src/lib.rs @@ -15,9 +15,6 @@ #![warn(missing_docs)] #![cfg_attr(test, allow(dead_code))] -#[macro_use] -extern crate slog; - use std::fmt; use std::net::SocketAddr; use std::ops; @@ -126,17 +123,6 @@ impl ops::Not for Side { } } -impl slog::Value for Side { - fn serialize( - &self, - _: &slog::Record<'_>, - key: slog::Key, - serializer: &mut dyn slog::Serializer, - ) -> slog::Result { - serializer.emit_arguments(key, &format_args!("{:?}", self)) - } -} - /// Whether a stream communicates data in both directions or only from the initiator #[derive(Debug, Copy, Clone, Eq, PartialEq, Ord, PartialOrd, Hash)] pub enum Dir { @@ -162,17 +148,6 @@ impl fmt::Display for Dir { } } -impl slog::Value for Dir { - fn serialize( - &self, - _: &slog::Record<'_>, - key: slog::Key, - serializer: &mut dyn slog::Serializer, - ) -> slog::Result { - serializer.emit_arguments(key, &format_args!("{:?}", self)) - } -} - /// Identifier for a stream within a particular connection #[derive(Debug, Copy, Clone, Eq, PartialEq, Ord, PartialOrd, Hash)] pub struct StreamId(#[doc(hidden)] pub u64); @@ -197,17 +172,6 @@ impl fmt::Display for StreamId { } } -impl slog::Value for StreamId { - fn serialize( - &self, - _: &slog::Record<'_>, - key: slog::Key, - serializer: &mut dyn slog::Serializer, - ) -> slog::Result { - serializer.emit_arguments(key, &format_args!("{:?}", self)) - } -} - impl StreamId { pub(crate) fn new(initiator: Side, dir: Dir, index: u64) -> Self { StreamId(index << 2 | (dir as u64) << 1 | initiator as u64) diff --git a/quinn-proto/src/packet.rs b/quinn-proto/src/packet.rs index 3b3fe647c..76049259c 100644 --- a/quinn-proto/src/packet.rs +++ b/quinn-proto/src/packet.rs @@ -2,7 +2,6 @@ use std::{cmp::Ordering, io, ops::Range, str}; use bytes::{BigEndian, Buf, BufMut, ByteOrder, Bytes, BytesMut}; use err_derive::Error; -use slog; use crate::coding::{self, BufExt, BufMutExt}; use crate::shared::ConnectionId; @@ -743,17 +742,6 @@ impl From for u8 { } } -impl slog::Value for LongHeaderType { - fn serialize( - &self, - _: &slog::Record<'_>, - key: slog::Key, - serializer: &mut dyn slog::Serializer, - ) -> slog::Result { - serializer.emit_arguments(key, &format_args!("{:?}", self)) - } -} - /// Long packet types with uniform header structure #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum LongType { @@ -761,17 +749,6 @@ pub enum LongType { ZeroRtt, } -impl slog::Value for LongType { - fn serialize( - &self, - _: &slog::Record<'_>, - key: slog::Key, - serializer: &mut dyn slog::Serializer, - ) -> slog::Result { - serializer.emit_arguments(key, &format_args!("{:?}", self)) - } -} - #[derive(Debug, Error, Clone, Eq, PartialEq, Ord, PartialOrd, Hash)] pub enum PacketDecodeError { #[error(display = "unsupported version")] @@ -814,17 +791,6 @@ impl SpaceId { } } -impl slog::Value for SpaceId { - fn serialize( - &self, - _: &slog::Record<'_>, - key: slog::Key, - serializer: &mut dyn slog::Serializer, - ) -> slog::Result { - serializer.emit_arguments(key, &format_args!("{:?}", self)) - } -} - #[cfg(test)] mod tests { use super::*; diff --git a/quinn-proto/src/shared.rs b/quinn-proto/src/shared.rs index aba6029a5..459de9fca 100644 --- a/quinn-proto/src/shared.rs +++ b/quinn-proto/src/shared.rs @@ -6,7 +6,7 @@ use std::{cmp, fmt}; use bytes::BytesMut; use err_derive::Error; use rand::{Rng, RngCore}; -use slog::Logger; +use tracing::warn; use crate::packet::PartialDecode; use crate::{crypto, VarInt, MAX_CID_SIZE, RESET_TOKEN_SIZE}; @@ -174,7 +174,7 @@ impl Default for TransportConfig { } impl TransportConfig { - pub(crate) fn validate(&self, log: &Logger) -> Result<(), ConfigError> { + pub(crate) fn validate(&self) -> Result<(), ConfigError> { if let Some((name, _)) = [ ("stream_window_bidi", self.stream_window_bidi), ("stream_window_uni", self.stream_window_uni), @@ -194,10 +194,8 @@ impl TransportConfig { } if self.idle_timeout != 0 && u64::from(self.keep_alive_interval) >= self.idle_timeout { warn!( - log, "keep-alive interval {} is ineffective due to lower idle timeout {}", - self.keep_alive_interval, - self.idle_timeout + self.keep_alive_interval, self.idle_timeout ); } Ok(()) @@ -360,11 +358,6 @@ pub struct ClientConfig { /// Cryptographic configuration to use pub crypto: C, - - /// Diagnostic logger - /// - /// If unset, the endpoint's logger is used. - pub log: Option, } /// Errors in the configuration of an endpoint @@ -491,17 +484,6 @@ impl fmt::Display for ConnectionId { } } -impl slog::Value for ConnectionId { - fn serialize( - &self, - _: &slog::Record<'_>, - key: slog::Key, - serializer: &mut dyn slog::Serializer, - ) -> slog::Result { - serializer.emit_arguments(key, &format_args!("{}", self)) - } -} - /// Explicit congestion notification codepoint #[repr(u8)] #[derive(Debug, Copy, Clone, Eq, PartialEq)] diff --git a/quinn-proto/src/streams.rs b/quinn-proto/src/streams.rs index 0f8766500..0a285b7c0 100644 --- a/quinn-proto/src/streams.rs +++ b/quinn-proto/src/streams.rs @@ -3,7 +3,7 @@ use std::collections::hash_map; use bytes::Bytes; use err_derive::Error; use fnv::FnvHashMap; -use slog::Logger; +use tracing::debug; use crate::assembler::Assembler; use crate::frame; @@ -309,7 +309,6 @@ impl Recv { pub fn ingest( &mut self, - log: &Logger, frame: frame::Stream, received: u64, max_data: u64, @@ -318,7 +317,7 @@ impl Recv { let end = frame.offset + frame.data.len() as u64; if let Some(final_offset) = self.final_offset() { if end > final_offset || (frame.fin && end != final_offset) { - debug!(log, "final offset error"; "frame end" => end, "final offset" => final_offset); + debug!(end, final_offset, "final offset error"); return Err(TransportError::FINAL_OFFSET_ERROR("")); } } @@ -327,9 +326,7 @@ impl Recv { let new_bytes = end.saturating_sub(prev_end); let stream_max_data = self.bytes_read + receive_window; if end > stream_max_data || received + new_bytes > max_data { - debug!(log, "flow control error"; - "stream" => frame.id.0, "recvd" => received, "new bytes" => new_bytes, - "max data" => max_data, "end" => end, "stream max data" => stream_max_data); + debug!(stream = %frame.id, received, new_bytes, max_data, end, stream_max_data, "flow control error"); return Err(TransportError::FLOW_CONTROL_ERROR("")); } diff --git a/quinn-proto/src/tests/mod.rs b/quinn-proto/src/tests/mod.rs index ff84f2456..99e701fd7 100644 --- a/quinn-proto/src/tests/mod.rs +++ b/quinn-proto/src/tests/mod.rs @@ -8,6 +8,7 @@ use bytes::Bytes; use hex_literal::hex; use rand::RngCore; use rustls::internal::msgs::enums::AlertDescription; +use tracing::info; use super::*; mod util; @@ -15,14 +16,9 @@ use util::*; #[test] fn version_negotiate_server() { - let log = logger(); + let _guard = subscribe(); let client_addr = "[::2]:7890".parse().unwrap(); - let mut server = Endpoint::new( - log.new(o!("peer" => "server")), - Default::default(), - Some(Arc::new(server_config())), - ) - .unwrap(); + let mut server = Endpoint::new(Default::default(), Some(Arc::new(server_config()))).unwrap(); let now = Instant::now(); let event = server.handle( now, @@ -47,10 +43,9 @@ fn version_negotiate_server() { #[test] fn version_negotiate_client() { - let log = logger(); + let _guard = subscribe(); let server_addr = "[::2]:7890".parse().unwrap(); let mut client = Endpoint::new( - log.new(o!("peer" => "client")), Arc::new(EndpointConfig { local_cid_len: 0, ..Default::default() @@ -86,6 +81,7 @@ fn version_negotiate_client() { #[test] fn lifecycle() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, server_ch) = pair.connect(); assert_matches!(pair.client_conn_mut(client_ch).poll(), None); @@ -93,7 +89,7 @@ fn lifecycle() { assert!(pair.server_conn_mut(server_ch).using_ecn()); const REASON: &[u8] = b"whee"; - info!(pair.log, "closing"); + info!("closing"); pair.client.connections.get_mut(&client_ch).unwrap().close( pair.time, VarInt(42), @@ -113,6 +109,7 @@ fn lifecycle() { #[test] fn stateless_retry() { + let _guard = subscribe(); let mut pair = Pair::new( Default::default(), ServerConfig { @@ -125,6 +122,7 @@ fn stateless_retry() { #[test] fn server_stateless_reset() { + let _guard = subscribe(); let mut reset_key = vec![0; 64]; let mut rng = rand::thread_rng(); rng.fill_bytes(&mut reset_key); @@ -136,19 +134,14 @@ fn server_stateless_reset() { let mut pair = Pair::new(endpoint_config.clone(), server_config()); let (client_ch, _) = pair.connect(); - pair.server.endpoint = Endpoint::new( - pair.log.new(o!("side" => "Server")), - endpoint_config, - Some(Arc::new(server_config())), - ) - .unwrap(); + pair.server.endpoint = Endpoint::new(endpoint_config, Some(Arc::new(server_config()))).unwrap(); // Send something big enough to allow room for a smaller stateless reset. pair.client.connections.get_mut(&client_ch).unwrap().close( pair.time, VarInt(42), (&[0xab; 128][..]).into(), ); - info!(pair.log, "resetting"); + info!("resetting"); pair.drive(); assert_matches!( pair.client_conn_mut(client_ch).poll(), @@ -160,6 +153,7 @@ fn server_stateless_reset() { #[test] fn client_stateless_reset() { + let _guard = subscribe(); let mut reset_key = vec![0; 64]; let mut rng = rand::thread_rng(); rng.fill_bytes(&mut reset_key); @@ -171,19 +165,14 @@ fn client_stateless_reset() { let mut pair = Pair::new(endpoint_config.clone(), server_config()); let (_, server_ch) = pair.connect(); - pair.client.endpoint = Endpoint::new( - pair.log.new(o!("side" => "Client")), - endpoint_config, - Some(Arc::new(server_config())), - ) - .unwrap(); + pair.client.endpoint = Endpoint::new(endpoint_config, Some(Arc::new(server_config()))).unwrap(); // Send something big enough to allow room for a smaller stateless reset. pair.server.connections.get_mut(&server_ch).unwrap().close( pair.time, VarInt(42), (&[0xab; 128][..]).into(), ); - info!(pair.log, "resetting"); + info!("resetting"); pair.drive(); assert_matches!( pair.server_conn_mut(server_ch).poll(), @@ -195,6 +184,7 @@ fn client_stateless_reset() { #[test] fn finish_stream() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, server_ch) = pair.connect(); @@ -225,6 +215,7 @@ fn finish_stream() { #[test] fn reset_stream() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, server_ch) = pair.connect(); @@ -234,7 +225,7 @@ fn reset_stream() { pair.client_conn_mut(client_ch).write(s, MSG).unwrap(); pair.drive(); - info!(pair.log, "resetting stream"); + info!("resetting stream"); const ERROR: VarInt = VarInt(42); pair.client_conn_mut(client_ch).reset(s, ERROR); pair.drive(); @@ -253,6 +244,7 @@ fn reset_stream() { #[test] fn stop_stream() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, server_ch) = pair.connect(); @@ -261,7 +253,7 @@ fn stop_stream() { pair.client_conn_mut(client_ch).write(s, MSG).unwrap(); pair.drive(); - info!(pair.log, "stopping stream"); + info!("stopping stream"); const ERROR: VarInt = VarInt(42); pair.server_conn_mut(server_ch) .stop_sending(s, ERROR) @@ -290,8 +282,9 @@ fn stop_stream() { #[test] fn reject_self_signed_cert() { + let _guard = subscribe(); let mut pair = Pair::default(); - info!(pair.log, "connecting"); + info!("connecting"); let client_ch = pair.begin_connect(ClientConfig::default()); pair.drive(); assert_matches!(pair.client_conn_mut(client_ch).poll(), @@ -301,6 +294,7 @@ fn reject_self_signed_cert() { #[test] fn congestion() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, _) = pair.connect(); @@ -329,6 +323,7 @@ fn congestion() { #[test] fn high_latency_handshake() { + let _guard = subscribe(); let mut pair = Pair::default(); pair.latency = Duration::from_micros(200 * 1000); let (client_ch, server_ch) = pair.connect(); @@ -340,6 +335,7 @@ fn high_latency_handshake() { #[test] fn zero_rtt_happypath() { + let _guard = subscribe(); let mut pair = Pair::new( Default::default(), ServerConfig { @@ -364,7 +360,7 @@ fn zero_rtt_happypath() { Ipv6Addr::LOCALHOST.into(), CLIENT_PORTS.lock().unwrap().next().unwrap(), ); - info!(pair.log, "resuming session"); + info!("resuming session"); let client_ch = pair.begin_connect(config.clone()); assert!(pair.client_conn_mut(client_ch).has_0rtt()); let s = pair.client_conn_mut(client_ch).open(Dir::Uni).unwrap(); @@ -382,6 +378,7 @@ fn zero_rtt_happypath() { #[test] fn zero_rtt_rejection() { + let _guard = subscribe(); let mut server_config = server_config(); Arc::get_mut(&mut server_config.crypto) .unwrap() @@ -419,7 +416,7 @@ fn zero_rtt_rejection() { Arc::get_mut(&mut client_config.crypto) .unwrap() .set_protocols(&["bar".into()]); - info!(pair.log, "resuming session"); + info!("resuming session"); let client_ch = pair.begin_connect(client_config); assert!(pair.client_conn_mut(client_ch).has_0rtt()); let s = pair.client_conn_mut(client_ch).open(Dir::Uni).unwrap(); @@ -444,6 +441,7 @@ fn zero_rtt_rejection() { #[test] fn alpn_success() { + let _guard = subscribe(); let mut server_config = server_config(); Arc::get_mut(&mut server_config.crypto) .unwrap() @@ -470,6 +468,7 @@ fn alpn_success() { #[test] fn stream_id_backpressure() { + let _guard = subscribe(); let server = ServerConfig { transport: Arc::new(TransportConfig { stream_window_uni: 1, @@ -536,6 +535,7 @@ fn stream_id_backpressure() { #[test] fn key_update() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, server_ch) = pair.connect(); let s = pair @@ -580,6 +580,7 @@ fn key_update() { #[test] fn key_update_reordered() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, server_ch) = pair.connect(); let s = pair @@ -592,16 +593,16 @@ fn key_update_reordered() { const MSG1: &[u8] = b"1"; pair.client_conn_mut(client_ch).write(s, MSG1).unwrap(); - pair.client.drive(&pair.log, pair.time, pair.server.addr); + pair.client.drive(pair.time, pair.server.addr); assert!(!pair.client.outbound.is_empty()); pair.client.delay_outbound(); pair.client_conn_mut(client_ch).initiate_key_update(); - info!(pair.log, "updated keys"); + info!("updated keys"); const MSG2: &[u8] = b"two"; pair.client_conn_mut(client_ch).write(s, MSG2).unwrap(); - pair.client.drive(&pair.log, pair.time, pair.server.addr); + pair.client.drive(pair.time, pair.server.addr); pair.client.finish_delay(); pair.drive(); @@ -623,9 +624,10 @@ fn key_update_reordered() { #[test] fn initial_retransmit() { + let _guard = subscribe(); let mut pair = Pair::default(); let client_ch = pair.begin_connect(client_config()); - pair.client.drive(&pair.log, pair.time, pair.server.addr); + pair.client.drive(pair.time, pair.server.addr); pair.client.outbound.clear(); // Drop initial pair.drive(); assert_matches!( @@ -636,8 +638,9 @@ fn initial_retransmit() { #[test] fn instant_close() { + let _guard = subscribe(); let mut pair = Pair::default(); - info!(pair.log, "connecting"); + info!("connecting"); let client_ch = pair.begin_connect(client_config()); pair.client .connections @@ -656,8 +659,9 @@ fn instant_close() { #[test] fn instant_close_2() { + let _guard = subscribe(); let mut pair = Pair::default(); - info!(pair.log, "connecting"); + info!("connecting"); let client_ch = pair.begin_connect(client_config()); // Unlike `instant_close`, the server sees a valid Initial packet first. pair.drive_client(); @@ -678,6 +682,7 @@ fn instant_close_2() { #[test] fn idle_timeout() { + let _guard = subscribe(); const IDLE_TIMEOUT: u64 = 10; let server = ServerConfig { transport: Arc::new(TransportConfig { @@ -719,6 +724,7 @@ fn idle_timeout() { #[test] fn server_busy() { + let _guard = subscribe(); let mut pair = Pair::new( Default::default(), ServerConfig { @@ -748,6 +754,7 @@ fn server_busy() { #[test] fn server_hs_retransmit() { + let _guard = subscribe(); let mut pair = Pair::default(); let client_ch = pair.begin_connect(client_config()); pair.step(); @@ -762,6 +769,7 @@ fn server_hs_retransmit() { #[test] fn migration() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, server_ch) = pair.connect(); pair.client.addr = SocketAddr::new( @@ -775,6 +783,7 @@ fn migration() { } fn test_flow_control(config: TransportConfig, window_size: usize) { + let _guard = subscribe(); let mut pair = Pair::new( Default::default(), ServerConfig { @@ -898,6 +907,7 @@ fn conn_flow_control() { #[test] fn stop_opens_bidi() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_conn, server_conn) = pair.connect(); let s = pair.client_conn_mut(client_conn).open(Dir::Bi).unwrap(); @@ -927,6 +937,7 @@ fn stop_opens_bidi() { #[test] fn implicit_open() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_conn, server_conn) = pair.connect(); let s1 = pair.client_conn_mut(client_conn).open(Dir::Uni).unwrap(); @@ -946,6 +957,7 @@ fn implicit_open() { #[test] fn zero_length_cid() { + let _guard = subscribe(); let mut pair = Pair::new( Arc::new(EndpointConfig { local_cid_len: 0, @@ -955,7 +967,7 @@ fn zero_length_cid() { ); let (client_ch, server_ch) = pair.connect(); // Ensure we can reconnect after a previous connection is cleaned up - info!(pair.log, "closing"); + info!("closing"); pair.client .connections .get_mut(&client_ch) @@ -972,6 +984,7 @@ fn zero_length_cid() { #[test] fn keep_alive() { + let _guard = subscribe(); const IDLE_TIMEOUT: u64 = 10_000; let server = ServerConfig { transport: Arc::new(TransportConfig { @@ -998,6 +1011,7 @@ fn keep_alive() { #[test] fn finish_stream_flow_control_reordered() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, server_ch) = pair.connect(); @@ -1006,19 +1020,19 @@ fn finish_stream_flow_control_reordered() { const MSG: &[u8] = b"hello"; pair.client_conn_mut(client_ch).write(s, MSG).unwrap(); pair.drive_client(); // Send stream data - pair.server.drive(&pair.log, pair.time, pair.client.addr); // Receive + pair.server.drive(pair.time, pair.client.addr); // Receive // Issue flow control credit assert_matches!( pair.server_conn_mut(server_ch).read_unordered(s), Ok(Some((ref data, 0))) if data == MSG ); - pair.server.drive(&pair.log, pair.time, pair.client.addr); + pair.server.drive(pair.time, pair.client.addr); pair.server.delay_outbound(); // Delay it pair.client_conn_mut(client_ch).finish(s).unwrap(); pair.drive_client(); // Send FIN - pair.server.drive(&pair.log, pair.time, pair.client.addr); // Acknowledge + pair.server.drive(pair.time, pair.client.addr); // Acknowledge pair.server.finish_delay(); // Add flow control packets after pair.drive(); @@ -1037,6 +1051,7 @@ fn finish_stream_flow_control_reordered() { #[test] fn handshake_1rtt_handling() { + let _guard = subscribe(); let mut pair = Pair::default(); let client_ch = pair.begin_connect(client_config()); pair.drive_client(); @@ -1044,7 +1059,7 @@ fn handshake_1rtt_handling() { let server_ch = pair.server.assert_accept(); // Server now has 1-RTT keys, but remains in Handshake state until the TLS CFIN has // authenticated the client. Delay the final client handshake flight so that doesn't happen yet. - pair.client.drive(&pair.log, pair.time, pair.server.addr); + pair.client.drive(pair.time, pair.server.addr); pair.client.delay_outbound(); // Send some 1-RTT data which will be received first. @@ -1052,7 +1067,7 @@ fn handshake_1rtt_handling() { const MSG: &[u8] = b"hello"; pair.client_conn_mut(client_ch).write(s, MSG).unwrap(); pair.client_conn_mut(client_ch).finish(s).unwrap(); - pair.client.drive(&pair.log, pair.time, pair.server.addr); + pair.client.drive(pair.time, pair.server.addr); // Add the handshake flight back on. pair.client.finish_delay(); @@ -1068,6 +1083,7 @@ fn handshake_1rtt_handling() { #[test] fn stop_before_finish() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, server_ch) = pair.connect(); @@ -1076,7 +1092,7 @@ fn stop_before_finish() { pair.client_conn_mut(client_ch).write(s, MSG).unwrap(); pair.drive(); - info!(pair.log, "stopping stream"); + info!("stopping stream"); const ERROR: VarInt = VarInt(42); pair.server_conn_mut(server_ch) .stop_sending(s, ERROR) @@ -1091,6 +1107,7 @@ fn stop_before_finish() { #[test] fn stop_during_finish() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, server_ch) = pair.connect(); @@ -1100,7 +1117,7 @@ fn stop_during_finish() { pair.drive(); assert_matches!(pair.server_conn_mut(server_ch).accept(Dir::Uni), Some(stream) if stream == s); - info!(pair.log, "stopping and finishing stream"); + info!("stopping and finishing stream"); const ERROR: VarInt = VarInt(42); pair.server_conn_mut(server_ch) .stop_sending(s, ERROR) @@ -1117,6 +1134,7 @@ fn stop_during_finish() { // Ensure we can recover from loss of tail packets when the congestion window is full #[test] fn congested_tail_loss() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, _) = pair.connect(); @@ -1147,6 +1165,7 @@ fn congested_tail_loss() { #[test] fn datagram_send_recv() { + let _guard = subscribe(); let mut pair = Pair::default(); let (client_ch, server_ch) = pair.connect(); assert_matches!(pair.server_conn_mut(server_ch).poll(), None); @@ -1172,6 +1191,7 @@ fn datagram_send_recv() { #[test] fn datagram_window() { + let _guard = subscribe(); const WINDOW: usize = 100; let server = ServerConfig { transport: Arc::new(TransportConfig { @@ -1236,6 +1256,7 @@ fn datagram_window() { #[test] fn datagram_unsupported() { + let _guard = subscribe(); let server = ServerConfig { transport: Arc::new(TransportConfig { datagram_receive_buffer_size: None, diff --git a/quinn-proto/src/tests/util.rs b/quinn-proto/src/tests/util.rs index 3fc4205f4..e9101d49e 100644 --- a/quinn-proto/src/tests/util.rs +++ b/quinn-proto/src/tests/util.rs @@ -4,19 +4,18 @@ use std::net::{Ipv6Addr, SocketAddr, UdpSocket}; use std::ops::RangeFrom; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; -use std::{cmp, env, fmt, mem, str}; +use std::{cmp, env, mem, str}; use fnv::FnvHashMap; use lazy_static::lazy_static; use rustls::KeyLogFile; -use slog::{Drain, Logger, KV}; +use tracing::{info_span, trace}; use super::*; use crate::crypto::rustls::{CertificateChain, PrivateKey}; use crate::timer::TimerKind; pub struct Pair { - pub log: Logger, pub server: TestEndpoint, pub client: TestEndpoint, pub time: Instant, @@ -29,14 +28,8 @@ pub struct Pair { impl Pair { pub fn new(endpoint_config: Arc, server_config: ServerConfig) -> Self { - let log = logger(); - let server = Endpoint::new( - log.new(o!("side" => "Server")), - endpoint_config.clone(), - Some(Arc::new(server_config)), - ) - .unwrap(); - let client = Endpoint::new(log.new(o!("side" => "Client")), endpoint_config, None).unwrap(); + let server = Endpoint::new(endpoint_config.clone(), Some(Arc::new(server_config))).unwrap(); + let client = Endpoint::new(endpoint_config, None).unwrap(); let server_addr = SocketAddr::new( Ipv6Addr::LOCALHOST.into(), @@ -47,9 +40,8 @@ impl Pair { CLIENT_PORTS.lock().unwrap().next().unwrap(), ); Self { - log, - server: TestEndpoint::new(Side::Server, server, server_addr), - client: TestEndpoint::new(Side::Client, client, client_addr), + server: TestEndpoint::new(server, server_addr), + client: TestEndpoint::new(client, client_addr), time: Instant::now(), latency: Duration::new(0, 0), spins: 0, @@ -71,14 +63,14 @@ impl Pair { Some(t) if Some(t) == client_t => { if t != self.time { self.time = self.time.max(t); - trace!(self.log, "advancing to {:?} for client", self.time); + trace!("advancing to {:?} for client", self.time); } true } Some(t) if Some(t) == server_t => { if t != self.time { self.time = self.time.max(t); - trace!(self.log, "advancing to {:?} for server", self.time); + trace!("advancing to {:?} for server", self.time); } true } @@ -93,8 +85,9 @@ impl Pair { } pub fn drive_client(&mut self) { - trace!(self.log, "client running"); - self.client.drive(&self.log, self.time, self.server.addr); + let span = info_span!("client"); + let _guard = span.enter(); + self.client.drive(self.time, self.server.addr); for x in self.client.outbound.drain(..) { if x.contents[0] & packet::LONG_HEADER_FORM == 0 { let spin = x.contents[0] & packet::SPIN_BIT != 0; @@ -113,8 +106,9 @@ impl Pair { } pub fn drive_server(&mut self) { - trace!(self.log, "server running"); - self.server.drive(&self.log, self.time, self.client.addr); + let span = info_span!("server"); + let _guard = span.enter(); + self.server.drive(self.time, self.client.addr); for x in self.server.outbound.drain(..) { if let Some(ref socket) = self.server.socket { socket.send_to(&x.contents, x.destination).unwrap(); @@ -128,7 +122,7 @@ impl Pair { } pub fn connect(&mut self) -> (ConnectionHandle, ConnectionHandle) { - info!(self.log, "connecting"); + info!("connecting"); let client_ch = self.begin_connect(client_config()); self.drive(); let server_ch = self.server.assert_accept(); @@ -145,6 +139,8 @@ impl Pair { /// Just start connecting the client pub fn begin_connect(&mut self, config: ClientConfig) -> ConnectionHandle { + let span = info_span!("client"); + let _guard = span.enter(); let (client_ch, client_conn) = self .client .connect(config, self.server.addr, "localhost") @@ -169,7 +165,6 @@ impl Default for Pair { } pub struct TestEndpoint { - side: Side, pub endpoint: Endpoint, pub addr: SocketAddr, socket: Option, @@ -183,7 +178,7 @@ pub struct TestEndpoint { } impl TestEndpoint { - fn new(side: Side, endpoint: Endpoint, addr: SocketAddr) -> Self { + fn new(endpoint: Endpoint, addr: SocketAddr) -> Self { let socket = if env::var_os("SSLKEYLOGFILE").is_some() { let socket = UdpSocket::bind(addr).expect("failed to bind UDP socket"); socket @@ -194,7 +189,6 @@ impl TestEndpoint { None }; Self { - side, endpoint, addr, socket, @@ -208,7 +202,7 @@ impl TestEndpoint { } } - pub fn drive(&mut self, log: &Logger, now: Instant, remote: SocketAddr) { + pub fn drive(&mut self, now: Instant, remote: SocketAddr) { if let Some(ref socket) = self.socket { loop { let mut buf = [0; 8192]; @@ -248,12 +242,7 @@ impl TestEndpoint { for (timer, setting) in &mut self.timers { if let Some(time) = *setting { if time <= now { - trace!( - log, - "{side:?} {timer:?} timeout", - side = self.side, - timer = timer - ); + trace!("{:?} timeout", timer); *setting = None; conn.handle_timeout(now, timer); } @@ -277,22 +266,11 @@ impl TestEndpoint { while let Some(x) = conn.poll_timers() { self.timers[x.timer] = match x.update { TimerSetting::Stop => { - trace!( - log, - "{side:?} {timer:?} stop", - side = self.side, - timer = x.timer - ); + trace!("{:?} stop", x.timer); None } TimerSetting::Start(time) => { - trace!( - log, - "{side:?} {timer:?} set to expire at {:?}", - time, - side = self.side, - timer = x.timer, - ); + trace!("{:?} set to expire at {:?}", x.timer, time); Some(time) } }; @@ -351,40 +329,26 @@ impl ::std::ops::DerefMut for TestEndpoint { } } -pub fn logger() -> Logger { - Logger::root(TestDrain.fuse(), o!()) +pub fn subscribe() -> tracing::subscriber::DefaultGuard { + let sub = tracing_subscriber::FmtSubscriber::builder() + .with_max_level(tracing::Level::TRACE) + .with_writer(|| TestWriter) + .finish(); + tracing::subscriber::set_default(sub) } -struct TestDrain; +struct TestWriter; -impl Drain for TestDrain { - type Ok = (); - type Err = io::Error; - fn log(&self, record: &slog::Record<'_>, values: &slog::OwnedKVList) -> Result<(), io::Error> { - let mut vals = Vec::new(); - values.serialize(&record, &mut TestSerializer(&mut vals))?; - record - .kv() - .serialize(&record, &mut TestSerializer(&mut vals))?; - println!( - "{} {}{}", - record.level(), - record.msg(), - str::from_utf8(&vals).unwrap() +impl Write for TestWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + print!( + "{}", + str::from_utf8(buf).expect("tried to log invalid UTF-8") ); - Ok(()) + Ok(buf.len()) } -} - -struct TestSerializer<'a, W>(&'a mut W); - -impl<'a, W> slog::Serializer for TestSerializer<'a, W> -where - W: Write + 'a, -{ - fn emit_arguments(&mut self, key: slog::Key, val: &fmt::Arguments<'_>) -> slog::Result { - write!(self.0, ", {}: {}", key, val).unwrap(); - Ok(()) + fn flush(&mut self) -> io::Result<()> { + io::stdout().flush() } } diff --git a/quinn-proto/src/timer.rs b/quinn-proto/src/timer.rs index 163015f35..3aa2d8168 100644 --- a/quinn-proto/src/timer.rs +++ b/quinn-proto/src/timer.rs @@ -136,14 +136,3 @@ impl IndexMut for TimerTable { &mut self.data[index.0 as usize] } } - -impl slog::Value for Timer { - fn serialize( - &self, - _: &slog::Record<'_>, - key: slog::Key, - serializer: &mut dyn slog::Serializer, - ) -> slog::Result { - serializer.emit_arguments(key, &format_args!("{:?}", self.0)) - } -} diff --git a/quinn-proto/src/transport_error.rs b/quinn-proto/src/transport_error.rs index 41b4b51a2..ea64e8645 100644 --- a/quinn-proto/src/transport_error.rs +++ b/quinn-proto/src/transport_error.rs @@ -1,7 +1,6 @@ use std::fmt; use bytes::{Buf, BufMut}; -use slog; use crate::coding::{self, BufExt, BufMutExt}; use crate::frame; @@ -42,17 +41,6 @@ impl From for Error { } } -impl slog::Value for Error { - fn serialize( - &self, - _: &slog::Record<'_>, - key: slog::Key, - serializer: &mut dyn slog::Serializer, - ) -> slog::Result { - serializer.emit_arguments(key, &format_args!("{:?}", self)) - } -} - /// Transport-level error code #[derive(Copy, Clone, Eq, PartialEq)] pub struct Code(u64); @@ -121,17 +109,6 @@ macro_rules! errors { } } -impl slog::Value for Code { - fn serialize( - &self, - _: &slog::Record<'_>, - key: slog::Key, - serializer: &mut dyn slog::Serializer, - ) -> slog::Result { - serializer.emit_arguments(key, &format_args!("{:?}", self)) - } -} - errors! { NO_ERROR(0x0) "the connection is being closed abruptly in the absence of any error"; INTERNAL_ERROR(0x1) "the endpoint encountered an internal error and cannot continue with the connection"; diff --git a/quinn/Cargo.toml b/quinn/Cargo.toml index 51dab9748..7af0cc388 100644 --- a/quinn/Cargo.toml +++ b/quinn/Cargo.toml @@ -29,7 +29,7 @@ libc = "0.2.49" mio = "0.6" proto = { package = "quinn-proto", path = "../quinn-proto", version = "0.4.0" } rustls = { version = "0.16", features = ["quic"] } -slog = "2.1" +tracing = "0.1.10" tokio-net = { version = "0.2.0-alpha.5", default-features = false } tokio-timer = "0.3.0-alpha.5" tokio-io = "0.2.0-alpha.5" @@ -43,7 +43,8 @@ directories = "2.0.0" failure = "0.1" rand = "0.7" rcgen = "0.7" -slog-term = "2" +tracing-subscriber = "0.1.5" +tracing-futures = { version = "0.1.0", default-features = false, features = ["std-future"] } structopt = "0.3.0" tokio = "0.2.0-alpha.5" unwrap = "1.2.1" diff --git a/quinn/benches/bench.rs b/quinn/benches/bench.rs index 9195e46ef..582cc7154 100644 --- a/quinn/benches/bench.rs +++ b/quinn/benches/bench.rs @@ -5,8 +5,9 @@ use std::thread; use bytes::Bytes; use criterion::{criterion_group, criterion_main, Criterion, Throughput}; use futures::StreamExt; -use slog::Drain; use tokio::runtime::current_thread::Runtime; +use tracing::info_span; +use tracing_futures::Instrument as _; use quinn::{ClientConfigBuilder, Endpoint, ServerConfigBuilder}; @@ -77,19 +78,10 @@ fn throughput(c: &mut Criterion) { struct Context { server_config: quinn::ServerConfig, client_config: quinn::ClientConfig, - log: slog::Logger, } impl Context { fn new() -> Self { - let decorator = slog_term::TermDecorator::new().stderr().build(); - let drain = slog_term::FullFormat::new(decorator) - .use_original_order() - .build() - .fuse(); - let drain = std::sync::Mutex::new(drain).fuse(); - let log = slog::Logger::root(drain, slog::o!()); - let cert = rcgen::generate_simple_self_signed(vec!["localhost".into()]).unwrap(); let key = quinn::PrivateKey::from_der(&cert.serialize_private_key_der()).unwrap(); let cert = quinn::Certificate::from_der(&cert.serialize_der().unwrap()).unwrap(); @@ -111,7 +103,6 @@ impl Context { Self { server_config: server_config.build(), client_config: client_config.build(), - log, } } @@ -119,45 +110,47 @@ impl Context { let sock = UdpSocket::bind(SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), 0)).unwrap(); let addr = sock.local_addr().unwrap(); let config = self.server_config.clone(); - let log = self.log.new(slog::o!("side" => "Server")); let handle = thread::spawn(move || { let mut endpoint = Endpoint::builder(); - endpoint.logger(log); endpoint.listen(config); let (driver, _, mut incoming) = endpoint.with_socket(sock).unwrap(); let mut runtime = Runtime::new().unwrap(); - runtime.spawn(async { driver.await.unwrap() }); - runtime.spawn(async move { - let quinn::NewConnection { - driver, - mut uni_streams, - .. - } = incoming.next().await.unwrap().await.unwrap(); - tokio::spawn(async move { - match driver.await { - Ok(()) => panic!("unexpected success"), - Err(quinn::ConnectionError::ApplicationClosed { .. }) => {} - Err(e) => panic!("{}", e), + runtime.spawn(async { driver.instrument(info_span!("server")).await.unwrap() }); + runtime.spawn( + async move { + let quinn::NewConnection { + driver, + mut uni_streams, + .. + } = incoming.next().await.unwrap().await.unwrap(); + tokio::spawn(async move { + match driver.await { + Ok(()) => panic!("unexpected success"), + Err(quinn::ConnectionError::ApplicationClosed { .. }) => {} + Err(e) => panic!("{}", e), + } + }); + while let Some(Ok(mut stream)) = uni_streams.next().await { + while let Some(_) = stream.read_unordered().await.unwrap() {} } - }); - while let Some(Ok(mut stream)) = uni_streams.next().await { - while let Some(_) = stream.read_unordered().await.unwrap() {} } - }); + .instrument(info_span!("server")), + ); runtime.run().unwrap(); }); (addr, handle) } pub fn make_client(&self, server_addr: SocketAddr) -> (quinn::Connection, Runtime) { - let mut endpoint = Endpoint::builder(); - endpoint.logger(self.log.new(slog::o!("side" => "Client"))); - let (endpoint_driver, endpoint, _) = endpoint + let (endpoint_driver, endpoint, _) = Endpoint::builder() .bind(&SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), 0)) .unwrap(); let mut runtime = Runtime::new().unwrap(); runtime.spawn(async move { - endpoint_driver.await.unwrap(); + endpoint_driver + .instrument(info_span!("client")) + .await + .unwrap(); }); let quinn::NewConnection { driver, connection, .. @@ -165,11 +158,12 @@ impl Context { .block_on( endpoint .connect_with(self.client_config.clone(), &server_addr, "localhost") - .unwrap(), + .unwrap() + .instrument(info_span!("client")), ) .unwrap(); runtime.spawn(async move { - driver.await.unwrap(); + driver.instrument(info_span!("client")).await.unwrap(); }); (connection, runtime) } diff --git a/quinn/examples/client.rs b/quinn/examples/client.rs index d66c73075..11b21f398 100644 --- a/quinn/examples/client.rs +++ b/quinn/examples/client.rs @@ -1,7 +1,5 @@ #[macro_use] extern crate failure; -#[macro_use] -extern crate slog; use std::fs; use std::io::{self, Write}; @@ -11,9 +9,9 @@ use std::time::{Duration, Instant}; use failure::Error; use futures::TryFutureExt; -use slog::{Drain, Logger}; use structopt::StructOpt; use tokio::runtime::current_thread::Runtime; +use tracing::{error, info}; use url::Url; mod common; @@ -44,16 +42,15 @@ struct Opt { } fn main() { + tracing::subscriber::set_global_default( + tracing_subscriber::FmtSubscriber::builder() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .finish(), + ) + .unwrap(); let opt = Opt::from_args(); let code = { - let decorator = slog_term::TermDecorator::new().stderr().build(); - let drain = slog_term::FullFormat::new(decorator) - .use_original_order() - .build() - .fuse(); - // We use a mutex-protected drain for simplicity; this example is single-threaded anyway. - let drain = std::sync::Mutex::new(drain).fuse(); - if let Err(e) = run(Logger::root(drain, o!()), opt) { + if let Err(e) = run(opt) { eprintln!("ERROR: {}", e); 1 } else { @@ -63,7 +60,7 @@ fn main() { ::std::process::exit(code); } -fn run(log: Logger, options: Opt) -> Result<()> { +fn run(options: Opt) -> Result<()> { let url = options.url; let remote = (url.host_str().unwrap(), url.port().unwrap_or(4433)) .to_socket_addrs()? @@ -73,7 +70,6 @@ fn run(log: Logger, options: Opt) -> Result<()> { let mut endpoint = quinn::Endpoint::builder(); let mut client_config = quinn::ClientConfigBuilder::default(); client_config.protocols(common::ALPN_QUIC_HTTP); - endpoint.logger(log.clone()); if options.keylog { client_config.enable_keylog(); } @@ -87,10 +83,10 @@ fn run(log: Logger, options: Opt) -> Result<()> { client_config.add_certificate_authority(quinn::Certificate::from_der(&cert)?)?; } Err(ref e) if e.kind() == io::ErrorKind::NotFound => { - info!(log, "local server certificate not found"); + info!("local server certificate not found"); } Err(e) => { - error!(log, "failed to open local server certificate: {}", e); + error!("failed to open local server certificate: {}", e); } } } diff --git a/quinn/examples/server.rs b/quinn/examples/server.rs index 2f7428c1e..817953a98 100644 --- a/quinn/examples/server.rs +++ b/quinn/examples/server.rs @@ -1,7 +1,5 @@ #[macro_use] extern crate failure; -#[macro_use] -extern crate slog; use std::net::SocketAddr; use std::path::{self, Path, PathBuf}; @@ -10,9 +8,10 @@ use std::{ascii, fs, io, str}; use failure::{Error, ResultExt}; use futures::{StreamExt, TryFutureExt}; -use slog::{Drain, Logger}; use structopt::{self, StructOpt}; use tokio::runtime::Runtime; +use tracing::{error, info, info_span}; +use tracing_futures::Instrument as _; mod common; @@ -42,15 +41,15 @@ struct Opt { } fn main() { + tracing::subscriber::set_global_default( + tracing_subscriber::FmtSubscriber::builder() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .finish(), + ) + .unwrap(); let opt = Opt::from_args(); let code = { - let decorator = slog_term::TermDecorator::new().stderr().build(); - let drain = slog_term::FullFormat::new(decorator) - .use_original_order() - .build(); - // We use a mutex-protected drain for simplicity; this example is single-threaded anyway. - let drain = std::sync::Mutex::new(drain).fuse(); - if let Err(e) = run(Logger::root(drain, o!()), opt) { + if let Err(e) = run(opt) { eprintln!("ERROR: {}", e); 1 } else { @@ -60,7 +59,7 @@ fn main() { ::std::process::exit(code); } -fn run(log: Logger, options: Opt) -> Result<()> { +fn run(options: Opt) -> Result<()> { let server_config = quinn::ServerConfig { transport: Arc::new(quinn::TransportConfig { stream_window_uni: 0, @@ -101,7 +100,7 @@ fn run(log: Logger, options: Opt) -> Result<()> { let (cert, key) = match fs::read(&cert_path).and_then(|x| Ok((x, fs::read(&key_path)?))) { Ok(x) => x, Err(ref e) if e.kind() == io::ErrorKind::NotFound => { - info!(log, "generating self-signed certificate"); + info!("generating self-signed certificate"); let cert = rcgen::generate_simple_self_signed(vec!["localhost".into()]).unwrap(); let key = cert.serialize_private_key_der(); let cert = cert.serialize_der().unwrap(); @@ -120,7 +119,6 @@ fn run(log: Logger, options: Opt) -> Result<()> { } let mut endpoint = quinn::Endpoint::builder(); - endpoint.logger(log.clone()); endpoint.listen(server_config.build()); let root = Arc::::from(options.root); @@ -130,18 +128,17 @@ fn run(log: Logger, options: Opt) -> Result<()> { let (endpoint_driver, mut incoming) = { let (driver, endpoint, incoming) = endpoint.bind(&options.listen)?; - info!(log, "listening on {}", endpoint.local_addr()?); + info!("listening on {}", endpoint.local_addr()?); (driver, incoming) }; let runtime = Runtime::new()?; runtime.spawn(async move { while let Some(conn) = incoming.next().await { - info!(log, "connection incoming"); - let log2 = log.clone(); + info!("connection incoming"); tokio::spawn( - handle_connection(root.clone(), log.clone(), conn).unwrap_or_else(move |e| { - error!(log2, "connection failed: {reason}", reason = e.to_string()) + handle_connection(root.clone(), conn).unwrap_or_else(move |e| { + error!("connection failed: {reason}", reason = e.to_string()) }), ); } @@ -151,17 +148,20 @@ fn run(log: Logger, options: Opt) -> Result<()> { Ok(()) } -async fn handle_connection(root: Arc, log: Logger, conn: quinn::Connecting) -> Result<()> { +async fn handle_connection(root: Arc, conn: quinn::Connecting) -> Result<()> { let quinn::NewConnection { driver, connection, mut bi_streams, .. } = conn.await?; - info!(log, "connection established"; - "remote_id" => %connection.remote_id(), - "address" => %connection.remote_address(), - "protocol" => connection.protocol().map_or_else(|| "".into(), |x| String::from_utf8_lossy(&x).into_owned())); + let span = info_span!( + "connection", + remote = %connection.remote_address(), + protocol = %connection.protocol().map_or_else(|| "".into(), |x| String::from_utf8_lossy(&x).into_owned()) + ); + let _guard = span.enter(); + info!("established"); // We ignore errors from the driver because they'll be reported by the `streams` handler anyway. tokio::spawn(driver.unwrap_or_else(|_| ())); @@ -170,7 +170,7 @@ async fn handle_connection(root: Arc, log: Logger, conn: quinn::Connecting while let Some(stream) = bi_streams.next().await { let stream = match stream { Err(quinn::ConnectionError::ApplicationClosed { .. }) => { - info!(log, "connection closed"); + info!("connection closed"); return Ok(()); } Err(e) => { @@ -178,11 +178,10 @@ async fn handle_connection(root: Arc, log: Logger, conn: quinn::Connecting } Ok(s) => s, }; - let log2 = log.clone(); tokio::spawn( - handle_request(root.clone(), log.clone(), stream).unwrap_or_else(move |e| { - error!(log2, "request failed: {reason}", reason = e.to_string()) - }), + handle_request(root.clone(), stream) + .unwrap_or_else(move |e| error!("failed: {reason}", reason = e.to_string())) + .instrument(info_span!("request")), ); } Ok(()) @@ -190,7 +189,6 @@ async fn handle_connection(root: Arc, log: Logger, conn: quinn::Connecting async fn handle_request( root: Arc, - log: Logger, (mut send, recv): (quinn::SendStream, quinn::RecvStream), ) -> Result<()> { let req = recv @@ -202,14 +200,10 @@ async fn handle_request( let part = ascii::escape_default(x).collect::>(); escaped.push_str(str::from_utf8(&part).unwrap()); } - info!(log, "got request"; "content" => escaped); + info!(content = %escaped); // Execute the request let resp = process_get(&root, &req).unwrap_or_else(|e| { - error!( - log, - "failed to process request: {reason}", - reason = e.to_string() - ); + error!("failed: {}", e); format!("failed to process request: {}\n", e) .into_bytes() .into() @@ -222,7 +216,7 @@ async fn handle_request( send.finish() .await .map_err(|e| format_err!("failed to shutdown stream: {}", e))?; - info!(log, "request complete"); + info!("complete"); Ok(()) } diff --git a/quinn/src/builders.rs b/quinn/src/builders.rs index ed99e2e5f..d21077dd3 100644 --- a/quinn/src/builders.rs +++ b/quinn/src/builders.rs @@ -7,7 +7,6 @@ use err_derive::Error; use proto::crypto::rustls::{Certificate, CertificateChain, PrivateKey}; use proto::{ClientConfig, EndpointConfig, ServerConfig}; use rustls::TLSError; -use slog::Logger; use crate::endpoint::{Endpoint, EndpointDriver, EndpointRef, Incoming}; use crate::udp::UdpSocket; @@ -16,7 +15,6 @@ use crate::udp::UdpSocket; #[derive(Clone, Debug)] pub struct EndpointBuilder { reactor: Option, - logger: Logger, server_config: Option, config: EndpointConfig, client_config: ClientConfig, @@ -54,13 +52,8 @@ impl EndpointBuilder { let addr = socket.local_addr().map_err(EndpointError::Socket)?; let socket = UdpSocket::from_std(socket, &reactor).map_err(EndpointError::Socket)?; let rc = EndpointRef::new( - self.logger.clone(), socket, - proto::Endpoint::new( - self.logger, - Arc::new(self.config), - self.server_config.map(Arc::new), - )?, + proto::Endpoint::new(Arc::new(self.config), self.server_config.map(Arc::new))?, addr.is_ipv6(), ); Ok(( @@ -83,10 +76,6 @@ impl EndpointBuilder { self.reactor = Some(handle); self } - pub fn logger(&mut self, logger: Logger) -> &mut Self { - self.logger = logger; - self - } /// Set the default configuration used for outgoing connections. /// @@ -101,7 +90,6 @@ impl Default for EndpointBuilder { fn default() -> Self { Self { reactor: None, - logger: Logger::root(slog::Discard, o!()), server_config: None, config: EndpointConfig::default(), client_config: ClientConfig::default(), @@ -226,12 +214,6 @@ impl ClientConfigBuilder { self } - /// Set the logger to use - pub fn logger(&mut self, logger: Logger) -> &mut Self { - self.config.log = Some(logger); - self - } - /// Begin connecting from `endpoint` to `addr`. pub fn build(self) -> ClientConfig { self.config diff --git a/quinn/src/connection.rs b/quinn/src/connection.rs index 8e7dc4695..61c8aea70 100644 --- a/quinn/src/connection.rs +++ b/quinn/src/connection.rs @@ -12,8 +12,8 @@ use futures::channel::{mpsc, oneshot}; use futures::task::{Context, Waker}; use futures::{Future, FutureExt, Poll, StreamExt}; use proto::{ConnectionError, ConnectionHandle, ConnectionId, Dir, StreamId, TimerUpdate}; -use slog::Logger; use tokio_timer::{delay, Delay}; +use tracing::{info_span, trace}; use crate::broadcast::{self, Broadcast}; use crate::streams::{RecvStream, SendStream, WriteError}; @@ -165,6 +165,9 @@ impl Future for ConnectionDriver { fn poll(mut self: Pin<&mut Self>, cx: &mut Context) -> Poll { let conn = &mut *self.0.lock().unwrap(); + let span = info_span!("drive", id = conn.handle.0); + let _guard = span.enter(); + loop { let now = Instant::now(); let mut keep_going = false; @@ -305,13 +308,6 @@ impl Connection { pub fn force_key_update(&self) { self.0.lock().unwrap().inner.initiate_key_update() } - - /// Replace the diagnostic logger - pub fn set_logger(&self, log: Logger) { - let mut conn = self.0.lock().unwrap(); - conn.log = log.clone(); - conn.inner.set_logger(log); - } } /// A stream of unidirectional QUIC streams initiated by a remote peer. @@ -522,14 +518,12 @@ pub struct ConnectionRef(Arc>); impl ConnectionRef { pub(crate) fn new( - log: Logger, handle: ConnectionHandle, conn: proto::Connection, endpoint_events: mpsc::UnboundedSender<(ConnectionHandle, EndpointEvent)>, conn_events: mpsc::UnboundedReceiver, ) -> Self { Self(Arc::new(Mutex::new(ConnectionInner { - log, epoch: Instant::now(), inner: conn, driver: None, @@ -585,7 +579,6 @@ impl std::ops::Deref for ConnectionRef { } pub struct ConnectionInner { - log: Logger, epoch: Instant, pub(crate) inner: proto::Connection, driver: Option, @@ -723,7 +716,7 @@ impl ConnectionInner { match delay.poll_unpin(cx) { Poll::Ready(()) => { *slot = None; - trace!(self.log, "{timer:?} timeout", timer = timer); + trace!("{:?} timeout", timer); self.inner.handle_timeout(now, timer); // Timeout call may have queued sends keep_going = true; @@ -745,11 +738,11 @@ impl ConnectionInner { update: proto::TimerSetting::Start(time), } => match self.timers[timer] { ref mut x @ None => { - trace!(self.log, "{timer:?} timer start", timer=timer; "time" => ?time.duration_since(self.epoch)); + trace!(time = ?time.duration_since(self.epoch), "{:?} timer start", timer); *x = Some(delay(time)); } Some(ref mut x) => { - trace!(self.log, "{timer:?} timer reset", timer=timer; "time" => ?time.duration_since(self.epoch)); + trace!(time = ?time.duration_since(self.epoch), "{:?} timer reset", timer); x.reset(time); } }, @@ -758,7 +751,7 @@ impl ConnectionInner { update: proto::TimerSetting::Stop, } => { if self.timers[timer].take().is_some() { - trace!(self.log, "{timer:?} timer stop", timer = timer); + trace!("{:?} timer stop", timer); } } } diff --git a/quinn/src/endpoint.rs b/quinn/src/endpoint.rs index a0aea30a6..e1a9915f9 100644 --- a/quinn/src/endpoint.rs +++ b/quinn/src/endpoint.rs @@ -12,7 +12,6 @@ use futures::channel::mpsc; use futures::task::{Context, Waker}; use futures::{Future, FutureExt, Poll, StreamExt}; use proto::{self as proto, ClientConfig, ConnectError, ConnectionHandle, DatagramEvent}; -use slog::Logger; use crate::builders::EndpointBuilder; use crate::connection::{Connecting, ConnectionDriver, ConnectionRef}; @@ -68,9 +67,8 @@ impl Endpoint { } else { *addr }; - let log = config.log.clone(); let (ch, conn) = endpoint.inner.connect(config, addr, server_name)?; - Ok(Connecting::new(endpoint.create_connection(log, ch, conn))) + Ok(Connecting::new(endpoint.create_connection(ch, conn))) } /// Switch to a new UDP socket @@ -172,7 +170,6 @@ impl Drop for EndpointDriver { #[derive(Debug)] pub(crate) struct EndpointInner { - log: Logger, socket: UdpSocket, inner: proto::Endpoint, outgoing: VecDeque, @@ -202,7 +199,7 @@ impl EndpointInner { Poll::Ready(Ok((n, addr, ecn))) => { match self.inner.handle(now, addr, ecn, (&buf[0..n]).into()) { Some((handle, DatagramEvent::NewConnection(conn))) => { - let conn = ConnectionDriver(self.create_connection(None, handle, conn)); + let conn = ConnectionDriver(self.create_connection(handle, conn)); if !self.incoming_live { conn.0.lock().unwrap().implicit_close(); } @@ -316,7 +313,6 @@ impl EndpointInner { fn create_connection( &mut self, - log: Option, handle: ConnectionHandle, conn: proto::Connection, ) -> ConnectionRef { @@ -329,13 +325,7 @@ impl EndpointInner { .unwrap(); } self.connections.insert(handle, send); - ConnectionRef::new( - log.unwrap_or_else(|| self.log.clone()), - handle, - conn, - self.sender.clone(), - recv, - ) + ConnectionRef::new(handle, conn, self.sender.clone(), recv) } } @@ -390,10 +380,9 @@ impl Drop for Incoming { pub(crate) struct EndpointRef(Arc>); impl EndpointRef { - pub(crate) fn new(log: Logger, socket: UdpSocket, inner: proto::Endpoint, ipv6: bool) -> Self { + pub(crate) fn new(socket: UdpSocket, inner: proto::Endpoint, ipv6: bool) -> Self { let (sender, events) = mpsc::unbounded(); Self(Arc::new(Mutex::new(EndpointInner { - log, socket, inner, ipv6, diff --git a/quinn/src/lib.rs b/quinn/src/lib.rs index 123f6a539..f5fa6f0ce 100644 --- a/quinn/src/lib.rs +++ b/quinn/src/lib.rs @@ -45,9 +45,6 @@ //! encryption alone. #![warn(missing_docs)] -#[macro_use] -extern crate slog; - mod broadcast; mod builders; mod platform; diff --git a/quinn/src/tests.rs b/quinn/src/tests.rs index bc308fef3..3c7d19a65 100644 --- a/quinn/src/tests.rs +++ b/quinn/src/tests.rs @@ -1,11 +1,12 @@ use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr, UdpSocket}; use std::sync::Arc; use std::time::{Duration, Instant}; -use std::{fmt, io, str}; +use std::{io, str}; use futures::{future, FutureExt, StreamExt, TryFutureExt}; -use slog::{o, Drain, Logger, KV}; use tokio; +use tracing::{info, info_span}; +use tracing_futures::Instrument as _; use super::{ ClientConfigBuilder, Endpoint, EndpointDriver, Incoming, NewConnection, RecvStream, SendStream, @@ -14,9 +15,8 @@ use super::{ #[test] fn handshake_timeout() { - let mut client = Endpoint::builder(); - client.logger(logger()); - let (client_driver, client, _) = client + let _guard = subscribe(); + let (client_driver, client, _) = Endpoint::builder() .bind(&SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)) .unwrap(); @@ -55,8 +55,8 @@ fn handshake_timeout() { #[test] fn drop_endpoint() { - let endpoint = Endpoint::builder(); - let (driver, endpoint, _) = endpoint + let _guard = subscribe(); + let (driver, endpoint, _) = Endpoint::builder() .bind(&SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)) .unwrap(); @@ -86,6 +86,7 @@ fn drop_endpoint() { #[test] fn drop_endpoint_driver() { + let _guard = subscribe(); let endpoint = Endpoint::builder(); let (_, endpoint, _) = endpoint .bind(&SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)) @@ -101,6 +102,7 @@ fn drop_endpoint_driver() { #[test] fn close_endpoint() { + let _guard = subscribe(); let endpoint = Endpoint::builder(); let (_driver, endpoint, incoming) = endpoint .bind(&SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)) @@ -143,7 +145,8 @@ fn local_addr() { #[test] fn read_after_close() { - let (_, driver, endpoint, mut incoming) = endpoint(); + let _guard = subscribe(); + let (driver, endpoint, mut incoming) = endpoint(); let mut runtime = tokio::runtime::current_thread::Runtime::new().unwrap(); runtime.spawn(driver.unwrap_or_else(|e| panic!("{}", e))); const MSG: &[u8] = b"goodbye!"; @@ -184,10 +187,9 @@ fn read_after_close() { } /// Construct an endpoint suitable for connecting to itself -fn endpoint() -> (Logger, EndpointDriver, Endpoint, Incoming) { +fn endpoint() -> (EndpointDriver, Endpoint, Incoming) { let mut endpoint = Endpoint::builder(); - let log = logger(); let mut server_config = ServerConfigBuilder::default(); let cert = rcgen::generate_simple_self_signed(vec!["localhost".into()]).unwrap(); let key = crate::PrivateKey::from_der(&cert.serialize_private_key_der()).unwrap(); @@ -197,20 +199,19 @@ fn endpoint() -> (Logger, EndpointDriver, Endpoint, Incoming) { endpoint.listen(server_config.build()); let mut client_config = ClientConfigBuilder::default(); - client_config.logger(log.new(o!("side" => "Client"))); client_config.add_certificate_authority(cert).unwrap(); endpoint.default_client_config(client_config.build()); - endpoint.logger(log.new(o!("side" => "Server"))); let (x, y, z) = endpoint .bind(&SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)) .unwrap(); - (log, x, y, z) + (x, y, z) } #[test] fn zero_rtt() { - let (log, driver, endpoint, incoming) = endpoint(); + let _guard = subscribe(); + let (driver, endpoint, incoming) = endpoint(); let mut runtime = tokio::runtime::current_thread::Runtime::new().unwrap(); runtime.spawn(driver.unwrap_or_else(|e| panic!("{}", e))); @@ -265,7 +266,7 @@ fn zero_rtt() { }); driver.unwrap_or_else(|_| ()).await }); - info!(log, "initial connection complete"); + info!("initial connection complete"); let ( NewConnection { connection, @@ -333,11 +334,11 @@ fn echo_dualstack() { } fn run_echo(client_addr: SocketAddr, server_addr: SocketAddr) { + let _guard = subscribe(); let mut runtime = tokio::runtime::current_thread::Runtime::new().unwrap(); { // We don't use the `endpoint` helper here because we want two different endpoints with // different addresses. - let log = logger(); let mut server_config = ServerConfigBuilder::default(); let cert = rcgen::generate_simple_self_signed(vec!["localhost".into()]).unwrap(); let key = crate::PrivateKey::from_der(&cert.serialize_private_key_der()).unwrap(); @@ -346,7 +347,6 @@ fn run_echo(client_addr: SocketAddr, server_addr: SocketAddr) { server_config.certificate(cert_chain, key).unwrap(); let mut server = Endpoint::builder(); - server.logger(log.new(o!("side" => "Server"))); server.listen(server_config.build()); let server_sock = UdpSocket::bind(server_addr).unwrap(); let server_addr = server_sock.local_addr().unwrap(); @@ -356,35 +356,48 @@ fn run_echo(client_addr: SocketAddr, server_addr: SocketAddr) { client_config.add_certificate_authority(cert).unwrap(); client_config.enable_keylog(); let mut client = Endpoint::builder(); - client.logger(log.new(o!("side" => "Client"))); client.default_client_config(client_config.build()); let (client_driver, client, _) = client.bind(&client_addr).unwrap(); - runtime.spawn(server_driver.unwrap_or_else(|e| panic!("server driver failed: {}", e))); - runtime.spawn(client_driver.unwrap_or_else(|e| panic!("client driver failed: {}", e))); + runtime.spawn( + server_driver + .unwrap_or_else(|e| panic!("server driver failed: {}", e)) + .instrument(info_span!("server endpoint")), + ); + runtime.spawn( + client_driver + .unwrap_or_else(|e| panic!("client driver failed: {}", e)) + .instrument(info_span!("client endpoint")), + ); runtime.spawn(async move { let incoming = server_incoming.next().await.unwrap(); - let new_conn = incoming.await.unwrap(); + let new_conn = incoming.instrument(info_span!("server")).await.unwrap(); tokio::spawn( new_conn .bi_streams .take_while(|x| future::ready(x.is_ok())) .for_each(|s| echo(s.unwrap())), ); - new_conn.driver.unwrap_or_else(|_| ()).await + new_conn + .driver + .unwrap_or_else(|_| ()) + .instrument(info_span!("server")) + .await }); - info!(log, "connecting from {} to {}", client_addr, server_addr); + info!("connecting from {} to {}", client_addr, server_addr); runtime.block_on(async move { let new_conn = client .connect(&server_addr, "localhost") .unwrap() + .instrument(info_span!("client")) .await .expect("connect"); tokio::spawn( new_conn .driver - .unwrap_or_else(|e| eprintln!("outgoing connection lost: {}", e)), + .unwrap_or_else(|e| eprintln!("outgoing connection lost: {}", e)) + .instrument(info_span!("client")), ); let (mut send, recv) = new_conn.connection.open_bi().await.expect("stream open"); send.write_all(b"foo").await.expect("write"); @@ -406,39 +419,25 @@ async fn echo((mut send, recv): (SendStream, RecvStream)) { let _ = send.finish().await; } -fn logger() -> Logger { - Logger::root(TestDrain.fuse(), o!()) +pub fn subscribe() -> tracing::subscriber::DefaultGuard { + let sub = tracing_subscriber::FmtSubscriber::builder() + .with_env_filter("quinn=trace") + .with_writer(|| TestWriter) + .finish(); + tracing::subscriber::set_default(sub) } -struct TestDrain; +struct TestWriter; -impl Drain for TestDrain { - type Ok = (); - type Err = io::Error; - fn log(&self, record: &slog::Record<'_>, values: &slog::OwnedKVList) -> Result<(), io::Error> { - let mut vals = Vec::new(); - values.serialize(&record, &mut TestSerializer(&mut vals))?; - record - .kv() - .serialize(&record, &mut TestSerializer(&mut vals))?; - println!( - "{} {}{}", - record.level(), - record.msg(), - str::from_utf8(&vals).unwrap() +impl std::io::Write for TestWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + print!( + "{}", + str::from_utf8(buf).expect("tried to log invalid UTF-8") ); - Ok(()) - } -} - -struct TestSerializer<'a, W>(&'a mut W); - -impl<'a, W> slog::Serializer for TestSerializer<'a, W> -where - W: io::Write + 'a, -{ - fn emit_arguments(&mut self, key: slog::Key, val: &fmt::Arguments<'_>) -> slog::Result { - write!(self.0, ", {}: {}", key, val).unwrap(); - Ok(()) + Ok(buf.len()) + } + fn flush(&mut self) -> io::Result<()> { + io::stdout().flush() } } diff --git a/quinn/tests/many_connections.rs b/quinn/tests/many_connections.rs index cf6b098fd..075d58d94 100644 --- a/quinn/tests/many_connections.rs +++ b/quinn/tests/many_connections.rs @@ -1,12 +1,9 @@ use std::sync::{Arc, Mutex}; -use std::time::Instant; -use std::{fmt, io, str}; use crc::crc32; use futures::{future, FutureExt, StreamExt, TryFutureExt, TryStreamExt}; use quinn::{ConnectionError, ReadError, WriteError}; use rand::{self, RngCore}; -use slog::{Drain, Logger, KV}; use tokio::runtime::current_thread::{self, Runtime}; use unwrap::unwrap; @@ -17,6 +14,13 @@ struct Shared { #[test] #[ignore] fn connect_n_nodes_to_1_and_send_1mb_data() { + tracing::subscriber::set_global_default( + tracing_subscriber::FmtSubscriber::builder() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .finish(), + ) + .unwrap(); + let mut runtime = unwrap!(Runtime::new()); let shared = Arc::new(Mutex::new(Shared { errors: vec![] })); @@ -30,15 +34,12 @@ fn connect_n_nodes_to_1_and_send_1mb_data() { let expected_messages = 50; - let epoch = Instant::now(); let shared2 = shared.clone(); let read_incoming_data = incoming_conns .filter_map(|connect| connect.map(|x| x.ok())) .take(expected_messages as u64) .for_each(move |new_conn| { let conn = new_conn.connection; - let logs = LogBuffer::new(); - conn.set_logger(Logger::root(logs.clone().fuse(), slog::o!())); current_thread::spawn(new_conn.driver.unwrap_or_else(|_| ())); let shared = shared2.clone(); @@ -52,11 +53,6 @@ fn connect_n_nodes_to_1_and_send_1mb_data() { }) }) .unwrap_or_else(move |e| { - let logs = logs.buffer.lock().unwrap(); - eprintln!("======== incoming connection failed: {}\nlogs:", e); - for (time, line) in &*logs { - eprintln!("{:?} {}", *time - epoch, line); - } shared.lock().unwrap().errors.push(e); }); current_thread::spawn(task); @@ -65,11 +61,9 @@ fn connect_n_nodes_to_1_and_send_1mb_data() { }); runtime.spawn(read_incoming_data); - let mut client_cfg = configure_connector(&listener_cert); + let client_cfg = configure_connector(&listener_cert); for _ in 0..expected_messages { - let logs = LogBuffer::new(); - client_cfg.log = Some(Logger::root(logs.clone().fuse(), slog::o!())); let data = random_data_with_hash(1024 * 1024); let shared = shared.clone(); let task = unwrap!(endpoint.connect_with(client_cfg.clone(), &listener_addr, "localhost")) @@ -88,17 +82,10 @@ fn connect_n_nodes_to_1_and_send_1mb_data() { | quinn::ConnectionError::Reset => {} // TODO: Determine why packet loss during connection close leads to this timing out // even though valid stateless reset packets are sent. - _ => { - let logs = logs.buffer.lock().unwrap(); - eprintln!("======== outgoing connection failed: {}\nlogs:", e); - for (time, line) in &*logs { - eprintln!("{:?} {}", *time - epoch, line); - } - if let quinn::ConnectionError::TimedOut = e { - } else { - shared.lock().unwrap().errors.push(e); - } - } + _ => match e { + quinn::ConnectionError::TimedOut => {} + _ => shared.lock().unwrap().errors.push(e), + }, } }); runtime.spawn(task); @@ -215,49 +202,3 @@ fn random_vec(size: usize) -> Vec { rand::thread_rng().fill_bytes(&mut ret[..]); ret } - -#[derive(Clone)] -struct LogBuffer { - buffer: Arc>>, -} - -impl LogBuffer { - fn new() -> Self { - Self { - buffer: Arc::new(Mutex::new(Vec::new())), - } - } -} - -impl Drain for LogBuffer { - type Ok = (); - type Err = (); - - fn log(&self, record: &slog::Record, _values: &slog::OwnedKVList) -> Result<(), ()> { - let mut kv = Vec::new(); - record - .kv() - .serialize(&record, &mut TestSerializer(&mut kv)) - .unwrap(); - let line = format!( - "{} {}{}", - record.level(), - record.msg(), - str::from_utf8(&kv).unwrap() - ); - self.buffer.lock().unwrap().push((Instant::now(), line)); - Ok(()) - } -} - -struct TestSerializer<'a, W>(&'a mut W); - -impl<'a, W> slog::Serializer for TestSerializer<'a, W> -where - W: io::Write + 'a, -{ - fn emit_arguments(&mut self, key: slog::Key, val: &fmt::Arguments<'_>) -> slog::Result { - write!(self.0, ", {}: {}", key, val).unwrap(); - Ok(()) - } -}