Using postcard
This commit is contained in:
@@ -166,11 +166,11 @@ impl StrategyEngine {
|
||||
None => {}
|
||||
|
||||
Some(RiskMessage::GetWatchList) => {
|
||||
let mut v = vec![1];
|
||||
|
||||
v.extend(engine.config.lock().await.watchlist.to_com());
|
||||
|
||||
self.risk.send_raw(&v).await?;
|
||||
self.risk
|
||||
.send(&&RiskEngineMessage::WatchList(
|
||||
engine.watch_list.lock().await.clone(),
|
||||
))
|
||||
.await?;
|
||||
}
|
||||
|
||||
Some(RiskMessage::Log(mut log)) => {
|
||||
|
||||
@@ -1,16 +1,14 @@
|
||||
use std::marker::PhantomData;
|
||||
|
||||
use pulse_wire::PulseWire;
|
||||
|
||||
use serde::Deserialize;
|
||||
use pulse_wire::map_postcard_err;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tokio::{
|
||||
io::{AsyncReadExt, AsyncWriteExt},
|
||||
process::{Child, ChildStdout},
|
||||
sync::Mutex,
|
||||
};
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct Plugin<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> {
|
||||
pub struct Plugin<S: Serialize, R: for<'de> Deserialize<'de>, M: for<'de> Deserialize<'de>> {
|
||||
pub manifest: Mutex<M>,
|
||||
pub stdout: Mutex<ChildStdout>,
|
||||
pub process: Mutex<Child>,
|
||||
@@ -18,7 +16,7 @@ pub struct Plugin<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> {
|
||||
pub _p: (PhantomData<S>, PhantomData<R>),
|
||||
}
|
||||
|
||||
impl<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> Plugin<S, R, M> {
|
||||
impl<S: Serialize, R: for<'de> Deserialize<'de>, M: for<'de> Deserialize<'de>> Plugin<S, R, M> {
|
||||
pub fn new(mut child: Child, manifest: M) -> Self {
|
||||
Self {
|
||||
stdout: Mutex::new(child.stdout.take().expect("Failed to obtain child stdout")),
|
||||
@@ -45,11 +43,12 @@ impl<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> Plugin<S, R, M> {
|
||||
|
||||
stdout.read_exact(&mut buffer).await?;
|
||||
|
||||
Ok(Some(R::from_com(&mut buffer)))
|
||||
Ok(Some(map_postcard_err(postcard::from_bytes(&buffer))?))
|
||||
}
|
||||
|
||||
pub async fn send(&self, msg: &S) -> tokio::io::Result<()> {
|
||||
self.send_raw(&msg.to_com()).await
|
||||
self.send_raw(&map_postcard_err(postcard::to_allocvec(msg))?)
|
||||
.await
|
||||
}
|
||||
|
||||
pub async fn send_raw(&self, msg: &[u8]) -> tokio::io::Result<()> {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
use crate::engine::Engine;
|
||||
use pulse_wire::prelude::*;
|
||||
use pulse_wire::{map_postcard_err, prelude::*};
|
||||
use std::{
|
||||
collections::HashMap,
|
||||
sync::{Arc, Weak},
|
||||
@@ -104,7 +104,7 @@ impl TerminalServer {
|
||||
|
||||
reader.read_exact(&mut buffer).await?;
|
||||
|
||||
match TerminalClientMessage::from_com(&mut buffer) {
|
||||
match map_postcard_err(postcard::from_bytes(&buffer))? {
|
||||
TerminalClientMessage::ExecuteCommand(command) => {
|
||||
let command = command.as_str();
|
||||
|
||||
@@ -126,7 +126,7 @@ impl TerminalServer {
|
||||
self: &Arc<Self>,
|
||||
message: pulse_wire::terminal::TerminalServerMessage,
|
||||
) -> tokio::io::Result<()> {
|
||||
let msg = message.to_com();
|
||||
let msg = map_postcard_err(postcard::to_allocvec(&message))?;
|
||||
|
||||
let mut clients = self.clients.lock().await;
|
||||
|
||||
@@ -158,7 +158,7 @@ impl TerminalServer {
|
||||
format!("Client({id}) does not exist"),
|
||||
)
|
||||
})?,
|
||||
&message.to_com(),
|
||||
&map_postcard_err(postcard::to_allocvec(&message))?,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
+77
-55
@@ -1,4 +1,5 @@
|
||||
use pulse_wire::prelude::*;
|
||||
use pulse_ui::state::State;
|
||||
use pulse_wire::{map_postcard_err, prelude::*};
|
||||
|
||||
use tokio::{
|
||||
io::{AsyncReadExt, AsyncWriteExt},
|
||||
@@ -29,7 +30,7 @@ impl TerminalClient {
|
||||
&mut self,
|
||||
message: pulse_wire::terminal::TerminalClientMessage,
|
||||
) -> tokio::io::Result<()> {
|
||||
let msg = message.to_com();
|
||||
let msg = map_postcard_err(postcard::to_allocvec(&message))?;
|
||||
self.writer.write(&msg.len().to_le_bytes()).await?;
|
||||
self.writer.write(&msg).await?;
|
||||
self.writer.flush().await?;
|
||||
@@ -42,7 +43,7 @@ impl TerminalClient {
|
||||
|
||||
std::mem::swap(&mut self.reader, &mut reader);
|
||||
|
||||
let mut reader = reader.expect("Reader failed to swap");
|
||||
let reader = reader.expect("Reader failed to swap");
|
||||
|
||||
let watch_list = app.watch_list.clone();
|
||||
let active_positions = app.active_positions.clone();
|
||||
@@ -52,61 +53,82 @@ impl TerminalClient {
|
||||
let status = app.status.clone();
|
||||
let inspect = app.inspect.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
loop {
|
||||
let mut len_buf = [0u8; size_of::<usize>()];
|
||||
reader
|
||||
.read_exact(&mut len_buf)
|
||||
.await
|
||||
.expect("Failed to get header length");
|
||||
|
||||
let len = usize::from_le_bytes(len_buf);
|
||||
|
||||
let mut buffer = vec![0u8; len];
|
||||
|
||||
reader
|
||||
.read_exact(&mut buffer)
|
||||
.await
|
||||
.expect("Failed to read socket");
|
||||
|
||||
match TerminalServerMessage::from_com(&mut buffer) {
|
||||
TerminalServerMessage::WatchListUpdated(v) => {
|
||||
*watch_list.lock().await = v;
|
||||
}
|
||||
|
||||
TerminalServerMessage::PositionsUpdated(v) => {
|
||||
*active_positions.lock().await = v;
|
||||
}
|
||||
|
||||
TerminalServerMessage::StrategyUpdated(v) => {
|
||||
*market_overview.lock().await = Some(v);
|
||||
}
|
||||
|
||||
TerminalServerMessage::SignalsUpdated(v) => {
|
||||
*signals.lock().await = v;
|
||||
}
|
||||
|
||||
TerminalServerMessage::Inspect(v) => {
|
||||
*inspect.lock().await = v;
|
||||
}
|
||||
|
||||
TerminalServerMessage::StatusUpdated(v) => {
|
||||
*status.lock().await = Some(v);
|
||||
}
|
||||
|
||||
TerminalServerMessage::SetLogs(v) => {
|
||||
*logs.lock().await = v;
|
||||
}
|
||||
|
||||
TerminalServerMessage::AddLog(v) => {
|
||||
logs.lock().await.push(v);
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
tokio::spawn(Self::run_client(
|
||||
reader,
|
||||
watch_list,
|
||||
active_positions,
|
||||
logs,
|
||||
signals,
|
||||
market_overview,
|
||||
status,
|
||||
inspect,
|
||||
));
|
||||
|
||||
app.sock = Some(self);
|
||||
|
||||
app
|
||||
}
|
||||
|
||||
pub async fn run_client(
|
||||
mut reader: OwnedReadHalf,
|
||||
|
||||
watch_list: State<Vec<MarketItem>>,
|
||||
active_positions: State<Vec<Position>>,
|
||||
logs: State<Vec<EventLog>>,
|
||||
signals: State<Vec<Signal>>,
|
||||
market_overview: State<Option<Strategy>>,
|
||||
status: State<Option<Status>>,
|
||||
inspect: State<InspectTarget>,
|
||||
) -> tokio::io::Result<()> {
|
||||
let mut len_buf = [0u8; size_of::<usize>()];
|
||||
reader
|
||||
.read_exact(&mut len_buf)
|
||||
.await
|
||||
.expect("Failed to get header length");
|
||||
|
||||
let len = usize::from_le_bytes(len_buf);
|
||||
|
||||
let mut buffer = vec![0u8; len];
|
||||
|
||||
reader
|
||||
.read_exact(&mut buffer)
|
||||
.await
|
||||
.expect("Failed to read socket");
|
||||
|
||||
match map_postcard_err(postcard::from_bytes(&buffer))? {
|
||||
TerminalServerMessage::WatchListUpdated(v) => {
|
||||
*watch_list.lock().await = v;
|
||||
}
|
||||
|
||||
TerminalServerMessage::PositionsUpdated(v) => {
|
||||
*active_positions.lock().await = v;
|
||||
}
|
||||
|
||||
TerminalServerMessage::StrategyUpdated(v) => {
|
||||
*market_overview.lock().await = Some(v);
|
||||
}
|
||||
|
||||
TerminalServerMessage::SignalsUpdated(v) => {
|
||||
*signals.lock().await = v;
|
||||
}
|
||||
|
||||
TerminalServerMessage::Inspect(v) => {
|
||||
*inspect.lock().await = v;
|
||||
}
|
||||
|
||||
TerminalServerMessage::StatusUpdated(v) => {
|
||||
*status.lock().await = Some(v);
|
||||
}
|
||||
|
||||
TerminalServerMessage::SetLogs(v) => {
|
||||
*logs.lock().await = v;
|
||||
}
|
||||
|
||||
TerminalServerMessage::AddLog(v) => {
|
||||
logs.lock().await.push(v);
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user