tycho_simulation/book/mod.rs
1//! Book feeds: a provider's complete set of priced pairs, republished whenever it changes.
2
3use std::{collections::HashMap, fmt, sync::Arc};
4
5use chrono::{DateTime, Utc};
6use tracing::debug;
7use tycho_common::{
8 models::{token::Token, Chain},
9 simulation::protocol_sim::ProtocolSim,
10 Bytes,
11};
12
13use crate::{
14 protocol::models::ProtocolComponent,
15 snapshot_feed::{
16 errors::FeedError, SnapshotFeedEvent, SnapshotFeedOutcome, SnapshotFeedStream,
17 SnapshotFeedStreams, SnapshotFeedWatch,
18 },
19};
20
21#[cfg_attr(
22 not(test),
23 expect(
24 dead_code,
25 reason = "nothing in this crate implements a book feed yet; the layer is exercised by its own tests"
26 )
27)]
28pub(crate) mod component;
29#[cfg_attr(
30 not(test),
31 expect(
32 dead_code,
33 reason = "nothing in this crate implements a book feed yet; the layer is exercised by its own tests"
34 )
35)]
36pub(crate) mod levels;
37pub mod quote_tokens;
38#[cfg_attr(
39 not(test),
40 expect(
41 dead_code,
42 reason = "nothing in this crate implements a book feed yet; the layer is exercised by its own tests"
43 )
44)]
45pub(crate) mod sim;
46#[cfg_attr(
47 not(test),
48 expect(
49 dead_code,
50 reason = "nothing in this crate implements a book feed yet; the layer is exercised by its own tests"
51 )
52)]
53pub(crate) mod tvl;
54
55/// One pair's book: the component that identifies the pair and the ready-to-simulate state
56/// holding the provider's price levels for it.
57///
58/// The state is shared, so keeping one past the snapshot it came from costs a reference count;
59/// cloning a whole book copies the component.
60#[derive(Clone, Debug)]
61pub struct Book {
62 pub component: ProtocolComponent,
63 pub state: Arc<dyn ProtocolSim>,
64 /// When the provider last updated this book, on the provider's clock; `None` when the
65 /// provider reports none. The snapshot's [`ReceivedAt`] anchor is the feed's own.
66 pub updated_at: Option<DateTime<Utc>>,
67}
68
69/// One provider's complete set of books at one instant, so a pair absent from a snapshot is one
70/// the provider has stopped serving. The books are `Arc`'d so cloning the snapshot (e.g. out of a
71/// watch channel) never copies states.
72///
73/// This is the [`SnapshotFeed::Snapshot`](crate::snapshot_feed::SnapshotFeed) payload of every
74/// feed in this module, wrapped in an `Option` whose `None` means there is no servable snapshot:
75/// none received yet, or the last one withdrawn as stale (see `max_snapshot_age` on both feed
76/// configs).
77///
78/// `A` is what the snapshot is anchored to. Every feed in this module publishes
79/// `BookSnapshot<ReceivedAt>`; a feed whose snapshots belong to a block rather than to a moment
80/// publishes a block anchor instead, and the anchor type tells its consumers which.
81#[derive(Clone, Debug)]
82pub struct BookSnapshot<A> {
83 pub anchor: A,
84 pub books: Arc<HashMap<String, Book>>,
85}
86
87impl BookSnapshot<ReceivedAt> {
88 /// `books` as the complete set a feed received just now.
89 pub fn received_now(books: HashMap<String, Book>) -> Self {
90 BookSnapshot { anchor: ReceivedAt(Utc::now()), books: Arc::new(books) }
91 }
92}
93
94/// When the feed received a snapshot, on the feed's clock.
95#[derive(Clone, Copy, Debug, PartialEq, Eq)]
96pub struct ReceivedAt(pub DateTime<Utc>);
97
98/// What every book feed needs, whatever transport it runs on: the chain, the tokens it may serve
99/// pairs of, and the liquidity floor below which a book is not worth publishing. Every feed
100/// builder takes it by value as its first argument and stores it whole; build a second config to
101/// give one provider a different universe or floor. The token map is reference-counted, so
102/// cloning the config for several feeds shares one map.
103///
104/// The transport's own tuning is separate:
105/// [`snapshot_feed::ws::WsFeedConfig`](crate::snapshot_feed::ws::WsFeedConfig)
106/// and [`snapshot_feed::http::HttpFeedConfig`](crate::snapshot_feed::http::HttpFeedConfig).
107#[derive(Clone, derive_more::Debug)]
108pub struct BookFeedConfig {
109 pub chain: Chain,
110 /// The universe of tradable tokens with their metadata; pairs involving tokens missing here
111 /// are not served.
112 #[debug("{} tokens", tokens.len())]
113 pub tokens: Arc<HashMap<Bytes, Token>>,
114 /// Minimum book TVL in USD; 100 USD is a sensible floor for the supported providers.
115 pub min_tvl_usd: f64,
116}
117
118impl BookFeedConfig {
119 /// The two tokens of a pair, or `None` when either is outside the universe and the pair is
120 /// therefore not served.
121 #[expect(
122 dead_code,
123 reason = "nothing in this crate implements a book feed yet; the layer is exercised by its own tests"
124 )]
125 pub(crate) fn pair_tokens(&self, a: &Bytes, b: &Bytes) -> Option<(&Token, &Token)> {
126 Some((self.tokens.get(a)?, self.tokens.get(b)?))
127 }
128
129 /// Whether a book with `tvl_usd` clears the floor; names `book_label` in the log line when
130 /// it does not.
131 #[expect(
132 dead_code,
133 reason = "nothing in this crate implements a book feed yet; the layer is exercised by its own tests"
134 )]
135 pub(crate) fn clears_min_tvl(&self, tvl_usd: f64, book_label: impl fmt::Display) -> bool {
136 let clears = tvl_usd >= self.min_tvl_usd;
137 if !clears {
138 debug!(
139 book = %book_label,
140 tvl_usd,
141 min_tvl_usd = self.min_tvl_usd,
142 "filtering out book below the TVL floor"
143 );
144 }
145 clears
146 }
147}
148
149/// Several book feeds, merged — [`SnapshotFeedStreams`] with the types every book feed uses.
150///
151/// [`add`](SnapshotFeedStreams::add) each feed under the label you want its events to carry —
152/// usually as each one is built, since a provider without credentials or without support for the
153/// chain simply never joins the set — then read [`BookFeedEvent`]s until the set runs dry:
154///
155/// ```
156/// # use futures::StreamExt as _;
157/// # use tycho_simulation::book::{BookFeedEvent, BookFeedStreams};
158/// # use tycho_simulation::snapshot_feed::SnapshotFeedOutcome;
159/// # async fn consume() {
160/// // Anchored like the snapshots it carries; `add` infers this from the feed.
161/// let mut feeds: BookFeedStreams = BookFeedStreams::new();
162///
163/// while let Some((provider, event)) = feeds.next().await {
164/// match event {
165/// BookFeedEvent::Published(snapshot) => { /* quote from `snapshot.books` */ }
166/// BookFeedEvent::Withdrawn => { /* stop quoting `provider` until it publishes again */ }
167/// // `provider` is out of the set now, whichever of the three ended it.
168/// BookFeedEvent::Ended(SnapshotFeedOutcome::Failed(error)) => { /* it gave up */ }
169/// BookFeedEvent::Ended(outcome) => { /* ran out, or died of a bug */ }
170/// }
171/// }
172/// # }
173/// ```
174pub type BookFeedStreams<A = ReceivedAt> = SnapshotFeedStreams<BookSnapshot<A>, FeedError>;
175
176/// One book feed, read as a stream of [`BookFeedEvent`]s. See [`SnapshotFeedStream`].
177pub type BookFeedStream<A = ReceivedAt> = SnapshotFeedStream<BookSnapshot<A>, FeedError>;
178
179/// What happened on a book feed. See [`BookFeedStreams`].
180pub type BookFeedEvent<A = ReceivedAt> = SnapshotFeedEvent<BookSnapshot<A>, FeedError>;
181
182/// One book feed, held as the newest book set a consumer can read without awaiting — for one
183/// that prices on demand rather than reacting to every update, with
184/// [`ended`](SnapshotFeedWatch::ended) for why a venue stopped. See [`SnapshotFeedWatch`].
185pub type BookFeedWatch<A = ReceivedAt> = SnapshotFeedWatch<BookSnapshot<A>, FeedError>;
186
187/// How a book feed ended. See [`SnapshotFeedWatch::ended`].
188pub type BookFeedOutcome = SnapshotFeedOutcome<FeedError>;