Working subscribe and unsubscribe

This commit is contained in:
2026-08-07 21:11:46 +02:00
parent af528e589b
commit 6032b56ab8
3 changed files with 17 additions and 11 deletions
+11 -5
View File
@@ -1,7 +1,10 @@
mod pump_fun; mod pump_fun;
use rust_decimal::Decimal; use rust_decimal::Decimal;
use std::{collections::HashMap, sync::Arc}; use std::{
collections::{HashMap, HashSet},
sync::Arc,
};
use tokio::sync::{Mutex, mpsc}; use tokio::sync::{Mutex, mpsc};
use crate::{data::Event, launchpad::pump_fun::PumpFun}; use crate::{data::Event, launchpad::pump_fun::PumpFun};
@@ -13,6 +16,7 @@ type ClientMutex = Mutex<
pub struct Client { pub struct Client {
pub helius: ClientMutex, pub helius: ClientMutex,
pub solana: ClientMutex, pub solana: ClientMutex,
pub subscribed: Mutex<HashSet<String>>,
} }
pub struct Executor { pub struct Executor {
@@ -72,6 +76,8 @@ impl Executor {
.await? .await?
.0, .0,
), ),
subscribed: Mutex::new(HashSet::new()),
}), }),
event_tx: Mutex::new(tx), event_tx: Mutex::new(tx),
@@ -133,11 +139,11 @@ impl Executor {
Ok(rx) Ok(rx)
} }
pub async fn subscribe(&self, mint: &str) -> anyhow::Result<()> { pub async fn subscribe(&self, mint: &str) {
Ok(()) self.client.subscribed.lock().await.insert(mint.to_string());
} }
pub async fn unsubscribe(&self, mint: &str) -> anyhow::Result<()> { pub async fn unsubscribe(&self, mint: &str) {
Ok(()) self.client.subscribed.lock().await.remove(mint);
} }
} }
+3 -1
View File
@@ -120,7 +120,9 @@ impl Launchpad for PumpFun {
if let Some(trade) = if let Some(trade) =
parse_trade_event(&data, notification.params.result.value.signature.clone()) 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?;
}
} }
} }
} }
+3 -5
View File
@@ -64,9 +64,7 @@ impl MomentumVelocityStrategy {
/// Internal helper to safely handle unsubscribing and cleaning state /// Internal helper to safely handle unsubscribing and cleaning state
async fn cleanup_and_unsubscribe(&mut self, bot: &Arc<Bot>, mint: &str) -> anyhow::Result<()> { async fn cleanup_and_unsubscribe(&mut self, bot: &Arc<Bot>, mint: &str) -> anyhow::Result<()> {
debug!("[{}] Cleaning up state and unsubscribing", mint); debug!("[{}] Cleaning up state and unsubscribing", mint);
if let Err(e) = bot.executor.unsubscribe(mint).await { bot.executor.unsubscribe(mint).await;
warn!("[{}] Unsubscribe request failed: {:?}", mint, e);
}
self.trackers.remove(mint); self.trackers.remove(mint);
self.active_subscriptions.retain(|m| m != mint); self.active_subscriptions.retain(|m| m != mint);
Ok(()) Ok(())
@@ -103,7 +101,7 @@ impl Strategy for MomentumVelocityStrategy {
"[{}] Capacity reached ({}/{}). Evicting un-bought token from queue.", "[{}] Capacity reached ({}/{}). Evicting un-bought token from queue.",
mint_to_remove, MAX_SUBSCRIBED_TOKENS, MAX_SUBSCRIBED_TOKENS 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); self.trackers.remove(&mint_to_remove);
} }
} else { } else {
@@ -116,7 +114,7 @@ impl Strategy for MomentumVelocityStrategy {
} }
info!("[{}] Subscribing and creating tracker.", token.mint); 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.active_subscriptions.push_back(token.mint.clone());
self.trackers.insert( self.trackers.insert(