Skip to content

Commit 2bb7671

Browse files
fix(sdk): bound a proxied connection by one deadline and keep the TLS config error's cause
Through a SOCKS5 proxy the connector got the whole connect budget and TLS then started a second one of the same length, so a delayed CONNECT followed by a stalled handshake could take twice `connect_timeout`. The channel's own connect timeout now covers the tunnel and TLS together; the connector fixes its phase deadlines when it is called, before tonic's clock starts, so a stuck proxy is still reported as a ProxyError. An invalid TLS configuration now keeps tonic's cause (for example the invalid DNS name) in the status message and source chain instead of only "transport error". Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
1 parent 324dbc9 commit 2bb7671

3 files changed

Lines changed: 137 additions & 41 deletions

File tree

‎packages/rs-dapi-client/src/transport/proxy.rs‎

Lines changed: 27 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -146,24 +146,19 @@ impl Socks5Proxy {
146146
format!("{} {auth}", self.endpoint)
147147
}
148148

149-
/// Opens a tunnel to `uri`'s host through the proxy within
150-
/// `connect_timeout`: the proxy gets [PROXY_ANSWER_TIMEOUT] (or less, if
151-
/// the budget is shorter) to accept and negotiate, and running out of it
152-
/// is the proxy's failure; the CONNECT reply gets what is left, and
153-
/// running out of that is the destination's.
149+
/// Opens a tunnel to `uri`'s host through the proxy by `deadline`: the
150+
/// proxy gets [PROXY_ANSWER_TIMEOUT] (or less, if the deadline is
151+
/// closer) to accept and negotiate, and running out of it is the proxy's
152+
/// failure; the CONNECT reply gets what is left, and running out of that
153+
/// is the destination's.
154154
async fn connect(
155155
self,
156156
uri: Uri,
157-
connect_timeout: Duration,
157+
started: tokio::time::Instant,
158+
deadline: tokio::time::Instant,
158159
) -> Result<Box<dyn ProxiedStream>, BoxError> {
159160
let target = target(&uri)?;
160-
let now = tokio::time::Instant::now();
161-
// A budget too large to add (an "unlimited" Duration::MAX) is
162-
// capped rather than overflowing.
163-
let deadline = now
164-
.checked_add(connect_timeout)
165-
.unwrap_or_else(|| now + UNLIMITED);
166-
let negotiated_by = deadline.min(now + PROXY_ANSWER_TIMEOUT);
161+
let negotiated_by = deadline.min(started + PROXY_ANSWER_TIMEOUT);
167162
let proxy_error = |source: BoxError| ProxyError {
168163
proxy: self.endpoint.to_string(),
169164
source,
@@ -439,6 +434,13 @@ impl<S: AsyncRead + AsyncWrite + Send + Unpin> ProxiedStream for S {}
439434

440435
/// The tonic connector that tunnels every connection through the proxy
441436
/// within the connect budget (see [Socks5Proxy::connect]).
437+
///
438+
/// The channel's own connect timeout, set to the same budget, bounds the
439+
/// tunnel and the TLS handshake after it together. Its clock starts when
440+
/// the connection future is first polled, after [call](tower_service::Service::call)
441+
/// returns, and it polls the connection before checking its own deadline, so
442+
/// the connector's deadlines, fixed in `call`, always fire first: a proxy
443+
/// that stalls the negotiation is still reported as a [ProxyError].
442444
#[derive(Clone)]
443445
pub(crate) struct Socks5Connector {
444446
pub(crate) proxy: Socks5Proxy,
@@ -456,8 +458,18 @@ impl tower_service::Service<Uri> for Socks5Connector {
456458

457459
fn call(&mut self, uri: Uri) -> Self::Future {
458460
let proxy = self.proxy.clone();
459-
let connect_timeout = self.connect_timeout;
460-
Box::pin(async move { proxy.connect(uri, connect_timeout).await.map(TokioIo::new) })
461+
let started = tokio::time::Instant::now();
462+
// A budget too large to add (an "unlimited" Duration::MAX) is
463+
// capped rather than overflowing.
464+
let deadline = started
465+
.checked_add(self.connect_timeout)
466+
.unwrap_or_else(|| started + UNLIMITED);
467+
Box::pin(async move {
468+
proxy
469+
.connect(uri, started, deadline)
470+
.await
471+
.map(TokioIo::new)
472+
})
461473
}
462474
}
463475

‎packages/rs-dapi-client/src/transport/tonic_channel.rs‎

Lines changed: 60 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,8 @@ use crate::{request_settings::AppliedRequestSettings, Uri};
44
use dapi_grpc::core::v0::core_client::CoreClient;
55
use dapi_grpc::platform::v0::platform_client::PlatformClient;
66
use dapi_grpc::tonic::transport::{Certificate, Channel, ClientTlsConfig};
7+
use std::error::Error as _;
8+
use std::sync::Arc;
79
use std::time::Duration;
810

911
/// Platform Client using gRPC transport.
@@ -54,22 +56,8 @@ pub fn create_channel(
5456

5557
let proxy = settings.and_then(|settings| settings.proxy.clone());
5658
let connect_timeout = settings.and_then(AppliedRequestSettings::effective_connect_timeout);
57-
if let Some(settings) = settings {
58-
match (connect_timeout, &proxy) {
59-
// Through a proxy the connector enforces the budget itself, so
60-
// it can tell a stuck proxy from a slow destination; TLS, which
61-
// runs after it, gets the same budget of its own (a connection
62-
// can so take up to twice the budget; the executor's attempt
63-
// deadline still bounds the whole attempt).
64-
(Some(timeout), Some(_)) => tls_config = tls_config.timeout(timeout),
65-
(Some(timeout), None) => builder = builder.connect_timeout(timeout),
66-
(None, _) => {}
67-
}
68-
69-
if let Some(pem) = settings.ca_certificate.as_ref() {
70-
let cert = Certificate::from_pem(pem);
71-
tls_config = tls_config.ca_certificate(cert);
72-
};
59+
if let Some(pem) = settings.and_then(|settings| settings.ca_certificate.as_ref()) {
60+
tls_config = tls_config.ca_certificate(Certificate::from_pem(pem));
7361
}
7462

7563
// Ping only while a request is in flight: a connection whose network path
@@ -81,21 +69,48 @@ pub fn create_channel(
8169
.keep_alive_timeout(HTTP2_KEEP_ALIVE_TIMEOUT)
8270
.keep_alive_while_idle(false);
8371

84-
builder = builder.tls_config(tls_config).map_err(|e| {
85-
TransportError::Grpc(dapi_grpc::tonic::Status::invalid_argument(format!(
86-
"invalid TLS configuration: {e}"
87-
)))
88-
})?;
72+
builder = builder.tls_config(tls_config).map_err(invalid_tls_config)?;
8973

9074
Ok(match proxy {
91-
None => builder.connect_lazy(),
92-
Some(proxy) => builder.connect_with_connector_lazy(Socks5Connector {
93-
proxy,
94-
connect_timeout: connect_timeout.unwrap_or(PROXY_CONNECT_TIMEOUT),
95-
}),
75+
None => {
76+
if let Some(timeout) = connect_timeout {
77+
builder = builder.connect_timeout(timeout);
78+
}
79+
builder.connect_lazy()
80+
}
81+
Some(proxy) => {
82+
// One absolute deadline for the whole connection: tonic applies
83+
// it around the connector, so it covers the SOCKS5 tunnel and the
84+
// TLS handshake after it. The connector times its own phases
85+
// against the same budget from a start that is never later than
86+
// tonic's, so it still tells a stuck proxy from a slow
87+
// destination before this deadline fires.
88+
let connect_timeout = connect_timeout.unwrap_or(PROXY_CONNECT_TIMEOUT);
89+
builder
90+
.connect_timeout(connect_timeout)
91+
.connect_with_connector_lazy(Socks5Connector {
92+
proxy,
93+
connect_timeout,
94+
})
95+
}
9696
})
9797
}
9898

99+
/// A TLS configuration tonic rejects, as an `InvalidArgument` status that
100+
/// keeps the cause: tonic's own message names only the error kind, the cause
101+
/// (an invalid server name, a malformed certificate) is in its source chain.
102+
fn invalid_tls_config(error: dapi_grpc::tonic::transport::Error) -> TransportError {
103+
let mut message = format!("invalid TLS configuration: {error}");
104+
let mut cause = error.source();
105+
while let Some(current) = cause {
106+
message.push_str(&format!(": {current}"));
107+
cause = current.source();
108+
}
109+
let mut status = dapi_grpc::tonic::Status::invalid_argument(message);
110+
status.set_source(Arc::new(error));
111+
TransportError::Grpc(status)
112+
}
113+
99114
/// The host of a URI without the brackets an IPv6 literal carries in it
100115
/// (`[2001:db8::1]` becomes `2001:db8::1`); any other host is unchanged.
101116
pub(crate) fn unbracketed(host: &str) -> &str {
@@ -135,4 +150,23 @@ mod tests {
135150
}));
136151
create_channel(uri, Some(&settings)).expect("through a proxy");
137152
}
153+
154+
#[test]
155+
fn should_keep_the_cause_of_an_invalid_tls_configuration() {
156+
// Not a valid TLS server name.
157+
let uri: Uri = "https://foo..bar:1443".parse().expect("uri");
158+
let Err(TransportError::Grpc(status)) = create_channel(uri, None) else {
159+
panic!("an invalid server name must be rejected");
160+
};
161+
assert_eq!(status.code(), dapi_grpc::tonic::Code::InvalidArgument);
162+
assert!(
163+
status.message().contains("invalid dns name"),
164+
"{}",
165+
status.message()
166+
);
167+
assert!(
168+
status.source().is_some(),
169+
"the tonic error stays in the source chain"
170+
);
171+
}
138172
}

‎packages/rs-dapi-client/tests/socks5_proxy.rs‎

Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -514,6 +514,56 @@ async fn should_hold_a_stalled_route_against_the_node() {
514514
assert!(is_banned(&client, uri));
515515
}
516516

517+
#[tokio::test]
518+
async fn should_bound_the_tunnel_and_tls_by_one_connect_deadline() {
519+
// Answers the greeting, holds the CONNECT reply for most of the budget,
520+
// grants it, and then never answers the TLS ClientHello. One budget
521+
// covers the tunnel and TLS together, so the connection fails when it
522+
// runs out instead of TLS starting a full budget of its own.
523+
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
524+
let slow = listener.local_addr().expect("address");
525+
tokio::spawn(async move {
526+
let mut held = Vec::new();
527+
while let Ok((mut stream, _)) = listener.accept().await {
528+
let mut greeting = [0u8; 3];
529+
if stream.read_exact(&mut greeting).await.is_err()
530+
|| stream.write_all(&[0x05, 0x00]).await.is_err()
531+
{
532+
continue;
533+
}
534+
// CONNECT to an IPv4 address: header, address and port.
535+
let mut request = [0u8; 10];
536+
if stream.read_exact(&mut request).await.is_err() {
537+
continue;
538+
}
539+
tokio::time::sleep(Duration::from_millis(1500)).await;
540+
if stream
541+
.write_all(&[0x05, SUCCEEDED, 0x00, 0x01, 0, 0, 0, 0, 0, 0])
542+
.await
543+
.is_ok()
544+
{
545+
held.push(stream);
546+
}
547+
}
548+
});
549+
let uri = "https://127.0.0.1:1";
550+
let client = client_with_connect_timeout(
551+
uri,
552+
tcp_proxy(slow, Socks5Auth::None),
553+
Duration::from_secs(2),
554+
);
555+
556+
let started = std::time::Instant::now();
557+
let error = get_status(&client).await;
558+
let elapsed = started.elapsed();
559+
560+
// With a TLS timeout of its own this took 1.5 s + 2 s.
561+
assert!(elapsed >= Duration::from_secs(2), "{elapsed:?}: {error:?}");
562+
assert!(elapsed < Duration::from_secs(3), "{elapsed:?}: {error:?}");
563+
assert!(!transport_error(&error).is_proxy_failure(), "{error:?}");
564+
assert!(is_banned(&client, uri));
565+
}
566+
517567
/// Starts a proxy that reads the greeting (and, if `auth`, the RFC 1929
518568
/// request), writes `reply`, and then holds the connection open.
519569
async fn start_stalling_proxy(reply: &'static [u8], auth: bool) -> SocketAddr {

0 commit comments

Comments
 (0)