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 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 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 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 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 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 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 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 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 #[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#[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#[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 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 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#[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 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 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 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 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 let tx_slot: Arc<Mutex<Option<mpsc::Sender<OrderbookData>>>> =
641 Arc::new(Mutex::new(Some(tx)));
642 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 *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}