tycho-simulation 0.450.0

Provides tools for interacting with protocol states, calculating spot prices, and quoting token swaps.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
use std::{
    collections::{HashMap, HashSet},
    time::{Duration, Instant, SystemTime, UNIX_EPOCH},
};

use async_stream::stream;
use num_bigint::BigUint;
use tokio_stream::{Stream, StreamExt};
use tycho_common::{models::token::Token, Bytes};

use super::{
    config::{
        default_denied_pamms, default_served_pamms, PriceLevelStreamConfig,
        DEFAULT_AUTO_DETECTED_GAS_COST,
    },
    titan::{self, ConnectionSettings, TITAN_PRICE_LEVEL_URL, TITAN_PRICE_LEVEL_URL_ENV},
    tracker::{FreshnessTracker, Now, TrackerSettings, DEFAULT_STALE_AFTER},
};
use crate::protocol::models::Update;

/// Static attribute under which each emitted component carries its pAMM venue address.
pub const PAMM_ADDRESS_ATTRIBUTE: &str = "pamm_address";

/// The longest [`stale_after`](PriceLevelStreamBuilder::stale_after) the builder accepts;
/// longer values are capped to it. Quotes target the block being built, so serving a ladder for
/// longer than this is never intended, and deadlines this far ahead stay representable on the
/// monotonic clock.
pub const MAX_STALE_AFTER: Duration = Duration::from_secs(3600);

/// Builds a stream of [`Update`]s from the Titan pAMM price level WebSocket.
///
/// A new builder serves no pAMMs: register the known venues via
/// [`with_known_pamms`](Self::with_known_pamms), individual ones via
/// [`add_pamm`](Self::add_pamm), or opt into serving unknown streamed venues via
/// [`auto_detect`](Self::auto_detect); [`with_tokens`](Self::with_tokens) provides the token
/// metadata pairs are interpreted with.
///
/// One component is emitted per (pAMM, token pair), identified by the concatenation
/// `pamm ++ token0 ++ token1` (tokens sorted ascending), under the protocol system
/// `fallback:{pamm}` — or `pricelevelstream:{pamm}` after
/// [`without_fallback_router`](Self::without_fallback_router). The venue address is exposed
/// through the [`PAMM_ADDRESS_ATTRIBUTE`] static attribute for downstream encoding.
pub struct PriceLevelStreamBuilder {
    registry: HashMap<Bytes, PriceLevelStreamConfig>,
    denied: HashSet<Bytes>,
    tokens: HashMap<Bytes, Token>,
    url: Option<String>,
    auto_detect: bool,
    auto_detected_gas_cost: Option<BigUint>,
    connection: ConnectionSettings,
    /// Whether components are emitted under the `fallback:` family, executed through
    /// `TychoFallbackRouter`, instead of the direct `pricelevelstream:` family.
    fallback_router: bool,
    /// See [`stale_after`](Self::stale_after).
    stale_after: Duration,
    /// See [`without_quote_guard`](Self::without_quote_guard).
    quote_guard: bool,
}

impl Default for PriceLevelStreamBuilder {
    fn default() -> Self {
        Self {
            registry: HashMap::new(),
            denied: HashSet::new(),
            tokens: HashMap::new(),
            url: None,
            auto_detect: false,
            auto_detected_gas_cost: None,
            connection: ConnectionSettings::default(),
            fallback_router: true,
            stale_after: DEFAULT_STALE_AFTER,
            quote_guard: true,
        }
    }
}

impl PriceLevelStreamBuilder {
    pub fn new() -> Self {
        Self::default()
    }

    /// Enables serving pAMMs that are not registered via
    /// [`with_known_pamms`](Self::with_known_pamms) or [`add_pamm`](Self::add_pamm)
    /// (disabled by default).
    ///
    /// When enabled, any unknown streamed venue — except denied ones (see
    /// [`deny_pamm`](Self::deny_pamm)) — is served under its full lowercase hex address
    /// as the name, with the default gas cost. A venue's protocol system therefore changes from
    /// the address form (`pricelevelstream:{0xaddress}`) to a name (`pricelevelstream:{name}`)
    /// once it gets registered — via [`add_pamm`](Self::add_pamm) or a release's
    /// [`default_served_pamms`] recognizing it; the name-independent identifiers — the component id
    /// and the [`PAMM_ADDRESS_ATTRIBUTE`] — stay stable across such renames.
    pub fn auto_detect(mut self, enabled: bool) -> Self {
        self.auto_detect = enabled;
        self
    }

