Strategy split ws
This commit is contained in:
Generated
+15
-14
@@ -2432,9 +2432,9 @@ checksum = "e6d5a32815ae3f33302d95fdcb2ce17862f8c65363dcfd29360480ba1001fc9c"
|
||||
|
||||
[[package]]
|
||||
name = "futures"
|
||||
version = "0.3.32"
|
||||
version = "0.3.33"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "8b147ee9d1f6d097cef9ce628cd2ee62288d963e16fb287bd9286455b241382d"
|
||||
checksum = "a88cf1f829d945f548cf8fec32c61b1f202b6d93b45848602fc02af4b12ad218"
|
||||
dependencies = [
|
||||
"futures-channel",
|
||||
"futures-core",
|
||||
@@ -2447,9 +2447,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "futures-channel"
|
||||
version = "0.3.32"
|
||||
version = "0.3.33"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "07bbe89c50d7a535e539b8c17bc0b49bdb77747034daa8087407d655f3f7cc1d"
|
||||
checksum = "262590f4fe6afeb0bc83be1daa64e52657fe185690a958af7f3ad0e92085c5ae"
|
||||
dependencies = [
|
||||
"futures-core",
|
||||
"futures-sink",
|
||||
@@ -2457,15 +2457,15 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "futures-core"
|
||||
version = "0.3.32"
|
||||
version = "0.3.33"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "7e3450815272ef58cec6d564423f6e755e25379b217b0bc688e295ba24df6b1d"
|
||||
checksum = "2cd50c473c80f6d7c3670a752354b8e569b1a7cbfdc0419ec88e5edad85e0dc7"
|
||||
|
||||
[[package]]
|
||||
name = "futures-executor"
|
||||
version = "0.3.32"
|
||||
version = "0.3.33"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "baf29c38818342a3b26b5b923639e7b1f4a61fc5e76102d4b1981c6dc7a7579d"
|
||||
checksum = "6754879cc9f2c66f88c6e5c35344bb0bdb0708b0352b1201815667c7eabc7458"
|
||||
dependencies = [
|
||||
"futures-core",
|
||||
"futures-task",
|
||||
@@ -2480,9 +2480,9 @@ checksum = "4577ecaa3c4f96589d473f679a71b596316f6641bc350038b962a5daf0085d7a"
|
||||
|
||||
[[package]]
|
||||
name = "futures-macro"
|
||||
version = "0.3.32"
|
||||
version = "0.3.33"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "e835b70203e41293343137df5c0664546da5745f82ec9b84d40be8336958447b"
|
||||
checksum = "2d6d3cde68c518367be28956066ddfef33813991b77a55005a69dae04bf3b10b"
|
||||
dependencies = [
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
@@ -2497,15 +2497,15 @@ checksum = "e34418ac499d6305c2fb5ad0ed2f6ac998c5f8ca209b4510f7f94242c647e307"
|
||||
|
||||
[[package]]
|
||||
name = "futures-task"
|
||||
version = "0.3.32"
|
||||
version = "0.3.33"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "037711b3d59c33004d3856fbdc83b99d4ff37a24768fa1be9ce3538a1cde4393"
|
||||
checksum = "b231ed28831efb4a61a08580c4bc233ec56bc009f4cd8f52da2c3cb97df0c109"
|
||||
|
||||
[[package]]
|
||||
name = "futures-util"
|
||||
version = "0.3.32"
|
||||
version = "0.3.33"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "389ca41296e6190b48053de0321d02a77f32f8a5d2461dd38762c0593805c6d6"
|
||||
checksum = "a77a90a256fce34da66415271e30f94ee91c57b04b8a2c042d9cf3220179deaa"
|
||||
dependencies = [
|
||||
"futures-channel",
|
||||
"futures-core",
|
||||
@@ -3959,6 +3959,7 @@ dependencies = [
|
||||
"anyhow",
|
||||
"chrono",
|
||||
"crossterm",
|
||||
"futures",
|
||||
"hypersdk",
|
||||
"postcard",
|
||||
"pulse-sdk",
|
||||
|
||||
@@ -25,6 +25,7 @@ chrono = "0.4.45"
|
||||
serde_json = "1"
|
||||
rand = "0.8.7"
|
||||
toml = "1.1.3"
|
||||
futures = "0.3.33"
|
||||
|
||||
[workspace]
|
||||
members = ["pulse-ui", "pulse-sdk"]
|
||||
|
||||
@@ -7,6 +7,7 @@ use crate::{
|
||||
store::{accounts::AccountList, config::Config},
|
||||
terminal::TerminalServer,
|
||||
};
|
||||
use hypersdk::hypercore::ws::ConnectionStream;
|
||||
use pulse_sdk::prelude::*;
|
||||
use std::{collections::HashMap, sync::Arc};
|
||||
use tokio::{sync::Mutex, task::JoinHandle};
|
||||
@@ -19,8 +20,12 @@ pub struct WatchList {
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct Engine {
|
||||
// engine
|
||||
pub terminal_server: Arc<TerminalServer>,
|
||||
pub strategy_engine: Arc<StrategyEngine>,
|
||||
pub ws_stream: Arc<Mutex<ConnectionStream>>,
|
||||
|
||||
// data
|
||||
pub config: Arc<Mutex<Config>>,
|
||||
pub accounts: Arc<Mutex<AccountList>>,
|
||||
pub watch_list: Arc<Mutex<WatchList>>,
|
||||
@@ -31,16 +36,19 @@ impl Engine {
|
||||
pub async fn new() -> tokio::io::Result<Arc<Self>> {
|
||||
let config = Config::new().await?;
|
||||
|
||||
let strategy = StrategyEngine::new(&config.strategy).await?;
|
||||
let (ws_handle, ws_stream) = hypersdk::hypercore::mainnet_ws().split();
|
||||
|
||||
let strategy = StrategyEngine::new(&config.strategy, ws_handle).await?;
|
||||
|
||||
let accounts = Arc::new(Mutex::new(AccountList::new().await?));
|
||||
let config = Arc::new(Mutex::new(config));
|
||||
|
||||
Ok(Arc::new_cyclic(|engine| Self {
|
||||
terminal_server: TerminalServer::new(engine.clone()),
|
||||
strategy_engine: strategy.initialize(engine.clone()),
|
||||
config,
|
||||
accounts,
|
||||
ws_stream: Arc::new(Mutex::new(ws_stream)),
|
||||
terminal_server: TerminalServer::new(engine.clone()),
|
||||
strategy_engine: strategy.initialize(engine.clone()),
|
||||
watch_list: Arc::new(Mutex::new(WatchList {
|
||||
name_to_index: HashMap::new(),
|
||||
items: Vec::new(),
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
use anyhow::Context;
|
||||
use hypersdk::hypercore::{self, CandleInterval, Subscription, WebSocket};
|
||||
use hypersdk::hypercore::{self, CandleInterval, Subscription, ws::ConnectionHandle};
|
||||
use pulse_sdk::prelude::*;
|
||||
use std::{
|
||||
collections::HashSet,
|
||||
@@ -19,12 +19,12 @@ use crate::{
|
||||
pub struct StrategyEngine {
|
||||
pub engine: Weak<Engine>,
|
||||
pub strategy: Mutex<StrategyChild>,
|
||||
pub ws: WebSocket,
|
||||
pub ws_handle: ConnectionHandle,
|
||||
pub subscriptions: Mutex<HashSet<Subscription>>,
|
||||
}
|
||||
|
||||
impl StrategyEngine {
|
||||
pub async fn new(strategy_id: &str) -> tokio::io::Result<Self> {
|
||||
pub async fn new(strategy_id: &str, ws_handle: ConnectionHandle) -> tokio::io::Result<Self> {
|
||||
let strategy = pulse_strategy(strategy_id)?;
|
||||
|
||||
let (strategy, strategy_manifest) = get_manifest(
|
||||
@@ -42,7 +42,7 @@ impl StrategyEngine {
|
||||
Ok(Self {
|
||||
strategy: Mutex::new(StrategyChild::new(strategy, strategy_manifest)),
|
||||
engine: Weak::new(),
|
||||
ws: hypercore::mainnet_ws(),
|
||||
ws_handle,
|
||||
subscriptions: Mutex::new(HashSet::new()),
|
||||
})
|
||||
}
|
||||
@@ -142,18 +142,18 @@ impl StrategyEngine {
|
||||
}
|
||||
|
||||
Some(StrategyMessage::Subscribe(subscription)) => {
|
||||
self.ws.subscribe(subscription.clone());
|
||||
self.ws_handle.subscribe(subscription.clone());
|
||||
self.subscriptions.lock().await.insert(subscription);
|
||||
}
|
||||
|
||||
Some(StrategyMessage::Unsubscribe(subscription)) => {
|
||||
self.subscriptions.lock().await.remove(&subscription);
|
||||
self.ws.unsubscribe(subscription);
|
||||
self.ws_handle.unsubscribe(subscription);
|
||||
}
|
||||
|
||||
Some(StrategyMessage::UnsubscribeAll) => {
|
||||
for sub in self.subscriptions.lock().await.drain() {
|
||||
self.ws.unsubscribe(sub);
|
||||
self.ws_handle.unsubscribe(sub);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user