Skip to main content

tycho_simulation/price_level_stream/
stream.rs

1use std::{
2    collections::{hash_map::Entry, HashMap, HashSet},
3    time::Duration,
4};
5
6use chrono::Utc;
7use num_bigint::BigUint;
8use tokio_stream::{Stream, StreamExt};
9use tycho_common::{
10    models::{token::Token, Chain},
11    simulation::protocol_sim::ProtocolSim,
12    Bytes,
13};
14
15use super::{
16    config::{
17        default_denied_pamms, default_served_pamms, PriceLevelStreamConfig,
18        DEFAULT_AUTO_DETECTED_GAS_COST,
19    },
20    state::{PriceLevelStreamQuote, PriceLevelStreamState},
21    titan::{
22        self, ConnectionSettings, TitanPairLevels, TitanPammLevels, TitanPriceLevel,
23        TitanPriceLevelMessage, TITAN_PRICE_LEVEL_URL, TITAN_PRICE_LEVEL_URL_ENV,
24    },
25};
26use crate::protocol::models::{ProtocolComponent, Update};
27
28/// Static attribute under which each emitted component carries its pAMM venue address.
29pub const PAMM_ADDRESS_ATTRIBUTE: &str = "pamm_address";
30
31/// Builds a stream of [`Update`]s from the Titan pAMM price level WebSocket.
32///
33/// A new builder serves no pAMMs: register the known venues via
34/// [`with_known_pamms`](Self::with_known_pamms), individual ones via
35/// [`add_pamm`](Self::add_pamm), or opt into serving unknown streamed venues via
36/// [`auto_detect`](Self::auto_detect); [`with_tokens`](Self::with_tokens) provides the token
37/// metadata pairs are interpreted with.
38///
39/// One component is emitted per (pAMM, token pair), identified by the concatenation
40/// `pamm ++ token0 ++ token1` (tokens sorted ascending), under the protocol system
41/// `fallback:{pamm}` — or `pricelevelstream:{pamm}` after
42/// [`without_fallback_router`](Self::without_fallback_router). The venue address is exposed
43/// through the [`PAMM_ADDRESS_ATTRIBUTE`] static attribute for downstream encoding.
44pub struct PriceLevelStreamBuilder {
45    registry: HashMap<Bytes, PriceLevelStreamConfig>,
46    denied: HashSet<Bytes>,
47    tokens: HashMap<Bytes, Token>,
48    url: Option<String>,
49    auto_detect: bool,
50    auto_detected_gas_cost: Option<BigUint>,
51    connection: ConnectionSettings,
52    /// Whether components are emitted under the `fallback:` family, executed through
53    /// `TychoFallbackRouter`, instead of the direct `pricelevelstream:` family.
54    fallback_router: bool,
55}
56
57impl Default for PriceLevelStreamBuilder {
58    fn default() -> Self {
59        Self {
60            registry: HashMap::new(),
61            denied: HashSet::new(),
62            tokens: HashMap::new(),
63            url: None,
64            auto_detect: false,
65            auto_detected_gas_cost: None,
66            connection: ConnectionSettings::default(),
67            fallback_router: true,
68        }
69    }
70}
71
72impl PriceLevelStreamBuilder {
73    pub fn new() -> Self {
74        Self::default()
75    }
76
77    /// Enables serving pAMMs that are not registered via
78    /// [`with_known_pamms`](Self::with_known_pamms) or [`add_pamm`](Self::add_pamm)
79    /// (disabled by default).
80    ///
81    /// When enabled, any unknown streamed venue — except denied ones (see
82    /// [`deny_pamm`](Self::deny_pamm)) — is served under its full lowercase hex address
83    /// as the name, with the default gas cost. A venue's protocol system therefore changes from
84    /// the address form (`pricelevelstream:{0xaddress}`) to a name (`pricelevelstream:{name}`)
85    /// once it gets registered — via [`add_pamm`](Self::add_pamm) or a release's
86    /// [`default_served_pamms`] recognizing it; the name-independent identifiers — the component id
87    /// and the [`PAMM_ADDRESS_ATTRIBUTE`] — stay stable across such renames.
88    pub fn auto_detect(mut self, enabled: bool) -> Self {
89        self.auto_detect = enabled;
90        self
91    }
92
93    /// Overrides the per-swap gas cost that auto-detected pAMMs (see
94    /// [`auto_detect`](Self::auto_detect)) are served with. Defaults to the maximum over the
95    /// known venue profiles, as the conservative choice. Registered venues are unaffected —
96    /// their gas cost comes from their [`PriceLevelStreamConfig`].
97    pub fn auto_detected_gas_cost(mut self, gas_cost: BigUint) -> Self {
98        self.auto_detected_gas_cost = Some(gas_cost);
99        self
100    }
101
102    /// Overrides the stream endpoint, e.g. to connect to a closer Titan region than the default
103    /// (see <https://docs.titanbuilder.xyz/propamms/takers>). Without it, the
104    /// `TITAN_PAMM_PRICE_LEVEL_URL` environment variable is used when set, else the built-in
105    /// default.
106    pub fn endpoint(mut self, url: impl Into<String>) -> Self {
107        self.url = Some(url.into());
108        self
109    }
110
111    /// Overrides how long a single connection attempt may take before it is aborted and retried
112    /// (default: 10s), so a hung TCP/TLS handshake cannot block the stream forever.
113    pub fn connect_timeout(mut self, timeout: Duration) -> Self {
114        self.connection.connect_timeout = timeout;
115        self
116    }
117
118    /// Overrides the longest gap between Titan messages tolerated before the connection is
119    /// treated as dead and re-established (default: 30s). Titan pushes several updates per
120    /// second, so a multi-second silence means a stalled or half-open connection.
121    pub fn read_idle_timeout(mut self, timeout: Duration) -> Self {
122        self.connection.read_idle_timeout = timeout;
123        self
124    }
125
126    /// Overrides the cap on the exponential reconnect backoff of `2^attempt` seconds
127    /// (default: 32s).
128    pub fn max_backoff(mut self, max_backoff: Duration) -> Self {
129        self.connection.max_backoff = max_backoff;
130        self
131    }
132
133    /// Registers a pAMM to be served under the given configuration, overriding any default,
134    /// denied, or auto-detected one for the same address.
135    ///
136    /// Between [`add_pamm`](Self::add_pamm) and [`deny_pamm`](Self::deny_pamm) for the same
137    /// address, the later call wins; the defaults applied by
138    /// [`with_known_pamms`](Self::with_known_pamms) never override either, in any call order.
139    pub fn add_pamm(mut self, config: PriceLevelStreamConfig) -> Self {
140        self.denied.remove(&config.address);
141        self.registry
142            .insert(config.address.clone(), config);
143        self
144    }
145
146    /// Excludes a venue from being served: drops its current registration (default or explicit)
147    /// and blocks auto-detecting it.
148    ///
149    /// Between [`add_pamm`](Self::add_pamm) and [`deny_pamm`](Self::deny_pamm) for the same
150    /// address, the later call wins; the defaults applied by
151    /// [`with_known_pamms`](Self::with_known_pamms) never override either, in any call order —
152    /// so denying a venue from the default set works whether the denial comes before or after
153    /// [`with_known_pamms`](Self::with_known_pamms).
154    pub fn deny_pamm(mut self, address: Bytes) -> Self {
155        self.registry.remove(&address);
156        self.denied.insert(address);
157        self
158    }
159
160    /// Applies what is known about the streamed venues: registers the known-good ones
161    /// ([`default_served_pamms`]) to be served and denies the known-bad ones
162    /// ([`default_denied_pamms`]) — venues that stream quotes but whose swaps are not executable.
163    ///
164    /// These defaults never override an explicit [`add_pamm`](Self::add_pamm) or
165    /// [`deny_pamm`](Self::deny_pamm) for the same address, regardless of call order.
166    pub fn with_known_pamms(mut self) -> Self {
167        for config in default_served_pamms() {
168            if self.denied.contains(&config.address) {
169                continue;
170            }
171            self.registry
172                .entry(config.address.clone())
173                .or_insert(config);
174        }
175        for address in default_denied_pamms() {
176            if self.registry.contains_key(&address) {
177                continue;
178            }
179            self.denied.insert(address);
180        }
181        self
182    }
183
184    /// Provides the token metadata used to build components and interpret amounts. Pairs whose
185    /// tokens are missing here are skipped.
186    pub fn with_tokens(mut self, tokens: HashMap<Bytes, Token>) -> Self {
187        self.tokens = tokens;
188        self
189    }
190
191    /// Keeps every venue on the direct `pricelevelstream:{name}` path, so swaps execute on the
192    /// venues themselves and a stale maker quote reverts the route.
193    ///
194    /// By default components are emitted under `fallback:{name}`, so tycho-execution routes
195    /// their swaps through `TychoFallbackRouter`. Opt out when the direct call is what you want
196    /// to measure or execute.
197    pub fn without_fallback_router(mut self) -> Self {
198        self.fallback_router = false;
199        self
200    }
201
202    /// Consumes the builder and opens the stream.
203    ///
204    /// Components are emitted under `fallback:{name}`, so tycho-execution routes their swaps
205    /// through `TychoFallbackRouter`, which retries a reverted pAMM swap — a stale maker quote
206    /// reverts in any simulation against a mined block — on the fallback pool the solver names.
207    /// [`without_fallback_router`](Self::without_fallback_router) keeps them on the direct
208    /// `pricelevelstream:` path.
209    ///
210    /// The connection is established lazily on first poll and maintained (with reconnects) for as
211    /// long as the stream is polled; it never terminates on its own, and dropping the stream
212    /// closes the connection. Frames that contain no served pAMM produce no update.
213    ///
214    /// Each streamed frame is a complete snapshot of everything Titan currently streams, so
215    /// every update carries the full set of the frame's pair states, with `new_pairs` /
216    /// `removed_pairs` derived by diffing against the previous frame — a pair (or a whole
217    /// venue) the stream stops serving is removed. Frames older than an already processed one
218    /// are skipped, so updates never move backwards in block number. Pairs whose tokens are
219    /// missing from the provided token metadata are skipped.
220    pub fn build(self) -> impl Stream<Item = Update> + Send {
221        let Self {
222            registry,
223            denied,
224            tokens,
225            url,
226            auto_detect,
227            auto_detected_gas_cost,
228            connection,
229            fallback_router,
230        } = self;
231        if registry.is_empty() && !auto_detect {
232            tracing::warn!(
233                "No pAMMs registered and auto-detection is off; the stream will never produce \
234                 an update"
235            );
236        }
237        if tokens.is_empty() {
238            tracing::warn!(
239                "No token metadata provided; every streamed pair will be skipped and the stream \
240                 will never produce an update"
241            );
242        }
243        let url = url.unwrap_or_else(|| {
244            std::env::var(TITAN_PRICE_LEVEL_URL_ENV)
245                .unwrap_or_else(|_| TITAN_PRICE_LEVEL_URL.to_string())
246        });
247        let auto_detected_gas_cost =
248            auto_detected_gas_cost.unwrap_or_else(|| BigUint::from(DEFAULT_AUTO_DETECTED_GAS_COST));
249
250        let mut tracker = SnapshotTracker::new(
251            registry,
252            denied,
253            tokens,
254            auto_detect,
255            auto_detected_gas_cost,
256            fallback_router,
257        );
258        titan::messages(url, connection).filter_map(move |message| tracker.process(message))
259    }
260}
261
262/// Turns Titan frames into [`Update`]s, tracking the previously emitted components so pair
263/// additions and removals can be diffed against the last snapshot.
264struct SnapshotTracker {
265    registry: HashMap<Bytes, PriceLevelStreamConfig>,
266    /// Venues excluded from auto-detection. The builder keeps this disjoint from the registry:
267    /// denying removes any registration and registering removes any denial.
268    denied: HashSet<Bytes>,
269    tokens: HashMap<Bytes, Token>,
270    /// Whether frames from pAMMs absent from the registry get an address-named configuration
271    /// synthesized (and cached in the registry) instead of being skipped.
272    auto_detect: bool,
273    /// The per-swap gas cost synthesized auto-detected configurations are served with.
274    auto_detected_gas_cost: BigUint,
275    /// Whether components are emitted under the `fallback:` family, so their swaps execute
276    /// through `TychoFallbackRouter` instead of the venue directly.
277    via_fallback_router: bool,
278    /// Components of the last emitted snapshot, across all pAMMs. A frame is a complete
279    /// snapshot of everything Titan currently streams, so removals are diffed globally: a
280    /// known component a frame does not re-emit is gone — including when its venue vanishes
281    /// from the stream entirely.
282    components: HashMap<String, ProtocolComponent>,
283    /// The newest block number processed so far. Frames targeting an older block (e.g.
284    /// delivered around a reconnect) are stale and skipped wholesale — processing one would
285    /// emit superseded states and churn the global diff.
286    newest_block: u64,
287}
288
289impl SnapshotTracker {
290    fn new(
291        registry: HashMap<Bytes, PriceLevelStreamConfig>,
292        denied: HashSet<Bytes>,
293        tokens: HashMap<Bytes, Token>,
294        auto_detect: bool,
295        auto_detected_gas_cost: BigUint,
296        via_fallback_router: bool,
297    ) -> Self {
298        Self {
299            registry,
300            denied,
301            tokens,
302            auto_detect,
303            auto_detected_gas_cost,
304            via_fallback_router,
305            components: HashMap::new(),
306            newest_block: 0,
307        }
308    }
309
310    /// Processes one frame into an [`Update`], or `None` if the frame targets an older block
311    /// than an already processed one or contains nothing relevant (no registered pAMM with at
312    /// least one known pair or a pair removal).
313    fn process(&mut self, message: TitanPriceLevelMessage) -> Option<Update> {
314        if message.block_number < self.newest_block {
315            tracing::warn!(
316                block_number = message.block_number,
317                newest_block = self.newest_block,
318                "Skipping out-of-order price level frame"
319            );
320            return None;
321        }
322        self.newest_block = message.block_number;
323
324        let mut states: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
325        let mut new_pairs = HashMap::new();
326        // The frame is a complete snapshot: every known component is presumed gone until the
327        // frame re-emits it below.
328        let mut previous = std::mem::take(&mut self.components);
329
330        for TitanPammLevels { pamm, pairs } in message.pamms {
331            let config = match self.registry.entry(pamm.clone()) {
332                Entry::Occupied(entry) => &*entry.into_mut(),
333                Entry::Vacant(entry) => {
334                    if !self.auto_detect {
335                        tracing::debug!(%pamm, "Skipping unregistered pAMM");
336                        continue;
337                    }
338                    if self.denied.contains(&pamm) {
339                        tracing::debug!(%pamm, "Skipping denied pAMM");
340                        continue;
341                    }
342                    tracing::info!(%pamm, "Serving auto-detected pAMM");
343                    &*entry.insert(PriceLevelStreamConfig::auto_detected(
344                        pamm.clone(),
345                        self.auto_detected_gas_cost.clone(),
346                    ))
347                }
348            };
349
350            // Merge the frame's per-direction ladders into one entry per unordered token pair.
351            let mut merged_pairs: HashMap<(Bytes, Bytes), (Vec<_>, Vec<_>)> = HashMap::new();
352            for TitanPairLevels { token_in, token_out, order_book } in pairs {
353                if !self.tokens.contains_key(&token_in) || !self.tokens.contains_key(&token_out) {
354                    tracing::debug!(%token_in, %token_out, "Skipping pair with unknown token");
355                    continue;
356                }
357                let sells_token0 = token_in < token_out;
358                let key = if sells_token0 {
359                    (token_in.clone(), token_out.clone())
360                } else {
361                    (token_out.clone(), token_in.clone())
362                };
363                let quotes = order_book
364                    .into_iter()
365                    .map(|TitanPriceLevel { amount_in, amount_out }| {
366                        PriceLevelStreamQuote::new(amount_in, amount_out)
367                    })
368                    .collect();
369                let entry = merged_pairs.entry(key).or_default();
370                if sells_token0 {
371                    entry.0 = quotes;
372                } else {
373                    entry.1 = quotes;
374                }
375            }
376
377            for ((token0, token1), (quotes_0_to_1, quotes_1_to_0)) in merged_pairs {
378                let id = component_id(&config.address, &token0, &token1);
379                let id_string = id.to_string();
380                let component = previous
381                    .remove(&id_string)
382                    .unwrap_or_else(|| {
383                        let component = build_component(
384                            &self.tokens,
385                            config,
386                            id,
387                            &token0,
388                            &token1,
389                            self.via_fallback_router,
390                        );
391                        new_pairs.insert(id_string.clone(), component.clone());
392                        component
393                    });
394
395                let state = PriceLevelStreamState::new(
396                    token0,
397                    token1,
398                    quotes_0_to_1,
399                    quotes_1_to_0,
400                    config.gas_cost.clone(),
401                );
402
403                states.insert(id_string.clone(), Box::new(state));
404                self.components
405                    .insert(id_string, component);
406            }
407        }
408
409        // Every re-emitted pair was moved back into `self.components` above — whatever remains
410        // is gone: the pair, or its whole venue, is no longer streamed.
411        let removed_pairs = previous;
412
413        if states.is_empty() && new_pairs.is_empty() && removed_pairs.is_empty() {
414            return None;
415        }
416
417        Some(
418            // Quotes target the block currently being built, hence partial. Sync states stay
419            // empty (like the RFQ path) because no full block header is available.
420            Update::new(message.block_number, states, new_pairs)
421                .set_is_partial(true)
422                .set_removed_pairs(removed_pairs),
423        )
424    }
425}
426
427fn build_component(
428    tokens: &HashMap<Bytes, Token>,
429    config: &PriceLevelStreamConfig,
430    id: Bytes,
431    token0: &Bytes,
432    token1: &Bytes,
433    via_router: bool,
434) -> ProtocolComponent {
435    let protocol_system =
436        if via_router { config.fallback_protocol_system() } else { config.protocol_system() };
437    ProtocolComponent::new(
438        id,
439        protocol_system.clone(),
440        protocol_system,
441        // Titan builds Ethereum L1 blocks; the stream carries no other chains.
442        Chain::Ethereum,
443        vec![tokens[token0].clone(), tokens[token1].clone()],
444        vec![config.address.clone()],
445        HashMap::from([(PAMM_ADDRESS_ATTRIBUTE.to_string(), config.address.clone())]),
446        Bytes::default(),
447        Utc::now().naive_utc(),
448    )
449}
450
451/// The component identity of a (pAMM, pair) combination: `pamm ++ token0 ++ token1`.
452fn component_id(pamm: &Bytes, token0: &Bytes, token1: &Bytes) -> Bytes {
453    Bytes::from([pamm.as_ref(), token0.as_ref(), token1.as_ref()].concat())
454}
455
456#[cfg(test)]
457mod tests {
458    use std::str::FromStr;
459
460    use num_bigint::BigUint;
461
462    use super::*;
463
464    const PAMM: &str = "0x5979458912f80b96d30d4220af8e2e4925a33320";
465    const WBTC: &str = "0x2260fac5e5542a773aa44fbcfedf7c193bc2c599";
466    const USDC: &str = "0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48";
467    const WETH: &str = "0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2";
468
469    fn token(address: &str, symbol: &str, decimals: u32) -> Token {
470        Token::new(
471            &Bytes::from_str(address).unwrap(),
472            symbol,
473            decimals,
474            0,
475            &[Some(10_000)],
476            Chain::Ethereum,
477            100,
478        )
479    }
480
481    fn tokens() -> HashMap<Bytes, Token> {
482        [token(WBTC, "WBTC", 8), token(USDC, "USDC", 6), token(WETH, "WETH", 18)]
483            .into_iter()
484            .map(|token| (token.address.clone(), token))
485            .collect()
486    }
487
488    fn tracker() -> SnapshotTracker {
489        let config = PriceLevelStreamConfig::new(
490            "fermiswap",
491            Bytes::from_str(PAMM).unwrap(),
492            BigUint::from(120_000u64),
493        );
494        SnapshotTracker::new(
495            HashMap::from([(config.address.clone(), config)]),
496            HashSet::new(),
497            tokens(),
498            false,
499            BigUint::from(DEFAULT_AUTO_DETECTED_GAS_COST),
500            false,
501        )
502    }
503
504    fn level(amount_in: u64, amount_out: u64) -> TitanPriceLevel {
505        TitanPriceLevel {
506            amount_in: BigUint::from(amount_in),
507            amount_out: BigUint::from(amount_out),
508        }
509    }
510
511    fn pair_levels(
512        token_in: &str,
513        token_out: &str,
514        order_book: Vec<TitanPriceLevel>,
515    ) -> TitanPairLevels {
516        TitanPairLevels {
517            token_in: Bytes::from_str(token_in).unwrap(),
518            token_out: Bytes::from_str(token_out).unwrap(),
519            order_book,
520        }
521    }
522
523    fn message(block_number: u64, pairs: Vec<TitanPairLevels>) -> TitanPriceLevelMessage {
524        TitanPriceLevelMessage {
525            block_number,
526            pamms: vec![TitanPammLevels { pamm: Bytes::from_str(PAMM).unwrap(), pairs }],
527        }
528    }
529
530    fn wbtc_usdc_pairs() -> Vec<TitanPairLevels> {
531        vec![
532            pair_levels(WBTC, USDC, vec![level(100_000_000, 100_000_000_000)]),
533            pair_levels(USDC, WBTC, vec![level(100_000_000_000, 99_000_000)]),
534        ]
535    }
536
537    fn expected_id() -> String {
538        // pamm ++ token0 ++ token1 with WBTC < USDC.
539        format!("{PAMM}{}{}", &WBTC[2..], &USDC[2..])
540    }
541
542    #[test]
543    fn first_snapshot_emits_new_pair_with_both_directions() {
544        let mut tracker = tracker();
545        let Update {
546            block_number_or_timestamp,
547            is_partial,
548            sync_states,
549            states,
550            new_pairs,
551            removed_pairs,
552        } = tracker
553            .process(message(100, wbtc_usdc_pairs()))
554            .expect("update expected");
555
556        assert_eq!(block_number_or_timestamp, 100);
557        assert!(is_partial);
558        assert!(sync_states.is_empty());
559        assert!(removed_pairs.is_empty());
560
561        let id = expected_id();
562        let component = &new_pairs[&id];
563        assert_eq!(component.protocol_system, "pricelevelstream:fermiswap");
564        assert_eq!(
565            component.static_attributes[PAMM_ADDRESS_ATTRIBUTE],
566            Bytes::from_str(PAMM).unwrap()
567        );
568
569        let PriceLevelStreamState { token0, token1, quotes_0_to_1, quotes_1_to_0, gas_cost } =
570            states[&id]
571                .as_any()
572                .downcast_ref::<PriceLevelStreamState>()
573                .expect("price level state");
574        assert_eq!(token0, &Bytes::from_str(WBTC).unwrap());
575        assert_eq!(token1, &Bytes::from_str(USDC).unwrap());
576        assert_eq!(quotes_0_to_1.len(), 1);
577        assert_eq!(quotes_1_to_0.len(), 1);
578        assert_eq!(quotes_0_to_1[0].amount_in, BigUint::from(100_000_000u64));
579        assert_eq!(gas_cost, &BigUint::from(120_000u64));
580    }
581
582    #[test]
583    fn repeated_snapshot_is_not_a_new_pair() {
584        let mut tracker = tracker();
585        tracker
586            .process(message(100, wbtc_usdc_pairs()))
587            .expect("update expected");
588        let update = tracker
589            .process(message(101, wbtc_usdc_pairs()))
590            .expect("update expected");
591
592        assert!(update.new_pairs.is_empty());
593        assert!(update.removed_pairs.is_empty());
594        assert!(update
595            .states
596            .contains_key(&expected_id()));
597    }
598
599    #[test]
600    fn dropped_pair_is_removed() {
601        let mut tracker = tracker();
602        tracker
603            .process(message(100, wbtc_usdc_pairs()))
604            .expect("update expected");
605        let weth_usdc =
606            vec![pair_levels(WETH, USDC, vec![level(1_000_000_000_000_000_000, 3_000_000_000)])];
607        let update = tracker
608            .process(message(101, weth_usdc))
609            .expect("update expected");
610
611        assert_eq!(update.removed_pairs.len(), 1);
612        assert!(update
613            .removed_pairs
614            .contains_key(&expected_id()));
615        assert_eq!(update.new_pairs.len(), 1);
616        assert_eq!(update.states.len(), 1);
617    }
618
619    #[test]
620    fn out_of_order_frame_is_skipped() {
621        let mut tracker = tracker();
622        tracker
623            .process(message(101, wbtc_usdc_pairs()))
624            .expect("update expected");
625
626        // A frame for an older block is stale: no update, and the caches stay untouched even
627        // though the frame's snapshot differs completely.
628        let stale =
629            vec![pair_levels(WETH, USDC, vec![level(1_000_000_000_000_000_000, 3_000_000_000)])];
630        assert!(tracker
631            .process(message(100, stale))
632            .is_none());
633
634        // The next current frame diffs against the pre-stale state: nothing was added or
635        // removed in between.
636        let update = tracker
637            .process(message(102, wbtc_usdc_pairs()))
638            .expect("update expected");
639        assert!(update.new_pairs.is_empty());
640        assert!(update.removed_pairs.is_empty());
641    }
642
643    #[test]
644    fn vanished_pamm_has_its_pairs_removed() {
645        let mut tracker = tracker();
646        tracker
647            .process(message(100, wbtc_usdc_pairs()))
648            .expect("update expected");
649
650        // The next frame no longer contains the pAMM at all: a complete snapshot without a
651        // venue means the venue is gone, pairs and all.
652        let update = tracker
653            .process(TitanPriceLevelMessage { block_number: 101, pamms: vec![] })
654            .expect("update expected");
655        assert!(update.states.is_empty());
656        assert!(update.new_pairs.is_empty());
657        assert_eq!(update.removed_pairs.len(), 1);
658        assert!(update
659            .removed_pairs
660            .contains_key(&expected_id()));
661
662        // Nothing served and nothing changed: no update.
663        assert!(tracker
664            .process(TitanPriceLevelMessage { block_number: 102, pamms: vec![] })
665            .is_none());
666
667        // A venue that reappears is a new pair again.
668        let update = tracker
669            .process(message(103, wbtc_usdc_pairs()))
670            .expect("update expected");
671        assert!(update
672            .new_pairs
673            .contains_key(&expected_id()));
674    }
675
676    #[test]
677    fn unregistered_pamm_produces_no_update_without_auto_detection() {
678        let mut tracker = SnapshotTracker::new(
679            HashMap::new(),
680            HashSet::new(),
681            tokens(),
682            false,
683            BigUint::from(DEFAULT_AUTO_DETECTED_GAS_COST),
684            false,
685        );
686        assert!(tracker
687            .process(message(100, wbtc_usdc_pairs()))
688            .is_none());
689    }
690
691    #[test]
692    fn denied_pamm_is_not_auto_detected() {
693        let denied = HashSet::from([Bytes::from_str(PAMM).unwrap()]);
694        let mut tracker = SnapshotTracker::new(
695            HashMap::new(),
696            denied,
697            tokens(),
698            true,
699            BigUint::from(DEFAULT_AUTO_DETECTED_GAS_COST),
700            false,
701        );
702        assert!(tracker
703            .process(message(100, wbtc_usdc_pairs()))
704            .is_none());
705    }
706
707    #[test]
708    fn explicit_add_and_deny_are_last_wins() {
709        let address = Bytes::from_str(PAMM).unwrap();
710        let custom =
711            || PriceLevelStreamConfig::new("custom", Bytes::from_str(PAMM).unwrap(), 1u64.into());
712
713        let builder = PriceLevelStreamBuilder::new()
714            .add_pamm(custom())
715            .deny_pamm(address.clone());
716        assert!(!builder.registry.contains_key(&address));
717        assert!(builder.denied.contains(&address));
718
719        let builder = PriceLevelStreamBuilder::new()
720            .deny_pamm(address.clone())
721            .add_pamm(custom());
722        assert_eq!(builder.registry[&address].protocol, "custom");
723        assert!(builder.denied.is_empty());
724    }
725
726    #[test]
727    fn defaults_never_override_explicit_calls() {
728        // Denying a venue from the default set works in either call order.
729        let fermiswap_router = Bytes::from_str(PAMM).unwrap();
730        for builder in [
731            PriceLevelStreamBuilder::new()
732                .deny_pamm(fermiswap_router.clone())
733                .with_known_pamms(),
734            PriceLevelStreamBuilder::new()
735                .with_known_pamms()
736                .deny_pamm(fermiswap_router.clone()),
737        ] {
738            assert!(!builder
739                .registry
740                .contains_key(&fermiswap_router));
741            assert!(builder
742                .denied
743                .contains(&fermiswap_router));
744            // The other defaults are unaffected.
745            assert!(!builder.registry.is_empty());
746        }
747
748        // Registering a venue from the default deny set works in either call order. Any one of
749        // them exercises that; the set is empty while every streamed venue is executable.
750        let Some(denied_venue) = default_denied_pamms().pop() else { return };
751        let custom = || PriceLevelStreamConfig::new("custom", denied_venue.clone(), 1u64.into());
752        for builder in [
753            PriceLevelStreamBuilder::new()
754                .add_pamm(custom())
755                .with_known_pamms(),
756            PriceLevelStreamBuilder::new()
757                .with_known_pamms()
758                .add_pamm(custom()),
759        ] {
760            assert_eq!(builder.registry[&denied_venue].protocol, "custom");
761            assert!(!builder.denied.contains(&denied_venue));
762        }
763    }
764
765    #[test]
766    fn auto_detected_pamm_is_served_under_its_address() {
767        let mut tracker = SnapshotTracker::new(
768            HashMap::new(),
769            HashSet::new(),
770            tokens(),
771            true,
772            BigUint::from(DEFAULT_AUTO_DETECTED_GAS_COST),
773            false,
774        );
775        let update = tracker
776            .process(message(100, wbtc_usdc_pairs()))
777            .expect("update expected");
778
779        let component = &update.new_pairs[&expected_id()];
780        assert_eq!(component.protocol_system, format!("pricelevelstream:{PAMM}"));
781        let state = update.states[&expected_id()]
782            .as_any()
783            .downcast_ref::<PriceLevelStreamState>()
784            .expect("price level state");
785        assert_eq!(state.gas_cost, BigUint::from(DEFAULT_AUTO_DETECTED_GAS_COST));
786
787        // The synthesized config is cached: the next snapshot is not a new pair again.
788        let update = tracker
789            .process(message(101, wbtc_usdc_pairs()))
790            .expect("update expected");
791        assert!(update.new_pairs.is_empty());
792    }
793
794    #[test]
795    fn auto_detected_gas_cost_override_applies() {
796        let mut tracker = SnapshotTracker::new(
797            HashMap::new(),
798            HashSet::new(),
799            tokens(),
800            true,
801            BigUint::from(42_000u64),
802            false,
803        );
804        let update = tracker
805            .process(message(100, wbtc_usdc_pairs()))
806            .expect("update expected");
807
808        let state = update.states[&expected_id()]
809            .as_any()
810            .downcast_ref::<PriceLevelStreamState>()
811            .expect("price level state");
812        assert_eq!(state.gas_cost, BigUint::from(42_000u64));
813    }
814
815    #[test]
816    fn with_known_pamms_registers_known_venues() {
817        // PAMM is the FermiSwap router, one of the default venues.
818        let fermiswap_router = Bytes::from_str(PAMM).unwrap();
819
820        let builder = PriceLevelStreamBuilder::new();
821        assert!(builder.registry.is_empty());
822        assert!(builder.denied.is_empty());
823
824        let builder = builder.with_known_pamms();
825        assert_eq!(builder.registry[&fermiswap_router].protocol, "fermiswap");
826        // The known-bad venues get denied alongside, and never overlap the served defaults.
827        assert_eq!(
828            builder.denied,
829            default_denied_pamms()
830                .into_iter()
831                .collect()
832        );
833        assert!(builder.denied.is_disjoint(
834            &builder
835                .registry
836                .keys()
837                .cloned()
838                .collect()
839        ));
840
841        // An `add_pamm` entry wins over the default for the same address, in either call order.
842        let custom =
843            || PriceLevelStreamConfig::new("custom", fermiswap_router.clone(), BigUint::from(1u64));
844        for builder in [
845            PriceLevelStreamBuilder::new()
846                .add_pamm(custom())
847                .with_known_pamms(),
848            PriceLevelStreamBuilder::new()
849                .with_known_pamms()
850                .add_pamm(custom()),
851        ] {
852            assert_eq!(builder.registry[&fermiswap_router].protocol, "custom");
853            assert_eq!(builder.registry[&fermiswap_router].gas_cost, BigUint::from(1u64));
854        }
855    }
856
857    /// Components are emitted under `fallback:{name}`, so their swaps execute through
858    /// `TychoFallbackRouter`; identity and attributes are the same as on the direct path.
859    #[test]
860    fn venues_are_served_under_the_fallback_family() {
861        let config = PriceLevelStreamConfig::new(
862            "fermiswap",
863            Bytes::from_str(PAMM).unwrap(),
864            BigUint::from(120_000u64),
865        );
866        let mut tracker = SnapshotTracker::new(
867            HashMap::from([(config.address.clone(), config)]),
868            HashSet::new(),
869            tokens(),
870            false,
871            BigUint::from(DEFAULT_AUTO_DETECTED_GAS_COST),
872            true,
873        );
874
875        let update = tracker
876            .process(message(100, wbtc_usdc_pairs()))
877            .expect("update expected");
878
879        let component = &update.new_pairs[&expected_id()];
880        assert_eq!(component.protocol_system, "fallback:fermiswap");
881        assert_eq!(
882            component.static_attributes[PAMM_ADDRESS_ATTRIBUTE],
883            Bytes::from_str(PAMM).unwrap()
884        );
885    }
886
887    /// Auto-detected, address-named venues take the fallback family too.
888    #[test]
889    fn auto_detected_venue_is_served_under_the_fallback_family() {
890        let mut tracker = SnapshotTracker::new(
891            HashMap::new(),
892            HashSet::new(),
893            tokens(),
894            true,
895            BigUint::from(DEFAULT_AUTO_DETECTED_GAS_COST),
896            true,
897        );
898
899        let update = tracker
900            .process(message(100, wbtc_usdc_pairs()))
901            .expect("update expected");
902
903        let component = &update.new_pairs[&expected_id()];
904        assert_eq!(component.protocol_system, format!("fallback:{PAMM}"));
905    }
906
907    /// Off the fallback router, a venue keeps the direct `pricelevelstream:{name}` family.
908    #[test]
909    fn without_fallback_router_keeps_the_direct_family() {
910        let mut tracker = tracker();
911
912        let update = tracker
913            .process(message(100, wbtc_usdc_pairs()))
914            .expect("update expected");
915
916        assert_eq!(update.new_pairs[&expected_id()].protocol_system, "pricelevelstream:fermiswap");
917    }
918
919    /// The fallback router path is the default; `without_fallback_router` is the way off it.
920    #[test]
921    fn fallback_router_is_on_unless_opted_out() {
922        assert!(PriceLevelStreamBuilder::new().fallback_router);
923        assert!(
924            !PriceLevelStreamBuilder::new()
925                .without_fallback_router()
926                .fallback_router
927        );
928    }
929
930    /// The families this stream emits are the ones tycho-execution resolves an encoder for. A
931    /// drift between the two makes every route through a pAMM fail to encode.
932    #[test]
933    fn families_match_the_execution_side_prefixes() {
934        use tycho_execution::encoding::evm::{FALLBACK_PREFIX, PRICE_LEVEL_STREAM_PREFIX};
935
936        use super::super::config::{FALLBACK_FAMILY, PRICE_LEVEL_STREAM_FAMILY};
937
938        assert_eq!(format!("{PRICE_LEVEL_STREAM_FAMILY}:"), PRICE_LEVEL_STREAM_PREFIX);
939        assert_eq!(format!("{FALLBACK_FAMILY}:"), FALLBACK_PREFIX);
940    }
941
942    #[test]
943    fn unknown_tokens_are_skipped() {
944        let mut tracker = tracker();
945        let unknown = vec![pair_levels(
946            "0x1111111111111111111111111111111111111111",
947            USDC,
948            vec![level(1, 1)],
949        )];
950        assert!(tracker
951            .process(message(100, unknown))
952            .is_none());
953    }
954}