    /// Overrides the per-swap gas cost that auto-detected pAMMs (see
    /// [`auto_detect`](Self::auto_detect)) are served with. Defaults to the maximum over the
    /// known venue profiles, as the conservative choice. Registered venues are unaffected —
    /// their gas cost comes from their [`PriceLevelStreamConfig`].
    pub fn auto_detected_gas_cost(mut self, gas_cost: BigUint) -> Self {
        self.auto_detected_gas_cost = Some(gas_cost);
        self
    }

    /// Overrides the stream endpoint, e.g. to connect to a closer Titan region than the default
    /// (see <https://docs.titanbuilder.xyz/propamms/takers>). Without it, the
    /// `TITAN_PAMM_PRICE_LEVEL_URL` environment variable is used when set, else the built-in
    /// default.
    pub fn endpoint(mut self, url: impl Into<String>) -> Self {
        self.url = Some(url.into());
        self
    }

    /// Overrides how long a single connection attempt may take before it is aborted and retried
    /// (default: 10s), so a hung TCP/TLS handshake cannot block the stream forever.
    pub fn connect_timeout(mut self, timeout: Duration) -> Self {
        self.connection.connect_timeout = timeout;
        self
    }

    /// Overrides the longest gap between parsed Titan frames tolerated before the connection is
    /// treated as dead and re-established (default: 10s). Titan pushes one frame per second and
    /// sends no keepalives, so a multi-second silence means a stalled or half-open connection.
    /// Control frames and unparsable text do not reset this timeout, and neither does the time
    /// the consumer spends between polls: the gap is measured while the stream waits on the
    /// socket.
    pub fn read_idle_timeout(mut self, timeout: Duration) -> Self {
        self.connection.read_idle_timeout = timeout;
        self
    }

    /// Overrides the cap on the exponential reconnect backoff of `2^attempt` seconds
    /// (default: 32s).
    pub fn max_backoff(mut self, max_backoff: Duration) -> Self {
        self.connection.max_backoff = max_backoff;
        self
    }

    /// Registers a pAMM to be served under the given configuration, overriding any default,
    /// denied, or auto-detected one for the same address.
    ///
    /// Between [`add_pamm`](Self::add_pamm) and [`deny_pamm`](Self::deny_pamm) for the same
    /// address, the later call wins; the defaults applied by
    /// [`with_known_pamms`](Self::with_known_pamms) never override either, in any call order.
    pub fn add_pamm(mut self, config: PriceLevelStreamConfig) -> Self {
        self.denied.remove(&config.address);
        self.registry
            .insert(config.address.clone(), config);
        self
    }

    /// Excludes a venue from being served: drops its current registration (default or explicit)
    /// and blocks auto-detecting it.
    ///
    /// Between [`add_pamm`](Self::add_pamm) and [`deny_pamm`](Self::deny_pamm) for the same
    /// address, the later call wins; the defaults applied by
    /// [`with_known_pamms`](Self::with_known_pamms) never override either, in any call order —
    /// so denying a venue from the default set works whether the denial comes before or after
    /// [`with_known_pamms`](Self::with_known_pamms).
    pub fn deny_pamm(mut self, address: Bytes) -> Self {
        self.registry.remove(&address);
        self.denied.insert(address);
        self
    }

    /// Applies what is known about the streamed venues: registers the known-good ones
    /// ([`default_served_pamms`]) to be served and denies the known-bad ones
    /// ([`default_denied_pamms`]) — venues that stream quotes but whose swaps are not executable.
    ///
    /// These defaults never override an explicit [`add_pamm`](Self::add_pamm) or
    /// [`deny_pamm`](Self::deny_pamm) for the same address, regardless of call order.
    pub fn with_known_pamms(mut self) -> Self {
        for config in default_served_pamms() {
            if self.denied.contains(&config.address) {
                continue;
            }
            self.registry
                .entry(config.address.clone())
                .or_insert(config);
        }
        for address in default_denied_pamms() {
            if self.registry.contains_key(&address) {
                continue;
            }
            self.denied.insert(address);
        }
        self
    }

