From 225aa6830cdb7ccabf42fbe7c4278c29bac0ee4d Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 6 Aug 2026 18:24:17 +0200 Subject: [PATCH] Subscribe and unsubscribe --- src/bot.rs | 67 +++++++++++++++++++++++++++++++++++++++++++++++++++- src/types.rs | 48 +++++++++++++++++++++++++++++++++++++ 2 files changed, 114 insertions(+), 1 deletion(-) diff --git a/src/bot.rs b/src/bot.rs index e9f81af..c8b97ed 100644 --- a/src/bot.rs +++ b/src/bot.rs @@ -1,19 +1,83 @@ +use std::collections::HashMap; + use anyhow::Context; use futures_util::{SinkExt, StreamExt}; use serde_json::json; use tokio::net::TcpStream; use tokio_tungstenite::{MaybeTlsStream, WebSocketStream, connect_async}; -use crate::{account::AccountManager, executor::Executor, types::NewToken}; +use crate::{ + account::AccountManager, + executor::Executor, + types::{Mode, NewToken, Token, Trade}, +}; pub struct Bot { pub ws: WebSocketStream>, pub accounts: AccountManager, pub executor: Box, + pub tokens: HashMap, } impl Bot { pub async fn on_new_coin(&mut self, token: NewToken) -> anyhow::Result<()> { + self.tokens.insert( + token.mint.clone(), + Token { + mode: Mode::Observing, + execute_next: false, + }, + ); + + self.subscribe(&token.mint).await?; + + Ok(()) + } + + pub async fn on_trade(&mut self, trade: Trade) -> anyhow::Result<()> { + let Some(token) = self.tokens.get_mut(&trade.mint) else { + log::error!("Token not found on trade: {:?}", trade.mint); + return Ok(()); + }; + + match &token.mode { + Mode::Observing => {} + + Mode::WaitingForEntry => {} + + Mode::WaitingForExit => {} + } + + Ok(()) + } + + pub async fn subscribe(&mut self, mint: &str) -> anyhow::Result<()> { + self.ws + .send(tokio_tungstenite::tungstenite::Message::Text( + json!({ + "method": "subscribeTokenTrade", + "keys": [mint] + }) + .to_string() + .into(), + )) + .await?; + + Ok(()) + } + + pub async fn unsubscribe(&mut self, mint: &str) -> anyhow::Result<()> { + self.ws + .send(tokio_tungstenite::tungstenite::Message::Text( + json!({ + "method": "unsubscribeTokenTrade", + "keys": [mint] + }) + .to_string() + .into(), + )) + .await?; + Ok(()) } } @@ -32,6 +96,7 @@ impl Bot { ws: connect_async("wss://pumpdev.io/ws").await?.0, executor: Box::new(account.executor()), accounts, + tokens: HashMap::new(), }) } diff --git a/src/types.rs b/src/types.rs index d97e303..28d218f 100644 --- a/src/types.rs +++ b/src/types.rs @@ -70,3 +70,51 @@ pub struct NewToken { #[serde(rename = "solAmount")] pub sol_amount: f64, } + +pub enum Mode { + Observing, + WaitingForEntry, + WaitingForExit, +} + +pub struct Token { + pub mode: Mode, + pub execute_next: bool, +} + +#[derive(Debug, Clone, Deserialize)] +pub enum TradeType { + #[serde(rename = "buy")] + Buy, + #[serde(rename = "sell")] + Sell, +} + +#[derive(Debug, Clone, Deserialize)] +pub struct Trade { + pub signature: String, + pub mint: String, + pub trader: String, + + #[serde(rename = "txType")] + pub tx_type: TradeType, + + /// Value in lamports + #[serde(rename = "solAmount")] + pub sol_amount: f64, + + /// Value in raw token units + #[serde(rename = "tokenAmount")] + pub token_amount: f64, + + #[serde(rename = "marketCapSol")] + pub market_cap_sol: f64, + + #[serde(rename = "vTokensInBondingCurve")] + pub v_tokens_in_bonding_curve: f64, + + #[serde(rename = "vSolInBondingCurve")] + pub v_sol_in_bonding_curve: f64, + + pub pool: String, +}