Copy bytes from an async reader into WebSocket binary frames.
(
mut reader: R,
mut ws_sink: futures::stream::SplitSink<axum::extract::ws::WebSocket, Message>,
)
| 91 | |
| 92 | /// Copy bytes from an async reader into WebSocket binary frames. |
| 93 | async fn copy_reader_to_ws<R>( |
| 94 | mut reader: R, |
| 95 | mut ws_sink: futures::stream::SplitSink<axum::extract::ws::WebSocket, Message>, |
| 96 | ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> |
| 97 | where |
| 98 | R: AsyncRead + Unpin, |
| 99 | { |
| 100 | let mut buf = vec![0u8; 32 * 1024]; |
| 101 | loop { |
| 102 | match reader.read(&mut buf).await { |
| 103 | Ok(0) => { |
| 104 | let _ = ws_sink.close().await; |
| 105 | break; |
| 106 | } |
| 107 | Ok(n) => { |
| 108 | if ws_sink |
| 109 | .send(Message::Binary(buf[..n].to_vec().into())) |
| 110 | .await |
| 111 | .is_err() |
| 112 | { |
| 113 | break; |
| 114 | } |
| 115 | } |
| 116 | Err(e) => { |
| 117 | debug!(error = %e, "WS tunnel: read error"); |
| 118 | let _ = ws_sink.close().await; |
| 119 | break; |
| 120 | } |
| 121 | } |
| 122 | } |
| 123 | Ok(()) |
| 124 | } |
| 125 | |
| 126 | /// Copy bytes from WebSocket binary/text frames into an async writer. |
| 127 | async fn copy_ws_to_writer<W>( |
no test coverage detected