    /// Provides the token metadata used to build components and interpret amounts. Pairs whose
    /// tokens are missing here are skipped.
    pub fn with_tokens(mut self, tokens: HashMap<Bytes, Token>) -> Self {
        self.tokens = tokens;
        self
    }

    /// Keeps every venue on the direct `pricelevelstream:{name}` path, so swaps execute on the
    /// venues themselves and a stale maker quote reverts the route.
    ///
    /// By default components are emitted under `fallback:{name}`, so tycho-execution routes
    /// their swaps through `TychoFallbackRouter`. Opt out when the direct call is what you want
    /// to measure or execute.
    pub fn without_fallback_router(mut self) -> Self {
        self.fallback_router = false;
        self
    }

    /// Overrides how long a component stays served after the last accepted frame that carried
    /// it (default: 24s, two slots). A component no accepted frame has carried for this long
    /// turns stale and is emitted in `removed_pairs`; the next accepted frame carrying it
    /// re-adds it in `new_pairs`. Frames whose `timestamp` is this old or older are rejected.
    /// Independent of this setting, a state refuses to quote once its frame is one slot old
    /// (see [`without_quote_guard`](Self::without_quote_guard)).
    ///
    /// Values above [`MAX_STALE_AFTER`] are capped to it. A value shorter than the age frames
    /// arrive with rejects every frame as `too_old`, which the
    /// `price_level_stream_frames_rejected_total` counter and a WARN log show.
    pub fn stale_after(mut self, duration: Duration) -> Self {
        self.stale_after = duration.min(MAX_STALE_AFTER);
        self
    }

    /// Emits states that never refuse to quote.
    ///
    /// By default every emitted state refuses `spot_price`, `get_amount_out` and `get_limits`
    /// once its frame is one slot ([`QUOTE_TTL`](super::state::QUOTE_TTL)) old: Titan quotes
    /// the block being built, and a venue rejects a fill against an older ladder as stale, so
    /// such a quote is not executable. Opt out only for a consumer that quotes a state more
    /// than one slot after it arrived by design — a batch simulator, a validation harness — and
    /// that accepts a quote the venue may no longer fill. The component still turns stale and
    /// is removed after [`stale_after`](Self::stale_after), which then becomes the only bound
    /// on how old a quoted ladder can be. Never disable it on a live router.
    pub fn without_quote_guard(mut self) -> Self {
        self.quote_guard = false;
        self
    }

