Event listener
This commit is contained in:
+9
-3
@@ -2,7 +2,7 @@ use std::sync::Arc;
|
|||||||
|
|
||||||
use anyhow::Context;
|
use anyhow::Context;
|
||||||
|
|
||||||
use tokio::sync::{Mutex, watch};
|
use tokio::sync::{Mutex, mpsc, watch};
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
data::{
|
data::{
|
||||||
@@ -44,6 +44,8 @@ impl Bot {
|
|||||||
self: &Arc<Self>,
|
self: &Arc<Self>,
|
||||||
mut shutdown: watch::Receiver<bool>,
|
mut shutdown: watch::Receiver<bool>,
|
||||||
) -> anyhow::Result<()> {
|
) -> anyhow::Result<()> {
|
||||||
|
let mut rx = self.executor.listen().await?;
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
_ = shutdown.changed() => {
|
_ = shutdown.changed() => {
|
||||||
@@ -62,7 +64,7 @@ impl Bot {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
result = self.tick() => {
|
result = self.tick(&mut rx) => {
|
||||||
if result? {
|
if result? {
|
||||||
log::warn!("Websocket closed.");
|
log::warn!("Websocket closed.");
|
||||||
break;
|
break;
|
||||||
@@ -74,7 +76,8 @@ impl Bot {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn tick(self: &Arc<Self>) -> anyhow::Result<bool> {
|
pub async fn tick(self: &Arc<Self>, rx: &mut mpsc::Receiver<u32>) -> anyhow::Result<bool> {
|
||||||
|
if let Some(data) = rx.recv().await {
|
||||||
// drop(ws);
|
// drop(ws);
|
||||||
// self.strategy
|
// self.strategy
|
||||||
// .lock()
|
// .lock()
|
||||||
@@ -90,6 +93,9 @@ impl Bot {
|
|||||||
// .await?;
|
// .await?;
|
||||||
|
|
||||||
Ok(false)
|
Ok(false)
|
||||||
|
} else {
|
||||||
|
Ok(true)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+12
-1
@@ -1,10 +1,11 @@
|
|||||||
mod pump_fun;
|
mod pump_fun;
|
||||||
|
|
||||||
use crate::launchpad::pump_fun::PumpFun;
|
use crate::launchpad::pump_fun::PumpFun;
|
||||||
|
use anyhow::Context;
|
||||||
use helius::Helius;
|
use helius::Helius;
|
||||||
use rust_decimal::Decimal;
|
use rust_decimal::Decimal;
|
||||||
use std::{collections::HashMap, sync::Arc};
|
use std::{collections::HashMap, sync::Arc};
|
||||||
use tokio::sync::Mutex;
|
use tokio::sync::{Mutex, mpsc};
|
||||||
|
|
||||||
pub struct Executor {
|
pub struct Executor {
|
||||||
pub client: Helius,
|
pub client: Helius,
|
||||||
@@ -89,4 +90,14 @@ impl Executor {
|
|||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub async fn listen(self: &Arc<Self>) -> anyhow::Result<mpsc::Receiver<u32>> {
|
||||||
|
let (tx, rx) = mpsc::channel(10);
|
||||||
|
|
||||||
|
let ws = self.client.ws().context("Failed to get WebSocket")?;
|
||||||
|
|
||||||
|
tokio::spawn(async move {});
|
||||||
|
|
||||||
|
Ok(rx)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user