From 07611375badf3fda88bcdbd83d16a6580110235a Mon Sep 17 00:00:00 2001 From: Prometheus1400 Date: Sat, 29 Nov 2025 09:18:41 -0600 Subject: [PATCH 1/6] change communication module to comm --- cli/src/actors/client.rs | 8 ++++---- cli/src/main.rs | 6 +++--- core/src/{communication.rs => comm.rs} | 0 core/src/lib.rs | 2 +- daemon/src/actors/client_connection.rs | 16 ++++++++-------- daemon/src/daemon.rs | 8 ++++---- 6 files changed, 20 insertions(+), 20 deletions(-) rename core/src/{communication.rs => comm.rs} (100%) diff --git a/cli/src/actors/client.rs b/cli/src/actors/client.rs index e5958b8..554cdb2 100644 --- a/cli/src/actors/client.rs +++ b/cli/src/actors/client.rs @@ -3,7 +3,7 @@ use std::time::Duration; use bytes::Bytes; use handle_macro::Handle; use remux_core::{ - communication, + comm, events::{CliEvent, DaemonEvent}, }; use tokio::{io::AsyncReadExt, net::UnixStream, sync::mpsc, time::interval}; @@ -88,7 +88,7 @@ impl Client { if let Some(index) = index { let selected_session = self.daemon_state.session_ids[index]; debug!("sending session selection: {selected_session}"); - communication::send_event(&mut self.stream, CliEvent::SwitchSession(selected_session)).await.unwrap(); + comm::send_event(&mut self.stream, CliEvent::SwitchSession(selected_session)).await.unwrap(); } debug!("Returning to normal state"); self.client_state = ClientState::Normal; @@ -97,7 +97,7 @@ impl Client { } } }, - res = communication::recv_daemon_event(&mut self.stream) => { + res = comm::recv_daemon_event(&mut self.stream) => { match res { Ok(event) => { match event { @@ -145,7 +145,7 @@ impl Client { for event in self.input_parser.process(&stdin_buf[..n]) { match event { ParsedEvent::DaemonAction(cli_event) => { - communication::send_event(&mut self.stream, cli_event).await?; + comm::send_event(&mut self.stream, cli_event).await?; }, ParsedEvent::LocalAction(local_action) => { match local_action { diff --git a/cli/src/main.rs b/cli/src/main.rs index 8c29001..4286dcc 100644 --- a/cli/src/main.rs +++ b/cli/src/main.rs @@ -10,7 +10,7 @@ mod widgets; use clap::Parser; use ratatui::crossterm::terminal::{disable_raw_mode, enable_raw_mode}; use remux_core::{ - communication, + comm, daemon_utils::get_sock_path, messages::{RequestMessage, ResponseBody, ResponseMessage}, }; @@ -74,7 +74,7 @@ async fn connect() -> Result { #[instrument(skip(stream))] async fn handle_session_command(mut stream: UnixStream, command: SessionCommands) -> Result<()> { let req: RequestMessage = command.into(); - let res: ResponseMessage = communication::send_and_recv(&mut stream, &req).await?; + let res: ResponseMessage = comm::send_and_recv(&mut stream, &req).await?; match res.body { ResponseBody::SessionsList { sessions } => { println!("{sessions:?}"); @@ -94,7 +94,7 @@ async fn run(command: Commands) -> Result<()> { #[instrument(skip(stream, attach_message))] async fn attach(mut stream: UnixStream, attach_message: RequestMessage) -> Result<()> { debug!("Sending attach request"); - communication::write_message(&mut stream, &attach_message) + comm::write_message(&mut stream, &attach_message) .await .map_err(|source| Error::SendRequestMessage { message: attach_message, diff --git a/core/src/communication.rs b/core/src/comm.rs similarity index 100% rename from core/src/communication.rs rename to core/src/comm.rs diff --git a/core/src/lib.rs b/core/src/lib.rs index 9dfa723..2d34584 100644 --- a/core/src/lib.rs +++ b/core/src/lib.rs @@ -1,4 +1,4 @@ -pub mod communication; +pub mod comm; pub mod constants; pub mod daemon_utils; pub mod error; diff --git a/daemon/src/actors/client_connection.rs b/daemon/src/actors/client_connection.rs index dee41cd..1ecc476 100644 --- a/daemon/src/actors/client_connection.rs +++ b/daemon/src/actors/client_connection.rs @@ -1,6 +1,6 @@ use bytes::Bytes; use handle_macro::Handle; -use remux_core::{communication, events::DaemonEvent}; +use remux_core::{comm, events::DaemonEvent}; use tokio::{net::UnixStream, sync::mpsc}; use tracing::Instrument; @@ -75,11 +75,11 @@ impl ClientConnection { SuccessAttachToSession(session_id) => { debug!("Client: SuccessAttachToSession"); self.state = ClientConnectionState::Attached(session_id); - communication::send_event(&mut self.stream, DaemonEvent::ActiveSession(session_id)).await.unwrap(); + comm::send_event(&mut self.stream, DaemonEvent::ActiveSession(session_id)).await.unwrap(); } FailedAttachToSession{..} => { debug!("Client: FailedAttachToSession"); - communication::send_event(&mut self.stream, DaemonEvent::Disconnected).await.unwrap(); + comm::send_event(&mut self.stream, DaemonEvent::Disconnected).await.unwrap(); } DetachFromSession{..} => { debug!("Client: DetachFromSession"); @@ -87,23 +87,23 @@ impl ClientConnection { } Disconnect => { trace!("Client: Disconnect"); - communication::send_event(&mut self.stream, DaemonEvent::Disconnected).await.unwrap(); + comm::send_event(&mut self.stream, DaemonEvent::Disconnected).await.unwrap(); } SessionOutput(bytes) => { trace!("Client: SessionOutput"); - communication::send_event(&mut self.stream, DaemonEvent::Raw(bytes)).await.unwrap(); + comm::send_event(&mut self.stream, DaemonEvent::Raw(bytes)).await.unwrap(); } NewSession(session_id) => { trace!("Client: NewSession"); - communication::send_event(&mut self.stream, DaemonEvent::NewSession(session_id)).await.unwrap(); + comm::send_event(&mut self.stream, DaemonEvent::NewSession(session_id)).await.unwrap(); } CurrentSessions(session_ids) => { trace!("Client: NewSession"); - communication::send_event(&mut self.stream, DaemonEvent::CurrentSessions(session_ids)).await.unwrap(); + comm::send_event(&mut self.stream, DaemonEvent::CurrentSessions(session_ids)).await.unwrap(); } } }, - res = communication::recv_cli_event(&mut self.stream), if matches!(self.state, ClientConnectionState::Attached(_)) => { + res = comm::recv_cli_event(&mut self.stream), if matches!(self.state, ClientConnectionState::Attached(_)) => { match res { Ok(event) => { match event { diff --git a/daemon/src/daemon.rs b/daemon/src/daemon.rs index e1e2bd6..07d42a5 100644 --- a/daemon/src/daemon.rs +++ b/daemon/src/daemon.rs @@ -1,7 +1,7 @@ use std::fs::{File, remove_file}; use remux_core::{ - communication, + comm, daemon_utils::{get_sock_path, lock_daemon_file}, messages::RequestBody, }; @@ -45,7 +45,7 @@ impl RemuxDaemon { loop { let (stream, _) = listener.accept().await?; info!("accepting connection"); - if let Err(e) = handle_communication(self.session_manager_handle.clone(), stream).await { + if let Err(e) = handle_comm(self.session_manager_handle.clone(), stream).await { error!("{e}"); } } @@ -53,8 +53,8 @@ impl RemuxDaemon { } #[instrument(skip(session_manager_handle, stream))] -async fn handle_communication(session_manager_handle: SessionManagerHandle, mut stream: UnixStream) -> Result<()> { - let req = communication::read_req(&mut stream).await?; +async fn handle_comm(session_manager_handle: SessionManagerHandle, mut stream: UnixStream) -> Result<()> { + let req = comm::read_req(&mut stream).await?; match req.body { RequestBody::Attach { session_id } => { debug!("running new client actor"); From 5a1a070f9b4fedf19ad11ce09ea4563c2fe30f32 Mon Sep 17 00:00:00 2001 From: Prometheus1400 Date: Sat, 29 Nov 2025 12:39:47 -0600 Subject: [PATCH 2/6] wip --- core/src/messages.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/core/src/messages.rs b/core/src/messages.rs index ff0c611..5537ce4 100644 --- a/core/src/messages.rs +++ b/core/src/messages.rs @@ -59,6 +59,7 @@ pub enum RequestBody { #[derive(Serialize, Deserialize, Debug, PartialEq, Display)] #[serde(tag = "type")] pub enum ResponseBody { + AttachResponse, #[display("sessions: {sessions:?}")] SessionsList { sessions: Vec }, } From bc9f5a0f8bb9e82284ec9ef92af63bb7457cb96b Mon Sep 17 00:00:00 2001 From: Prometheus1400 Date: Sun, 30 Nov 2025 22:44:38 -0600 Subject: [PATCH 3/6] wip --- Cargo.lock | 1 + cli/src/actors/client.rs | 12 +- cli/src/actors/lua.rs | 3 +- cli/src/actors/ui.rs | 3 +- cli/src/args.rs | 76 +++++++----- cli/src/error.rs | 5 +- cli/src/main.rs | 70 ++++++----- cli/src/states/mod.rs | 1 - core/Cargo.toml | 1 + core/src/comm.rs | 113 ++++++++++-------- core/src/error.rs | 12 ++ core/src/lib.rs | 2 + core/src/messages.rs | 65 ---------- core/src/messages/mod.rs | 8 ++ core/src/messages/request.rs | 75 ++++++++++++ core/src/messages/response.rs | 59 +++++++++ core/src/messages/traits.rs | 9 ++ core/src/rand.rs | 6 + .../daemon_state.rs => core/src/states.rs | 6 +- daemon/src/actors/client_connection.rs | 73 ++++++++--- daemon/src/actors/session_manager.rs | 39 +++--- daemon/src/daemon.rs | 19 ++- daemon/src/error.rs | 4 +- 23 files changed, 433 insertions(+), 229 deletions(-) delete mode 100644 core/src/messages.rs create mode 100644 core/src/messages/mod.rs create mode 100644 core/src/messages/request.rs create mode 100644 core/src/messages/response.rs create mode 100644 core/src/messages/traits.rs create mode 100644 core/src/rand.rs rename cli/src/states/daemon_state.rs => core/src/states.rs (86%) diff --git a/Cargo.lock b/Cargo.lock index 2a205b6..fbe0739 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -976,6 +976,7 @@ dependencies = [ "bytes", "derive_more", "fs2", + "rand", "serde", "serde_json", "thiserror 2.0.17", diff --git a/cli/src/actors/client.rs b/cli/src/actors/client.rs index 554cdb2..045021b 100644 --- a/cli/src/actors/client.rs +++ b/cli/src/actors/client.rs @@ -5,6 +5,7 @@ use handle_macro::Handle; use remux_core::{ comm, events::{CliEvent, DaemonEvent}, + states::DaemonState, }; use tokio::{io::AsyncReadExt, net::UnixStream, sync::mpsc, time::interval}; use tracing::{Instrument, debug}; @@ -13,7 +14,6 @@ use crate::{ actors::ui::{UI, UIHandle}, input_parser::{Action, InputParser, ParsedEvent}, prelude::*, - states::daemon_state::DaemonState, utils::DisplayableVec, }; @@ -43,12 +43,12 @@ pub struct Client { } impl Client { #[instrument(skip(stream))] - pub fn spawn(stream: UnixStream) -> Result { - Client::new(stream)?.run() + pub fn spawn(stream: UnixStream, daemon_state: DaemonState) -> Result { + Client::new(stream, daemon_state)?.run() } #[instrument(skip(stream))] - fn new(stream: UnixStream) -> Result { + fn new(stream: UnixStream, daemon_state: DaemonState) -> Result { let (tx, rx) = mpsc::channel(100); let (ui_stdin_tx, ui_stdin_rx) = mpsc::channel(100); let handle = ClientHandle { tx }; @@ -59,8 +59,8 @@ impl Client { rx, ui_stdin_tx, ui_handle, - daemon_state: DaemonState::default(), - sync_daemon_state: false, + daemon_state, + sync_daemon_state: true, input_parser: InputParser::new(), client_state: ClientState::Normal, }) diff --git a/cli/src/actors/lua.rs b/cli/src/actors/lua.rs index 4a331f3..fb9c83b 100644 --- a/cli/src/actors/lua.rs +++ b/cli/src/actors/lua.rs @@ -8,12 +8,13 @@ use std::{ }; use mlua::Lua as MLua; +use remux_core::states::DaemonState; use tokio::runtime::Handle; use crate::{ actors::ui::UIHandle, prelude::*, - states::{daemon_state::DaemonState, status_line_state::StatusLineState}, + states::{status_line_state::StatusLineState}, }; pub enum LuaEvent { diff --git a/cli/src/actors/ui.rs b/cli/src/actors/ui.rs index 1e0a964..01d3f2b 100644 --- a/cli/src/actors/ui.rs +++ b/cli/src/actors/ui.rs @@ -11,6 +11,7 @@ use crossterm::{ }; use handle_macro::Handle; use ratatui::{Terminal, prelude::CrosstermBackend}; +use remux_core::states::DaemonState; use tokio::{ sync::{broadcast, mpsc}, time::interval, @@ -25,7 +26,7 @@ use crate::{ lua::{Lua, LuaHandle}, }, prelude::*, - states::{daemon_state::DaemonState, status_line_state::StatusLineState}, + states::{status_line_state::StatusLineState}, utils::DisplayableVec, widgets::{BasicSelector, FuzzySelector, Selector}, }; diff --git a/cli/src/args.rs b/cli/src/args.rs index 7ccecea..b0e0ab6 100644 --- a/cli/src/args.rs +++ b/cli/src/args.rs @@ -1,5 +1,5 @@ use clap::{Parser, Subcommand}; -use remux_core::messages::{RequestBody, RequestMessage}; +use remux_core::messages::{RequestBody, RequestBuilder, CliRequestMessage, request}; #[derive(Parser, Debug)] pub struct Args { @@ -24,34 +24,54 @@ pub enum SessionCommands { List, } -#[allow(clippy::from_over_into)] -impl Into for Commands { - fn into(self) -> RequestBody { +impl Commands { + pub fn into_request(self) -> CliRequestMessage { match self { - Self::Attach { session_id } => RequestBody::Attach { session_id }, - Self::Session { action } => action.into(), + Self::Attach { session_id } => RequestBuilder::default() + .body(request::Attach { + session_id, + create: true, + }) + .build(), + Self::Session { .. } => todo!(), } } } -#[allow(clippy::from_over_into)] -impl Into for SessionCommands { - fn into(self) -> RequestBody { - match self { - SessionCommands::List => RequestBody::SessionsList, - } - } -} -#[allow(clippy::from_over_into)] -impl Into for Commands { - fn into(self) -> RequestMessage { - let body: RequestBody = self.into(); - RequestMessage::body(body) - } -} -#[allow(clippy::from_over_into)] -impl Into for SessionCommands { - fn into(self) -> RequestMessage { - let body: RequestBody = self.into(); - RequestMessage::body(body) - } -} + +// #[allow(clippy::from_over_into)] +// impl Into> for Commands { +// fn into(self) -> RequestMessage { +// match self { +// Self::Attach { session_id } => RequestBuilder::default() +// .body(request::Attach { +// session_id, +// create: true, +// }) +// .build(), +// Self::Session { action } => action.into(), +// } +// } +// } + +// #[allow(clippy::from_over_into)] +// impl Into for SessionCommands { +// fn into(self) -> RequestBody { +// match self { +// SessionCommands::List => RequestBody::SessionsList, +// } +// } +// } +// #[allow(clippy::from_over_into)] +// impl Into for Commands { +// fn into(self) -> RequestBody { +// let body: RequestBody = self.into(); +// RequestBuilder::default().body(body).build() +// } +// } +// #[allow(clippy::from_over_into)] +// impl Into for SessionCommands { +// fn into(self) -> RequestBody { +// let body: RequestBody = self.into(); +// RequestBuilder::default().body(body).build() +// } +// } diff --git a/cli/src/error.rs b/cli/src/error.rs index 895ef4e..200d2b6 100644 --- a/cli/src/error.rs +++ b/cli/src/error.rs @@ -1,5 +1,4 @@ use bytes::Bytes; -use remux_core::messages::RequestMessage; use thiserror::Error; use tokio::sync::mpsc::error::SendError; @@ -30,9 +29,9 @@ pub enum Error { source: std::io::Error, }, - #[error("Error sending message {message}: {source}")] + #[error("Error sending message {message:?}: {source}")] SendRequestMessage { - message: RequestMessage, + message: String, source: remux_core::error::Error, }, diff --git a/cli/src/main.rs b/cli/src/main.rs index 4286dcc..93ad54d 100644 --- a/cli/src/main.rs +++ b/cli/src/main.rs @@ -12,13 +12,16 @@ use ratatui::crossterm::terminal::{disable_raw_mode, enable_raw_mode}; use remux_core::{ comm, daemon_utils::get_sock_path, - messages::{RequestMessage, ResponseBody, ResponseMessage}, + messages::{ + CliRequestMessage, RequestBuilder, + request::{self, Attach}, + }, }; use tokio::net::UnixStream; use crate::{ actors::Client, - args::{Args, Commands, SessionCommands}, + args::{Args, Commands}, error::{Error, Result}, prelude::*, }; @@ -71,39 +74,52 @@ async fn connect() -> Result { }) } -#[instrument(skip(stream))] -async fn handle_session_command(mut stream: UnixStream, command: SessionCommands) -> Result<()> { - let req: RequestMessage = command.into(); - let res: ResponseMessage = comm::send_and_recv(&mut stream, &req).await?; - match res.body { - ResponseBody::SessionsList { sessions } => { - println!("{sessions:?}"); - } - } - Ok(()) -} +// #[instrument(skip(stream))] +// async fn handle_session_command(mut stream: UnixStream, command: SessionCommands) -> Result<()> { +// let req: Request = command.into(); +// let res: Response = comm::send_and_recv(&mut stream, &req).await?; +// assert_eq!(res.status, ResponseStatus::Ok); +// match res.body { +// ResponseBody::SessionsList { sessions } => { +// println!("{sessions:?}"); +// } +// _ => {} +// } +// Ok(()) +// } +#[instrument] async fn run(command: Commands) -> Result<()> { let stream = connect().await?; + debug!("Running command: {:?}", command); match command { - a @ Commands::Attach { .. } => attach(stream, a.into()).await, - Commands::Session { action } => handle_session_command(stream, action).await, + Commands::Attach { session_id } => { + attach( + stream, + RequestBuilder::default() + .body(request::Attach { + session_id, + create: true, + }) + .build(), + ) + .await + } + // Commands::Session { action } => handle_session_command(stream, action).await, + _ => todo!(), } } -#[instrument(skip(stream, attach_message))] -async fn attach(mut stream: UnixStream, attach_message: RequestMessage) -> Result<()> { - debug!("Sending attach request"); - comm::write_message(&mut stream, &attach_message) - .await - .map_err(|source| Error::SendRequestMessage { - message: attach_message, - source, - })?; - debug!("Sent attach request successfully"); +#[instrument(skip(stream, attach_request))] +async fn attach(mut stream: UnixStream, attach_request: CliRequestMessage) -> Result<()> { + debug!("Sending attach request: {:?}", attach_request); + let res = comm::send_and_recv_message(&mut stream, &attach_request).await?; + debug!("Recieved attach response: {:?}", res); + debug!("Recieved initial daemon state: {:?}", res.initial_daemon_state); + enable_raw_mode()?; - debug!("raw mode enabled"); - if let Ok(task) = Client::spawn(stream) { + debug!("Enabled raw mode"); + if let Ok(task) = Client::spawn(stream, res.initial_daemon_state) { match task.await { Ok(Err(e)) => { error!("Error joining client task: {e}"); diff --git a/cli/src/states/mod.rs b/cli/src/states/mod.rs index 3f5c944..acc87c5 100644 --- a/cli/src/states/mod.rs +++ b/cli/src/states/mod.rs @@ -1,2 +1 @@ -pub mod daemon_state; pub mod status_line_state; diff --git a/core/Cargo.toml b/core/Cargo.toml index 0206789..4ba58ff 100644 --- a/core/Cargo.toml +++ b/core/Cargo.toml @@ -14,3 +14,4 @@ tokio.workspace = true bincode = "2.0.1" fs2 = "0.4.3" +rand = "0.9.2" diff --git a/core/src/comm.rs b/core/src/comm.rs index ab0ab93..b187147 100644 --- a/core/src/comm.rs +++ b/core/src/comm.rs @@ -1,12 +1,13 @@ -use serde::{Serialize, de::DeserializeOwned}; +use serde::{Deserialize, Serialize, de::DeserializeOwned}; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, net::UnixStream, }; use crate::{ + error::ResponseError, events::{CliEvent, DaemonEvent}, - messages::{Message, RequestMessage, ResponseMessage}, + messages::{CliRequestMessage, Message, RequestBody, ResponseMessage, ResponseResult}, prelude::*, }; @@ -37,48 +38,43 @@ async fn recv_event(stream: &mut UnixStream) -> Result { Ok(serde_json::from_slice(&message_bytes)?) } -pub async fn send_and_recv(stream: &mut UnixStream, message: &Req) -> Result -where - Req: Message + Serialize, - Res: Message + DeserializeOwned, -{ - write_message(stream, message).await?; - read_message(stream).await -} - -/// writes a serializable message and returns the request id -pub async fn write_message(stream: &mut UnixStream, message: &M) -> Result -where - M: Message + Serialize, -{ +pub async fn send_message(stream: &mut UnixStream, message: &impl Message) -> Result<()> { let bytes = serde_json::to_vec(message)?; let num_bytes = bytes.len() as u32; let _written = stream.write(&num_bytes.to_be_bytes()).await?; let _written = stream.write(&bytes).await?; - Ok(message.get_id()) + Ok(()) } -pub async fn read_message(stream: &mut UnixStream) -> Result -where - M: Message + DeserializeOwned, -{ +pub async fn read_message(stream: &mut UnixStream) -> Result { + println!("reading message"); let mut num_bytes = [0u8; 4]; stream.read_exact(&mut num_bytes).await?; let num_bytes = u32::from_be_bytes(num_bytes); - let mut message_bytes = vec![0u8; num_bytes as usize]; stream.read_exact(&mut message_bytes).await?; - - Ok(serde_json::from_slice(&message_bytes)?) + println!("reading message2"); + let res = serde_json::from_slice(&message_bytes)?; + println!("here"); + Ok(res) } -pub async fn read_req(stream: &mut UnixStream) -> Result { - read_message(stream).await -} - -pub async fn read_res(stream: &mut UnixStream) -> Result { - read_message(stream).await +pub async fn send_and_recv_message(stream: &mut UnixStream, req: &CliRequestMessage) -> Result +where + B: RequestBody + Serialize + for<'de> Deserialize<'de>, +{ + let req_id = req.id; + send_message(stream, req).await?; + let res: ResponseMessage = read_message(stream).await?; + let res_id = res.id; + // if req_id != res_id { + // return Err(Error::Response(ResponseError::UnexpectedId { expected: req_id, actual: res_id })); + // } + match res.result { + ResponseResult::Success(body) => Ok(body), + ResponseResult::Failure(msg) => Err(Error::Response(ResponseError::Status(msg))), + } } #[cfg(test)] @@ -89,7 +85,15 @@ mod test { use tokio::net::UnixListener; use super::*; - use crate::{constants::TEMP_SOCK_DIR, messages::RequestBody}; + use crate::{ + constants::TEMP_SOCK_DIR, + messages::{ + RequestBuilder, ResponseBuilder, + request::{self, DaemonRequestMessage, DaemonRequestMessageBody}, + response, + }, + states::DaemonState, + }; #[tokio::test] async fn test_tcp_message() -> Result<()> { @@ -102,31 +106,40 @@ mod test { let listener = UnixListener::bind(temp_dir)?; let addr = listener.local_addr()?; + let attach = request::Attach { + session_id: 1, + create: true, + }; + let cli_req = RequestBuilder::default().body(attach.clone()).build(); + let daemon_req = DaemonRequestMessage { + id: cli_req.id, + body: DaemonRequestMessageBody::Attach(attach), + }; + + let attach_response = response::Attach { + initial_daemon_state: DaemonState::default(), + }; + let res = ResponseBuilder::default() + .result(ResponseResult::Success(attach_response.clone())) + .build(); + // Spawn server - let server: tokio::task::JoinHandle> = tokio::spawn(async move { - let (mut socket, _) = listener.accept().await?; - let msg1 = read_req(&mut socket).await?; - assert_eq!( - msg1, - RequestMessage { - id: msg1.get_id(), - body: RequestBody::Attach { session_id: 1 } - } - ); - - Ok(()) + let server: tokio::task::JoinHandle> = tokio::spawn({ + let res = res.clone(); + async move { + let (mut socket, _) = listener.accept().await?; + let msg1 = read_message::(&mut socket).await.unwrap(); + assert_eq!(msg1, daemon_req); + send_message(&mut socket, &res).await.unwrap(); + Ok(()) + } }); // Connect client let mut client = UnixStream::connect(addr.as_pathname().unwrap()).await.unwrap(); - write_message( - &mut client, - &RequestMessage::body(RequestBody::Attach { session_id: 1 }), - ) - .await - .unwrap(); + let res1 = send_and_recv_message(&mut client, &cli_req).await.unwrap(); + assert_eq!(res1, attach_response); server.await.unwrap()?; - Ok(()) } } diff --git a/core/src/error.rs b/core/src/error.rs index b753723..95b66ad 100644 --- a/core/src/error.rs +++ b/core/src/error.rs @@ -12,4 +12,16 @@ pub enum Error { #[error("Serialization error: {0}")] SerializationError(#[from] serde_json::Error), + + #[error("Response Error: {0}")] + Response(ResponseError), +} + +#[derive(Error, Debug)] +pub enum ResponseError { + #[error("UnexpectedId: expected({expected}) actual({actual})")] + UnexpectedId{expected:u32, actual:u32}, + #[error("Bad Status: {0}")] + Status(String) } + diff --git a/core/src/lib.rs b/core/src/lib.rs index 2d34584..30e616e 100644 --- a/core/src/lib.rs +++ b/core/src/lib.rs @@ -1,4 +1,6 @@ pub mod comm; +pub mod states; +pub mod rand; pub mod constants; pub mod daemon_utils; pub mod error; diff --git a/core/src/messages.rs b/core/src/messages.rs deleted file mode 100644 index 5537ce4..0000000 --- a/core/src/messages.rs +++ /dev/null @@ -1,65 +0,0 @@ -use derive_more::Display; -use serde::{Deserialize, Serialize}; - -pub use crate::error::{Error, Result}; - -pub trait Message { - fn get_id(&self) -> u32; -} - -#[derive(Serialize, Deserialize, Debug, PartialEq, Display)] -#[display("request({id}, {body})")] -pub struct RequestMessage { - pub id: u32, - pub body: RequestBody, -} -impl RequestMessage { - pub fn new(id: u32, body: RequestBody) -> Self { - Self { id, body } - } - pub fn body(body: RequestBody) -> Self { - let id = 1; // TODO: make this randomly generated - Self { id, body } - } -} -impl Message for RequestMessage { - fn get_id(&self) -> u32 { - self.id - } -} - -#[derive(Serialize, Deserialize, Debug, PartialEq, Display)] -#[display("response({id}, {body})")] -pub struct ResponseMessage { - id: u32, - pub body: ResponseBody, -} -impl ResponseMessage { - pub fn new(id: u32, body: ResponseBody) -> Self { - Self { id, body } - } -} -impl Message for ResponseMessage { - fn get_id(&self) -> u32 { - self.id - } -} - -#[derive(Serialize, Deserialize, Debug, PartialEq, Display)] -#[serde(tag = "type")] -pub enum RequestBody { - #[display("Attach: {{session_id: {session_id}}}")] - Attach { - session_id: u32, - }, - // session commands - SessionsList, -} - -#[derive(Serialize, Deserialize, Debug, PartialEq, Display)] -#[serde(tag = "type")] -pub enum ResponseBody { - AttachResponse, - #[display("sessions: {sessions:?}")] - SessionsList { sessions: Vec }, -} diff --git a/core/src/messages/mod.rs b/core/src/messages/mod.rs new file mode 100644 index 0000000..d4b1a17 --- /dev/null +++ b/core/src/messages/mod.rs @@ -0,0 +1,8 @@ +pub mod request; +pub mod response; +mod traits; + +pub use request::{RequestBuilder, CliRequestMessage}; +pub use response::{ResponseMessage, ResponseBuilder, ResponseResult}; +pub use traits::RequestBody; +pub use traits::Message; diff --git a/core/src/messages/request.rs b/core/src/messages/request.rs new file mode 100644 index 0000000..8f65550 --- /dev/null +++ b/core/src/messages/request.rs @@ -0,0 +1,75 @@ +use serde::{Deserialize, Serialize}; + +use crate::{ + messages::{response, traits::{Message, RequestBody}}, + rand, +}; + +// --------- serialized from the client --------- // + +#[derive(Clone, Serialize, Deserialize, Debug, PartialEq)] +pub struct CliRequestMessage { + pub id: u32, + pub body: T, +} +impl Deserialize<'de>> Message for CliRequestMessage {} + +// --------- deserialized in the daemon --------- // + +#[derive(Serialize, Deserialize, Debug, PartialEq)] +pub struct DaemonRequestMessage { + pub id: u32, + pub body: DaemonRequestMessageBody, +} +#[derive(Serialize, Deserialize, Debug, PartialEq)] +#[serde(untagged)] +pub enum DaemonRequestMessageBody { + Attach(Attach) +} +impl Message for DaemonRequestMessage {} + +// --------- message bodies --------- // + +#[derive(Clone, Serialize, Deserialize, Debug, PartialEq)] +pub struct Attach { + pub session_id: u32, + pub create: bool, +} +impl RequestBody for Attach { + type ResponseBody = response::Attach; +} + +// --------- builder --------- // + +pub struct BodyUnset; +pub type BodySet = T; + +#[derive(Debug)] +pub struct RequestBuilder { + id: u32, + body: BodyState, +} + +impl Default for RequestBuilder { + fn default() -> Self { + Self { + id: rand::generate_id(), + body: BodyUnset, + } + } +} + +impl RequestBuilder { + pub fn body(self, body: T) -> RequestBuilder> { + RequestBuilder { id: self.id, body } + } +} + +impl RequestBuilder> { + pub fn build(self) -> CliRequestMessage { + CliRequestMessage { + id: self.id, + body: self.body, + } + } +} diff --git a/core/src/messages/response.rs b/core/src/messages/response.rs new file mode 100644 index 0000000..ac0fcc0 --- /dev/null +++ b/core/src/messages/response.rs @@ -0,0 +1,59 @@ +use serde::{Deserialize, Serialize}; + +use crate::{messages::traits::Message, rand, states::DaemonState}; + +#[derive(Clone, Serialize, Deserialize, Debug, PartialEq)] +pub struct ResponseMessage { + pub id: u32, + pub result: ResponseResult, +} +impl Deserialize<'de>> Message for ResponseMessage {} + +#[derive(Clone, Serialize, Deserialize, Debug, PartialEq)] +#[serde(tag = "type")] +pub enum ResponseResult { + Success(T), + Failure(String), +} + +// --------- message bodies --------- // + +#[derive(Clone, Serialize, Deserialize, Debug, PartialEq)] +pub struct Attach { + pub initial_daemon_state: DaemonState +} + + +// --------- builder --------- // + +pub struct ResultUnset; +pub type ResultSet = ResponseResult; + +pub struct ResponseBuilder { + id: u32, + result: ResultState, +} + +impl Default for ResponseBuilder { + fn default() -> Self { + Self { + id: rand::generate_id(), + result: ResultUnset, + } + } +} + +impl ResponseBuilder { + pub fn result(self, result: ResponseResult) -> ResponseBuilder> { + ResponseBuilder { id: self.id, result } + } +} + +impl ResponseBuilder> { + pub fn build(self) -> ResponseMessage { + ResponseMessage { + id: self.id, + result: self.result + } + } +} diff --git a/core/src/messages/traits.rs b/core/src/messages/traits.rs new file mode 100644 index 0000000..587ab4e --- /dev/null +++ b/core/src/messages/traits.rs @@ -0,0 +1,9 @@ +use std::fmt::Debug; + +use serde::{Deserialize, Serialize, de::DeserializeOwned}; + +pub trait Message: Serialize + DeserializeOwned + for<'de> Deserialize<'de> {} + +pub trait RequestBody: { + type ResponseBody: Serialize + DeserializeOwned + for<'de> Deserialize<'de>; +} diff --git a/core/src/rand.rs b/core/src/rand.rs new file mode 100644 index 0000000..c608eca --- /dev/null +++ b/core/src/rand.rs @@ -0,0 +1,6 @@ +use rand::Rng; + +pub fn generate_id() -> u32 { + let mut rng = rand::rng(); + rng.random() +} diff --git a/cli/src/states/daemon_state.rs b/core/src/states.rs similarity index 86% rename from cli/src/states/daemon_state.rs rename to core/src/states.rs index 4dc822a..0a4699d 100644 --- a/cli/src/states/daemon_state.rs +++ b/core/src/states.rs @@ -1,5 +1,7 @@ -// clients view of the state -#[derive(Debug, Clone, Default)] +/// comprehensive summary of the state of the daemon +use serde::{Deserialize, Serialize}; + +#[derive(Default, Serialize, Deserialize, Debug, Clone, PartialEq)] pub struct DaemonState { pub session_ids: Vec, pub active_session: Option, diff --git a/daemon/src/actors/client_connection.rs b/daemon/src/actors/client_connection.rs index 1ecc476..05fea73 100644 --- a/daemon/src/actors/client_connection.rs +++ b/daemon/src/actors/client_connection.rs @@ -1,15 +1,20 @@ use bytes::Bytes; use handle_macro::Handle; -use remux_core::{comm, events::DaemonEvent}; +use remux_core::{ + comm, + events::DaemonEvent, + messages::{ResponseBuilder, ResponseResult, response}, + states::DaemonState, +}; use tokio::{net::UnixStream, sync::mpsc}; use tracing::Instrument; use crate::{actors::session_manager::SessionManagerHandle, layout::SplitDirection, prelude::*}; #[allow(unused)] -#[derive(Handle)] +#[derive(Handle, Debug)] pub enum ClientConnectionEvent { - AttachToSession(u32), + // AttachToSession(u32), SuccessAttachToSession(u32), FailedAttachToSession(u32), DetachFromSession(u32), @@ -17,14 +22,19 @@ pub enum ClientConnectionEvent { NewSession(u32), CurrentSessions(Vec), Disconnect, + + // variants related to initialization phase + InitialAttach(u32), // invoked directly by the daemon + // this variant is unique in that it responds to client by sending a message not an event + InitialAttachResult(Result), } use ClientConnectionEvent::*; #[allow(unused)] +#[derive(Debug)] enum ClientConnectionState { Unattached, - Attaching(u32), - Attached(u32), + Attached, } pub struct ClientConnection { @@ -37,15 +47,20 @@ pub struct ClientConnection { } impl ClientConnection { #[instrument(skip(stream, session_manager_handle))] - pub fn spawn(stream: UnixStream, session_manager_handle: SessionManagerHandle) -> Result { + pub fn spawn( + stream: UnixStream, + session_manager_handle: SessionManagerHandle, + connecting_session_id: u32, + ) -> Result { let client = Self::new(stream, session_manager_handle); - client.run() + client.run(connecting_session_id) } #[instrument(skip(stream, session_manager_handle))] fn new(stream: UnixStream, session_manager_handle: SessionManagerHandle) -> Self { let (tx, rx) = mpsc::channel(10); let handle = ClientConnectionHandle { tx }; let id: u32 = rand::random(); + Self { id, stream, @@ -56,25 +71,41 @@ impl ClientConnection { } } #[instrument(skip(self), fields(client_id = self.id))] - fn run(mut self) -> crate::error::Result { + fn run(mut self, initial_session_id: u32) -> crate::error::Result { let span = tracing::Span::current(); let handle_clone = self.handle.clone(); + let _task = tokio::spawn({ async move { let handle = self.handle.clone(); + self.session_manager_handle.client_connect(self.id, handle.clone(), initial_session_id, true).await?; loop { use remux_core::events::CliEvent; tokio::select! { Some(event) = self.rx.recv() => { match event { - AttachToSession(session_id) => { - debug!("Client: AttachToSession"); - self.session_manager_handle.client_connect(self.id, handle.clone(), session_id, true).await.unwrap(); - self.state = ClientConnectionState::Attaching(session_id); + InitialAttachResult(result) if matches!(self.state, ClientConnectionState::Unattached) => { + debug!("Client: InitialAttachResult({result:?})"); + match result { + Ok(daemon_state) => { + let res = ResponseBuilder::default().result(ResponseResult::Success(response::Attach{initial_daemon_state: daemon_state})).build(); + debug!("Sending response: {res:?}"); + comm::send_message(&mut self.stream, &res).await.unwrap(); + self.state = ClientConnectionState::Attached; + } + Err(e) => { + comm::send_message(&mut self.stream, &ResponseBuilder::default().result(ResponseResult::Failure::<()>(e.to_string())).build()).await.unwrap(); + } + } } + // AttachToSession(session_id) => { + // debug!("Client: AttachToSession"); + // self.session_manager_handle.client_connect(self.id, handle.clone(), session_id, true).await.unwrap(); + // self.state = ClientConnectionState::Attaching; + // } SuccessAttachToSession(session_id) => { debug!("Client: SuccessAttachToSession"); - self.state = ClientConnectionState::Attached(session_id); + self.state = ClientConnectionState::Attached; comm::send_event(&mut self.stream, DaemonEvent::ActiveSession(session_id)).await.unwrap(); } FailedAttachToSession{..} => { @@ -94,16 +125,21 @@ impl ClientConnection { comm::send_event(&mut self.stream, DaemonEvent::Raw(bytes)).await.unwrap(); } NewSession(session_id) => { - trace!("Client: NewSession"); - comm::send_event(&mut self.stream, DaemonEvent::NewSession(session_id)).await.unwrap(); + // trace!("Client: NewSession"); + // comm::send_event(&mut self.stream, DaemonEvent::NewSession(session_id)).await.unwrap(); } CurrentSessions(session_ids) => { - trace!("Client: NewSession"); - comm::send_event(&mut self.stream, DaemonEvent::CurrentSessions(session_ids)).await.unwrap(); + // trace!("Client: NewSession"); + // comm::send_event(&mut self.stream, DaemonEvent::CurrentSessions(session_ids)).await.unwrap(); + } + _ => { + error!("Unhandled or invalid event '{:?}' for current state '{:?}'", event, self.state); + panic!("Unhandled or invalid event '{:?}' for current state '{:?}'", event, self.state); + // TODO: better error handling } } }, - res = comm::recv_cli_event(&mut self.stream), if matches!(self.state, ClientConnectionState::Attached(_)) => { + res = comm::recv_cli_event(&mut self.stream), if matches!(self.state, ClientConnectionState::Attached) => { match res { Ok(event) => { match event { @@ -151,6 +187,7 @@ impl ClientConnection { } } } + Ok::<(), Error>(()) }.instrument(span) }); diff --git a/daemon/src/actors/session_manager.rs b/daemon/src/actors/session_manager.rs index 75ec7a5..5f9429b 100644 --- a/daemon/src/actors/session_manager.rs +++ b/daemon/src/actors/session_manager.rs @@ -2,6 +2,7 @@ use std::{collections::HashMap, vec}; use bytes::Bytes; use handle_macro::Handle; +use remux_core::states::DaemonState; use tokio::sync::mpsc; use tracing::Instrument; @@ -65,6 +66,7 @@ pub struct SessionManager { clients: HashMap, session_to_client_mapping: HashMap>, // support multiple clients attached to same session client_to_session_mapping: HashMap, // one client can only attach to one session + daemon_state: DaemonState, } impl SessionManager { #[instrument] @@ -84,6 +86,7 @@ impl SessionManager { clients: HashMap::new(), session_to_client_mapping: HashMap::new(), client_to_session_mapping: HashMap::new(), + daemon_state: DaemonState::default(), } } @@ -152,30 +155,36 @@ impl SessionManager { session_id: u32, create_session: bool, ) -> Result<()> { - // session doesn't exist either send client error or create it + // session doesn't exist and client not trying to create it + if !self.sessions.contains_key(&session_id) && !create_session { + client_handle + .initial_attach_result(Err(Error::Custom("no such session".to_owned()))) + .await?; + return Ok(()); + } + + // session didn't exist if !self.sessions.contains_key(&session_id) { - if create_session { - let new_session = Session::spawn(session_id, self.handle.clone()).unwrap(); - self.sessions.insert(session_id, new_session); - } else { - client_handle.failed_attach_to_session(session_id).await?; - } + let new_session = Session::spawn(session_id, self.handle.clone()).unwrap(); + self.sessions.insert(session_id, new_session); } + // session exists self.clients.insert(client_id, client_handle.clone()); let clients = self.session_to_client_mapping.entry(session_id).or_insert(vec![]); clients.push(client_id); - for c in self.clients.values() { - c.new_session(session_id).await?; - } + // for c in self.clients.values() { + // c.new_session(session_id).await?; + // } self.client_to_session_mapping.insert(client_id, session_id); - client_handle - .current_sessions(self.sessions.keys().copied().collect()) - .await?; + // client_handle + // .current_sessions(self.sessions.keys().copied().collect()) + // .await?; let session_handle = self.sessions.get_mut(&session_id).expect("session should exist here"); - client_handle.success_attach_to_session(session_id).await?; - session_handle.user_connection().await + client_handle.initial_attach_result(Ok(self.daemon_state.clone())).await?; + Ok(()) + // session_handle.user_connection().await } async fn handle_client_disconnect(&mut self, client_id: u32) -> Result<()> { let client_handle = self.clients.remove(&client_id); diff --git a/daemon/src/daemon.rs b/daemon/src/daemon.rs index 07d42a5..8e6e6c3 100644 --- a/daemon/src/daemon.rs +++ b/daemon/src/daemon.rs @@ -3,7 +3,6 @@ use std::fs::{File, remove_file}; use remux_core::{ comm, daemon_utils::{get_sock_path, lock_daemon_file}, - messages::RequestBody, }; use tokio::net::{UnixListener, UnixStream}; @@ -45,7 +44,7 @@ impl RemuxDaemon { loop { let (stream, _) = listener.accept().await?; info!("accepting connection"); - if let Err(e) = handle_comm(self.session_manager_handle.clone(), stream).await { + if let Err(e) = handle_message(self.session_manager_handle.clone(), stream).await { error!("{e}"); } } @@ -53,16 +52,16 @@ impl RemuxDaemon { } #[instrument(skip(session_manager_handle, stream))] -async fn handle_comm(session_manager_handle: SessionManagerHandle, mut stream: UnixStream) -> Result<()> { - let req = comm::read_req(&mut stream).await?; +async fn handle_message(session_manager_handle: SessionManagerHandle, mut stream: UnixStream) -> Result<()> { + use remux_core::messages::request::{self, DaemonRequestMessage, DaemonRequestMessageBody}; + + let req: DaemonRequestMessage = comm::read_message(&mut stream).await?; + debug!("Handling message: {req:?}"); match req.body { - RequestBody::Attach { session_id } => { + DaemonRequestMessageBody::Attach(request::Attach { session_id, create }) => { debug!("running new client actor"); - let client = ClientConnection::spawn(stream, session_manager_handle).unwrap(); - client.attach_to_session(session_id).await.unwrap(); - } - RequestBody::SessionsList => { - todo!() + let _client = ClientConnection::spawn(stream, session_manager_handle, session_id).unwrap(); + // client.attach_to_session(session_id).await.unwrap(); } }; Ok(()) diff --git a/daemon/src/error.rs b/daemon/src/error.rs index 8b45387..480d7a1 100644 --- a/daemon/src/error.rs +++ b/daemon/src/error.rs @@ -44,7 +44,7 @@ pub enum Error { #[derive(Error, Debug)] pub enum EventSendError { #[error("Client send error: {0}")] - Client(SendError), + Client(Box>), #[error("Session Manager send error: {0}")] SessionManager(SendError), #[error("Session send error: {0}")] @@ -59,7 +59,7 @@ pub enum EventSendError { impl From> for Error { fn from(e: SendError) -> Self { - Self::EventSend(EventSendError::Client(e)) + Self::EventSend(EventSendError::Client(Box::new(e))) } } impl From> for Error { From 4be17fddf9743862e2338131c53890b5fea745dc Mon Sep 17 00:00:00 2001 From: Prometheus1400 Date: Tue, 2 Dec 2025 07:58:24 -0600 Subject: [PATCH 4/6] wip --- daemon/src/actors/client_connection.rs | 13 ++++------ daemon/src/actors/pty.rs | 4 ++-- daemon/src/actors/session_manager.rs | 33 +++++++++++++++++--------- 3 files changed, 29 insertions(+), 21 deletions(-) diff --git a/daemon/src/actors/client_connection.rs b/daemon/src/actors/client_connection.rs index 05fea73..3b08d2c 100644 --- a/daemon/src/actors/client_connection.rs +++ b/daemon/src/actors/client_connection.rs @@ -19,10 +19,11 @@ pub enum ClientConnectionEvent { FailedAttachToSession(u32), DetachFromSession(u32), SessionOutput(Bytes), - NewSession(u32), - CurrentSessions(Vec), Disconnect, + // client side state update events + NewSession(u32), + // variants related to initialization phase InitialAttach(u32), // invoked directly by the daemon // this variant is unique in that it responds to client by sending a message not an event @@ -125,12 +126,8 @@ impl ClientConnection { comm::send_event(&mut self.stream, DaemonEvent::Raw(bytes)).await.unwrap(); } NewSession(session_id) => { - // trace!("Client: NewSession"); - // comm::send_event(&mut self.stream, DaemonEvent::NewSession(session_id)).await.unwrap(); - } - CurrentSessions(session_ids) => { - // trace!("Client: NewSession"); - // comm::send_event(&mut self.stream, DaemonEvent::CurrentSessions(session_ids)).await.unwrap(); + trace!("Client: NewSession"); + comm::send_event(&mut self.stream, DaemonEvent::NewSession(session_id)).await.unwrap(); } _ => { error!("Unhandled or invalid event '{:?}' for current state '{:?}'", event, self.state); diff --git a/daemon/src/actors/pty.rs b/daemon/src/actors/pty.rs index f053d83..f6a77be 100644 --- a/daemon/src/actors/pty.rs +++ b/daemon/src/actors/pty.rs @@ -42,13 +42,13 @@ pub struct Pty { rect: Rect, } impl Pty { - #[instrument(skip(pane_handle))] + #[instrument(skip(pane_handle, rect))] pub fn spawn(pane_handle: PaneHandle, rect: Rect) -> Result { let pty = Pty::new(pane_handle, rect); pty.run() } - #[instrument(skip(pane_handle))] + #[instrument(skip(pane_handle, rect))] fn new(pane_handle: PaneHandle, rect: Rect) -> Self { let (tx, rx) = mpsc::channel::(10); let (pty_tx, pty_rx) = mpsc::unbounded_channel::(); diff --git a/daemon/src/actors/session_manager.rs b/daemon/src/actors/session_manager.rs index 5f9429b..607d8d4 100644 --- a/daemon/src/actors/session_manager.rs +++ b/daemon/src/actors/session_manager.rs @@ -148,6 +148,18 @@ impl SessionManager { Ok(handle_clone) } + /// creates a new session and handles updating the state and notifying clients about the update + async fn create_session(&mut self, session_id: u32) -> Result<&SessionHandle> { + let new_session = Session::spawn(session_id, self.handle.clone())?; + self.sessions.insert(session_id, new_session); + let session_handle_ref = self.sessions.get(&session_id).unwrap(); + self.daemon_state.add_session(session_id); + for c in self.clients.values() { + c.new_session(session_id).await?; + } + Ok(session_handle_ref) + } + async fn handle_client_connect( &mut self, client_id: u32, @@ -165,26 +177,25 @@ impl SessionManager { // session didn't exist if !self.sessions.contains_key(&session_id) { - let new_session = Session::spawn(session_id, self.handle.clone()).unwrap(); - self.sessions.insert(session_id, new_session); + self.create_session(session_id).await?; } // session exists - self.clients.insert(client_id, client_handle.clone()); let clients = self.session_to_client_mapping.entry(session_id).or_insert(vec![]); + for c in self.clients.values() { + c.new_session(session_id).await?; + } + self.clients.insert(client_id, client_handle.clone()); clients.push(client_id); - // for c in self.clients.values() { - // c.new_session(session_id).await?; - // } self.client_to_session_mapping.insert(client_id, session_id); - // client_handle - // .current_sessions(self.sessions.keys().copied().collect()) - // .await?; let session_handle = self.sessions.get_mut(&session_id).expect("session should exist here"); - client_handle.initial_attach_result(Ok(self.daemon_state.clone())).await?; + client_handle + .initial_attach_result(Ok(self.daemon_state.clone())) + .await?; + client_handle.success_attach_to_session(session_id).await?; + session_handle.redraw().await?; Ok(()) - // session_handle.user_connection().await } async fn handle_client_disconnect(&mut self, client_id: u32) -> Result<()> { let client_handle = self.clients.remove(&client_id); From 8558738b899de8adeba7f17a2b7e6eec729bbfd0 Mon Sep 17 00:00:00 2001 From: Prometheus1400 Date: Tue, 2 Dec 2025 16:33:04 -0800 Subject: [PATCH 5/6] wip --- core/src/messages/request.rs | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/core/src/messages/request.rs b/core/src/messages/request.rs index 8f65550..a4380b1 100644 --- a/core/src/messages/request.rs +++ b/core/src/messages/request.rs @@ -1,7 +1,10 @@ use serde::{Deserialize, Serialize}; use crate::{ - messages::{response, traits::{Message, RequestBody}}, + messages::{ + response, + traits::{Message, RequestBody}, + }, rand, }; @@ -24,7 +27,7 @@ pub struct DaemonRequestMessage { #[derive(Serialize, Deserialize, Debug, PartialEq)] #[serde(untagged)] pub enum DaemonRequestMessageBody { - Attach(Attach) + Attach(Attach), } impl Message for DaemonRequestMessage {} From 9caf9091110436d438b150bed53f3bc62f33ad7d Mon Sep 17 00:00:00 2001 From: Prometheus1400 Date: Tue, 2 Dec 2025 16:33:45 -0800 Subject: [PATCH 6/6] format --- cli/src/actors/lua.rs | 6 +----- cli/src/actors/ui.rs | 2 +- cli/src/args.rs | 2 +- core/src/comm.rs | 4 ++-- core/src/error.rs | 5 ++--- core/src/lib.rs | 4 ++-- core/src/messages/mod.rs | 7 +++---- core/src/messages/response.rs | 5 ++--- core/src/messages/traits.rs | 2 +- 9 files changed, 15 insertions(+), 22 deletions(-) diff --git a/cli/src/actors/lua.rs b/cli/src/actors/lua.rs index fb9c83b..5d30f68 100644 --- a/cli/src/actors/lua.rs +++ b/cli/src/actors/lua.rs @@ -11,11 +11,7 @@ use mlua::Lua as MLua; use remux_core::states::DaemonState; use tokio::runtime::Handle; -use crate::{ - actors::ui::UIHandle, - prelude::*, - states::{status_line_state::StatusLineState}, -}; +use crate::{actors::ui::UIHandle, prelude::*, states::status_line_state::StatusLineState}; pub enum LuaEvent { Kill, diff --git a/cli/src/actors/ui.rs b/cli/src/actors/ui.rs index 01d3f2b..daf2706 100644 --- a/cli/src/actors/ui.rs +++ b/cli/src/actors/ui.rs @@ -26,7 +26,7 @@ use crate::{ lua::{Lua, LuaHandle}, }, prelude::*, - states::{status_line_state::StatusLineState}, + states::status_line_state::StatusLineState, utils::DisplayableVec, widgets::{BasicSelector, FuzzySelector, Selector}, }; diff --git a/cli/src/args.rs b/cli/src/args.rs index b0e0ab6..3dc5383 100644 --- a/cli/src/args.rs +++ b/cli/src/args.rs @@ -1,5 +1,5 @@ use clap::{Parser, Subcommand}; -use remux_core::messages::{RequestBody, RequestBuilder, CliRequestMessage, request}; +use remux_core::messages::{CliRequestMessage, RequestBody, RequestBuilder, request}; #[derive(Parser, Debug)] pub struct Args { diff --git a/core/src/comm.rs b/core/src/comm.rs index b187147..9541d64 100644 --- a/core/src/comm.rs +++ b/core/src/comm.rs @@ -117,8 +117,8 @@ mod test { }; let attach_response = response::Attach { - initial_daemon_state: DaemonState::default(), - }; + initial_daemon_state: DaemonState::default(), + }; let res = ResponseBuilder::default() .result(ResponseResult::Success(attach_response.clone())) .build(); diff --git a/core/src/error.rs b/core/src/error.rs index 95b66ad..04c7359 100644 --- a/core/src/error.rs +++ b/core/src/error.rs @@ -20,8 +20,7 @@ pub enum Error { #[derive(Error, Debug)] pub enum ResponseError { #[error("UnexpectedId: expected({expected}) actual({actual})")] - UnexpectedId{expected:u32, actual:u32}, + UnexpectedId { expected: u32, actual: u32 }, #[error("Bad Status: {0}")] - Status(String) + Status(String), } - diff --git a/core/src/lib.rs b/core/src/lib.rs index 30e616e..433c153 100644 --- a/core/src/lib.rs +++ b/core/src/lib.rs @@ -1,9 +1,9 @@ pub mod comm; -pub mod states; -pub mod rand; pub mod constants; pub mod daemon_utils; pub mod error; pub mod events; pub mod messages; mod prelude; +pub mod rand; +pub mod states; diff --git a/core/src/messages/mod.rs b/core/src/messages/mod.rs index d4b1a17..07627d3 100644 --- a/core/src/messages/mod.rs +++ b/core/src/messages/mod.rs @@ -2,7 +2,6 @@ pub mod request; pub mod response; mod traits; -pub use request::{RequestBuilder, CliRequestMessage}; -pub use response::{ResponseMessage, ResponseBuilder, ResponseResult}; -pub use traits::RequestBody; -pub use traits::Message; +pub use request::{CliRequestMessage, RequestBuilder}; +pub use response::{ResponseBuilder, ResponseMessage, ResponseResult}; +pub use traits::{Message, RequestBody}; diff --git a/core/src/messages/response.rs b/core/src/messages/response.rs index ac0fcc0..a58ad58 100644 --- a/core/src/messages/response.rs +++ b/core/src/messages/response.rs @@ -20,10 +20,9 @@ pub enum ResponseResult { #[derive(Clone, Serialize, Deserialize, Debug, PartialEq)] pub struct Attach { - pub initial_daemon_state: DaemonState + pub initial_daemon_state: DaemonState, } - // --------- builder --------- // pub struct ResultUnset; @@ -53,7 +52,7 @@ impl ResponseBuilder> { pub fn build(self) -> ResponseMessage { ResponseMessage { id: self.id, - result: self.result + result: self.result, } } } diff --git a/core/src/messages/traits.rs b/core/src/messages/traits.rs index 587ab4e..b0f8fd6 100644 --- a/core/src/messages/traits.rs +++ b/core/src/messages/traits.rs @@ -4,6 +4,6 @@ use serde::{Deserialize, Serialize, de::DeserializeOwned}; pub trait Message: Serialize + DeserializeOwned + for<'de> Deserialize<'de> {} -pub trait RequestBody: { +pub trait RequestBody { type ResponseBody: Serialize + DeserializeOwned + for<'de> Deserialize<'de>; }