use std::{collections::HashMap, sync::Arc};
use alloy::{
primitives::{Address, U256, keccak256},
sol,
sol_types::{SolCall, private::primitives::aliases::I24},
};
use nautilus_core::{UnixNanos, hex};
use nautilus_model::{
defi::{
data::block::BlockPosition,
pool_analysis::{
position::PoolPosition,
snapshot::{PoolAnalytics, PoolSnapshot, PoolState},
},
tick_map::tick::PoolTick,
},
identifiers::InstrumentId,
};
use thiserror::Error;
use super::base::{BaseContract, ContractCall, Multicall3};
use crate::rpc::{error::BlockchainRpcClientError, http::BlockchainHttpRpcClient};
sol! {
#[sol(rpc)]
contract UniswapV3Pool {
struct Slot0Data {
uint160 sqrtPriceX96;
int24 tick;
uint16 observationIndex;
uint16 observationCardinality;
uint16 observationCardinalityNext;
uint32 feeProtocol;
bool unlocked;
}
struct TickInfo {
uint128 liquidityGross;
int128 liquidityNet;
uint256 feeGrowthOutside0X128;
uint256 feeGrowthOutside1X128;
int56 tickCumulativeOutside;
uint160 secondsPerLiquidityOutsideX128;
uint32 secondsOutside;
bool initialized;
}
struct PositionInfo {
uint128 liquidity;
uint256 feeGrowthInside0LastX128;
uint256 feeGrowthInside1LastX128;
uint128 tokensOwed0;
uint128 tokensOwed1;
}
function slot0() external view returns (Slot0Data memory);
function liquidity() external view returns (uint128);
function feeGrowthGlobal0X128() external view returns (uint256);
function feeGrowthGlobal1X128() external view returns (uint256);
function protocolFees() external view returns (uint128 token0, uint128 token1);
function ticks(int24 tick) external view returns (TickInfo memory);
function positions(bytes32 key) external view returns (PositionInfo memory);
}
}
const PANCAKESWAP_V3_PROTOCOL_FEE_LANE_SIZE: u32 = 65_536;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FeeProtocolEncoding {
UniswapV3Packed,
PancakeSwapV3BasisPoints,
}
#[derive(Debug, Error)]
pub enum UniswapV3PoolError {
#[error("RPC error: {0}")]
RpcError(#[from] BlockchainRpcClientError),
#[error("Failed to decode {field} for pool {pool}: {reason} (raw data: {raw_data})")]
DecodingError {
field: String,
pool: Address,
reason: String,
raw_data: String,
},
#[error("Call failed for {field} at pool {pool}: {reason}")]
CallFailed {
field: String,
pool: Address,
reason: String,
},
#[error("Tick {tick} is not initialized in pool {pool}")]
TickNotInitialized { tick: i32, pool: Address },
}
#[derive(Debug)]
pub struct UniswapV3PoolContract {
base: BaseContract,
}
impl UniswapV3PoolContract {
#[must_use]
pub fn new(client: Arc<BlockchainHttpRpcClient>, multicall_calls_per_rpc_request: u32) -> Self {
Self {
base: BaseContract::new_with_multicall_limit(client, multicall_calls_per_rpc_request),
}
}
pub async fn get_global_state(
&self,
pool_address: &Address,
block: Option<u64>,
fee_protocol_encoding: FeeProtocolEncoding,
) -> Result<PoolState, UniswapV3PoolError> {
let calls = vec![
ContractCall {
target: *pool_address,
allow_failure: false,
call_data: UniswapV3Pool::slot0Call {}.abi_encode(),
},
ContractCall {
target: *pool_address,
allow_failure: false,
call_data: UniswapV3Pool::liquidityCall {}.abi_encode(),
},
ContractCall {
target: *pool_address,
allow_failure: false,
call_data: UniswapV3Pool::feeGrowthGlobal0X128Call {}.abi_encode(),
},
ContractCall {
target: *pool_address,
allow_failure: false,
call_data: UniswapV3Pool::feeGrowthGlobal1X128Call {}.abi_encode(),
},
ContractCall {
target: *pool_address,
allow_failure: false,
call_data: UniswapV3Pool::protocolFeesCall {}.abi_encode(),
},
];
let results = self.base.execute_multicall(calls, block).await?;
if results.len() != 5 {
return Err(UniswapV3PoolError::CallFailed {
field: "global_state_multicall".to_string(),
pool: *pool_address,
reason: format!("Expected 5 results, received {}", results.len()),
});
}
let slot0 =
UniswapV3Pool::slot0Call::abi_decode_returns(&results[0].returnData).map_err(|e| {
UniswapV3PoolError::DecodingError {
field: "slot0".to_string(),
pool: *pool_address,
reason: e.to_string(),
raw_data: hex::encode(&results[0].returnData),
}
})?;
let liquidity = UniswapV3Pool::liquidityCall::abi_decode_returns(&results[1].returnData)
.map_err(|e| UniswapV3PoolError::DecodingError {
field: "liquidity".to_string(),
pool: *pool_address,
reason: e.to_string(),
raw_data: hex::encode(&results[1].returnData),
})?;
let fee_growth_0 =
UniswapV3Pool::feeGrowthGlobal0X128Call::abi_decode_returns(&results[2].returnData)
.map_err(|e| UniswapV3PoolError::DecodingError {
field: "feeGrowthGlobal0X128".to_string(),
pool: *pool_address,
reason: e.to_string(),
raw_data: hex::encode(&results[2].returnData),
})?;
let fee_growth_1 =
UniswapV3Pool::feeGrowthGlobal1X128Call::abi_decode_returns(&results[3].returnData)
.map_err(|e| UniswapV3PoolError::DecodingError {
field: "feeGrowthGlobal1X128".to_string(),
pool: *pool_address,
reason: e.to_string(),
raw_data: hex::encode(&results[3].returnData),
})?;
let protocol_fees = UniswapV3Pool::protocolFeesCall::abi_decode_returns(
&results[4].returnData,
)
.map_err(|e| UniswapV3PoolError::DecodingError {
field: "protocolFees".to_string(),
pool: *pool_address,
reason: e.to_string(),
raw_data: hex::encode(&results[4].returnData),
})?;
let mut state = PoolState {
current_tick: slot0.tick.as_i32(),
price_sqrt_ratio_x96: slot0.sqrtPriceX96,
liquidity,
protocol_fees_token0: U256::from(protocol_fees.token0),
protocol_fees_token1: U256::from(protocol_fees.token1),
fee_protocol: 0,
fee_protocol0_basis_points: None,
fee_protocol1_basis_points: None,
fee_growth_global_0: fee_growth_0,
fee_growth_global_1: fee_growth_1,
};
match fee_protocol_encoding {
FeeProtocolEncoding::UniswapV3Packed => {
let fee_protocol = u8::try_from(slot0.feeProtocol).map_err(|e| {
UniswapV3PoolError::DecodingError {
field: "slot0.feeProtocol".to_string(),
pool: *pool_address,
reason: e.to_string(),
raw_data: slot0.feeProtocol.to_string(),
}
})?;
state.set_uniswap_v3_fee_protocol(fee_protocol);
}
FeeProtocolEncoding::PancakeSwapV3BasisPoints => {
let (fee_protocol0, fee_protocol1) =
split_pancakeswap_v3_fee_protocol(slot0.feeProtocol);
state.set_protocol_fee_basis_points(fee_protocol0, fee_protocol1);
}
}
Ok(state)
}
pub async fn get_tick(
&self,
pool_address: &Address,
tick: i32,
block: Option<u64>,
) -> Result<PoolTick, UniswapV3PoolError> {
let tick_i24 = I24::try_from(tick).map_err(|_| UniswapV3PoolError::CallFailed {
field: "tick".to_string(),
pool: *pool_address,
reason: format!("Tick {tick} out of range for int24"),
})?;
let call_data = UniswapV3Pool::ticksCall { tick: tick_i24 }.abi_encode();
let raw_response = self
.base
.execute_call(pool_address, &call_data, block)
.await?;
let tick_info =
UniswapV3Pool::ticksCall::abi_decode_returns(&raw_response).map_err(|e| {
UniswapV3PoolError::DecodingError {
field: format!("ticks({tick})"),
pool: *pool_address,
reason: e.to_string(),
raw_data: hex::encode(&raw_response),
}
})?;
Ok(PoolTick::new(
tick,
tick_info.liquidityGross,
tick_info.liquidityNet,
tick_info.feeGrowthOutside0X128,
tick_info.feeGrowthOutside1X128,
tick_info.initialized,
0, ))
}
pub async fn batch_get_ticks(
&self,
pool_address: &Address,
ticks: &[i32],
block: Option<u64>,
) -> Result<HashMap<i32, PoolTick>, UniswapV3PoolError> {
let calls: Vec<ContractCall> = ticks
.iter()
.map(|&tick| {
let tick_i24 = I24::try_from(tick).map_err(|_| UniswapV3PoolError::CallFailed {
field: format!("ticks({tick})"),
pool: *pool_address,
reason: "tick is out of range for int24".to_string(),
})?;
Ok(ContractCall {
target: *pool_address,
allow_failure: false,
call_data: UniswapV3Pool::ticksCall { tick: tick_i24 }.abi_encode(),
})
})
.collect::<Result<_, UniswapV3PoolError>>()?;
let results = self.base.execute_multicall(calls, block).await?;
decode_tick_results(pool_address, ticks, &results)
}
#[must_use]
pub fn compute_position_key(owner: &Address, tick_lower: i32, tick_upper: i32) -> [u8; 32] {
let mut packed = Vec::with_capacity(26);
packed.extend_from_slice(owner.as_slice());
let tick_lower_bytes = tick_lower.to_be_bytes();
packed.extend_from_slice(&tick_lower_bytes[1..4]);
let tick_upper_bytes = tick_upper.to_be_bytes();
packed.extend_from_slice(&tick_upper_bytes[1..4]);
keccak256(&packed).into()
}
pub async fn batch_get_positions(
&self,
pool_address: &Address,
positions: &[(Address, i32, i32)],
block: Option<u64>,
) -> Result<Vec<PoolPosition>, UniswapV3PoolError> {
let calls: Vec<ContractCall> = positions
.iter()
.map(|(owner, tick_lower, tick_upper)| {
let position_key = Self::compute_position_key(owner, *tick_lower, *tick_upper);
ContractCall {
target: *pool_address,
allow_failure: false,
call_data: UniswapV3Pool::positionsCall {
key: position_key.into(),
}
.abi_encode(),
}
})
.collect();
let results = self.base.execute_multicall(calls, block).await?;
decode_position_results(pool_address, positions, &results)
}
#[expect(clippy::too_many_arguments)]
pub async fn fetch_snapshot(
&self,
pool_address: &Address,
instrument_id: InstrumentId,
tick_values: &[i32],
position_keys: &[(Address, i32, i32)],
block_position: BlockPosition,
ts_event: UnixNanos,
ts_init: UnixNanos,
fee_protocol_encoding: FeeProtocolEncoding,
) -> Result<PoolSnapshot, UniswapV3PoolError> {
let block = Some(block_position.number);
let global_state = self
.get_global_state(pool_address, block, fee_protocol_encoding)
.await?;
let ticks_map = self
.batch_get_ticks(pool_address, tick_values, block)
.await?;
let positions = self
.batch_get_positions(pool_address, position_keys, block)
.await?;
Ok(PoolSnapshot::new(
instrument_id,
global_state,
positions,
ticks_map.into_values().collect(),
PoolAnalytics::default(),
block_position,
ts_event,
ts_init,
))
}
}
fn decode_tick_results(
pool_address: &Address,
ticks: &[i32],
results: &[Multicall3::Result],
) -> Result<HashMap<i32, PoolTick>, UniswapV3PoolError> {
if results.len() != ticks.len() {
return Err(UniswapV3PoolError::CallFailed {
field: "ticks".to_string(),
pool: *pool_address,
reason: format!(
"expected {} multicall results, received {}",
ticks.len(),
results.len()
),
});
}
let mut tick_infos = HashMap::with_capacity(ticks.len());
for (&tick_value, result) in ticks.iter().zip(results) {
if !result.success {
return Err(UniswapV3PoolError::CallFailed {
field: format!("ticks({tick_value})"),
pool: *pool_address,
reason: "multicall subcall failed".to_string(),
});
}
let tick_info =
UniswapV3Pool::ticksCall::abi_decode_returns(&result.returnData).map_err(|e| {
UniswapV3PoolError::DecodingError {
field: format!("ticks({tick_value})"),
pool: *pool_address,
reason: e.to_string(),
raw_data: hex::encode(&result.returnData),
}
})?;
tick_infos.insert(
tick_value,
PoolTick::new(
tick_value,
tick_info.liquidityGross,
tick_info.liquidityNet,
tick_info.feeGrowthOutside0X128,
tick_info.feeGrowthOutside1X128,
tick_info.initialized,
0,
),
);
}
Ok(tick_infos)
}
fn decode_position_results(
pool_address: &Address,
positions: &[(Address, i32, i32)],
results: &[Multicall3::Result],
) -> Result<Vec<PoolPosition>, UniswapV3PoolError> {
if results.len() != positions.len() {
return Err(UniswapV3PoolError::CallFailed {
field: "positions".to_string(),
pool: *pool_address,
reason: format!(
"expected {} multicall results, received {}",
positions.len(),
results.len()
),
});
}
positions
.iter()
.zip(results)
.map(|((owner, tick_lower, tick_upper), result)| {
let field = format!("positions({owner}, {tick_lower}, {tick_upper})");
if !result.success {
return Err(UniswapV3PoolError::CallFailed {
field,
pool: *pool_address,
reason: "multicall subcall failed".to_string(),
});
}
let info = UniswapV3Pool::positionsCall::abi_decode_returns(&result.returnData)
.map_err(|e| UniswapV3PoolError::DecodingError {
field,
pool: *pool_address,
reason: e.to_string(),
raw_data: hex::encode(&result.returnData),
})?;
Ok(PoolPosition {
owner: *owner,
tick_lower: *tick_lower,
tick_upper: *tick_upper,
liquidity: info.liquidity,
fee_growth_inside_0_last: info.feeGrowthInside0LastX128,
fee_growth_inside_1_last: info.feeGrowthInside1LastX128,
tokens_owed_0: info.tokensOwed0,
tokens_owed_1: info.tokensOwed1,
total_amount0_deposited: U256::ZERO,
total_amount1_deposited: U256::ZERO,
total_amount0_collected: 0,
total_amount1_collected: 0,
})
})
.collect()
}
const fn split_pancakeswap_v3_fee_protocol(fee_protocol: u32) -> (u32, u32) {
(
fee_protocol % PANCAKESWAP_V3_PROTOCOL_FEE_LANE_SIZE,
fee_protocol >> 16,
)
}
#[cfg(test)]
mod tests {
use alloy::primitives::Bytes;
use rstest::rstest;
use super::*;
#[rstest]
fn split_pancakeswap_v3_fee_protocol_returns_token_basis_points() {
let fee_protocol = 3_200 + (4_000 << 16);
assert_eq!(
split_pancakeswap_v3_fee_protocol(fee_protocol),
(3_200, 4_000)
);
assert_eq!(split_pancakeswap_v3_fee_protocol(0), (0, 0));
}
#[rstest]
fn decode_tick_results_rejects_failed_subcall() {
let results = vec![Multicall3::Result {
success: false,
returnData: Bytes::new(),
}];
let error = decode_tick_results(&Address::ZERO, &[10], &results).unwrap_err();
assert!(matches!(
error,
UniswapV3PoolError::CallFailed { field, .. } if field == "ticks(10)"
));
}
#[rstest]
fn decode_tick_results_rejects_missing_result() {
let error = decode_tick_results(&Address::ZERO, &[10], &[]).unwrap_err();
assert!(matches!(
error,
UniswapV3PoolError::CallFailed { field, .. } if field == "ticks"
));
}
#[rstest]
fn decode_position_results_rejects_undecodable_result() {
let results = vec![Multicall3::Result {
success: true,
returnData: Bytes::new(),
}];
let error = decode_position_results(&Address::ZERO, &[(Address::ZERO, -10, 10)], &results)
.unwrap_err();
assert!(matches!(error, UniswapV3PoolError::DecodingError { .. }));
}
#[rstest]
fn decode_position_results_rejects_failed_subcall() {
let results = vec![Multicall3::Result {
success: false,
returnData: Bytes::new(),
}];
let error = decode_position_results(&Address::ZERO, &[(Address::ZERO, -10, 10)], &results)
.unwrap_err();
assert!(matches!(
error,
UniswapV3PoolError::CallFailed { field, .. }
if field == "positions(0x0000000000000000000000000000000000000000, -10, 10)"
));
}
#[rstest]
fn decode_position_results_rejects_missing_result() {
let error =
decode_position_results(&Address::ZERO, &[(Address::ZERO, -10, 10)], &[]).unwrap_err();
assert!(matches!(
error,
UniswapV3PoolError::CallFailed { field, .. } if field == "positions"
));
}
}