Skip to main content

guilder_client_hyperliquid/ws/
manager.rs

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;
17/// Grace period after unsubscribe during which messages for that coin are silently dropped.
18const 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    /// Returns the coin/user symbol for this subscription, used for tracking
99    /// recent unsubscriptions so in-flight messages are silently dropped.
100    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
206/// Special error message sent to all subscribers during graceful shutdown.
207/// Sync loops detect this and exit immediately without reconnecting.
208pub 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    /// Graceful shutdown: close all broadcast channels so subscribers exit
219    /// without attempting to reconnect. The manager task exits after this.
220    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    /// Unsubscribe all subscriptions that match the given coin/user symbol.
255    /// Used by bridges during shutdown to cleanly drain subscriptions before exiting.
256    pub(crate) fn unsubscribe_by_coin(&self, coin: &str) {
257        // We can't inspect the manager's subscription map from here, so we
258        // send a dedicated command. But for now, the bridge knows its own
259        // subscriptions and calls `unsubscribe` per-subscription. This method
260        // is a convenience for the case where we know the coin but not the
261        // exact subscription type (e.g. L2Book vs Trades vs ActiveAssetCtx).
262        // We send Release for all known subscription variants for this coin.
263        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    /// Unsubscribe all user-related subscriptions for a given user address.
274    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    /// Graceful shutdown: close all broadcast channels with a shutdown marker
286    /// so subscribers exit immediately without reconnecting.
287    /// The manager task will process this command, close all subscriptions,
288    /// and then exit.
289    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    // Coins recently unsubscribed — messages for these during grace period are silently dropped.
303    let mut unsubscribed_coins: HashMap<String, Instant> = HashMap::new();
304    let mut backoff_secs = 1_u64;
305
306    loop {
307        // Clean up expired unsubscribe entries on each outer loop iteration.
308        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                // Track the coin(s) this subscription was for so in-flight messages
470                // during shutdown don't trigger spurious warnings.
471                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            // Close all broadcast channels with a shutdown marker so subscribers
487            // exit immediately without attempting to reconnect.
488            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            // Include a truncated payload so venue error frames (channel
542            // "error") carry their rejection reason — without it the warn is
543            // undiagnosable (albatross 09-27: {"channel":"error"} with no body).
544            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                // Check if this message belongs to a recently unsubscribed coin.
566                // In-flight messages during the grace period are silently dropped.
567                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                    // Messages for unsubscribed coins are expected during shutdown.
583                    // Always log at DEBUG to avoid spam - the grace period check
584                    // above is for cleanup, not for deciding whether to warn.
585                    tracing::debug!(message = %message_label(msg), "WS message for unsubscribed symbol dropped");
586                }
587            }
588        }
589    }
590}
591
592/// Extract the coin/user symbol from a message, for unsubscribe-grace lookup.
593fn 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            // User messages don't carry a coin — they're identified by user address.
604            // The unsubscribed_coins map stores user addresses too.
605            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
642/// Marker sent to all subscribers when the manager tears the connection down
643/// for a RECONNECT (transport reset / idle watchdog / stream end). Consumers
644/// MUST treat this as "a reconnect is happening, subscriptions will be
645/// replayed" — NOT a per-stream error; the manager logs the cause once at
646/// warn and the loop replays all subscriptions on the fresh connection.
647/// Mirrors the SHUTDOWN_MARKER pattern in core/src/orderbook/sync.rs.
648pub const WS_RECONNECTING_MARKER: &str = "__ws_reconnecting__";
649
650/// True when an error string is a transport-level reconnect marker rather
651/// than a per-stream failure. Consumers use this to downgrade logging
652/// (debug) vs genuine stream errors (warn).
653pub 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                    // This is expected during shutdown when the WS manager closes
703                    // all broadcast channels. Log at DEBUG level to avoid spam.
704                    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}