diff --git a/src/testutil.rs b/src/testutil.rs index 65874819..f55e42e8 100644 --- a/src/testutil.rs +++ b/src/testutil.rs @@ -19,6 +19,11 @@ pub(crate) fn make_node_addr(val: u8) -> NodeAddr { pub(crate) struct LogCapture(std::sync::Arc>>); impl LogCapture { + /// Every captured line, each prefixed with its level. + pub(crate) fn lines(&self) -> Vec { + self.0.lock().unwrap().clone() + } + /// Only the captured lines emitted at WARN. pub(crate) fn warnings(&self) -> Vec { self.0 diff --git a/src/transport/tcp/mod.rs b/src/transport/tcp/mod.rs index 9b9d87bc..e85940ff 100644 --- a/src/transport/tcp/mod.rs +++ b/src/transport/tcp/mod.rs @@ -382,7 +382,9 @@ impl TcpTransport { let socket_addrs: Vec<_> = resolve_socket_addrs(addr).await?.collect(); let timeout_ms = self.config.connect_timeout_ms(); - let stream = match connect_to_any_addr(&socket_addrs, timeout_ms).await { + let connected = + connect_to_any_addr(self.transport_id, addr, &socket_addrs, timeout_ms).await; + let stream = match connected { Ok(stream) => stream, Err(error @ TransportError::ConnectionRefused) => { self.stats.record_connect_refused(); @@ -531,17 +533,11 @@ impl TcpTransport { // Resolve address (may involve DNS for hostnames) let socket_addrs: Vec<_> = resolve_socket_addrs(&remote_addr).await?.collect(); - let stream = match connect_to_any_addr(&socket_addrs, timeout_ms).await { + // A refusal is logged with its OS error by connect_to_any_addr. + let connected = + connect_to_any_addr(transport_id, &remote_addr, &socket_addrs, timeout_ms).await; + let stream = match connected { Ok(stream) => stream, - Err(error @ TransportError::ConnectionRefused) => { - debug!( - transport_id = %transport_id, - remote_addr = %remote_addr, - error = %error, - "Background TCP connect refused" - ); - return Err(error); - } Err(error @ TransportError::Timeout) => { debug!( transport_id = %transport_id, @@ -1093,14 +1089,18 @@ async fn tcp_receive_loop( // Socket Configuration Helpers // ============================================================================ +/// Connect to the first reachable of a peer's resolved addresses within one timeout. async fn connect_to_any_addr( + transport_id: TransportId, + remote_addr: &TransportAddr, socket_addrs: &[SocketAddr], timeout_ms: u64, ) -> Result { if socket_addrs.is_empty() { - return Err(TransportError::InvalidAddress( - "DNS resolution returned no addresses".into(), - )); + return Err(TransportError::InvalidAddress(format!( + "DNS resolution returned no addresses for {}", + remote_addr + ))); } // Try candidates in resolver order within one overall connection timeout. @@ -1112,7 +1112,14 @@ async fn connect_to_any_addr( { Ok(Ok(stream)) => Ok(stream), Ok(Err(error)) => { - trace!(error = %error, "TCP connection candidates failed"); + // Tokio returns the last candidate's error; keep it in the log, + // since the returned variant does not carry it. + debug!( + transport_id = %transport_id, + remote_addr = %remote_addr, + error = %error, + "TCP connect failed" + ); Err(TransportError::ConnectionRefused) } Err(_) => Err(TransportError::Timeout), @@ -1268,9 +1275,11 @@ mod tests { // No listener can bind port zero: binding it allocates an ephemeral port. let bad_addr = "127.0.0.1:0".parse().unwrap(); - let stream = connect_to_any_addr(&[bad_addr, good_addr], 1_000) - .await - .expect("second TCP candidate should connect"); + let remote = TransportAddr::from_string(&format!("localhost:{}", good_addr.port())); + let stream = + connect_to_any_addr(TransportId::new(1), &remote, &[bad_addr, good_addr], 1_000) + .await + .expect("second TCP candidate should connect"); let (accepted, _) = timeout(Duration::from_secs(1), listener.accept()) .await .unwrap() @@ -1281,16 +1290,46 @@ mod tests { #[tokio::test] async fn test_connect_candidates_fail() { + let id = TransportId::new(1); + let remote = TransportAddr::from_string("localhost:0"); assert!(matches!( - connect_to_any_addr(&[], 1_000).await, - Err(TransportError::InvalidAddress(_)) + connect_to_any_addr(id, &remote, &[], 1_000).await, + Err(TransportError::InvalidAddress(ref message)) if message.ends_with("localhost:0") )); assert!(matches!( - connect_to_any_addr(&["127.0.0.1:0".parse().unwrap()], 1_000).await, + connect_to_any_addr(id, &remote, &["127.0.0.1:0".parse().unwrap()], 1_000).await, Err(TransportError::ConnectionRefused) )); } + #[tokio::test] + async fn test_connect_failure_logs_os_error_with_peer() { + let remote = TransportAddr::from_string("localhost:0"); + let refused = ["127.0.0.1:0".parse().unwrap()]; + // Register the log callsite before installing the capture. Other + // tests reach it in parallel, and one that registers it first while + // this capture is the only subscriber caches its own thread's "never" + // interest; installing the capture rebuilds already-registered + // callsites, so registering first makes the capture see the event. + let _ = connect_to_any_addr(TransportId::new(7), &remote, &refused, 1_000).await; + let (logs, guard) = crate::testutil::capture_logs_scoped(); + let result = connect_to_any_addr(TransportId::new(7), &remote, &refused, 1_000).await; + drop(guard); + + assert!(matches!(result, Err(TransportError::ConnectionRefused))); + // The returned variant always reads "connection refused"; the log + // line must carry the OS error and the peer instead. + let lines = logs.lines(); + assert!( + lines.iter().any(|line| line.starts_with("DEBUG") + && line.contains("TCP connect failed") + && line.contains("transport_id=transport:7") + && line.contains("remote_addr=localhost:0") + && line.contains("os error")), + "expected a debug line with the OS error and peer, got {lines:?}", + ); + } + async fn send_to_hostname(background: bool) { // Prefer the last localhost address so dual-stack hosts also // exercise fallback through the full transport connection path.