Skip to main content

polyester/services/
market_data.rs

1use super::ServiceContext;
2use super::unary;
3use crate::codecs::decode::{
4    candles_columns_from_proto, candles_from_proto, depth_enum_for_levels,
5    market_overview_list_from_proto, market_trades_from_proto, orderbook_from_proto,
6    spot_config_from_proto,
7};
8use crate::connect::marketdata::v1::MarketDataServiceClient;
9use crate::connect::marketoverview::v1::MarketOverviewServiceClient;
10use crate::connect::orderbook::v1::OrderbookServiceClient;
11use crate::errors::{Error, Result};
12use crate::models::{
13    Candle, CandlesResult, GetCandlesOpts, GetTradesOpts, MarketOverviewList, MarketTradesResult,
14    OrderbookData, SpotConfig,
15};
16use crate::models::{MarketOverviewEntry, MarketTrade, OrderBookDeltaUpdate};
17use crate::proto::marketdata::v1::{
18    GetCandlesColumnsRequest, GetCandlesRequest, GetSpotConfigRequest, GetTradesRequest, Timeframe,
19};
20use crate::proto::marketoverview::v1::ListMarketOverviewRequest;
21use crate::proto::orderbook::v1::GetOrderBookRequest;
22use buffa_types::google::protobuf::Timestamp;
23
24#[derive(Clone)]
25pub struct MarketDataService {
26    ctx: ServiceContext,
27}
28
29impl MarketDataService {
30    pub fn new(ctx: ServiceContext) -> Self {
31        Self { ctx }
32    }
33
34    fn client(&self) -> MarketDataServiceClient<crate::transport::SharedTransport> {
35        MarketDataServiceClient::new(
36            self.ctx.factory.transport(),
37            self.ctx.factory.connect_config(),
38        )
39    }
40
41    fn resolve_symbol_id(
42        &self,
43        symbol: Option<&str>,
44        symbol_id: Option<u32>,
45        label: &str,
46    ) -> Result<u32> {
47        if let Some(id) = symbol_id.filter(|id| *id != 0) {
48            return Ok(id);
49        }
50        let Some(symbol) = symbol.filter(|s| !s.is_empty()) else {
51            return Err(Error::validation(format!(
52                "{label} requires symbol or symbol_id"
53            )));
54        };
55        self.ctx
56            .catalogs
57            .symbol_id_for_symbol(symbol)
58            .ok_or_else(|| {
59                Error::validation(format!(
60                    "unknown symbol {symbol}; call hydrate_catalogs / get_spot_config first"
61                ))
62            })
63    }
64
65    fn require_quantity_scale(&self, symbol_id: u32, label: &str) -> Result<u32> {
66        self.ctx
67            .catalogs
68            .base_quantity_scale_for_symbol_id(symbol_id)
69            .ok_or_else(|| {
70                Error::validation(format!(
71                    "{label} requires a catalog quantity scale for symbol_id {symbol_id}; \
72                     wait_for_catalogs or hydrate the spot catalog first"
73                ))
74            })
75    }
76
77    fn timestamp_field(secs: Option<i64>) -> buffa::MessageField<Timestamp> {
78        match secs {
79            Some(seconds) => Timestamp {
80                seconds,
81                nanos: 0,
82                ..Default::default()
83            }
84            .into(),
85            None => buffa::MessageField::none(),
86        }
87    }
88
89    /// Spot pair catalog (Go `SpotConfig` raw-map escape hatch).
90    pub async fn get_spot_config(&self) -> Result<SpotConfig> {
91        let resp = unary::await_public(
92            self.client()
93                .get_spot_config(GetSpotConfigRequest::default()),
94        )
95        .await?
96        .into_owned();
97        Ok(spot_config_from_proto(&resp))
98    }
99
100    /// Recent public trades for a symbol (resolves `symbol_id` via catalogs after hydrate).
101    pub async fn get_trades(&self, symbol: &str, limit: Option<u32>) -> Result<MarketTradesResult> {
102        self.get_trades_with(GetTradesOpts {
103            symbol: Some(symbol.to_owned()),
104            limit,
105            ..Default::default()
106        })
107        .await
108    }
109
110    pub async fn get_trades_with(&self, opts: GetTradesOpts) -> Result<MarketTradesResult> {
111        let symbol_id =
112            self.resolve_symbol_id(opts.symbol.as_deref(), opts.symbol_id, "get_trades")?;
113        let quantity_scale = self.require_quantity_scale(symbol_id, "get_trades")?;
114        let req = GetTradesRequest {
115            symbol_id,
116            limit: opts.limit.unwrap_or(0),
117            start_time: Self::timestamp_field(opts.start),
118            end_time: Self::timestamp_field(opts.end),
119            page_token: opts.page_token.unwrap_or_default(),
120            ..Default::default()
121        };
122        let resp = unary::await_public(self.client().get_trades(req))
123            .await?
124            .into_owned();
125        Ok(market_trades_from_proto(&resp, quantity_scale))
126    }
127
128    /// Candle series for a symbol, ordered newest-first by `ts_sec`.
129    ///
130    /// `interval` accepts values like `"1m"`, `"MIN_1"`, `"5m"`. When
131    /// incomplete candles are requested, the current open candle is prepended.
132    pub async fn get_candles(
133        &self,
134        symbol: &str,
135        interval: &str,
136        limit: Option<u32>,
137    ) -> Result<CandlesResult> {
138        self.get_candles_with(GetCandlesOpts {
139            symbol: Some(symbol.to_owned()),
140            timeframe: interval.to_owned(),
141            limit,
142            ..Default::default()
143        })
144        .await
145    }
146
147    pub async fn get_candles_with(&self, opts: GetCandlesOpts) -> Result<CandlesResult> {
148        let (req, volume_scale) = self.build_candles_request(&opts)?;
149        let resp = unary::await_public(self.client().get_candles(req))
150            .await?
151            .into_owned();
152        candles_from_proto(&resp, volume_scale)
153    }
154
155    /// Latest candle for a symbol/timeframe, or `None` when the market has no rows.
156    pub async fn get_current_candle(
157        &self,
158        symbol: &str,
159        timeframe: &str,
160    ) -> Result<Option<Candle>> {
161        let result = self
162            .get_candles_with(GetCandlesOpts {
163                symbol: Some(symbol.to_owned()),
164                timeframe: timeframe.to_owned(),
165                limit: Some(1),
166                include_incomplete: true,
167                ..Default::default()
168            })
169            .await?;
170        Ok(newest_candle(result))
171    }
172
173    /// Columnar OHLCV candles decoded into row-oriented [`CandlesResult`].
174    pub async fn get_candles_columns(&self, opts: GetCandlesOpts) -> Result<CandlesResult> {
175        let (base, volume_scale) = self.build_candles_request(&opts)?;
176        let req = GetCandlesColumnsRequest {
177            symbol_id: base.symbol_id,
178            timeframe: base.timeframe,
179            limit: base.limit,
180            start_time: base.start_time,
181            end_time: base.end_time,
182            include_incomplete: base.include_incomplete,
183            include_reference: base.include_reference,
184            page_token: base.page_token,
185            ..Default::default()
186        };
187        let resp = unary::await_public(self.client().get_candles_columns(req))
188            .await?
189            .into_owned();
190        candles_columns_from_proto(&resp, volume_scale)
191    }
192
193    fn build_candles_request(&self, opts: &GetCandlesOpts) -> Result<(GetCandlesRequest, u32)> {
194        let symbol_id =
195            self.resolve_symbol_id(opts.symbol.as_deref(), opts.symbol_id, "get_candles")?;
196        let timeframe_label = if opts.timeframe.is_empty() {
197            "1m"
198        } else {
199            opts.timeframe.as_str()
200        };
201        let timeframe = parse_timeframe(timeframe_label)?;
202        let volume_scale = self.require_quantity_scale(symbol_id, "get_candles")?;
203        let req = GetCandlesRequest {
204            symbol_id,
205            timeframe: timeframe.into(),
206            limit: opts.limit.unwrap_or(0),
207            start_time: Self::timestamp_field(opts.start),
208            end_time: Self::timestamp_field(opts.end),
209            include_incomplete: opts.include_incomplete,
210            page_token: opts.page_token.clone().unwrap_or_default(),
211            ..Default::default()
212        };
213        Ok((req, volume_scale))
214    }
215
216    /// Subscribe to public spot trades for a symbol (requires `realtime` feature + hydrated catalogs).
217    pub async fn subscribe_trades(
218        &self,
219        symbol: &str,
220    ) -> Result<crate::realtime::TypedSubscription<MarketTrade>> {
221        let symbol_id = self
222            .ctx
223            .catalogs
224            .symbol_id_for_symbol(symbol)
225            .ok_or_else(|| {
226                Error::validation(format!(
227                    "unknown symbol {symbol}; call hydrate_catalogs / get_spot_config first"
228                ))
229            })?;
230        let quantity_scale = self.require_quantity_scale(symbol_id, "subscribe_trades")?;
231        let channel = format!("public:spot:market:trades:{symbol_id}:proto");
232        self.ctx
233            .realtime
234            .subscribe_proto(
235                &channel,
236                crate::codecs::decode::market_trade_from_bytes(quantity_scale),
237            )
238            .await
239    }
240
241    /// Subscribe to public candle updates (requires `realtime` feature + hydrated catalogs).
242    pub async fn subscribe_candles(
243        &self,
244        symbol: &str,
245        timeframe: &str,
246    ) -> Result<crate::realtime::TypedSubscription<Candle>> {
247        let symbol_id = self
248            .ctx
249            .catalogs
250            .symbol_id_for_symbol(symbol)
251            .ok_or_else(|| {
252                Error::validation(format!(
253                    "unknown symbol {symbol}; call hydrate_catalogs / get_spot_config first"
254                ))
255            })?;
256        // Live Centrifugo channels use human labels (`1m`), not REST enum names (`MIN_1`).
257        let resolved = parse_timeframe(timeframe)?;
258        let channel_tf = crate::codecs::decode::timeframe_label(resolved);
259        if channel_tf.is_empty() {
260            return Err(Error::validation(format!(
261                "unsupported candle interval {timeframe:?}"
262            )));
263        }
264        let volume_scale = self.require_quantity_scale(symbol_id, "subscribe_candles")?;
265        let channel = format!("public:spot:market:candles:{channel_tf}:{symbol_id}:proto");
266        let decode = crate::codecs::decode::candle_point_from_bytes(
267            symbol_id,
268            channel_tf.to_owned(),
269            volume_scale,
270        );
271        self.ctx.realtime.subscribe_proto(&channel, decode).await
272    }
273
274    /// Channel segment for candle subscriptions after alias normalization.
275    #[cfg(test)]
276    pub(crate) fn candle_channel_timeframe(timeframe: &str) -> Result<&'static str> {
277        let resolved = parse_timeframe(timeframe)?;
278        let label = crate::codecs::decode::timeframe_label(resolved);
279        if label.is_empty() {
280            return Err(Error::validation(format!(
281                "unsupported candle interval {timeframe:?}"
282            )));
283        }
284        Ok(label)
285    }
286}
287
288fn parse_timeframe(interval: &str) -> Result<Timeframe> {
289    let key = interval.trim().to_ascii_lowercase().replace('_', "");
290    let tf = match key.as_str() {
291        "1s" | "sec1" => Timeframe::Sec1,
292        "1m" | "min1" => Timeframe::Min1,
293        "5m" | "min5" => Timeframe::Min5,
294        "15m" | "min15" => Timeframe::Min15,
295        "30m" | "min30" => Timeframe::Min30,
296        "1h" | "hour1" => Timeframe::Hour1,
297        "4h" | "hour4" => Timeframe::Hour4,
298        "12h" | "hour12" => Timeframe::Hour12,
299        "1d" | "day1" => Timeframe::Day1,
300        "1w" | "week1" => Timeframe::Week1,
301        "1mo" | "month1" => Timeframe::Month1,
302        _ => {
303            return Err(Error::validation(format!(
304                "unsupported candle interval {interval:?}"
305            )));
306        }
307    };
308    Ok(tf)
309}
310
311fn newest_candle(result: CandlesResult) -> Option<Candle> {
312    result.candles.into_iter().next()
313}
314
315/// Options for [`MarketOverviewService::list`].
316#[derive(Debug, Clone, Default)]
317pub struct ListMarketOverviewOptions {
318    pub symbols: Option<Vec<String>>,
319    pub limit: Option<u32>,
320    pub include_sparklines: bool,
321}
322
323impl From<Option<u32>> for ListMarketOverviewOptions {
324    fn from(limit: Option<u32>) -> Self {
325        Self {
326            limit,
327            ..Default::default()
328        }
329    }
330}
331
332/// Options for managed market-overview subscriptions.
333#[derive(Debug, Clone, Default)]
334pub struct MarketOverviewCreateSubscriptionOptions {
335    pub symbols: Option<Vec<String>>,
336    pub limit: Option<u32>,
337    pub include_sparklines: bool,
338}
339
340#[derive(Clone)]
341pub struct MarketOverviewService {
342    ctx: ServiceContext,
343}
344
345impl MarketOverviewService {
346    pub fn new(ctx: ServiceContext) -> Self {
347        Self { ctx }
348    }
349
350    pub async fn list(
351        &self,
352        opts: impl Into<ListMarketOverviewOptions>,
353    ) -> Result<MarketOverviewList> {
354        let opts = opts.into();
355        let req = ListMarketOverviewRequest {
356            symbols: opts.symbols.unwrap_or_default(),
357            limit: opts.limit.unwrap_or_default(),
358            include_sparklines: opts.include_sparklines,
359            ..Default::default()
360        };
361        let client = MarketOverviewServiceClient::new(
362            self.ctx.factory.transport(),
363            self.ctx.factory.connect_config(),
364        );
365        let resp = unary::await_public(client.list_market_overview(req))
366            .await?
367            .into_owned();
368        Ok(market_overview_list_from_proto(&resp))
369    }
370
371    /// Subscribe to public market overview batches.
372    pub async fn subscribe(
373        &self,
374    ) -> Result<crate::realtime::TypedSubscription<MarketOverviewList>> {
375        self.ctx
376            .realtime
377            .subscribe_proto(
378                "public:spot:market_overview:updates:proto",
379                crate::codecs::decode::market_overview_batch_from_bytes,
380            )
381            .await
382    }
383
384    /// Snapshot-then-stream merged overview rows.
385    pub async fn create_subscription(
386        &self,
387        opts: MarketOverviewCreateSubscriptionOptions,
388    ) -> Result<crate::marketoverview::Subscription> {
389        use crate::realtime::{SnapshotThenStream, SnapshotThenStreamConfig};
390        use std::collections::HashMap;
391        use std::sync::atomic::AtomicBool;
392        use std::sync::{Arc, Mutex};
393        use tokio::sync::mpsc;
394
395        let limit = opts.limit.filter(|n| *n > 0).unwrap_or(50);
396        let symbols = opts.symbols.clone();
397        let include_sparklines = opts.include_sparklines;
398        let channel = "public:spot:market_overview:updates:proto".to_owned();
399
400        let by_symbol_id: Arc<Mutex<HashMap<u32, MarketOverviewEntry>>> =
401            Arc::new(Mutex::new(HashMap::new()));
402        let closed = Arc::new(AtomicBool::new(false));
403        let last_error: Arc<Mutex<Option<crate::Error>>> = Arc::new(Mutex::new(None));
404        let (tx, rx) = mpsc::channel::<Vec<MarketOverviewEntry>>(50);
405        let tx_slot: Arc<Mutex<Option<mpsc::Sender<Vec<MarketOverviewEntry>>>>> =
406            Arc::new(Mutex::new(Some(tx)));
407        let stream_slot: Arc<
408            Mutex<Option<SnapshotThenStream<MarketOverviewList, MarketOverviewList>>>,
409        > = Arc::new(Mutex::new(None));
410
411        let emit = {
412            let by_symbol_id = by_symbol_id.clone();
413            let closed = closed.clone();
414            let last_error = last_error.clone();
415            let tx_slot = tx_slot.clone();
416            let stream_slot = stream_slot.clone();
417            Arc::new(move || {
418                if closed.load(std::sync::atomic::Ordering::SeqCst) {
419                    return;
420                }
421                let rows: Vec<MarketOverviewEntry> =
422                    crate::realtime::lock_unpoisoned(&by_symbol_id)
423                        .values()
424                        .cloned()
425                        .collect();
426                let Some(tx) = crate::realtime::lock_unpoisoned(&tx_slot).as_ref().cloned() else {
427                    return;
428                };
429                if !crate::realtime::try_enqueue(
430                    &tx,
431                    rows,
432                    &closed,
433                    &last_error,
434                    "market overview subscription queue full; consumer too slow",
435                ) && let Some(err @ Error::QueueOverflow(_)) =
436                    crate::realtime::lock_unpoisoned(&last_error).clone()
437                {
438                    let _ = crate::realtime::lock_unpoisoned(&tx_slot).take();
439                    if let Some(stream) = crate::realtime::lock_unpoisoned(&stream_slot).as_ref() {
440                        stream.fail(err);
441                    }
442                }
443            }) as Arc<dyn Fn() + Send + Sync>
444        };
445
446        let apply_rows = {
447            let by_symbol_id = by_symbol_id.clone();
448            Arc::new(move |rows: Vec<MarketOverviewEntry>| {
449                let mut map = crate::realtime::lock_unpoisoned(&by_symbol_id);
450                for row in rows {
451                    map.insert(row.symbol_id, row);
452                }
453            }) as Arc<dyn Fn(Vec<MarketOverviewEntry>) + Send + Sync>
454        };
455
456        let svc = self.clone();
457        let fetch_symbols = symbols.clone();
458        let stream = SnapshotThenStream::new(SnapshotThenStreamConfig {
459            client: self.ctx.realtime.clone(),
460            channel,
461            decode: Arc::new(crate::codecs::decode::market_overview_batch_from_bytes),
462            fetch_snapshot: Arc::new(move || {
463                let svc = svc.clone();
464                let symbols = fetch_symbols.clone();
465                Box::pin(async move {
466                    svc.list(ListMarketOverviewOptions {
467                        symbols,
468                        limit: Some(limit),
469                        include_sparklines,
470                    })
471                    .await
472                })
473            }),
474            read_publication: Arc::new(|batch: MarketOverviewList| vec![batch]),
475            apply_snapshot: {
476                let apply_rows = apply_rows.clone();
477                let emit = emit.clone();
478                let by_symbol_id = by_symbol_id.clone();
479                Arc::new(
480                    move |snapshot: MarketOverviewList, buffered: Vec<MarketOverviewList>| {
481                        crate::realtime::lock_unpoisoned(&by_symbol_id).clear();
482                        apply_rows(snapshot.markets);
483                        for batch in buffered {
484                            apply_rows(batch.markets);
485                        }
486                        emit();
487                    },
488                )
489            },
490            apply_live_publications: {
491                let apply_rows = apply_rows.clone();
492                let emit = emit.clone();
493                Arc::new(move |batches: Vec<MarketOverviewList>| {
494                    for batch in batches {
495                        apply_rows(batch.markets);
496                    }
497                    emit();
498                })
499            },
500            max_buffered: 2000,
501            on_reconnect: None,
502            on_snapshot_refresh: None,
503            on_error: None,
504        });
505        *crate::realtime::lock_unpoisoned(&stream_slot) = Some(stream.clone());
506
507        let subscription = crate::marketoverview::Subscription::new(
508            rx,
509            stream.clone(),
510            closed,
511            last_error,
512            tx_slot,
513        );
514        if let Err(err) = stream.start().await {
515            subscription.close();
516            return Err(err);
517        }
518        Ok(subscription)
519    }
520}
521
522/// Options for managed orderbook subscriptions.
523#[derive(Debug, Clone, Default)]
524pub struct CreateSubscriptionOptions {
525    pub symbol: String,
526    pub symbol_id: Option<u32>,
527    pub depth: Option<u32>,
528    pub bucket: Option<String>,
529}
530
531#[derive(Clone)]
532pub struct OrderbookService {
533    ctx: ServiceContext,
534}
535
536impl OrderbookService {
537    pub fn new(ctx: ServiceContext) -> Self {
538        Self { ctx }
539    }
540
541    /// Snapshot orderbook for `symbol`. `depth` maps like Go (`None` / `0` → depth 5 bucket).
542    pub async fn get(&self, symbol: &str, depth: Option<u32>) -> Result<OrderbookData> {
543        let depth_levels = depth.unwrap_or(0);
544        let depth_enum = if depth_levels == 0 {
545            crate::proto::orderbook::v1::Depth::DepthUnspecified
546        } else {
547            depth_enum_for_levels(depth_levels)
548        };
549        // Record the requested depth for the model; unspecified defaults to 50 server-side.
550        let reported_depth = if depth_levels == 0 { 50 } else { depth_levels };
551        let req = GetOrderBookRequest {
552            symbol: symbol.to_owned(),
553            depth: depth_enum.into(),
554            ..Default::default()
555        };
556        let quantity_scale = self
557            .ctx
558            .catalogs
559            .base_quantity_scale_for_symbol(symbol)
560            .ok_or_else(|| {
561                Error::validation(format!(
562                    "orderbook get requires a catalog quantity scale for {symbol}; \
563                     wait_for_catalogs or hydrate the spot catalog first"
564                ))
565            })?;
566        let client = OrderbookServiceClient::new(
567            self.ctx.factory.transport(),
568            self.ctx.factory.connect_config(),
569        );
570        let resp = unary::await_public(client.get_order_book(req))
571            .await?
572            .into_owned();
573        orderbook_from_proto(&resp, symbol, reported_depth, quantity_scale)
574    }
575
576    /// Subscribe to public orderbook delta updates (requires `realtime` feature).
577    pub async fn subscribe_deltas(
578        &self,
579        symbol_id: u32,
580        depth: Option<u32>,
581    ) -> Result<crate::realtime::TypedSubscription<OrderBookDeltaUpdate>> {
582        let ws_depth = depth.unwrap_or(50).clamp(1, 500);
583        let channel = format!("public:spot:orderbook:deltas:depth:{ws_depth}:{symbol_id}:proto");
584        self.ctx
585            .realtime
586            .subscribe_proto(&channel, crate::codecs::decode::orderbook_delta_from_bytes)
587            .await
588    }
589
590    /// Snapshot-then-stream orderbook merging (requires `realtime` feature).
591    pub async fn create_subscription(
592        &self,
593        opts: CreateSubscriptionOptions,
594    ) -> Result<crate::orderbook::Subscription> {
595        use crate::orderbook::{
596            BookSide, apply_delta, build_orderbook_data, levels_from_orderbook_side,
597            parse_bucket_ticks,
598        };
599        use crate::realtime::{SnapshotThenStream, SnapshotThenStreamConfig};
600        use std::sync::atomic::AtomicBool;
601        use std::sync::{Arc, Mutex};
602        use tokio::sync::mpsc;
603
604        let symbol = opts.symbol;
605        let depth = opts.depth.unwrap_or(50);
606        let ws_depth = depth.clamp(1, 500);
607        let resolved_symbol_id = opts
608            .symbol_id
609            .or_else(|| self.ctx.catalogs.symbol_id_for_symbol(&symbol));
610        let Some(symbol_id) = resolved_symbol_id.filter(|id| *id != 0) else {
611            return Err(Error::validation(format!(
612                "symbol_id is required for orderbook subscriptions ({symbol:?})"
613            )));
614        };
615        let channel = format!("public:spot:orderbook:deltas:depth:{ws_depth}:{symbol_id}:proto");
616        let quantity_scale = self
617            .ctx
618            .catalogs
619            .base_quantity_scale_for_symbol(&symbol)
620            .ok_or_else(|| {
621                Error::validation(format!(
622                    "orderbook subscription requires a catalog quantity scale for {symbol}; \
623                     wait_for_catalogs or hydrate the spot catalog first"
624                ))
625            })?;
626        let bucket_ticks = Arc::new(Mutex::new(parse_bucket_ticks(
627            opts.bucket.as_deref().unwrap_or(""),
628        )?));
629
630        let state = Arc::new(Mutex::new(BookState {
631            bids: BookSide::new(),
632            asks: BookSide::new(),
633            book_seq: 0,
634        }));
635        let closed = Arc::new(AtomicBool::new(false));
636        let last_error: Arc<Mutex<Option<crate::Error>>> = Arc::new(Mutex::new(None));
637        let (tx, rx) = mpsc::channel::<OrderbookData>(200);
638        // Dropping this slot closes the consumer channel so recv() cannot hang
639        // forever after close/error.
640        let tx_slot: Arc<Mutex<Option<mpsc::Sender<OrderbookData>>>> =
641            Arc::new(Mutex::new(Some(tx)));
642        // Filled after stream construction so synchronous emit callbacks can
643        // terminate the underlying socket on local queue overflow.
644        let stream_slot: Arc<
645            Mutex<Option<SnapshotThenStream<OrderbookData, OrderBookDeltaUpdate>>>,
646        > = Arc::new(Mutex::new(None));
647
648        let emit = {
649            let state = state.clone();
650            let bucket_ticks = bucket_ticks.clone();
651            let closed = closed.clone();
652            let last_error = last_error.clone();
653            let tx_slot = tx_slot.clone();
654            let stream_slot = stream_slot.clone();
655            let symbol = symbol.clone();
656            Arc::new(move || {
657                if closed.load(std::sync::atomic::Ordering::SeqCst) {
658                    return;
659                }
660                let (bids, asks, book_seq) = {
661                    let s = crate::realtime::lock_unpoisoned(&state);
662                    (s.bids.clone(), s.asks.clone(), s.book_seq)
663                };
664                let ticks = *crate::realtime::lock_unpoisoned(&bucket_ticks);
665                let data = match build_orderbook_data(
666                    &symbol,
667                    ws_depth,
668                    book_seq,
669                    &bids,
670                    &asks,
671                    ticks,
672                    quantity_scale,
673                ) {
674                    Ok(data) => data,
675                    Err(err) => {
676                        closed.store(true, std::sync::atomic::Ordering::SeqCst);
677                        *crate::realtime::lock_unpoisoned(&last_error) = Some(err.clone());
678                        let _ = crate::realtime::lock_unpoisoned(&tx_slot).take();
679                        if let Some(stream) =
680                            crate::realtime::lock_unpoisoned(&stream_slot).as_ref()
681                        {
682                            stream.fail(err);
683                        }
684                        return;
685                    }
686                };
687                let Some(tx) = crate::realtime::lock_unpoisoned(&tx_slot).as_ref().cloned() else {
688                    return;
689                };
690                if !crate::realtime::try_enqueue(
691                    &tx,
692                    data,
693                    &closed,
694                    &last_error,
695                    "orderbook subscription queue full; consumer too slow",
696                ) && let Some(err @ Error::QueueOverflow(_)) =
697                    crate::realtime::lock_unpoisoned(&last_error).clone()
698                {
699                    let _ = crate::realtime::lock_unpoisoned(&tx_slot).take();
700                    if let Some(stream) = crate::realtime::lock_unpoisoned(&stream_slot).as_ref() {
701                        stream.fail(err);
702                    }
703                }
704            }) as Arc<dyn Fn() + Send + Sync>
705        };
706
707        let handle_delta = {
708            let state = state.clone();
709            let emit = emit.clone();
710            let stream_slot = stream_slot.clone();
711            Arc::new(move |delta: OrderBookDeltaUpdate| {
712                let needs_refresh = {
713                    let mut s = crate::realtime::lock_unpoisoned(&state);
714                    let BookState {
715                        bids,
716                        asks,
717                        book_seq,
718                    } = &mut *s;
719                    let (new_seq, needs_refresh) = apply_delta(bids, asks, *book_seq, &delta);
720                    *book_seq = new_seq;
721                    needs_refresh
722                };
723                if needs_refresh {
724                    if let Some(stream) = crate::realtime::lock_unpoisoned(&stream_slot).as_ref() {
725                        stream.request_refresh();
726                    }
727                    return false;
728                }
729                emit();
730                true
731            }) as Arc<dyn Fn(OrderBookDeltaUpdate) -> bool + Send + Sync>
732        };
733
734        let svc = self.clone();
735        let fetch_symbol = symbol.clone();
736        let stream = SnapshotThenStream::new(SnapshotThenStreamConfig {
737            client: self.ctx.realtime.clone(),
738            channel,
739            decode: Arc::new(crate::codecs::decode::orderbook_delta_from_bytes),
740            fetch_snapshot: Arc::new(move || {
741                let svc = svc.clone();
742                let symbol = fetch_symbol.clone();
743                Box::pin(async move { svc.get(&symbol, Some(ws_depth)).await })
744            }),
745            read_publication: Arc::new(|delta: OrderBookDeltaUpdate| vec![delta]),
746            apply_snapshot: {
747                let state = state.clone();
748                let handle_delta = handle_delta.clone();
749                let emit = emit.clone();
750                let stream_slot = stream_slot.clone();
751                let last_error = last_error.clone();
752                Arc::new(
753                    move |snapshot: OrderbookData, buffered: Vec<OrderBookDeltaUpdate>| {
754                        let parsed_seq = match snapshot.book_seq.parse::<u64>() {
755                            Ok(seq) => seq,
756                            Err(_) => {
757                                // Fail toward refresh — never treat garbage as seq 0 silently.
758                                *crate::realtime::lock_unpoisoned(&last_error) =
759                                    Some(Error::realtime(
760                                        "orderbook snapshot book_seq is not a valid u64".to_owned(),
761                                    ));
762                                if let Some(stream) =
763                                    crate::realtime::lock_unpoisoned(&stream_slot).as_ref()
764                                {
765                                    stream.request_refresh();
766                                }
767                                return;
768                            }
769                        };
770                        {
771                            let mut s = crate::realtime::lock_unpoisoned(&state);
772                            s.bids = levels_from_orderbook_side(&snapshot.bids);
773                            s.asks = levels_from_orderbook_side(&snapshot.asks);
774                            s.book_seq = parsed_seq;
775                        }
776                        let mut applied_all = true;
777                        for delta in buffered {
778                            if !handle_delta(delta) {
779                                applied_all = false;
780                                break;
781                            }
782                        }
783                        if applied_all {
784                            emit();
785                        }
786                    },
787                )
788            },
789            apply_live_publications: {
790                let handle_delta = handle_delta.clone();
791                Arc::new(move |deltas: Vec<OrderBookDeltaUpdate>| {
792                    for delta in deltas {
793                        if !handle_delta(delta) {
794                            break;
795                        }
796                    }
797                })
798            },
799            max_buffered: 200,
800            on_reconnect: None,
801            on_snapshot_refresh: None,
802            on_error: None,
803        });
804        *crate::realtime::lock_unpoisoned(&stream_slot) = Some(stream.clone());
805
806        let subscription = crate::orderbook::Subscription::new(
807            rx,
808            stream.clone(),
809            closed,
810            bucket_ticks,
811            emit,
812            last_error,
813            tx_slot,
814        );
815        if let Err(err) = stream.start().await {
816            subscription.close();
817            return Err(err);
818        }
819        Ok(subscription)
820    }
821}
822
823struct BookState {
824    bids: crate::orderbook::BookSide,
825    asks: crate::orderbook::BookSide,
826    book_seq: u64,
827}
828
829#[cfg(test)]
830mod tests {
831    use super::*;
832    use std::sync::atomic::AtomicBool;
833    use std::sync::{Arc, Mutex};
834    use tokio::sync::mpsc;
835
836    #[test]
837    fn candle_channel_normalizes_aliases_to_human_label() {
838        for alias in ["1m", "MIN_1", "min1", "Min_1"] {
839            assert_eq!(
840                MarketDataService::candle_channel_timeframe(alias).unwrap(),
841                "1m",
842                "alias {alias}"
843            );
844        }
845        assert_eq!(
846            MarketDataService::candle_channel_timeframe("1h").unwrap(),
847            "1h"
848        );
849    }
850
851    #[test]
852    fn current_candle_selects_first_newest_row() {
853        let candle = |ts_sec| Candle {
854            ts_sec,
855            open: "1".into(),
856            high: "1".into(),
857            low: "1".into(),
858            close: "1".into(),
859            volume: "1".into(),
860            symbol_id: 1,
861            timeframe: "1m".into(),
862        };
863        let newest = newest_candle(CandlesResult {
864            symbol_id: 1,
865            timeframe: "1m".into(),
866            candles: vec![candle(20), candle(10)],
867            next_page_token: String::new(),
868        })
869        .expect("current candle");
870        assert_eq!(newest.ts_sec, 20);
871    }
872
873    #[tokio::test]
874    async fn orderbook_close_unblocks_recv() {
875        let (tx, rx) = mpsc::channel::<crate::models::OrderbookData>(2);
876        let tx_slot = Arc::new(Mutex::new(Some(tx)));
877        let closed = Arc::new(AtomicBool::new(false));
878        let last_error = Arc::new(Mutex::new(None));
879        let stream =
880            crate::realtime::SnapshotThenStream::new(crate::realtime::SnapshotThenStreamConfig {
881                client: crate::realtime::Client::new(
882                    "wss://example.invalid",
883                    "https://example.invalid",
884                    None,
885                    None,
886                ),
887                channel: "public:test".into(),
888                decode: Arc::new(|_: &[u8]| {
889                    Ok(crate::models::OrderBookDeltaUpdate {
890                        symbol_id: 1,
891                        book_seq_start: 1,
892                        book_seq_end: 1,
893                        reset: false,
894                        bids: vec![],
895                        asks: vec![],
896                    })
897                }),
898                fetch_snapshot: Arc::new(|| {
899                    Box::pin(async {
900                        Ok(crate::models::OrderbookData {
901                            symbol: "BTC-USDT".into(),
902                            depth: 1,
903                            book_seq: "1".into(),
904                            bids: vec![],
905                            asks: vec![],
906                        })
907                    })
908                }),
909                read_publication: Arc::new(|d| vec![d]),
910                apply_snapshot: Arc::new(|_, _| {}),
911                apply_live_publications: Arc::new(|_| {}),
912                max_buffered: 10,
913                on_reconnect: None,
914                on_snapshot_refresh: None,
915                on_error: None,
916            });
917        let mut sub = crate::orderbook::Subscription::new(
918            rx,
919            stream,
920            closed,
921            Arc::new(Mutex::new(0)),
922            Arc::new(|| {}),
923            last_error,
924            tx_slot,
925        );
926        sub.close();
927        let finished =
928            tokio::time::timeout(std::time::Duration::from_secs(1), sub.updates().recv())
929                .await
930                .expect("recv must not hang after close");
931        assert!(finished.is_none());
932    }
933}