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};
use tokio_tungstenite::{connect_async, tungstenite::Message};
use tracing::{info, warn};
use tycho_common::Bytes;
pub(super) const TITAN_PRICE_LEVEL_URL: &str = "wss://eu.rpc.titanbuilder.xyz/ws/pamm_price_levels";
#[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(30),
max_backoff: Duration::from_secs(32),
}
}
}
fn backoff(attempt: u32, max_backoff: Duration) -> Duration {
let exponential = 2u64
.checked_pow(attempt)
.map(Duration::from_secs)
.unwrap_or(Duration::MAX);
exponential.min(max_backoff)
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(super) struct TitanPriceLevelMessage {
pub block_number: 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");
loop {
let message = match timeout(settings.read_idle_timeout, ws_stream.next())
.await
{
Ok(Some(message)) => message,
Ok(None) => {
warn!("Titan price level stream ended; reconnecting");
break;
}
Err(_elapsed) => {
warn!(
idle_secs = settings.read_idle_timeout.as_secs(),
"No Titan message within idle timeout; reconnecting"
);
break;
}
};
match message {
Ok(Message::Text(text)) => {
attempt = 0;
match serde_json::from_str::<TitanPriceLevelMessage>(text.as_str())
{
Ok(message) => yield message,
Err(e) => {
warn!(error = %e, "Failed to parse Titan price level message")
}
}
}
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");
break;
}
Err(e) => {
warn!(error = %e, "Titan price level stream read error; reconnecting");
break;
}
}
}
}
Ok(Err(e)) => {
warn!(error = %e, "Failed to connect to Titan price level stream; retrying");
}
Err(_elapsed) => {
warn!(
timeout_secs = settings.connect_timeout.as_secs(),
"Titan price level connect timed out; retrying"
);
}
}
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;
use super::*;
const SAMPLE_MESSAGE: &str = r#"{
"slot": 14581462,
"blockNumber": 25345763,
"timestamp": 1781801564588230787,
"pamms": [
{
"pamm": "0x5979458912f80b96d30d4220af8e2e4925a33320",
"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.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_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.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);
}
}