1use super::inbound::HyperliquidWsSubscriptionResponse;
2use super::transport::{HyperliquidWs, WsTransport};
3use super::{HyperliquidWsInboundMessage, HyperliquidWsOutboundMessage};
4use std::time::{Duration, Instant};
5use tokio::sync::{broadcast, mpsc, oneshot};
6use tokio::time::{self, Duration as TokioDuration};
7use tracing::warn;
8
9use std::collections::HashMap;
10
11const HEARTBEAT_INTERVAL: TokioDuration = TokioDuration::from_secs(25);
12const IDLE_WATCHDOG_INTERVAL: TokioDuration = TokioDuration::from_secs(5);
13const MAX_IDLE_BEFORE_RECONNECT: TokioDuration = TokioDuration::from_secs(40);
14const BACKOFF_MAX_SECS: u64 = 30;
15const SEND_SPACING: TokioDuration = TokioDuration::from_millis(40);
16const FANOUT_CAPACITY: usize = 1024;
17const UNSUBSCRIBE_GRACE_PERIOD: Duration = Duration::from_secs(5);
19
20#[derive(Clone, Debug, Hash, Eq, PartialEq)]
21pub(crate) enum HyperliquidSubscription {
22 L2Book { coin: String },
23 Trades { coin: String },
24 ActiveAssetCtx { coin: String },
25 UserEvents { user_addr: String },
26 OrderUpdates { user_addr: String },
27 NonFundingLedger { user_addr: String },
28}
29
30impl HyperliquidSubscription {
31 pub(crate) fn label(&self) -> String {
32 match self {
33 Self::L2Book { coin } => format!("l2Book:{coin}"),
34 Self::Trades { coin } => format!("trades:{coin}"),
35 Self::ActiveAssetCtx { coin } => format!("activeAssetCtx:{coin}"),
36 Self::UserEvents { user_addr } => format!("user:{user_addr}"),
37 Self::OrderUpdates { user_addr } => format!("orderUpdates:{user_addr}"),
38 Self::NonFundingLedger { user_addr } => {
39 format!("userNonFundingLedgerUpdates:{user_addr}")
40 }
41 }
42 }
43
44 pub(crate) fn subscribe_message(&self) -> HyperliquidWsOutboundMessage {
45 match self {
46 Self::L2Book { coin } => {
47 HyperliquidWsOutboundMessage::SubscribeL2Book { coin: coin.clone() }
48 }
49 Self::Trades { coin } => {
50 HyperliquidWsOutboundMessage::SubscribeTrades { coin: coin.clone() }
51 }
52 Self::ActiveAssetCtx { coin } => {
53 HyperliquidWsOutboundMessage::SubscribeActiveAssetCtx { coin: coin.clone() }
54 }
55 Self::UserEvents { user_addr } => {
56 HyperliquidWsOutboundMessage::SubscribeUserEvents {
57 user_addr: user_addr.clone(),
58 }
59 }
60 Self::OrderUpdates { user_addr } => {
61 HyperliquidWsOutboundMessage::SubscribeOrderUpdates {
62 user_addr: user_addr.clone(),
63 }
64 }
65 Self::NonFundingLedger { user_addr } => {
66 HyperliquidWsOutboundMessage::SubcribeNonFundingLedger {
67 user_addr: user_addr.clone(),
68 }
69 }
70 }
71 }
72
73 pub(crate) fn unsubscribe_message(&self) -> HyperliquidWsOutboundMessage {
74 HyperliquidWsOutboundMessage::Unsubscribe {
75 subscription: self.subscription_payload(),
76 }
77 }
78
79 fn subscription_payload(&self) -> serde_json::Value {
80 match self {
81 Self::L2Book { coin } => serde_json::json!({"type": "l2Book", "coin": coin}),
82 Self::Trades { coin } => serde_json::json!({"type": "trades", "coin": coin}),
83 Self::ActiveAssetCtx { coin } => {
84 serde_json::json!({"type": "activeAssetCtx", "coin": coin})
85 }
86 Self::UserEvents { user_addr } => {
87 serde_json::json!({"type": "user", "user": user_addr})
88 }
89 Self::OrderUpdates { user_addr } => {
90 serde_json::json!({"type": "orderUpdates", "user": user_addr})
91 }
92 Self::NonFundingLedger { user_addr } => {
93 serde_json::json!({"type": "userNonFundingLedgerUpdates", "user": user_addr})
94 }
95 }
96 }
97
98 fn unsubscribe_symbol(&self) -> Option<String> {
101 match self {
102 Self::L2Book { coin }
103 | Self::Trades { coin }
104 | Self::ActiveAssetCtx { coin } => Some(coin.clone()),
105 Self::UserEvents { user_addr }
106 | Self::OrderUpdates { user_addr }
107 | Self::NonFundingLedger { user_addr } => Some(user_addr.clone()),
108 }
109 }
110
111 fn matches_message(&self, msg: &HyperliquidWsInboundMessage, manager_user: Option<&str>) -> bool {
112 match (self, msg) {
113 (Self::L2Book { coin }, HyperliquidWsInboundMessage::L2Book(book)) => book.coin == *coin,
114 (Self::Trades { coin }, HyperliquidWsInboundMessage::Trades(trades)) => {
115 trades.first().map(|trade| trade.coin.as_str()) == Some(coin.as_str())
116 }
117 (Self::ActiveAssetCtx { coin }, HyperliquidWsInboundMessage::ActiveAssetCtx(ctx)) => {
118 ctx.coin == *coin
119 }
120 (Self::UserEvents { user_addr }, HyperliquidWsInboundMessage::User(_)) => {
121 Some(user_addr.as_str()) == manager_user
122 }
123 (Self::OrderUpdates { user_addr }, HyperliquidWsInboundMessage::OrderUpdates(_)) => {
124 Some(user_addr.as_str()) == manager_user
125 }
126 (
127 Self::NonFundingLedger { user_addr },
128 HyperliquidWsInboundMessage::NonFundingLedger(_),
129 ) => Some(user_addr.as_str()) == manager_user,
130 (expected, HyperliquidWsInboundMessage::SubscriptionResponse(resp)) => {
131 subscription_from_response(resp).as_ref() == Some(expected)
132 }
133 _ => false,
134 }
135 }
136}
137
138fn subscription_from_response(
139 resp: &HyperliquidWsSubscriptionResponse,
140) -> Option<HyperliquidSubscription> {
141 match resp.subscription.get("type")?.as_str()? {
142 "l2Book" => Some(HyperliquidSubscription::L2Book {
143 coin: resp.subscription.get("coin")?.as_str()?.to_string(),
144 }),
145 "trades" => Some(HyperliquidSubscription::Trades {
146 coin: resp.subscription.get("coin")?.as_str()?.to_string(),
147 }),
148 "activeAssetCtx" => Some(HyperliquidSubscription::ActiveAssetCtx {
149 coin: resp.subscription.get("coin")?.as_str()?.to_string(),
150 }),
151 "user" => Some(HyperliquidSubscription::UserEvents {
152 user_addr: resp.subscription.get("user")?.as_str()?.to_string(),
153 }),
154 "orderUpdates" => Some(HyperliquidSubscription::OrderUpdates {
155 user_addr: resp.subscription.get("user")?.as_str()?.to_string(),
156 }),
157 "userNonFundingLedgerUpdates" => Some(HyperliquidSubscription::NonFundingLedger {
158 user_addr: resp.subscription.get("user")?.as_str()?.to_string(),
159 }),
160 _ => None,
161 }
162}
163
164#[derive(Clone)]
165pub(crate) struct WsSendRateLimiter {
166 request_tx: mpsc::UnboundedSender<oneshot::Sender<()>>,
167}
168
169impl WsSendRateLimiter {
170 pub(crate) fn new() -> Self {
171 let (request_tx, mut request_rx) = mpsc::unbounded_channel::<oneshot::Sender<()>>();
172
173 tokio::spawn(async move {
174 while let Some(reply_tx) = request_rx.recv().await {
175 let _ = reply_tx.send(());
176 time::sleep(SEND_SPACING).await;
177 }
178 });
179
180 Self { request_tx }
181 }
182
183 async fn acquire(&self) -> Result<(), String> {
184 let (reply_tx, reply_rx) = oneshot::channel();
185 self.request_tx
186 .send(reply_tx)
187 .map_err(|_| "WS send rate limiter is not running".to_string())?;
188 reply_rx
189 .await
190 .map_err(|_| "WS send rate limiter dropped permit".to_string())
191 }
192}
193
194#[derive(Clone)]
195pub(crate) struct HyperliquidWsManager {
196 cmd_tx: mpsc::UnboundedSender<ManagerCommand>,
197}
198
199struct ManagedSubscription {
200 ref_count: usize,
201 sender: broadcast::Sender<Result<HyperliquidWsInboundMessage, String>>,
202 window_messages: u64,
203 total_messages: u64,
204}
205
206pub const SHUTDOWN_MARKER: &str = "__orderbook_shutting_down__";
209
210enum ManagerCommand {
211 Acquire {
212 subscription: HyperliquidSubscription,
213 response_tx: oneshot::Sender<broadcast::Receiver<Result<HyperliquidWsInboundMessage, String>>>,
214 },
215 Release {
216 subscription: HyperliquidSubscription,
217 },
218 Shutdown,
221}
222
223impl HyperliquidWsManager {
224 pub(crate) fn new(
225 user_addr: Option<String>,
226 send_limiter: WsSendRateLimiter,
227 ws_url: &'static str,
228 ) -> Self {
229 let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
230 tokio::spawn(run_manager(cmd_rx, user_addr, send_limiter, ws_url));
231 Self { cmd_tx }
232 }
233
234 pub(crate) async fn subscribe(
235 &self,
236 subscription: HyperliquidSubscription,
237 ) -> Result<broadcast::Receiver<Result<HyperliquidWsInboundMessage, String>>, String> {
238 let (response_tx, response_rx) = oneshot::channel();
239 self.cmd_tx
240 .send(ManagerCommand::Acquire {
241 subscription,
242 response_tx,
243 })
244 .map_err(|_| "websocket manager is not running".to_string())?;
245 response_rx
246 .await
247 .map_err(|_| "websocket manager dropped subscribe request".to_string())
248 }
249
250 pub(crate) fn unsubscribe(&self, subscription: HyperliquidSubscription) {
251 let _ = self.cmd_tx.send(ManagerCommand::Release { subscription });
252 }
253
254 pub(crate) fn unsubscribe_by_coin(&self, coin: &str) {
257 let variants = vec![
264 HyperliquidSubscription::L2Book { coin: coin.to_string() },
265 HyperliquidSubscription::Trades { coin: coin.to_string() },
266 HyperliquidSubscription::ActiveAssetCtx { coin: coin.to_string() },
267 ];
268 for sub in variants {
269 let _ = self.cmd_tx.send(ManagerCommand::Release { subscription: sub });
270 }
271 }
272
273 pub(crate) fn unsubscribe_user(&self, user_addr: &str) {
275 let variants = vec![
276 HyperliquidSubscription::UserEvents { user_addr: user_addr.to_string() },
277 HyperliquidSubscription::OrderUpdates { user_addr: user_addr.to_string() },
278 HyperliquidSubscription::NonFundingLedger { user_addr: user_addr.to_string() },
279 ];
280 for sub in variants {
281 let _ = self.cmd_tx.send(ManagerCommand::Release { subscription: sub });
282 }
283 }
284
285 pub(crate) fn shutdown(&self) {
290 let _ = self.cmd_tx.send(ManagerCommand::Shutdown);
291 }
292}
293
294async fn run_manager(
295 mut cmd_rx: mpsc::UnboundedReceiver<ManagerCommand>,
296 user_addr: Option<String>,
297 send_limiter: WsSendRateLimiter,
298 ws_url: &'static str,
299) {
300 let mut ws = HyperliquidWs::new(ws_url);
301 let mut subscriptions: HashMap<HyperliquidSubscription, ManagedSubscription> = HashMap::new();
302 let mut unsubscribed_coins: HashMap<String, Instant> = HashMap::new();
304 let mut backoff_secs = 1_u64;
305
306 loop {
307 unsubscribed_coins.retain(|_, since| since.elapsed() < UNSUBSCRIBE_GRACE_PERIOD);
309
310 while subscriptions.is_empty() {
311 let Some(cmd) = cmd_rx.recv().await else {
312 return;
313 };
314 handle_command(
315 cmd,
316 &mut subscriptions,
317 &mut unsubscribed_coins,
318 &mut ws,
319 &send_limiter,
320 )
321 .await;
322 }
323
324 if let Err(err) = ws.connect().await {
325 warn!(error = ?err, "WS connect failed, backing off");
326 time::sleep(Duration::from_secs(backoff_secs)).await;
327 backoff_secs = (backoff_secs * 2).min(BACKOFF_MAX_SECS);
328 continue;
329 }
330
331
332 if let Err(err) = replay_subscriptions(&mut ws, &subscriptions, &send_limiter).await {
333 warn!(error = %err, "WS replay failed, reconnecting");
334 let _ = ws.close().await;
335 time::sleep(Duration::from_secs(backoff_secs)).await;
336 backoff_secs = (backoff_secs * 2).min(BACKOFF_MAX_SECS);
337 continue;
338 }
339
340 backoff_secs = 1;
341 let mut heartbeat = time::interval(HEARTBEAT_INTERVAL);
342 heartbeat.set_missed_tick_behavior(time::MissedTickBehavior::Delay);
343 let mut idle_watchdog = time::interval(IDLE_WATCHDOG_INTERVAL);
344 idle_watchdog.set_missed_tick_behavior(time::MissedTickBehavior::Skip);
345 let mut last_inbound_activity = Instant::now();
346
347 loop {
348 tokio::select! {
349 maybe_cmd = cmd_rx.recv() => {
350 let Some(cmd) = maybe_cmd else {
351 let _ = ws.close().await;
352 return;
353 };
354 handle_command(
355 cmd,
356 &mut subscriptions,
357 &mut unsubscribed_coins,
358 &mut ws,
359 &send_limiter,
360 )
361 .await;
362 if subscriptions.is_empty() {
363 let _ = ws.close().await;
364 break;
365 }
366 }
367 _ = heartbeat.tick() => {
368 if let Err(err) = send_with_limit(&mut ws, HyperliquidWsOutboundMessage::Ping, &send_limiter).await {
369 warn!(error = %err, "WS heartbeat failed, reconnecting");
370 fanout_error(&subscriptions, WS_RECONNECTING_MARKER.to_string());
371 let _ = ws.close().await;
372 break;
373 }
374 }
375 _ = idle_watchdog.tick() => {
376 let idle_for = last_inbound_activity.elapsed();
377 if idle_for >= MAX_IDLE_BEFORE_RECONNECT {
378 warn!(
379 idle_for_ms = idle_for.as_millis(),
380 "WS idle watchdog triggered, reconnecting"
381 );
382 fanout_error(&subscriptions, WS_RECONNECTING_MARKER.to_string());
383 let _ = ws.close().await;
384 break;
385 }
386 }
387 inbound = ws.recv() => {
388 match inbound {
389 Some(Ok(HyperliquidWsInboundMessage::Pong)) => {
390 last_inbound_activity = Instant::now();
391 }
392 Some(Ok(msg)) => {
393 last_inbound_activity = Instant::now();
394 dispatch_message(&mut subscriptions, &unsubscribed_coins, &msg, user_addr.as_deref())
395 }
396 Some(Err(err)) => {
397 warn!(error = %err, "WS recv error, reconnecting");
398 fanout_error(&subscriptions, WS_RECONNECTING_MARKER.to_string());
399 let _ = ws.close().await;
400 break;
401 }
402 None => {
403 warn!("WS stream ended, reconnecting");
404 fanout_error(&subscriptions, WS_RECONNECTING_MARKER.to_string());
405 let _ = ws.close().await;
406 break;
407 }
408 }
409 }
410 }
411 }
412 }
413}
414
415async fn handle_command(
416 cmd: ManagerCommand,
417 subscriptions: &mut HashMap<HyperliquidSubscription, ManagedSubscription>,
418 unsubscribed_coins: &mut HashMap<String, Instant>,
419 ws: &mut HyperliquidWs,
420 send_limiter: &WsSendRateLimiter,
421) {
422 match cmd {
423 ManagerCommand::Acquire {
424 subscription,
425 response_tx,
426 } => {
427 if let Some(managed) = subscriptions.get_mut(&subscription) {
428 managed.ref_count += 1;
429 let _ = response_tx.send(managed.sender.subscribe());
430 return;
431 }
432
433 let (sender, receiver) = broadcast::channel(FANOUT_CAPACITY);
434 subscriptions.insert(
435 subscription.clone(),
436 ManagedSubscription {
437 ref_count: 1,
438 sender,
439 window_messages: 0,
440 total_messages: 0,
441 },
442 );
443 let _ = response_tx.send(receiver);
444
445 if ws.is_connected() {
446 if let Err(err) =
447 send_with_limit(ws, subscription.subscribe_message(), send_limiter).await
448 {
449 if let Some(managed) = subscriptions.get(&subscription) {
450 let _ = managed.sender.send(Err(format!(
451 "failed to subscribe {:?}: {err}",
452 subscription
453 )));
454 }
455 }
456 }
457 }
458 ManagerCommand::Release { subscription } => {
459 let remove = match subscriptions.get_mut(&subscription) {
460 Some(managed) if managed.ref_count > 1 => {
461 managed.ref_count -= 1;
462 false
463 }
464 Some(_) => true,
465 None => false,
466 };
467
468 if remove {
469 if let Some(coin) = subscription.unsubscribe_symbol() {
472 unsubscribed_coins.insert(coin, Instant::now());
473 }
474
475 subscriptions.remove(&subscription);
476 if ws.is_connected() {
477 if let Err(err) =
478 send_with_limit(ws, subscription.unsubscribe_message(), send_limiter).await
479 {
480 warn!(error = %err, subscription = ?subscription, "WS unsubscribe failed");
481 }
482 }
483 }
484 }
485 ManagerCommand::Shutdown => {
486 for (sub, managed) in subscriptions.drain() {
489 if let Some(coin) = sub.unsubscribe_symbol() {
490 unsubscribed_coins.insert(coin, Instant::now());
491 }
492 let _ = managed.sender.send(Err(SHUTDOWN_MARKER.to_string()));
493 }
494 if ws.is_connected() {
495 let _ = ws.close().await;
496 }
497 }
498 }
499}
500
501async fn replay_subscriptions(
502 ws: &mut HyperliquidWs,
503 subscriptions: &HashMap<HyperliquidSubscription, ManagedSubscription>,
504 send_limiter: &WsSendRateLimiter,
505) -> Result<(), String> {
506 for subscription in subscriptions.keys() {
507 send_with_limit(ws, subscription.subscribe_message(), send_limiter).await?;
508 }
509 Ok(())
510}
511
512async fn send_with_limit(
513 ws: &mut HyperliquidWs,
514 msg: HyperliquidWsOutboundMessage,
515 send_limiter: &WsSendRateLimiter,
516) -> Result<(), String> {
517 send_limiter.acquire().await?;
518 ws.send(msg).await.map_err(|e| e.to_string())
519}
520
521fn dispatch_message(
522 subscriptions: &mut HashMap<HyperliquidSubscription, ManagedSubscription>,
523 unsubscribed_coins: &HashMap<String, Instant>,
524 msg: &HyperliquidWsInboundMessage,
525 manager_user: Option<&str>,
526) {
527 match msg {
528 HyperliquidWsInboundMessage::SubscriptionResponse(resp) => {
529 if let Some(subscription) = subscription_from_response(resp) {
530 if let Some(managed) = subscriptions.get(&subscription) {
531 if !resp.success.unwrap_or(true) {
532 let _ = managed.sender.send(Err(format!(
533 "subscription rejected: {:?}",
534 subscription
535 )));
536 }
537 }
538 }
539 }
540 HyperliquidWsInboundMessage::Unknown { channel, data } => {
541 let payload = data.to_string();
545 let payload = if payload.len() > 300 {
546 let mut cut = String::from(&payload[..300]);
547 cut.push_str("…");
548 cut
549 } else {
550 payload
551 };
552 warn!(channel = %channel, payload = %payload, "WS message ignored as unknown");
553 }
554 _ => {
555 let mut matched = 0usize;
556 for (subscription, managed) in subscriptions.iter_mut() {
557 if subscription.matches_message(msg, manager_user) {
558 matched += 1;
559 managed.window_messages += 1;
560 managed.total_messages += 1;
561 let _ = managed.sender.send(Ok(msg.clone()));
562 }
563 }
564 if matched == 0 {
565 let sym = message_symbol(msg);
568 let recently_unsubscribed = sym
569 .as_ref()
570 .map(|s| {
571 unsubscribed_coins
572 .get(s)
573 .is_some_and(|since| since.elapsed() < UNSUBSCRIBE_GRACE_PERIOD)
574 })
575 .unwrap_or(false);
576 if recently_unsubscribed {
577 tracing::debug!(
578 message = %message_label(msg),
579 "WS message dropped for recently unsubscribed symbol"
580 );
581 } else {
582 tracing::debug!(message = %message_label(msg), "WS message for unsubscribed symbol dropped");
586 }
587 }
588 }
589 }
590}
591
592fn message_symbol(msg: &HyperliquidWsInboundMessage) -> Option<String> {
594 match msg {
595 HyperliquidWsInboundMessage::L2Book(book) => Some(book.coin.clone()),
596 HyperliquidWsInboundMessage::ActiveAssetCtx(ctx) => Some(ctx.coin.clone()),
597 HyperliquidWsInboundMessage::Trades(trades) => {
598 trades.first().map(|t| t.coin.clone())
599 }
600 HyperliquidWsInboundMessage::User(_)
601 | HyperliquidWsInboundMessage::OrderUpdates(_)
602 | HyperliquidWsInboundMessage::NonFundingLedger(_) => {
603 None
606 }
607 _ => None,
608 }
609}
610
611pub(crate) fn message_label(msg: &HyperliquidWsInboundMessage) -> String {
612 match msg {
613 HyperliquidWsInboundMessage::Pong => "pong".to_string(),
614 HyperliquidWsInboundMessage::L2Book(book) => format!("l2Book:{}", book.coin),
615 HyperliquidWsInboundMessage::ActiveAssetCtx(ctx) => format!("activeAssetCtx:{}", ctx.coin),
616 HyperliquidWsInboundMessage::Trades(trades) => format!(
617 "trades:{}",
618 trades.first().map(|trade| trade.coin.as_str()).unwrap_or("<empty>")
619 ),
620 HyperliquidWsInboundMessage::User(_) => "user".to_string(),
621 HyperliquidWsInboundMessage::OrderUpdates(_) => "orderUpdates".to_string(),
622 HyperliquidWsInboundMessage::NonFundingLedger(_) => "userNonFundingLedgerUpdates".to_string(),
623 HyperliquidWsInboundMessage::SubscriptionResponse(resp) => format!(
624 "subscriptionResponse:{}",
625 subscription_from_response(resp)
626 .map(|sub| sub.label())
627 .unwrap_or_else(|| "<unknown>".to_string())
628 ),
629 HyperliquidWsInboundMessage::Unknown { channel, .. } => format!("unknown:{channel}"),
630 }
631}
632
633fn fanout_error(
634 subscriptions: &HashMap<HyperliquidSubscription, ManagedSubscription>,
635 error: String,
636) {
637 for managed in subscriptions.values() {
638 let _ = managed.sender.send(Err(error.clone()));
639 }
640}
641
642pub const WS_RECONNECTING_MARKER: &str = "__ws_reconnecting__";
649
650pub fn is_reconnect_marker(e: &str) -> bool {
654 e == WS_RECONNECTING_MARKER
655 || e == "websocket connection closed"
656 || e == "websocket stream ended"
657 || e.starts_with("websocket idle watchdog")
658 || e.starts_with("websocket heartbeat failed")
659 || e == "WS recv error: websocket connection closed"
660}
661
662pub(crate) fn managed_stream<T, F>(
663 manager: HyperliquidWsManager,
664 subscription: HyperliquidSubscription,
665 parse: F,
666) -> impl futures_core::Stream<Item = Result<T, String>> + Send + 'static
667where
668 T: Send + 'static,
669 F: Fn(HyperliquidWsInboundMessage) -> Vec<Result<T, String>> + Send + Sync + 'static,
670{
671 async_stream::stream! {
672 let mut receiver = match manager.subscribe(subscription.clone()).await {
673 Ok(receiver) => receiver,
674 Err(err) => {
675 yield Err(err);
676 return;
677 }
678 };
679
680 struct ReleaseOnDrop {
681 manager: HyperliquidWsManager,
682 subscription: HyperliquidSubscription,
683 }
684
685 impl Drop for ReleaseOnDrop {
686 fn drop(&mut self) {
687 self.manager.unsubscribe(self.subscription.clone());
688 }
689 }
690
691 let _release_on_drop = ReleaseOnDrop { manager, subscription };
692
693 loop {
694 match receiver.recv().await {
695 Ok(Ok(message)) => {
696 for item in parse(message) {
697 yield item;
698 }
699 }
700 Ok(Err(err)) => yield Err(err),
701 Err(broadcast::error::RecvError::Closed) => {
702 tracing::debug!(subscription = %_release_on_drop.subscription.label(), "managed stream receiver closed");
705 yield Err("websocket subscription closed".to_string());
706 return;
707 }
708 Err(broadcast::error::RecvError::Lagged(skipped)) => {
709 warn!(
710 subscription = %_release_on_drop.subscription.label(),
711 skipped = skipped,
712 "managed stream receiver lagged"
713 );
714 yield Err(format!("websocket subscription lagged by {skipped} messages"));
715 }
716 }
717 }
718 }
719}