    /// Consumes the builder and opens the stream.
    ///
    /// Components are emitted under `fallback:{name}`, so tycho-execution routes their swaps
    /// through `TychoFallbackRouter`, which retries a reverted pAMM swap — a stale maker quote
    /// reverts in any simulation against a mined block — on the fallback pool the solver names.
    /// [`without_fallback_router`](Self::without_fallback_router) keeps them on the direct
    /// `pricelevelstream:` path.
    ///
    /// The connection is established lazily on first poll and maintained (with reconnects) for as
    /// long as the stream is polled; it never terminates on its own, and dropping the stream
    /// closes the connection.
    ///
    /// Every accepted frame yields an update with the states of the served pairs it carries,
    /// with `new_pairs` for pairs not currently served. The update does not mention pairs the
    /// frame does not carry, so consumers keep their previous state. A component no accepted
    /// frame has carried for [`stale_after`](Self::stale_after) turns stale and is emitted in
    /// `removed_pairs`, together with every other component turning stale at that instant, and
    /// re-added by the next accepted frame carrying it. Frames that are too old, from the
    /// future, out of order, or whose block regresses or jumps more than one block per elapsed
    /// slot plus 2 are rejected without effect. Frames that contain no served pAMM produce no
    /// update. Pairs whose tokens are missing from the provided token metadata are skipped.
    ///
    /// Every update is stamped with the block its frame targets, a removal with the newest
    /// accepted block. Block numbers never decrease while something is served. Once every
    /// component has turned stale, the next accepted frame is judged as a first frame and may
    /// carry a lower block than the removal did: that is how the stream recovers from a frame
    /// with an implausible block, so consumers must not rely on the block number to order
    /// updates across such a gap. See the [module documentation](super) for the full contract.
    pub fn build(self) -> impl Stream<Item = Update> + Send {
        let Self {
            registry,
            denied,
            tokens,
            url,
            auto_detect,
            auto_detected_gas_cost,
            connection,
            fallback_router,
            stale_after,
            quote_guard,
        } = self;
        if registry.is_empty() && !auto_detect {
            tracing::warn!(
                "No pAMMs registered and auto-detection is off; the stream will never produce \
                 an update"
            );
        }
        if tokens.is_empty() {
            tracing::warn!(
                "No token metadata provided; every streamed pair will be skipped and the stream \
                 will never produce an update"
            );
        }
        let url = url.unwrap_or_else(|| {
            std::env::var(TITAN_PRICE_LEVEL_URL_ENV)
                .unwrap_or_else(|_| TITAN_PRICE_LEVEL_URL.to_string())
        });
        let auto_detected_gas_cost =
            auto_detected_gas_cost.unwrap_or_else(|| BigUint::from(DEFAULT_AUTO_DETECTED_GAS_COST));
        let mut tracker = FreshnessTracker::new(TrackerSettings {
            registry,
            denied,
            tokens,
            auto_detect,
            auto_detected_gas_cost,
            stale_after,
            via_fallback_router: fallback_router,
            quote_guard,
        });

        stream! {
            let frames = titan::messages(url, connection);
            tokio::pin!(frames);
            loop {
                let deadline = tracker.stale_deadline();
                let sleep_until_deadline = tokio::time::sleep_until(
                    deadline.map_or_else(tokio::time::Instant::now, tokio::time::Instant::from_std),
                );
                let update = tokio::select! {
                    Some(frame) = frames.next() => tracker.on_frame(frame, now()),
                    () = sleep_until_deadline, if deadline.is_some() => {
                        tracker.on_stale_deadline(timer_now())
                    }
                };
                if let Some(update) = update {
                    yield update;
                }
            }
        }
    }
}

/// The clocks a frame is judged against: the wall clock for Titan's `timestamp`, and the
/// timer's clock for every deadline. This is the only place the stream reads a clock; the
/// tracker only ever sees the values it is handed.
fn now() -> Now {
    Now { wall_nanos: wall_clock_nanos(), monotonic: timer_now() }
}

/// Nanoseconds since the Unix epoch. A system clock unrepresentable as unix nanoseconds yields
/// 0, which rejects every frame as `in_future`.
fn wall_clock_nanos() -> u64 {
    let since_epoch = SystemTime::now()
        .duration_since(UNIX_EPOCH)
        .ok()
        .and_then(|since_epoch| u64::try_from(since_epoch.as_nanos()).ok());
    match since_epoch {
        Some(wall_nanos) => wall_nanos,
        None => {
            tracing::error!(
                "System clock unrepresentable as unix nanoseconds; every price level frame will \
                 be rejected as in_future"
            );
            0
        }
    }
}

/// The clock the deadline timer runs on. Deadlines are set and swept with it so that a sweep
/// fired by the timer finds the component due, also under `tokio::time::pause()`, where the
/// runtime's virtual time and `std::time::Instant` diverge.
fn timer_now() -> Instant {
    tokio::time::Instant::now().into_std()
}

#[cfg(test)]
mod tests {
    use std::{
        pin::Pin,
        str::FromStr,
        sync::{
            atomic::{AtomicBool, Ordering},
            Arc,
        },
        time::Duration,
    };

    use futures::{future::BoxFuture, SinkExt};
    use num_bigint::BigUint;
    use rstest::rstest;
    use tokio_tungstenite::tungstenite::Message;

    use super::{
        super::{
            config::{default_denied_pamms, PriceLevelStreamConfig},
            state::PriceLevelStreamState,
            telemetry::{
                recorded::{counter_value, record_async},
                RECONNECTS,
            },
            test_support::{
                fermiswap, frame_text, frame_then_repeat, tokens, wall_nanos_now, FakeConnection,
                FakeTitan, PAMM,
            },
        },
        *,
    };

