use std::time::Duration;
use async_stream::stream;
use futures::{Stream, StreamExt};
use num_bigint::BigUint;
use serde::{Deserialize, Deserializer};
use tokio::time::{sleep, timeout, timeout_at, Instant};
use tokio_tungstenite::{connect_async, tungstenite::Message};
use tracing::{info, warn};
use tycho_common::Bytes;
use super::telemetry::{self, ReconnectReason, RejectReason};
pub(super) const TITAN_PRICE_LEVEL_URL: &str =
"wss://eu.data.titanbuilder.xyz/ws/pamm_price_levels";
pub(super) const TITAN_PRICE_LEVEL_URL_ENV: &str = "TITAN_PAMM_PRICE_LEVEL_URL";
#[derive(Clone, Copy, Debug)]
pub(super) struct ConnectionSettings {
pub connect_timeout: Duration,
pub read_idle_timeout: Duration,
pub max_backoff: Duration,
}
impl Default for ConnectionSettings {
fn default() -> Self {
Self {
connect_timeout: Duration::from_secs(10),
read_idle_timeout: Duration::from_secs(10),
max_backoff: Duration::from_secs(32),
}
}
}
fn backoff(attempt: u32, max_backoff: Duration) -> Duration {
2u64.checked_pow(attempt)
.map_or(Duration::MAX, Duration::from_secs)
.min(max_backoff)
}
fn deadline_after(idle_timeout: Duration) -> Instant {
Instant::now() + idle_timeout.min(Duration::from_secs(u32::MAX.into()))
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(super) struct TitanPriceLevelMessage {
pub block_number: u64,
pub timestamp: u64,
pub pamms: Vec<TitanPammLevels>,
}
#[derive(Debug, Deserialize)]
pub(super) struct TitanPammLevels {
pub pamm: Bytes,
pub pairs: Vec<TitanPairLevels>,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(super) struct TitanPairLevels {
pub token_in: Bytes,
pub token_out: Bytes,
pub order_book: Vec<TitanPriceLevel>,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(super) struct TitanPriceLevel {
#[serde(deserialize_with = "quantity")]
pub amount_in: BigUint,
#[serde(deserialize_with = "quantity")]
pub amount_out: BigUint,
}
fn quantity<'de, D>(deserializer: D) -> Result<BigUint, D::Error>
where
D: Deserializer<'de>,
{
let text = String::deserialize(deserializer)?;
let digits = text.strip_prefix("0x").ok_or_else(|| {
serde::de::Error::custom(format!("quantity must be a 0x-prefixed hex string: {text}"))
})?;
BigUint::parse_bytes(digits.as_bytes(), 16)
.ok_or_else(|| serde::de::Error::custom(format!("invalid quantity: {text}")))
}
pub(super) fn messages(
url: String,
settings: ConnectionSettings,
) -> impl Stream<Item = TitanPriceLevelMessage> + Send {
stream! {
let mut attempt: u32 = 0;
loop {
match timeout(settings.connect_timeout, connect_async(url.as_str())).await {
Ok(Ok((mut ws_stream, _))) => {
info!(%url, "Connected to Titan pAMM price level stream");
let mut idle_deadline = deadline_after(settings.read_idle_timeout);
loop {
let message = match timeout_at(idle_deadline, ws_stream.next()).await {
Ok(Some(message)) => message,
Ok(None) => {
warn!("Titan price level stream ended; reconnecting");
telemetry::record_reconnect(ReconnectReason::Ended);
break;
}
Err(_elapsed) => {
warn!(
idle_secs = settings.read_idle_timeout.as_secs(),
"No parsed Titan frame within idle timeout; reconnecting"
);
telemetry::record_reconnect(ReconnectReason::IdleTimeout);
break;
}
};
match message {
Ok(Message::Text(text)) => {
match serde_json::from_str::<TitanPriceLevelMessage>(text.as_str())
{
Ok(message) => {
attempt = 0;
yield message;
idle_deadline = deadline_after(settings.read_idle_timeout);
}
Err(e) => {
warn!(error = %e, "Failed to parse Titan price level message");
telemetry::record_frame_rejected(RejectReason::ParseError);
}
}
}
Ok(Message::Binary(bytes)) => {
warn!(len = bytes.len(), "Ignoring unexpected binary Titan frame");
}
Ok(Message::Ping(_)) | Ok(Message::Pong(_)) => {}
Ok(Message::Frame(_)) => {}
Ok(Message::Close(frame)) => {
warn!(?frame, "Titan price level stream closed by server; reconnecting");
telemetry::record_reconnect(ReconnectReason::Closed);
break;
}
Err(e) => {
warn!(error = %e, "Titan price level stream read error; reconnecting");
telemetry::record_reconnect(ReconnectReason::ReadError);
break;
}
}
}
}
Ok(Err(e)) => {
warn!(error = %e, "Failed to connect to Titan price level stream; retrying");
telemetry::record_reconnect(ReconnectReason::ConnectFailed);
}
Err(_elapsed) => {
warn!(
timeout_secs = settings.connect_timeout.as_secs(),
"Titan price level connect timed out; retrying"
);
telemetry::record_reconnect(ReconnectReason::ConnectTimeout);
}
}
attempt = attempt.saturating_add(1);
let backoff = backoff(attempt, settings.max_backoff);
warn!(seconds = backoff.as_secs(), attempt, "Backing off before reconnecting to Titan");
sleep(backoff).await;
}
}
}
#[cfg(test)]
mod tests {
use std::{str::FromStr, sync::atomic::Ordering};
use futures::SinkExt;
use rstest::rstest;
use super::{
super::{
telemetry::{
recorded::{counter_value, record_async},
FRAMES_REJECTED, RECONNECTS,
},
test_support::{frame_text, frame_then_repeat, wall_nanos_now, FakeTitan},
},
*,
};
const SAMPLE_MESSAGE: &str = r#"{
"slot": 14581462,
"blockNumber": 25345763,
"timestamp": 1781801564588230787,
"pamms": [
{
"pamm": "0x5979458912f80b96d30d4220af8e2e4925a33320",
"maker": 12,
"pairs": [
{
"tokenIn": "0x2260fac5e5542a773aa44fbcfedf7c193bc2c599",
"tokenOut": "0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48",
"orderBook": [
{
"amountIn": "0x989680",
"amountOut": "0x174b67393",
"variant": "Simulated"
},
{
"amountIn": "0x1312d00",
"amountOut": "0x2e968e726",
"variant": "Interpolated"
}
]
}
]
}
]
}"#;
#[test]
fn parses_documented_sample_message() {
let message: TitanPriceLevelMessage = serde_json::from_str(SAMPLE_MESSAGE).unwrap();
assert_eq!(message.block_number, 25345763);
assert_eq!(message.timestamp, 1781801564588230787);
assert_eq!(message.pamms.len(), 1);
let pamm = &message.pamms[0];
assert_eq!(
pamm.pamm,
Bytes::from_str("0x5979458912f80b96d30d4220af8e2e4925a33320").unwrap()
);
assert_eq!(pamm.pairs.len(), 1);
let pair = &pamm.pairs[0];
assert_eq!(
pair.token_in,
Bytes::from_str("0x2260fac5e5542a773aa44fbcfedf7c193bc2c599").unwrap()
);
assert_eq!(pair.order_book.len(), 2);
assert_eq!(pair.order_book[0].amount_in, BigUint::from(0x989680u64));
assert_eq!(pair.order_book[0].amount_out, BigUint::from(0x174b67393u64));
}
#[test]
fn rejects_frame_without_timestamp() {
let json = r#"{"slot": 1, "blockNumber": 2, "pamms": []}"#;
assert!(serde_json::from_str::<TitanPriceLevelMessage>(json).is_err());
}
#[test]
fn rejects_quantities_that_are_not_hex_strings() {
for json in [
r#"{"amountIn": "0xzz", "amountOut": "0x1"}"#,
r#"{"amountIn": "0x", "amountOut": "0x1"}"#,
r#"{"amountIn": "1000", "amountOut": "0x1"}"#,
r#"{"amountIn": 1000, "amountOut": "0x1"}"#,
] {
assert!(serde_json::from_str::<TitanPriceLevel>(json).is_err(), "accepted: {json}");
}
}
const CAPTURED_MESSAGE: &str =
include_str!("test_responses/pamm_price_levels_1784126589047308938.json");
#[test]
fn parses_captured_live_message() {
let message: TitanPriceLevelMessage =
serde_json::from_str(CAPTURED_MESSAGE).expect("valid JSON");
assert_eq!(message.block_number, 25538727);
assert_eq!(message.timestamp, 1784126589047308938);
assert_eq!(message.pamms.len(), 2);
let fermiswap = &message.pamms[0];
assert_eq!(
fermiswap.pamm,
Bytes::from_str("0x5979458912f80b96d30d4220af8e2e4925a33320").unwrap()
);
assert_eq!(fermiswap.pairs.len(), 16);
let kipseli = &message.pamms[1];
assert_eq!(
kipseli.pamm,
Bytes::from_str("0x71e790dd841c8a9061487cb3e78c288e75ce0b3d").unwrap()
);
assert_eq!(kipseli.pairs.len(), 4);
let first = &fermiswap.pairs[0];
assert_eq!(
first.token_in,
Bytes::from_str("0x2260fac5e5542a773aa44fbcfedf7c193bc2c599").unwrap()
);
assert_eq!(first.order_book[0].amount_in, BigUint::from(0xc350u64));
assert_eq!(first.order_book[0].amount_out, BigUint::from(0x1f27427u64));
for pamm in &message.pamms {
for pair in &pamm.pairs {
assert!(pair.order_book.len() >= 64, "unexpectedly short ladder");
for level in &pair.order_book {
assert!(level.amount_in > BigUint::ZERO);
assert!(level.amount_out > BigUint::ZERO);
}
}
}
}
#[test]
fn backoff_grows_exponentially_up_to_the_cap() {
let max_backoff = ConnectionSettings::default().max_backoff;
assert_eq!(backoff(1, max_backoff), Duration::from_secs(2));
assert_eq!(backoff(4, max_backoff), Duration::from_secs(16));
assert_eq!(backoff(5, max_backoff), max_backoff);
assert_eq!(backoff(100, max_backoff), max_backoff);
assert_eq!(backoff(u32::MAX, max_backoff), max_backoff);
}
fn fast_settings() -> ConnectionSettings {
ConnectionSettings {
connect_timeout: Duration::from_secs(1),
read_idle_timeout: Duration::from_millis(100),
max_backoff: Duration::from_millis(10),
}
}
fn frame() -> Message {
Message::Text(frame_text(100, wall_nanos_now()).into())
}
async fn next_frame(stream: &mut (impl Stream<Item = TitanPriceLevelMessage> + Unpin)) {
tokio::time::timeout(Duration::from_secs(2), stream.next())
.await
.expect("a frame within two seconds")
.expect("the stream never ends");
}
#[rstest]
#[case::ping_only(Message::Ping(Vec::new().into()), false)]
#[case::malformed_text(Message::Text("nonsense".into()), true)]
fn non_frame_traffic_does_not_count_as_liveness(
#[case] filler: Message,
#[case] rejected_as_parse_error: bool,
) {
let (connections, snapshot) = record_async(async {
let fake =
FakeTitan::spawn(frame_then_repeat(frame(), filler, Duration::from_millis(10)))
.await;
let stream = messages(fake.url(), fast_settings());
tokio::pin!(stream);
next_frame(&mut stream).await;
next_frame(&mut stream).await;
fake.connections.load(Ordering::SeqCst)
});
assert!(connections >= 2, "no reconnect on non-frame traffic");
assert!(counter_value(&snapshot, RECONNECTS, &[("reason", "idle_timeout")]) >= 1);
let parse_errors = counter_value(&snapshot, FRAMES_REJECTED, &[("reason", "parse_error")]);
assert_eq!(parse_errors > 0, rejected_as_parse_error);
}
#[test]
fn server_close_frame_reconnects() {
let (connections, snapshot) = record_async(async {
let fake = FakeTitan::spawn(|_, mut socket| async move {
if socket.send(frame()).await.is_err() {
return;
}
let _ = socket.send(Message::Close(None)).await;
})
.await;
let stream = messages(fake.url(), fast_settings());
tokio::pin!(stream);
next_frame(&mut stream).await;
next_frame(&mut stream).await;
fake.connections.load(Ordering::SeqCst)
});
assert_eq!(connections, 2);
assert_eq!(counter_value(&snapshot, RECONNECTS, &[("reason", "closed")]), 1);
}
#[test]
fn slow_consumer_does_not_trigger_the_idle_timeout() {
let (connections, snapshot) = record_async(async {
let fake =
FakeTitan::spawn(frame_then_repeat(frame(), frame(), Duration::from_millis(20)))
.await;
let stream = messages(fake.url(), fast_settings());
tokio::pin!(stream);
next_frame(&mut stream).await;
tokio::time::sleep(Duration::from_millis(300)).await;
next_frame(&mut stream).await;
fake.connections.load(Ordering::SeqCst)
});
assert_eq!(connections, 1, "reconnected while the consumer was not polling");
assert_eq!(counter_value(&snapshot, RECONNECTS, &[("reason", "idle_timeout")]), 0);
}
#[tokio::test]
async fn parsed_frames_keep_the_connection_alive() {
let fake =
FakeTitan::spawn(frame_then_repeat(frame(), frame(), Duration::from_millis(30))).await;
let settings =
ConnectionSettings { read_idle_timeout: Duration::from_millis(150), ..fast_settings() };
let stream = messages(fake.url(), settings);
tokio::pin!(stream);
for _ in 0..10 {
next_frame(&mut stream).await;
}
assert_eq!(fake.connections.load(Ordering::SeqCst), 1);
}
}