use crate::{
dex_connector::{slippage_price, string_to_decimal, DexConnector},
dex_request::{DexError, DexRequest, HttpMethod},
dex_websocket::DexWebSocket,
BalanceResponse, CreateOrderResponse, FilledOrder, FilledOrdersResponse, OrderSide,
TickerResponse,
};
use async_trait::async_trait;
use debot_utils::parse_to_decimal;
use ecdsa::SigningKey;
use ethers::{prelude::Signer, types::transaction::eip712::EIP712Domain};
use ethers::{
prelude::*,
types::transaction::eip712::{Eip712DomainType, TypedData, Types},
};
use futures::{
stream::{SplitSink, SplitStream},
SinkExt, StreamExt,
};
use generic_array::GenericArray;
use k256::Secp256k1;
use k256::SecretKey;
use rust_decimal::Decimal;
use serde::{Deserialize, Serialize};
use serde_json::json;
use sha3::{Digest, Keccak256};
use std::{
collections::{BTreeMap, HashMap},
sync::{
atomic::{AtomicBool, Ordering},
Arc,
},
time::{Duration, SystemTime, UNIX_EPOCH},
};
use tokio::signal::unix::SignalKind;
use tokio::sync::Mutex;
use tokio::sync::RwLock;
use tokio::time::sleep;
use tokio::{net::TcpStream, task::JoinHandle};
use tokio::{select, signal::unix::signal};
use tokio_tungstenite::tungstenite::protocol::Message;
use tokio_tungstenite::MaybeTlsStream;
use tokio_tungstenite::WebSocketStream;
struct Config {
agent_private_key: String,
evm_wallet_address: String,
market_ids: Vec<String>,
}
struct TradeResult {
pub filled_side: OrderSide,
pub filled_size: Decimal,
pub filled_value: Decimal,
pub filled_fee: Decimal,
}
#[derive(Default)]
struct MarketInfo {
pub last_trade_price: Option<Decimal>,
pub market_price: Option<Decimal>,
pub min_order: Option<Decimal>,
pub min_tick: Option<Decimal>,
pub decimals: Option<u32>,
pub max_leverage: Option<u32>,
pub asset_index: Option<u32>,
}
pub struct HyperliquidConnector {
config: Config,
request: DexRequest,
web_socket: DexWebSocket,
running: Arc<AtomicBool>,
read_socket: Arc<Mutex<Option<SplitStream<WebSocketStream<MaybeTlsStream<TcpStream>>>>>>,
task_handle_read_message: Arc<Mutex<Option<JoinHandle<()>>>>,
task_handle_read_sigterm: Arc<Mutex<Option<JoinHandle<()>>>>,
trade_results: Arc<RwLock<HashMap<String, HashMap<String, TradeResult>>>>,
market_info: Arc<RwLock<HashMap<String, MarketInfo>>>,
nonce: Arc<Mutex<u128>>,
wallet: Wallet<SigningKey<Secp256k1>>,
}
#[derive(Deserialize, Debug)]
struct WebSocketMessage {
push: Option<PushData>,
}
#[derive(Deserialize, Debug)]
struct PushData {
channel: String,
#[serde(rename = "pub")]
pub_data: PubData,
}
#[derive(Deserialize, Debug)]
struct PubData {
data: MarketData,
}
#[derive(Deserialize, Debug)]
struct MarketData {
id: Option<String>,
min_tick: Option<String>,
min_order: Option<String>,
last_trade_price: Option<String>,
market_price: Option<String>,
}
#[derive(Deserialize, Debug)]
struct Fill {
order_id: String,
side: String,
price: String,
size: String,
market_id: String,
fee: String,
}
#[derive(Deserialize, Debug)]
struct AccountData {
fills: Vec<Fill>,
}
#[derive(Deserialize, Debug)]
struct AccountPubData {
data: AccountData,
}
#[derive(Deserialize, Debug)]
struct AccountPushData {
channel: String,
#[serde(rename = "pub")]
pub_data: AccountPubData,
}
#[derive(Deserialize, Debug)]
struct AccountWebSocketMessage {
push: Option<AccountPushData>,
}
impl HyperliquidConnector {
pub async fn new(
rest_endpoint: &str,
web_socket_endpoint: &str,
agent_private_key: &str,
evm_wallet_address: &str,
market_ids: &[String],
) -> Result<Self, DexError> {
let request = DexRequest::new(rest_endpoint.to_owned()).await?;
let web_socket = DexWebSocket::new(web_socket_endpoint.to_owned());
let config = Config {
agent_private_key: agent_private_key.to_owned(),
evm_wallet_address: evm_wallet_address.to_owned(),
market_ids: market_ids.to_vec(),
};
let private_key_bytes = hex::decode(agent_private_key)
.map_err(|e| DexError::Other(format!("Failed to decode private key: {}", e)))?;
let private_key_bytes = GenericArray::from_slice(&private_key_bytes);
let secret_key = SecretKey::from_bytes(private_key_bytes)
.map_err(|e| DexError::Other(format!("Failed to create secret key: {}", e)))?;
let signing_key = SigningKey::from(&secret_key);
let wallet = Wallet::from(signing_key);
Ok(HyperliquidConnector {
config,
request,
web_socket,
trade_results: Arc::new(RwLock::new(HashMap::new())),
market_info: Arc::new(RwLock::new(HashMap::new())),
running: Arc::new(AtomicBool::new(false)),
read_socket: Arc::new(Mutex::new(None)),
task_handle_read_message: Arc::new(Mutex::new(None)),
task_handle_read_sigterm: Arc::new(Mutex::new(None)),
nonce: Arc::new(Mutex::new(
SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("Time went backwards")
.as_millis(),
)),
wallet,
})
}
pub async fn start_web_socket(&self) -> Result<(), DexError> {
log::info!("start_web_socket");
let web_socket = self.web_socket.clone();
let (mut write, read) = match web_socket.connect().await {
Ok((write, read)) => (write, read),
Err(_) => {
return Err(DexError::Other(
"Failed to connect to WebSocket".to_string(),
))
}
};
let mut read_lock = self.read_socket.lock().await;
*read_lock = Some(read);
self.running.store(true, Ordering::SeqCst);
self.subscribe_to_channels(&mut write, &self.config.market_ids)
.await
.unwrap();
log::debug!("subscription is done");
let running_clone = self.running.clone();
let read_clone = self.read_socket.clone();
let write_clone = Arc::new(Mutex::new(write));
let market_info_clone = self.market_info.clone();
let trade_results_clone = self.trade_results.clone();
let handle = tokio::spawn(async move {
log::debug!("WebSocket message handling task started");
let mut message_counter = 0;
while running_clone.load(Ordering::SeqCst) {
let mut read_guard = read_clone.lock().await;
if let Some(read_stream) = read_guard.as_mut() {
tokio::select! {
message = read_stream.next() => match message {
Some(Ok(msg)) => {
message_counter = 0;
if msg == "{}".into() {
write_clone
.lock()
.await
.send(Message::Text(msg.to_string()))
.await
.unwrap();
log::trace!("Responsed to the ping")
} else {
log::trace!("Received message: {:?}", msg);
if let Err(e) = Self::handle_websocket_message(
msg,
market_info_clone.clone(),
trade_results_clone.clone(),
)
.await
{
log::error!("Error handling WebSocket message: {:?}", e);
break;
}
}
}
Some(Err(e)) => {
log::error!("Failed to read: {:?}", e);
break;
}
None => {
log::info!("WebSocket stream ended");
break;
}
},
_ = tokio::time::sleep(tokio::time::Duration::from_secs(10)) => {
if !running_clone.load(Ordering::SeqCst) {
log::info!("Running flag changed, shutting down...");
break;
}
message_counter += 1;
if message_counter >= 10 {
log::error!("No message has been received for some time");
break;
}
},
}
}
}
running_clone.store(false, Ordering::SeqCst);
log::info!("WebSocket message handling task ended");
});
let mut task_handle = self.task_handle_read_message.lock().await;
*task_handle = Some(handle);
let mut sigterm =
signal(SignalKind::terminate()).expect("Failed to create SIGTERM listener");
let running_clone = self.running.clone();
let handle = tokio::spawn(async move {
log::debug!("SIGTERM handling task started");
loop {
select! {
_ = sigterm.recv() => {
log::info!("SIGTERM received, shutting down...");
running_clone.store(false, Ordering::SeqCst);
break;
},
_ = tokio::time::sleep(tokio::time::Duration::from_secs(1)) => {
if !running_clone.load(Ordering::SeqCst) {
log::info!("Running flag changed, shutting down...");
break;
}
},
}
}
});
let mut task_handle = self.task_handle_read_sigterm.lock().await;
*task_handle = Some(handle);
Ok(())
}
pub async fn stop_web_socket(&self) -> Result<(), DexError> {
log::info!("stop_web_socket");
self.running.store(false, Ordering::SeqCst);
Ok(())
}
async fn subscribe_to_channels(
&self,
socket: &mut SplitSink<WebSocketStream<MaybeTlsStream<TcpStream>>, Message>,
market_ids: &[String],
) -> Result<(), DexError> {
let mut channels = Vec::new();
for market_id in market_ids {
channels.push(format!("market:{}", market_id));
}
for (idx, channel) in channels.iter().enumerate() {
let data = serde_json::json!({
"subscribe": {
"channel": channel,
"name": "js",
},
"id": idx + 1
});
socket.send(Message::Text(data.to_string())).await.unwrap();
}
Ok(())
}
async fn handle_websocket_message(
msg: Message,
market_info: Arc<RwLock<HashMap<String, MarketInfo>>>,
trade_results: Arc<RwLock<HashMap<String, HashMap<String, TradeResult>>>>,
) -> Result<(), DexError> {
match msg {
Message::Text(text) => {
for line in text.split('\n') {
if line.is_empty() {
continue;
}
if let Ok(account_message) =
serde_json::from_str::<AccountWebSocketMessage>(line)
{
if let Some(account_push_data) = account_message.push {
if account_push_data.channel.starts_with("account@") {
Self::process_account_data(
&account_push_data.pub_data.data,
trade_results.clone(),
)
.await;
}
}
} else if let Ok(market_message) =
serde_json::from_str::<WebSocketMessage>(line)
{
if let Some(push_data) = market_message.push {
if push_data.channel.starts_with("market:") {
let market_id = match push_data.pub_data.data.id {
Some(v) => v,
None => return Ok(()),
};
let last_trade_price =
string_to_decimal(push_data.pub_data.data.last_trade_price)
.ok();
let market_price =
string_to_decimal(push_data.pub_data.data.market_price).ok();
let min_order =
string_to_decimal(push_data.pub_data.data.min_order).ok();
let min_tick =
string_to_decimal(push_data.pub_data.data.min_tick).ok();
log::trace!("last_trade_price = {:?}", last_trade_price);
let mut market_info_guard = market_info.write().await;
let market_info_entry = market_info_guard
.entry(market_id.to_owned())
.or_insert_with(|| MarketInfo::default());
if let Some(price) = last_trade_price {
market_info_entry.last_trade_price = Some(price);
}
if last_trade_price.is_some() {
market_info_entry.last_trade_price = last_trade_price;
}
if market_price.is_some() {
market_info_entry.market_price = market_price;
}
if min_order.is_some() {
market_info_entry.min_order = min_order;
}
if min_tick.is_some() {
market_info_entry.min_tick = min_tick;
}
}
}
}
}
}
_ => {
log::warn!("Message is empty");
}
}
Ok(())
}
async fn process_account_data(
data: &AccountData,
trade_results: Arc<RwLock<HashMap<String, HashMap<String, TradeResult>>>>,
) {
for fill in &data.fills {
log::debug!("fill: {:?}", fill);
let filled_price = match parse_to_decimal(&fill.price) {
Ok(v) => v,
Err(_) => {
log::error!("Invalid filled_price: {}", fill.price);
return;
}
};
let filled_side = match fill.side.as_str() {
"long" => OrderSide::Long,
"short" => OrderSide::Short,
_ => return,
};
let filled_size = match parse_to_decimal(&fill.size) {
Ok(v) => v,
Err(_) => {
log::error!("Invalid filled_size: {}", fill.size);
return;
}
};
let filled_fee = match parse_to_decimal(&fill.fee) {
Ok(v) => -v,
Err(_) => {
log::error!("Invalid filled_fee: {}", fill.fee);
return;
}
};
let filled_value = filled_price * filled_size;
let trade_result = TradeResult {
filled_side,
filled_size,
filled_value,
filled_fee,
};
let mut trade_results_guard = trade_results.write().await;
trade_results_guard
.entry(fill.market_id.clone())
.or_default()
.insert(fill.order_id.to_owned(), trade_result);
}
}
}
impl HyperliquidDefaultPayload {
fn new(r#type: &str, user: &str) -> Self {
Self {
r#type: r#type.to_owned(),
user: Some(user.to_owned()),
}
}
}
#[derive(Serialize, Debug)]
struct HyperliquidDefaultPayload {
r#type: String,
#[serde(skip_serializing_if = "Option::is_none")]
user: Option<String>,
}
impl HyperliquidCommonResponse {
fn is_success(&self) -> Result<(), DexError> {
match &self.response.data {
Some(data) => {
if data.statuses.len() > 1 {
match &data.statuses[0].error {
Some(e) => Err(DexError::Other(e.to_string())),
None => Ok(()),
}
} else {
Ok(())
}
}
None => Ok(()),
}
}
}
#[derive(Deserialize, Debug)]
struct HyperliquidCommonResponse {
success: String,
response: HyperliquidCommonResponseBody,
}
#[derive(Deserialize, Debug)]
struct HyperliquidCommonResponseBody {
r#type: String,
data: Option<HyperliquidCommonResponseData>,
}
#[derive(Deserialize, Debug)]
struct HyperliquidCommonResponseData {
statuses: Vec<HyperliquidCommonResponseStatus>,
}
#[derive(Deserialize, Debug)]
struct HyperliquidCommonResponseStatus {
error: Option<String>,
}
impl HyperliquidUpdateLeveragePayload {
fn new(asset: u32, leverage: u32) -> Self {
Self {
r#type: "updateLeverage".to_owned(),
asset,
is_cross: false,
leverage,
}
}
}
#[derive(Serialize, Debug)]
struct HyperliquidUpdateLeveragePayload {
r#type: String,
asset: u32,
#[serde(rename = "isCross")]
is_cross: bool,
leverage: u32,
}
#[derive(Deserialize, Debug)]
struct HyperliquidUpdateLeverageResponse {
status: String,
}
#[derive(Deserialize, Debug)]
struct HyperliquidRetrieveUserStateResponse {
margin_summary: Option<HyperliquidMarginSummary>,
}
#[derive(Deserialize, Debug)]
struct HyperliquidMarginSummary {
#[serde(rename = "accountValue")]
account_value: String,
#[serde(rename = "totalRawUsd")]
total_rawusd: String,
}
impl HyperliquidCreateOrderPayload {
fn new(asset: u32, is_buy: bool, price: Decimal, size: Decimal, order_type: &str) -> Self {
Self {
r#type: "order".to_owned(),
grouping: "na".to_owned(),
orders: vec![HyperliquidOrder::new(
asset, is_buy, price, size, order_type,
)],
}
}
}
#[derive(Serialize, Debug)]
struct HyperliquidCreateOrderPayload {
r#type: String,
grouping: String,
orders: Vec<HyperliquidOrder>,
}
impl HyperliquidOrder {
fn new(asset: u32, is_buy: bool, price: Decimal, size: Decimal, order_type: &str) -> Self {
Self {
asset,
is_buy,
limit_px: price,
sz: size,
reduce_only: false,
order_type: HyperliquidOrderType::new(order_type),
}
}
}
#[derive(Serialize, Debug)]
struct HyperliquidOrder {
asset: u32,
#[serde(rename = "isBuy")]
is_buy: bool,
#[serde(rename = "limitPx")]
limit_px: Decimal,
sz: Decimal,
#[serde(rename = "reduceOnly")]
reduce_only: bool,
#[serde(rename = "orderType")]
order_type: HyperliquidOrderType,
}
impl HyperliquidOrderType {
fn new(tif: &str) -> Self {
Self {
tif: tif.to_owned(),
}
}
}
#[derive(Serialize, Debug)]
struct HyperliquidOrderType {
tif: String,
}
#[derive(Deserialize, Debug)]
struct HyperliquidCreateOrderResponse {
status: String,
response: HyperliquidOrderResponseBody,
}
#[derive(Deserialize, Debug)]
struct HyperliquidOrderResponseBody {
r#type: String,
data: HyperliquidOrderResponseData,
}
#[derive(Deserialize, Debug)]
struct HyperliquidOrderResponseData {
statuses: Vec<HyperliquidOrderStatus>,
}
#[derive(Deserialize, Debug)]
struct HyperliquidOrderStatus {
resting: Option<HyperliquidOrderStatusDetail>,
filled: Option<HyperliquidOrderStatusDetail>,
}
#[derive(Deserialize, Debug)]
struct HyperliquidOrderStatusDetail {
oid: u32,
}
#[derive(Serialize, Debug)]
struct HyperliquidCancelOrderPayload {
r#type: String,
cancels: Vec<HyperliquidCancelOrder>,
}
#[derive(Serialize, Debug)]
struct HyperliquidCancelOrder {
assert: u32,
oid: u32,
}
#[derive(Deserialize, Debug)]
struct HyperliquidRetriveUserOpenOrder {
coin: String,
oid: u32,
side: String,
sz: String,
dir: Option<String>,
}
#[derive(Deserialize, Debug)]
struct HyperliquidRetriveUserOpenOrderResponse {
open_orders: Vec<HyperliquidRetriveUserOpenOrder>,
}
#[derive(Deserialize, Debug)]
struct HyperliquidRetriveUserFillResponse {
filled_positions: Vec<HyperliquidRetriveUserOpenOrder>,
}
#[derive(Serialize, Debug)]
struct HyperliquidRetrieveUserFillsPayload {
r#type: String,
user: String,
#[serde(rename = "startTime")]
start_time: u128,
}
#[derive(Deserialize, Debug)]
struct HyperliquidRetriveMarketMetadataResponse {
universe: Vec<HyperliquidRetriveMarketMetadata>,
}
#[derive(Deserialize, Debug)]
struct HyperliquidRetriveMarketMetadata {
name: String,
#[serde(rename = "zDecimals")]
decimals: u32,
#[serde(rename = " maxLeverage")]
max_leverage: u32,
#[serde(rename = "onlyIsolated")]
only_isolated: bool,
}
#[async_trait]
impl DexConnector for HyperliquidConnector {
async fn start(&self) -> Result<(), DexError> {
self.start_web_socket().await?;
sleep(Duration::from_secs(5)).await;
Ok(())
}
async fn stop(&self) -> Result<(), DexError> {
self.stop_web_socket().await?;
Ok(())
}
async fn set_leverage(&self, symbol: &str, leverage: u32) -> Result<(), DexError> {
let request_url = "/exchange";
let asset = self.get_asset_index(symbol).await?;
let action = HyperliquidUpdateLeveragePayload {
r#type: "updateLeverage".to_owned(),
asset,
is_cross: false,
leverage,
};
let res = self
.handle_request_with_action::<HyperliquidCommonResponse, HyperliquidUpdateLeveragePayload>(
request_url.to_string(),
&action,
)
.await?;
res.is_success()
}
async fn get_ticker(&self, symbol: &str) -> Result<TickerResponse, DexError> {
if !self.running.load(Ordering::SeqCst) {
return Err(DexError::NoConnection);
}
let market_info_guard = self.market_info.read().await;
let price;
let min_tick;
let min_order;
match market_info_guard.get(symbol) {
Some(v) => {
match v.market_price {
Some(v) => price = v,
None => return Err(DexError::Other("No price available".to_string())),
}
match v.min_tick {
Some(v) => min_tick = v,
None => return Err(DexError::Other("No min_tick available".to_string())),
}
match v.min_order {
Some(v) => min_order = v,
None => return Err(DexError::Other("No min_order available".to_string())),
}
}
None => return Err(DexError::Other("No market info available".to_string())),
};
Ok(TickerResponse {
symbol: symbol.to_owned(),
price,
min_tick,
min_order,
})
}
async fn get_filled_orders(&self, symbol: &str) -> Result<FilledOrdersResponse, DexError> {
let mut response: Vec<FilledOrder> = vec![];
let trade_results_guard = self.trade_results.read().await;
let orders = match trade_results_guard.get(symbol) {
Some(v) => v,
None => return Ok(FilledOrdersResponse::default()),
};
Ok(FilledOrdersResponse { orders: response })
}
async fn get_balance(&self) -> Result<BalanceResponse, DexError> {
todo!()
}
async fn clear_filled_order(&self, symbol: &str, order_id: &str) -> Result<(), DexError> {
let mut trade_results_guard = self.trade_results.write().await;
if let Some(orders) = trade_results_guard.get_mut(symbol) {
if orders.contains_key(order_id) {
orders.remove(order_id);
} else {
return Err(DexError::Other(format!(
"filled order(order_id:{}({})) does not exist",
order_id, symbol
)));
}
} else {
return Err(DexError::Other(format!(
"filled order(symbol:{}({})) does not exist",
symbol, order_id
)));
}
Ok(())
}
async fn create_order(
&self,
symbol: &str,
size: Decimal,
side: OrderSide,
price: Option<Decimal>,
) -> Result<CreateOrderResponse, DexError> {
let request_url = "/exchange";
let (price, r#type, time_in_force) = match price {
Some(v) => (v, "limit", "post_only"),
None => {
let price = self.get_worst_price(symbol, &side).await?;
(price, "market", "good_till_cancel")
}
};
let side_str = format!("{}", side);
let rounded_price;
let rounded_size;
{
let market_info_guard = self.market_info.read().await;
let (min_tick, min_order) = match market_info_guard.get(symbol) {
Some(v) => (v.min_tick, v.min_order),
None => return Err(DexError::Other("No price available".to_string())),
};
let min_tick = match min_tick {
Some(v) => v,
None => return Err(DexError::Other("No min_tick available".to_string())),
};
let min_order = match min_order {
Some(v) => v,
None => return Err(DexError::Other("No min_order available".to_string())),
};
rounded_price = self.round_price(price, min_tick, side);
rounded_size = self.floor_size(size, min_order);
log::debug!(
"{}, {}, {:?}({}), {:?}({})",
symbol,
price,
rounded_price,
min_tick,
rounded_size,
min_order
);
}
if rounded_size.is_zero() {
return Ok(CreateOrderResponse {
order_id: String::new(),
ordered_price: Decimal::new(0, 0),
ordered_size: Decimal::new(0, 0),
});
}
let action = HyperliquidCreateOrderPayload {
r#type: todo!(),
grouping: todo!(),
orders: todo!(),
};
let res = self
.handle_request_with_action::<HyperliquidCreateOrderResponse, HyperliquidCreateOrderPayload>(
request_url.to_string(),
&action,
)
.await?;
}
async fn cancel_order(&self, symbol: &str, order_id: &str) -> Result<(), DexError> {
let request_url = "/exchange";
let action = HyperliquidCancelOrderPayload {
r#type: todo!(),
cancels: todo!(),
};
let res = self
.handle_request_with_action::<HyperliquidCommonResponse, HyperliquidCancelOrderPayload>(
request_url.to_string(),
&action,
)
.await?;
}
async fn cancel_all_orders(&self, symbol: Option<String>) -> Result<(), DexError> {
Ok(())
}
async fn close_all_positions(&self, symbol: Option<String>) -> Result<(), DexError> {
let current_positions = self.get_positions().await?;
for position in current_positions {
}
Ok(())
}
}
#[derive(Serialize, Deserialize)]
struct PhantomAgent {
source: String,
connection_id: Vec<u8>,
}
impl PhantomAgent {
fn new(hash: &[u8], is_mainnet: bool) -> Self {
Self {
source: if is_mainnet {
"a".to_string()
} else {
"b".to_string()
},
connection_id: hash.to_vec(),
}
}
}
impl HyperliquidConnector {
fn action_hash<A: Serialize>(action: &A, vault_address: Option<&str>, nonce: u128) -> Vec<u8> {
let mut buf = Vec::new();
action
.serialize(&mut rmp_serde::Serializer::new(&mut buf))
.unwrap();
buf.extend_from_slice(&nonce.to_be_bytes());
if let Some(address) = vault_address {
buf.push(1);
buf.extend_from_slice(&hex::decode(address.trim_start_matches("0x")).unwrap());
} else {
buf.push(0);
}
let mut hasher = Keccak256::new();
hasher.update(&buf);
hasher.finalize().to_vec()
}
async fn sign_l1_action<A: Serialize>(
wallet: &Wallet<SigningKey<Secp256k1>>,
action: &A,
nonce: u128,
vault_address: Option<&str>,
is_mainnet: bool,
) -> Result<Signature, Box<dyn std::error::Error + Send + Sync>> {
let hash = Self::action_hash(action, vault_address, nonce);
let phantom_agent = PhantomAgent::new(&hash, is_mainnet);
let phantom_agent_value = serde_json::to_value(phantom_agent)?;
let mut types = Types::new();
types.insert(
"EIP712Domain".to_string(),
vec![
Eip712DomainType {
name: "name".to_string(),
r#type: "string".to_string(),
},
Eip712DomainType {
name: "version".to_string(),
r#type: "string".to_string(),
},
Eip712DomainType {
name: "chainId".to_string(),
r#type: "uint256".to_string(),
},
Eip712DomainType {
name: "verifyingContract".to_string(),
r#type: "address".to_string(),
},
],
);
types.insert(
"Agent".to_string(),
vec![
Eip712DomainType {
name: "source".to_string(),
r#type: "string".to_string(),
},
Eip712DomainType {
name: "connectionId".to_string(),
r#type: "bytes32".to_string(),
},
],
);
let domain = EIP712Domain {
name: Some("Exchange".to_string()),
version: Some("1".to_string()),
chain_id: Some(if is_mainnet { 1.into() } else { 4.into() }),
verifying_contract: Some("0x0000000000000000000000000000000000000000".parse()?),
salt: None,
};
let typed_data = TypedData {
types,
domain,
primary_type: "Agent".to_string(),
message: BTreeMap::from([
("source".to_string(), phantom_agent_value["source"].clone()),
(
"connectionId".to_string(),
phantom_agent_value["connectionId"].clone(),
),
]),
};
let signature = wallet.sign_typed_data(&typed_data).await?;
Ok(signature)
}
async fn handle_request_with_action<T, U>(
&self,
request_url: String,
action: &U,
) -> Result<T, DexError>
where
T: for<'de> Deserialize<'de>,
U: Serialize,
{
let mut nonce_lock = self.nonce.lock().await;
*nonce_lock += 1;
let nonce = *nonce_lock;
log::debug!("nonce = {}", nonce);
let signature = Self::sign_l1_action(&self.wallet, action, nonce, None, true).await;
match signature {
Ok(signature) => {
let action_string = serde_json::to_string(action).unwrap_or_default();
let json_payload = json!({
"action": action_string,
"nonce": *nonce_lock,
"signature": signature,
"vaultAddress": Option::<String>::None,
});
self.request
.handle_request::<T, U>(
HttpMethod::Post,
request_url,
&HashMap::new(),
json_payload.to_string(),
)
.await
.map_err(|e| DexError::Other(e.to_string()))
}
Err(e) => Err(DexError::Other(e.to_string())),
}
}
async fn get_positions(&self) -> Result<Vec<HyperliquidRetriveUserFillResponse>, DexError> {
todo!()
}
async fn get_orders(&self) -> Result<Vec<HyperliquidRetriveUserOpenOrderResponse>, DexError> {
todo!()
}
async fn retrive_market_metadata(&self) -> Result<(), DexError> {
let request_url = "/info";
let action = HyperliquidDefaultPayload {
r#type: "meta".to_owned(),
user: None,
};
let res = self
.handle_request_with_action::<HyperliquidRetriveMarketMetadataResponse, HyperliquidDefaultPayload>(
request_url.to_string(),
&action,
)
.await?;
let mut market_info_guard = self.market_info.write().await;
for (asset_index, metadata) in res.universe.into_iter().enumerate() {
let market_id = format!("{}-USD", metadata.name);
let market_info_entry =
market_info_guard
.entry(market_id)
.or_insert_with(|| MarketInfo {
last_trade_price: None,
market_price: None,
min_order: None,
min_tick: None,
decimals: Some(metadata.decimals),
max_leverage: Some(metadata.max_leverage),
asset_index: Some(asset_index as u32),
});
}
Ok(())
}
async fn get_worst_price(&self, symbol: &str, side: &OrderSide) -> Result<Decimal, DexError> {
let market_info_guard = self.market_info.read().await;
let last_price = match market_info_guard.get(symbol) {
Some(v) => match v.market_price {
Some(v) => v,
None => return Err(DexError::Other("Price is None".to_string())),
},
None => return Err(DexError::Other("No price available".to_string())),
};
let worst_price = slippage_price(last_price, *side == OrderSide::Long);
Ok(worst_price)
}
async fn get_asset_index(&self, symbol: &str) -> Result<u32, DexError> {
let market_info_guard = self.market_info.read().await;
match market_info_guard.get(symbol) {
Some(v) => match v.asset_index {
Some(v) => Ok(v),
None => {
return Err(DexError::Other(format!(
"Asset index is not available for {}",
symbol
)))
}
},
None => return Err(DexError::Other(format!("Unknown symbol: {}", symbol))),
}
}
}