    #[test]
    fn explicit_add_and_deny_are_last_wins() {
        let address = Bytes::from_str(PAMM).unwrap();
        let custom =
            || PriceLevelStreamConfig::new("custom", Bytes::from_str(PAMM).unwrap(), 1u64.into());

        let builder = PriceLevelStreamBuilder::new()
            .add_pamm(custom())
            .deny_pamm(address.clone());
        assert!(!builder.registry.contains_key(&address));
        assert!(builder.denied.contains(&address));

        let builder = PriceLevelStreamBuilder::new()
            .deny_pamm(address.clone())
            .add_pamm(custom());
        assert_eq!(builder.registry[&address].protocol, "custom");
        assert!(builder.denied.is_empty());
    }

    #[test]
    fn quote_guard_is_on_unless_opted_out() {
        assert!(PriceLevelStreamBuilder::new().quote_guard);
        assert!(
            !PriceLevelStreamBuilder::new()
                .without_quote_guard()
                .quote_guard
        );
    }

    /// A builder with short timings: `stale_after` 2 s, `read_idle_timeout` 100 ms,
    /// `max_backoff` 20 ms.
    fn fast_builder(fake: &FakeTitan) -> PriceLevelStreamBuilder {
        PriceLevelStreamBuilder::new()
            .endpoint(fake.url())
            .add_pamm(fermiswap())
            .with_tokens(tokens())
            .stale_after(STALE_AFTER)
            .connect_timeout(Duration::from_secs(1))
            .read_idle_timeout(Duration::from_millis(100))
            .max_backoff(Duration::from_millis(20))
    }

    /// The `stale_after` of [`fast_builder`]: long enough that a loaded runner does not reject
    /// the first frame as too old, short enough that a removal arrives within [`WAIT`].
    const STALE_AFTER: Duration = Duration::from_secs(2);

    /// How long a test waits for one update.
    const WAIT: Duration = Duration::from_secs(5);

    async fn next_within(
        stream: &mut Pin<&mut impl Stream<Item = Update>>,
        limit: Duration,
    ) -> Option<Update> {
        tokio::time::timeout(limit, stream.next())
            .await
            .ok()
            .flatten()
    }

    /// Waits for the first update and checks that it adds the one served component.
    async fn expect_first_update(stream: &mut Pin<&mut impl Stream<Item = Update>>) -> Update {
        let first = next_within(stream, WAIT)
            .await
            .expect("first update");
        assert_eq!(first.new_pairs.len(), 1);
        assert!(first.removed_pairs.is_empty());
        first
    }

    /// Waits up to [`WAIT`] for an update that removes components, skipping the updates that
    /// only refresh served ones; an update that re-adds a component before the removal fails.
    async fn expect_removal(stream: &mut Pin<&mut impl Stream<Item = Update>>) -> Update {
        tokio::time::timeout(WAIT, async {
            loop {
                let update = stream
                    .next()
                    .await
                    .expect("stream ended");
                if !update.removed_pairs.is_empty() {
                    return update;
                }
                assert!(update.new_pairs.is_empty(), "component re-added before the removal");
            }
        })
        .await
        .expect("removal")
    }

    fn assert_removal_only(update: &Update, expected_removed: usize) {
        assert!(update.states.is_empty());
        assert!(update.new_pairs.is_empty());
        assert!(update.sync_states.is_empty());
        assert!(update.is_partial);
        assert_eq!(update.removed_pairs.len(), expected_removed);
    }

    fn fresh_frame() -> Message {
        Message::Text(frame_text(100, wall_nanos_now()).into())
    }

    /// Sends one fresh frame on the first connection only; later connections stay silent.
    async fn first_connection_sends_one_frame(index: usize, mut socket: FakeConnection) {
        if index == 0 {
            let _ = socket.send(fresh_frame()).await;
        }
        std::future::pending::<()>().await;
    }

    #[tokio::test]
    async fn silence_past_stale_after_removes_every_served_component() {
        let fake = FakeTitan::spawn(first_connection_sends_one_frame).await;
        let stream = fast_builder(&fake).build();
        tokio::pin!(stream);

        expect_first_update(&mut stream).await;
        // The stream reconnects on idle timeout, but the later connections send no frame, so no
        // component deadline moves.
        let removal = next_within(&mut stream, WAIT)
            .await
            .expect("removal");

        assert_removal_only(&removal, 1);
        assert_eq!(removal.block_number_or_timestamp, 100);
        assert!(fake.connections.load(Ordering::SeqCst) >= 2, "no reconnect on idle timeout");
    }

