Copy bytes from WebSocket binary/text frames into an async writer.
(
mut ws_source: futures::stream::SplitStream<axum::extract::ws::WebSocket>,
mut writer: W,
)
| 125 | |
| 126 | /// Copy bytes from WebSocket binary/text frames into an async writer. |
| 127 | async fn copy_ws_to_writer<W>( |
| 128 | mut ws_source: futures::stream::SplitStream<axum::extract::ws::WebSocket>, |
| 129 | mut writer: W, |
| 130 | ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> |
| 131 | where |
| 132 | W: AsyncWrite + Unpin, |
| 133 | { |
| 134 | while let Some(msg) = ws_source.next().await { |
| 135 | match msg { |
| 136 | Ok(Message::Binary(data)) => { |
| 137 | if writer.write_all(&data).await.is_err() { |
| 138 | break; |
| 139 | } |
| 140 | } |
| 141 | Ok(Message::Text(text)) => { |
| 142 | // Some proxies send text frames — treat as binary. |
| 143 | if writer.write_all(text.as_bytes()).await.is_err() { |
| 144 | break; |
| 145 | } |
| 146 | } |
| 147 | Ok(Message::Close(_)) => break, |
| 148 | Ok(Message::Ping(_) | Message::Pong(_)) => { |
| 149 | // Handled automatically by axum's WebSocket. |
| 150 | } |
| 151 | Err(e) => { |
| 152 | debug!(error = %e, "WS tunnel: ws read error"); |
| 153 | break; |
| 154 | } |
| 155 | } |
| 156 | } |
| 157 | let _ = writer.shutdown().await; |
| 158 | Ok(()) |
| 159 | } |
no test coverage detected