use std::{
collections::{BTreeSet, HashMap, HashSet},
time::{Duration, Instant},
};
use chrono::Utc;
use num_bigint::BigUint;
use tycho_common::{
models::{token::Token, Chain},
simulation::protocol_sim::ProtocolSim,
Bytes,
};
use super::{
config::PriceLevelStreamConfig,
state::{PriceLevelStreamQuote, PriceLevelStreamState, QUOTE_TTL},
stream::PAMM_ADDRESS_ATTRIBUTE,
telemetry::{self, RejectReason},
titan::{TitanPairLevels, TitanPammLevels, TitanPriceLevel, TitanPriceLevelMessage},
SLOT,
};
use crate::protocol::models::{ProtocolComponent, Update};
pub(super) const DEFAULT_STALE_AFTER: Duration = Duration::from_secs(2 * SLOT.as_secs());
pub(super) const NANOS_PER_SECOND: u64 = 1_000_000_000;
const MAX_FUTURE_SKEW_NANOS: u64 = SLOT.as_secs() * NANOS_PER_SECOND;
const BLOCK_JUMP_SLACK: u64 = 2;
const MAX_UNREGISTERED_LOGGED: usize = 64;
const MAX_AUTO_DETECTED: usize = 64;
#[derive(Clone, Copy, Debug)]
pub(super) struct Now {
pub wall_nanos: u64,
pub monotonic: Instant,
}
struct ServedComponent {
component: ProtocolComponent,
venue_name: String,
stale_at: Instant,
}
#[derive(Clone, Copy, Debug)]
struct Frontier {
block: u64,
accepted_at: Instant,
}
#[derive(Clone, Copy)]
enum ServingState {
Unserved,
Serving { deadline: Instant, frontier: Frontier },
}
impl ServingState {
fn telemetry_state(self) -> telemetry::ServingState {
match self {
ServingState::Unserved => telemetry::ServingState::Unserved,
ServingState::Serving { deadline: _, frontier: _ } => telemetry::ServingState::Serving,
}
}
}
pub(super) struct TrackerSettings {
pub registry: HashMap<Bytes, PriceLevelStreamConfig>,
pub denied: HashSet<Bytes>,
pub tokens: HashMap<Bytes, Token>,
pub auto_detect: bool,
pub auto_detected_gas_cost: BigUint,
pub stale_after: Duration,
pub via_fallback_router: bool,
pub quote_guard: bool,
}
pub(super) struct FreshnessTracker {
registry: HashMap<Bytes, PriceLevelStreamConfig>,
denied: HashSet<Bytes>,
tokens: HashMap<Bytes, Token>,
auto_detect: bool,
auto_detected_gas_cost: BigUint,
stale_after: Duration,
via_fallback_router: bool,
quote_guard: bool,
served: HashMap<String, ServedComponent>,
serving_state: ServingState,
newest_timestamp_nanos: u64,
rejecting: bool,
logged_unregistered: HashSet<Bytes>,
auto_detected: usize,
}
impl FreshnessTracker {
pub(super) fn new(settings: TrackerSettings) -> Self {
let TrackerSettings {
registry,
denied,
tokens,
auto_detect,
auto_detected_gas_cost,
stale_after,
via_fallback_router,
quote_guard,
} = settings;
telemetry::record_serving_state(telemetry::ServingState::Unserved);
for config in registry.values() {
telemetry::record_served_components(&config.protocol, 0);
telemetry::record_last_seen(&config.protocol, 0);
}
Self {
registry,
denied,
tokens,
auto_detect,
auto_detected_gas_cost,
stale_after,
via_fallback_router,
quote_guard,
served: HashMap::new(),
serving_state: ServingState::Unserved,
newest_timestamp_nanos: 0,
rejecting: false,
logged_unregistered: HashSet::new(),
auto_detected: 0,
}
}
fn check_frame(
&self,
frame: &TitanPriceLevelMessage,
now: Now,
) -> Result<Duration, RejectReason> {
let frame_age = Duration::from_nanos(
now.wall_nanos
.saturating_sub(frame.timestamp),
);
if frame_age >= self.stale_after {
return Err(RejectReason::TooOld);
}
if frame.timestamp >
now.wall_nanos
.saturating_add(MAX_FUTURE_SKEW_NANOS)
{
return Err(RejectReason::InFuture);
}
if frame.timestamp < self.newest_timestamp_nanos {
return Err(RejectReason::OutOfOrder);
}
let ServingState::Serving { deadline: _, frontier } = self.serving_state else {
return Ok(frame_age);
};
if frame.block_number < frontier.block {
return Err(RejectReason::BlockRegression);
}
let elapsed_slots = now
.monotonic
.saturating_duration_since(frontier.accepted_at)
.as_secs() /
SLOT.as_secs();
let allowed = frontier
.block
.saturating_add(elapsed_slots)
.saturating_add(BLOCK_JUMP_SLACK);
if frame.block_number > allowed {
return Err(RejectReason::BlockJump);
}
Ok(frame_age)
}
pub(super) fn on_frame(&mut self, frame: TitanPriceLevelMessage, now: Now) -> Option<Update> {
let frame_age = match self.check_frame(&frame, now) {
Ok(frame_age) => frame_age,
Err(reason) => {
self.log_rejection(reason, &frame);
telemetry::record_frame_rejected(reason);
return None;
}
};
self.rejecting = false;
telemetry::record_frame_accepted();
let age_nanos = i128::from(now.wall_nanos) - i128::from(frame.timestamp);
telemetry::record_frame_age(age_nanos as f64 / NANOS_PER_SECOND as f64);
self.newest_timestamp_nanos = frame.timestamp;
let stale_at = now.monotonic +
self.stale_after
.saturating_sub(frame_age);
let quotable_until = self
.quote_guard
.then(|| now.monotonic + QUOTE_TTL.saturating_sub(frame_age));
let frame_unix_seconds = frame.timestamp / NANOS_PER_SECOND;
let mut states: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
let mut new_pairs = HashMap::new();
for TitanPammLevels { pamm, pairs } in frame.pamms {
self.admit(&pamm);
let Some(config) = self.registry.get(&pamm) else {
continue;
};
telemetry::record_last_seen(&config.protocol, frame_unix_seconds);
for ((token0, token1), (quotes_0_to_1, quotes_1_to_0)) in
merge_pairs(&self.tokens, pairs)
{
let id = component_id(&config.address, &token0, &token1);
let id_string = id.to_string();
match self.served.get_mut(&id_string) {
Some(served) => served.stale_at = stale_at,
None => {
let component = build_component(
&self.tokens,
config,
id,
&token0,
&token1,
self.via_fallback_router,
);
new_pairs.insert(id_string.clone(), component.clone());
self.served.insert(
id_string.clone(),
ServedComponent {
component,
venue_name: config.protocol.clone(),
stale_at,
},
);
}
}
let mut state = PriceLevelStreamState::new(
token0,
token1,
quotes_0_to_1,
quotes_1_to_0,
config.gas_cost.clone(),
);
if let Some(until) = quotable_until {
state = state.with_quotable_until(until);
}
states.insert(id_string, Box::new(state));
}
}
self.settle(Frontier { block: frame.block_number, accepted_at: now.monotonic });
if states.is_empty() {
return None;
}
Some(Update::new(frame.block_number, states, new_pairs).set_is_partial(true))
}
fn admit(&mut self, pamm: &Bytes) {
if self.registry.contains_key(pamm) {
return;
}
if self.denied.contains(pamm) {
tracing::debug!(%pamm, "Skipping denied pAMM");
return;
}
if self.auto_detect && self.auto_detected < MAX_AUTO_DETECTED {
tracing::info!(%pamm, "Serving auto-detected pAMM");
let config = PriceLevelStreamConfig::auto_detected(
pamm.clone(),
self.auto_detected_gas_cost.clone(),
);
telemetry::record_served_components(&config.protocol, 0);
telemetry::record_last_seen(&config.protocol, 0);
self.registry
.insert(pamm.clone(), config);
self.auto_detected += 1;
return;
}
telemetry::record_unregistered_pamm();
if self.logged_unregistered.len() < MAX_UNREGISTERED_LOGGED &&
self.logged_unregistered
.insert(pamm.clone())
{
if self.auto_detect {
tracing::warn!(
%pamm,
cap = MAX_AUTO_DETECTED,
"Skipping unknown pAMM: the auto-detection cap is reached; register it via \
add_pamm to serve it"
);
} else {
tracing::info!(
%pamm,
"Skipping unregistered pAMM; register it via add_pamm to serve it"
);
}
}
}
pub(super) fn on_stale_deadline(&mut self, now: Instant) -> Option<Update> {
let ServingState::Serving { deadline: _, frontier } = self.serving_state else {
return None;
};
let mut removed = HashMap::new();
let mut venues = BTreeSet::new();
for (id, served) in self
.served
.extract_if(|_, served| served.stale_at <= now)
{
telemetry::record_stale_removal(&served.venue_name);
venues.insert(served.venue_name);
removed.insert(id, served.component);
}
if removed.is_empty() {
return None;
}
let component_ids: Vec<&String> = removed.keys().collect();
tracing::warn!(
removed = removed.len(),
venues = ?venues,
components = ?component_ids,
stale_after_secs = self.stale_after.as_secs(),
"Removing price level components: no accepted frame carried them within stale_after"
);
let update = Update::new(frontier.block, HashMap::new(), HashMap::new())
.set_is_partial(true)
.set_removed_pairs(removed);
self.settle(frontier);
Some(update)
}
pub(super) fn stale_deadline(&self) -> Option<Instant> {
match self.serving_state {
ServingState::Serving { deadline, frontier: _ } => Some(deadline),
ServingState::Unserved => None,
}
}
fn settle(&mut self, frontier: Frontier) {
let deadline = self
.served
.values()
.map(|served| served.stale_at)
.min();
self.serving_state = match deadline {
Some(deadline) => ServingState::Serving { deadline, frontier },
None => ServingState::Unserved,
};
telemetry::record_serving_state(self.serving_state.telemetry_state());
let mut per_venue: HashMap<&str, usize> = self
.registry
.values()
.map(|config| (config.protocol.as_str(), 0))
.collect();
for served in self.served.values() {
*per_venue
.entry(served.venue_name.as_str())
.or_default() += 1;
}
for (venue, count) in per_venue {
telemetry::record_served_components(venue, count);
}
}
fn log_rejection(&mut self, reason: RejectReason, frame: &TitanPriceLevelMessage) {
let reason: &'static str = reason.into();
if self.rejecting {
tracing::debug!(
reason,
block_number = frame.block_number,
timestamp = frame.timestamp,
"Rejecting price level frame"
);
return;
}
self.rejecting = true;
let newest_block = match self.serving_state {
ServingState::Serving { deadline: _, frontier } => Some(frontier.block),
ServingState::Unserved => None,
};
tracing::warn!(
reason,
block_number = frame.block_number,
timestamp = frame.timestamp,
newest_block = ?newest_block,
newest_timestamp_nanos = self.newest_timestamp_nanos,
"Rejecting price level frame; further rejections logged at debug until a frame is \
accepted"
);
}
}
fn merge_pairs(
tokens: &HashMap<Bytes, Token>,
pairs: Vec<TitanPairLevels>,
) -> HashMap<(Bytes, Bytes), (Vec<PriceLevelStreamQuote>, Vec<PriceLevelStreamQuote>)> {
let mut merged: HashMap<(Bytes, Bytes), (Vec<_>, Vec<_>)> = HashMap::new();
for TitanPairLevels { token_in, token_out, order_book } in pairs {
if !tokens.contains_key(&token_in) || !tokens.contains_key(&token_out) {
tracing::debug!(%token_in, %token_out, "Skipping pair with unknown token");
continue;
}
let sells_token0 = token_in < token_out;
let key = if sells_token0 {
(token_in.clone(), token_out.clone())
} else {
(token_out.clone(), token_in.clone())
};
let quotes = order_book
.into_iter()
.map(|TitanPriceLevel { amount_in, amount_out }| {
PriceLevelStreamQuote::new(amount_in, amount_out)
})
.collect();
let entry = merged.entry(key).or_default();
if sells_token0 {
entry.0 = quotes;
} else {
entry.1 = quotes;
}
}
merged
}
fn build_component(
tokens: &HashMap<Bytes, Token>,
config: &PriceLevelStreamConfig,
id: Bytes,
token0: &Bytes,
token1: &Bytes,
via_router: bool,
) -> ProtocolComponent {
let protocol_system =
if via_router { config.fallback_protocol_system() } else { config.protocol_system() };
ProtocolComponent::new(
id,
protocol_system.clone(),
protocol_system,
Chain::Ethereum,
vec![tokens[token0].clone(), tokens[token1].clone()],
vec![config.address.clone()],
HashMap::from([(PAMM_ADDRESS_ATTRIBUTE.to_string(), config.address.clone())]),
Bytes::default(),
Utc::now().naive_utc(),
)
}
fn component_id(pamm: &Bytes, token0: &Bytes, token1: &Bytes) -> Bytes {
Bytes::from([pamm.as_ref(), token0.as_ref(), token1.as_ref()].concat())
}
#[cfg(test)]
mod tests {
use std::{
str::FromStr,
time::{Duration, Instant},
};
use rstest::rstest;
use super::{
super::{
config::DEFAULT_AUTO_DETECTED_GAS_COST,
state::QUOTE_TTL,
telemetry::{
recorded::{counter_value, gauge_value, histogram_values, record_async},
FRAMES_ACCEPTED, FRAMES_REJECTED, FRAME_AGE, LAST_SEEN, SERVED_COMPONENTS,
SERVING_STATE, STALE_REMOVALS, UNREGISTERED_PAMM_ENTRIES,
},
test_support::*,
},
*,
};
const BASE_WALL_NANOS: u64 = 1_788_624_558_000_000_000;
struct Clock {
start: Instant,
}
impl Clock {
fn new() -> Self {
Self { start: Instant::now() }
}
fn at(&self, seconds: u64) -> Now {
Now {
wall_nanos: BASE_WALL_NANOS + seconds * NANOS_PER_SECOND,
monotonic: self.start + Duration::from_secs(seconds),
}
}
}
fn settings(configs: Vec<PriceLevelStreamConfig>) -> TrackerSettings {
TrackerSettings {
registry: configs
.into_iter()
.map(|config| (config.address.clone(), config))
.collect(),
denied: HashSet::new(),
tokens: tokens(),
auto_detect: false,
auto_detected_gas_cost: BigUint::from(DEFAULT_AUTO_DETECTED_GAS_COST),
stale_after: DEFAULT_STALE_AFTER,
via_fallback_router: false,
quote_guard: true,
}
}
fn tracker_serving(
configs: Vec<PriceLevelStreamConfig>,
via_fallback_router: bool,
) -> FreshnessTracker {
FreshnessTracker::new(TrackerSettings { via_fallback_router, ..settings(configs) })
}
fn tracker() -> FreshnessTracker {
tracker_serving(vec![fermiswap()], false)
}
fn level(amount_in: u64, amount_out: u64) -> TitanPriceLevel {
TitanPriceLevel {
amount_in: BigUint::from(amount_in),
amount_out: BigUint::from(amount_out),
}
}
fn pair_levels(
token_in: &str,
token_out: &str,
order_book: Vec<TitanPriceLevel>,
) -> TitanPairLevels {
TitanPairLevels {
token_in: Bytes::from_str(token_in).unwrap(),
token_out: Bytes::from_str(token_out).unwrap(),
order_book,
}
}
fn message_at(
block_number: u64,
seconds: u64,
pairs: Vec<TitanPairLevels>,
) -> TitanPriceLevelMessage {
TitanPriceLevelMessage {
block_number,
timestamp: BASE_WALL_NANOS + seconds * NANOS_PER_SECOND,
pamms: vec![TitanPammLevels { pamm: Bytes::from_str(PAMM).unwrap(), pairs }],
}
}
fn message(block_number: u64, pairs: Vec<TitanPairLevels>) -> TitanPriceLevelMessage {
message_at(block_number, 0, pairs)
}
fn wbtc_usdc_pairs() -> Vec<TitanPairLevels> {
vec![
pair_levels(WBTC, USDC, vec![level(100_000_000, 100_000_000_000)]),
pair_levels(USDC, WBTC, vec![level(100_000_000_000, 99_000_000)]),
]
}
fn expected_id() -> String {
format!("{PAMM}{}{}", &WBTC[2..], &USDC[2..])
}
fn state_of(update: &Update) -> &PriceLevelStreamState {
update.states[&expected_id()]
.as_any()
.downcast_ref::<PriceLevelStreamState>()
.expect("price level state")
}
fn quotable_until(update: &Update) -> Instant {
state_of(update)
.quotable_until()
.expect("quotable_until set")
}
#[test]
fn first_frame_emits_new_pair_with_both_directions() {
let clock = Clock::new();
let mut tracker = tracker();
let Update {
block_number_or_timestamp,
is_partial,
sync_states,
states,
new_pairs,
removed_pairs,
} = tracker
.on_frame(message(100, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
assert_eq!(block_number_or_timestamp, 100);
assert!(is_partial);
assert!(sync_states.is_empty());
assert!(removed_pairs.is_empty());
let id = expected_id();
let component = &new_pairs[&id];
assert_eq!(component.protocol_system, "pricelevelstream:fermiswap");
assert_eq!(
component.static_attributes[PAMM_ADDRESS_ATTRIBUTE],
Bytes::from_str(PAMM).unwrap()
);
let state = states[&id]
.as_any()
.downcast_ref::<PriceLevelStreamState>()
.expect("price level state");
assert_eq!(state.token0, Bytes::from_str(WBTC).unwrap());
assert_eq!(state.token1, Bytes::from_str(USDC).unwrap());
assert_eq!(state.quotes_0_to_1.len(), 1);
assert_eq!(state.quotes_1_to_0.len(), 1);
assert_eq!(state.quotes_0_to_1[0].amount_in, BigUint::from(100_000_000u64));
assert_eq!(state.gas_cost, BigUint::from(120_000u64));
assert!(state.quotable_until().is_some());
}
#[test]
fn repeated_frame_is_not_a_new_pair() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message(100, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
let update = tracker
.on_frame(message(101, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
assert!(update.new_pairs.is_empty());
assert!(update.removed_pairs.is_empty());
assert!(update
.states
.contains_key(&expected_id()));
}
#[test]
fn block_regression_is_rejected() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(101, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
let regressed =
vec![pair_levels(WETH, USDC, vec![level(1_000_000_000_000_000_000, 3_000_000_000)])];
assert!(tracker
.on_frame(message_at(100, 1, regressed), clock.at(1))
.is_none());
let update = tracker
.on_frame(message_at(102, 2, wbtc_usdc_pairs()), clock.at(2))
.expect("update expected");
assert!(update.new_pairs.is_empty());
}
#[rstest]
#[case::at_the_limit(0, false)]
#[case::inside_the_limit(1, true)]
fn frame_age_at_the_stale_after_limit(#[case] built_at: u64, #[case] accepted: bool) {
let clock = Clock::new();
let mut tracker = tracker();
let update = tracker.on_frame(message_at(100, built_at, wbtc_usdc_pairs()), clock.at(24));
assert_eq!(update.is_some(), accepted);
}
#[rstest]
#[case::beyond_the_limit(13, false)]
#[case::at_the_limit(12, true)]
fn frame_timestamp_at_the_future_skew_limit(#[case] built_at: u64, #[case] accepted: bool) {
let clock = Clock::new();
let mut tracker = tracker();
let update = tracker.on_frame(message_at(100, built_at, wbtc_usdc_pairs()), clock.at(0));
assert_eq!(update.is_some(), accepted);
}
#[test]
fn older_timestamp_frame_is_rejected() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 5, wbtc_usdc_pairs()), clock.at(5))
.expect("update expected");
assert!(tracker
.on_frame(message_at(100, 4, wbtc_usdc_pairs()), clock.at(5))
.is_none());
}
#[test]
fn equal_timestamp_frame_is_accepted() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 5, wbtc_usdc_pairs()), clock.at(5))
.expect("update expected");
let changed = vec![
pair_levels(WBTC, USDC, vec![level(100_000_000, 101_000_000_000)]),
pair_levels(USDC, WBTC, vec![level(100_000_000_000, 99_000_000)]),
];
let update = tracker
.on_frame(message_at(100, 5, changed), clock.at(5))
.expect("update expected");
assert!(update
.states
.contains_key(&expected_id()));
}
#[test]
fn block_jump_beyond_the_elapsed_bound_is_rejected() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
assert!(tracker
.on_frame(message_at(103, 1, wbtc_usdc_pairs()), clock.at(1))
.is_none());
assert!(tracker
.on_frame(message_at(u64::MAX / 2, 1, wbtc_usdc_pairs()), clock.at(1))
.is_none());
assert!(tracker
.on_frame(message_at(101, 2, wbtc_usdc_pairs()), clock.at(2))
.is_some());
}
#[test]
fn block_jump_within_the_elapsed_bound_is_accepted() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
assert!(tracker
.on_frame(message_at(107, 60, wbtc_usdc_pairs()), clock.at(60))
.is_some());
}
#[test]
fn rejections_are_counted_by_reason() {
let ((), snapshot) = record_async(async {
let clock = Clock::new();
let mut tracker = tracker();
tracker.on_frame(message_at(100, 5, wbtc_usdc_pairs()), clock.at(5));
tracker.on_frame(message_at(100, 4, wbtc_usdc_pairs()), clock.at(20));
tracker.on_frame(message_at(100, 5, wbtc_usdc_pairs()), clock.at(29)); tracker.on_frame(message_at(100, 45, wbtc_usdc_pairs()), clock.at(29)); tracker.on_frame(message_at(99, 6, wbtc_usdc_pairs()), clock.at(29)); tracker.on_frame(message_at(200, 6, wbtc_usdc_pairs()), clock.at(29)); });
assert_eq!(counter_value(&snapshot, FRAMES_ACCEPTED, &[]), 1);
for reason in ["too_old", "in_future", "out_of_order", "block_regression", "block_jump"] {
assert_eq!(
counter_value(&snapshot, FRAMES_REJECTED, &[("reason", reason)]),
1,
"{reason}"
);
}
}
#[test]
fn accepted_frame_records_its_signed_age() {
let ((), snapshot) = record_async(async {
let clock = Clock::new();
let mut tracker = tracker();
tracker.on_frame(message_at(100, 7, wbtc_usdc_pairs()), clock.at(9));
tracker.on_frame(message_at(100, 12, wbtc_usdc_pairs()), clock.at(9));
tracker.on_frame(message_at(100, 6, wbtc_usdc_pairs()), clock.at(9));
});
assert_eq!(histogram_values(&snapshot, FRAME_AGE, &[]), vec![2.0, -3.0]);
}
#[test]
fn auto_detection_stops_at_its_cap() {
let (registered, snapshot) = record_async(async {
let clock = Clock::new();
let mut tracker =
FreshnessTracker::new(TrackerSettings { auto_detect: true, ..settings(vec![]) });
for index in 0..(MAX_AUTO_DETECTED + 2) as u64 {
let address = Bytes::from_str(&format!("0x{index:040x}")).unwrap();
let frame = TitanPriceLevelMessage {
block_number: 100,
timestamp: BASE_WALL_NANOS + index * 1_000,
pamms: vec![TitanPammLevels { pamm: address, pairs: wbtc_usdc_pairs() }],
};
let served = tracker
.on_frame(frame, clock.at(1))
.is_some();
assert_eq!(served, (index as usize) < MAX_AUTO_DETECTED, "venue {index}");
}
tracker.registry.len()
});
assert_eq!(registered, MAX_AUTO_DETECTED);
assert_eq!(counter_value(&snapshot, UNREGISTERED_PAMM_ENTRIES, &[]), 2);
}
#[test]
fn quote_guard_off_emits_states_that_never_expire() {
let clock = Clock::new();
let mut tracker = FreshnessTracker::new(TrackerSettings {
quote_guard: false,
..settings(vec![fermiswap()])
});
let update = tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(15))
.expect("update expected");
assert!(state_of(&update)
.quotable_until()
.is_none());
assert_eq!(
tracker
.stale_deadline()
.expect("serving"),
clock.at(24).monotonic
);
}
#[test]
fn quotable_until_is_one_block_after_the_frame_was_built() {
let clock = Clock::new();
let mut tracker = tracker();
let update = tracker
.on_frame(message_at(100, 7, wbtc_usdc_pairs()), clock.at(9))
.expect("update expected");
assert_eq!(
quotable_until(&update),
clock.at(9).monotonic + (QUOTE_TTL - Duration::from_secs(2))
);
}
#[test]
fn quotable_until_never_exceeds_one_block_for_a_future_stamped_frame() {
let clock = Clock::new();
let mut tracker = tracker();
let update = tracker
.on_frame(message_at(100, 12, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
assert_eq!(quotable_until(&update), clock.at(0).monotonic + QUOTE_TTL);
}
#[test]
fn frame_older_than_one_block_yields_an_unquotable_state() {
let clock = Clock::new();
let mut tracker = tracker();
let update = tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(15))
.expect("update expected");
assert_eq!(quotable_until(&update), clock.at(15).monotonic);
assert_eq!(
tracker
.stale_deadline()
.expect("serving"),
clock.at(24).monotonic
);
}
#[test]
fn omitted_pair_is_not_removed_by_the_frame() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
let weth_usdc =
vec![pair_levels(WETH, USDC, vec![level(1_000_000_000_000_000_000, 3_000_000_000)])];
let update = tracker
.on_frame(message_at(101, 1, weth_usdc), clock.at(1))
.expect("update expected");
assert!(update.removed_pairs.is_empty());
assert_eq!(update.new_pairs.len(), 1);
assert_eq!(update.states.len(), 1);
assert!(!update
.states
.contains_key(&expected_id()));
}
#[test]
fn omitted_venue_is_not_removed_by_the_frame() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
let empty = TitanPriceLevelMessage {
block_number: 101,
timestamp: BASE_WALL_NANOS + NANOS_PER_SECOND,
pamms: vec![],
};
assert!(tracker
.on_frame(empty, clock.at(1))
.is_none());
assert!(tracker.stale_deadline().is_some());
}
#[test]
fn stale_deadline_fired_a_second_early_removes_nothing() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
assert!(tracker
.on_stale_deadline(clock.at(23).monotonic)
.is_none());
assert!(tracker.stale_deadline().is_some());
}
#[test]
fn component_turns_stale_at_its_deadline() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
let update = tracker
.on_stale_deadline(clock.at(24).monotonic)
.expect("removal expected");
assert!(update.states.is_empty());
assert!(update.new_pairs.is_empty());
assert_eq!(update.removed_pairs.len(), 1);
assert!(update
.removed_pairs
.contains_key(&expected_id()));
assert_eq!(update.block_number_or_timestamp, 100);
assert!(update.is_partial);
assert!(tracker.stale_deadline().is_none());
}
#[test]
fn stale_deadline_fired_with_nothing_served_removes_nothing() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
tracker
.on_stale_deadline(clock.at(24).monotonic)
.expect("removal expected");
assert!(tracker
.on_stale_deadline(clock.at(25).monotonic)
.is_none());
}
#[test]
fn fresh_frame_after_the_stale_removal_re_adds_the_component() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
tracker
.on_stale_deadline(clock.at(24).monotonic)
.expect("removal expected");
let update = tracker
.on_frame(message_at(102, 26, wbtc_usdc_pairs()), clock.at(26))
.expect("update expected");
assert!(update
.new_pairs
.contains_key(&expected_id()));
assert!(update.removed_pairs.is_empty());
}
#[test]
fn stale_removal_is_counted_per_venue() {
let ((), snapshot) = record_async(async {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
tracker
.on_stale_deadline(clock.at(24).monotonic)
.expect("removal expected");
});
assert_eq!(counter_value(&snapshot, STALE_REMOVALS, &[("venue", "fermiswap")]), 1);
}
#[test]
fn deadline_is_shortened_by_the_frame_age_at_acceptance() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(3))
.expect("update expected");
assert_eq!(
tracker
.stale_deadline()
.expect("serving"),
clock.at(24).monotonic
);
}
#[test]
fn replayed_frame_does_not_extend_the_deadline() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
tracker
.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(10))
.expect("update expected");
assert_eq!(
tracker
.stale_deadline()
.expect("serving"),
clock.at(24).monotonic
);
}
#[test]
fn omitted_pair_turns_stale_alone_while_the_rest_stays_served() {
let clock = Clock::new();
let mut tracker = tracker();
let both = vec![
pair_levels(WBTC, USDC, vec![level(100_000_000, 100_000_000_000)]),
pair_levels(USDC, WBTC, vec![level(100_000_000_000, 99_000_000)]),
pair_levels(WETH, USDC, vec![level(1_000_000_000_000_000_000, 3_000_000_000)]),
];
tracker
.on_frame(message_at(100, 0, both), clock.at(0))
.expect("update expected");
for second in 1..=23 {
let weth_usdc = vec![pair_levels(
WETH,
USDC,
vec![level(1_000_000_000_000_000_000, 3_000_000_000)],
)];
tracker.on_frame(message_at(100 + second / 12, second, weth_usdc), clock.at(second));
}
let update = tracker
.on_stale_deadline(clock.at(24).monotonic)
.expect("removal expected");
assert_eq!(update.removed_pairs.len(), 1);
assert!(update
.removed_pairs
.contains_key(&expected_id()));
assert_eq!(
tracker
.stale_deadline()
.expect("serving"),
clock.at(47).monotonic
);
}
#[test]
fn components_from_one_frame_turn_stale_together() {
let clock = Clock::new();
let mut tracker = tracker();
let both = vec![
pair_levels(WBTC, USDC, vec![level(100_000_000, 100_000_000_000)]),
pair_levels(WETH, USDC, vec![level(1_000_000_000_000_000_000, 3_000_000_000)]),
];
tracker
.on_frame(message_at(100, 0, both), clock.at(0))
.expect("update expected");
let update = tracker
.on_stale_deadline(clock.at(24).monotonic)
.expect("removal expected");
assert_eq!(update.removed_pairs.len(), 2);
}
#[test]
fn pair_set_oscillation_emits_no_removal() {
let clock = Clock::new();
let mut tracker = tracker();
let wide = || {
vec![
pair_levels(WBTC, USDC, vec![level(100_000_000, 100_000_000_000)]),
pair_levels(USDC, WBTC, vec![level(100_000_000_000, 99_000_000)]),
pair_levels(WETH, USDC, vec![level(1_000_000_000_000_000_000, 3_000_000_000)]),
]
};
let narrow =
|| vec![pair_levels(WETH, USDC, vec![level(1_000_000_000_000_000_000, 3_000_000_000)])];
tracker
.on_frame(message_at(100, 0, wide()), clock.at(0))
.expect("update expected");
for second in 1..=10u64 {
let pairs = if second % 2 == 0 { wide() } else { narrow() };
let update = tracker
.on_frame(message_at(100, second, pairs), clock.at(second))
.expect("update expected");
assert!(update.removed_pairs.is_empty(), "second {second}");
assert!(update.new_pairs.is_empty(), "second {second}");
}
}
#[test]
fn implausible_first_block_is_forgotten_at_stale_removal() {
let clock = Clock::new();
let mut tracker = tracker();
tracker
.on_frame(message_at(u64::MAX / 2, 0, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
for second in 1..=23 {
assert!(tracker
.on_frame(message_at(100, second, wbtc_usdc_pairs()), clock.at(second))
.is_none());
}
let removal = tracker
.on_stale_deadline(clock.at(24).monotonic)
.expect("removal expected");
assert_eq!(removal.removed_pairs.len(), 1);
assert_eq!(removal.block_number_or_timestamp, u64::MAX / 2);
let update = tracker
.on_frame(message_at(100, 25, wbtc_usdc_pairs()), clock.at(25))
.expect("update expected");
assert_eq!(update.block_number_or_timestamp, 100);
assert!(update
.new_pairs
.contains_key(&expected_id()));
}
#[test]
fn implausible_block_on_a_frame_serving_nothing() {
let clock = Clock::new();
let mut tracker = tracker();
let unknown = vec![pair_levels(
"0x1111111111111111111111111111111111111111",
USDC,
vec![level(1, 1)],
)];
assert!(tracker
.on_frame(message_at(u64::MAX / 2, 0, unknown), clock.at(0))
.is_none());
let update = tracker
.on_frame(message_at(100, 1, wbtc_usdc_pairs()), clock.at(1))
.expect("update expected");
assert_eq!(update.block_number_or_timestamp, 100);
assert!(update
.new_pairs
.contains_key(&expected_id()));
}
#[test]
fn recovery_re_adds_only_the_pairs_the_frame_carries() {
let clock = Clock::new();
let mut tracker = tracker();
let both = vec![
pair_levels(WBTC, USDC, vec![level(100_000_000, 100_000_000_000)]),
pair_levels(WETH, USDC, vec![level(1_000_000_000_000_000_000, 3_000_000_000)]),
];
tracker
.on_frame(message_at(100, 0, both), clock.at(0))
.expect("update expected");
tracker
.on_stale_deadline(clock.at(24).monotonic)
.expect("removal expected");
let weth_usdc =
vec![pair_levels(WETH, USDC, vec![level(1_000_000_000_000_000_000, 3_000_000_000)])];
let update = tracker
.on_frame(message_at(102, 25, weth_usdc), clock.at(25))
.expect("update expected");
assert_eq!(update.new_pairs.len(), 1);
assert!(!update
.new_pairs
.contains_key(&expected_id()));
assert!(update.removed_pairs.is_empty());
}
#[test]
fn one_direction_frame_re_adds_the_pair_with_the_other_direction_unquotable() {
let clock = Clock::new();
let mut tracker = tracker();
let one_way = vec![pair_levels(WBTC, USDC, vec![level(100_000_000, 100_000_000_000)])];
let update = tracker
.on_frame(message_at(100, 0, one_way), clock.at(0))
.expect("update expected");
let state = update.states[&expected_id()]
.as_any()
.downcast_ref::<PriceLevelStreamState>()
.expect("price level state");
assert_eq!(state.quotes_0_to_1.len(), 1);
assert!(state.quotes_1_to_0.is_empty());
let usdc = token(USDC, "USDC", 6);
let wbtc = token(WBTC, "WBTC", 8);
assert!(state
.get_amount_out(BigUint::from(1_000_000u64), &usdc, &wbtc)
.is_err());
}
#[test]
fn new_tracker_gauges_every_registered_venue_at_zero() {
let ((), snapshot) = record_async(async {
let tracker = tracker();
drop(tracker);
});
assert_eq!(gauge_value(&snapshot, SERVED_COMPONENTS, &[("venue", "fermiswap")]), 0.0);
assert_eq!(gauge_value(&snapshot, LAST_SEEN, &[("venue", "fermiswap")]), 0.0);
assert_eq!(gauge_value(&snapshot, SERVING_STATE, &[]), 0.0);
}
#[test]
fn accepted_frame_gauges_its_venue_as_served() {
let ((), snapshot) = record_async(async {
let clock = Clock::new();
let mut tracker = tracker();
tracker.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0));
});
assert_eq!(gauge_value(&snapshot, SERVED_COMPONENTS, &[("venue", "fermiswap")]), 1.0);
assert_eq!(
gauge_value(&snapshot, LAST_SEEN, &[("venue", "fermiswap")]),
(BASE_WALL_NANOS / NANOS_PER_SECOND) as f64
);
assert_eq!(gauge_value(&snapshot, SERVING_STATE, &[]), 1.0);
}
#[test]
fn stale_removal_gauges_its_venue_back_to_zero() {
let ((), snapshot) = record_async(async {
let clock = Clock::new();
let mut tracker = tracker();
tracker.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0));
tracker.on_stale_deadline(clock.at(24).monotonic);
});
assert_eq!(gauge_value(&snapshot, SERVED_COMPONENTS, &[("venue", "fermiswap")]), 0.0);
assert_eq!(gauge_value(&snapshot, SERVING_STATE, &[]), 0.0);
}
#[test]
fn unregistered_pamm_produces_no_update_without_auto_detection() {
let clock = Clock::new();
let mut tracker = tracker_serving(vec![], false);
assert!(tracker
.on_frame(message(100, wbtc_usdc_pairs()), clock.at(0))
.is_none());
}
#[test]
fn denied_pamm_is_not_auto_detected() {
let clock = Clock::new();
let denied = HashSet::from([Bytes::from_str(PAMM).unwrap()]);
let mut tracker = FreshnessTracker::new(TrackerSettings {
denied,
auto_detect: true,
..settings(vec![])
});
assert!(tracker
.on_frame(message(100, wbtc_usdc_pairs()), clock.at(0))
.is_none());
}
#[test]
fn auto_detected_pamm_is_served_under_its_address() {
let clock = Clock::new();
let mut tracker =
FreshnessTracker::new(TrackerSettings { auto_detect: true, ..settings(vec![]) });
let update = tracker
.on_frame(message(100, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
let component = &update.new_pairs[&expected_id()];
assert_eq!(component.protocol_system, format!("pricelevelstream:{PAMM}"));
let state = update.states[&expected_id()]
.as_any()
.downcast_ref::<PriceLevelStreamState>()
.expect("price level state");
assert_eq!(state.gas_cost, BigUint::from(DEFAULT_AUTO_DETECTED_GAS_COST));
let update = tracker
.on_frame(message(101, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
assert!(update.new_pairs.is_empty());
}
#[test]
fn auto_detected_gas_cost_override_applies() {
let clock = Clock::new();
let mut tracker = FreshnessTracker::new(TrackerSettings {
auto_detect: true,
auto_detected_gas_cost: BigUint::from(42_000u64),
..settings(vec![])
});
let update = tracker
.on_frame(message(100, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
let state = update.states[&expected_id()]
.as_any()
.downcast_ref::<PriceLevelStreamState>()
.expect("price level state");
assert_eq!(state.gas_cost, BigUint::from(42_000u64));
}
#[test]
fn unregistered_venues_counter_and_log_cap() {
let (logged, snapshot) = record_async(async {
let clock = Clock::new();
let mut tracker = tracker_serving(vec![], false);
for round in 0..2u64 {
for index in 0..70u64 {
let address = Bytes::from_str(&format!("0x{index:040x}")).unwrap();
let frame = TitanPriceLevelMessage {
block_number: 100,
timestamp: BASE_WALL_NANOS + (round * 70 + index) * 1_000,
pamms: vec![TitanPammLevels { pamm: address, pairs: wbtc_usdc_pairs() }],
};
tracker.on_frame(frame, clock.at(1));
}
}
tracker.logged_unregistered.len()
});
assert_eq!(counter_value(&snapshot, UNREGISTERED_PAMM_ENTRIES, &[]), 140);
assert_eq!(logged, MAX_UNREGISTERED_LOGGED);
}
#[test]
fn unknown_tokens_are_skipped() {
let clock = Clock::new();
let mut tracker = tracker();
let unknown = vec![pair_levels(
"0x1111111111111111111111111111111111111111",
USDC,
vec![level(1, 1)],
)];
assert!(tracker
.on_frame(message(100, unknown), clock.at(0))
.is_none());
}
#[test]
fn no_update_both_removes_and_adds_an_id() {
let clock = Clock::new();
let mut tracker = tracker();
let mut updates = Vec::new();
updates.extend(tracker.on_frame(message_at(100, 0, wbtc_usdc_pairs()), clock.at(0)));
updates.extend(tracker.on_stale_deadline(clock.at(24).monotonic));
updates.extend(tracker.on_frame(message_at(102, 25, wbtc_usdc_pairs()), clock.at(25)));
assert_eq!(updates.len(), 3);
for update in &updates {
for id in update.removed_pairs.keys() {
assert!(!update.new_pairs.contains_key(id), "{id} removed and added");
assert!(!update.states.contains_key(id), "{id} removed with a state");
}
}
}
#[test]
fn venues_are_served_under_the_fallback_family() {
let clock = Clock::new();
let mut tracker = tracker_serving(vec![fermiswap()], true);
let update = tracker
.on_frame(message(100, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
let component = &update.new_pairs[&expected_id()];
assert_eq!(component.protocol_system, "fallback:fermiswap");
assert_eq!(
component.static_attributes[PAMM_ADDRESS_ATTRIBUTE],
Bytes::from_str(PAMM).unwrap()
);
}
#[test]
fn auto_detected_venue_is_served_under_the_fallback_family() {
let clock = Clock::new();
let mut tracker = FreshnessTracker::new(TrackerSettings {
auto_detect: true,
via_fallback_router: true,
..settings(vec![])
});
let update = tracker
.on_frame(message(100, wbtc_usdc_pairs()), clock.at(0))
.expect("update expected");
assert_eq!(update.new_pairs[&expected_id()].protocol_system, format!("fallback:{PAMM}"));
}
}