Refactoring to helius
This commit is contained in:
+56
-117
@@ -1,63 +1,27 @@
|
||||
use std::{collections::HashMap, 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 helius::Helius;
|
||||
|
||||
use tokio::sync::{Mutex, watch};
|
||||
|
||||
use crate::{
|
||||
account::AccountManager,
|
||||
account::{Account, AccountManager},
|
||||
executor::{ExecutorWrapper, pump_fun::PumpDev},
|
||||
strategy::{Strategy, veloc::MomentumVelocityStrategy},
|
||||
tradelog::TradeLog,
|
||||
};
|
||||
|
||||
pub struct Bot {
|
||||
pub ws: Mutex<WebSocketStream<MaybeTlsStream<TcpStream>>>,
|
||||
pub accounts: Mutex<AccountManager>,
|
||||
pub executor: Mutex<ExecutorWrapper>,
|
||||
pub account_manager: Mutex<AccountManager>,
|
||||
pub current_account: Mutex<Account>,
|
||||
|
||||
pub strategy: Mutex<Box<dyn Strategy>>,
|
||||
|
||||
pub executor: Mutex<ExecutorWrapper>,
|
||||
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(())
|
||||
}
|
||||
pub helius: Helius,
|
||||
}
|
||||
|
||||
impl Bot {
|
||||
@@ -71,55 +35,24 @@ impl Bot {
|
||||
.clone();
|
||||
|
||||
Ok(Arc::new(Self {
|
||||
ws: Mutex::new(connect_async("wss://pumpdev.io/ws").await?.0),
|
||||
helius: Helius::new_async(&accounts.api_key, helius::types::Cluster::Devnet).await?,
|
||||
|
||||
account_manager: Mutex::new(accounts),
|
||||
current_account: Mutex::new(account),
|
||||
strategy: Mutex::new(Box::new(MomentumVelocityStrategy::new())),
|
||||
|
||||
trade_log: Mutex::new(Vec::new()),
|
||||
executor: Mutex::new(ExecutorWrapper {
|
||||
executor: Box::new(PumpDev { account }),
|
||||
executor: Box::new(PumpDev {}),
|
||||
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(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(PumpDev { account });
|
||||
|
||||
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?;
|
||||
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = shutdown.changed() => {
|
||||
@@ -151,44 +84,50 @@ impl Bot {
|
||||
}
|
||||
|
||||
pub async fn tick(self: &Arc<Self>) -> anyhow::Result<bool> {
|
||||
let mut ws = self.ws.lock().await;
|
||||
// drop(ws);
|
||||
// self.strategy
|
||||
// .lock()
|
||||
// .await
|
||||
// .on_new_coin(self.clone(), token)
|
||||
// .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);
|
||||
// drop(ws);
|
||||
// self.strategy
|
||||
// .lock()
|
||||
// .await
|
||||
// .on_trade(self.clone(), trade)
|
||||
// .await?;
|
||||
|
||||
self.strategy
|
||||
.lock()
|
||||
.await
|
||||
.on_new_coin(self.clone(), token)
|
||||
.await?;
|
||||
}
|
||||
Ok(false)
|
||||
}
|
||||
}
|
||||
|
||||
Ok(crate::types::PumpDevEvent::Trade(trade)) => {
|
||||
drop(ws);
|
||||
impl Bot {
|
||||
pub async fn subscribe(&self, mint: &str) -> anyhow::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
self.strategy
|
||||
.lock()
|
||||
.await
|
||||
.on_trade(self.clone(), trade)
|
||||
.await?;
|
||||
}
|
||||
pub async fn unsubscribe(&self, mint: &str) -> anyhow::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
Ok(_event) => {
|
||||
// println!("{:?}", event);
|
||||
}
|
||||
pub async fn refresh_account(self: &Arc<Self>) -> anyhow::Result<()> {
|
||||
self.strategy
|
||||
.lock()
|
||||
.await
|
||||
.execute_sell_all(self.clone())
|
||||
.await?;
|
||||
|
||||
Err(err) => {
|
||||
log::error!("{err}, MSG -> {text}");
|
||||
}
|
||||
}
|
||||
}
|
||||
let accounts = self.account_manager.lock().await;
|
||||
|
||||
Ok(false)
|
||||
} else {
|
||||
Ok(true)
|
||||
}
|
||||
let account = accounts
|
||||
.accounts
|
||||
.get(&accounts.active)
|
||||
.context("Failed to get account")?
|
||||
.clone();
|
||||
|
||||
*self.current_account.lock().await = account;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user