|
5 | 5 | //! WebSocket transport
|
6 | 6 |
|
7 | 7 | use std::fmt;
|
| 8 | +use std::pin::Pin; |
8 | 9 | use std::sync::Arc;
|
| 10 | +use std::task::{Context, Poll}; |
9 | 11 | use std::time::Duration;
|
10 | 12 |
|
| 13 | +use async_utility::futures_util::stream::SplitSink; |
11 | 14 | use async_wsocket::futures_util::{Sink, SinkExt, Stream, StreamExt, TryStreamExt};
|
12 | 15 | use async_wsocket::{ConnectionMode, Message, WebSocket};
|
13 | 16 | use nostr::util::BoxedFuture;
|
@@ -94,9 +97,43 @@ impl WebSocketTransport for DefaultWebsocketTransport {
|
94 | 97 |
|
95 | 98 | // Split sink and stream
|
96 | 99 | let (tx, rx) = socket.split();
|
97 |
| - let sink: BoxSink = Box::new(tx.sink_map_err(TransportError::backend)) as BoxSink; |
| 100 | + |
| 101 | + // NOTE: don't use sink_map_err here, as it may cause panics! |
| 102 | + // Issue: https://github.com/rust-nostr/nostr/issues/984 |
| 103 | + let sink: BoxSink = Box::new(TransportSink(tx)) as BoxSink; |
98 | 104 | let stream: BoxStream = Box::new(rx.map_err(TransportError::backend)) as BoxStream;
|
| 105 | + |
99 | 106 | Ok((sink, stream))
|
100 | 107 | })
|
101 | 108 | }
|
102 | 109 | }
|
| 110 | + |
| 111 | +struct TransportSink(SplitSink<WebSocket, Message>); |
| 112 | + |
| 113 | +impl Sink<Message> for TransportSink { |
| 114 | + type Error = TransportError; |
| 115 | + |
| 116 | + fn poll_ready(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> { |
| 117 | + Pin::new(&mut self.0) |
| 118 | + .poll_ready_unpin(cx) |
| 119 | + .map_err(TransportError::backend) |
| 120 | + } |
| 121 | + |
| 122 | + fn start_send(mut self: Pin<&mut Self>, item: Message) -> Result<(), Self::Error> { |
| 123 | + Pin::new(&mut self.0) |
| 124 | + .start_send_unpin(item) |
| 125 | + .map_err(TransportError::backend) |
| 126 | + } |
| 127 | + |
| 128 | + fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> { |
| 129 | + Pin::new(&mut self.0) |
| 130 | + .poll_flush_unpin(cx) |
| 131 | + .map_err(TransportError::backend) |
| 132 | + } |
| 133 | + |
| 134 | + fn poll_close(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> { |
| 135 | + Pin::new(&mut self.0) |
| 136 | + .poll_close_unpin(cx) |
| 137 | + .map_err(TransportError::backend) |
| 138 | + } |
| 139 | +} |
0 commit comments