From 23e58902655e003a3e6c111b6f5017b16e010acc Mon Sep 17 00:00:00 2001 From: kingchenc Date: Fri, 22 May 2026 04:21:09 +0200 Subject: [PATCH] C4: track a closed flag on the Binance stream BinanceKlineStream had no closed-state flag, so after the server closed the connection (Ok(None)) a caller could keep calling next_event and poll a dead socket. A closed: bool is now set when the server closes or sends a Close frame; next_event short-circuits to Ok(None) once set, and is_closed() exposes the state. --- crates/wickra-data/src/live/binance.rs | 28 +++++++++++++++++++++++--- 1 file changed, 25 insertions(+), 3 deletions(-) diff --git a/crates/wickra-data/src/live/binance.rs b/crates/wickra-data/src/live/binance.rs index 570a7bf7..10c92d58 100644 --- a/crates/wickra-data/src/live/binance.rs +++ b/crates/wickra-data/src/live/binance.rs @@ -91,6 +91,9 @@ pub struct BinanceKlineStream { socket: WebSocketStream>, /// Interval requested at connect time. Used to tag every event. interval: Interval, + /// `true` once the server has closed the stream. A closed stream is never + /// polled again — `next_event` short-circuits to `Ok(None)`. + closed: bool, } /// Wire-format representation of an incoming Binance kline tick. Public so callers @@ -159,17 +162,33 @@ impl BinanceKlineStream { ); let url = url::Url::parse(&url).map_err(|e| Error::Malformed(e.to_string()))?; let (socket, _) = tokio_tungstenite::connect_async(url.as_str()).await?; - Ok(Self { socket, interval }) + Ok(Self { + socket, + interval, + closed: false, + }) + } + + /// Whether the server has closed the stream. Once closed, every further + /// [`next_event`](Self::next_event) call yields `Ok(None)` immediately. + pub fn is_closed(&self) -> bool { + self.closed } /// Receive the next kline event. Yields `Ok(None)` when the server closes /// the connection cleanly. pub async fn next_event(&mut self) -> Result> { + if self.closed { + return Ok(None); + } loop { let msg = match self.socket.next().await { Some(Ok(m)) => m, Some(Err(e)) => return Err(Error::from(e)), - None => return Ok(None), + None => { + self.closed = true; + return Ok(None); + } }; match msg { Message::Text(text) => { @@ -189,7 +208,10 @@ impl BinanceKlineStream { self.socket.send(Message::Pong(payload)).await?; } Message::Pong(_) | Message::Frame(_) => {} - Message::Close(_) => return Ok(None), + Message::Close(_) => { + self.closed = true; + return Ok(None); + } } } }