New token streaming

This commit is contained in:
2026-08-07 20:13:54 +02:00
parent c7d28d6cc7
commit 665e4a5832
6 changed files with 214 additions and 33 deletions
Generated
+2
View File
@@ -3654,6 +3654,8 @@ version = "0.1.0"
dependencies = [ dependencies = [
"anyhow", "anyhow",
"async-trait", "async-trait",
"base64 0.22.1",
"bs58",
"chrono", "chrono",
"csv", "csv",
"env_logger", "env_logger",
+2
View File
@@ -11,6 +11,8 @@ anyhow = "1.0.104"
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"
base64 = "0.22.1"
bs58 = "0.5.1"
csv = "1.4.0" csv = "1.4.0"
log = "0.4.33" log = "0.4.33"
+8 -15
View File
@@ -6,6 +6,7 @@ use tokio::sync::{Mutex, mpsc, watch};
use crate::{ use crate::{
data::{ data::{
NewToken,
account::{Account, AccountManager}, account::{Account, AccountManager},
tradelog::TradeLog, tradelog::TradeLog,
}, },
@@ -76,21 +77,13 @@ impl Bot {
Ok(()) Ok(())
} }
pub async fn tick(self: &Arc<Self>, rx: &mut mpsc::Receiver<u32>) -> anyhow::Result<bool> { pub async fn tick(self: &Arc<Self>, rx: &mut mpsc::Receiver<NewToken>) -> anyhow::Result<bool> {
if let Some(data) = rx.recv().await { if let Some(token) = rx.recv().await {
// drop(ws); self.strategy
// self.strategy .lock()
// .lock() .await
// .await .on_new_coin(self.clone(), token)
// .on_new_coin(self.clone(), token) .await?;
// .await?;
// drop(ws);
// self.strategy
// .lock()
// .await
// .on_trade(self.clone(), trade)
// .await?;
Ok(false) Ok(false)
} else { } else {
-2
View File
@@ -14,8 +14,6 @@ pub struct NewToken {
pub uri: String, pub uri: String,
#[serde(rename = "marketCapSol")] #[serde(rename = "marketCapSol")]
pub market_cap_sol: f64, pub market_cap_sol: f64,
#[serde(rename = "solAmount")]
pub sol_amount: f64,
} }
#[derive(Debug, Clone, Deserialize)] #[derive(Debug, Clone, Deserialize)]
+8 -16
View File
@@ -1,11 +1,10 @@
mod pump_fun; mod pump_fun;
use crate::launchpad::pump_fun::PumpFun;
use futures_util::StreamExt;
use rust_decimal::Decimal; use rust_decimal::Decimal;
use std::{collections::HashMap, sync::Arc}; use std::{collections::HashMap, sync::Arc};
use tokio::sync::{Mutex, mpsc}; use tokio::sync::{Mutex, mpsc};
use tokio_tungstenite::tungstenite::Message;
use crate::{data::NewToken, launchpad::pump_fun::PumpFun};
pub struct Client( pub struct Client(
pub Mutex< pub Mutex<
@@ -44,6 +43,8 @@ pub trait Launchpad: Send + Sync {
slippage: u16, slippage: u16,
) -> anyhow::Result<()>; ) -> anyhow::Result<()>;
async fn listen(client: Arc<Client>, tx: mpsc::Sender<NewToken>) -> anyhow::Result<()>;
fn get_positions(&self) -> HashMap<String, Decimal>; fn get_positions(&self) -> HashMap<String, Decimal>;
} }
@@ -52,7 +53,7 @@ impl Executor {
Ok(Arc::new(Self { Ok(Arc::new(Self {
client: Arc::new(Client(Mutex::new( client: Arc::new(Client(Mutex::new(
tokio_tungstenite::connect_async(format!( tokio_tungstenite::connect_async(format!(
"wss://devnet.helius-rpc.com/?api-key={api_key}" "wss://mainnet.helius-rpc.com/?api-key={api_key}"
)) ))
.await? .await?
.0, .0,
@@ -105,19 +106,10 @@ impl Executor {
Ok(()) Ok(())
} }
pub async fn listen(self: &Arc<Self>) -> anyhow::Result<mpsc::Receiver<u32>> { pub async fn listen(self: &Arc<Self>) -> anyhow::Result<mpsc::Receiver<NewToken>> {
let (tx, rx) = mpsc::channel(10); let (tx, rx) = mpsc::channel(100);
let client = self.client.clone(); tokio::spawn(PumpFun::listen(self.client.clone(), tx));
tokio::spawn(async move {
match client.0.lock().await.next().await {
Some(Ok(Message::Text(msg))) => {}
Some(Ok(msg)) => {}
Some(Err(e)) => {}
None => {}
}
});
Ok(rx) Ok(rx)
} }
+194
View File
@@ -1,10 +1,21 @@
use std::collections::HashMap; 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 rust_decimal::Decimal;
use serde::Deserialize;
use tokio::sync::mpsc;
use tokio_tungstenite::tungstenite::Message;
use crate::data::NewToken;
use crate::launchpad::Client; use crate::launchpad::Client;
use crate::launchpad::Launchpad; use crate::launchpad::Launchpad;
const PUMP_FUN_PROGRAM: &str = "6EF8rrecthR5Dkzon8Nwu78hRvfCKubJ14M5uBEwF6P";
const CREATE_EVENT_DISCRIMINATOR: [u8; 8] = [27, 114, 169, 77, 222, 235, 99, 118];
pub struct PumpFun { pub struct PumpFun {
pub positions: HashMap<String, Decimal>, pub positions: HashMap<String, Decimal>,
} }
@@ -44,7 +55,190 @@ impl Launchpad for PumpFun {
Ok(()) Ok(())
} }
async fn listen(client: Arc<Client>, tx: mpsc::Sender<NewToken>) -> anyhow::Result<()> {
let mut ws = client.0.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;
}
let Some(event) = logs.iter().find_map(|log| {
log.strip_prefix("Program data: ")
.and_then(|encoded| STANDARD.decode(encoded).ok())
.and_then(|data| parse_create_event(&data))
}) else {
continue;
};
let token = NewToken {
mint: event.mint,
trader_public_key: event.user,
name: event.name,
symbol: event.symbol,
uri: event.uri,
market_cap_sol: event.market_cap_sol,
};
log::info!(
"[PUMP.FUN] New token: {} (${}) - {}",
token.name,
token.symbol,
token.mint
);
if tx.send(token).await.is_err() {
break;
}
}
Ok(())
}
fn get_positions<'a>(&'a self) -> HashMap<String, Decimal> { fn get_positions<'a>(&'a self) -> HashMap<String, Decimal> {
self.positions.clone() self.positions.clone()
} }
} }
struct CreateEvent {
name: String,
symbol: String,
uri: String,
mint: String,
user: String,
market_cap_sol: f64,
}
fn parse_create_event(data: &[u8]) -> Option<CreateEvent> {
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 user = 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(CreateEvent {
name,
symbol,
uri,
mint,
user,
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))
}
#[derive(Deserialize)]
struct LogsNotification {
method: String,
params: NotificationParams,
}
#[derive(Deserialize)]
struct NotificationParams {
result: NotificationResult,
}
#[derive(Deserialize)]
struct NotificationResult {
value: LogsValue,
}
#[derive(Deserialize)]
struct LogsValue {
logs: Vec<String>,
}