    #[tokio::test]
    async fn repeated_immediate_closes_remove_within_stale_after() {
        let fake = FakeTitan::spawn(|index, mut socket| async move {
            if index == 0 {
                let _ = socket.send(fresh_frame()).await;
                tokio::time::sleep(Duration::from_millis(20)).await;
            }
            let _ = socket.close(None).await;
        })
        .await;
        let stream = fast_builder(&fake).build();
        tokio::pin!(stream);

        expect_first_update(&mut stream).await;
        let removal = next_within(&mut stream, WAIT)
            .await
            .expect("removal");

        assert_removal_only(&removal, 1);
    }

    #[test]
    fn refused_reconnects_remove_within_stale_after() {
        let (removal, snapshot) = record_async(async {
            let mut fake = FakeTitan::spawn(|_, mut socket| async move {
                let _ = socket.send(fresh_frame()).await;
                tokio::time::sleep(Duration::from_millis(20)).await;
                let _ = socket.close(None).await;
            })
            .await;
            let stream = fast_builder(&fake).build();
            tokio::pin!(stream);

            expect_first_update(&mut stream).await;
            // From here every connect is refused at TCP level. The fake closes 20 ms in and the
            // backoff is capped at 20 ms, so one reconnect can be accepted before `shutdown`
            // drops the listener; its fresh frame refreshes the component without re-adding it.
            fake.shutdown();
            expect_removal(&mut stream).await
        });

        assert_removal_only(&removal, 1);
        assert!(
            counter_value(&snapshot, RECONNECTS, &[("reason", "connect_failed")]) >= 1,
            "no connect was refused"
        );
    }

    #[tokio::test]
    async fn replayed_frames_remove_within_stale_after_and_never_re_add() {
        // One frame, stamped once, replayed every 50 ms forever.
        let replay = fresh_frame();
        let fake =
            FakeTitan::spawn(frame_then_repeat(replay.clone(), replay, Duration::from_millis(50)))
                .await;
        let stream = fast_builder(&fake)
            .read_idle_timeout(Duration::from_secs(5))
            .build();
        tokio::pin!(stream);

        expect_first_update(&mut stream).await;
        // Replays are accepted while fresh but cannot extend the deadline; the deadline fires
        // even though a frame is ready on every poll.
        let removal = expect_removal(&mut stream).await;

        assert_removal_only(&removal, 1);
        // Every later replay is too old to be accepted: nothing comes back.
        assert!(next_within(&mut stream, Duration::from_millis(500))
            .await
            .is_none());
    }

    #[tokio::test]
    async fn fresh_frame_after_removal_re_adds_the_component() {
        let fake = FakeTitan::spawn(|_, mut socket| async move {
            let _ = socket.send(fresh_frame()).await;
            tokio::time::sleep(STALE_AFTER + Duration::from_millis(500)).await;
            let _ = socket
                .send(Message::Text(frame_text(101, wall_nanos_now()).into()))
                .await;
            std::future::pending::<()>().await;
        })
        .await;
        let stream = fast_builder(&fake)
            .read_idle_timeout(Duration::from_secs(5))
            .build();
        tokio::pin!(stream);

        expect_first_update(&mut stream).await;
        let removal = next_within(&mut stream, WAIT)
            .await
            .expect("removal");
        let re_added = next_within(&mut stream, WAIT)
            .await
            .expect("re-add");

        assert_removal_only(&removal, 1);
        assert_eq!(re_added.new_pairs.len(), 1);
        assert!(re_added.removed_pairs.is_empty());
        assert_eq!(re_added.block_number_or_timestamp, 101);
    }

    #[tokio::test]
    async fn frames_are_forwarded_without_waiting_on_timers() {
        let fake = FakeTitan::spawn(fresh_frame_every(Duration::from_millis(50))).await;
        let stream = fast_builder(&fake)
            .stale_after(Duration::from_secs(24))
            .build();
        tokio::pin!(stream);

        expect_first_update(&mut stream).await;
        for _ in 0..5 {
            let update = next_within(&mut stream, Duration::from_secs(1))
                .await
                .expect("steady-state frame");
            assert!(update.removed_pairs.is_empty());
        }
    }

