Merge pull request #3 from selimaj-dev/working-strategy

Working strategy
This commit is contained in:
2026-08-07 06:44:38 +02:00
committed by GitHub
10 changed files with 864 additions and 79 deletions
Generated
+122
View File
@@ -22,6 +22,15 @@ dependencies = [
"memchr", "memchr",
] ]
[[package]]
name = "android_system_properties"
version = "0.1.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ae221649c9976a6f6c56ae1facf410f3ddb33cc661c4b7b61020a912d4237fbc"
dependencies = [
"libc",
]
[[package]] [[package]]
name = "anstream" name = "anstream"
version = "1.0.0" version = "1.0.0"
@@ -262,6 +271,20 @@ dependencies = [
"rand_core 0.10.1", "rand_core 0.10.1",
] ]
[[package]]
name = "chrono"
version = "0.4.45"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1aa79e62e7697b8e29b513a68abacf485adcd1fe8284a4316c5ae868e6633327"
dependencies = [
"iana-time-zone",
"js-sys",
"num-traits",
"serde",
"wasm-bindgen",
"windows-link",
]
[[package]] [[package]]
name = "cmake" name = "cmake"
version = "0.1.58" version = "0.1.58"
@@ -337,6 +360,27 @@ dependencies = [
"hybrid-array", "hybrid-array",
] ]
[[package]]
name = "csv"
version = "1.4.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "52cd9d68cf7efc6ddfaaee42e7288d3a99d613d4b50f76ce9827ae0c6e14f938"
dependencies = [
"csv-core",
"itoa",
"ryu",
"serde_core",
]
[[package]]
name = "csv-core"
version = "0.1.13"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "704a3c26996a80471189265814dbc2c257598b96b8a7feae2d31ace646bb9782"
dependencies = [
"memchr",
]
[[package]] [[package]]
name = "data-encoding" name = "data-encoding"
version = "2.11.1" version = "2.11.1"
@@ -726,6 +770,30 @@ dependencies = [
"windows-registry", "windows-registry",
] ]
[[package]]
name = "iana-time-zone"
version = "0.1.65"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e31bc9ad994ba00e440a8aa5c9ef0ec67d5cb5e5cb0cc7f8b744a35b389cc470"
dependencies = [
"android_system_properties",
"core-foundation-sys",
"iana-time-zone-haiku",
"js-sys",
"log",
"wasm-bindgen",
"windows-core",
]
[[package]]
name = "iana-time-zone-haiku"
version = "0.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f31827a206f56af32e590ba56d5d2d085f558508192593743f16b2306495269f"
dependencies = [
"cc",
]
[[package]] [[package]]
name = "icu_collections" name = "icu_collections"
version = "2.2.0" version = "2.2.0"
@@ -1453,11 +1521,22 @@ dependencies = [
"num-traits", "num-traits",
"rand 0.8.7", "rand 0.8.7",
"rkyv", "rkyv",
"rust_decimal_macros",
"serde", "serde",
"serde_json", "serde_json",
"wasm-bindgen", "wasm-bindgen",
] ]
[[package]]
name = "rust_decimal_macros"
version = "1.40.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "74a5a6f027e892c7a035c6fddb50435a1fbf5a734ffc0c2a9fed4d0221440519"
dependencies = [
"quote",
"syn 2.0.119",
]
[[package]] [[package]]
name = "rustc-hash" name = "rustc-hash"
version = "2.1.3" version = "2.1.3"
@@ -1567,6 +1646,12 @@ version = "1.0.23"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cf54715a573b99ac80df0bc206da022bcd442c974952c7b9720069370852e21f" checksum = "cf54715a573b99ac80df0bc206da022bcd442c974952c7b9720069370852e21f"
[[package]]
name = "ryu"
version = "1.0.23"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f"
[[package]] [[package]]
name = "same-file" name = "same-file"
version = "1.0.6" version = "1.0.6"
@@ -1724,6 +1809,8 @@ version = "0.1.0"
dependencies = [ dependencies = [
"anyhow", "anyhow",
"async-trait", "async-trait",
"chrono",
"csv",
"env_logger", "env_logger",
"futures-util", "futures-util",
"log", "log",
@@ -2259,6 +2346,41 @@ dependencies = [
"windows-sys 0.61.2", "windows-sys 0.61.2",
] ]
[[package]]
name = "windows-core"
version = "0.62.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b8e83a14d34d0623b51dce9581199302a221863196a1dde71a7663a4c2be9deb"
dependencies = [
"windows-implement",
"windows-interface",
"windows-link",
"windows-result",
"windows-strings",
]
[[package]]
name = "windows-implement"
version = "0.60.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "053e2e040ab57b9dc951b72c264860db7eb3b0200ba345b4e4c3b14f67855ddf"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.119",
]
[[package]]
name = "windows-interface"
version = "0.59.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3f316c4a2570ba26bbec722032c4099d8c8bc095efccdc15688708623367e358"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.119",
]
[[package]] [[package]]
name = "windows-link" name = "windows-link"
version = "0.2.1" version = "0.2.1"
+4 -7
View File
@@ -10,15 +10,12 @@ futures-util = "0.3.33"
serde = { version = "1.0.229", features = ["derive", "serde_derive"] } serde = { version = "1.0.229", features = ["derive", "serde_derive"] }
serde_json = "1.0.151" serde_json = "1.0.151"
tokio = { version = "1.53.1", features = [ tokio = { version = "1.53.1", features = ["macros", "rt-multi-thread", "sync", "fs", "io-std"] }
"macros",
"rt-multi-thread",
"sync",
"fs",
] }
tokio-tungstenite = { version = "0.30.0", features = ["native-tls"] } tokio-tungstenite = { version = "0.30.0", features = ["native-tls"] }
log = "0.4.33" log = "0.4.33"
env_logger = "0.11.11" env_logger = "0.11.11"
rust_decimal = "1.42.1" rust_decimal = { version = "1.42.1", features = ["macros"] }
reqwest = { version = "0.13.4", features = ["json"] } reqwest = { version = "0.13.4", features = ["json"] }
async-trait = "0.1.91" async-trait = "0.1.91"
csv = "1.4.0"
chrono = { version = "0.4.45", features = ["serde"] }
+139 -30
View File
@@ -1,25 +1,67 @@
use std::{collections::HashMap, sync::Arc};
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, watch},
};
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::ExecutorWrapper,
strategy::{Strategy, veloc::MomentumVelocityStrategy},
tradelog::TradeLog,
};
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<ExecutorWrapper>,
pub strategy: Mutex<Box<dyn Strategy>>,
pub trade_log: Mutex<Vec<TradeLog>>,
} }
impl Bot { impl Bot {
pub async fn on_new_coin(&mut self, token: NewToken) -> anyhow::Result<()> { 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(()) Ok(())
} }
} }
impl Bot { impl Bot {
pub async fn new() -> anyhow::Result<Self> { pub async fn new() -> anyhow::Result<Arc<Self>> {
let accounts = AccountManager::get().await?; let accounts = AccountManager::get().await?;
let account = accounts let account = accounts
@@ -28,28 +70,42 @@ impl Bot {
.context("Failed to get account")? .context("Failed to get account")?
.clone(); .clone();
Ok(Self { Ok(Arc::new(Self {
ws: connect_async("wss://pumpdev.io/ws").await?.0, ws: Mutex::new(connect_async("wss://pumpdev.io/ws").await?.0),
executor: Mutex::new(ExecutorWrapper {
executor: Box::new(account.executor()), executor: Box::new(account.executor()),
accounts, positions: HashMap::new(),
}) }),
accounts: Mutex::new(accounts),
strategy: Mutex::new(Box::new(MomentumVelocityStrategy::new())),
trade_log: Mutex::new(Vec::new()),
}))
} }
pub async fn refresh_account(&mut self) -> anyhow::Result<()> { pub async fn refresh_account(self: &Arc<Self>) -> anyhow::Result<()> {
let account = self self.strategy
.lock()
.await
.execute_sell_all(self.clone())
.await?;
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.executor = 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(),
)) ))
@@ -58,23 +114,34 @@ impl Bot {
Ok(()) Ok(())
} }
pub async fn start(&mut self) -> anyhow::Result<()> { pub async fn start(
self: &Arc<Self>,
mut shutdown: watch::Receiver<bool>,
) -> anyhow::Result<()> {
self.initialize_websocket_subscribe().await?; self.initialize_websocket_subscribe().await?;
while let Some(msg) = self.ws.next().await { loop {
let msg = msg?; tokio::select! {
if let tokio_tungstenite::tungstenite::Message::Text(text) = msg { _ = shutdown.changed() => {
match serde_json::from_str::<crate::types::PumpDevEvent>(&text) { if *shutdown.borrow() {
Ok(crate::types::PumpDevEvent::Create(token)) => { log::info!("Shutdown signal received.");
self.on_new_coin(token).await?;
self.strategy
.lock()
.await
.execute_sell_all(self.clone())
.await?;
log::info!("Bye!");
break;
}
} }
Ok(event) => { result = self.tick() => {
println!("{:?}", event); if result? {
} log::warn!("Websocket closed.");
break;
Err(err) => {
log::error!("{err}");
} }
} }
} }
@@ -82,4 +149,46 @@ impl Bot {
Ok(()) 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);
self.strategy
.lock()
.await
.on_new_coin(self.clone(), token)
.await?;
}
Ok(crate::types::PumpDevEvent::Trade(trade)) => {
drop(ws);
self.strategy
.lock()
.await
.on_trade(self.clone(), trade)
.await?;
}
Ok(_event) => {
// println!("{:?}", event);
}
Err(err) => {
log::error!("{err}, MSG -> {text}");
}
}
}
Ok(false)
} else {
Ok(true)
}
}
} }
+89 -4
View File
@@ -1,15 +1,22 @@
pub mod pump_fun; pub mod pump_fun;
use std::collections::HashMap;
use rust_decimal::Decimal; use rust_decimal::Decimal;
use crate::account::Account; use crate::account::Account;
pub struct ExecutorWrapper {
pub executor: Box<dyn Executor>,
pub positions: HashMap<String, Decimal>,
}
#[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: &str,
amount: Decimal, amount: Decimal,
priority: Decimal, priority: Decimal,
slippage: u16, slippage: u16,
@@ -19,7 +26,7 @@ pub trait Executor {
async fn sell( async fn sell(
&self, &self,
mint: String, mint: &str,
amount: Decimal, amount: Decimal,
priority: Decimal, priority: Decimal,
slippage: u16, slippage: u16,
@@ -29,7 +36,7 @@ pub trait Executor {
async fn sell_percent( async fn sell_percent(
&self, &self,
mint: String, mint: &str,
amount: u8, amount: u8,
priority: Decimal, priority: Decimal,
slippage: u16, slippage: u16,
@@ -45,3 +52,81 @@ impl 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(())
}
}
+37 -31
View File
@@ -39,7 +39,7 @@ impl PumpDev {
async fn trade( async fn trade(
&self, &self,
action: &str, action: &str,
mint: String, mint: &str,
amount: String, amount: String,
priority: Decimal, priority: Decimal,
slippage: u16, slippage: u16,
@@ -88,55 +88,61 @@ impl PumpDev {
impl Executor for PumpDev { impl Executor for PumpDev {
async fn buy( async fn buy(
&self, &self,
mint: String, mint: &str,
amount: Decimal, amount: Decimal,
priority: Decimal, priority: Decimal,
slippage: u16, slippage: u16,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
self.trade( log::info!("BUY {mint} {amount} SOL");
"buy", Ok(())
mint, // self.trade(
amount.round_dp(3).to_string(), // "buy",
priority, // mint,
slippage, // amount.round_dp(3).to_string(),
true, // priority,
) // slippage,
.await // true,
// )
// .await
} }
async fn sell( async fn sell(
&self, &self,
mint: String, mint: &str,
amount: Decimal, amount: Decimal,
priority: Decimal, priority: Decimal,
slippage: u16, slippage: u16,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
self.trade( log::info!("SELL {mint} {amount} SOL");
"sell", Ok(())
mint, // self.trade(
amount.round_dp(3).to_string(), // "sell",
priority, // mint,
slippage, // amount.round_dp(3).to_string(),
false, // priority,
) // slippage,
.await // false,
// )
// .await
} }
async fn sell_percent( async fn sell_percent(
&self, &self,
mint: String, mint: &str,
amount: u8, amount: u8,
priority: Decimal, priority: Decimal,
slippage: u16, slippage: u16,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
self.trade( log::info!("SELL {mint} {amount}%");
"sell", Ok(())
mint, // self.trade(
format!("{amount}%"), // "sell",
priority, // mint,
slippage, // format!("{amount}%"),
false, // priority,
) // slippage,
.await // false,
// )
// .await
} }
} }
+63 -2
View File
@@ -1,17 +1,78 @@
pub mod account; pub mod account;
pub mod bot; pub mod bot;
pub mod executor; pub mod executor;
pub mod strategy;
pub mod tradelog;
pub mod types; pub mod types;
use std::sync::Arc;
use chrono::Local;
use tokio::{
io::{self, AsyncBufReadExt, BufReader},
sync::watch,
};
use crate::bot::Bot; use crate::bot::Bot;
async fn try_shutdown_listener(bot: Arc<Bot>, tx: watch::Sender<bool>) -> anyhow::Result<()> {
let mut stdin = BufReader::new(io::stdin());
let mut line = String::new();
loop {
stdin.read_line(&mut line).await?;
match line.to_lowercase().trim() {
"exit" | "shutdown" => {
let _ = tx.send(true);
break Ok(());
}
"save" => {
log::info!("Saving trades");
let now = Local::now();
let formatted_time = now.format("%m-%d-%H-%M").to_string();
let filename = format!("sol-hun-{}.csv", formatted_time);
let trade_log = bot.trade_log.lock().await;
log::info!("Saving {} trades", trade_log.len());
let mut writer = csv::Writer::from_path(&filename)?;
for trade in trade_log.iter() {
writer.serialize(trade)?;
}
writer.flush()?;
log::info!("Saved trades to {}", filename);
}
_ => {}
}
}
}
async fn shutdown_listener(bot: Arc<Bot>, tx: watch::Sender<bool>) {
if let Err(e) = try_shutdown_listener(bot, tx).await {
log::error!("{e}");
}
}
#[tokio::main] #[tokio::main]
async fn main() -> anyhow::Result<()> { async fn main() -> anyhow::Result<()> {
let mut builder = env_logger::Builder::from_default_env(); let mut builder = env_logger::Builder::from_default_env();
builder.filter_level(log::LevelFilter::Info); builder.filter_level(log::LevelFilter::Info);
builder.init(); builder.init();
let mut bot = Bot::new().await?; let (shutdown_tx, shutdown_rx) = watch::channel(false);
bot.start().await let bot = Bot::new().await?;
tokio::spawn(shutdown_listener(bot.clone(), shutdown_tx));
bot.start(shutdown_rx).await
} }
+17
View File
@@ -0,0 +1,17 @@
pub mod veloc;
use std::sync::Arc;
use crate::{
bot::Bot,
types::{NewToken, Trade},
};
#[async_trait::async_trait]
pub trait Strategy: Send + Sync {
async fn execute_sell_all(&mut self, bot: Arc<Bot>) -> anyhow::Result<()>;
async fn on_new_coin(&mut self, bot: Arc<Bot>, token: NewToken) -> anyhow::Result<()>;
async fn on_trade(&mut self, bot: Arc<Bot>, trade: Trade) -> anyhow::Result<()>;
}
+275
View File
@@ -0,0 +1,275 @@
use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::Arc;
use std::time::{Duration, Instant};
use log::{debug, info, trace, warn};
use rust_decimal::{Decimal, dec};
use crate::bot::Bot;
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;
struct TokenTracker {
created_at: Instant,
unique_buyers: HashSet<String>,
net_sol_flow: f64,
trade_count: usize,
}
struct OpenPosition {
trade: TradeLog,
highest_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>,
positions: HashMap<String, OpenPosition>,
// Tracks active token subscriptions to enforce <= 5 limit
active_subscriptions: VecDeque<String>,
}
impl MomentumVelocityStrategy {
pub fn new() -> 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),
}
}
fn calculate_price_sol(&self, trade: &Trade) -> f64 {
if trade.v_tokens_in_bonding_curve == 0.0 {
return 0.0;
}
trade.v_sol_in_bonding_curve / trade.v_tokens_in_bonding_curve
}
/// 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);
}
self.trackers.remove(mint);
self.active_subscriptions.retain(|m| m != mint);
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
}
async fn on_new_coin(&mut self, bot: Arc<Bot>, token: NewToken) -> anyhow::Result<()> {
trace!("[NEW COIN] Event received for token: {}", token.mint);
if self.positions.contains_key(&token.mint) || self.trackers.contains_key(&token.mint) {
trace!(
"[{}] Already tracking or holding position. Skipping.",
token.mint
);
return Ok(());
}
// Evict oldest tracked token that DOES NOT have an active open position
while self.active_subscriptions.len() >= MAX_SUBSCRIBED_TOKENS {
let eviction_index = self
.active_subscriptions
.iter()
.position(|mint| !self.positions.contains_key(mint));
if let Some(idx) = eviction_index {
if let Some(mint_to_remove) = self.active_subscriptions.remove(idx) {
info!(
"[{}] Capacity reached ({}/{}). Evicting un-bought token from queue.",
mint_to_remove, MAX_SUBSCRIBED_TOKENS, MAX_SUBSCRIBED_TOKENS
);
let _ = bot.unsubscribe(&mint_to_remove).await;
self.trackers.remove(&mint_to_remove);
}
} else {
warn!(
"[QUEUE FULL] All {} slots are occupied by active positions. Cannot track {}",
MAX_SUBSCRIBED_TOKENS, token.mint
);
return Ok(());
}
}
info!("[{}] Subscribing and creating tracker.", token.mint);
bot.subscribe(&token.mint).await?;
self.active_subscriptions.push_back(token.mint.clone());
self.trackers.insert(
token.mint,
TokenTracker {
created_at: Instant::now(),
unique_buyers: HashSet::new(),
net_sol_flow: 0.0,
trade_count: 0,
},
);
Ok(())
}
async fn on_trade(&mut self, bot: Arc<Bot>, trade: Trade) -> anyhow::Result<()> {
let mint = &trade.mint;
let current_price = self.calculate_price_sol(&trade);
// -------------------------------------------------------------
// 1. Manage Active Positions (Take Profit / Stop Loss / Stall)
// -------------------------------------------------------------
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;
if current_price > pos.highest_price_sol {
pos.highest_price_sol = current_price;
pos.last_high_time = Instant::now();
trace!("[{}] New high reached: {:.9} SOL", mint, current_price);
}
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 => {
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);
}
self.cleanup_and_unsubscribe(&bot, mint).await?;
}
return Ok(());
}
// -------------------------------------------------------------
// 2. Evaluate Potential Buys
// -------------------------------------------------------------
if let Some(tracker) = self.trackers.get_mut(mint) {
let elapsed = tracker.created_at.elapsed();
if elapsed > self.max_tracking_duration {
info!(
"[{}] Tracking window expired ({:?} > {:?}). Cleaning up.",
mint, elapsed, self.max_tracking_duration
);
self.cleanup_and_unsubscribe(&bot, mint).await?;
return Ok(());
}
tracker.trade_count += 1;
match trade.tx_type {
TradeType::Buy => {
tracker.unique_buyers.insert(trade.trader.clone());
tracker.net_sol_flow += trade.sol_amount;
}
TradeType::Sell => {
tracker.net_sol_flow -= trade.sol_amount;
}
}
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;
// Log detailed status of buy criteria evaluation on every trade
debug!(
"[{}] Trade #{} ({:?}) | Buyers: {}/{} [{}] | Net Flow: {:.3}/{:.3} SOL [{}] | Curve SOL: {:.2} < 60 [{}]",
mint,
tracker.trade_count,
trade.tx_type,
tracker.unique_buyers.len(),
self.min_unique_buyers,
if has_enough_buyers { "PASS" } else { "FAIL" },
tracker.net_sol_flow,
self.min_net_sol_flow,
if has_volume_surge { "PASS" } else { "FAIL" },
trade.v_sol_in_bonding_curve,
if is_early_curve { "PASS" } else { "FAIL" }
);
if has_enough_buyers && has_volume_surge && is_early_curve {
info!(
"🚀 BUY SIGNAL TRIGGERED for {}! Unique Buyers: {}, Net Flow: {:.3} SOL, Curve SOL: {:.2}",
mint,
tracker.unique_buyers.len(),
tracker.net_sol_flow,
trade.v_sol_in_bonding_curve
);
bot.executor
.lock()
.await
.buy(mint, BUY_AMOUNT_SOL, PRIORITY, SLIPPAGE)
.await?;
self.positions.insert(
mint.clone(),
OpenPosition {
trade: TradeLog::new(
mint.clone(),
current_price,
tracker.unique_buyers.len(),
tracker.net_sol_flow,
trade.v_sol_in_bonding_curve,
),
highest_price_sol: current_price,
last_high_time: Instant::now(),
},
);
self.trackers.remove(mint);
}
} else {
trace!("[{}] Received trade for untracked mint.", mint);
}
Ok(())
}
}
+70
View File
@@ -0,0 +1,70 @@
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum ExitReason {
TakeProfit,
StopLoss,
TrailingStop,
MomentumStalled,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TradeLog {
pub mint: String,
pub opened_at: chrono::DateTime<chrono::Utc>,
pub closed_at: Option<chrono::DateTime<chrono::Utc>>,
pub duration_secs: Option<i64>,
pub entry_price_sol: f64,
pub exit_price_sol: Option<f64>,
pub pnl_percent: Option<f64>,
pub exit_reason: Option<ExitReason>,
pub unique_buyers: usize,
pub net_sol_flow: f64,
pub curve_sol: f64,
}
impl TradeLog {
pub fn new(
mint: String,
entry_price_sol: f64,
unique_buyers: usize,
net_sol_flow: f64,
curve_sol: f64,
) -> Self {
Self {
mint,
opened_at: chrono::Utc::now(),
closed_at: None,
entry_price_sol,
exit_price_sol: None,
pnl_percent: None,
duration_secs: None,
exit_reason: None,
unique_buyers,
net_sol_flow,
curve_sol,
}
}
pub fn close(&mut self, exit_price_sol: f64, reason: ExitReason) {
let now = chrono::Utc::now();
let pnl = ((exit_price_sol - self.entry_price_sol) / self.entry_price_sol) * 100.0;
self.exit_price_sol = Some(exit_price_sol);
self.pnl_percent = Some(pnl);
self.closed_at = Some(now);
self.duration_secs = Some(now.signed_duration_since(self.opened_at).num_seconds());
self.exit_reason = Some(reason);
}
}
+47 -4
View File
@@ -6,7 +6,9 @@ pub enum PumpDevEvent {
Connected { client_id: u64, message: String }, Connected { client_id: u64, message: String },
ConnectionStatus { connected: bool, timestamp: u64 }, ConnectionStatus { connected: bool, timestamp: u64 },
Subscribed { method: String }, Subscribed { method: String },
Unsubscribed { method: String, keys: Vec<String> },
Create(NewToken), Create(NewToken),
Trade(Trade),
} }
impl<'de> Deserialize<'de> for PumpDevEvent { impl<'de> Deserialize<'de> for PumpDevEvent {
@@ -15,11 +17,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));
} }
@@ -32,12 +37,12 @@ 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 },
#[serde(rename = "unsubscribed")]
Unsubscribed { method: String, keys: Vec<String> },
} }
match serde_json::from_value(value).map_err(serde::de::Error::custom)? { match serde_json::from_value(value).map_err(serde::de::Error::custom)? {
@@ -52,6 +57,9 @@ impl<'de> Deserialize<'de> for PumpDevEvent {
timestamp, timestamp,
}), }),
Tagged::Subscribed { method } => Ok(PumpDevEvent::Subscribed { method }), Tagged::Subscribed { method } => Ok(PumpDevEvent::Subscribed { method }),
Tagged::Unsubscribed { method, keys } => {
Ok(PumpDevEvent::Unsubscribed { method, keys })
}
} }
} }
} }
@@ -70,3 +78,38 @@ pub struct NewToken {
#[serde(rename = "solAmount")] #[serde(rename = "solAmount")]
pub sol_amount: f64, 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,
}