hubcio commented on code in PR #4193:
URL: https://github.com/apache/iggy/pull/4193#discussion_r4016273497


##########
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:
   `DupFailed` preserves `io::Error`, and the accept callback logs it, so 
`EMFILE` remains visible. separate failure counters would help monitoring, but 
the error details are not lost.



##########
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 shares the existing TCP/WS queue, while TLS handshakes run in spawned 
tasks. queue contention under connection churn needs measurement before 
changing scheduling.



##########
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:
   the owner-bit rule is shared with the adjacent setup variants. a common 
comment would clarify it, but documenting fd and `Arc` drops repeats standard 
ownership behavior.



##########
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:
   `SharedTlsServerConfig` lets `shard` carry the configuration without adding 
a direct rustls dependency. replacing the local `Arc<rustls::ServerConfig>` 
signatures would add no type-safety benefit.



-- 
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]

Reply via email to