diff --git a/src/bot.rs b/src/bot.rs index d58870a..c078352 100644 --- a/src/bot.rs +++ b/src/bot.rs @@ -127,8 +127,8 @@ impl Bot { ws = self.ws.lock().await; } - Ok(event) => { - println!("{:?}", event); + Ok(_event) => { + // println!("{:?}", event); } Err(err) => { diff --git a/src/executor/pump_fun.rs b/src/executor/pump_fun.rs index f448e70..186e0c0 100644 --- a/src/executor/pump_fun.rs +++ b/src/executor/pump_fun.rs @@ -93,7 +93,7 @@ impl Executor for PumpDev { priority: Decimal, slippage: u16, ) -> anyhow::Result<()> { - println!("BUY {mint} {amount} SOL"); + log::info!("BUY {mint} {amount} SOL"); Ok(()) // self.trade( // "buy", @@ -113,7 +113,7 @@ impl Executor for PumpDev { priority: Decimal, slippage: u16, ) -> anyhow::Result<()> { - println!("SELL {mint} {amount} SOL"); + log::info!("SELL {mint} {amount} SOL"); Ok(()) // self.trade( // "sell", @@ -133,7 +133,7 @@ impl Executor for PumpDev { priority: Decimal, slippage: u16, ) -> anyhow::Result<()> { - println!("SELL {mint} {amount}%"); + log::info!("SELL {mint} {amount}%"); Ok(()) // self.trade( // "sell", diff --git a/src/strategy/veloc.rs b/src/strategy/veloc.rs index 50ae19d..3d8243a 100644 --- a/src/strategy/veloc.rs +++ b/src/strategy/veloc.rs @@ -2,6 +2,7 @@ use std::collections::{HashMap, HashSet, VecDeque}; use std::sync::Arc; use std::time::{Duration, Instant}; +use log::{debug, info, trace, warn}; use rust_decimal::{Decimal, dec}; use crate::bot::Bot; @@ -41,8 +42,8 @@ pub struct MomentumVelocityStrategy { impl MomentumVelocityStrategy { pub fn new() -> Self { Self { - min_unique_buyers: 4, - min_net_sol_flow: 1.5, + min_unique_buyers: 2, + min_net_sol_flow: , max_tracking_duration: Duration::from_secs(45), trackers: HashMap::new(), positions: HashMap::new(), @@ -59,7 +60,10 @@ impl MomentumVelocityStrategy { /// Internal helper to safely handle unsubscribing and cleaning state async fn cleanup_and_unsubscribe(&mut self, bot: &Arc, mint: &str) -> anyhow::Result<()> { - bot.unsubscribe(mint).await?; + debug!("[{}] Cleaning up state and unsubscribing", mint); + if let Err(e) = bot.unsubscribe(mint).await { + warn!("[{}] Unsubscribe request failed: {:?}", mint, e); + } self.trackers.remove(mint); self.active_subscriptions.retain(|m| m != mint); Ok(()) @@ -69,28 +73,45 @@ impl MomentumVelocityStrategy { #[async_trait::async_trait] impl Strategy for MomentumVelocityStrategy { async fn on_new_coin(&mut self, bot: Arc, token: NewToken) -> anyhow::Result<()> { - // If we hold an active position in this mint already, skip tracking setup - if self.positions.contains_key(&token.mint) { + trace!("[NEW COIN] Event received for token: {}", token.mint); + + if self.positions.contains_key(&token.mint) || self.trackers.contains_key(&token.mint) { + trace!( + "[{}] Already tracking or holding position. Skipping.", + token.mint + ); return Ok(()); } - // Enforce maximum 5 active subscriptions + // Evict oldest tracked token that DOES NOT have an active open position while self.active_subscriptions.len() >= MAX_SUBSCRIBED_TOKENS { - if let Some(oldest_mint) = self.active_subscriptions.pop_front() { - // Do not drop subscription if we currently hold an open position in it - if self.positions.contains_key(&oldest_mint) { - continue; + let eviction_index = self + .active_subscriptions + .iter() + .position(|mint| !self.positions.contains_key(mint)); + + if let Some(idx) = eviction_index { + if let Some(mint_to_remove) = self.active_subscriptions.remove(idx) { + info!( + "[{}] Capacity reached ({}/{}). Evicting un-bought token from queue.", + mint_to_remove, MAX_SUBSCRIBED_TOKENS, MAX_SUBSCRIBED_TOKENS + ); + let _ = bot.unsubscribe(&mint_to_remove).await; + self.trackers.remove(&mint_to_remove); } - bot.unsubscribe(&oldest_mint).await?; - self.trackers.remove(&oldest_mint); + } else { + warn!( + "[QUEUE FULL] All {} slots are occupied by active positions. Cannot track {}", + MAX_SUBSCRIBED_TOKENS, token.mint + ); + return Ok(()); } } - // Subscribe to new token + info!("[{}] Subscribing and creating tracker.", token.mint); bot.subscribe(&token.mint).await?; self.active_subscriptions.push_back(token.mint.clone()); - // Initialize tracking self.trackers.insert( token.mint, TokenTracker { @@ -117,19 +138,36 @@ impl Strategy for MomentumVelocityStrategy { if current_price > pos.highest_price_sol { pos.highest_price_sol = current_price; pos.last_high_time = Instant::now(); + trace!("[{}] New high reached: {:.9} SOL", mint, current_price); } let drop_from_peak = (pos.highest_price_sol - current_price) / pos.highest_price_sol; - let should_sell = match () { - _ if price_change_pct >= 0.40 => true, // Take Profit: +40% - _ if price_change_pct <= -0.15 => true, // Hard Stop Loss: -15% - _ if drop_from_peak >= 0.12 && price_change_pct > 0.10 => true, // Trailing stop: 12% drop from peak - _ if pos.last_high_time.elapsed() >= Duration::from_secs(25) => true, // Momentum stalled for 25s - _ => false, + let (should_sell, reason) = match () { + _ if price_change_pct >= 0.40 => ( + true, + format!("Take Profit (+{:.1}%)", price_change_pct * 100.0), + ), + _ if price_change_pct <= -0.15 => ( + true, + format!("Hard Stop Loss ({:.1}%)", price_change_pct * 100.0), + ), + _ if drop_from_peak >= 0.12 && price_change_pct > 0.10 => ( + true, + format!( + "Trailing Stop (Peak drop: {:.1}%, gain: +{:.1}%)", + drop_from_peak * 100.0, + price_change_pct * 100.0 + ), + ), + _ if pos.last_high_time.elapsed() >= Duration::from_secs(25) => { + (true, format!("Momentum Stalled (no high for 25s)")) + } + _ => (false, String::new()), }; if should_sell { + info!("[{}] EXECUTING SELL. Reason: {}", mint, reason); bot.executor .lock() .await @@ -147,17 +185,20 @@ impl Strategy for MomentumVelocityStrategy { // 2. Evaluate Potential Buys // ------------------------------------------------------------- if let Some(tracker) = self.trackers.get_mut(mint) { - // Unsubscribe & remove if evaluation window expired without a signal - if tracker.created_at.elapsed() > self.max_tracking_duration { + let elapsed = tracker.created_at.elapsed(); + if elapsed > self.max_tracking_duration { + info!( + "[{}] Tracking window expired ({:?} > {:?}). Cleaning up.", + mint, elapsed, self.max_tracking_duration + ); self.cleanup_and_unsubscribe(&bot, mint).await?; return Ok(()); } - // Update trade metrics tracker.trade_count += 1; match trade.tx_type { TradeType::Buy => { - tracker.unique_buyers.insert(trade.trader); + tracker.unique_buyers.insert(trade.trader.clone()); tracker.net_sol_flow += trade.sol_amount; } TradeType::Sell => { @@ -169,14 +210,37 @@ impl Strategy for MomentumVelocityStrategy { let has_volume_surge = tracker.net_sol_flow >= self.min_net_sol_flow; let is_early_curve = trade.v_sol_in_bonding_curve < 60.0; + // Log detailed status of buy criteria evaluation on every trade + debug!( + "[{}] Trade #{} ({:?}) | Buyers: {}/{} [{}] | Net Flow: {:.3}/{:.3} SOL [{}] | Curve SOL: {:.2} < 60 [{}]", + mint, + tracker.trade_count, + trade.tx_type, + tracker.unique_buyers.len(), + self.min_unique_buyers, + if has_enough_buyers { "PASS" } else { "FAIL" }, + tracker.net_sol_flow, + self.min_net_sol_flow, + if has_volume_surge { "PASS" } else { "FAIL" }, + trade.v_sol_in_bonding_curve, + if is_early_curve { "PASS" } else { "FAIL" } + ); + if has_enough_buyers && has_volume_surge && is_early_curve { + info!( + "🚀 BUY SIGNAL TRIGGERED for {}! Unique Buyers: {}, Net Flow: {:.3} SOL, Curve SOL: {:.2}", + mint, + tracker.unique_buyers.len(), + tracker.net_sol_flow, + trade.v_sol_in_bonding_curve + ); + bot.executor .lock() .await .buy(mint, BUY_AMOUNT_SOL, PRIORITY, SLIPPAGE) .await?; - // Register open position self.positions.insert( mint.clone(), OpenPosition { @@ -186,9 +250,10 @@ impl Strategy for MomentumVelocityStrategy { }, ); - // Stop tracking buy metrics (subscription remains active while holding position) self.trackers.remove(mint); } + } else { + trace!("[{}] Received trade for untracked mint.", mint); } Ok(()) diff --git a/src/types.rs b/src/types.rs index 918b9d5..fc2a030 100644 --- a/src/types.rs +++ b/src/types.rs @@ -6,6 +6,7 @@ pub enum PumpDevEvent { Connected { client_id: u64, message: String }, ConnectionStatus { connected: bool, timestamp: u64 }, Subscribed { method: String }, + Unsubscribed { method: String, keys: Vec }, Create(NewToken), Trade(Trade), } @@ -40,6 +41,8 @@ impl<'de> Deserialize<'de> for PumpDevEvent { ConnectionStatus { connected: bool, timestamp: u64 }, #[serde(rename = "subscribed")] Subscribed { method: String }, + #[serde(rename = "unsubscribed")] + Unsubscribed { method: String, keys: Vec }, } match serde_json::from_value(value).map_err(serde::de::Error::custom)? { @@ -54,6 +57,9 @@ impl<'de> Deserialize<'de> for PumpDevEvent { timestamp, }), Tagged::Subscribed { method } => Ok(PumpDevEvent::Subscribed { method }), + Tagged::Unsubscribed { method, keys } => { + Ok(PumpDevEvent::Unsubscribed { method, keys }) + } } } }