    #[tokio::test]
    async fn no_connection_before_first_poll() {
        let fake = FakeTitan::spawn(first_connection_sends_one_frame).await;
        let stream = fast_builder(&fake).build();
        tokio::pin!(stream);

        tokio::time::sleep(Duration::from_millis(150)).await;

        assert_eq!(fake.connections.load(Ordering::SeqCst), 0, "connected before first poll");
    }

    #[tokio::test]
    async fn drop_closes_the_socket() {
        let server_saw_close = Arc::new(AtomicBool::new(false));
        let fake = {
            let server_saw_close = server_saw_close.clone();
            FakeTitan::spawn(move |_, mut socket| {
                let server_saw_close = server_saw_close.clone();
                async move {
                    let _ = socket.send(fresh_frame()).await;
                    // Read until the client goes away.
                    while let Some(Ok(message)) = socket.next().await {
                        if matches!(message, Message::Close(_)) {
                            break;
                        }
                    }
                    server_saw_close.store(true, Ordering::SeqCst);
                }
            })
            .await
        };
        // `Box::pin` rather than `tokio::pin!`, so that `drop` below drops the stream itself
        // rather than a `Pin<&mut _>` pointing at a value that outlives the assertion.
        let mut stream = Box::pin(
            fast_builder(&fake)
                .read_idle_timeout(Duration::from_secs(5))
                .build(),
        );
        expect_first_update(&mut stream.as_mut()).await;

        drop(stream);
        tokio::time::sleep(Duration::from_millis(200)).await;

        assert!(server_saw_close.load(Ordering::SeqCst), "socket not closed on drop");
        assert_eq!(fake.connections.load(Ordering::SeqCst), 1, "reconnected after drop");
    }

    /// `Instant + Duration` panics on overflow, so an absurd `stale_after` must be capped before
    /// the first accepted frame computes a deadline from it.
    #[tokio::test]
    async fn stale_after_above_the_cap_still_serves() {
        assert_eq!(
            PriceLevelStreamBuilder::new()
                .stale_after(Duration::MAX)
                .stale_after,
            MAX_STALE_AFTER
        );
        let fake = FakeTitan::spawn(first_connection_sends_one_frame).await;
        let stream = fast_builder(&fake)
            .stale_after(Duration::MAX)
            .build();
        tokio::pin!(stream);

        expect_first_update(&mut stream).await;
    }

    #[tokio::test]
    async fn without_quote_guard_emits_states_that_never_expire() {
        let fake = FakeTitan::spawn(first_connection_sends_one_frame).await;
        let stream = fast_builder(&fake)
            .without_quote_guard()
            .build();
        tokio::pin!(stream);

        let first = expect_first_update(&mut stream).await;

        let state = first
            .states
            .values()
            .next()
            .expect("one state")
            .as_any()
            .downcast_ref::<PriceLevelStreamState>()
            .expect("price level state");
        assert!(state.quotable_until().is_none());
    }

    /// A [`FakeTitan`] handler that sends a freshly stamped frame every `interval` until the
    /// socket closes.
    fn fresh_frame_every(
        interval: Duration,
    ) -> impl Fn(usize, FakeConnection) -> BoxFuture<'static, ()> + Send + Sync + 'static {
        move |_, mut socket| {
            Box::pin(async move {
                while socket.send(fresh_frame()).await.is_ok() {
                    tokio::time::sleep(interval).await;
                }
            })
        }
    }

    /// The fallback router path is the default; `without_fallback_router` is the way off it.
    #[test]
    fn fallback_router_is_on_unless_opted_out() {
        assert!(PriceLevelStreamBuilder::new().fallback_router);
        assert!(
            !PriceLevelStreamBuilder::new()
                .without_fallback_router()
                .fallback_router
        );
    }

