Refactored
This commit is contained in:
+1
-9
@@ -6,7 +6,7 @@ use tokio::sync::{Mutex, mpsc, watch};
|
||||
|
||||
use crate::{
|
||||
data::{
|
||||
Event, NewToken,
|
||||
Event,
|
||||
account::{Account, AccountManager},
|
||||
tradelog::TradeLog,
|
||||
},
|
||||
@@ -105,14 +105,6 @@ impl Bot {
|
||||
}
|
||||
|
||||
impl Bot {
|
||||
pub async fn subscribe(&self, mint: &str) -> anyhow::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn unsubscribe(&self, mint: &str) -> anyhow::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn refresh_account(self: &Arc<Self>) -> anyhow::Result<()> {
|
||||
self.strategy
|
||||
.lock()
|
||||
|
||||
+33
-13
@@ -6,17 +6,19 @@ use tokio::sync::{Mutex, mpsc};
|
||||
|
||||
use crate::{data::Event, launchpad::pump_fun::PumpFun};
|
||||
|
||||
pub struct Client(
|
||||
pub Mutex<
|
||||
tokio_tungstenite::WebSocketStream<
|
||||
tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>,
|
||||
>,
|
||||
>,
|
||||
);
|
||||
type ClientMutex = Mutex<
|
||||
tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>,
|
||||
>;
|
||||
|
||||
pub struct Client {
|
||||
pub helius: ClientMutex,
|
||||
pub solana: ClientMutex,
|
||||
}
|
||||
|
||||
pub struct Executor {
|
||||
pub client: Arc<Client>,
|
||||
pub new_tokens_client: Arc<Client>,
|
||||
|
||||
pub event_tx: Mutex<mpsc::Sender<Event>>,
|
||||
|
||||
pub pump_fun: Mutex<PumpFun>,
|
||||
}
|
||||
@@ -51,20 +53,28 @@ pub trait Launchpad: Send + Sync {
|
||||
|
||||
impl Executor {
|
||||
pub async fn new(api_key: &str) -> anyhow::Result<Arc<Self>> {
|
||||
let (tx, _) = mpsc::channel(1);
|
||||
|
||||
tx.closed().await;
|
||||
|
||||
Ok(Arc::new(Self {
|
||||
client: Arc::new(Client(Mutex::new(
|
||||
client: Arc::new(Client {
|
||||
helius: Mutex::new(
|
||||
tokio_tungstenite::connect_async(format!(
|
||||
"wss://mainnet.helius-rpc.com/?api-key={api_key}"
|
||||
))
|
||||
.await?
|
||||
.0,
|
||||
))),
|
||||
),
|
||||
|
||||
new_tokens_client: Arc::new(Client(Mutex::new(
|
||||
solana: Mutex::new(
|
||||
tokio_tungstenite::connect_async(format!("wss://api.mainnet-beta.solana.com"))
|
||||
.await?
|
||||
.0,
|
||||
))),
|
||||
),
|
||||
}),
|
||||
|
||||
event_tx: Mutex::new(tx),
|
||||
|
||||
pump_fun: Mutex::new(PumpFun::new()),
|
||||
}))
|
||||
@@ -116,8 +126,18 @@ impl Executor {
|
||||
pub async fn listen(self: &Arc<Self>) -> anyhow::Result<mpsc::Receiver<Event>> {
|
||||
let (tx, rx) = mpsc::channel(100);
|
||||
|
||||
tokio::spawn(PumpFun::listen(self.new_tokens_client.clone(), tx));
|
||||
*self.event_tx.lock().await = tx.clone();
|
||||
|
||||
tokio::spawn(PumpFun::listen(self.client.clone(), tx));
|
||||
|
||||
Ok(rx)
|
||||
}
|
||||
|
||||
pub async fn subscribe(&self, mint: &str) -> anyhow::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn unsubscribe(&self, mint: &str) -> anyhow::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -57,7 +57,7 @@ impl Launchpad for PumpFun {
|
||||
}
|
||||
|
||||
async fn listen(client: Arc<Client>, tx: mpsc::Sender<Event>) -> anyhow::Result<()> {
|
||||
let mut ws = client.0.lock().await;
|
||||
let mut ws = client.solana.lock().await;
|
||||
|
||||
ws.send(
|
||||
serde_json::json!({
|
||||
@@ -118,13 +118,6 @@ impl Launchpad for PumpFun {
|
||||
market_cap_sol: event.market_cap_sol,
|
||||
};
|
||||
|
||||
log::info!(
|
||||
"[PUMP.FUN] New token: {} (${}) - {}",
|
||||
token.name,
|
||||
token.symbol,
|
||||
token.mint
|
||||
);
|
||||
|
||||
if tx.send(Event::NewToken(token)).await.is_err() {
|
||||
break;
|
||||
}
|
||||
|
||||
@@ -64,7 +64,7 @@ impl MomentumVelocityStrategy {
|
||||
/// Internal helper to safely handle unsubscribing and cleaning state
|
||||
async fn cleanup_and_unsubscribe(&mut self, bot: &Arc<Bot>, mint: &str) -> anyhow::Result<()> {
|
||||
debug!("[{}] Cleaning up state and unsubscribing", mint);
|
||||
if let Err(e) = bot.unsubscribe(mint).await {
|
||||
if let Err(e) = bot.executor.unsubscribe(mint).await {
|
||||
warn!("[{}] Unsubscribe request failed: {:?}", mint, e);
|
||||
}
|
||||
self.trackers.remove(mint);
|
||||
@@ -103,7 +103,7 @@ impl Strategy for MomentumVelocityStrategy {
|
||||
"[{}] Capacity reached ({}/{}). Evicting un-bought token from queue.",
|
||||
mint_to_remove, MAX_SUBSCRIBED_TOKENS, MAX_SUBSCRIBED_TOKENS
|
||||
);
|
||||
let _ = bot.unsubscribe(&mint_to_remove).await;
|
||||
let _ = bot.executor.unsubscribe(&mint_to_remove).await;
|
||||
self.trackers.remove(&mint_to_remove);
|
||||
}
|
||||
} else {
|
||||
@@ -116,7 +116,7 @@ impl Strategy for MomentumVelocityStrategy {
|
||||
}
|
||||
|
||||
info!("[{}] Subscribing and creating tracker.", token.mint);
|
||||
bot.subscribe(&token.mint).await?;
|
||||
bot.executor.subscribe(&token.mint).await?;
|
||||
self.active_subscriptions.push_back(token.mint.clone());
|
||||
|
||||
self.trackers.insert(
|
||||
|
||||
Reference in New Issue
Block a user