numinnex commented on code in PR #4193:
URL: https://github.com/apache/iggy/pull/4193#discussion_r4015885236
##########
core/integration/tests/server/specific.rs:
##########
@@ -141,3 +164,99 @@ async fn restart_offset_skip(harness: &mut TestHarness) {
async fn segment_rotation_scenario(harness: &TestHarness) {
segment_rotation_race_scenario::run(harness).await;
}
+
+async fn assert_tls_client_sharding_and_cleanup(
+ harness: &TestHarness,
+ clients: &[IggyClient],
+ transport: &str,
+) {
+ let mut client_ids = HashSet::with_capacity(clients.len());
+ for client in clients {
+ client_ids.insert(client.get_me().await.unwrap().client_id);
+ }
+ assert_eq!(client_ids.len(), clients.len(), "client IDs must be unique");
+
+ // The wire exposes only the sequence tail. Match it to the router's
+ // full ID and thread name to prove execution placement after handoff.
+ let install_marker = format!("installing delegated {transport} client fd");
+ tokio::time::timeout(TLS_CLEANUP_TIMEOUT, async {
+ loop {
+ let logs = harness.server().stdout_plain();
+ let mut counts = [0; TLS_TEST_SHARDS];
+ let mut installed = HashSet::new();
+ for line in logs.lines().filter(|line|
line.contains(&install_marker)) {
+ let fields: Vec<_> = line.split_whitespace().collect();
+ let full_id: u128 = fields
+ .iter()
+ .find_map(|field| field.strip_prefix("client_id="))
+ .expect("install log must carry client_id")
+ .parse()
+ .unwrap();
+ let wire_id = u32::try_from(full_id &
u128::from(u32::MAX)).unwrap();
+ if !client_ids.contains(&wire_id) {
+ continue;
+ }
+ let owner = usize::try_from(full_id >> 112).unwrap();
+ assert!(owner < TLS_TEST_SHARDS, "invalid owner: {line}");
+ assert!(
+ fields.contains(&format!("shard={owner}").as_str()),
+ "{line}"
+ );
+ assert!(
+ fields.contains(&format!("shard-{owner}").as_str()),
+ "{line}"
+ );
+ assert!(installed.insert(wire_id), "client installed twice:
{line}");
+ counts[owner] += 1;
+ }
+ if installed == client_ids {
+ assert_eq!(
Review Comment:
Exact equality pins the round-robin *phase*, not just the distribution.
`ship_client_fd` advances the shared `client_rr` before `peer_addr()` and
`dup_fd()` can fail (`coordinator.rs:397-401`), so a single transient accept
failure on any of the four client transports shifts the phase and this
assertion fails. `TLS_TEST_SHARDS` is also hard-coded to 4, independently of
what the server actually bound.
##########
core/integration/tests/server/specific.rs:
##########
@@ -141,3 +164,99 @@ async fn restart_offset_skip(harness: &mut TestHarness) {
async fn segment_rotation_scenario(harness: &TestHarness) {
segment_rotation_race_scenario::run(harness).await;
}
+
+async fn assert_tls_client_sharding_and_cleanup(
+ harness: &TestHarness,
+ clients: &[IggyClient],
+ transport: &str,
+) {
+ let mut client_ids = HashSet::with_capacity(clients.len());
+ for client in clients {
+ client_ids.insert(client.get_me().await.unwrap().client_id);
+ }
+ assert_eq!(client_ids.len(), clients.len(), "client IDs must be unique");
+
+ // The wire exposes only the sequence tail. Match it to the router's
+ // full ID and thread name to prove execution placement after handoff.
+ let install_marker = format!("installing delegated {transport} client fd");
+ tokio::time::timeout(TLS_CLEANUP_TIMEOUT, async {
+ loop {
+ let logs = harness.server().stdout_plain();
+ let mut counts = [0; TLS_TEST_SHARDS];
+ let mut installed = HashSet::new();
+ for line in logs.lines().filter(|line|
line.contains(&install_marker)) {
+ let fields: Vec<_> = line.split_whitespace().collect();
+ let full_id: u128 = fields
+ .iter()
+ .find_map(|field| field.strip_prefix("client_id="))
+ .expect("install log must carry client_id")
+ .parse()
+ .unwrap();
+ let wire_id = u32::try_from(full_id &
u128::from(u32::MAX)).unwrap();
+ if !client_ids.contains(&wire_id) {
+ continue;
+ }
+ let owner = usize::try_from(full_id >> 112).unwrap();
+ assert!(owner < TLS_TEST_SHARDS, "invalid owner: {line}");
+ assert!(
+ fields.contains(&format!("shard={owner}").as_str()),
+ "{line}"
+ );
+ assert!(
+ fields.contains(&format!("shard-{owner}").as_str()),
+ "{line}"
+ );
+ assert!(installed.insert(wire_id), "client installed twice:
{line}");
+ counts[owner] += 1;
+ }
+ if installed == client_ids {
+ assert_eq!(
+ counts,
+ [TLS_TEST_CLIENTS / TLS_TEST_SHARDS; TLS_TEST_SHARDS]
+ );
+ break;
+ }
+ tokio::time::sleep(TLS_POLL_INTERVAL).await;
+ }
+ })
+ .await
+ .expect("every encrypted client must be installed on its owning shard
thread");
+
+ for client in clients {
+ client.disconnect().await.unwrap();
+ }
+ // Polling opens a separate SDK connection. A fresh observer has none,
+ // so an exact count also detects leaked partition-client sessions.
+ let observer = harness.root_client().await.unwrap();
+ let observer_id = observer.get_me().await.unwrap().client_id;
+ tokio::time::timeout(TLS_CLEANUP_TIMEOUT, async {
+ loop {
+ let connected = observer.get_clients().await.unwrap();
+ if connected.len() == 1 {
Review Comment:
This can pass vacuously. `IggyShard::list_all_clients`
(`core/shard/src/lib.rs:1979-1987`, `:2014-2024`) under-reports warn-only when
a shard's `ListClients` try_send is rejected or the 3s gather times out, so
`len() == 1` is satisfiable while another shard still holds a leaked session —
and the loop breaks on the first such iteration. Pre-PR every TLS session lived
on shard 0, so no cross-shard gather was involved.
##########
core/integration/tests/server/specific.rs:
##########
@@ -20,33 +20,56 @@ use
crate::server::scenarios::{reconnect_after_restart_scenario, restart_offset_
use crate::server::scenarios::{
segment_rotation_race_scenario, tcp_tls_scenario, websocket_tls_scenario,
};
+use iggy::prelude::*;
+use integration::harness::TestHarness;
use integration::iggy_harness;
+use std::collections::HashSet;
+use std::time::Duration;
+use tokio::io::AsyncWriteExt;
+use tokio::net::TcpStream;
+
+const TLS_TEST_SHARDS: usize = 4;
+const TLS_TEST_CLIENTS: usize = TLS_TEST_SHARDS * 2;
+const TLS_CLEANUP_TIMEOUT: Duration = Duration::from_secs(10);
+const TLS_POLL_INTERVAL: Duration = Duration::from_millis(20);
#[iggy_harness(
test_client_transport = TcpTlsGenerated,
- server(tls = generated)
+ server(tls = generated, sharding.cpu_allocation = "0..4", logging.level =
"info")
Review Comment:
`sharding.cpu_allocation = "0..4"` lands in `extra_envs`, and
`harness/handle/server.rs:374-376` applies it with `insert` — hard-overriding
the adaptive `"all"` fallback at `:297-313`. Below 4 usable CPUs
`validators.rs:218-223` rejects `Range(0,4)` and the server refuses to boot,
reporting a cpu_allocation error that names nothing about TLS. At exactly 4
CPUs it passes with zero margin while three tests each demand four shards under
concurrent nextest. Same string at `:51` and `:64`.
##########
core/message_bus/src/lib.rs:
##########
@@ -403,42 +403,26 @@ pub type AcceptedQuicClientFn = std::rc::Rc<dyn
Fn(AcceptedQuicConn)>;
/// requires a dupable plaintext fd) stays well-defined.
pub type AcceptedWsClientFn = std::rc::Rc<dyn Fn(compio::net::TcpStream)>;
+/// Listener TLS configuration shared with the shard owning each connection.
+/// Per-connection TLS state is created only on the destination runtime.
+pub type SharedTlsServerConfig = std::sync::Arc<rustls::ServerConfig>;
Review Comment:
Applied at 5 sites while ~20 on the same delegation path still spell
`Arc<rustls::ServerConfig>` (`installer/tcp_tls.rs:54`, `installer/wss.rs:53`,
`client_listener/{tcp_tls,wss}.rs:67,73,93,99`, `transports/tcp_tls.rs:143`,
`transports/wss.rs:101`). The alias also hides the `Arc` at the one boundary
where the shared-nothing invariant makes it load-bearing: a reader of `config:
SharedTlsServerConfig` (`shard/src/lib.rs:742`) cannot see that it is
atomically refcounted, crosses threads, and carries shared interior state.
##########
core/shard/src/coordinator.rs:
##########
@@ -358,31 +344,78 @@ impl ShardZeroCoordinator {
/// fails or `dup(2)` fails. Returns [`SendError::RoutingFailed`]
/// when the target shard's inbox refuses the setup frame.
pub fn delegate_ws_client(&self, stream: TcpStream) -> Result<u128,
SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::Ws, |fd, meta| {
+ LifecycleFrame::ClientWsConnectionSetup { fd, meta }
+ })
+ }
+
+ /// Delegate a raw TCP socket and its listener's TLS configuration.
+ /// The destination shard owns the TLS handshake and encrypted I/O.
+ ///
+ /// # Errors
+ ///
+ /// Returns the same lookup, duplication and routing errors as
+ /// [`Self::delegate_client`]. Both socket handles close on failure.
+ pub fn delegate_tcp_tls_client(
+ &self,
+ stream: TcpStream,
+ config: SharedTlsServerConfig,
+ ) -> Result<u128, SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::TcpTls, |fd, meta| {
+ LifecycleFrame::ClientTcpTlsConnectionSetup { fd, meta, config }
+ })
+ }
+
+ /// Delegate WSS before TLS or WebSocket state is created. Both
+ /// handshakes and all subsequent I/O run on the destination shard.
+ ///
+ /// # Errors
+ ///
+ /// Returns the same lookup, duplication and routing errors as
+ /// [`Self::delegate_client`]. Both socket handles close on failure.
+ pub fn delegate_wss_client(
+ &self,
+ stream: TcpStream,
+ config: SharedTlsServerConfig,
+ ) -> Result<u128, SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::Wss, |fd, meta| {
+ LifecycleFrame::ClientWssConnectionSetup { fd, meta, config }
+ })
+ }
+
+ #[must_use]
+ pub const fn total_shards(&self) -> u16 {
+ self.total_shards
+ }
+
+ fn ship_client_fd(
+ &self,
+ stream: TcpStream,
+ transport: ClientTransportKind,
+ build_frame: impl FnOnce(fd_transfer::DupedFd, ClientConnMeta) ->
LifecycleFrame,
+ ) -> Result<u128, SendError> {
let target = self.next_client_target();
Review Comment:
One `client_rr` cursor serves all four client transports. With k transport
classes interleaved on one cursor and N shards, each class reaches only
N/gcd(k,N) distinct shards — so an alternating accept pattern on 4 shards pins
TLS connections to a subset, for the life of each connection. The matrix test
at `:951` cannot observe this: it rotates all four transports and asserts the
aggregate sequence, not per-transport placement.
##########
core/shard/src/lib.rs:
##########
@@ -734,6 +734,21 @@ pub enum LifecycleFrame {
fd: DupedFd,
meta: ClientConnMeta,
},
+ /// Delegate TCP-TLS before reading TLS bytes. The destination wraps the
Review Comment:
Both new variants omit the routing invariant their neighbours state — that
the owning shard is encoded in the top 16 bits of `meta.client_id` — which the
reply path depends on and which this PR's own test asserts by hand-rolling `>>
112` (`specific.rs:199`). Neither names the receiving installer the way
`ClientWsConnectionSetup` names `install_client_ws_fd`, nor records that
dropping the frame releases both the fd and a config refcount.
##########
core/message_bus/src/transports/tcp_tls.rs:
##########
@@ -218,7 +218,11 @@ impl TransportConn for TcpTlsTransportConn {
let mut tls = match self.state {
ConnState::Established(tls) => *tls,
ConnState::Pending { stream, role } => {
- match compio::time::timeout(handshake_grace, handshake(role,
stream)).await {
+ let outcome = futures::select_biased! {
+ () = ctx.shutdown.wait().fuse() => return,
Review Comment:
This arm returns with no log and no comment, where every sibling exit in
this match warns with `%label, %peer`. The router has already emitted an INFO
`installing delegated TCP-TLS client fd` for the connection, which then
disappears with no record of why.
##########
core/message_bus/src/installer/replica.rs:
##########
@@ -369,8 +369,6 @@ pub fn install_replica_conn<C: TransportConn>(
// `notify_connection_lost` stands down whenever a live registry entry
// exists - but the loser also drives `replica_dispatch_loop`, which
// must never hand the winner's replica id to `on_message`.
- // `compio::runtime::JoinHandle::drop` does not cancel the spawned
- // task, so we have to tell the tasks to stand down in-band.
let install_aborted = Rc::new(Cell::new(false));
Review Comment:
The rationale for `install_aborted` was deleted rather than corrected. The
flag is still required: `drain_rejected_registration`
(`installer/common.rs:38-51`) awaits both handles under a timeout rather than
dropping them, so the loser's tasks run cooperatively to completion and must
stand down in-band or they hand the winner's replica id to `on_message`. The
three sibling `JoinHandle` comments in this PR were corrected in place; this
one now carries no stated reason at all.
##########
core/shard/src/coordinator.rs:
##########
@@ -358,31 +344,78 @@ impl ShardZeroCoordinator {
/// fails or `dup(2)` fails. Returns [`SendError::RoutingFailed`]
/// when the target shard's inbox refuses the setup frame.
pub fn delegate_ws_client(&self, stream: TcpStream) -> Result<u128,
SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::Ws, |fd, meta| {
+ LifecycleFrame::ClientWsConnectionSetup { fd, meta }
+ })
+ }
+
+ /// Delegate a raw TCP socket and its listener's TLS configuration.
+ /// The destination shard owns the TLS handshake and encrypted I/O.
+ ///
+ /// # Errors
+ ///
+ /// Returns the same lookup, duplication and routing errors as
+ /// [`Self::delegate_client`]. Both socket handles close on failure.
+ pub fn delegate_tcp_tls_client(
+ &self,
+ stream: TcpStream,
+ config: SharedTlsServerConfig,
+ ) -> Result<u128, SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::TcpTls, |fd, meta| {
+ LifecycleFrame::ClientTcpTlsConnectionSetup { fd, meta, config }
+ })
+ }
+
+ /// Delegate WSS before TLS or WebSocket state is created. Both
+ /// handshakes and all subsequent I/O run on the destination shard.
+ ///
+ /// # Errors
+ ///
+ /// Returns the same lookup, duplication and routing errors as
+ /// [`Self::delegate_client`]. Both socket handles close on failure.
+ pub fn delegate_wss_client(
+ &self,
+ stream: TcpStream,
+ config: SharedTlsServerConfig,
+ ) -> Result<u128, SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::Wss, |fd, meta| {
+ LifecycleFrame::ClientWssConnectionSetup { fd, meta, config }
+ })
+ }
+
+ #[must_use]
+ pub const fn total_shards(&self) -> u16 {
+ self.total_shards
+ }
+
+ fn ship_client_fd(
Review Comment:
No in-flight cap on un-handshaken client delegations, where the replica
plane gates every accept on `try_acquire_replica_handshake_slot`
(`message_bus/src/lib.rs:893`, cap 256). Each delegated connection commits an
inbox slot, ~18 KiB of ring buffer, two compio tasks, a `clients()` entry and a
`client_meta` entry before one TLS byte and before LOGIN, held up to
`handshake_grace`. A TLS-only node previously had no unauthenticated path into
a worker shard at all.
##########
core/shard/src/coordinator.rs:
##########
@@ -358,31 +344,78 @@ impl ShardZeroCoordinator {
/// fails or `dup(2)` fails. Returns [`SendError::RoutingFailed`]
/// when the target shard's inbox refuses the setup frame.
pub fn delegate_ws_client(&self, stream: TcpStream) -> Result<u128,
SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::Ws, |fd, meta| {
+ LifecycleFrame::ClientWsConnectionSetup { fd, meta }
+ })
+ }
+
+ /// Delegate a raw TCP socket and its listener's TLS configuration.
+ /// The destination shard owns the TLS handshake and encrypted I/O.
+ ///
+ /// # Errors
+ ///
+ /// Returns the same lookup, duplication and routing errors as
+ /// [`Self::delegate_client`]. Both socket handles close on failure.
+ pub fn delegate_tcp_tls_client(
+ &self,
+ stream: TcpStream,
+ config: SharedTlsServerConfig,
+ ) -> Result<u128, SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::TcpTls, |fd, meta| {
+ LifecycleFrame::ClientTcpTlsConnectionSetup { fd, meta, config }
+ })
+ }
+
+ /// Delegate WSS before TLS or WebSocket state is created. Both
+ /// handshakes and all subsequent I/O run on the destination shard.
+ ///
+ /// # Errors
+ ///
+ /// Returns the same lookup, duplication and routing errors as
+ /// [`Self::delegate_client`]. Both socket handles close on failure.
+ pub fn delegate_wss_client(
+ &self,
+ stream: TcpStream,
+ config: SharedTlsServerConfig,
+ ) -> Result<u128, SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::Wss, |fd, meta| {
+ LifecycleFrame::ClientWssConnectionSetup { fd, meta, config }
+ })
+ }
+
+ #[must_use]
+ pub const fn total_shards(&self) -> u16 {
+ self.total_shards
+ }
+
+ fn ship_client_fd(
+ &self,
+ stream: TcpStream,
+ transport: ClientTransportKind,
+ build_frame: impl FnOnce(fd_transfer::DupedFd, ClientConnMeta) ->
LifecycleFrame,
+ ) -> Result<u128, SendError> {
let target = self.next_client_target();
let client_id = self.mint_client_id(target);
let peer_addr = stream.peer_addr().map_err(SendError::DupFailed)?;
Review Comment:
`peer_addr()` and `dup_fd()` failures share one `SendError::DupFailed` and
record no metric, now on all four client planes. `fcntl(F_DUPFD_CLOEXEC)`
failing is EMFILE/ENFILE — the fd-exhaustion signal — and it is
indistinguishable here from getpeername on a reset peer. Only the `try_send`
branch below reaches `record_frame_drop`; successful delegations are counted
nowhere, which is why the new integration test has to scrape stdout to observe
placement.
##########
core/message_bus/src/client_listener/tcp_tls.rs:
##########
@@ -86,11 +82,7 @@ pub fn bind(
// cannot enable it accidentally.
cfg.max_early_data_size = 0;
Review Comment:
Only `max_early_data_size` is overridden; the remaining rustls defaults now
cross shards with this config. `session_storage` is
`ServerSessionMemoryCache::new(256)` behind a `std::sync::Mutex`, and with
`NeverProducesTickets` + `send_tls13_tickets = 2` the stateful branch runs two
blocking `put()`s per successful TLS 1.3 handshake. That lock was
single-threaded while TLS terminated on shard 0; it is now contended across N
shard threads on a runtime with no blocking pool, where a futex sleep parks the
whole shard.
##########
core/shard/src/coordinator.rs:
##########
@@ -358,31 +344,78 @@ impl ShardZeroCoordinator {
/// fails or `dup(2)` fails. Returns [`SendError::RoutingFailed`]
/// when the target shard's inbox refuses the setup frame.
pub fn delegate_ws_client(&self, stream: TcpStream) -> Result<u128,
SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::Ws, |fd, meta| {
+ LifecycleFrame::ClientWsConnectionSetup { fd, meta }
+ })
+ }
+
+ /// Delegate a raw TCP socket and its listener's TLS configuration.
+ /// The destination shard owns the TLS handshake and encrypted I/O.
+ ///
+ /// # Errors
+ ///
+ /// Returns the same lookup, duplication and routing errors as
+ /// [`Self::delegate_client`]. Both socket handles close on failure.
+ pub fn delegate_tcp_tls_client(
+ &self,
+ stream: TcpStream,
+ config: SharedTlsServerConfig,
+ ) -> Result<u128, SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::TcpTls, |fd, meta| {
+ LifecycleFrame::ClientTcpTlsConnectionSetup { fd, meta, config }
+ })
+ }
+
+ /// Delegate WSS before TLS or WebSocket state is created. Both
+ /// handshakes and all subsequent I/O run on the destination shard.
+ ///
+ /// # Errors
+ ///
+ /// Returns the same lookup, duplication and routing errors as
+ /// [`Self::delegate_client`]. Both socket handles close on failure.
+ pub fn delegate_wss_client(
+ &self,
+ stream: TcpStream,
+ config: SharedTlsServerConfig,
+ ) -> Result<u128, SendError> {
+ self.ship_client_fd(stream, ClientTransportKind::Wss, |fd, meta| {
+ LifecycleFrame::ClientWssConnectionSetup { fd, meta, config }
+ })
+ }
+
+ #[must_use]
+ pub const fn total_shards(&self) -> u16 {
+ self.total_shards
+ }
+
+ fn ship_client_fd(
+ &self,
+ stream: TcpStream,
+ transport: ClientTransportKind,
+ build_frame: impl FnOnce(fd_transfer::DupedFd, ClientConnMeta) ->
LifecycleFrame,
+ ) -> Result<u128, SendError> {
let target = self.next_client_target();
let client_id = self.mint_client_id(target);
let peer_addr = stream.peer_addr().map_err(SendError::DupFailed)?;
let fd = fd_transfer::dup_fd(&stream).map_err(SendError::DupFailed)?;
- let meta = ClientConnMeta::new(client_id, peer_addr,
ClientTransportKind::Ws);
- let setup = LifecycleFrame::ClientWsConnectionSetup { fd, meta };
+ let meta = ClientConnMeta::new(client_id, peer_addr, transport);
+ let setup = build_frame(fd, meta);
if let Err(e) = self.senders[target as
usize].try_send(ShardFrame::lifecycle(setup)) {
Review Comment:
Setup frames ride the main lane (`self.senders[target]` derefs to `inner`,
`core/shard/src/lib.rs:580-586`) — the same bounded queue as
`ShardFrame::Consensus`, reply lane unused, no priority. The queue does not
need to fill to hurt: every setup frame the destination pump dequeues costs
that shard tens of microseconds during which no consensus frame on it moves.
Pre-PR that cost landed on shard 0 only; it now lands on the shards carrying
consensus.
##########
core/message_bus/src/client_listener/tcp_tls.rs:
##########
@@ -86,11 +82,7 @@ pub fn bind(
// cannot enable it accidentally.
cfg.max_early_data_size = 0;
- let listener = bind_reusable_tcp_listener(addr)
- .map_err(|_| IggyError::CannotBindToSocket(addr.to_string()))?;
- let actual = listener
- .local_addr()
- .map_err(|e| IggyError::IoError(e.to_string()))?;
+ let (listener, actual) = bind_nodelay_listener(addr)?;
Review Comment:
This moves TCP-TLS onto `bind_nodelay_listener`, but `wss::bind`
(`client_listener/wss.rs:75`) still uses `bind_reusable_tcp_listener`. The two
TLS listeners were symmetric before this PR, and WSS is now the only client
plane with no listener-level Nagle setting — its sole mechanism is the
warn-and-continue per-connection call at `installer/wss.rs:58-63`. The updated
helper doc at `mod.rs:101-102` and the new test at `mod.rs:151-164` both
exclude WSS.
##########
core/message_bus/src/transports/wss.rs:
##########
@@ -174,12 +174,14 @@ impl TransportConn for WssTransportConn {
// both inner handshake futures are large (rustls + tungstenite
// state machines) and would otherwise push `run`'s state past
// the `clippy::large_futures` threshold.
- let tls_stream = match compio::time::timeout(
- handshake_grace,
- Box::pin(tls_handshake(role, self.stream)),
- )
- .await
- {
+ let tls_outcome = futures::select_biased! {
+ () = ctx.shutdown.wait().fuse() => return,
Review Comment:
Same shape as the TCP-TLS arm: a silent `return` with no log, in a function
where every other exit warns with `%label, %peer`. The second `select_biased!`
at `:218` repeats it, and that one drops an already-established `tls_stream`.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]