    #[test]
    fn defaults_never_override_explicit_calls() {
        // Denying a venue from the default set works in either call order.
        let fermiswap_router = Bytes::from_str(PAMM).unwrap();
        for builder in [
            PriceLevelStreamBuilder::new()
                .deny_pamm(fermiswap_router.clone())
                .with_known_pamms(),
            PriceLevelStreamBuilder::new()
                .with_known_pamms()
                .deny_pamm(fermiswap_router.clone()),
        ] {
            assert!(!builder
                .registry
                .contains_key(&fermiswap_router));
            assert!(builder
                .denied
                .contains(&fermiswap_router));
            // The other defaults are unaffected.
            assert!(!builder.registry.is_empty());
        }

        // Registering a venue from the default deny set works in either call order. Any one of
        // them exercises that; the set is empty while every streamed venue is executable.
        let Some(denied_venue) = default_denied_pamms().pop() else { return };
        let custom = || PriceLevelStreamConfig::new("custom", denied_venue.clone(), 1u64.into());
        for builder in [
            PriceLevelStreamBuilder::new()
                .add_pamm(custom())
                .with_known_pamms(),
            PriceLevelStreamBuilder::new()
                .with_known_pamms()
                .add_pamm(custom()),
        ] {
            assert_eq!(builder.registry[&denied_venue].protocol, "custom");
            assert!(!builder.denied.contains(&denied_venue));
        }
    }

    #[test]
    fn with_known_pamms_registers_known_venues() {
        // PAMM is the FermiSwap router, one of the default venues.
        let fermiswap_router = Bytes::from_str(PAMM).unwrap();

        let builder = PriceLevelStreamBuilder::new();
        assert!(builder.registry.is_empty());
        assert!(builder.denied.is_empty());

        let builder = builder.with_known_pamms();
        assert_eq!(builder.registry[&fermiswap_router].protocol, "fermiswap");
        // The known-bad venues get denied alongside, and never overlap the served defaults.
        assert_eq!(
            builder.denied,
            default_denied_pamms()
                .into_iter()
                .collect()
        );
        assert!(builder.denied.is_disjoint(
            &builder
                .registry
                .keys()
                .cloned()
                .collect()
        ));

        // An `add_pamm` entry wins over the default for the same address, in either call order.
        let custom =
            || PriceLevelStreamConfig::new("custom", fermiswap_router.clone(), BigUint::from(1u64));
        for builder in [
            PriceLevelStreamBuilder::new()
                .add_pamm(custom())
                .with_known_pamms(),
            PriceLevelStreamBuilder::new()
                .with_known_pamms()
                .add_pamm(custom()),
        ] {
            assert_eq!(builder.registry[&fermiswap_router].protocol, "custom");
            assert_eq!(builder.registry[&fermiswap_router].gas_cost, BigUint::from(1u64));
        }
    }

    /// The families this stream emits are the ones tycho-execution resolves an encoder for. A
    /// drift between the two makes every route through a pAMM fail to encode.
    #[test]
    fn families_match_the_execution_side_prefixes() {
        use tycho_execution::encoding::evm::{FALLBACK_PREFIX, PRICE_LEVEL_STREAM_PREFIX};

        use super::super::config::{FALLBACK_FAMILY, PRICE_LEVEL_STREAM_FAMILY};

        assert_eq!(format!("{PRICE_LEVEL_STREAM_FAMILY}:"), PRICE_LEVEL_STREAM_PREFIX);
        assert_eq!(format!("{FALLBACK_FAMILY}:"), FALLBACK_PREFIX);
    }

    /// Traffic that is not a parsed frame keeps neither the socket nor the component alive: the
    /// idle timeout reconnects, every reconnect resends the same frame, which refreshes the
    /// component but cannot move its deadline, and the removal follows.
    #[rstest]
    #[case::ping_only(Message::Ping(Vec::new().into()))]
    #[case::malformed_text(Message::Text("nonsense".into()))]
    #[tokio::test]
    async fn non_frame_traffic_removes_within_stale_after(#[case] filler: Message) {
        let fake =
            FakeTitan::spawn(frame_then_repeat(fresh_frame(), filler, Duration::from_millis(10)))
                .await;
        let stream = fast_builder(&fake).build();
        tokio::pin!(stream);

        expect_first_update(&mut stream).await;
        let removal = expect_removal(&mut stream).await;

        assert_removal_only(&removal, 1);
        assert!(fake.connections.load(Ordering::SeqCst) >= 2, "no reconnect on idle timeout");
    }
}