#![cfg(feature = "lighter-sdk")]
const ORDER_TYPE_LIMIT: u32 = 0;
const ORDER_TYPE_IOC: u32 = 1;
const ORDER_TYPE_TRIGGER: u32 = 2;
const ORDER_TYPE_FOK: u32 = 2;
const SIDE_SELL: u32 = 0; const SIDE_BUY: u32 = 1;
const TIME_IN_FORCE_IOC: u32 = 0; const TIME_IN_FORCE_GTC: u32 = 1;
const PRICE_SCALE_FACTOR: u64 = 10; const BASE_AMOUNT_SCALE: u64 = 100000;
use crate::{
dex_connector::{string_to_decimal, DexConnector},
dex_request::{DexError, HttpMethod},
dex_websocket::DexWebSocket,
BalanceResponse, CanceledOrder, CanceledOrdersResponse, CombinedBalanceResponse,
CreateOrderResponse, FilledOrder, FilledOrdersResponse, LastTrade, LastTradesResponse,
OpenOrder, OpenOrdersResponse, OrderSide, TickerResponse, TpSl, TriggerOrderStyle,
};
use async_trait::async_trait;
use chrono::{DateTime, Duration as ChronoDuration, Utc};
use reqwest::Client;
use rust_decimal::prelude::{FromStr, ToPrimitive};
use rust_decimal::Decimal;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::{
collections::HashMap,
sync::{
atomic::{AtomicBool, AtomicU64, Ordering},
Arc, Mutex,
},
time::{Duration, Instant},
};
use tokio::sync::RwLock;
use tokio::task::JoinHandle;
fn is_buy_for_tpsl(position_side: OrderSide) -> bool {
matches!(position_side, OrderSide::Short)
}
struct MaintenanceInfo {
next_start: Option<DateTime<Utc>>,
}
#[derive(Clone, Debug)]
struct MarketInfo {
canonical_symbol: String,
market_id: u32,
}
#[derive(Default, Debug)]
struct MarketCache {
by_symbol: HashMap<String, MarketInfo>,
by_id: HashMap<u32, MarketInfo>,
}
#[derive(Debug)]
enum OutboundMessage {
Control(tokio_tungstenite::tungstenite::Message), }
impl OutboundMessage {
fn into_message(self) -> tokio_tungstenite::tungstenite::Message {
match self {
OutboundMessage::Control(msg) => msg,
}
}
fn is_pong(&self) -> bool {
match self {
OutboundMessage::Control(tokio_tungstenite::tungstenite::Message::Pong(_)) => true,
_ => false,
}
}
}
fn normalize_symbol(symbol: &str) -> String {
let upper = symbol.trim().to_ascii_uppercase();
let mut normalized = upper
.replace("-PERP", "")
.replace("_PERP", "")
.replace(".PERP", "")
.replace("-USD", "")
.replace("_USD", "")
.replace("/USD", "")
.replace("-USDC", "")
.replace("_USDC", "")
.replace("/USDC", "");
if normalized.ends_with("-PERP") {
normalized = normalized.trim_end_matches("-PERP").to_string();
}
normalized
}
#[cfg(feature = "lighter-sdk")]
use libc::{c_char, c_int, c_longlong};
use secp256k1::{Message, Secp256k1, SecretKey};
use sha3::{Digest, Keccak256};
#[cfg(feature = "lighter-sdk")]
use std::ffi::{CStr, CString};
use tokio::time::sleep;
use tokio_tungstenite;
#[cfg(feature = "lighter-sdk")]
#[repr(C)]
pub struct StrOrErr {
pub str: *mut c_char,
pub err: *mut c_char,
}
#[cfg(feature = "lighter-sdk")]
extern "C" {
fn CreateClient(
url: *const c_char,
private_key: *const c_char,
chain_id: c_int,
api_key_index: c_int,
account_index: c_longlong,
) -> *mut c_char;
fn CheckClient(api_key_index: c_int, account_index: c_longlong) -> *mut c_char;
fn GetClientPubKey(api_key_index: c_int, account_index: c_longlong) -> *mut c_char;
fn SignCreateOrder(
market_index: c_int,
client_order_index: c_longlong,
base_amount: c_longlong,
price: c_int,
is_ask: c_int,
order_type: c_int,
time_in_force: c_int,
reduce_only: c_int,
trigger_price: c_int,
order_expiry: c_longlong,
nonce: c_longlong,
) -> StrOrErr;
fn SignChangePubKey(new_pubkey: *const c_char, nonce: c_longlong) -> StrOrErr;
fn SignCancelAllOrders(time_in_force: c_int, time: c_longlong, nonce: c_longlong) -> StrOrErr;
fn SignMessageWithEVM(private_key: *const c_char, message: *const c_char) -> StrOrErr;
}
static API_CALL_COUNTER: AtomicU64 = AtomicU64::new(0);
static API_CALL_TRACKER: std::sync::LazyLock<Mutex<Vec<(Instant, String)>>> =
std::sync::LazyLock::new(|| Mutex::new(Vec::new()));
fn track_api_call(endpoint: &str, method: &str) {
let call_count = API_CALL_COUNTER.fetch_add(1, Ordering::SeqCst) + 1;
let now = Instant::now();
{
let mut tracker = API_CALL_TRACKER.lock().unwrap();
tracker.retain(|(time, _)| now.duration_since(*time) < Duration::from_secs(60));
tracker.push((now, format!("{} {}", method, endpoint)));
let recent_calls = tracker.len();
log::info!(
"[API_TRACKER] #{} {} {} | Recent calls (60s): {} | Rate: {:.1}/min",
call_count,
method,
endpoint,
recent_calls,
recent_calls as f64
);
if recent_calls > 45 {
log::warn!(
"[API_TRACKER] ⚠️ Approaching rate limit: {}/60 calls in last 60s",
recent_calls
);
}
}
}
#[derive(Clone)]
pub struct LighterConnector {
api_key_public: String, api_key_index: u32, api_private_key_hex: String, #[cfg(feature = "lighter-sdk")]
evm_wallet_private_key: Option<String>, account_index: u32, base_url: String,
websocket_url: String,
_l1_address: String, client: Client,
filled_orders: Arc<RwLock<HashMap<String, Vec<FilledOrder>>>>,
canceled_orders: Arc<RwLock<HashMap<String, Vec<CanceledOrder>>>>,
cached_server_pubkey: Arc<tokio::sync::RwLock<Option<(String, std::time::Instant)>>>,
is_running: Arc<AtomicBool>,
cleanup_started: Arc<AtomicBool>,
cleanup_handle: Arc<tokio::sync::Mutex<Option<JoinHandle<()>>>>,
_ws: Option<DexWebSocket>, current_price: Arc<RwLock<Option<(Decimal, u64)>>>, current_volume: Arc<RwLock<Option<Decimal>>>,
order_book: Arc<RwLock<Option<LighterOrderBook>>>,
maintenance: Arc<RwLock<MaintenanceInfo>>,
cached_open_orders: Arc<RwLock<HashMap<String, Vec<OpenOrder>>>>, connection_epoch: Arc<AtomicU64>,
market_cache: Arc<RwLock<MarketCache>>,
tracked_symbols: Vec<String>,
}
#[derive(Deserialize, Debug, Clone)]
struct LighterOrderBook {
bids: Vec<LighterOrderBookEntry>,
asks: Vec<LighterOrderBookEntry>,
}
#[derive(Deserialize, Debug, Clone)]
struct LighterOrderBookEntry {
price: String,
size: String,
}
#[derive(Deserialize, Debug)]
#[allow(dead_code)]
struct LighterAccountResponse {
code: i32,
total: i32,
accounts: Vec<LighterAccountInfo>,
}
#[derive(Deserialize, Debug)]
#[allow(dead_code)]
struct LighterAccountInfo {
account_index: i64,
available_balance: String,
collateral: String,
total_asset_value: String,
positions: Vec<LighterPosition>,
}
#[derive(Deserialize, Debug)]
#[allow(dead_code)]
struct LighterPosition {
market_id: u8,
symbol: String,
position: String,
sign: i8,
open_order_count: u32,
avg_entry_price: String,
}
#[derive(Deserialize, Debug)]
#[allow(dead_code)]
struct LighterTradesResponse {
code: i32,
trades: Vec<LighterTrade>,
}
#[derive(Deserialize, Debug)]
#[allow(dead_code)]
struct LighterTrade {
trade_id: u64,
price: String,
size: String,
usd_amount: String,
market_id: u8,
}
#[derive(Deserialize, Debug)]
#[allow(dead_code)]
struct LighterExchangeStats {
code: i32,
order_book_stats: Vec<LighterOrderBookStats>,
daily_usd_volume: f64,
daily_trades_count: u32,
}
#[derive(Deserialize, Debug)]
#[allow(dead_code)]
struct LighterOrderBookStats {
symbol: String,
last_trade_price: f64,
daily_trades_count: u32,
daily_base_token_volume: f64,
daily_quote_token_volume: f64,
daily_price_change: f64,
}
#[derive(Deserialize, Debug)]
#[allow(dead_code)]
struct LighterFundingRates {
code: i32,
funding_rates: Vec<LighterFundingRate>,
}
#[derive(Deserialize, Debug)]
#[allow(dead_code)]
struct LighterFundingRate {
market_id: u32,
exchange: String,
symbol: String,
rate: f64,
}
#[allow(dead_code)]
#[derive(Deserialize, Debug)]
struct LighterNonceResponse {
nonce: u64,
}
#[derive(Deserialize, Debug)]
struct ApiKeyInfo {
#[serde(rename = "account_index")]
#[allow(dead_code)]
account_index: u32,
#[serde(rename = "api_key_index")]
#[allow(dead_code)]
api_key_index: u32,
#[allow(dead_code)]
nonce: u32,
#[serde(rename = "public_key")]
public_key: String,
}
#[derive(Deserialize, Debug)]
struct ApiKeyResponse {
#[allow(dead_code)]
code: u32,
#[serde(rename = "api_keys")]
api_keys: Vec<ApiKeyInfo>,
}
#[allow(dead_code)]
#[derive(Deserialize, Debug)]
struct LighterOrderResponse {
order_id: String,
price: String,
amount: String,
}
#[allow(dead_code)]
#[derive(Serialize, Debug)]
struct LighterTx {
tx_type: String,
ticker: String,
amount: String,
price: Option<String>,
order_type: String,
time_in_force: String,
}
#[allow(dead_code)]
#[derive(Serialize, Debug)]
struct LighterSignedEnvelope {
sig: String,
nonce: u64,
tx: LighterTx,
}
impl LighterConnector {
async fn refresh_market_cache(&self) -> Result<(), DexError> {
let funding = self.get_funding_rates().await?;
let mut cache = MarketCache::default();
for entry in funding.funding_rates {
let normalized = normalize_symbol(&entry.symbol);
if normalized.is_empty() {
continue;
}
let info = MarketInfo {
canonical_symbol: normalized.clone(),
market_id: entry.market_id,
};
cache.by_symbol.insert(normalized.clone(), info.clone());
cache.by_id.insert(entry.market_id, info);
}
if cache.by_symbol.is_empty() {
return Err(DexError::Other(
"Funding rates response did not contain any markets".to_string(),
));
}
*self.market_cache.write().await = cache;
Ok(())
}
async fn resolve_market_info(&self, symbol: &str) -> Result<MarketInfo, DexError> {
let normalized = normalize_symbol(symbol);
{
let cache = self.market_cache.read().await;
if let Some(info) = cache.by_symbol.get(&normalized) {
return Ok(info.clone());
}
}
self.refresh_market_cache().await?;
let cache = self.market_cache.read().await;
cache
.by_symbol
.get(&normalized)
.cloned()
.ok_or_else(|| DexError::Other(format!("Unknown symbol: {}", symbol)))
}
#[cfg(feature = "lighter-sdk")]
async fn create_go_client(&self) -> Result<(), DexError> {
unsafe {
let url = CString::new(self.base_url.as_str())
.map_err(|e| DexError::Other(format!("Invalid URL: {}", e)))?;
let private_key_hex = self
.api_private_key_hex
.strip_prefix("0x")
.unwrap_or(&self.api_private_key_hex);
if private_key_hex.len() != 80 {
return Err(DexError::Other(format!(
"API private key must be 40 bytes (80 hex chars), got: {}",
private_key_hex.len()
)));
}
let private_key = CString::new(private_key_hex)
.map_err(|e| DexError::Other(format!("Invalid private key: {}", e)))?;
let result = CreateClient(
url.as_ptr(),
private_key.as_ptr(),
304, self.api_key_index as c_int,
self.account_index as c_longlong,
);
if !result.is_null() {
let error_cstr = CStr::from_ptr(result);
let error_msg = error_cstr.to_string_lossy().to_string();
libc::free(result as *mut libc::c_void);
return Err(DexError::Other(format!(
"CreateClient error: {}",
error_msg
)));
}
let server_pubkey = match self.get_server_public_key().await {
Ok(pubkey) => pubkey,
Err(e) => {
log::error!("Failed to get server public key: {}", e);
return Err(e);
}
};
let go_pubkey_result = GetClientPubKey(
self.api_key_index as c_int,
self.account_index as c_longlong,
);
let go_derived_pubkey = if !go_pubkey_result.is_null() {
let pubkey_cstr = CStr::from_ptr(go_pubkey_result);
let pubkey_str = pubkey_cstr.to_string_lossy().to_string();
libc::free(go_pubkey_result as *mut libc::c_void);
Some(pubkey_str)
} else {
log::error!("Failed to get public key from Go client");
None
};
if let Some(go_key) = &go_derived_pubkey {
let srv = server_pubkey
.to_lowercase()
.trim_start_matches("0x")
.to_string();
let loc = go_key.to_lowercase().trim_start_matches("0x").to_string();
if loc != srv {
log::debug!(
"API key mismatch detected (account={}, index={}). server={}…{} vs local={}…{} — will attempt ChangePubKey",
self.account_index,
self.api_key_index,
&srv[..8], &srv[srv.len()-8..],
&loc[..8], &loc[loc.len()-8..]
);
} else {
}
}
let check_result = CheckClient(
self.api_key_index as c_int,
self.account_index as c_longlong,
);
if !check_result.is_null() {
let error_cstr = CStr::from_ptr(check_result);
let error_msg = error_cstr.to_string_lossy().to_string();
libc::free(check_result as *mut libc::c_void);
log::error!("API key validation failed: {}", error_msg);
if error_msg.contains("ownPubKey:") && error_msg.contains("PublicKey:") {
if let Some(own_start) = error_msg.find("ownPubKey: ") {
if let Some(own_end) = error_msg[own_start + 11..].find(" ") {
let own_key = &error_msg[own_start + 11..own_start + 11 + own_end];
log::error!(
" Our derived public key (first 8): {}",
&own_key[..std::cmp::min(8, own_key.len())]
);
log::error!(
" Our derived public key (last 8): {}",
&own_key[std::cmp::max(0, own_key.len().saturating_sub(8))..]
);
}
}
if let Some(resp_start) = error_msg.find("PublicKey:") {
if let Some(resp_end) = error_msg[resp_start + 10..].find("}") {
let resp_key = &error_msg[resp_start + 10..resp_start + 10 + resp_end];
log::error!(
" Server expected public key (first 8): {}",
&resp_key[..std::cmp::min(8, resp_key.len())]
);
log::error!(
" Server expected public key (last 8): {}",
&resp_key[std::cmp::max(0, resp_key.len().saturating_sub(8))..]
);
}
}
}
#[cfg(feature = "lighter-sdk")]
if let (Some(_), Some(_)) = (&go_derived_pubkey, &self.evm_wallet_private_key) {
return Err(DexError::ApiKeyRegistrationRequired);
} else {
return Err(DexError::Other(format!(
"API key validation failed: {}",
error_msg
)));
}
#[cfg(not(feature = "lighter-sdk"))]
return Err(DexError::Other(format!(
"API key validation failed: {}",
error_msg
)));
}
Ok(())
}
}
pub fn start_auto_cleanup(&self, cleanup_interval_hours: u64) {
if self.cleanup_started.swap(true, Ordering::SeqCst) {
log::warn!("[AUTO_CLEANUP] already started; ignoring.");
return;
}
log::info!(
"[AUTO_CLEANUP] Starting background task (interval: {}h)",
cleanup_interval_hours
);
let filled_orders = Arc::clone(&self.filled_orders);
let canceled_orders = Arc::clone(&self.canceled_orders);
let is_running = Arc::clone(&self.is_running);
let cleanup_started = Arc::clone(&self.cleanup_started);
let cleanup_handle = Arc::clone(&self.cleanup_handle);
let handle = tokio::spawn(async move {
let mut interval =
tokio::time::interval(Duration::from_secs(cleanup_interval_hours * 3600));
interval.tick().await;
while is_running.load(Ordering::Relaxed) {
interval.tick().await;
let mut filled_removed = 0usize;
let mut canceled_removed = 0usize;
{
let mut filled = filled_orders.write().await;
for (symbol, orders) in filled.iter_mut() {
const KEEP_FILLED_PER_SYMBOL: usize = 50;
if orders.len() > KEEP_FILLED_PER_SYMBOL {
let remove_count = orders.len() - KEEP_FILLED_PER_SYMBOL;
orders.drain(0..remove_count);
filled_removed += remove_count;
log::debug!(
"🗑️ [AUTO_CLEANUP] Removed {} old filled orders for {} (kept {})",
remove_count,
symbol,
KEEP_FILLED_PER_SYMBOL
);
}
}
filled.retain(|_, orders| !orders.is_empty());
}
{
let mut canceled = canceled_orders.write().await;
let cutoff_secs = (Utc::now() - ChronoDuration::hours(24)).timestamp() as u64;
for (symbol, orders) in canceled.iter_mut() {
let initial_len = orders.len();
orders.retain(|order| order.canceled_timestamp > cutoff_secs);
let removed = initial_len.saturating_sub(orders.len());
canceled_removed += removed;
if removed > 0 {
log::debug!(
"🗑️ [AUTO_CLEANUP] Removed {} old canceled orders for {}",
removed,
symbol
);
}
}
canceled.retain(|_, orders| !orders.is_empty());
}
let total_removed = filled_removed + canceled_removed;
if total_removed > 0 {
log::info!(
"🗑️ [AUTO_CLEANUP] removed total={} (filled={}, canceled={})",
total_removed,
filled_removed,
canceled_removed
);
}
}
cleanup_started.store(false, Ordering::SeqCst);
let mut guard = cleanup_handle.lock().await;
*guard = None;
log::info!("🛑 [AUTO_CLEANUP] task exited, ready for restart");
});
let cleanup_handle_for_storage = Arc::clone(&self.cleanup_handle);
tokio::spawn(async move {
let mut guard = cleanup_handle_for_storage.lock().await;
*guard = Some(handle);
});
}
#[cfg(not(feature = "lighter-sdk"))]
async fn create_go_client(&self) -> Result<(), DexError> {
Err(DexError::Other(
"Lighter Go SDK not available. Build with --features lighter-sdk to enable."
.to_string(),
))
}
#[cfg(feature = "lighter-sdk")]
async fn call_go_sign_create_order(
&self,
market_index: i32,
client_order_index: i64,
base_amount: i64,
price: i32,
is_ask: i32,
order_type: i32,
time_in_force: i32,
reduce_only: i32,
trigger_price: i32,
order_expiry: i64,
nonce: i64,
) -> Result<String, DexError> {
self.create_go_client().await?;
unsafe {
let result = SignCreateOrder(
market_index,
client_order_index,
base_amount,
price,
is_ask,
order_type,
time_in_force,
reduce_only,
trigger_price,
order_expiry,
nonce,
);
if !result.err.is_null() {
let error_cstr = CStr::from_ptr(result.err);
let error_msg = error_cstr.to_string_lossy().to_string();
libc::free(result.err as *mut libc::c_void);
if !result.str.is_null() {
libc::free(result.str as *mut libc::c_void);
}
return Err(DexError::Other(format!("Go SDK error: {}", error_msg)));
}
if result.str.is_null() {
return Err(DexError::Other("Go SDK returned null result".to_string()));
}
let result_cstr = CStr::from_ptr(result.str);
let json_str = result_cstr.to_string_lossy().to_string();
libc::free(result.str as *mut libc::c_void);
Ok(json_str)
}
}
#[cfg(not(feature = "lighter-sdk"))]
async fn call_go_sign_create_order(
&self,
_market_index: i32,
_client_order_index: i64,
_base_amount: i64,
_price: i32,
_is_ask: i32,
_order_type: i32,
_time_in_force: i32,
_reduce_only: i32,
_trigger_price: i32,
_order_expiry: i64,
_nonce: i64,
) -> Result<String, DexError> {
Err(DexError::Other(
"Lighter Go SDK not available. Build with --features lighter-sdk to enable."
.to_string(),
))
}
#[cfg(feature = "lighter-sdk")]
pub fn new(
api_key_public: String,
api_key_index: u32,
api_private_key_hex: String,
evm_wallet_private_key: Option<String>,
account_index: u32,
base_url: String,
websocket_url: String,
tracked_symbols: Vec<String>,
) -> Result<Self, DexError> {
let l1_address = "N/A".to_string();
log::debug!(
"Creating LighterConnector with API key index: {}, account: {}",
api_key_index,
account_index
);
Ok(Self {
api_key_public,
api_key_index,
api_private_key_hex,
evm_wallet_private_key,
account_index,
base_url: base_url.clone(),
websocket_url: websocket_url.clone(),
_l1_address: l1_address,
client: Client::new(),
filled_orders: Arc::new(RwLock::new(HashMap::new())),
canceled_orders: Arc::new(RwLock::new(HashMap::new())),
cached_server_pubkey: Arc::new(tokio::sync::RwLock::new(None)),
is_running: Arc::new(AtomicBool::new(false)),
cleanup_started: Arc::new(AtomicBool::new(false)),
cleanup_handle: Arc::new(tokio::sync::Mutex::new(None)),
_ws: Some(DexWebSocket::new(websocket_url)),
current_price: Arc::new(RwLock::new(None)),
current_volume: Arc::new(RwLock::new(None)),
order_book: Arc::new(RwLock::new(None)),
maintenance: Arc::new(RwLock::new(MaintenanceInfo { next_start: None })),
cached_open_orders: Arc::new(RwLock::new(HashMap::new())),
connection_epoch: Arc::new(AtomicU64::new(0)),
market_cache: Arc::new(RwLock::new(MarketCache::default())),
tracked_symbols,
})
}
#[cfg(not(feature = "lighter-sdk"))]
pub fn new(
api_key_public: String,
api_key_index: u32,
api_private_key_hex: String,
_evm_wallet_private_key: Option<String>,
account_index: u32,
base_url: String,
websocket_url: String,
tracked_symbols: Vec<String>,
) -> Result<Self, DexError> {
let l1_address = "N/A".to_string();
log::debug!(
"Creating LighterConnector with API key index: {}, account: {}",
api_key_index,
account_index
);
Ok(Self {
api_key_public,
api_key_index,
api_private_key_hex,
account_index,
base_url: base_url.clone(),
websocket_url: websocket_url.clone(),
_l1_address: l1_address,
client: Client::new(),
filled_orders: Arc::new(RwLock::new(HashMap::new())),
canceled_orders: Arc::new(RwLock::new(HashMap::new())),
cached_server_pubkey: Arc::new(tokio::sync::RwLock::new(None)),
is_running: Arc::new(AtomicBool::new(false)),
cleanup_started: Arc::new(AtomicBool::new(false)),
cleanup_handle: Arc::new(tokio::sync::Mutex::new(None)),
_ws: Some(DexWebSocket::new(websocket_url)),
current_price: Arc::new(RwLock::new(None)),
current_volume: Arc::new(RwLock::new(None)),
order_book: Arc::new(RwLock::new(None)),
maintenance: Arc::new(RwLock::new(MaintenanceInfo { next_start: None })),
cached_open_orders: Arc::new(RwLock::new(HashMap::new())),
connection_epoch: Arc::new(AtomicU64::new(0)),
market_cache: Arc::new(RwLock::new(MarketCache::default())),
tracked_symbols,
})
}
async fn get_server_public_key_cached(&self) -> Result<String, DexError> {
{
let cache = self.cached_server_pubkey.read().await;
if let Some((pubkey, timestamp)) = &*cache {
if timestamp.elapsed() < std::time::Duration::from_secs(300) {
log::debug!("[API_CACHE] Using cached server public key, no API call needed");
return Ok(pubkey.clone());
}
}
}
let endpoint = format!(
"/api/v1/apikeys?account_index={}&api_key_index={}",
self.account_index, self.api_key_index
);
log::debug!("Getting server public key from: {}", endpoint);
let response: ApiKeyResponse = self
.make_request(&endpoint, crate::dex_request::HttpMethod::Get, None)
.await?;
if response.api_keys.is_empty() {
return Err(DexError::Other("No API keys found on server".to_string()));
}
let server_pubkey = response.api_keys[0].public_key.clone();
{
let mut cache = self.cached_server_pubkey.write().await;
*cache = Some((server_pubkey.clone(), std::time::Instant::now()));
}
Ok(server_pubkey)
}
async fn create_order_native_with_type(
&self,
market_id: u32,
side: u32,
tif: u32,
base_amount: u64,
price: u64,
client_order_id: Option<String>,
order_type: u32,
reduce_only: bool,
expiry_secs: Option<u64>,
) -> Result<CreateOrderResponse, DexError> {
let timestamp = chrono::Utc::now().timestamp_millis() as u64;
let _client_id = client_order_id.unwrap_or_else(|| format!("rust-native-{}", timestamp));
let nonce = self.get_nonce().await?;
log::debug!(
"Creating native order: market_id={}, side={}, base_amount={}, price={}, calculated_price_usd={:.1}",
market_id,
side,
base_amount,
price,
price as f64 / 10.0
);
let client_order_index = timestamp; let order_type_param = order_type as u64; let time_in_force = tif as u64; let reduce_only_param = if reduce_only { 1u64 } else { 0u64 };
let trigger_price = 0u64; let order_expiry = if order_type == ORDER_TYPE_IOC || tif == 0 {
0i64 } else {
let expiry_duration_ms = if let Some(expiry_secs) = expiry_secs {
expiry_secs * 1000 } else {
24 * 60 * 60 * 1000 };
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as i64;
now_ms + (expiry_duration_ms as i64)
};
let actual_market_id = market_id as u64;
let actual_base_amount = base_amount;
let actual_price = price;
let actual_side = side as u64;
let _tx_data = [
actual_market_id, client_order_index, actual_base_amount, actual_price, actual_side, order_type_param, time_in_force, reduce_only_param, trigger_price, order_expiry as u64, nonce, ];
let private_key_hex = self
.api_private_key_hex
.strip_prefix("0x")
.unwrap_or(&self.api_private_key_hex);
let private_key_bytes = hex::decode(private_key_hex)
.map_err(|e| DexError::Other(format!("Invalid private key hex: {}", e)))?;
let mut key_bytes = [0u8; 40];
let copy_len = std::cmp::min(private_key_bytes.len(), 40);
key_bytes[..copy_len].copy_from_slice(&private_key_bytes[..copy_len]);
let go_result = self
.call_go_sign_create_order(
actual_market_id as i32,
client_order_index as i64,
actual_base_amount as i64,
actual_price as i32,
actual_side as i32,
order_type_param as i32,
time_in_force as i32,
reduce_only as i32,
trigger_price as i32,
order_expiry as i64,
nonce as i64,
)
.await?;
log::debug!("=== GO SDK RESULT ===");
log::debug!("Go SDK JSON: {}", go_result);
let tx_info = go_result;
let form_data = format!(
"tx_type=14&tx_info={}&price_protection=false",
urlencoding::encode(&tx_info)
);
log::debug!("=== REQUEST DEBUG ===");
log::debug!("Timestamp: {}", timestamp);
log::debug!("TX Info JSON: {}", tx_info);
log::debug!("Form data: {}", form_data);
track_api_call("POST /api/v1/sendTx", "POST");
let response = self
.client
.post(&format!("{}/api/v1/sendTx", self.base_url))
.header("Content-Type", "application/x-www-form-urlencoded")
.body(form_data)
.send()
.await
.map_err(|e| DexError::Other(format!("HTTP request failed: {}", e)))?;
let status = response.status();
let response_text = response
.text()
.await
.map_err(|e| DexError::Other(format!("Failed to read response: {}", e)))?;
log::debug!(
"Native order response: HTTP {}, Body: {}",
status,
response_text
);
if status.is_success() {
log::debug!("Native order submitted successfully!");
let order_id = client_order_index.to_string();
log::debug!(
"Created order: order_id={}, client_order_index={}, side={}, size={}",
order_id,
client_order_index,
side,
Decimal::new(base_amount as i64, 5)
);
Ok(CreateOrderResponse {
order_id,
ordered_price: Decimal::new(price as i64, 6),
ordered_size: Decimal::new(base_amount as i64, 5),
})
} else {
Err(DexError::Other(format!(
"Order failed: HTTP {}, {}",
status, response_text
)))
}
}
async fn create_order_native_with_trigger(
&self,
market_id: u32,
side: u32,
tif: u32,
base_amount: u64,
price: u64,
trigger_price: u64,
client_order_id: Option<String>,
order_type: u32,
reduce_only: bool,
expiry_secs: Option<u64>,
) -> Result<CreateOrderResponse, DexError> {
let timestamp = chrono::Utc::now().timestamp_millis() as u64;
let _client_id = client_order_id.unwrap_or_else(|| format!("rust-trigger-{}", timestamp));
let nonce = self.get_nonce().await?;
log::debug!(
"Creating trigger order: market_id={}, side={}, base_amount={}, price={}, trigger_price={}, order_type={}",
market_id,
side,
base_amount,
price,
trigger_price,
order_type
);
let client_order_index = timestamp;
let order_type_param = order_type as u64;
let time_in_force = tif as u64;
let reduce_only_param = if reduce_only { 1u64 } else { 0u64 };
let trigger_price_param = trigger_price;
let order_expiry = if order_type == ORDER_TYPE_TRIGGER
|| order_type == 4
|| order_type == 3
|| order_type == 5
{
let expiry_duration_ms = if let Some(expiry_secs) = expiry_secs {
let min_expiry_secs = 60;
std::cmp::max(min_expiry_secs, expiry_secs) * 1000
} else {
28 * 24 * 60 * 60 * 1000 };
(chrono::Utc::now().timestamp_millis() as u64 + expiry_duration_ms) as i64
} else {
let expiry_duration_ms = if let Some(expiry_secs) = expiry_secs {
expiry_secs * 1000
} else {
24 * 60 * 60 * 1000 };
(chrono::Utc::now().timestamp_millis() as u64 + expiry_duration_ms) as i64
};
let actual_market_id = market_id as u64;
let actual_base_amount = base_amount;
let actual_price = price;
let actual_side = side as u64;
let private_key_hex = self
.api_private_key_hex
.strip_prefix("0x")
.unwrap_or(&self.api_private_key_hex);
let private_key_bytes = hex::decode(private_key_hex)
.map_err(|e| DexError::Other(format!("Invalid private key hex: {}", e)))?;
let mut key_bytes = [0u8; 40];
let copy_len = std::cmp::min(private_key_bytes.len(), 40);
key_bytes[..copy_len].copy_from_slice(&private_key_bytes[..copy_len]);
let go_result = self
.call_go_sign_create_order(
actual_market_id as i32,
client_order_index as i64,
actual_base_amount as i64,
actual_price as i32,
actual_side as i32,
order_type_param as i32,
time_in_force as i32,
reduce_only_param as i32,
trigger_price_param as i32,
order_expiry,
nonce as i64,
)
.await;
let signature = match go_result {
Ok(sig) => sig,
Err(e) => {
log::error!("Failed to sign trigger order via Go SDK: {}", e);
return Err(DexError::Other(format!(
"Signature generation failed: {}",
e
)));
}
};
let form_data = format!(
"tx_type=14&tx_info={}&price_protection=false",
urlencoding::encode(&signature)
);
log::debug!("Trigger order form data: {}", form_data);
let client = &self.client;
let url = format!("{}/api/v1/sendTx", self.base_url);
let response = client
.post(&url)
.header("Content-Type", "application/x-www-form-urlencoded")
.body(form_data)
.send()
.await
.map_err(|e| DexError::Other(e.to_string()))?;
let status = response.status();
let response_text = response
.text()
.await
.map_err(|e| DexError::Other(e.to_string()))?;
log::debug!("Trigger order response: HTTP {}, {}", status, response_text);
if status.is_success() {
let order_id = format!("trigger-{}-{}", timestamp, market_id);
log::info!(
"✅ [TRIGGER_ORDER] Successfully created trigger order: {} (type={}, trigger_price={})",
order_id,
order_type,
trigger_price_param
);
Ok(CreateOrderResponse {
order_id,
ordered_price: Decimal::new(price as i64, 6),
ordered_size: Decimal::new(base_amount as i64, 5),
})
} else {
Err(DexError::Other(format!(
"Trigger order failed: HTTP {}, {}",
status, response_text
)))
}
}
#[allow(dead_code)]
async fn send_order_via_sdk(
&self,
market_id: u32,
side: u32,
tif: u32,
base_amount: u64,
price: u64,
client_order_id: Option<String>,
) -> Result<CreateOrderResponse, DexError> {
let timestamp = chrono::Utc::now().timestamp_millis() as u64;
let client_id = client_order_id.unwrap_or_else(|| format!("rust-order-{}", timestamp));
log::debug!(
"Delegating order to Python SDK: market_id={}, side={}, base_amount={}, price={}",
market_id,
side,
base_amount,
price
);
let output = std::process::Command::new("./venv/bin/python")
.arg("sdk_send_order.py")
.arg(&format!("--market-id={}", market_id))
.arg(&format!("--side={}", side))
.arg(&format!("--tif={}", tif))
.arg(&format!("--base-amt={}", base_amount))
.arg(&format!("--price={}", price))
.arg(&format!("--client-id={}", client_id))
.env("LIGHTER_ACCOUNT_INDEX", &self.account_index.to_string())
.env("LIGHTER_API_KEY_INDEX", &self.api_key_index.to_string())
.env("LIGHTER_PRIVATE_API_KEY", &self.api_private_key_hex)
.current_dir(".")
.output()
.map_err(|e| DexError::Other(format!("Failed to execute SDK script: {}", e)))?;
let stdout = String::from_utf8_lossy(&output.stdout);
let stderr = String::from_utf8_lossy(&output.stderr);
if !output.status.success() {
log::error!("SDK delegation failed. stderr: {}", stderr);
return Err(DexError::Other(format!("SDK execution failed: {}", stderr)));
}
if !stderr.is_empty() {
log::warn!("SDK delegation warnings: {}", stderr);
}
let response: serde_json::Value = serde_json::from_str(&stdout)
.map_err(|e| DexError::Other(format!("Failed to parse SDK response: {}", e)))?;
if let Some(true) = response.get("success").and_then(|v| v.as_bool()) {
log::debug!("Order successfully sent via SDK");
let order_id = response
.get("tx_hash")
.and_then(|v| v.as_str())
.unwrap_or(&client_id)
.to_string();
Ok(CreateOrderResponse {
order_id,
ordered_price: Decimal::new(price as i64, 6), ordered_size: Decimal::new(base_amount as i64, 5), })
} else if let Some(error) = response.get("error") {
log::error!("SDK order failed: {}", error);
Err(DexError::Other(format!("SDK order error: {}", error)))
} else {
Err(DexError::Other(
"Unexpected SDK response format".to_string(),
))
}
}
#[allow(dead_code)]
async fn get_nonce(&self) -> Result<u64, DexError> {
self.get_nonce_with_key(&self.api_key_public).await
}
async fn get_nonce_with_key(&self, api_key: &str) -> Result<u64, DexError> {
let url = format!(
"{}/api/v1/nextNonce?account_index={}&api_key_index={}",
self.base_url, self.account_index, self.api_key_index
);
log::debug!("Getting nonce from: {}", url);
log::debug!("Using API key: {}", api_key);
track_api_call("/api/v1/nextNonce", "GET");
let response = self
.client
.get(&url)
.header("X-API-KEY", api_key)
.send()
.await
.map_err(|e| DexError::Other(format!("Failed to get nonce: {}", e)))?;
if !response.status().is_success() {
let status = response.status();
let error_body = response
.text()
.await
.unwrap_or_else(|_| "Failed to read error response".to_string());
log::error!(
"Nonce request failed: HTTP {}, Body: {}",
status,
error_body
);
return Err(DexError::Other(format!(
"Failed to get nonce: HTTP {}, Body: {}",
status, error_body
)));
}
let nonce_response: LighterNonceResponse = response
.json()
.await
.map_err(|e| DexError::Other(format!("Failed to parse nonce response: {}", e)))?;
Ok(nonce_response.nonce)
}
#[allow(dead_code)]
async fn discover_account_index(&self) -> Result<u32, DexError> {
Ok(self.account_index)
}
async fn get_server_public_key(&self) -> Result<String, DexError> {
self.get_server_public_key_cached().await
}
async fn make_request<T>(
&self,
endpoint: &str,
method: HttpMethod,
body: Option<&str>,
) -> Result<T, DexError>
where
T: for<'de> serde::Deserialize<'de>,
{
let method_str = match method {
HttpMethod::Get => "GET",
HttpMethod::Post => "POST",
HttpMethod::Put => "PUT",
HttpMethod::Delete => "DELETE",
};
track_api_call(endpoint, method_str);
let url = format!("{}{}", self.base_url, endpoint);
let mut request = match method {
HttpMethod::Get => self.client.get(&url),
HttpMethod::Post => self.client.post(&url),
HttpMethod::Put => self.client.put(&url),
HttpMethod::Delete => self.client.delete(&url),
};
request = request.header("X-API-KEY", &self.api_key_public);
if let Some(body_content) = body {
request = request
.header("Content-Type", "application/json")
.body(body_content.to_string());
}
let response = request
.send()
.await
.map_err(|e| DexError::Other(format!("Request failed: {}", e)))?;
let status = response.status();
if !status.is_success() {
let error_text = response
.text()
.await
.unwrap_or_else(|_| "Unknown error".to_string());
return Err(DexError::Other(format!("HTTP {}: {}", status, error_text)));
}
response
.json()
.await
.map_err(|e| DexError::Other(format!("Failed to parse response: {}", e)))
}
#[cfg(feature = "lighter-sdk")]
async fn register_api_key(
&self,
evm_private_key: &str,
go_public_key: &str,
server_public_key: &str,
) -> Result<(), String> {
log::debug!(
"Attempting ChangePubKey: server='{}' -> local='{}'",
server_public_key,
go_public_key
);
let nonce = self
.get_nonce_with_key(server_public_key)
.await
.map_err(|e| format!("Failed to get nonce: {:?}", e))?;
log::debug!("Got nonce for ChangePubKey: {}", nonce);
let new_pubkey = if go_public_key.starts_with("0x") {
go_public_key.to_string()
} else {
format!("0x{}", go_public_key)
};
log::debug!("New public key to register: {}", new_pubkey);
let sign_result = unsafe {
SignChangePubKey(
std::ffi::CString::new(new_pubkey.clone()).unwrap().as_ptr(),
nonce as c_longlong,
)
};
if !sign_result.err.is_null() {
let error_msg = unsafe { std::ffi::CStr::from_ptr(sign_result.err).to_string_lossy() };
return Err(format!("Failed to sign ChangePubKey: {}", error_msg));
}
let tx_info_str = unsafe { std::ffi::CStr::from_ptr(sign_result.str).to_string_lossy() };
log::debug!("SignChangePubKey result: {}", tx_info_str);
let mut tx_info: serde_json::Value = serde_json::from_str(&tx_info_str)
.map_err(|e| format!("Failed to parse tx_info: {}", e))?;
let message_to_sign = tx_info["MessageToSign"]
.as_str()
.ok_or("MessageToSign not found in tx_info")?
.to_string();
log::debug!("MessageToSign: {}", message_to_sign);
tx_info.as_object_mut().unwrap().remove("MessageToSign");
let evm_signature =
self.sign_message_with_lighter_go_evm(evm_private_key, &message_to_sign)?;
log::debug!("EVM signature: {}", evm_signature);
if let Ok(recovered_addr) =
self.recover_address_from_signature(&message_to_sign, &evm_signature)
{
if let Ok(expected_addr) = self.get_account_l1_address().await {
log::debug!("L1 Address Comparison:");
log::debug!(
" Expected (account {}): {}",
self.account_index,
expected_addr
);
log::debug!(" Recovered from EVM sig: {}", recovered_addr);
if expected_addr.to_lowercase() == recovered_addr.to_lowercase() {
log::debug!(" ✓ Addresses match - signature should be valid");
} else {
log::error!(" ✗ Addresses MISMATCH - signature will fail validation");
log::error!(" This explains the L1 signature failure (code 21504)");
}
} else {
log::warn!("Could not retrieve expected L1 address for comparison");
}
}
tx_info["L1Sig"] = serde_json::Value::String(evm_signature);
let response = self.send_change_api_key_request(&tx_info.to_string()).await;
match response {
Ok(_v) => {
let srv_short = &server_public_key[..8];
let srv_end = &server_public_key[server_public_key.len() - 8..];
let new_short = &new_pubkey.trim_start_matches("0x")[..8];
let new_end_start = new_pubkey.len().saturating_sub(10); let new_end = &new_pubkey[new_end_start..];
log::debug!(
"ChangePubKey succeeded (account={}, index={}). Server public key updated from {}…{} to {}…{}",
self.account_index,
self.api_key_index,
srv_short, srv_end,
new_short, new_end
);
Ok(())
}
Err(e) => {
let srv_short = &server_public_key[..8];
let srv_end = &server_public_key[server_public_key.len() - 8..];
log::error!(
"ChangePubKey failed (account={}, index={}) -> {}. Server key remains {}…{}",
self.account_index,
self.api_key_index,
e,
srv_short,
srv_end
);
Err(e)
}
}
}
#[cfg(not(feature = "lighter-sdk"))]
async fn register_api_key(
&self,
_evm_private_key: &str,
_go_public_key: &str,
_server_public_key: &str,
) -> Result<(), String> {
Err("API key registration requires lighter-sdk feature".to_string())
}
#[cfg(feature = "lighter-sdk")]
fn sign_message_with_lighter_go_evm(
&self,
evm_private_key: &str,
message: &str,
) -> Result<String, String> {
log::debug!("Using lighter-go SignMessageWithEVM for EVM signature");
let private_key_cstr = std::ffi::CString::new(evm_private_key)
.map_err(|e| format!("Failed to create CString for private key: {}", e))?;
let message_cstr = std::ffi::CString::new(message)
.map_err(|e| format!("Failed to create CString for message: {}", e))?;
let sign_result =
unsafe { SignMessageWithEVM(private_key_cstr.as_ptr(), message_cstr.as_ptr()) };
if !sign_result.err.is_null() {
let error_msg = unsafe { std::ffi::CStr::from_ptr(sign_result.err).to_string_lossy() };
return Err(format!("EVM signature failed: {}", error_msg));
}
let signature = unsafe { std::ffi::CStr::from_ptr(sign_result.str).to_string_lossy() };
let signature_with_prefix = if signature.starts_with("0x") {
signature.to_string()
} else {
format!("0x{}", signature)
};
if signature_with_prefix.len() == 132 {
let mut sig_bytes = hex::decode(&signature_with_prefix[2..])
.map_err(|e| format!("Failed to decode signature hex: {}", e))?;
if sig_bytes.len() == 65 {
if sig_bytes[64] == 0 {
log::debug!("Converting v from 0 to 27");
sig_bytes[64] = 27;
} else if sig_bytes[64] == 1 {
log::debug!("Converting v from 1 to 28");
sig_bytes[64] = 28;
}
let corrected_signature = format!("0x{}", hex::encode(sig_bytes));
log::debug!("EVM signature v-corrected: {}", corrected_signature);
if let Ok(recovered_addr) =
self.recover_address_from_signature(message, &corrected_signature)
{
log::debug!("EVM signature recovery check - Address: {}", recovered_addr);
} else {
log::warn!("Failed to recover address from EVM signature for verification");
}
return Ok(corrected_signature);
}
}
Ok(signature_with_prefix)
}
#[cfg(not(feature = "lighter-sdk"))]
fn sign_message_with_lighter_go_evm(
&self,
_evm_private_key: &str,
_message: &str,
) -> Result<String, String> {
Err("EVM signing with lighter-go requires lighter-sdk feature".to_string())
}
async fn get_account_l1_address(&self) -> Result<String, DexError> {
let url = format!(
"{}/api/v1/account?account_index={}",
self.base_url, self.account_index
);
log::debug!("Getting account details from: {}", url);
let response = self
.client
.get(&url)
.header("X-API-KEY", &self.api_key_public)
.send()
.await
.map_err(|e| DexError::Other(format!("Failed to get account details: {}", e)))?;
let status = response.status();
let response_text = response
.text()
.await
.map_err(|e| DexError::Other(format!("Failed to read response: {}", e)))?;
log::debug!(
"Account details response: HTTP {}, Body: {}",
status,
response_text
);
if !status.is_success() {
return Err(DexError::Other(format!(
"HTTP {}: {}",
status, response_text
)));
}
let account_data: serde_json::Value = serde_json::from_str(&response_text)
.map_err(|e| DexError::Other(format!("Failed to parse account response: {}", e)))?;
if let Some(l1_address) = account_data.get("l1Address").and_then(|v| v.as_str()) {
Ok(l1_address.to_string())
} else {
Err(DexError::Other(
"l1Address not found in account response".to_string(),
))
}
}
fn recover_address_from_signature(
&self,
message: &str,
signature: &str,
) -> Result<String, String> {
let signature_hex = signature.strip_prefix("0x").unwrap_or(signature);
let signature_bytes = hex::decode(signature_hex)
.map_err(|e| format!("Failed to decode signature hex: {}", e))?;
if signature_bytes.len() != 65 {
return Err(format!(
"Invalid signature length: {} (expected 65)",
signature_bytes.len()
));
}
let _r = &signature_bytes[0..32];
let _s = &signature_bytes[32..64];
let v = signature_bytes[64];
let recovery_id = match v {
27 => 0,
28 => 1,
0 | 1 => v, _ => return Err(format!("Invalid v value: {}", v)),
};
let prefix = format!("\x19Ethereum Signed Message:\n{}", message.len());
let mut hasher = Keccak256::new();
hasher.update(prefix.as_bytes());
hasher.update(message.as_bytes());
let message_hash = hasher.finalize();
let secp = Secp256k1::new();
let message_obj = Message::from_digest_slice(&message_hash)
.map_err(|e| format!("Invalid message hash: {}", e))?;
let recoverable_sig = secp256k1::ecdsa::RecoverableSignature::from_compact(
&signature_bytes[0..64],
secp256k1::ecdsa::RecoveryId::from_i32(recovery_id as i32)
.map_err(|e| format!("Invalid recovery id: {}", e))?,
)
.map_err(|e| format!("Failed to create recoverable signature: {}", e))?;
let public_key = secp
.recover_ecdsa(&message_obj, &recoverable_sig)
.map_err(|e| format!("Failed to recover public key: {}", e))?;
let public_key_bytes = public_key.serialize_uncompressed();
let mut hasher = Keccak256::new();
hasher.update(&public_key_bytes[1..]); let hash = hasher.finalize();
let address = format!("0x{}", hex::encode(&hash[12..]));
Ok(address)
}
fn _unused_sign_with_evm_key(
&self,
evm_private_key: &str,
message: &str,
) -> Result<String, String> {
let private_key_bytes = if evm_private_key.contains("=")
|| evm_private_key.contains("+")
|| evm_private_key.contains("/")
{
use base64::Engine;
base64::engine::general_purpose::STANDARD
.decode(evm_private_key)
.map_err(|e| format!("Failed to decode private key base64: {}", e))?
} else {
let private_key_hex = if evm_private_key.starts_with("0x") {
&evm_private_key[2..]
} else {
evm_private_key
};
hex::decode(private_key_hex)
.map_err(|e| format!("Failed to decode private key hex: {}", e))?
};
if private_key_bytes.len() != 32 {
return Err("Private key must be 32 bytes".to_string());
}
let secret_key = SecretKey::from_slice(&private_key_bytes)
.map_err(|e| format!("Failed to create secret key: {}", e))?;
let prefix = format!("\x19Ethereum Signed Message:\n{}", message.len());
let full_message = format!("{}{}", prefix, message);
let mut hasher = Keccak256::new();
hasher.update(full_message.as_bytes());
let message_hash = hasher.finalize();
let secp = Secp256k1::new();
let message_obj = Message::from_digest_slice(&message_hash)
.map_err(|e| format!("Failed to create message: {}", e))?;
let signature = secp.sign_ecdsa_recoverable(&message_obj, &secret_key);
let (recovery_id, compact_sig) = signature.serialize_compact();
let mut signature_bytes = [0u8; 65];
signature_bytes[0..64].copy_from_slice(&compact_sig);
signature_bytes[64] = (recovery_id.to_i32() + 27) as u8;
Ok(hex::encode(signature_bytes))
}
async fn send_change_api_key_request(
&self,
tx_info: &str,
) -> Result<serde_json::Value, String> {
let base_url = self.base_url.trim_end_matches('/');
let url = format!("{}/api/v1/sendTx", base_url);
let form_data = [
("tx_type", "8"), ("tx_info", tx_info),
];
log::debug!("Sending change API key request to: {}", url);
log::debug!("Form data: {:?}", form_data);
track_api_call("POST /api/v1/sendTx (change_api_key)", "POST");
let response = self
.client
.post(&url)
.form(&form_data)
.send()
.await
.map_err(|e| format!("Failed to send request: {}", e))?;
let status = response.status();
let response_text = response
.text()
.await
.map_err(|e| format!("Failed to read response: {}", e))?;
log::debug!(
"Change API key response: HTTP {}, Body: {}",
status,
response_text
);
if !status.is_success() {
return Err(format!("HTTP {}: {}", status, response_text));
}
serde_json::from_str(&response_text)
.map_err(|e| format!("Failed to parse response JSON: {}", e))
}
}
#[async_trait]
impl DexConnector for LighterConnector {
async fn start(&self) -> Result<(), DexError> {
self.is_running.store(true, Ordering::SeqCst);
log::debug!(
"Lighter connector started with WebSocket: {}",
self.websocket_url
);
#[cfg(feature = "lighter-sdk")]
{
match self.create_go_client().await {
Ok(()) => {}
Err(DexError::ApiKeyRegistrationRequired) => {
#[cfg(feature = "lighter-sdk")]
if let Some(evm_key) = &self.evm_wallet_private_key {
let go_pubkey_result = unsafe {
GetClientPubKey(
self.api_key_index as c_int,
self.account_index as c_longlong,
)
};
if !go_pubkey_result.is_null() {
let pubkey_cstr = unsafe { CStr::from_ptr(go_pubkey_result) };
let go_key = pubkey_cstr.to_string_lossy().to_string();
unsafe { libc::free(go_pubkey_result as *mut libc::c_void) };
log::debug!("API key registration required. Attempting to register...");
let server_pubkey =
self.get_server_public_key().await.map_err(|e| {
DexError::Other(format!(
"Failed to get server public key: {}",
e
))
})?;
self.register_api_key(evm_key, &go_key, &server_pubkey)
.await
.map_err(|e| {
DexError::Other(format!("API key registration failed: {}", e))
})?;
self.create_go_client().await?;
} else {
log::error!("Failed to get Go-derived public key for registration");
return Err(DexError::Other(
"Cannot get Go-derived public key".to_string(),
));
}
} else {
return Err(DexError::ApiKeyRegistrationRequired);
}
#[cfg(not(feature = "lighter-sdk"))]
return Err(DexError::ApiKeyRegistrationRequired);
}
Err(e) => return Err(e),
}
}
self.start_websocket().await?;
self.start_auto_cleanup(6);
log::info!("🗑️ [AUTO_CLEANUP] Started background cleanup task (every 6 hours)");
Ok(())
}
async fn stop(&self) -> Result<(), DexError> {
self.is_running.store(false, Ordering::SeqCst);
if let Some(handle) = self.cleanup_handle.lock().await.take() {
handle.abort();
self.cleanup_started.store(false, Ordering::SeqCst);
log::debug!("🛑 [AUTO_CLEANUP] task forcibly aborted on stop");
}
log::debug!("Lighter connector stopped");
Ok(())
}
async fn restart(&self, _max_retries: i32) -> Result<(), DexError> {
self.stop().await?;
sleep(Duration::from_secs(1)).await;
self.start().await
}
async fn set_leverage(&self, _symbol: &str, _leverage: u32) -> Result<(), DexError> {
log::warn!("Leverage setting not implemented for Lighter");
Ok(())
}
async fn get_ticker(
&self,
symbol: &str,
test_price: Option<Decimal>,
) -> Result<TickerResponse, DexError> {
if let Some(price) = test_price {
let min_tick = Self::calculate_min_tick(price, 3, false); return Ok(TickerResponse {
symbol: symbol.to_string(),
price,
min_tick: Some(min_tick),
min_order: None,
volume: Some(Decimal::ZERO),
num_trades: None,
open_interest: None,
funding_rate: None,
oracle_price: None,
});
}
let market_info = self.resolve_market_info(symbol).await?;
let canonical_symbol = market_info.canonical_symbol.clone();
let stats_data = self.get_exchange_stats().await.ok();
let funding_data = self.get_funding_rates().await.ok();
if let Some((ws_price, price_timestamp)) = *self.current_price.read().await {
let current_time = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs();
let price_age = current_time.saturating_sub(price_timestamp);
if price_age > 30 {
log::warn!(
"WebSocket price is stale ({}s old), falling back to REST API",
price_age
);
} else {
let min_tick = Self::calculate_min_tick(ws_price, 3, false);
let (volume, num_trades) = if let Some(stats) = &stats_data {
if let Some(market_stats) = stats
.order_book_stats
.iter()
.find(|s| normalize_symbol(&s.symbol) == canonical_symbol)
{
(
Some(
Decimal::from_f64_retain(market_stats.daily_base_token_volume)
.unwrap_or(Decimal::ZERO),
),
Some(market_stats.daily_trades_count as u64),
)
} else {
(Some(Decimal::ZERO), None)
}
} else {
(Some(Decimal::ZERO), None)
};
let funding_rate = if let Some(funding) = &funding_data {
funding
.funding_rates
.iter()
.find(|f| normalize_symbol(&f.symbol) == canonical_symbol)
.and_then(|f| Decimal::from_f64_retain(f.rate))
} else {
None
};
log::trace!(
"Using WebSocket price with API stats: price={}, volume={:?}, trades={:?}",
ws_price,
volume,
num_trades
);
return Ok(TickerResponse {
symbol: symbol.to_string(),
price: ws_price,
min_tick: Some(min_tick),
min_order: None,
volume,
num_trades,
open_interest: None,
funding_rate,
oracle_price: None,
});
}
}
log::warn!("WebSocket data not available, falling back to REST API");
let market_id = market_info.market_id;
let endpoint = format!("/api/v1/recentTrades?market_id={}&limit=100", market_id);
let url = format!("{}{}", self.base_url, endpoint);
let response = self
.client
.get(&url)
.header("X-API-KEY", &self.api_key_public)
.send()
.await
.map_err(|e| DexError::Other(format!("Request failed: {}", e)))?;
let status = response.status();
let response_text = response
.text()
.await
.map_err(|e| DexError::Other(format!("Failed to read response: {}", e)))?;
log::debug!(
"Trades API response (status: {}): {}",
status,
response_text
);
if !status.is_success() {
return Err(DexError::Other(format!(
"HTTP {}: {}",
status, response_text
)));
}
let trades_response: LighterTradesResponse = serde_json::from_str(&response_text)
.map_err(|e| DexError::Other(format!("Failed to parse response: {}", e)))?;
let price = if let Some(trade) = trades_response.trades.first() {
string_to_decimal(Some(trade.price.clone()))?
} else {
Decimal::new(50000, 0)
};
let min_tick = Self::calculate_min_tick(price, 3, false);
let funding_rate = if let Some(funding) = &funding_data {
funding
.funding_rates
.iter()
.find(|f| normalize_symbol(&f.symbol) == canonical_symbol)
.and_then(|f| Decimal::from_f64_retain(f.rate))
} else {
None
};
let (volume, num_trades) = if let Some(stats) = &stats_data {
if let Some(market_stats) = stats
.order_book_stats
.iter()
.find(|s| normalize_symbol(&s.symbol) == canonical_symbol)
{
(
Some(
Decimal::from_f64_retain(market_stats.daily_base_token_volume)
.unwrap_or(Decimal::ZERO),
),
Some(market_stats.daily_trades_count as u64),
)
} else {
let volume = trades_response
.trades
.iter()
.map(|trade| string_to_decimal(Some(trade.size.clone())))
.collect::<Result<Vec<_>, _>>()?
.iter()
.sum();
(Some(volume), Some(trades_response.trades.len() as u64))
}
} else {
let volume = trades_response
.trades
.iter()
.map(|trade| string_to_decimal(Some(trade.size.clone())))
.collect::<Result<Vec<_>, _>>()?
.iter()
.sum();
(Some(volume), Some(trades_response.trades.len() as u64))
};
Ok(TickerResponse {
symbol: symbol.to_string(),
price,
min_tick: Some(min_tick),
min_order: None,
volume,
num_trades,
open_interest: None,
funding_rate,
oracle_price: None,
})
}
async fn get_filled_orders(&self, symbol: &str) -> Result<FilledOrdersResponse, DexError> {
let orders = self.filled_orders.read().await;
let normalized = normalize_symbol(symbol);
let symbol_orders = orders
.get(symbol)
.or_else(|| orders.get(&normalized))
.cloned()
.unwrap_or_default();
Ok(FilledOrdersResponse {
orders: symbol_orders,
})
}
async fn get_canceled_orders(&self, symbol: &str) -> Result<CanceledOrdersResponse, DexError> {
let orders = self.canceled_orders.read().await;
let symbol_orders = orders.get(symbol).cloned().unwrap_or_default();
Ok(CanceledOrdersResponse {
orders: symbol_orders,
})
}
async fn get_balance(&self, symbol: Option<&str>) -> Result<BalanceResponse, DexError> {
let endpoint = format!("/api/v1/account?by=index&value={}", self.account_index);
let url = format!("{}{}", self.base_url, endpoint);
log::info!(
"get_balance called for symbol: {:?}, requesting URL: {}",
symbol,
url
);
track_api_call(&endpoint, "GET");
let response = self
.client
.get(&url)
.header("X-API-KEY", &self.api_key_public)
.send()
.await
.map_err(|e| DexError::Other(format!("Request failed: {}", e)))?;
let status = response.status();
let response_text = response
.text()
.await
.map_err(|e| DexError::Other(format!("Failed to read response: {}", e)))?;
log::info!(
"Account API response (status: {}): {}",
status,
response_text
);
if !status.is_success() {
return Err(DexError::Other(format!(
"HTTP {}: {}",
status, response_text
)));
}
let account_response: LighterAccountResponse = serde_json::from_str(&response_text)
.map_err(|e| DexError::Other(format!("Failed to parse response: {}", e)))?;
if account_response.accounts.is_empty() {
return Err(DexError::Other("No account found".to_string()));
}
let account = &account_response.accounts[0];
log::info!("Account balance info:");
log::info!(" - Account Index: {}", account.account_index);
log::info!(" - Available Balance: {} USD", account.available_balance);
log::info!(" - Collateral: {} USD", account.collateral);
log::info!(" - Total Asset Value: {} USD", account.total_asset_value);
log::info!(" - Positions count: {}", account.positions.len());
for (i, position) in account.positions.iter().enumerate() {
log::info!(
" Position [{}]: market_id={}, symbol={}, position={}, sign={}",
i,
position.market_id,
position.symbol,
position.position,
position.sign
);
}
if let Some(token_symbol) = symbol {
log::trace!("Looking for position with symbol: {}", token_symbol);
for position in &account.positions {
if position.symbol == token_symbol {
log::trace!(
"✓ Found position for {}: {} (sign: {})",
token_symbol,
position.position,
position.sign
);
let position_decimal = string_to_decimal(Some(position.position.clone()))?;
let entry_price = string_to_decimal(Some(position.avg_entry_price.clone()))?;
return Ok(BalanceResponse {
equity: position_decimal,
balance: position_decimal,
position_entry_price: Some(entry_price),
position_sign: Some(position.sign.into()),
});
}
}
log::trace!("✗ No position found for {}, returning zero", token_symbol);
return Ok(BalanceResponse {
equity: rust_decimal::Decimal::ZERO,
balance: rust_decimal::Decimal::ZERO,
position_entry_price: None,
position_sign: None,
});
}
log::trace!("No symbol specified, returning account-level USD balances");
let total_asset_value = string_to_decimal(Some(account.total_asset_value.clone()))?;
let available_balance = string_to_decimal(Some(account.available_balance.clone()))?;
log::info!(
"Account balances: total_asset_value={}, available_balance={}",
total_asset_value,
available_balance
);
Ok(BalanceResponse {
equity: total_asset_value, balance: available_balance, position_entry_price: None, position_sign: None,
})
}
async fn get_combined_balance(&self) -> Result<CombinedBalanceResponse, DexError> {
let endpoint = format!("/api/v1/account?by=index&value={}", self.account_index);
let url = format!("{}{}", self.base_url, endpoint);
log::info!("get_combined_balance called, requesting URL: {}", url);
let response = self
.client
.get(&url)
.header("X-API-KEY", &self.api_key_public)
.send()
.await
.map_err(|e| DexError::Other(format!("Request failed: {}", e)))?;
let status = response.status();
let response_text = response
.text()
.await
.map_err(|e| DexError::Other(format!("Failed to read response: {}", e)))?;
if !status.is_success() {
return Err(DexError::Other(format!(
"HTTP {}: {}",
status, response_text
)));
}
let account_response: LighterAccountResponse = serde_json::from_str(&response_text)
.map_err(|e| DexError::Other(format!("Failed to parse response: {}", e)))?;
if account_response.accounts.is_empty() {
return Err(DexError::Other("No account found".to_string()));
}
let account = &account_response.accounts[0];
let usd_balance = string_to_decimal(Some(account.available_balance.clone()))?;
let mut token_balances = std::collections::HashMap::new();
for position in &account.positions {
let position_decimal = string_to_decimal(Some(position.position.clone()))?;
let entry_price = string_to_decimal(Some(position.avg_entry_price.clone()))?;
token_balances.insert(
position.symbol.clone(),
BalanceResponse {
equity: position_decimal,
balance: position_decimal,
position_entry_price: Some(entry_price),
position_sign: Some(position.sign.into()),
},
);
}
log::debug!(
"Combined balance: USD={}, tokens={} positions",
usd_balance,
token_balances.len()
);
Ok(CombinedBalanceResponse {
usd_balance,
token_balances,
})
}
async fn get_open_orders(&self, symbol: &str) -> Result<OpenOrdersResponse, DexError> {
log::debug!(
"[WS_ORDER_TRACKING] get_open_orders called for symbol: {} (WebSocket-only)",
symbol
);
let orders_guard = self.cached_open_orders.read().await;
let orders = orders_guard.get(symbol).cloned().unwrap_or_default();
log::debug!(
"[WS_ORDER_TRACKING] Returning {} orders for {} from WebSocket tracking",
orders.len(),
symbol
);
Ok(OpenOrdersResponse { orders })
}
async fn get_last_trades(&self, symbol: &str) -> Result<LastTradesResponse, DexError> {
let market_id = self.resolve_market_info(symbol).await?.market_id;
let endpoint = format!("/api/v1/recentTrades?market_id={}&limit=10", market_id);
let url = format!("{}{}", self.base_url, endpoint);
let response = self
.client
.get(&url)
.header("X-API-KEY", &self.api_key_public)
.send()
.await
.map_err(|e| DexError::Other(format!("Request failed: {}", e)))?;
let status = response.status();
let response_text = response
.text()
.await
.map_err(|e| DexError::Other(format!("Failed to read response: {}", e)))?;
log::debug!(
"Last trades API response (status: {}): {}",
status,
response_text
);
if !status.is_success() {
return Err(DexError::Other(format!(
"HTTP {}: {}",
status, response_text
)));
}
let trades_response: LighterTradesResponse = serde_json::from_str(&response_text)
.map_err(|e| DexError::Other(format!("Failed to parse response: {}", e)))?;
let trades = trades_response
.trades
.into_iter()
.map(|t| LastTrade {
price: string_to_decimal(Some(t.price)).unwrap_or_default(),
})
.collect();
Ok(LastTradesResponse { trades })
}
async fn clear_filled_order(&self, symbol: &str, trade_id: &str) -> Result<(), DexError> {
let mut filled_orders = self.filled_orders.write().await;
if let Some(orders) = filled_orders.get_mut(symbol) {
let initial_len = orders.len();
orders.retain(|order| order.trade_id != trade_id);
if orders.len() < initial_len {
log::debug!(
"🗑️ [CLEAR_FILL] Removed trade_id {} for {}",
trade_id,
symbol
);
Ok(())
} else {
Err(DexError::Other(format!(
"Trade ID {} not found for symbol {}",
trade_id, symbol
)))
}
} else {
Err(DexError::Other(format!(
"No filled orders found for symbol {}",
symbol
)))
}
}
async fn clear_all_filled_orders(&self) -> Result<(), DexError> {
let mut filled_orders = self.filled_orders.write().await;
let total_cleared = filled_orders.values().map(|v| v.len()).sum::<usize>();
filled_orders.clear();
log::info!(
"🗑️ [CLEAR_ALL_FILLS] Cleared {} filled orders across all symbols",
total_cleared
);
Ok(())
}
async fn clear_canceled_order(&self, _symbol: &str, _order_id: &str) -> Result<(), DexError> {
Err(DexError::Other(
"clear_canceled_order not supported for Lighter - canceled orders are streamed via WebSocket only".to_string()
))
}
async fn clear_all_canceled_orders(&self) -> Result<(), DexError> {
Err(DexError::Other(
"clear_all_canceled_orders not supported for Lighter - canceled orders are streamed via WebSocket only".to_string()
))
}
async fn create_order(
&self,
symbol: &str,
size: Decimal,
side: OrderSide,
price: Option<Decimal>,
_spread: Option<i64>,
expiry_secs: Option<u64>,
) -> Result<CreateOrderResponse, DexError> {
let market_info = self.resolve_market_info(symbol).await?;
let market_id = market_info.market_id;
let side_value = match side {
OrderSide::Long => 0,
OrderSide::Short => 1,
};
let default_tif = 0;
let base_amount = (size * Decimal::new(100_000, 0))
.to_u64()
.ok_or_else(|| DexError::Other("Invalid size amount".to_string()))?;
let (price_value, order_type, tif) = if let Some(p) = price {
let (final_price, order_tif) = if let Some(spread_ticks) = _spread {
if spread_ticks < 0 {
let tif_value = match spread_ticks {
-1 => ORDER_TYPE_IOC, -2 => ORDER_TYPE_FOK, _ => {
log::warn!("Invalid TIF spread value: {}, using GTC", spread_ticks);
default_tif
}
};
(p, tif_value) } else {
let tick_size = Decimal::new(1, 1); let spread_amount = Decimal::from(spread_ticks) * tick_size;
(p + spread_amount, default_tif)
}
} else {
(p, default_tif)
};
let price_val = (final_price * Decimal::new(10, 0))
.to_u32()
.ok_or_else(|| DexError::Other("Invalid price".to_string()))?
as u64;
let tif_name = match order_tif {
0 => "GTC",
1 => "IOC",
2 => "FOK",
_ => "UNKNOWN",
};
log::debug!("Creating limit order: side={}, original_price={}, spread_param={:?}, final_price={}, TIF={} ({}), scaled_price={}, size={}, scaled_base_amount={}",
side_value, p, _spread, final_price, order_tif, tif_name, price_val, size, base_amount);
(price_val, ORDER_TYPE_LIMIT, order_tif)
} else {
let ticker = self.get_ticker(symbol, None).await?;
let current_price = ticker.price;
let protection_price = if side_value == 1 {
current_price * Decimal::new(800, 3) } else {
current_price * Decimal::new(1200, 3) };
let price_val = (protection_price * Decimal::new(10, 0))
.to_u32()
.ok_or_else(|| DexError::Other("Invalid protection price".to_string()))?
as u64;
log::debug!(
"Market order: current_price={}, protection_price={}, side={}",
current_price,
protection_price,
side_value
);
(price_val, ORDER_TYPE_IOC, 0u32) };
let result = self
.create_order_native_with_type(
market_id,
side_value,
tif,
base_amount,
price_value,
None,
order_type,
false,
expiry_secs,
)
.await;
if let Ok(ref response) = result {
let actual_price = price.unwrap_or(Decimal::ZERO);
self.update_order_tracking_after_create(
symbol,
&response.order_id,
side,
size,
actual_price,
)
.await;
}
result
}
async fn create_advanced_trigger_order(
&self,
symbol: &str,
size: Decimal,
side: OrderSide,
trigger_px: Decimal,
limit_px: Option<Decimal>,
order_style: TriggerOrderStyle,
slippage_bps: Option<u32>,
tpsl: TpSl,
reduce_only: bool,
expiry_secs: Option<u64>,
) -> Result<CreateOrderResponse, DexError> {
log::info!(
"🎯 [ADVANCED_TRIGGER_ORDER] Creating {} order for {}: style={:?}, trigger={}, limit={:?}, slippage_bps={:?}",
match tpsl { TpSl::Sl => "stop loss", TpSl::Tp => "take profit" },
symbol,
order_style,
trigger_px,
limit_px,
slippage_bps
);
let market_info = self.resolve_market_info(symbol).await?;
let market_id = market_info.market_id;
let side_value = if is_buy_for_tpsl(side) {
SIDE_BUY
} else {
SIDE_SELL
};
let (is_market, final_limit_price, order_type) = match order_style {
TriggerOrderStyle::Market => {
let order_type = match tpsl {
TpSl::Sl => 2, TpSl::Tp => 4, };
(true, trigger_px, order_type)
}
TriggerOrderStyle::MarketWithSlippageControl => {
if let Some(slippage) = slippage_bps {
let slippage_factor = Decimal::new(slippage as i64, 4);
let adjusted_price = match (side, tpsl) {
(OrderSide::Long, TpSl::Sl) => {
trigger_px * (Decimal::ONE - slippage_factor)
}
(OrderSide::Short, TpSl::Sl) => {
trigger_px * (Decimal::ONE + slippage_factor)
}
(OrderSide::Long, TpSl::Tp) => {
trigger_px * (Decimal::ONE - slippage_factor)
}
(OrderSide::Short, TpSl::Tp) => {
trigger_px * (Decimal::ONE + slippage_factor)
}
};
let order_type = match tpsl {
TpSl::Sl => 3, TpSl::Tp => 5, };
(false, adjusted_price, order_type)
} else {
let order_type = match tpsl {
TpSl::Sl => 2,
TpSl::Tp => 4,
};
(true, trigger_px, order_type)
}
}
TriggerOrderStyle::Limit => {
let limit_price = limit_px.ok_or_else(|| {
DexError::Other("limit_px required for Limit order style".into())
})?;
match (side, tpsl) {
(OrderSide::Long, TpSl::Sl) => {
if limit_price < trigger_px {
return Err(DexError::Other(
"For Buy Stop Loss, limit_px must be >= trigger_px".into(),
));
}
}
(OrderSide::Short, TpSl::Sl) => {
if limit_price > trigger_px {
return Err(DexError::Other(
"For Sell Stop Loss, limit_px must be <= trigger_px".into(),
));
}
}
(OrderSide::Long, TpSl::Tp) => {
if limit_price > trigger_px {
return Err(DexError::Other(
"For Buy Take Profit, limit_px must be <= trigger_px".into(),
));
}
}
(OrderSide::Short, TpSl::Tp) => {
if limit_price < trigger_px {
return Err(DexError::Other(
"For Sell Take Profit, limit_px must be >= trigger_px".into(),
));
}
}
}
let order_type = match tpsl {
TpSl::Sl => 3, TpSl::Tp => 5, };
(false, limit_price, order_type)
}
};
let base_amount = (size * Decimal::from(BASE_AMOUNT_SCALE))
.to_u64()
.ok_or_else(|| DexError::Other("base_amount overflow or invalid value".into()))?;
let trigger_price_native = (trigger_px * Decimal::from(PRICE_SCALE_FACTOR))
.to_u64()
.ok_or_else(|| {
DexError::Other("trigger_price_native overflow or invalid value".into())
})?;
let execution_price_native = if is_market {
0 } else {
(final_limit_price * Decimal::from(PRICE_SCALE_FACTOR))
.to_u64()
.ok_or_else(|| {
DexError::Other("execution_price_native overflow or invalid value".into())
})?
};
let time_in_force = if is_market {
TIME_IN_FORCE_IOC
} else {
TIME_IN_FORCE_GTC
};
log::debug!(
"Creating trigger order: market_id={}, side={}, base_amount={}, price={}, trigger_price={}, order_type={}",
market_id, side_value, base_amount, execution_price_native, trigger_price_native, order_type
);
self.create_order_native_with_trigger(
market_id,
side_value,
time_in_force,
base_amount,
execution_price_native,
trigger_price_native,
None,
order_type,
reduce_only,
expiry_secs,
)
.await
}
async fn cancel_order(&self, symbol: &str, order_id: &str) -> Result<(), DexError> {
log::warn!("cancel_order called for symbol: {}, order_id: {} - Individual order cancellation not implemented for Lighter, use cancel_all_orders instead", symbol, order_id);
Err(DexError::Other(
"Order cancellation not implemented yet".to_string(),
))
}
async fn cancel_all_orders(&self, _symbol: Option<String>) -> Result<(), DexError> {
log::info!(
"Starting cancel_all_orders for symbol: {:?}, API key index: {}, account index: {}",
_symbol,
self.api_key_index,
self.account_index
);
#[cfg(feature = "lighter-sdk")]
{
log::debug!("Getting nonce for cancel_all_orders");
let nonce = self.get_nonce().await?;
log::debug!("Retrieved nonce: {}", nonce);
let time_in_force = 0; let time = 0;
log::debug!(
"Signing cancel_all_orders transaction with time_in_force: {}, time: {}, nonce: {}",
time_in_force,
time,
nonce
);
unsafe {
let result = SignCancelAllOrders(time_in_force, time, nonce as i64);
if !result.err.is_null() {
let error_cstr = CStr::from_ptr(result.err);
let error_msg = error_cstr.to_string_lossy().to_string();
libc::free(result.err as *mut libc::c_void);
return Err(DexError::Other(format!(
"Cancel all orders failed: {}",
error_msg
)));
}
let result_cstr = CStr::from_ptr(result.str);
let tx_json = result_cstr.to_string_lossy().to_string();
libc::free(result.str as *mut libc::c_void);
let _timestamp = chrono::Utc::now().timestamp_millis() as u64;
let form_data = format!(
"tx_type=16&tx_info={}&price_protection=false",
urlencoding::encode(&tx_json)
);
log::debug!(
"Submitting cancel_all_orders transaction to API: {}/api/v1/sendTx",
self.base_url
);
track_api_call("POST /api/v1/sendTx (cancel_all_orders)", "POST");
let response = self
.client
.post(&format!("{}/api/v1/sendTx", self.base_url))
.header("Content-Type", "application/x-www-form-urlencoded")
.body(form_data)
.send()
.await
.map_err(|e| DexError::Other(format!("Network error: {}", e)))?;
let status = response.status();
let body = response.text().await.unwrap_or_default();
if !status.is_success() {
log::error!("Cancel all orders failed: HTTP {}, Body: {}", status, body);
return Err(DexError::Other(format!(
"Cancel all orders failed: HTTP {}, {}",
status, body
)));
}
log::info!(
"Cancel all orders submitted successfully, Response: {}",
body
);
tokio::time::sleep(Duration::from_millis(1000)).await;
log::info!(
"Cancel all orders submitted successfully - trusting server response without verification"
);
}
}
#[cfg(not(feature = "lighter-sdk"))]
{
log::error!("Cancel all orders called but lighter-sdk feature is not enabled. Symbol: {:?}, API key index: {}, Account index: {}",
_symbol, self.api_key_index, self.account_index);
}
{
let orders_guard = self.cached_open_orders.read().await;
if let Some(ref symbol) = _symbol {
if let Some(orders) = orders_guard.get(symbol) {
for order in orders {
log::debug!(
"[WS_ORDER_TRACKING] Marking order {} as cancelled for symbol {}",
order.order_id,
symbol
);
}
}
} else {
for (symbol, orders) in orders_guard.iter() {
for order in orders {
log::debug!(
"[WS_ORDER_TRACKING] Marking order {} as cancelled for symbol {}",
order.order_id,
symbol
);
}
}
}
}
{
let mut orders_guard = self.cached_open_orders.write().await;
if let Some(symbol) = _symbol {
orders_guard.remove(&symbol);
log::debug!("[WS_ORDER_TRACKING] Cleared orders for symbol: {}", symbol);
} else {
orders_guard.clear();
log::debug!("[WS_ORDER_TRACKING] Cleared all orders from tracking");
}
}
Ok(())
}
async fn cancel_orders(
&self,
symbol: Option<String>,
order_ids: Vec<String>,
) -> Result<(), DexError> {
log::warn!("cancel_orders called for symbol: {:?}, order_ids: {:?} - Individual order cancellation not implemented for Lighter, use cancel_all_orders instead", symbol, order_ids);
Err(DexError::Other(
"Individual order cancellation not implemented for Lighter. Use cancel_all_orders instead.".to_string(),
))
}
async fn close_all_positions(&self, _symbol: Option<String>) -> Result<(), DexError> {
let endpoint = format!("/api/v1/account?by=index&value={}", self.account_index);
let url = format!("{}{}", self.base_url, endpoint);
let response = self
.client
.get(&url)
.header("X-API-KEY", &self.api_key_public)
.send()
.await
.map_err(|e| DexError::Other(format!("Request failed: {}", e)))?;
let status = response.status();
let response_text = response
.text()
.await
.map_err(|e| DexError::Other(format!("Failed to read response: {}", e)))?;
if !status.is_success() {
return Err(DexError::Other(format!(
"HTTP {}: {}",
status, response_text
)));
}
let account_response: LighterAccountResponse = serde_json::from_str(&response_text)
.map_err(|e| DexError::Other(format!("Failed to parse response: {}", e)))?;
if account_response.accounts.is_empty() {
return Err(DexError::Other("No account found".to_string()));
}
let account = &account_response.accounts[0];
let mut has_positions = false;
for position in &account.positions {
if let Ok(pos_size) = position.position.parse::<f64>() {
if pos_size.abs() > 0.0 {
has_positions = true;
log::info!(
"Found open position: market_id={}, symbol={}, size={}",
position.market_id,
position.symbol,
position.position
);
}
}
}
if !has_positions {
log::info!("No open positions found (threshold: > 0.0), nothing to close");
for position in &account.positions {
if let Ok(pos_size) = position.position.parse::<f64>() {
if pos_size.abs() > 0.0 {
log::debug!("Small position below threshold: market_id={}, symbol={}, size={} (abs: {})",
position.market_id, position.symbol, position.position, pos_size.abs());
}
}
}
return Ok(());
}
for position in &account.positions {
if let Ok(pos_size) = position.position.parse::<f64>() {
if pos_size.abs() > 0.0 {
log::info!(
"Closing position: market_id={}, symbol={}, size={}, sign={}",
position.market_id,
position.symbol,
position.position,
position.sign
);
let order_side = if position.sign > 0 {
1 } else {
0 };
let market_id = position.market_id;
let pos_decimal =
rust_decimal::Decimal::from_str(&position.position.replace('-', ""))
.unwrap_or_else(|_| {
let pos_str = format!("{:.8}", pos_size.abs());
rust_decimal::Decimal::from_str(&pos_str)
.unwrap_or(rust_decimal::Decimal::ZERO)
});
let mut base_amount = (pos_decimal * rust_decimal::Decimal::new(100000, 0))
.to_u64()
.unwrap_or((pos_size.abs() * 100000.0) as u64);
if base_amount == 0 && pos_size.abs() > 0.0 {
base_amount = 1;
log::debug!(
"Position too small for conversion, using minimum base_amount=1"
);
}
log::debug!(
"Converting position {} to base_amount: {} (original: {}, decimal: {})",
position.position,
base_amount,
pos_size,
pos_decimal
);
log::info!(
"Placing reduce-only market order to close position: market_id={}, side={}, size={}",
market_id, order_side, pos_decimal
);
let ticker_price = match self.get_ticker(&position.symbol, None).await {
Ok(ticker) => ticker.price,
Err(e) => {
log::warn!(
"Failed to fetch ticker for {} while closing position: {}. Using fallback price",
position.symbol,
e
);
rust_decimal::Decimal::new(50000, 0)
}
};
let protection_price = if order_side == 1 {
ticker_price * rust_decimal::Decimal::new(700, 3) } else {
ticker_price * rust_decimal::Decimal::new(1300, 3) };
let current_price = (protection_price * rust_decimal::Decimal::new(10, 0))
.to_u64()
.unwrap_or(0);
match self
.create_order_native_with_type(
market_id as u32,
order_side as u32,
0, base_amount,
current_price,
None,
1, true, None, )
.await
{
Ok(response) => {
log::info!(
"Successfully submitted reduce-only close order for {} position in market {}: Order ID {}",
position.symbol,
market_id,
response.order_id
);
}
Err(e) => {
log::error!("Failed to close position in market {}: {}", market_id, e);
return Err(e);
}
}
}
}
}
log::info!("All position close orders submitted successfully");
Ok(())
}
async fn clear_last_trades(&self, _symbol: &str) -> Result<(), DexError> {
Ok(())
}
async fn is_upcoming_maintenance(&self, hours_ahead: i64) -> bool {
let info = self.maintenance.read().await;
if let Some(start) = info.next_start {
let now = Utc::now();
if now < start && (start - now) <= ChronoDuration::hours(hours_ahead) {
return true;
}
}
false
}
async fn sign_evm_65b(&self, message: &str) -> Result<String, DexError> {
use ethers::signers::{LocalWallet, Signer};
use std::str::FromStr;
let private_key = self
.evm_wallet_private_key
.as_ref()
.ok_or_else(|| DexError::Other("EVM wallet private key not set".to_string()))?;
let cleaned_key = private_key.strip_prefix("0x").unwrap_or(private_key);
let wallet = LocalWallet::from_str(cleaned_key)
.map_err(|e| DexError::Other(format!("Invalid private key: {}", e)))?;
let signature = wallet
.sign_message(message.as_bytes())
.await
.map_err(|e| DexError::Other(format!("Signing failed: {}", e)))?;
Ok(format!("0x{}", signature))
}
async fn sign_evm_65b_with_eip191(&self, message: &str) -> Result<String, DexError> {
let prefixed = format!("\x19Ethereum Signed Message:\n{}{}", message.len(), message);
self.sign_evm_65b(&prefixed).await
}
}
impl LighterConnector {
async fn update_order_tracking_after_create(
&self,
symbol: &str,
order_id: &str,
side: OrderSide,
size: Decimal,
price: Decimal,
) {
let mut orders_guard = self.cached_open_orders.write().await;
let orders = orders_guard
.entry(symbol.to_string())
.or_insert_with(Vec::new);
let new_order = OpenOrder {
order_id: order_id.to_string(),
symbol: symbol.to_string(),
side,
size,
price,
status: "open".to_string(),
};
orders.push(new_order);
log::debug!(
"[WS_ORDER_TRACKING] Added order {} to tracking for {} (total: {} orders)",
order_id,
symbol,
orders.len()
);
}
#[allow(dead_code)]
async fn update_order_tracking_after_cancel(&self, symbol: &str, order_id: &str) {
let mut orders_guard = self.cached_open_orders.write().await;
if let Some(orders) = orders_guard.get_mut(symbol) {
orders.retain(|order| order.order_id != order_id);
log::debug!(
"[WS_ORDER_TRACKING] Removed order {} from tracking for {} (remaining: {} orders)",
order_id,
symbol,
orders.len()
);
}
}
}
impl LighterConnector {
async fn get_exchange_stats(&self) -> Result<LighterExchangeStats, DexError> {
let url = format!("{}/api/v1/exchangeStats", self.base_url);
let response = self
.client
.get(&url)
.header("X-API-KEY", &self.api_key_public)
.send()
.await
.map_err(|e| DexError::Other(format!("Failed to get exchange stats: {}", e)))?;
let status = response.status();
let response_text = response.text().await.map_err(|e| {
DexError::Other(format!("Failed to read exchange stats response: {}", e))
})?;
if !status.is_success() {
return Err(DexError::Other(format!(
"Exchange stats HTTP {}: {}",
status, response_text
)));
}
log::trace!("Exchange stats response: {}", response_text);
serde_json::from_str(&response_text)
.map_err(|e| DexError::Other(format!("Failed to parse exchange stats: {}", e)))
}
async fn get_funding_rates(&self) -> Result<LighterFundingRates, DexError> {
let url = format!("{}/api/v1/funding-rates", self.base_url);
let response = self
.client
.get(&url)
.header("X-API-KEY", &self.api_key_public)
.send()
.await
.map_err(|e| DexError::Other(format!("Failed to get funding rates: {}", e)))?;
let status = response.status();
let response_text = response.text().await.map_err(|e| {
DexError::Other(format!("Failed to read funding rates response: {}", e))
})?;
if !status.is_success() {
return Err(DexError::Other(format!(
"Funding rates HTTP {}: {}",
status, response_text
)));
}
log::trace!("Funding rates response: {}", response_text);
serde_json::from_str(&response_text)
.map_err(|e| DexError::Other(format!("Failed to parse funding rates: {}", e)))
}
async fn start_websocket(&self) -> Result<(), DexError> {
let ws_url = self
.websocket_url
.replace("https://", "wss://")
.replace("http://", "ws://");
log::info!("Connecting to WebSocket: {}", ws_url);
let ws = match self._ws.as_ref() {
Some(ws) => ws,
None => return Err(DexError::Other("WebSocket not initialized".to_string())),
};
let (_sink, _stream) = ws
.connect()
.await
.map_err(|_| DexError::Other("Failed to connect to WebSocket".to_string()))?;
let primary_market_id = if let Some(symbol) = self.tracked_symbols.first() {
match self.resolve_market_info(symbol).await {
Ok(info) => info.market_id,
Err(e) => {
log::warn!(
"Failed to resolve primary symbol '{}' for WS order book subscription: {}. Falling back to market_id=1",
symbol,
e
);
1
}
}
} else {
1
};
let current_price = self.current_price.clone();
let current_volume = self.current_volume.clone();
let order_book = self.order_book.clone();
let filled_orders = self.filled_orders.clone();
let canceled_orders = self.canceled_orders.clone();
let is_running = self.is_running.clone();
let connection_epoch = self.connection_epoch.clone();
let account_index = self.account_index;
let market_cache = Arc::clone(&self.market_cache);
let default_symbol = self
.tracked_symbols
.first()
.cloned()
.unwrap_or_else(|| "BTC".to_string());
let ws_url_clone = ws_url.clone();
tokio::spawn(async move {
use rand::Rng;
const BACKOFF_MAX_SECS: u64 = 60;
const BACKOFF_BASE: f64 = 1.5;
let mut reconnect_attempt = 0u32;
let mut last_reconnect_time = std::time::SystemTime::now();
async fn reconnect_backoff(attempt: u32) {
let pow = BACKOFF_BASE.powi(attempt.min(12) as i32);
let base_secs = (pow as f64).min(BACKOFF_MAX_SECS as f64);
let jitter_ms: i64 = rand::thread_rng().gen_range(0..=250);
let dur = std::time::Duration::from_secs_f64(base_secs)
+ std::time::Duration::from_millis(jitter_ms as u64);
log::debug!(
"Reconnect backoff: attempt={}, delay={:.1}s",
attempt,
dur.as_secs_f64()
);
tokio::time::sleep(dur).await;
}
loop {
if !is_running.load(Ordering::SeqCst) {
log::info!("WebSocket task stopping due to is_running flag");
break;
}
let now = std::time::SystemTime::now();
if let Ok(elapsed) = now.duration_since(last_reconnect_time) {
if elapsed.as_secs() > 300 {
reconnect_attempt = 0;
log::debug!("Reset reconnect attempt counter after successful period");
}
}
if reconnect_attempt > 0 {
reconnect_backoff(reconnect_attempt).await;
}
log::info!(
"Attempting WebSocket connection to: {} (attempt: {})",
ws_url_clone,
reconnect_attempt + 1
);
last_reconnect_time = now;
use tokio_tungstenite::tungstenite::protocol::WebSocketConfig;
let mut config = WebSocketConfig::default();
config.max_message_size = Some(64 * 1024 * 1024); config.max_frame_size = Some(16 * 1024 * 1024); config.write_buffer_size = 128 * 1024; config.max_write_buffer_size = 1024 * 1024;
let connection_result = tokio_tungstenite::connect_async_with_config(
&ws_url_clone,
Some(config),
false,
)
.await;
match connection_result {
Ok((mut ws_stream, _)) => {
let current_epoch = connection_epoch.fetch_add(1, Ordering::SeqCst) + 1;
let (local_addr, peer_addr) = match ws_stream.get_ref() {
tokio_tungstenite::MaybeTlsStream::Rustls(tls_stream) => {
let tcp_stream = tls_stream.get_ref().0;
if let Err(e) = tcp_stream.set_nodelay(true) {
log::warn!("Failed to set TCP_NODELAY on TLS: {}", e);
}
match (tcp_stream.local_addr(), tcp_stream.peer_addr()) {
(Ok(local), Ok(peer)) => (local, peer),
_ => {
log::warn!(
"Failed to get TLS socket addresses for epoch {}",
current_epoch
);
("0.0.0.0:0".parse().unwrap(), "0.0.0.0:0".parse().unwrap())
}
}
}
tokio_tungstenite::MaybeTlsStream::Plain(tcp_stream) => {
if let Err(e) = tcp_stream.set_nodelay(true) {
log::warn!("Failed to set TCP_NODELAY on plain WS: {}", e);
}
match (tcp_stream.local_addr(), tcp_stream.peer_addr()) {
(Ok(local), Ok(peer)) => (local, peer),
_ => {
log::warn!(
"Failed to get plain socket addresses for epoch {}",
current_epoch
);
("0.0.0.0:0".parse().unwrap(), "0.0.0.0:0".parse().unwrap())
}
}
}
other => {
log::warn!(
"Unsupported WebSocket stream type {:?} for epoch {}",
other,
current_epoch
);
("0.0.0.0:0".parse().unwrap(), "0.0.0.0:0".parse().unwrap())
}
};
let epoch_prefix = format!("[{:03}]", current_epoch);
log::info!(
"{} WebSocket connected successfully: {} -> {}",
epoch_prefix,
local_addr,
peer_addr
);
let orderbook_market_id = primary_market_id;
let subscribe_orderbook = serde_json::json!({
"type": "subscribe",
"channel": format!("order_book/{}", orderbook_market_id)
});
let subscribe_account = serde_json::json!({
"type": "subscribe",
"channel": format!("account_all/{}", account_index)
});
log::info!(
"🔗 [WS_DEBUG] Sending subscriptions - orderbook: {}, account: {}",
subscribe_orderbook,
subscribe_account
);
if let Err(e) = ws_stream
.send(tokio_tungstenite::tungstenite::Message::Text(
subscribe_orderbook.to_string(),
))
.await
{
log::error!("Failed to send orderbook subscription: {}", e);
continue;
}
if let Err(e) = ws_stream
.send(tokio_tungstenite::tungstenite::Message::Text(
subscribe_account.to_string(),
))
.await
{
log::error!("Failed to send account subscription: {}", e);
continue;
}
log::info!("WebSocket subscriptions sent successfully");
use futures::sink::SinkExt;
use futures::stream::StreamExt;
let (write, mut read) = ws_stream.split();
let ws_writer_arc = Arc::new(tokio::sync::Mutex::new(write));
let ws_writer_for_reader = ws_writer_arc.clone();
let (tx_ctrl, mut rx_ctrl) =
tokio::sync::mpsc::channel::<OutboundMessage>(32);
let writer_is_running = is_running.clone();
let ws_writer_for_task = ws_writer_arc.clone();
let _writer_task = tokio::spawn(async move {
loop {
if !writer_is_running.load(Ordering::SeqCst) {
break;
}
let (msg, _) = tokio::select! {
Some(msg) = rx_ctrl.recv() => (msg, true),
else => break,
};
let send_start = std::time::Instant::now();
let is_pong = msg.is_pong();
let mut ws_write = ws_writer_for_task.lock().await;
if let Err(e) = ws_write.send(msg.into_message()).await {
log::error!("WebSocket send failed: {:?}", e);
break;
}
let send_duration = send_start.elapsed();
if is_pong {
let latency_ms = send_duration.as_millis();
if latency_ms > 100 {
log::warn!("High pong send latency: {}ms", latency_ms);
}
}
}
log::debug!("WebSocket writer task terminated");
});
const IDLE_PING_SECS: u64 = 20; const PONG_TIMEOUT_SECS: u64 = 8; const HEARTBEAT_CHECK_SECS: u64 = 5;
fn get_pong_payload(ping_payload: &[u8]) -> Vec<u8> {
ping_payload.to_vec()
}
use parking_lot::Mutex;
use std::sync::atomic::{AtomicBool, AtomicU64};
use std::time::{SystemTime, UNIX_EPOCH};
fn now_secs() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
}
let last_rx = std::sync::Arc::new(AtomicU64::new(now_secs()));
let last_tx = std::sync::Arc::new(AtomicU64::new(now_secs()));
let last_app_ping = std::sync::Arc::new(AtomicU64::new(now_secs()));
let last_server_ping = std::sync::Arc::new(AtomicU64::new(0));
let pending_client_ping = std::sync::Arc::new(AtomicBool::new(false));
let pending_app_pong = std::sync::Arc::new(AtomicBool::new(false));
let last_client_ping_payload =
std::sync::Arc::new(Mutex::new(Vec::<u8>::new()));
let ping_is_running = is_running.clone();
let _ping_last_rx = last_rx.clone();
let ping_last_tx = last_tx.clone();
let ping_last_app_ping = last_app_ping.clone();
let _ping_last_server_ping = last_server_ping.clone();
let ping_pending_client_ping = pending_client_ping.clone();
let ping_pending_app_pong = pending_app_pong.clone();
let ping_last_client_ping_payload = last_client_ping_payload.clone();
let ping_tx_ctrl = tx_ctrl.clone();
let _ping_task = tokio::spawn(async move {
let mut heartbeat_interval = tokio::time::interval(
std::time::Duration::from_secs(HEARTBEAT_CHECK_SECS),
);
heartbeat_interval
.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
_ = heartbeat_interval.tick() => {
if !ping_is_running.load(Ordering::SeqCst) {
break;
}
let now = now_secs();
let idle_tx = now.saturating_sub(ping_last_tx.load(Ordering::SeqCst));
if !ping_pending_client_ping.load(Ordering::SeqCst)
&& idle_tx >= IDLE_PING_SECS
{
let payload: [u8; 8] = (now as u64).to_be_bytes();
*ping_last_client_ping_payload.lock() = payload.to_vec();
let ping_msg = OutboundMessage::Control(
tokio_tungstenite::tungstenite::Message::Ping(payload.to_vec())
);
if let Err(e) = ping_tx_ctrl.send(ping_msg).await {
log::warn!("Failed to send client ping: {:?}", e);
break;
}
ping_pending_client_ping.store(true, Ordering::SeqCst);
ping_last_tx.store(now, Ordering::SeqCst);
}
let idle_app_ping = now.saturating_sub(ping_last_app_ping.load(Ordering::SeqCst));
if !ping_pending_app_pong.load(Ordering::SeqCst) && idle_app_ping >= IDLE_PING_SECS {
let app_ping = serde_json::json!({
"type": "ping",
"ts": now
});
let ping_msg = OutboundMessage::Control(
tokio_tungstenite::tungstenite::Message::Text(app_ping.to_string())
);
if let Err(e) = ping_tx_ctrl.send(ping_msg).await {
log::warn!("Failed to send application-layer ping: {:?}", e);
break;
}
ping_pending_app_pong.store(true, Ordering::SeqCst);
ping_last_app_ping.store(now, Ordering::SeqCst);
}
if ping_pending_client_ping.load(Ordering::SeqCst) {
let waited = now.saturating_sub(ping_last_tx.load(Ordering::SeqCst));
if waited >= PONG_TIMEOUT_SECS {
log::warn!("Control pong timeout ({}s), reconnecting", waited);
let close_msg = OutboundMessage::Control(
tokio_tungstenite::tungstenite::Message::Close(None)
);
let _ = ping_tx_ctrl.send(close_msg).await;
break;
}
}
if ping_pending_app_pong.load(Ordering::SeqCst) {
let waited = now.saturating_sub(ping_last_app_ping.load(Ordering::SeqCst));
if waited >= PONG_TIMEOUT_SECS {
log::warn!("Application pong timeout ({}s), reconnecting", waited);
let close_msg = OutboundMessage::Control(
tokio_tungstenite::tungstenite::Message::Close(None)
);
let _ = ping_tx_ctrl.send(close_msg).await;
break;
}
}
}
}
}
log::debug!("Heartbeat task ended");
});
log::debug!("Starting WebSocket message handling loop");
while let Some(message) = read.next().await {
if !is_running.load(Ordering::SeqCst) {
log::info!("WebSocket stopping due to is_running flag");
break;
}
match message {
Ok(message) => match message {
tokio_tungstenite::tungstenite::Message::Text(text) => {
let msg_start = std::time::Instant::now();
log::trace!("WebSocket text message: {}", text);
let now = now_secs();
last_rx.store(now, Ordering::SeqCst);
if let Ok(parsed) = serde_json::from_str::<Value>(&text) {
if let Some(msg_type) =
parsed.get("type").and_then(|t| t.as_str())
{
if msg_type == "ping" {
let mut pong =
serde_json::json!({"type": "pong"});
if let Some(ts) = parsed.get("ts") {
pong["ts"] = ts.clone();
}
if let Some(id) = parsed.get("id") {
pong["id"] = id.clone();
}
if let Some(nonce) = parsed.get("nonce") {
pong["nonce"] = nonce.clone();
}
if let Ok(mut ws_write) =
ws_writer_for_reader.try_lock()
{
if let Err(e) = ws_write.send(
tokio_tungstenite::tungstenite::Message::Text(pong.to_string())
).await {
log::error!("Failed to send application-layer pong: {:?}", e);
} else {
}
} else {
}
continue;
} else if msg_type == "pong" {
pending_app_pong.store(false, Ordering::SeqCst);
continue;
}
}
Self::handle_websocket_message(
parsed,
¤t_price,
¤t_volume,
&order_book,
&filled_orders,
&canceled_orders,
account_index,
&market_cache,
default_symbol.as_str(),
)
.await;
let total_duration = msg_start.elapsed();
if total_duration.as_millis() > 10 {
log::warn!(
"Slow message processing: {}ms (len={})",
total_duration.as_millis(),
text.len()
);
}
} else {
log::warn!(
"Failed to parse WebSocket message as JSON: {}",
text
);
}
}
tokio_tungstenite::tungstenite::Message::Ping(payload) => {
let now = now_secs();
last_server_ping.store(now, Ordering::SeqCst);
last_rx.store(now, Ordering::SeqCst);
let pong_payload = get_pong_payload(&payload);
let current_epoch = connection_epoch.load(Ordering::SeqCst);
if let Ok(mut ws_write) = ws_writer_for_reader.try_lock() {
if let Err(e) = ws_write
.send(
tokio_tungstenite::tungstenite::Message::Pong(
pong_payload.clone(),
),
)
.await
{
log::error!("🚨 [CRITICAL] Failed to send pong directly: {:?}", e);
break;
}
if let Err(e) =
futures::SinkExt::flush(&mut *ws_write).await
{
log::error!("🚨 [CRITICAL] Failed to flush after pong: {:?}", e);
break;
}
last_tx.store(now, Ordering::SeqCst);
let pong_epoch =
connection_epoch.load(Ordering::SeqCst);
if pong_epoch != current_epoch {
log::error!(
"Pong epoch mismatch: ping={} pong={}",
current_epoch,
pong_epoch
);
}
} else {
let mut pong_sent = false;
for retry in 0..4 {
tokio::time::sleep(
std::time::Duration::from_millis(50),
)
.await;
if let Ok(mut ws_write) =
ws_writer_for_reader.try_lock()
{
if let Err(e) = ws_write.send(tokio_tungstenite::tungstenite::Message::Pong(pong_payload.clone())).await {
log::error!("Failed to send pong on retry {}: {:?}", retry, e);
break;
}
if let Err(e) =
futures::SinkExt::flush(&mut *ws_write)
.await
{
log::error!("Failed to flush after pong retry {}: {:?}", retry, e);
break;
}
last_tx.store(now, Ordering::SeqCst);
pong_sent = true;
break;
}
}
if !pong_sent {
log::warn!("Pong timeout, closing connection");
let mut ws_write =
ws_writer_for_reader.lock().await;
let _ = ws_write.close().await;
break;
}
}
}
tokio_tungstenite::tungstenite::Message::Pong(payload) => {
let now = now_secs();
last_rx.store(now, Ordering::SeqCst);
let expected_payload =
last_client_ping_payload.lock().clone();
if !expected_payload.is_empty()
&& payload == expected_payload
{
pending_client_ping.store(false, Ordering::SeqCst);
} else {
}
}
tokio_tungstenite::tungstenite::Message::Close(frame) => {
log::warn!("WebSocket close frame received: {:?}", frame);
break;
}
tokio_tungstenite::tungstenite::Message::Binary(data) => {
log::debug!(
"Received binary WebSocket message: {} bytes",
data.len()
);
}
tokio_tungstenite::tungstenite::Message::Frame(_) => {
log::trace!("Received raw WebSocket frame");
}
},
Err(e) => {
log::error!(
"WebSocket error: {} (type: {:?}). Will attempt reconnection.",
e, std::any::type_name_of_val(&e)
);
break; }
}
}
reconnect_attempt += 1;
}
Err(e) => {
reconnect_attempt += 1;
let error_str = e.to_string();
if error_str.contains("429") || error_str.contains("Too Many Requests") {
log::error!(
"WebSocket connection failed with rate limit (429): {}. Attempt: {}. Using exponential backoff.",
e, reconnect_attempt
);
} else {
log::error!(
"Failed to connect to WebSocket: {}. Attempt: {}. Will retry with backoff.",
e, reconnect_attempt
);
}
}
}
}
log::info!("WebSocket task ended");
});
Ok(())
}
async fn handle_websocket_message(
message: Value,
current_price: &Arc<RwLock<Option<(Decimal, u64)>>>,
current_volume: &Arc<RwLock<Option<Decimal>>>,
order_book: &Arc<RwLock<Option<LighterOrderBook>>>,
filled_orders: &Arc<RwLock<HashMap<String, Vec<FilledOrder>>>>,
canceled_orders: &Arc<RwLock<HashMap<String, Vec<CanceledOrder>>>>,
account_index: u32,
market_cache: &Arc<RwLock<MarketCache>>,
default_symbol: &str,
) {
let msg_type = message.get("type").and_then(|t| t.as_str()).unwrap_or("");
match msg_type {
"subscribed/order_book" | "update/order_book" => {
if let Some(order_book_data) = message.get("order_book") {
if let Ok(ob) =
serde_json::from_value::<LighterOrderBook>(order_book_data.clone())
{
if let (Some(best_bid), Some(best_ask)) = (ob.bids.first(), ob.asks.first())
{
if let (Ok(bid_price), Ok(ask_price)) = (
string_to_decimal(Some(best_bid.price.clone())),
string_to_decimal(Some(best_ask.price.clone())),
) {
let mid_price = (bid_price + ask_price) / Decimal::from(2);
let timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs();
*current_price.write().await = Some((mid_price, timestamp));
log::trace!(
"Updated price from WebSocket: {} at {}",
mid_price,
timestamp
);
}
}
let total_volume: Decimal = ob
.bids
.iter()
.chain(ob.asks.iter())
.filter_map(|entry| string_to_decimal(Some(entry.size.clone())).ok())
.sum();
*current_volume.write().await = Some(total_volume);
*order_book.write().await = Some(ob);
}
}
}
"subscribed/account_all" | "update/account_all" => {
log::trace!(
"Received account message: type={}, message={:?}",
msg_type,
message
);
Self::handle_account_update(
&message,
filled_orders,
canceled_orders,
account_index as u64,
market_cache,
default_symbol,
)
.await;
}
_ => {
log::trace!("Unhandled WebSocket message type: {}", msg_type);
}
}
}
async fn handle_account_update(
data: &Value,
filled_orders: &Arc<RwLock<HashMap<String, Vec<FilledOrder>>>>,
canceled_orders: &Arc<RwLock<HashMap<String, Vec<CanceledOrder>>>>,
account_id: u64,
market_cache: &Arc<RwLock<MarketCache>>,
default_symbol: &str,
) {
log::trace!("handle_account_update called with data: {:?}", data);
if let Some(fills) = data.get("fills").and_then(|f| f.as_array()) {
log::info!(
"✅ [FILL_DETECTION] Found {} fills in account update",
fills.len()
);
let default_symbol = default_symbol.to_string();
let mut filled_map = filled_orders.write().await;
for fill in fills {
log::debug!("🔍 [FILL_DETECTION] Processing fill: {:?}", fill);
if let Ok(filled_order) = Self::parse_filled_order(fill, account_id) {
log::info!("✅ [FILL_DETECTION] Added filled order: order_id={}, size={:?}, value={:?}",
filled_order.order_id, filled_order.filled_size, filled_order.filled_value);
filled_map
.entry(default_symbol.clone())
.or_insert_with(Vec::new)
.push(filled_order);
} else {
log::warn!("Failed to parse filled order: {:?}", fill);
}
}
}
if let Some(trades) = data.get("trades").and_then(|t| t.as_object()) {
log::info!("✅ [FILL_DETECTION] Found trades object in account update");
let mut pending_inserts: Vec<(String, FilledOrder)> = Vec::new();
for (market_id, trade_array) in trades {
let market_id_num = match market_id.parse::<u32>() {
Ok(id) => id,
Err(_) => {
log::warn!(
"[FILL_DETECTION] Unable to parse market_id '{}' as u32",
market_id
);
continue;
}
};
let market_symbol = {
let cache = market_cache.read().await;
cache
.by_id
.get(&market_id_num)
.map(|info| info.canonical_symbol.clone())
};
let market_symbol = match market_symbol {
Some(symbol) => symbol,
None => {
log::warn!(
"[FILL_DETECTION] Market cache missing entry for market_id {}",
market_id_num
);
continue;
}
};
if let Some(trades_array) = trade_array.as_array() {
for trade in trades_array {
if let (Some(ask_id), Some(bid_id), Some(size_str), Some(price_str)) = (
trade.get("ask_id").and_then(|v| v.as_u64()),
trade.get("bid_id").and_then(|v| v.as_u64()),
trade.get("size").and_then(|v| v.as_str()),
trade.get("price").and_then(|v| v.as_str()),
) {
let order_id = if account_id == ask_id { ask_id } else { bid_id };
log::info!(
"✅ [FILL_DETECTION] Trade detected: order_id={}, size={}, price={}, market_id={}",
order_id, size_str, price_str, market_id_num
);
if let (Ok(size), Ok(price)) = (
size_str.parse::<rust_decimal::Decimal>(),
price_str.parse::<rust_decimal::Decimal>(),
) {
let filled_order = FilledOrder {
order_id: order_id.to_string(),
is_rejected: false,
trade_id: trade
.get("trade_id")
.and_then(|v| v.as_u64())
.unwrap_or(0)
.to_string(),
filled_side: if account_id == ask_id {
Some(OrderSide::Short)
} else {
Some(OrderSide::Long)
},
filled_size: Some(size),
filled_value: Some(size * price),
filled_fee: None,
};
pending_inserts.push((market_symbol.clone(), filled_order));
log::info!(
"✅ [FILL_DETECTION] Added filled order from trade: order_id={}",
order_id
);
}
}
}
}
}
if !pending_inserts.is_empty() {
let mut filled_map = filled_orders.write().await;
for (symbol_key, filled_order) in pending_inserts {
filled_map
.entry(symbol_key)
.or_insert_with(Vec::new)
.push(filled_order);
}
}
} else {
log::trace!("No 'fills' array or 'trades' object found in account data");
}
if let Some(cancels) = data.get("cancels").and_then(|c| c.as_array()) {
let default_symbol = default_symbol.to_string();
let mut canceled_map = canceled_orders.write().await;
for cancel in cancels {
if let Ok(canceled_order) = Self::parse_canceled_order(cancel) {
canceled_map
.entry(default_symbol.clone())
.or_insert_with(Vec::new)
.push(canceled_order);
}
}
}
}
fn parse_filled_order(_data: &Value, _account_id: u64) -> Result<FilledOrder, DexError> {
Err(DexError::Other(
"Filled order tracking not supported for Lighter DEX".to_string(),
))
}
fn parse_canceled_order(data: &Value) -> Result<CanceledOrder, DexError> {
Ok(CanceledOrder {
order_id: data
.get("order_id")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string(),
canceled_timestamp: data.get("timestamp").and_then(|v| v.as_u64()).unwrap_or(0),
})
}
fn calculate_min_tick(price: Decimal, sz_decimals: u32, is_spot: bool) -> Decimal {
let price_str = price.to_string();
let integer_part = price_str.split('.').next().unwrap_or("");
let integer_digits = if integer_part == "0" {
0
} else {
integer_part.len()
};
let scale_by_sig: u32 = if integer_digits >= 5 {
0
} else {
(5 - integer_digits) as u32
};
let max_decimals: u32 = if is_spot { 8u32 } else { 6u32 };
let scale_by_dec: u32 = max_decimals.saturating_sub(sz_decimals);
let scale: u32 = scale_by_sig.min(scale_by_dec);
Decimal::new(1, scale)
}
}
pub fn create_lighter_connector(
api_key_public: String,
api_key_index: u32,
api_private_key_hex: String,
evm_wallet_private_key: Option<String>,
account_index: u32,
base_url: String,
websocket_url: String,
tracked_symbols: Vec<String>,
) -> Result<Box<dyn DexConnector>, DexError> {
let connector = LighterConnector::new(
api_key_public,
api_key_index,
api_private_key_hex,
evm_wallet_private_key,
account_index,
base_url,
websocket_url,
tracked_symbols,
)?;
Ok(Box::new(connector))
}
#[cfg(test)]
mod tests {
use super::*;
use std::env;
#[tokio::test]
async fn test_get_open_orders() {
let api_key_public = match env::var("LIGHTER_PLAIN_PUBLIC_API_KEY") {
Ok(key) => key,
Err(_) => {
println!("Skipping test - LIGHTER_PLAIN_PUBLIC_API_KEY not set");
return;
}
};
let base_url = env::var("LIGHTER_BASE_URL")
.unwrap_or_else(|_| "https://mainnet.zklighter.elliot.ai".to_string());
let account_index = env::var("LIGHTER_ACCOUNT_INDEX")
.unwrap_or_else(|_| "0".to_string())
.parse::<u32>()
.unwrap_or(0);
let connector = match LighterConnector::new(
api_key_public,
0, "dummy_private_key".to_string(),
None, account_index,
base_url,
"dummy_websocket_url".to_string(),
) {
Ok(c) => c,
Err(e) => {
println!("Failed to create connector: {}", e);
return;
}
};
match connector.get_open_orders("BTC").await {
Ok(response) => {
println!(
"✅ get_open_orders success: {} orders found",
response.orders.len()
);
for (i, order) in response.orders.iter().enumerate() {
println!(" Order {}: {}", i, order.order_id);
}
}
Err(e) => {
println!("❌ get_open_orders failed: {}", e);
panic!("get_open_orders test failed: {}", e);
}
}
}
}