Subscribe and unsubscribe

This commit is contained in:
2026-08-06 18:24:17 +02:00
parent 851ad550c1
commit 225aa6830c
2 changed files with 114 additions and 1 deletions
+66 -1
View File
@@ -1,19 +1,83 @@
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;
use tokio_tungstenite::{MaybeTlsStream, WebSocketStream, connect_async}; 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 struct Bot {
pub ws: WebSocketStream<MaybeTlsStream<TcpStream>>, pub ws: WebSocketStream<MaybeTlsStream<TcpStream>>,
pub accounts: AccountManager, pub accounts: AccountManager,
pub executor: Box<dyn Executor>, pub executor: Box<dyn Executor>,
pub tokens: 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(&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(()) Ok(())
} }
} }
@@ -32,6 +96,7 @@ impl Bot {
ws: connect_async("wss://pumpdev.io/ws").await?.0, ws: connect_async("wss://pumpdev.io/ws").await?.0,
executor: Box::new(account.executor()), executor: Box::new(account.executor()),
accounts, accounts,
tokens: HashMap::new(),
}) })
} }
+48
View File
@@ -70,3 +70,51 @@ pub struct NewToken {
#[serde(rename = "solAmount")] #[serde(rename = "solAmount")]
pub sol_amount: f64, 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,
}