Merge pull request #4 from selimaj-dev/using-rpc

Using rpc
This commit is contained in:
2026-08-08 02:12:30 +02:00
committed by GitHub
14 changed files with 5226 additions and 795 deletions
Generated
+4492 -165
View File
File diff suppressed because it is too large Load Diff
+9 -6
View File
@@ -4,18 +4,21 @@ version = "0.1.0"
edition = "2024"
[dependencies]
anyhow = "1.0.104"
tokio = { version = "1.53.1", features = ["macros", "rt-multi-thread", "sync", "fs", "io-std"] }
tokio-tungstenite = { version = "0.30.0", features = ["native-tls"] }
futures-util = "0.3.33"
anyhow = "1.0.104"
serde = { version = "1.0.229", features = ["derive", "serde_derive"] }
serde_json = "1.0.151"
base64 = "0.22.1"
bs58 = "0.5.1"
csv = "1.4.0"
tokio = { version = "1.53.1", features = ["macros", "rt-multi-thread", "sync", "fs", "io-std"] }
tokio-tungstenite = { version = "0.30.0", features = ["native-tls"] }
log = "0.4.33"
env_logger = "0.11.11"
rust_decimal = { version = "1.42.1", features = ["macros"] }
reqwest = { version = "0.13.4", features = ["json"] }
async-trait = "0.1.91"
csv = "1.4.0"
rust_decimal = { version = "1.42.1", features = ["macros"] }
chrono = { version = "0.4.45", features = ["serde"] }
helius = "1.1.0"
+44 -111
View File
@@ -1,65 +1,27 @@
use std::{collections::HashMap, sync::Arc};
use std::sync::Arc;
use anyhow::Context;
use futures_util::{SinkExt, StreamExt};
use serde_json::json;
use tokio::{
net::TcpStream,
sync::{Mutex, watch},
};
use tokio_tungstenite::{MaybeTlsStream, WebSocketStream, connect_async};
use tokio::sync::{Mutex, mpsc, watch};
use crate::{
account::AccountManager,
executor::ExecutorWrapper,
strategy::{Strategy, veloc::MomentumVelocityStrategy},
data::{
Event,
account::{Account, AccountManager},
tradelog::TradeLog,
},
launchpad::Executor,
strategy::{Strategy, veloc::MomentumVelocityStrategy},
};
pub struct Bot {
pub ws: Mutex<WebSocketStream<MaybeTlsStream<TcpStream>>>,
pub accounts: Mutex<AccountManager>,
pub executor: Mutex<ExecutorWrapper>,
pub executor: Arc<Executor>,
pub strategy: Mutex<Box<dyn Strategy>>,
pub account_manager: Mutex<AccountManager>,
pub current_account: Mutex<Account>,
pub trade_log: Mutex<Vec<TradeLog>>,
}
impl Bot {
pub async fn subscribe(&self, mint: &str) -> anyhow::Result<()> {
self.ws
.lock()
.await
.send(tokio_tungstenite::tungstenite::Message::Text(
json!({
"method": "subscribeTokenTrade",
"keys": [mint]
})
.to_string()
.into(),
))
.await?;
Ok(())
}
pub async fn unsubscribe(&self, mint: &str) -> anyhow::Result<()> {
self.ws
.lock()
.await
.send(tokio_tungstenite::tungstenite::Message::Text(
json!({
"method": "unsubscribeTokenTrade",
"keys": [mint]
})
.to_string()
.into(),
))
.await?;
Ok(())
}
}
impl Bot {
pub async fn new() -> anyhow::Result<Arc<Self>> {
let accounts = AccountManager::get().await?;
@@ -71,54 +33,19 @@ impl Bot {
.clone();
Ok(Arc::new(Self {
ws: Mutex::new(connect_async("wss://pumpdev.io/ws").await?.0),
executor: Mutex::new(ExecutorWrapper {
executor: Box::new(account.executor()),
positions: HashMap::new(),
}),
accounts: Mutex::new(accounts),
executor: Executor::new(&accounts.api_key).await?,
strategy: Mutex::new(Box::new(MomentumVelocityStrategy::new())),
account_manager: Mutex::new(accounts),
current_account: Mutex::new(account),
trade_log: Mutex::new(Vec::new()),
}))
}
pub async fn refresh_account(self: &Arc<Self>) -> anyhow::Result<()> {
self.strategy
.lock()
.await
.execute_sell_all(self.clone())
.await?;
let accounts = self.accounts.lock().await;
let account = accounts
.accounts
.get(&accounts.active)
.context("Failed to get account")?
.clone();
self.executor.lock().await.executor = Box::new(account.executor());
Ok(())
}
pub async fn initialize_websocket_subscribe(&self) -> anyhow::Result<()> {
self.ws
.lock()
.await
.send(tokio_tungstenite::tungstenite::Message::Text(
json!({ "method": "subscribeNewToken" }).to_string().into(),
))
.await?;
Ok(())
}
pub async fn start(
self: &Arc<Self>,
mut shutdown: watch::Receiver<bool>,
) -> anyhow::Result<()> {
self.initialize_websocket_subscribe().await?;
let mut rx = self.executor.listen().await?;
loop {
tokio::select! {
@@ -138,7 +65,7 @@ impl Bot {
}
}
result = self.tick() => {
result = self.tick(&mut rx) => {
if result? {
log::warn!("Websocket closed.");
break;
@@ -150,15 +77,10 @@ impl Bot {
Ok(())
}
pub async fn tick(self: &Arc<Self>) -> anyhow::Result<bool> {
let mut ws = self.ws.lock().await;
if let Some(msg) = ws.next().await.transpose()? {
if let tokio_tungstenite::tungstenite::Message::Text(text) = msg {
match serde_json::from_str::<crate::types::PumpDevEvent>(&text) {
Ok(crate::types::PumpDevEvent::Create(token)) => {
drop(ws);
pub async fn tick(self: &Arc<Self>, rx: &mut mpsc::Receiver<Event>) -> anyhow::Result<bool> {
if let Some(event) = rx.recv().await {
match event {
Event::NewToken(token) => {
self.strategy
.lock()
.await
@@ -166,24 +88,13 @@ impl Bot {
.await?;
}
Ok(crate::types::PumpDevEvent::Trade(trade)) => {
drop(ws);
Event::Trade(trade) => {
self.strategy
.lock()
.await
.on_trade(self.clone(), trade)
.await?;
}
Ok(_event) => {
// println!("{:?}", event);
}
Err(err) => {
log::error!("{err}, MSG -> {text}");
}
}
}
Ok(false)
@@ -192,3 +103,25 @@ impl Bot {
}
}
}
impl Bot {
pub async fn refresh_account(self: &Arc<Self>) -> anyhow::Result<()> {
self.strategy
.lock()
.await
.execute_sell_all(self.clone())
.await?;
let accounts = self.account_manager.lock().await;
let account = accounts
.accounts
.get(&accounts.active)
.context("Failed to get account")?
.clone();
*self.current_account.lock().await = account;
Ok(())
}
}
+18 -36
View File
@@ -1,7 +1,6 @@
use std::{collections::HashMap, path::PathBuf};
use crate::executor::pump_fun::PumpDevAccount;
use anyhow::{Context, anyhow};
use anyhow::Context;
use serde::{Deserialize, Serialize};
pub fn get_accounts_path() -> anyhow::Result<PathBuf> {
@@ -13,47 +12,21 @@ pub fn get_accounts_path() -> anyhow::Result<PathBuf> {
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type")]
pub enum Account {
PumpDev(PumpDevAccount),
pub struct Account {
#[serde(rename = "publicKey")]
pub public_key: String,
#[serde(rename = "privateKey")]
pub private_key: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AccountManager {
#[serde(rename = "apiKey")]
pub api_key: String,
pub active: String,
pub accounts: HashMap<String, Account>,
}
impl Account {
pub async fn new() -> anyhow::Result<Self> {
let response = reqwest::Client::new()
.post("https://pumpdev.io/api/wallet/create")
.json(&serde_json::json!({}))
.send()
.await
.context("Failed to create PumpDev wallet")?;
if !response.status().is_success() {
let status = response.status();
let body: serde_json::Value = response.json().await.unwrap_or_default();
let message = body
.get("error")
.and_then(serde_json::Value::as_str)
.unwrap_or("Unknown error");
return Err(anyhow!(
"Failed to create PumpDev wallet ({status}): {message}"
));
}
let account: PumpDevAccount = response
.json()
.await
.context("Failed to parse PumpDev wallet")?;
Ok(Self::PumpDev(account))
}
}
impl AccountManager {
pub async fn get() -> anyhow::Result<Self> {
let path = get_accounts_path()?;
@@ -64,6 +37,8 @@ impl AccountManager {
tokio::fs::create_dir_all(path.parent().unwrap()).await?;
tokio::fs::write(&path, default.to_string()?).await?;
println!("Please open {path:?} and edit the data accordingly");
return Ok(default);
}
@@ -74,10 +49,17 @@ impl AccountManager {
pub async fn new() -> anyhow::Result<Self> {
let mut accounts = HashMap::new();
accounts.insert("default".to_string(), Account::new().await?);
accounts.insert(
"default".to_string(),
Account {
public_key: String::new(),
private_key: String::new(),
},
);
Ok(Self {
active: "default".to_string(),
api_key: String::new(),
accounts,
})
}
+46
View File
@@ -0,0 +1,46 @@
pub mod account;
pub mod tradelog;
use serde::Deserialize;
#[derive(Debug, Clone)]
pub enum Event {
NewToken(NewToken),
Trade(Trade),
}
#[derive(Debug, Clone, Deserialize)]
pub struct NewToken {
pub mint: String,
#[serde(rename = "traderPublicKey")]
pub trader_public_key: String,
pub name: String,
pub symbol: String,
pub uri: String,
#[serde(rename = "marketCapSol")]
pub market_cap_sol: f64,
}
#[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,
pub tx_type: TradeType,
pub sol_amount: f64,
pub token_amount: f64,
pub market_cap_sol: f64,
pub v_tokens_in_bonding_curve: f64,
pub v_sol_in_bonding_curve: f64,
}
-132
View File
@@ -1,132 +0,0 @@
pub mod pump_fun;
use std::collections::HashMap;
use rust_decimal::Decimal;
use crate::account::Account;
pub struct ExecutorWrapper {
pub executor: Box<dyn Executor>,
pub positions: HashMap<String, Decimal>,
}
#[allow(unused_variables)]
#[async_trait::async_trait]
pub trait Executor: Send + Sync {
async fn buy(
&self,
mint: &str,
amount: Decimal,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
Ok(())
}
async fn sell(
&self,
mint: &str,
amount: Decimal,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
Ok(())
}
async fn sell_percent(
&self,
mint: &str,
amount: u8,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
Ok(())
}
}
impl Account {
pub fn executor(self) -> impl Executor {
match self {
Self::PumpDev(account) => pump_fun::PumpDev::new(account),
}
}
}
impl ExecutorWrapper {
pub async fn buy(
&mut self,
mint: &str,
amount: Decimal,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
self.executor.buy(mint, amount, priority, slippage).await?;
let position = self
.positions
.entry(mint.to_string())
.or_insert(Decimal::ZERO);
*position += amount;
Ok(())
}
pub async fn sell(
&mut self,
mint: &str,
amount: Decimal,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
self.executor.sell(mint, amount, priority, slippage).await?;
if let Some(position) = self.positions.get_mut(mint) {
*position -= amount;
if *position <= Decimal::ZERO {
self.positions.remove(mint);
}
}
Ok(())
}
pub async fn sell_percent(
&mut self,
mint: &str,
amount: u8,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
let sell_amount = match self.positions.get(mint) {
Some(position) => *position * Decimal::from(amount) / Decimal::from(100),
None => return Ok(()),
};
self.executor
.sell_percent(mint, amount, priority, slippage)
.await?;
if let Some(position) = self.positions.get_mut(mint) {
*position -= sell_amount;
if *position <= Decimal::ZERO {
self.positions.remove(mint);
}
}
Ok(())
}
pub async fn sell_all(&mut self, priority: Decimal, slippage: u16) -> anyhow::Result<()> {
for (mint, _) in self.positions.drain() {
self.executor
.sell_percent(&mint, 100, priority, slippage)
.await?;
}
Ok(())
}
}
-148
View File
@@ -1,148 +0,0 @@
use anyhow::{Context, anyhow};
use reqwest::Client;
use rust_decimal::Decimal;
use serde::{Deserialize, Serialize};
use tokio::sync::Mutex;
use crate::executor::Executor;
const TRADE_LIGHTNING_URL: &str = "https://pumpdev.io/api/trade-lightning";
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PumpDevAccount {
#[serde(rename = "apiKey")]
pub api_key: String,
#[serde(rename = "publicKey")]
pub public_key: String,
#[serde(rename = "privateKey")]
pub private_key: String,
}
#[derive(Debug, Deserialize)]
struct TradeResponse {
signature: String,
}
pub struct PumpDev {
pub account: Mutex<PumpDevAccount>,
pub client: Client,
}
impl PumpDev {
pub fn new(account: PumpDevAccount) -> Self {
Self {
account: Mutex::new(account),
client: Client::new(),
}
}
async fn trade(
&self,
action: &str,
mint: &str,
amount: String,
priority: Decimal,
slippage: u16,
denominated_in_sol: bool,
) -> anyhow::Result<()> {
let account = self.account.lock().await;
let response = self
.client
.post(format!("{TRADE_LIGHTNING_URL}?api-key={}", account.api_key))
.json(&serde_json::json!({
"action": action,
"mint": mint,
"amount": amount,
"denominatedInSol": if denominated_in_sol { "true" } else { "false" },
"slippage": slippage,
"priorityFee": priority,
}))
.send()
.await
.context("Failed to send trade request")?;
let status = response.status();
let body = response.text().await?;
if !status.is_success() {
let error = serde_json::from_str::<serde_json::Value>(&body)
.ok()
.and_then(|value| value.get("error").cloned())
.and_then(|value| value.as_str().map(str::to_string))
.unwrap_or_else(|| body.trim().to_string());
return Err(anyhow!("Trade failed ({status}): {error}"));
}
let data: TradeResponse =
serde_json::from_str(&body).context("Failed to parse trade response")?;
log::info!("{action} executed: {}", data.signature);
Ok(())
}
}
#[async_trait::async_trait]
impl Executor for PumpDev {
async fn buy(
&self,
mint: &str,
amount: Decimal,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
log::info!("BUY {mint} {amount} SOL");
Ok(())
// self.trade(
// "buy",
// mint,
// amount.round_dp(3).to_string(),
// priority,
// slippage,
// true,
// )
// .await
}
async fn sell(
&self,
mint: &str,
amount: Decimal,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
log::info!("SELL {mint} {amount} SOL");
Ok(())
// self.trade(
// "sell",
// mint,
// amount.round_dp(3).to_string(),
// priority,
// slippage,
// false,
// )
// .await
}
async fn sell_percent(
&self,
mint: &str,
amount: u8,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
log::info!("SELL {mint} {amount}%");
Ok(())
// self.trade(
// "sell",
// mint,
// format!("{amount}%"),
// priority,
// slippage,
// false,
// )
// .await
}
}
+149
View File
@@ -0,0 +1,149 @@
mod pump_fun;
use rust_decimal::Decimal;
use std::{
collections::{HashMap, HashSet},
sync::Arc,
};
use tokio::sync::{Mutex, mpsc};
use crate::{data::Event, launchpad::pump_fun::PumpFun};
type ClientMutex = Mutex<
tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>,
>;
pub struct Client {
pub helius: ClientMutex,
pub solana: ClientMutex,
pub subscribed: Mutex<HashSet<String>>,
}
pub struct Executor {
pub client: Arc<Client>,
pub event_tx: Mutex<mpsc::Sender<Event>>,
pub pump_fun: Mutex<PumpFun>,
}
#[allow(unused_variables)]
#[async_trait::async_trait]
pub trait Launchpad: Send + Sync {
async fn buy(
&mut self,
client: &Client,
mint: &str,
amount: Decimal,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()>;
async fn sell(
&mut self,
client: &Client,
mint: &str,
amount: u8,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()>;
async fn listen(client: Arc<Client>, tx: mpsc::Sender<Event>) -> anyhow::Result<()>;
fn get_positions(&self) -> HashMap<String, Decimal>;
}
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 {
helius: Mutex::new(
tokio_tungstenite::connect_async(format!(
"wss://mainnet.helius-rpc.com/?api-key={api_key}"
))
.await?
.0,
),
solana: Mutex::new(
tokio_tungstenite::connect_async(format!("wss://api.mainnet-beta.solana.com"))
.await?
.0,
),
subscribed: Mutex::new(HashSet::new()),
}),
event_tx: Mutex::new(tx),
pump_fun: Mutex::new(PumpFun::new()),
}))
}
pub async fn buy(
self: &Arc<Self>,
mint: &str,
amount: Decimal,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
self.pump_fun
.lock()
.await
.buy(&self.client, mint, amount, priority, slippage)
.await
}
pub async fn sell(
self: &Arc<Self>,
mint: &str,
amount: u8,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
self.pump_fun
.lock()
.await
.sell(&self.client, mint, amount, priority, slippage)
.await
}
pub async fn sell_all(
self: &Arc<Self>,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
let mut pump = self.pump_fun.lock().await;
for (mint, _) in pump.get_positions() {
pump.sell(&self.client, &mint, 100, priority, slippage)
.await?;
}
Ok(())
}
pub async fn listen(self: &Arc<Self>) -> anyhow::Result<mpsc::Receiver<Event>> {
let (tx, rx) = mpsc::channel(100);
*self.event_tx.lock().await = tx.clone();
tokio::spawn(PumpFun::listen(self.client.clone(), tx));
Ok(rx)
}
pub async fn subscribe(&self, mint: &str) {
self.client.subscribed.lock().await.insert(mint.to_string());
}
pub async fn unsubscribe(&self, mint: &str) {
self.client.subscribed.lock().await.remove(mint);
}
}
+299
View File
@@ -0,0 +1,299 @@
use std::collections::HashMap;
use std::sync::Arc;
use base64::{Engine, engine::general_purpose::STANDARD};
use futures_util::SinkExt;
use futures_util::StreamExt;
use rust_decimal::Decimal;
use serde::Deserialize;
use tokio::sync::mpsc;
use tokio_tungstenite::tungstenite::Message;
use crate::data::Event;
use crate::data::NewToken;
use crate::data::Trade;
use crate::data::TradeType;
use crate::launchpad::Client;
use crate::launchpad::Launchpad;
const PUMP_FUN_PROGRAM: &str = "6EF8rrecthR5Dkzon8Nwu78hRvfCKubJ14M5uBEwF6P";
const CREATE_EVENT_DISCRIMINATOR: [u8; 8] = [27, 114, 169, 77, 222, 235, 99, 118];
const TRADE_EVENT_DISCRIMINATOR: [u8; 8] = [189, 219, 127, 211, 78, 230, 97, 238];
pub struct PumpFun {
pub positions: HashMap<String, Decimal>,
}
impl PumpFun {
pub fn new() -> Self {
Self {
positions: HashMap::new(),
}
}
}
#[allow(unused_variables)]
#[async_trait::async_trait]
impl Launchpad for PumpFun {
async fn buy(
&mut self,
client: &Client,
mint: &str,
amount: Decimal,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
Ok(())
}
async fn sell(
&mut self,
client: &Client,
mint: &str,
amount: u8,
priority: Decimal,
slippage: u16,
) -> anyhow::Result<()> {
Ok(())
}
async fn listen(client: Arc<Client>, tx: mpsc::Sender<Event>) -> anyhow::Result<()> {
let mut ws = client.solana.lock().await;
ws.send(
serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "logsSubscribe",
"params": [
{ "mentions": [PUMP_FUN_PROGRAM] },
{ "commitment": "processed" }
]
})
.to_string()
.into(),
)
.await?;
log::info!(
"[PUMP.FUN] Listening for new tokens from {}",
PUMP_FUN_PROGRAM
);
while let Some(message) = ws.next().await {
let Message::Text(text) = message? else {
continue;
};
let Ok(notification) = serde_json::from_str::<LogsNotification>(&text) else {
continue;
};
if notification.method != "logsNotification" {
continue;
}
let logs = notification.params.result.value.logs;
if !logs
.iter()
.any(|log| log.starts_with("Program log: Instruction: Create"))
{
continue;
}
for log in &logs {
let Some(encoded) = log.strip_prefix("Program data: ") else {
continue;
};
let Ok(data) = STANDARD.decode(encoded) else {
continue;
};
if let Some(token) = parse_create_event(&data) {
tx.send(Event::NewToken(token)).await?;
}
if let Some(trade) =
parse_trade_event(&data, notification.params.result.value.signature.clone())
{
if client.subscribed.lock().await.contains(&trade.mint) {
tx.send(Event::Trade(trade)).await?;
}
}
}
}
Ok(())
}
fn get_positions<'a>(&'a self) -> HashMap<String, Decimal> {
self.positions.clone()
}
}
fn parse_trade_event(data: &[u8], signature: String) -> Option<Trade> {
if data.len() < 8 || data[..8] != TRADE_EVENT_DISCRIMINATOR {
return None;
}
let mut offset = 8;
let mint = decode_pubkey(data, &mut offset)?;
let sol_amount = decode_u64(data, &mut offset)? as f64;
let token_amount = decode_u64(data, &mut offset)? as f64;
let is_buy = decode_bool(data, &mut offset)?;
let trader = decode_pubkey(data, &mut offset)?;
// timestamp exists but we don't need it currently
let _timestamp = decode_i64(data, &mut offset)?;
let v_sol_in_bonding_curve = decode_u64(data, &mut offset)? as f64;
let v_tokens_in_bonding_curve = decode_u64(data, &mut offset)? as f64;
// These exist in the event but aren't needed for your struct
let _real_sol_reserves = decode_u64(data, &mut offset)?;
let _real_token_reserves = decode_u64(data, &mut offset)?;
Some(Trade {
signature,
mint,
trader,
tx_type: if is_buy {
TradeType::Buy
} else {
TradeType::Sell
},
sol_amount,
token_amount,
// Pump.fun bonding curve market cap:
// token_price = virtual SOL / virtual tokens
// market cap is derived
market_cap_sol: if v_tokens_in_bonding_curve > 0.0 {
v_sol_in_bonding_curve / v_tokens_in_bonding_curve
} else {
0.0
},
v_tokens_in_bonding_curve,
v_sol_in_bonding_curve,
})
}
fn parse_create_event(data: &[u8]) -> Option<NewToken> {
if data.len() < 8 || data[..8] != CREATE_EVENT_DISCRIMINATOR {
return None;
}
let mut offset = 8;
let name = decode_string(data, &mut offset)?;
let symbol = decode_string(data, &mut offset)?;
let uri = decode_string(data, &mut offset)?;
let mint = decode_pubkey(data, &mut offset)?;
let _bonding_curve = decode_pubkey(data, &mut offset)?;
let trader_public_key = decode_pubkey(data, &mut offset)?;
let _creator = decode_pubkey(data, &mut offset)?;
// Virtual reserves were added in a later program version; fall back to 0.0 for legacy events.
let market_cap_sol = decode_market_cap(data, &mut offset).unwrap_or(0.0);
Some(NewToken {
name,
symbol,
uri,
mint,
trader_public_key,
market_cap_sol,
})
}
fn decode_market_cap(data: &[u8], offset: &mut usize) -> Option<f64> {
let _timestamp = decode_i64(data, offset)?;
let virtual_token_reserves = decode_u64(data, offset)?;
let virtual_sol_reserves = decode_u64(data, offset)?;
let _real_token_reserves = decode_u64(data, offset)?;
let token_total_supply = decode_u64(data, offset)?;
if virtual_token_reserves == 0 {
return Some(0.0);
}
Some(
virtual_sol_reserves as f64 / 1e9 * token_total_supply as f64
/ virtual_token_reserves as f64,
)
}
fn decode_string(data: &[u8], offset: &mut usize) -> Option<String> {
let len = u32::from_le_bytes(data.get(*offset..*offset + 4)?.try_into().ok()?) as usize;
*offset += 4;
let bytes = data.get(*offset..*offset + len)?;
*offset += len;
String::from_utf8(bytes.to_vec()).ok()
}
fn decode_pubkey(data: &[u8], offset: &mut usize) -> Option<String> {
let bytes = data.get(*offset..*offset + 32)?;
*offset += 32;
Some(bs58::encode(bytes).into_string())
}
fn decode_u64(data: &[u8], offset: &mut usize) -> Option<u64> {
let bytes = data.get(*offset..*offset + 8)?.try_into().ok()?;
*offset += 8;
Some(u64::from_le_bytes(bytes))
}
fn decode_i64(data: &[u8], offset: &mut usize) -> Option<i64> {
let bytes = data.get(*offset..*offset + 8)?.try_into().ok()?;
*offset += 8;
Some(i64::from_le_bytes(bytes))
}
fn decode_bool(data: &[u8], offset: &mut usize) -> Option<bool> {
if *offset >= data.len() {
return None;
}
let value = data[*offset] != 0;
*offset += 1;
Some(value)
}
#[derive(Debug, Deserialize)]
struct LogsNotification {
method: String,
params: NotificationParams,
}
#[derive(Debug, Deserialize)]
struct NotificationParams {
result: NotificationResult,
}
#[derive(Debug, Deserialize)]
struct NotificationResult {
value: LogsValue,
}
#[derive(Debug, Deserialize)]
struct LogsValue {
signature: String,
logs: Vec<String>,
}
+4 -6
View File
@@ -1,9 +1,7 @@
pub mod account;
pub mod bot;
pub mod executor;
pub mod data;
pub mod launchpad;
pub mod strategy;
pub mod tradelog;
pub mod types;
use std::sync::Arc;
@@ -29,8 +27,6 @@ async fn try_shutdown_listener(bot: Arc<Bot>, tx: watch::Sender<bool>) -> anyhow
}
"save" => {
log::info!("Saving trades");
let now = Local::now();
let formatted_time = now.format("%m-%d-%H-%M").to_string();
@@ -53,6 +49,8 @@ async fn try_shutdown_listener(bot: Arc<Bot>, tx: watch::Sender<bool>) -> anyhow
_ => {}
}
line.clear();
}
}
+1 -1
View File
@@ -4,7 +4,7 @@ use std::sync::Arc;
use crate::{
bot::Bot,
types::{NewToken, Trade},
data::{NewToken, Trade},
};
#[async_trait::async_trait]
+153 -64
View File
@@ -6,17 +6,33 @@ use log::{debug, info, trace, warn};
use rust_decimal::{Decimal, dec};
use crate::bot::Bot;
use crate::data::tradelog::{ExitReason, TradeLog};
use crate::data::{NewToken, Trade, TradeType};
use crate::strategy::Strategy;
use crate::tradelog::{ExitReason, TradeLog};
use crate::types::{NewToken, Trade, TradeType};
const BUY_AMOUNT_SOL: Decimal = dec!(0.2);
const PRIORITY: Decimal = dec!(0.0002);
const SLIPPAGE: u16 = 10;
const MAX_SUBSCRIBED_TOKENS: usize = 5;
const MAX_SUBSCRIBED_TOKENS: usize = 25;
const MAX_OPEN_POSITIONS: usize = 5;
// Strict entry filters
const MIN_UNIQUE_BUYERS: usize = 0;
const MIN_NET_SOL_FLOW: f64 = 0.0003;
const MIN_TRADE_COUNT: usize = 0;
const MAX_CURVE_SOL: f64 = 800.0;
const MIN_MOMENTUM_PCT: f64 = 0.001;
// Exit rules
const TAKE_PROFIT_PCT: f64 = 0.40;
const STOP_LOSS_PCT: f64 = -0.10;
const TRAILING_STOP_DROP_PCT: f64 = 0.12;
const TRAILING_STOP_MIN_GAIN_PCT: f64 = 0.10;
const STALL_DURATION: Duration = Duration::from_secs(2);
struct TokenTracker {
created_at: Instant,
first_price_sol: Option<f64>,
unique_buyers: HashSet<String>,
net_sol_flow: f64,
trade_count: usize,
@@ -26,12 +42,11 @@ struct OpenPosition {
trade: TradeLog,
highest_price_sol: f64,
last_price_sol: f64,
last_high_time: Instant,
}
pub struct MomentumVelocityStrategy {
min_unique_buyers: usize,
min_net_sol_flow: f64,
max_tracking_duration: Duration,
trackers: HashMap<String, TokenTracker>,
@@ -41,17 +56,21 @@ pub struct MomentumVelocityStrategy {
active_subscriptions: VecDeque<String>,
}
impl MomentumVelocityStrategy {
pub fn new() -> Self {
impl Default for MomentumVelocityStrategy {
fn default() -> Self {
Self {
min_unique_buyers: 1,
min_net_sol_flow: 0.001,
max_tracking_duration: Duration::from_secs(45),
trackers: HashMap::new(),
positions: HashMap::new(),
active_subscriptions: VecDeque::with_capacity(MAX_SUBSCRIBED_TOKENS),
}
}
}
impl MomentumVelocityStrategy {
pub fn new() -> Self {
Self::default()
}
fn calculate_price_sol(&self, trade: &Trade) -> f64 {
if trade.v_tokens_in_bonding_curve == 0.0 {
@@ -63,19 +82,53 @@ 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 {
warn!("[{}] Unsubscribe request failed: {:?}", mint, e);
}
bot.executor.unsubscribe(mint).await;
// if let Err(e) = bot.executor.unsubscribe(mint).await {
// warn!("[{}] Unsubscribe request failed: {:?}", mint, e);
// }
self.trackers.remove(mint);
self.active_subscriptions.retain(|m| m != mint);
Ok(())
}
async fn execute_exit(
&mut self,
bot: &Arc<Bot>,
mint: &str,
reason: ExitReason,
) -> anyhow::Result<()> {
let Some((entry_price, exit_price)) = self
.positions
.get(mint)
.map(|pos| (pos.trade.entry_price_sol, pos.last_price_sol))
else {
return Ok(());
};
let pnl = ((exit_price - entry_price) / entry_price) * 100.0;
info!(
"[{}] EXECUTING SELL {}% {:?}",
mint, pnl, reason,
);
bot.executor.sell(mint, 100, PRIORITY, SLIPPAGE).await?;
if let Some(mut pos) = self.positions.remove(mint) {
pos.trade.close(exit_price, reason);
bot.trade_log.lock().await.push(pos.trade);
}
self.cleanup_and_unsubscribe(bot, mint).await?;
Ok(())
}
}
#[async_trait::async_trait]
impl Strategy for MomentumVelocityStrategy {
async fn execute_sell_all(&mut self, bot: Arc<Bot>) -> anyhow::Result<()> {
bot.executor.lock().await.sell_all(PRIORITY, SLIPPAGE).await
bot.executor.sell_all(PRIORITY, SLIPPAGE).await
}
async fn on_new_coin(&mut self, bot: Arc<Bot>, token: NewToken) -> anyhow::Result<()> {
@@ -98,11 +151,11 @@ impl Strategy for MomentumVelocityStrategy {
if let Some(idx) = eviction_index {
if let Some(mint_to_remove) = self.active_subscriptions.remove(idx) {
info!(
debug!(
"[{}] Capacity reached ({}/{}). Evicting un-bought token from queue.",
mint_to_remove, MAX_SUBSCRIBED_TOKENS, MAX_SUBSCRIBED_TOKENS
);
let _ = bot.unsubscribe(&mint_to_remove).await;
bot.executor.unsubscribe(&mint_to_remove).await;
self.trackers.remove(&mint_to_remove);
}
} else {
@@ -114,14 +167,15 @@ impl Strategy for MomentumVelocityStrategy {
}
}
info!("[{}] Subscribing and creating tracker.", token.mint);
bot.subscribe(&token.mint).await?;
debug!("[{}] Subscribing and creating tracker.", token.mint);
bot.executor.subscribe(&token.mint).await;
self.active_subscriptions.push_back(token.mint.clone());
self.trackers.insert(
token.mint,
TokenTracker {
created_at: Instant::now(),
first_price_sol: None,
unique_buyers: HashSet::new(),
net_sol_flow: 0.0,
trade_count: 0,
@@ -136,11 +190,13 @@ impl Strategy for MomentumVelocityStrategy {
let current_price = self.calculate_price_sol(&trade);
// -------------------------------------------------------------
// 1. Manage Active Positions (Take Profit / Stop Loss / Stall)
// 1. Price-based exits for the mint of this trade.
// Updates the highest price / stall timer before any stall scan.
// -------------------------------------------------------------
let mut exits: Vec<(String, ExitReason)> = Vec::new();
if let Some(pos) = self.positions.get_mut(mint) {
let price_change_pct =
(current_price - pos.trade.entry_price_sol) / pos.trade.entry_price_sol;
pos.last_price_sol = current_price;
if current_price > pos.highest_price_sol {
pos.highest_price_sol = current_price;
@@ -148,48 +204,47 @@ impl Strategy for MomentumVelocityStrategy {
trace!("[{}] New high reached: {:.9} SOL", mint, current_price);
}
let price_change_pct =
(current_price - pos.trade.entry_price_sol) / pos.trade.entry_price_sol;
let drop_from_peak = (pos.highest_price_sol - current_price) / pos.highest_price_sol;
let should_sell = match () {
_ if price_change_pct >= 0.40 => Some(ExitReason::TakeProfit),
_ if price_change_pct <= -0.1 => Some(ExitReason::StopLoss),
_ if drop_from_peak >= 0.12 && price_change_pct > 0.10 => {
let reason = match () {
_ if price_change_pct >= TAKE_PROFIT_PCT => Some(ExitReason::TakeProfit),
_ if price_change_pct <= STOP_LOSS_PCT => Some(ExitReason::StopLoss),
_ if drop_from_peak >= TRAILING_STOP_DROP_PCT
&& price_change_pct > TRAILING_STOP_MIN_GAIN_PCT =>
{
Some(ExitReason::TrailingStop)
}
_ if pos.last_high_time.elapsed() >= Duration::from_secs(25) => {
Some(ExitReason::MomentumStalled)
}
_ => None,
};
if let Some(reason) = should_sell {
info!(
"[{}] EXECUTING SELL. Reason: {:?} {:.1}%",
mint,
reason,
price_change_pct * 100.0
);
bot.executor
.lock()
.await
.sell_percent(mint, 100, PRIORITY, SLIPPAGE)
.await?;
if let Some(mut pos) = self.positions.remove(mint) {
pos.trade.close(current_price, reason);
bot.trade_log.lock().await.push(pos.trade);
if let Some(reason) = reason {
exits.push((mint.clone(), reason));
}
self.cleanup_and_unsubscribe(&bot, mint).await?;
}
return Ok(());
}
// -------------------------------------------------------------
// 2. Evaluate Potential Buys
// 2. Stall exits for ALL open positions. Stalled positions stop
// producing trades, so this must NOT be gated on the incoming
// trade's mint — any trade evaluates every position.
// -------------------------------------------------------------
for (position_mint, pos) in self.positions.iter() {
if pos.last_high_time.elapsed() >= STALL_DURATION {
exits.push((position_mint.clone(), ExitReason::MomentumStalled));
}
}
// -------------------------------------------------------------
// 3. Execute any pending exits. Duplicates are no-ops since the
// position is removed on the first exit.
// -------------------------------------------------------------
for (exit_mint, reason) in exits {
self.execute_exit(&bot, &exit_mint, reason).await?;
}
// -------------------------------------------------------------
// 4. Evaluate Potential Buys
// -------------------------------------------------------------
if let Some(tracker) = self.trackers.get_mut(mint) {
let elapsed = tracker.created_at.elapsed();
@@ -203,6 +258,11 @@ impl Strategy for MomentumVelocityStrategy {
}
tracker.trade_count += 1;
if tracker.first_price_sol.is_none() {
tracker.first_price_sol = Some(current_price);
}
match trade.tx_type {
TradeType::Buy => {
tracker.unique_buyers.insert(trade.trader.clone());
@@ -213,39 +273,67 @@ impl Strategy for MomentumVelocityStrategy {
}
}
let has_enough_buyers = tracker.unique_buyers.len() >= self.min_unique_buyers;
let has_volume_surge = tracker.net_sol_flow >= self.min_net_sol_flow;
let v_sol = trade.v_sol_in_bonding_curve / 1_000_000_000.0;
let is_early_curve = v_sol < 60.0;
let price_change_pct = match tracker.first_price_sol {
Some(first_price) if first_price > 0.0 => {
(current_price - first_price) / first_price
}
_ => 0.0,
};
let has_enough_buyers = tracker.unique_buyers.len() >= MIN_UNIQUE_BUYERS;
let has_volume_surge = tracker.net_sol_flow >= MIN_NET_SOL_FLOW;
let has_min_trades = tracker.trade_count >= MIN_TRADE_COUNT;
let is_early_curve = v_sol < MAX_CURVE_SOL;
let is_buy_trade = matches!(trade.tx_type, TradeType::Buy);
let has_momentum = price_change_pct >= MIN_MOMENTUM_PCT;
let has_capacity = self.positions.len() < MAX_OPEN_POSITIONS;
// Log detailed status of buy criteria evaluation on every trade
debug!(
"[{}] Trade #{} ({:?}) | Buyers: {}/{} [{}] | Net Flow: {:.3}/{:.3} SOL [{}] | Curve SOL: {:.2} < 60 [{}]",
"[{}] Trade #{} ({:?}) | Buyers: {}/{} [{}] | Net Flow: {:.3}/{:.3} SOL [{}] | Trades: {}/{} [{}] | Curve SOL: {:.2} < {:.0} [{}] | Price Δ: {:.2}% >= {:.0}% [{}] | Capacity: {}/{} [{}]",
mint,
tracker.trade_count,
trade.tx_type,
tracker.unique_buyers.len(),
self.min_unique_buyers,
MIN_UNIQUE_BUYERS,
if has_enough_buyers { "PASS" } else { "FAIL" },
tracker.net_sol_flow,
self.min_net_sol_flow,
MIN_NET_SOL_FLOW,
if has_volume_surge { "PASS" } else { "FAIL" },
trade.v_sol_in_bonding_curve,
if is_early_curve { "PASS" } else { "FAIL" }
tracker.trade_count,
MIN_TRADE_COUNT,
if has_min_trades { "PASS" } else { "FAIL" },
v_sol,
MAX_CURVE_SOL,
if is_early_curve { "PASS" } else { "FAIL" },
price_change_pct * 100.0,
MIN_MOMENTUM_PCT * 100.0,
if has_momentum { "PASS" } else { "FAIL" },
self.positions.len(),
MAX_OPEN_POSITIONS,
if has_capacity { "PASS" } else { "FAIL" }
);
if has_enough_buyers && has_volume_surge && is_early_curve {
if has_capacity
&& has_enough_buyers
&& has_volume_surge
&& has_min_trades
&& is_buy_trade
&& has_momentum
&& is_early_curve
{
info!(
"🚀 BUY SIGNAL TRIGGERED for {}! Unique Buyers: {}, Net Flow: {:.3} SOL, Curve SOL: {:.2}",
"🚀 BUY SIGNAL TRIGGERED for {}! Unique Buyers: {}, Net Flow: {:.3} SOL, Trades: {}, Curve SOL: {:.2}, Price Δ: {:.2}%",
mint,
tracker.unique_buyers.len(),
tracker.net_sol_flow,
trade.v_sol_in_bonding_curve
tracker.trade_count,
trade.v_sol_in_bonding_curve,
price_change_pct * 100.0
);
bot.executor
.lock()
.await
.buy(mint, BUY_AMOUNT_SOL, PRIORITY, SLIPPAGE)
.await?;
@@ -260,6 +348,7 @@ impl Strategy for MomentumVelocityStrategy {
trade.v_sol_in_bonding_curve,
),
highest_price_sol: current_price,
last_price_sol: current_price,
last_high_time: Instant::now(),
},
);
-115
View File
@@ -1,115 +0,0 @@
use serde::{Deserialize, Deserializer};
use serde_json::Value;
#[derive(Debug)]
pub enum PumpDevEvent {
Connected { client_id: u64, message: String },
ConnectionStatus { connected: bool, timestamp: u64 },
Subscribed { method: String },
Unsubscribed { method: String, keys: Vec<String> },
Create(NewToken),
Trade(Trade),
}
impl<'de> Deserialize<'de> for PumpDevEvent {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
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("name").is_some() {
let token: NewToken =
serde_json::from_value(value).map_err(serde::de::Error::custom)?;
return Ok(PumpDevEvent::Create(token));
}
#[derive(Deserialize)]
#[serde(tag = "type")]
enum Tagged {
#[serde(rename = "connected")]
Connected {
#[serde(rename = "clientId")]
client_id: u64,
message: String,
},
#[serde(rename = "connectionStatus")]
ConnectionStatus { connected: bool, timestamp: u64 },
#[serde(rename = "subscribed")]
Subscribed { method: String },
#[serde(rename = "unsubscribed")]
Unsubscribed { method: String, keys: Vec<String> },
}
match serde_json::from_value(value).map_err(serde::de::Error::custom)? {
Tagged::Connected { client_id, message } => {
Ok(PumpDevEvent::Connected { client_id, message })
}
Tagged::ConnectionStatus {
connected,
timestamp,
} => Ok(PumpDevEvent::ConnectionStatus {
connected,
timestamp,
}),
Tagged::Subscribed { method } => Ok(PumpDevEvent::Subscribed { method }),
Tagged::Unsubscribed { method, keys } => {
Ok(PumpDevEvent::Unsubscribed { method, keys })
}
}
}
}
#[derive(Debug, Clone, Deserialize)]
pub struct NewToken {
pub mint: String,
#[serde(rename = "traderPublicKey")]
pub trader_public_key: String,
pub name: String,
pub symbol: String,
pub uri: String,
#[serde(rename = "marketCapSol")]
pub market_cap_sol: f64,
#[serde(rename = "solAmount")]
pub sol_amount: f64,
}
#[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,
#[serde(rename = "traderPublicKey")]
pub trader: String,
#[serde(rename = "txType")]
pub tx_type: TradeType,
#[serde(rename = "solAmount")]
pub sol_amount: f64,
#[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,
}