(
socket: TcpStream,
server_dir: ServerDir,
state_ref: StateRef,
carrier: &Senders,
end_flag: Arc<Notify>,
key: Option<Arc<SecretKey>>,
)
| 61 | } |
| 62 | |
| 63 | async fn handle_client( |
| 64 | socket: TcpStream, |
| 65 | server_dir: ServerDir, |
| 66 | state_ref: StateRef, |
| 67 | carrier: &Senders, |
| 68 | end_flag: Arc<Notify>, |
| 69 | key: Option<Arc<SecretKey>>, |
| 70 | ) -> crate::Result<()> { |
| 71 | log::debug!("New client connection"); |
| 72 | let socket = accept_client(socket, key).await?; |
| 73 | let (tx, rx) = socket.split(); |
| 74 | |
| 75 | client_rpc_loop(tx, rx, server_dir, state_ref, carrier, end_flag).await; |
| 76 | log::debug!("Client connection ended"); |
| 77 | Ok(()) |
| 78 | } |
| 79 | |
| 80 | async fn stream_history_events<Tx: Sink<ToClientMessage, Error = tako::Error> + Unpin + 'static>( |
| 81 | tx: &mut Tx, |
no test coverage detected