Working trade scanner

This commit is contained in:
2026-08-06 21:33:21 +02:00
parent 3e280908fe
commit 21b5aaa3b9
4 changed files with 50 additions and 35 deletions
+41 -25
View File
@@ -3,7 +3,7 @@ use std::collections::HashMap;
use anyhow::Context; use anyhow::Context;
use futures_util::{SinkExt, StreamExt}; use futures_util::{SinkExt, StreamExt};
use serde_json::json; use serde_json::json;
use tokio::net::TcpStream; use tokio::{net::TcpStream, sync::Mutex};
use tokio_tungstenite::{MaybeTlsStream, WebSocketStream, connect_async}; use tokio_tungstenite::{MaybeTlsStream, WebSocketStream, connect_async};
use crate::{ use crate::{
@@ -13,15 +13,15 @@ use crate::{
}; };
pub struct Bot { pub struct Bot {
pub ws: WebSocketStream<MaybeTlsStream<TcpStream>>, pub ws: Mutex<WebSocketStream<MaybeTlsStream<TcpStream>>>,
pub accounts: AccountManager, pub accounts: Mutex<AccountManager>,
pub executor: Box<dyn Executor>, pub executor: Mutex<Box<dyn Executor>>,
pub tokens: HashMap<String, Token>, pub tokens: Mutex<HashMap<String, Token>>,
} }
impl Bot { impl Bot {
pub async fn on_new_coin(&mut self, token: NewToken) -> anyhow::Result<()> { pub async fn on_new_coin(&self, token: NewToken) -> anyhow::Result<()> {
self.tokens.insert( self.tokens.lock().await.insert(
token.mint.clone(), token.mint.clone(),
Token { Token {
mode: Mode::Observing, mode: Mode::Observing,
@@ -34,8 +34,12 @@ impl Bot {
Ok(()) Ok(())
} }
pub async fn on_trade(&mut self, trade: Trade) -> anyhow::Result<()> { pub async fn on_trade(&self, trade: Trade) -> anyhow::Result<()> {
let Some(token) = self.tokens.get_mut(&trade.mint) else { log::info!("Trade: {trade:?}");
let mut tokens = self.tokens.lock().await;
let Some(token) = tokens.get_mut(&trade.mint) else {
log::error!("Token not found on trade: {:?}", trade.mint); log::error!("Token not found on trade: {:?}", trade.mint);
return Ok(()); return Ok(());
}; };
@@ -51,8 +55,10 @@ impl Bot {
Ok(()) Ok(())
} }
pub async fn subscribe(&mut self, mint: &str) -> anyhow::Result<()> { pub async fn subscribe(&self, mint: &str) -> anyhow::Result<()> {
self.ws self.ws
.lock()
.await
.send(tokio_tungstenite::tungstenite::Message::Text( .send(tokio_tungstenite::tungstenite::Message::Text(
json!({ json!({
"method": "subscribeTokenTrade", "method": "subscribeTokenTrade",
@@ -66,8 +72,10 @@ impl Bot {
Ok(()) Ok(())
} }
pub async fn unsubscribe(&mut self, mint: &str) -> anyhow::Result<()> { pub async fn unsubscribe(&self, mint: &str) -> anyhow::Result<()> {
self.ws self.ws
.lock()
.await
.send(tokio_tungstenite::tungstenite::Message::Text( .send(tokio_tungstenite::tungstenite::Message::Text(
json!({ json!({
"method": "unsubscribeTokenTrade", "method": "unsubscribeTokenTrade",
@@ -93,28 +101,31 @@ impl Bot {
.clone(); .clone();
Ok(Self { Ok(Self {
ws: connect_async("wss://pumpdev.io/ws").await?.0, ws: Mutex::new(connect_async("wss://pumpdev.io/ws").await?.0),
executor: Box::new(account.executor()), executor: Mutex::new(Box::new(account.executor())),
accounts, accounts: Mutex::new(accounts),
tokens: HashMap::new(), tokens: Mutex::new(HashMap::new()),
}) })
} }
pub async fn refresh_account(&mut self) -> anyhow::Result<()> { pub async fn refresh_account(&self) -> anyhow::Result<()> {
let account = self let accounts = self.accounts.lock().await;
let account = accounts
.accounts .accounts
.accounts .get(&accounts.active)
.get(&self.accounts.active)
.context("Failed to get account")? .context("Failed to get account")?
.clone(); .clone();
self.executor = Box::new(account.executor()); *self.executor.lock().await = Box::new(account.executor());
Ok(()) Ok(())
} }
pub async fn initialize_websocket_subscribe(&mut self) -> anyhow::Result<()> { pub async fn initialize_websocket_subscribe(&self) -> anyhow::Result<()> {
self.ws self.ws
.lock()
.await
.send(tokio_tungstenite::tungstenite::Message::Text( .send(tokio_tungstenite::tungstenite::Message::Text(
json!({ "method": "subscribeNewToken" }).to_string().into(), json!({ "method": "subscribeNewToken" }).to_string().into(),
)) ))
@@ -123,19 +134,24 @@ impl Bot {
Ok(()) Ok(())
} }
pub async fn start(&mut self) -> anyhow::Result<()> { pub async fn start(&self) -> anyhow::Result<()> {
self.initialize_websocket_subscribe().await?; self.initialize_websocket_subscribe().await?;
while let Some(msg) = self.ws.next().await { let mut ws = self.ws.lock().await;
let msg = msg?;
while let Some(msg) = ws.next().await.transpose()? {
if let tokio_tungstenite::tungstenite::Message::Text(text) = msg { if let tokio_tungstenite::tungstenite::Message::Text(text) = msg {
match serde_json::from_str::<crate::types::PumpDevEvent>(&text) { match serde_json::from_str::<crate::types::PumpDevEvent>(&text) {
Ok(crate::types::PumpDevEvent::Create(token)) => { Ok(crate::types::PumpDevEvent::Create(token)) => {
drop(ws);
self.on_new_coin(token).await?; self.on_new_coin(token).await?;
ws = self.ws.lock().await;
} }
Ok(crate::types::PumpDevEvent::Trade(trade)) => { Ok(crate::types::PumpDevEvent::Trade(trade)) => {
drop(ws);
self.on_trade(trade).await?; self.on_trade(trade).await?;
ws = self.ws.lock().await;
} }
Ok(event) => { Ok(event) => {
@@ -143,7 +159,7 @@ impl Bot {
} }
Err(err) => { Err(err) => {
log::error!("{err}"); log::error!("{err}, MSG -> {text}");
} }
} }
} }
+1 -1
View File
@@ -6,7 +6,7 @@ use crate::account::Account;
#[allow(unused_variables)] #[allow(unused_variables)]
#[async_trait::async_trait] #[async_trait::async_trait]
pub trait Executor { pub trait Executor: Send + Sync {
async fn buy( async fn buy(
&self, &self,
mint: String, mint: String,
+1 -1
View File
@@ -11,7 +11,7 @@ async fn main() -> anyhow::Result<()> {
builder.filter_level(log::LevelFilter::Info); builder.filter_level(log::LevelFilter::Info);
builder.init(); builder.init();
let mut bot = Bot::new().await?; let bot = Bot::new().await?;
bot.start().await bot.start().await
} }
+7 -8
View File
@@ -16,11 +16,14 @@ impl<'de> Deserialize<'de> for PumpDevEvent {
D: Deserializer<'de>, D: Deserializer<'de>,
{ {
let value = Value::deserialize(deserializer)?; let value = Value::deserialize(deserializer)?;
if value.get("txType").is_some() && value.get("name").is_none() {
let trade: Trade = serde_json::from_value(value).map_err(serde::de::Error::custom)?;
return Ok(PumpDevEvent::Trade(trade));
}
if value.get("txType").is_some() { if value.get("name").is_some() {
let token: NewToken = let token: NewToken =
serde_json::from_value(value).map_err(serde::de::Error::custom)?; serde_json::from_value(value).map_err(serde::de::Error::custom)?;
return Ok(PumpDevEvent::Create(token)); return Ok(PumpDevEvent::Create(token));
} }
@@ -33,10 +36,8 @@ impl<'de> Deserialize<'de> for PumpDevEvent {
client_id: u64, client_id: u64,
message: String, message: String,
}, },
#[serde(rename = "connectionStatus")] #[serde(rename = "connectionStatus")]
ConnectionStatus { connected: bool, timestamp: u64 }, ConnectionStatus { connected: bool, timestamp: u64 },
#[serde(rename = "subscribed")] #[serde(rename = "subscribed")]
Subscribed { method: String }, Subscribed { method: String },
} }
@@ -95,16 +96,16 @@ pub enum TradeType {
pub struct Trade { pub struct Trade {
pub signature: String, pub signature: String,
pub mint: String, pub mint: String,
#[serde(rename = "traderPublicKey")]
pub trader: String, pub trader: String,
#[serde(rename = "txType")] #[serde(rename = "txType")]
pub tx_type: TradeType, pub tx_type: TradeType,
/// Value in lamports
#[serde(rename = "solAmount")] #[serde(rename = "solAmount")]
pub sol_amount: f64, pub sol_amount: f64,
/// Value in raw token units
#[serde(rename = "tokenAmount")] #[serde(rename = "tokenAmount")]
pub token_amount: f64, pub token_amount: f64,
@@ -116,6 +117,4 @@ pub struct Trade {
#[serde(rename = "vSolInBondingCurve")] #[serde(rename = "vSolInBondingCurve")]
pub v_sol_in_bonding_curve: f64, pub v_sol_in_bonding_curve: f64,
pub pool: String,
} }