diff --git a/src/launchpad/mod.rs b/src/launchpad/mod.rs index 51ed05b..9a98d65 100644 --- a/src/launchpad/mod.rs +++ b/src/launchpad/mod.rs @@ -1,7 +1,10 @@ mod pump_fun; use rust_decimal::Decimal; -use std::{collections::HashMap, sync::Arc}; +use std::{ + collections::{HashMap, HashSet}, + sync::Arc, +}; use tokio::sync::{Mutex, mpsc}; use crate::{data::Event, launchpad::pump_fun::PumpFun}; @@ -13,6 +16,7 @@ type ClientMutex = Mutex< pub struct Client { pub helius: ClientMutex, pub solana: ClientMutex, + pub subscribed: Mutex>, } pub struct Executor { @@ -72,6 +76,8 @@ impl Executor { .await? .0, ), + + subscribed: Mutex::new(HashSet::new()), }), event_tx: Mutex::new(tx), @@ -133,11 +139,11 @@ impl Executor { Ok(rx) } - pub async fn subscribe(&self, mint: &str) -> anyhow::Result<()> { - Ok(()) + pub async fn subscribe(&self, mint: &str) { + self.client.subscribed.lock().await.insert(mint.to_string()); } - pub async fn unsubscribe(&self, mint: &str) -> anyhow::Result<()> { - Ok(()) + pub async fn unsubscribe(&self, mint: &str) { + self.client.subscribed.lock().await.remove(mint); } } diff --git a/src/launchpad/pump_fun.rs b/src/launchpad/pump_fun.rs index 7112bba..bbe8ab2 100644 --- a/src/launchpad/pump_fun.rs +++ b/src/launchpad/pump_fun.rs @@ -120,7 +120,9 @@ impl Launchpad for PumpFun { if let Some(trade) = parse_trade_event(&data, notification.params.result.value.signature.clone()) { - tx.send(Event::Trade(trade)).await?; + if client.subscribed.lock().await.contains(&trade.mint) { + tx.send(Event::Trade(trade)).await?; + } } } } diff --git a/src/strategy/veloc.rs b/src/strategy/veloc.rs index 61ed627..ccfef1e 100644 --- a/src/strategy/veloc.rs +++ b/src/strategy/veloc.rs @@ -64,9 +64,7 @@ impl MomentumVelocityStrategy { /// Internal helper to safely handle unsubscribing and cleaning state async fn cleanup_and_unsubscribe(&mut self, bot: &Arc, mint: &str) -> anyhow::Result<()> { debug!("[{}] Cleaning up state and unsubscribing", mint); - if let Err(e) = bot.executor.unsubscribe(mint).await { - warn!("[{}] Unsubscribe request failed: {:?}", mint, e); - } + bot.executor.unsubscribe(mint).await; self.trackers.remove(mint); self.active_subscriptions.retain(|m| m != mint); Ok(()) @@ -103,7 +101,7 @@ impl Strategy for MomentumVelocityStrategy { "[{}] Capacity reached ({}/{}). Evicting un-bought token from queue.", mint_to_remove, MAX_SUBSCRIBED_TOKENS, MAX_SUBSCRIBED_TOKENS ); - let _ = bot.executor.unsubscribe(&mint_to_remove).await; + bot.executor.unsubscribe(&mint_to_remove).await; self.trackers.remove(&mint_to_remove); } } else { @@ -116,7 +114,7 @@ impl Strategy for MomentumVelocityStrategy { } info!("[{}] Subscribing and creating tracker.", token.mint); - bot.executor.subscribe(&token.mint).await?; + bot.executor.subscribe(&token.mint).await; self.active_subscriptions.push_back(token.mint.clone()); self.trackers.insert(