| 130 | } |
| 131 | |
| 132 | async fn relay_messages( |
| 133 | wss_stream: WebSocketStream<TlsStream<TcpStream>>, |
| 134 | ws_address: SocketAddr, |
| 135 | ) -> Result<(), anyhow::Error> { |
| 136 | let (ws_stream, _ws_response) = |
| 137 | tokio_tungstenite::connect_async(format!("ws://{}", ws_address)).await?; |
| 138 | let (mut wss_sender, mut wss_receiver) = wss_stream.split(); |
| 139 | let (mut ws_sender, mut ws_receiver) = ws_stream.split(); |
| 140 | |
| 141 | /* Relay from WSS to WS */ |
| 142 | tokio::spawn(async move { |
| 143 | while let Some(writer) = wss_receiver.next().await { |
| 144 | if let Ok(msg) = writer { |
| 145 | if let Err(e) = ws_sender.send(msg.clone()).await { |
| 146 | log::debug!("Error sending message to WS server: {}", e); |
| 147 | break; |
| 148 | } |
| 149 | } |
| 150 | } |
| 151 | }); |
| 152 | |
| 153 | /* Relay from WS to WSS */ |
| 154 | tokio::spawn(async move { |
| 155 | while let Some(msg) = ws_receiver.next().await { |
| 156 | if let Ok(msg) = msg { |
| 157 | if let Err(e) = wss_sender.send(msg.clone()).await { |
| 158 | log::debug!("Error sending message to WSS client: {}", e); |
| 159 | break; |
| 160 | } |
| 161 | } |
| 162 | } |
| 163 | }); |
| 164 | Ok(()) |
| 165 | } |
| 166 | |
| 167 | /* Workaround: Using log crate right before plugin exit will not print */ |
| 168 | fn log_error(error: String) { |