diff --git a/src/core/backup/backup_engine.rs b/src/core/backup/backup_engine.rs index c94dd35..2492771 100644 --- a/src/core/backup/backup_engine.rs +++ b/src/core/backup/backup_engine.rs @@ -1,14 +1,16 @@ use crate::core::backup::progress_tracker::ProgressTracker; -use crate::core::gui::gui_message_handler::GuiMessageHandler; -use crate::core::infrastructure::actor_system::ActorSystem; use crate::core::infrastructure::app_config::AppConfig; +use crate::core::infrastructure::communication_manager::CommunicationManager; use crate::core::infrastructure::io_manager::IOManager; -use crate::interface::file_system::FileSystemTrait; -use crate::model::core::backup::backup_execution::*; -use crate::model::core::gui::message::GuiMessage; +use crate::interface::communication::command::CommandHandler; +use crate::interface::communication::query::QueryHandler; +use crate::interface::core::file_system::FileSystemTrait; +use crate::model::core::backup::execution::*; +use crate::model::core::backup::communication::*; use crate::model::error::system::SystemError; use crate::model::error::task::TaskError; use crate::model::error::Error; +use async_trait::async_trait; use crossbeam_queue::SegQueue; use dashmap::DashMap; use futures::future::join_all; @@ -20,13 +22,14 @@ use tokio::sync::oneshot; use tokio::task::JoinHandle; use tracing::error; use uuid::Uuid; +use crate::model::core::gui::communication::ExecutionErrors; pub struct BackupEngine { app_config: Arc, io_manager: Arc, - actor_system: Arc, + communication_manager: Arc, progress_tracker: Arc, - executions: Arc>, + executions: Arc>, running_executions: Arc, JoinHandle<()>)>>, } @@ -34,19 +37,29 @@ impl BackupEngine { pub fn new( app_config: Arc, io_manager: Arc, - actor_system: Arc, + communication_manager: Arc, progress_tracker: Arc, ) -> Self { Self { app_config, io_manager, - actor_system, + communication_manager, progress_tracker, executions: Arc::new(DashMap::new()), running_executions: Arc::new(DashMap::new()), } } + pub async fn register_services(self: Arc) { + let communication_manager = self.communication_manager.clone(); + communication_manager + .with_service(self) + .command::() + .query::() + .event::() + .build(); + } + pub async fn stop_all_executions(&self) { let keys: Vec = self .running_executions @@ -66,14 +79,14 @@ impl BackupEngine { } } - pub fn get_all_executions(&self) -> Vec<(Uuid, BackupExecution)> { + pub fn get_all_executions(&self) -> Vec<(Uuid, Execution)> { self.executions .iter() .map(|entry| (*entry.key(), entry.value().clone())) .collect() } - pub async fn add_execution(&self, execution: BackupExecution) { + pub async fn add_execution(&self, execution: Execution) { self.executions.insert(execution.uuid, execution); } @@ -81,14 +94,14 @@ impl BackupEngine { self.executions.remove(uuid); } - pub async fn start_execution(&self, uuid: Uuid) -> Result<(), Error> { - if self.running_executions.contains_key(&uuid) { + pub async fn start_execution(&self, uuid: &Uuid) -> Result<(), Error> { + if self.running_executions.contains_key(uuid) { Err(TaskError::IllegalRunState)? } let mut ref_mut = self .executions - .get_mut(&uuid) + .get_mut(uuid) .ok_or(TaskError::ExecutionNotFound)?; let execution = ref_mut.value_mut(); if execution.state != BackupState::Pending { @@ -100,14 +113,14 @@ impl BackupEngine { let execution = execution.clone(); let (tx, rx) = oneshot::channel(); let handle = tokio::spawn(async move { execution_runner.run(execution, rx, false).await }); - self.running_executions.insert(uuid, (tx, handle)); + self.running_executions.insert(*uuid, (tx, handle)); Ok(()) } - pub async fn suspend_execution(&self, uuid: Uuid) -> Result<(), Error> { + pub async fn suspend_execution(&self, uuid: &Uuid) -> Result<(), Error> { let mut ref_mut = self .executions - .get_mut(&uuid) + .get_mut(uuid) .ok_or(TaskError::ExecutionNotFound)?; let execution = ref_mut.value_mut(); if execution.state != BackupState::Running { @@ -118,7 +131,7 @@ impl BackupEngine { let (_, (shutdown, handle)) = self .running_executions - .remove(&uuid) + .remove(uuid) .ok_or(TaskError::ExecutionNotFound)?; shutdown .send(()) @@ -127,14 +140,14 @@ impl BackupEngine { Ok(()) } - pub async fn resume_execution(&self, uuid: Uuid) -> Result<(), Error> { + pub async fn resume_execution(&self, uuid: &Uuid) -> Result<(), Error> { if self.running_executions.contains_key(&uuid) { Err(TaskError::IllegalRunState)? } let mut ref_mut = self .executions - .get_mut(&uuid) + .get_mut(uuid) .ok_or(TaskError::ExecutionNotFound)?; let execution = ref_mut.value_mut(); if execution.state != BackupState::Suspended { @@ -146,21 +159,21 @@ impl BackupEngine { let execution = execution.clone(); let (tx, rx) = oneshot::channel(); let handle = tokio::spawn(async move { execution_runner.run(execution, rx, true).await }); - self.running_executions.insert(uuid, (tx, handle)); + self.running_executions.insert(*uuid, (tx, handle)); Ok(()) } fn to_execution_runner(&self) -> ExecutionRunner { let config = self.app_config.clone(); let io_manager = self.io_manager.clone(); - let actor_system = self.actor_system.clone(); + let communication_manager = self.communication_manager.clone(); let progress_tracker = self.progress_tracker.clone(); let executions = self.executions.clone(); let running_executions = self.running_executions.clone(); ExecutionRunner::new( config, io_manager, - actor_system, + communication_manager, progress_tracker, executions, running_executions, @@ -171,9 +184,9 @@ impl BackupEngine { struct ExecutionRunner { app_config: Arc, io_manager: Arc, - actor_system: Arc, + communication_manager: Arc, progress_tracker: Arc, - executions: Arc>, + executions: Arc>, running_executions: Arc, JoinHandle<()>)>>, } @@ -181,27 +194,22 @@ impl ExecutionRunner { pub fn new( app_config: Arc, io_manager: Arc, - actor_system: Arc, + communication_manager: Arc, progress_tracker: Arc, - executions: Arc>, + executions: Arc>, running_executions: Arc, JoinHandle<()>)>>, ) -> Self { Self { app_config, io_manager, - actor_system, + communication_manager, progress_tracker, executions, running_executions, } } - async fn run( - &self, - execution: BackupExecution, - mut shutdown: oneshot::Receiver<()>, - resume: bool, - ) { + async fn run(&self, execution: Execution, mut shutdown: oneshot::Receiver<()>, resume: bool) { let config = &self.app_config; let progress_tracker = &self.progress_tracker; @@ -253,17 +261,16 @@ impl ExecutionRunner { next_level.extend(worker_next_level); if !worker_errors.is_empty() { errors.extend(worker_errors.clone()); - if let Some(gui_ref) = self.actor_system.actor_of::() + let event = ExecutionErrors { + uuid: execution.uuid, + errors: worker_errors, + }; + if let Err(err) = self + .communication_manager + .publish_event::(event) + .await { - if let Err(err) = gui_ref - .tell(GuiMessage::ExecutionErrors { - uuid: execution.uuid, - errors: worker_errors, - }) - .await - { - error!("{}", err); - } + error!("{}", err); } } } @@ -317,7 +324,7 @@ impl Worker { async fn run( &self, - execution: BackupExecution, + execution: Execution, global_queue: Arc>, mut shutdown: oneshot::Receiver<()>, ) -> (Vec, Vec) { @@ -382,7 +389,7 @@ impl Worker { async fn process_entry( &self, - execution: &BackupExecution, + execution: &Execution, current_path: &Path, ) -> Result, Error> { let io_manager = &self.io_manager; @@ -391,7 +398,8 @@ impl Worker { let destination_root = &execution.destination_path; let source_path = current_path; - let destination_path = self.calculate_destination_path(source_path, source_root, destination_root)?; + let destination_path = + self.calculate_destination_path(source_path, source_root, destination_root)?; let destination_path = destination_path.as_path(); let is_symlink = io_manager.is_symlink(source_path).await.unwrap_or(false); @@ -413,7 +421,7 @@ impl Worker { async fn backup_directory( &self, - execution: &BackupExecution, + execution: &Execution, source_path: &Path, destination_path: &Path, ) -> Result, Error> { @@ -438,7 +446,7 @@ impl Worker { async fn backup_file( &self, - execution: &BackupExecution, + execution: &Execution, source_path: &Path, destination_path: &Path, ) -> Result, Error> { @@ -469,7 +477,7 @@ impl Worker { #[inline(always)] async fn process_symlink( &self, - execution: &BackupExecution, + execution: &Execution, source_path: &Path, destination_path: &Path, ) -> Result<(), Error> { @@ -484,7 +492,7 @@ impl Worker { async fn follow_symlink( &self, - execution: &BackupExecution, + execution: &Execution, source_path: &Path, destination_path: &Path, ) -> Result<(), Error> { @@ -541,7 +549,7 @@ impl Worker { async fn copy_symlink( &self, - execution: &BackupExecution, + execution: &Execution, source_path: &Path, destination_path: &Path, ) -> Result<(), Error> { @@ -646,3 +654,39 @@ impl Worker { Ok(destination_root.join(relative_path)) } } + +#[async_trait] +impl CommandHandler for BackupEngine { + async fn handle_command(&self, command: BackupCommand) -> Result<(), Error> { + match command { + BackupCommand::AddExecution(execution) => { + self.add_execution(execution).await; + } + BackupCommand::RemoveExecution(uuid) => { + self.remove_execution(&uuid).await; + } + BackupCommand::StartExecution(uuid) => { + self.start_execution(&uuid).await?; + } + BackupCommand::SuspendExecution(uuid) => { + self.suspend_execution(&uuid).await?; + } + BackupCommand::ResumeExecution(uuid) => { + self.resume_execution(&uuid).await?; + } + } + Ok(()) + } +} + +#[async_trait] +impl QueryHandler for BackupEngine { + async fn handle_query(&self, query: BackupQuery) -> Result { + match query { + BackupQuery::GetExecutions => { + let executions = self.get_all_executions(); + Ok(BackupQueryResponse::GetExecutions(executions)) + } + } + } +} diff --git a/src/core/backup/backup_service.rs b/src/core/backup/backup_service.rs index 399dd27..1ed12b3 100644 --- a/src/core/backup/backup_service.rs +++ b/src/core/backup/backup_service.rs @@ -1,13 +1,8 @@ use crate::core::backup::backup_engine::BackupEngine; use crate::core::backup::progress_tracker::ProgressTracker; -use crate::core::infrastructure::actor_system::ActorSystem; use crate::core::infrastructure::app_config::AppConfig; +use crate::core::infrastructure::communication_manager::CommunicationManager; use crate::core::infrastructure::io_manager::IOManager; -use crate::interface::actor::actor::Actor; -use crate::interface::actor::message::Message; -use crate::model::core::backup::message::*; -use crate::model::error::Error; -use async_trait::async_trait; use std::sync::Arc; pub struct BackupService { @@ -15,68 +10,27 @@ pub struct BackupService { } impl BackupService { - pub async fn init( + pub async fn new( app_config: Arc, io_manager: Arc, - actor_system: Arc, - ) { + communication_manager: Arc, + ) -> Self { let progress_tracker = Arc::new(ProgressTracker::new(io_manager.clone())); let backup_engine = Arc::new(BackupEngine::new( app_config, io_manager, - actor_system.clone(), + communication_manager, progress_tracker, )); - let backup_service = Self { - backup_engine, - }; - actor_system.spawn(backup_service).await; + Self { backup_engine } } -} - -#[async_trait] -impl Actor for BackupService { - type Message = BackupServiceMessage; - - async fn pre_start(&mut self) {} - async fn post_stop(&mut self) { - self.backup_engine.stop_all_executions().await; + pub async fn register_services(&self) { + let backup_engine = self.backup_engine.clone(); + backup_engine.register_services().await; } - async fn receive( - &mut self, - message: Self::Message, - ) -> Result<::Response, Error> { - match message { - BackupServiceMessage::ServiceCall(service_call) => match service_call { - ServiceCallMessage::AddExecution(execution) => { - self.backup_engine.add_execution(execution).await; - Ok(BackupServiceResponse::None) - } - ServiceCallMessage::RemoveExecution(uuid) => { - self.backup_engine.resume_execution(uuid).await?; - Ok(BackupServiceResponse::None) - } - ServiceCallMessage::StartExecution(uuid) => { - self.backup_engine.start_execution(uuid).await?; - Ok(BackupServiceResponse::None) - } - ServiceCallMessage::SuspendExecution(uuid) => { - self.backup_engine.suspend_execution(uuid).await?; - Ok(BackupServiceResponse::None) - } - ServiceCallMessage::ResumeExecution(uuid) => { - self.backup_engine.resume_execution(uuid).await?; - Ok(BackupServiceResponse::None) - } - ServiceCallMessage::GetExecutions => { - let execution = self.backup_engine.get_all_executions(); - Ok(BackupServiceResponse::ServiceCall( - ServiceCallResponse::GetExecutions(execution), - )) - } - }, - } + pub async fn shutdown(&self) { + self.backup_engine.stop_all_executions().await; } } diff --git a/src/core/backup/progress_tracker.rs b/src/core/backup/progress_tracker.rs index b01080d..b224c0f 100644 --- a/src/core/backup/progress_tracker.rs +++ b/src/core/backup/progress_tracker.rs @@ -1,5 +1,5 @@ use crate::core::infrastructure::io_manager::IOManager; -use crate::interface::file_system::FileSystemTrait; +use crate::interface::core::file_system::FileSystemTrait; use crate::model::core::backup::progress_data::ProgressData; use crate::model::error::Error; use crate::model::error::io::IOError; diff --git a/src/core/gui/gui_manager.rs b/src/core/gui/gui_manager.rs index 5e7ac53..2b931fe 100644 --- a/src/core/gui/gui_manager.rs +++ b/src/core/gui/gui_manager.rs @@ -1,6 +1,5 @@ -use crate::core::gui::gui_message_handler::GuiMessageHandler; -use crate::core::infrastructure::actor_system::ActorSystem; use crate::core::infrastructure::app_config::AppConfig; +use crate::core::infrastructure::communication_manager::CommunicationManager; use crate::model::error::misc::MiscError; use crate::model::error::Error; use crate::ui::execution_page::ExecutionPage; @@ -13,28 +12,26 @@ use std::sync::Arc; pub struct GuiManager { app_config: Arc, - actor_system: Arc, + communication_manager: Arc, } impl GuiManager { - pub fn new(app_config: Arc, actor_system: Arc) -> Self { + pub fn new( + app_config: Arc, + communication_manager: Arc, + ) -> Self { Self { app_config, - actor_system, + communication_manager, } } pub async fn start(&self) -> Result<(), Error> { let app_config = self.app_config.clone(); - let actor_system = self.actor_system.clone(); + let communication_manager = self.communication_manager.clone(); - let mut handler = GuiMessageHandler::new(); - let message_rx = handler.subscribe(); - actor_system.spawn(handler).await; - - let execution_page = - ExecutionPage::new(app_config.clone(), actor_system.clone(), message_rx); - let schedule_page = SchedulePage::new(app_config, actor_system)?; + let execution_page = ExecutionPage::new(app_config.clone(), communication_manager.clone())?; + let schedule_page = SchedulePage::new(app_config, communication_manager)?; let main_page = MainPage::new(execution_page, schedule_page); let icon_data = Assets::load_app_icon()?; diff --git a/src/core/gui/gui_message_handler.rs b/src/core/gui/gui_message_handler.rs deleted file mode 100644 index 3cd5bb8..0000000 --- a/src/core/gui/gui_message_handler.rs +++ /dev/null @@ -1,48 +0,0 @@ -use crate::interface::actor::actor::Actor; -use crate::interface::actor::message::Message; -use crate::model::core::gui::message::GuiMessage; -use crate::model::error::Error; -use async_trait::async_trait; -use std::sync::mpsc; - -pub struct GuiMessageHandler { - subscriber: Vec>, -} - -impl GuiMessageHandler { - pub fn new() -> Self { - Self { - subscriber: Vec::new(), - } - } - - pub fn subscribe(&mut self) -> mpsc::Receiver { - let (tx, rx) = mpsc::channel(); - self.subscriber.push(tx); - rx - } - - fn broadcast(&mut self, message: GuiMessage) { - self.subscriber - .retain(|tx| tx.send(message.clone()).is_ok()); - } -} - -#[async_trait] -impl Actor for GuiMessageHandler { - type Message = GuiMessage; - - async fn pre_start(&mut self) {} - - async fn post_stop(&mut self) { - self.subscriber.clear(); - } - - async fn receive( - &mut self, - message: Self::Message, - ) -> Result<::Response, Error> { - self.broadcast(message); - Ok(()) - } -} diff --git a/src/core/gui/mod.rs b/src/core/gui/mod.rs index 464dd49..166deaf 100644 --- a/src/core/gui/mod.rs +++ b/src/core/gui/mod.rs @@ -1,2 +1 @@ -pub mod gui_message_handler; pub mod gui_manager; diff --git a/src/core/infrastructure/actor_system.rs b/src/core/infrastructure/actor_system.rs deleted file mode 100644 index bfaf795..0000000 --- a/src/core/infrastructure/actor_system.rs +++ /dev/null @@ -1,56 +0,0 @@ -use crate::interface::actor::actor::Actor; -use crate::model::core::actor::actor_ref::ActorRef; -use crate::model::core::actor::actor_runtime::ActorRuntime; -use crate::model::error::system::SystemError; -use dashmap::DashMap; -use macros::log; -use std::any::{Any, TypeId}; -use tokio::sync::oneshot; - -pub struct ActorSystem { - actors: DashMap>, - shutdowns: DashMap>, -} - -impl ActorSystem { - pub fn new() -> Self { - Self { - actors: DashMap::new(), - shutdowns: DashMap::new(), - } - } - - pub async fn spawn(&self, actor: A) - where - A: Actor + 'static, - { - let actor_id = TypeId::of::(); - let (actor_runtime, actor_ref) = ActorRuntime::new(actor); - let shutdown = actor_runtime.run().await; - self.actors.insert(actor_id, Box::new(actor_ref)); - self.shutdowns.insert(actor_id, shutdown); - } - - pub fn shutdown(&self) { - let keys = self - .shutdowns - .iter() - .map(|x| x.key().clone()) - .collect::>(); - for key in keys { - if let Some((_, shutdown)) = self.shutdowns.remove(&key) { - if let Err(_) = shutdown.send(()) { - log!(SystemError::ShutdownSignalFailed); - } - } - } - } - - pub fn actor_of(&self) -> Option> { - let type_id = TypeId::of::(); - self.actors - .get(&type_id)? - .downcast_ref::>() - .cloned() - } -} diff --git a/src/core/infrastructure/communication_manager.rs b/src/core/infrastructure/communication_manager.rs new file mode 100644 index 0000000..8842751 --- /dev/null +++ b/src/core/infrastructure/communication_manager.rs @@ -0,0 +1,162 @@ +use crate::core::infrastructure::app_config::AppConfig; +use crate::interface::communication::command::*; +use crate::interface::communication::event::Event; +use crate::interface::communication::event::EventBroadcaster; +use crate::interface::communication::query::*; +use crate::model::core::infrastructure::event_broadcaster::TypedEventBroadcaster; +use crate::model::error::misc::MiscError; +use crate::model::error::Error; +use dashmap::DashMap; +use std::any::{Any, TypeId}; +use std::sync::Arc; +use tokio::sync::broadcast; + +pub struct CommunicationManager { + app_config: Arc, + command_handlers: DashMap, + query_handlers: DashMap, + event_broadcasters: DashMap>, +} + +impl CommunicationManager { + pub fn new(app_config: Arc) -> Self { + Self { + app_config, + command_handlers: DashMap::new(), + query_handlers: DashMap::new(), + event_broadcasters: DashMap::new(), + } + } + + pub fn with_service( + self: Arc, + service: Arc, + ) -> ServiceRegistrar { + ServiceRegistrar::new(service, self) + } + + pub fn register_command_handler( + &self, + handler: Arc + Send + Sync>, + ) { + let type_id = TypeId::of::(); + let boxed_handler: CommandHandlerFn = Box::new(move |command: Box| { + let handler = handler.clone(); + Box::pin(async move { + let command = *command + .downcast::() + .map_err(|_| MiscError::TypeMismatch)?; + handler.handle_command(command).await + }) as CommandFuture + }); + + self.command_handlers.insert(type_id, boxed_handler); + } + + pub async fn send_command(&self, command: C) -> Result<(), Error> { + let type_id = TypeId::of::(); + if let Some(handler) = self.command_handlers.get(&type_id) { + handler(Box::new(command)).await + } else { + Err(MiscError::HandlerNotFound)? + } + } + + pub fn register_query_handler( + &self, + handler: Arc + Send + Sync>, + ) { + let type_id = TypeId::of::(); + let boxed_handler: QueryHandlerFn = Box::new(move |query: Box| { + let handler = handler.clone(); + Box::pin(async move { + let query = *query.downcast::().map_err(|_| MiscError::TypeMismatch)?; + let response = handler.handle_query(query).await?; + Ok(Box::new(response) as Box) + }) as QueryFuture + }); + + self.query_handlers.insert(type_id, boxed_handler); + } + + pub async fn send_query(&self, query: Q) -> Result { + let type_id = TypeId::of::(); + if let Some(handler) = self.query_handlers.get(&type_id) { + let response = handler(Box::new(query)).await?; + Ok(*response + .downcast::() + .map_err(|_| MiscError::TypeMismatch)?) + } else { + Err(MiscError::HandlerNotFound)? + } + } + + pub fn register_event_type(&self) { + let channel_capacity = self.app_config.channel_capacity; + let type_id = TypeId::of::(); + let (tx, _) = broadcast::channel::(channel_capacity); + let broadcaster = TypedEventBroadcaster { sender: tx }; + self.event_broadcasters + .insert(type_id, Box::new(broadcaster)); + } + + pub fn subscribe_event(&self) -> Result, Error> { + let type_id = TypeId::of::(); + let broadcaster = self + .event_broadcasters + .get(&type_id) + .ok_or(MiscError::TypeNotRegistered)?; + let receiver_box = broadcaster.subscribe_typed(); + let receiver = *receiver_box + .downcast::>() + .map_err(|_| MiscError::TypeMismatch)?; + Ok(receiver) + } + + pub async fn publish_event(&self, event: E) -> Result<(), Error> { + let type_id = TypeId::of::(); + let broadcaster = self + .event_broadcasters + .get(&type_id) + .ok_or(MiscError::TypeNotRegistered)?; + broadcaster.broadcast_event(Box::new(event)) + } +} + +pub struct ServiceRegistrar { + service: Arc, + comm: Arc, +} + +impl ServiceRegistrar { + fn new(service: Arc, comm: Arc) -> Self { + Self { service, comm } + } + + pub fn command(self) -> Self + where + S: CommandHandler, + { + let handler: Arc + Send + Sync> = self.service.clone(); + self.comm.register_command_handler::(handler); + self + } + + pub fn query(self) -> Self + where + S: QueryHandler, + { + let handler: Arc + Send + Sync> = self.service.clone(); + self.comm.register_query_handler::(handler); + self + } + + pub fn event(self) -> Self { + self.comm.register_event_type::(); + self + } + + pub fn build(self) -> Arc { + self.comm + } +} diff --git a/src/core/infrastructure/io_manager.rs b/src/core/infrastructure/io_manager.rs index bf7ea72..c929ce6 100644 --- a/src/core/infrastructure/io_manager.rs +++ b/src/core/infrastructure/io_manager.rs @@ -1,5 +1,5 @@ use crate::core::infrastructure::app_config::AppConfig; -use crate::interface::file_system::FileSystemTrait; +use crate::interface::core::file_system::FileSystemTrait; use crate::platform::file_system::FileSystem; use std::ops::Deref; use std::sync::Arc; diff --git a/src/core/infrastructure/mod.rs b/src/core/infrastructure/mod.rs index 1ecb5f1..ec8dd39 100644 --- a/src/core/infrastructure/mod.rs +++ b/src/core/infrastructure/mod.rs @@ -1,4 +1,4 @@ pub mod app_config; +pub mod communication_manager; pub mod database_manager; pub mod io_manager; -pub mod actor_system; \ No newline at end of file diff --git a/src/core/schedule/schedule_manager.rs b/src/core/schedule/schedule_manager.rs index bfef95e..7c9200f 100644 --- a/src/core/schedule/schedule_manager.rs +++ b/src/core/schedule/schedule_manager.rs @@ -1,10 +1,13 @@ -use crate::core::backup::backup_service::BackupService; -use crate::core::infrastructure::actor_system::ActorSystem; +use crate::core::infrastructure::communication_manager::CommunicationManager; use crate::core::infrastructure::database_manager::DatabaseManager; +use crate::interface::communication::command::CommandHandler; +use crate::interface::communication::query::QueryHandler; use crate::interface::repository::schedule::ScheduleRepository; -use crate::model::core::backup::message::{BackupServiceMessage, ServiceCallMessage}; -use crate::model::core::schedule::backup_schedule::*; +use crate::model::core::backup::communication::BackupCommand; +use crate::model::core::schedule::communication::*; +use crate::model::core::schedule::schedule::*; use crate::model::error::Error; +use async_trait::async_trait; use chrono::{Duration, Months, Utc}; use dashmap::DashMap; use std::sync::Arc; @@ -12,14 +15,14 @@ use uuid::Uuid; pub struct ScheduleManager { database_manager: Arc, - actor_system: Arc, - schedules: DashMap, + communication_manager: Arc, + schedules: DashMap, } impl ScheduleManager { pub async fn new( database_manager: Arc, - actor_system: Arc, + communication_manager: Arc, ) -> Result { let schedules = DashMap::new(); let database_schedules = database_manager.get_all_backup_schedules().await?; @@ -28,35 +31,53 @@ impl ScheduleManager { } let schedule_manager = ScheduleManager { database_manager, - actor_system, + communication_manager, schedules, }; Ok(schedule_manager) } - pub async fn get_all_schedules(&self) -> Vec { + pub async fn register_services(self: Arc) { + let communication_manager = self.communication_manager.clone(); + communication_manager + .with_service(self) + .command::() + .query::() + .build(); + } + + pub async fn get_all_schedules(&self) -> Vec { self.schedules.iter().map(|x| x.value().clone()).collect() } - pub async fn create_schedule(&self, schedule: BackupSchedule) -> Result<(), Error> { + pub async fn create_schedule(&self, schedule: Schedule) -> Result<(), Error> { self.database_manager .create_backup_schedule(&schedule) .await?; self.schedules.insert(schedule.uuid, schedule); + self.communication_manager + .send_command(ScheduleTimerCommand::RefreshTimer) + .await?; Ok(()) } - pub async fn modify_schedule(&self, schedule: BackupSchedule) -> Result<(), Error> { + pub async fn modify_schedule(&self, schedule: Schedule) -> Result<(), Error> { self.database_manager .modify_backup_schedule(&schedule) .await?; self.schedules.insert(schedule.uuid, schedule); + self.communication_manager + .send_command(ScheduleTimerCommand::RefreshTimer) + .await?; Ok(()) } pub async fn remove_schedule(&self, uuid: Uuid) -> Result<(), Error> { self.database_manager.remove_backup_schedule(uuid).await?; self.schedules.remove(&uuid); + self.communication_manager + .send_command(ScheduleTimerCommand::RefreshTimer) + .await?; Ok(()) } @@ -67,6 +88,9 @@ impl ScheduleManager { .modify_backup_schedule(&schedule) .await?; self.schedules.insert(schedule.uuid, schedule); + self.communication_manager + .send_command(ScheduleTimerCommand::RefreshTimer) + .await?; } Ok(()) } @@ -78,6 +102,9 @@ impl ScheduleManager { .modify_backup_schedule(&schedule) .await?; self.schedules.insert(schedule.uuid, schedule); + self.communication_manager + .send_command(ScheduleTimerCommand::RefreshTimer) + .await?; } Ok(()) } @@ -89,6 +116,9 @@ impl ScheduleManager { .modify_backup_schedule(&schedule) .await?; self.schedules.insert(schedule.uuid, schedule); + self.communication_manager + .send_command(ScheduleTimerCommand::RefreshTimer) + .await?; } Ok(()) } @@ -108,13 +138,8 @@ impl ScheduleManager { continue; } let execution = schedule.to_execution(); - if let Some(service_ref) = self.actor_system.actor_of::() { - service_ref - .tell(BackupServiceMessage::ServiceCall( - ServiceCallMessage::AddExecution(execution), - )) - .await?; - } + let command = BackupCommand::AddExecution(execution); + self.communication_manager.send_command(command).await?; self.update_next_run_time(schedule); database_manager.modify_backup_schedule(schedule).await?; } @@ -123,7 +148,7 @@ impl ScheduleManager { Ok(()) } - fn update_next_run_time(&self, schedule: &mut BackupSchedule) { + fn update_next_run_time(&self, schedule: &mut Schedule) { if schedule.next_run_time.is_none() { return; } @@ -143,3 +168,48 @@ impl ScheduleManager { schedule.next_run_time = new_next_run_time; } } + +#[async_trait] +impl CommandHandler for ScheduleManager { + async fn handle_command(&self, command: ScheduleManagerCommand) -> Result<(), Error> { + match command { + ScheduleManagerCommand::AddSchedule(schedule) => { + self.create_schedule(schedule).await?; + } + ScheduleManagerCommand::ModifySchedule(schedule) => { + self.modify_schedule(schedule).await?; + } + ScheduleManagerCommand::RemoveSchedule(uuid) => { + self.remove_schedule(uuid).await?; + } + ScheduleManagerCommand::ActivateSchedule(uuid) => { + self.active_schedule(uuid).await?; + } + ScheduleManagerCommand::PauseSchedule(uuid) => { + self.pause_schedule(uuid).await?; + } + ScheduleManagerCommand::DisableSchedule(uuid) => { + self.disable_schedule(uuid).await?; + } + ScheduleManagerCommand::ExecuteReadySchedules => { + self.execute_ready_schedule().await?; + } + } + Ok(()) + } +} + +#[async_trait] +impl QueryHandler for ScheduleManager { + async fn handle_query( + &self, + query: ScheduleManagerQuery, + ) -> Result { + match query { + ScheduleManagerQuery::GetSchedules => { + let executions = self.get_all_schedules().await; + Ok(ScheduleManagerQueryResponse::GetSchedules(executions)) + } + } + } +} diff --git a/src/core/schedule/schedule_service.rs b/src/core/schedule/schedule_service.rs index 8996786..7dd28d1 100644 --- a/src/core/schedule/schedule_service.rs +++ b/src/core/schedule/schedule_service.rs @@ -1,123 +1,49 @@ -use crate::core::infrastructure::actor_system::ActorSystem; use crate::core::infrastructure::app_config::AppConfig; +use crate::core::infrastructure::communication_manager::CommunicationManager; use crate::core::infrastructure::database_manager::DatabaseManager; use crate::core::schedule::schedule_manager::ScheduleManager; use crate::core::schedule::schedule_timer::ScheduleTimer; -use crate::interface::actor::actor::Actor; -use crate::interface::actor::message::Message; -use crate::model::core::schedule::message::*; -use crate::model::error::system::SystemError; +use crate::interface::core::runnable::Runnable; use crate::model::error::Error; use async_trait::async_trait; -use macros::log; -use std::sync::{Arc, OnceLock}; -use tokio::sync::{mpsc, oneshot}; +use std::sync::Arc; +use tokio::sync::oneshot::Receiver; pub struct ScheduleService { schedule_manager: Arc, schedule_timer: Arc, - timer_refresh: OnceLock>, - shutdowns: Vec>, } impl ScheduleService { - pub async fn init( + pub async fn new( app_config: Arc, database_manager: Arc, - actor_system: Arc, - ) -> Result<(), Error> { + communication_manager: Arc, + ) -> Result { let schedule_manager = - Arc::new(ScheduleManager::new(database_manager, actor_system.clone()).await?); - let schedule_timer = Arc::new(ScheduleTimer::new(app_config, actor_system.clone())); + Arc::new(ScheduleManager::new(database_manager, communication_manager.clone()).await?); + let schedule_timer = Arc::new(ScheduleTimer::new(app_config, communication_manager)); let schedule_service = Self { schedule_manager, schedule_timer, - timer_refresh: OnceLock::new(), - shutdowns: Vec::new(), }; - actor_system.spawn(schedule_service).await; - Ok(()) + Ok(schedule_service) } - pub fn refresh_timer(&self) { - if let Some(timer_refresh) = self.timer_refresh.get() { - let _ = timer_refresh.send(()); - } + pub async fn register_services(&self) { + let schedule_manager = self.schedule_manager.clone(); + let schedule_timer = self.schedule_timer.clone(); + schedule_manager.register_services().await; + schedule_timer.register_services().await; } } #[async_trait] -impl Actor for ScheduleService { - type Message = ScheduleServiceMessage; - - async fn pre_start(&mut self) { +impl Runnable for ScheduleService { + async fn run_impl(self: Arc, shutdown_rx: Receiver<()>) { let schedule_timer = self.schedule_timer.clone(); - if let Ok((timer_refresh, shutdown)) = schedule_timer.run().await { - self.timer_refresh.get_or_init(|| timer_refresh); - self.shutdowns.push(shutdown); - } - } - - async fn post_stop(&mut self) { - let shutdowns = std::mem::take(&mut self.shutdowns); - for shutdown in shutdowns { - if let Err(_) = shutdown.send(()) { - log!(SystemError::ShutdownSignalFailed) - } - } - } - - async fn receive( - &mut self, - message: Self::Message, - ) -> Result<::Response, Error> { - match message { - ScheduleServiceMessage::UnitNotification(unit_notification) => { - match unit_notification { - UnitNotificationMessage::CheckSchedule => { - self.schedule_manager.execute_ready_schedule().await?; - Ok(ScheduleServiceResponse::None) - } - } - } - ScheduleServiceMessage::ServiceCall(service_call) => match service_call { - ServiceCallMessage::AddSchedule(schedule) => { - self.schedule_manager.create_schedule(schedule).await?; - self.refresh_timer(); - Ok(ScheduleServiceResponse::None) - } - ServiceCallMessage::ModifySchedule(schedule) => { - self.schedule_manager.modify_schedule(schedule).await?; - self.refresh_timer(); - Ok(ScheduleServiceResponse::None) - } - ServiceCallMessage::RemoveSchedule(uuid) => { - self.schedule_manager.remove_schedule(uuid).await?; - self.refresh_timer(); - Ok(ScheduleServiceResponse::None) - } - ServiceCallMessage::ActivateSchedule(uuid) => { - self.schedule_manager.active_schedule(uuid).await?; - self.refresh_timer(); - Ok(ScheduleServiceResponse::None) - } - ServiceCallMessage::PauseSchedule(uuid) => { - self.schedule_manager.pause_schedule(uuid).await?; - self.refresh_timer(); - Ok(ScheduleServiceResponse::None) - } - ServiceCallMessage::DisableSchedule(uuid) => { - self.schedule_manager.disable_schedule(uuid).await?; - self.refresh_timer(); - Ok(ScheduleServiceResponse::None) - } - ServiceCallMessage::GetSchedules => { - let schedules = self.schedule_manager.get_all_schedules().await; - Ok(ScheduleServiceResponse::ServiceCall( - ServiceCallResponse::GetSchedules(schedules), - )) - } - }, - } + let timer_shutdown = schedule_timer.run().await; + let _ = shutdown_rx.await; + let _ = timer_shutdown.send(()); } } diff --git a/src/core/schedule/schedule_timer.rs b/src/core/schedule/schedule_timer.rs index 2bfa16e..1844799 100644 --- a/src/core/schedule/schedule_timer.rs +++ b/src/core/schedule/schedule_timer.rs @@ -1,86 +1,52 @@ -use crate::core::infrastructure::actor_system::ActorSystem; use crate::core::infrastructure::app_config::AppConfig; -use crate::core::schedule::schedule_service::ScheduleService; -use crate::model::core::schedule::backup_schedule::ScheduleState; -use crate::model::core::schedule::message::*; -use crate::model::error::actor::ActorError; +use crate::core::infrastructure::communication_manager::CommunicationManager; +use crate::interface::communication::command::CommandHandler; +use crate::interface::core::runnable::Runnable; +use crate::model::core::schedule::communication::*; +use crate::model::core::schedule::schedule::ScheduleState; use crate::model::error::Error; +use async_trait::async_trait; use chrono::{Duration, Utc}; use std::sync::Arc; use tokio::select; -use tokio::sync::{mpsc, oneshot}; +use tokio::sync::Notify; +use tokio::sync::oneshot; use tokio::time::sleep; use tracing::error; pub struct ScheduleTimer { app_config: Arc, - actor_system: Arc, + communication_manager: Arc, + refresh_notify: Arc, } impl ScheduleTimer { - pub fn new(app_config: Arc, actor_system: Arc) -> Self { + pub fn new( + app_config: Arc, + communication_manager: Arc, + ) -> Self { ScheduleTimer { app_config, - actor_system, + communication_manager, + refresh_notify: Arc::new(Notify::new()), } } - pub async fn run( - self: Arc, - ) -> Result<(mpsc::UnboundedSender<()>, oneshot::Sender<()>), Error> { - let (refresh_tx, mut refresh_rx) = mpsc::unbounded_channel(); - let (shutdown_tx, mut shutdown_rx) = oneshot::channel(); - let service_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - tokio::spawn(async move { - loop { - let mut sleep_time = match self.calculate_sleep_duration().await { - Ok(Some(duration)) => duration, - Ok(None) => Duration::seconds(self.app_config.default_wakeup_time), - Err(err) => { - error!("{}", err); - Duration::seconds(self.app_config.default_wakeup_time) - } - }; - if sleep_time < Duration::seconds(0) { - sleep_time = Duration::seconds(0); - } - select! { - biased; - _ = &mut shutdown_rx => { break; } - _ = refresh_rx.recv() => { continue; } - _ = sleep(sleep_time.to_std().unwrap()) => {} - } - if let Err(err) = service_ref - .tell(ScheduleServiceMessage::UnitNotification( - UnitNotificationMessage::CheckSchedule, - )) - .await - { - error!("{}", err); - } - } - }); - Ok((refresh_tx, shutdown_tx)) + pub async fn register_services(self: Arc) { + let communication_manager = self.communication_manager.clone(); + communication_manager + .with_service(self) + .command::() + .build(); } async fn calculate_sleep_duration(&self) -> Result, Error> { let mut next_time = None; - let service_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let response = service_ref - .ask(ScheduleServiceMessage::ServiceCall( - ServiceCallMessage::GetSchedules, - )) + let communication_manager = self.communication_manager.clone(); + let response = communication_manager + .send_query(ScheduleManagerQuery::GetSchedules) .await?; - let ScheduleServiceResponse::ServiceCall(service_call) = response else { - return Ok(None); - }; - let ServiceCallResponse::GetSchedules(schedules) = service_call; + let ScheduleManagerQueryResponse::GetSchedules(schedules) = response; for schedule in schedules { if schedule.state != ScheduleState::Active { continue; @@ -105,3 +71,48 @@ impl ScheduleTimer { } } } + +#[async_trait] +impl Runnable for ScheduleTimer { + async fn run_impl(self: Arc, mut shutdown_rx: oneshot::Receiver<()>) { + let communication_manager = self.communication_manager.clone(); + + loop { + let mut sleep_time = match self.calculate_sleep_duration().await { + Ok(Some(duration)) => duration, + Ok(None) => Duration::seconds(self.app_config.default_wakeup_time), + Err(err) => { + error!("{}", err); + Duration::seconds(self.app_config.default_wakeup_time) + } + }; + if sleep_time < Duration::seconds(0) { + sleep_time = Duration::seconds(0); + } + select! { + biased; + _ = &mut shutdown_rx => { break; } + _ = self.refresh_notify.notified() => { continue; } + _ = sleep(sleep_time.to_std().unwrap()) => {} + } + if let Err(err) = communication_manager + .send_command(ScheduleManagerCommand::ExecuteReadySchedules) + .await + { + error!("{}", err); + } + } + } +} + +#[async_trait] +impl CommandHandler for ScheduleTimer { + async fn handle_command(&self, command: ScheduleTimerCommand) -> Result<(), Error> { + match command { + ScheduleTimerCommand::RefreshTimer => { + self.refresh_notify.notify_one(); + Ok(()) + } + } + } +} diff --git a/src/core/system.rs b/src/core/system.rs index 0d37c5f..c84f086 100644 --- a/src/core/system.rs +++ b/src/core/system.rs @@ -1,27 +1,30 @@ use crate::core::backup::backup_service::BackupService; use crate::core::gui::gui_manager::GuiManager; -use crate::core::infrastructure::actor_system::ActorSystem; use crate::core::infrastructure::app_config::AppConfig; +use crate::core::infrastructure::communication_manager::CommunicationManager; use crate::core::infrastructure::database_manager::DatabaseManager; use crate::core::infrastructure::io_manager::IOManager; use crate::core::schedule::schedule_service::ScheduleService; +use crate::interface::core::runnable::Runnable; use crate::model::error::Error; use crate::model::log::system::SystemLog; #[cfg(any(target_os = "windows", not(debug_assertions)))] use crate::platform::elevate; use crate::utils::logging::Logging; +use crossbeam_queue::SegQueue; use macros::log; #[cfg(not(debug_assertions))] use privilege::user::privileged; #[cfg(not(debug_assertions))] use std::process; use std::sync::Arc; +use tokio::sync::oneshot; pub struct System { - app_config: Arc, - io_manager: Arc, - database_manager: Arc, - actor_system: Arc, + backup_service: Arc, + schedule_service: Arc, + gui_manager: Arc, + shutdowns: SegQueue>, } impl System { @@ -29,12 +32,29 @@ impl System { let app_config = Arc::new(AppConfig::new()?); let io_manager = Arc::new(IOManager::new(app_config.clone())); let database_manager = Arc::new(DatabaseManager::new().await?); - let actor_system = Arc::new(ActorSystem::new()); + let communication_manager = Arc::new(CommunicationManager::new(app_config.clone())); + let backup_service = Arc::new( + BackupService::new( + app_config.clone(), + io_manager.clone(), + communication_manager.clone(), + ) + .await, + ); + let schedule_service = Arc::new( + ScheduleService::new( + app_config.clone(), + database_manager.clone(), + communication_manager.clone(), + ) + .await?, + ); + let gui_manager = Arc::new(GuiManager::new(app_config, communication_manager)); let system = Self { - app_config, - io_manager, - database_manager, - actor_system, + backup_service, + schedule_service, + gui_manager, + shutdowns: SegQueue::new(), }; Ok(system) } @@ -43,24 +63,22 @@ impl System { Logging::initialize().await; log!(SystemLog::Initializing); Self::elevate_privileges()?; - let app_config = self.app_config.clone(); - let io_manager = self.io_manager.clone(); - let database_manager = self.database_manager.clone(); - let actor_system = self.actor_system.clone(); - BackupService::init(app_config.clone(), io_manager.clone(), actor_system.clone()).await; - ScheduleService::init( - app_config.clone(), - database_manager.clone(), - actor_system.clone(), - ) - .await?; - let gui_manager = Arc::new(GuiManager::new(app_config.clone(), actor_system.clone())); + let backup_service = self.backup_service.clone(); + let schedule_service = self.schedule_service.clone(); + let gui_manager = self.gui_manager.clone(); + backup_service.register_services().await; + schedule_service.register_services().await; + let schedule_service_shutdown = schedule_service.run().await; + self.shutdowns.push(schedule_service_shutdown); log!(SystemLog::InitializeComplete); gui_manager.start().await } - pub fn shutdown(&self) { - self.actor_system.shutdown(); + pub async fn shutdown(&self) { + self.backup_service.shutdown().await; + while let Some(shutdown) = self.shutdowns.pop() { + let _ = shutdown.send(()); + } } fn elevate_privileges() -> Result<(), Error> { diff --git a/src/interface/actor/actor.rs b/src/interface/actor/actor.rs deleted file mode 100644 index 085220f..0000000 --- a/src/interface/actor/actor.rs +++ /dev/null @@ -1,14 +0,0 @@ -use crate::interface::actor::message::Message; -use crate::model::error::Error; -use async_trait::async_trait; - -#[async_trait] -pub trait Actor: Send + 'static { - type Message: Message; - async fn pre_start(&mut self); - async fn post_stop(&mut self); - async fn receive( - &mut self, - message: Self::Message, - ) -> Result<::Response, Error>; -} diff --git a/src/interface/actor/mod.rs b/src/interface/actor/mod.rs deleted file mode 100644 index 7df36ef..0000000 --- a/src/interface/actor/mod.rs +++ /dev/null @@ -1,2 +0,0 @@ -pub mod actor; -pub mod message; diff --git a/src/interface/communication/command.rs b/src/interface/communication/command.rs new file mode 100644 index 0000000..3ad66f4 --- /dev/null +++ b/src/interface/communication/command.rs @@ -0,0 +1,15 @@ +use crate::interface::communication::message::Message; +use crate::model::error::Error; +use async_trait::async_trait; +use std::any::Any; +use std::pin::Pin; + +pub type CommandFuture = Pin> + Send + 'static>>; +pub type CommandHandlerFn = Box) -> CommandFuture + Send + Sync>; + +pub trait Command: Message {} + +#[async_trait] +pub trait CommandHandler { + async fn handle_command(&self, command: C) -> Result<(), Error>; +} diff --git a/src/interface/communication/event.rs b/src/interface/communication/event.rs new file mode 100644 index 0000000..faeaa07 --- /dev/null +++ b/src/interface/communication/event.rs @@ -0,0 +1,9 @@ +use crate::model::error::Error; +use std::any::Any; + +pub trait Event: Send + Clone + 'static {} + +pub trait EventBroadcaster: Send + Sync { + fn subscribe_typed(&self) -> Box; + fn broadcast_event(&self, event: Box) -> Result<(), Error>; +} diff --git a/src/interface/actor/message.rs b/src/interface/communication/message.rs similarity index 100% rename from src/interface/actor/message.rs rename to src/interface/communication/message.rs diff --git a/src/interface/communication/mod.rs b/src/interface/communication/mod.rs new file mode 100644 index 0000000..fe42cba --- /dev/null +++ b/src/interface/communication/mod.rs @@ -0,0 +1,4 @@ +pub mod command; +pub mod event; +pub mod message; +pub mod query; diff --git a/src/interface/communication/query.rs b/src/interface/communication/query.rs new file mode 100644 index 0000000..3e7568f --- /dev/null +++ b/src/interface/communication/query.rs @@ -0,0 +1,15 @@ +use crate::interface::communication::message::Message; +use crate::model::error::Error; +use async_trait::async_trait; +use std::any::Any; +use std::pin::Pin; + +pub type QueryFuture = Pin, Error>> + Send + 'static>>; +pub type QueryHandlerFn = Box) -> QueryFuture + Send + Sync>; + +pub trait Query: Message {} + +#[async_trait] +pub trait QueryHandler { + async fn handle_query(&self, query: Q) -> Result; +} diff --git a/src/interface/file_system.rs b/src/interface/core/file_system.rs similarity index 99% rename from src/interface/file_system.rs rename to src/interface/core/file_system.rs index db45a6e..bcec57b 100644 --- a/src/interface/file_system.rs +++ b/src/interface/core/file_system.rs @@ -1,7 +1,7 @@ use crate::model::error::io::IOError; use crate::model::error::system::SystemError; use crate::model::error::Error; -use crate::model::core::backup::backup_execution::HashType; +use crate::model::core::backup::execution::HashType; use crate::platform::attributes::*; use crate::utils::file_hash::*; use async_trait::async_trait; diff --git a/src/interface/core/mod.rs b/src/interface/core/mod.rs new file mode 100644 index 0000000..a758f05 --- /dev/null +++ b/src/interface/core/mod.rs @@ -0,0 +1,2 @@ +pub mod file_system; +pub mod runnable; diff --git a/src/interface/core/runnable.rs b/src/interface/core/runnable.rs new file mode 100644 index 0000000..7b2bc27 --- /dev/null +++ b/src/interface/core/runnable.rs @@ -0,0 +1,16 @@ +use async_trait::async_trait; +use std::sync::Arc; +use tokio::sync::oneshot; + +#[async_trait] +pub trait Runnable: 'static { + async fn run(self: Arc) -> oneshot::Sender<()> { + let (shutdown_tx, shutdown_rx) = oneshot::channel(); + + tokio::spawn(self.run_impl(shutdown_rx)); + + shutdown_tx + } + + async fn run_impl(self: Arc, shutdown_rx: oneshot::Receiver<()>); +} diff --git a/src/interface/mod.rs b/src/interface/mod.rs index e5e1678..7a7d75d 100644 --- a/src/interface/mod.rs +++ b/src/interface/mod.rs @@ -1,3 +1,3 @@ -pub mod actor; -pub mod file_system; +pub mod communication; +pub mod core; pub mod repository; diff --git a/src/interface/repository/schedule.rs b/src/interface/repository/schedule.rs index efd66dc..15ed47a 100644 --- a/src/interface/repository/schedule.rs +++ b/src/interface/repository/schedule.rs @@ -1,5 +1,5 @@ use crate::core::infrastructure::database_manager::DatabaseManager; -use crate::model::core::schedule::backup_schedule::BackupSchedule; +use crate::model::core::schedule::schedule::Schedule; use crate::model::error::Error; use crate::model::error::database::DatabaseError; use crate::model::error::misc::MiscError; @@ -8,11 +8,11 @@ use uuid::Uuid; pub trait ScheduleRepository { async fn create_backup_schedule_table(&self) -> Result<(), Error>; - async fn create_backup_schedule(&self, backup_schedule: &BackupSchedule) -> Result<(), Error>; - async fn modify_backup_schedule(&self, backup_schedule: &BackupSchedule) -> Result<(), Error>; + async fn create_backup_schedule(&self, backup_schedule: &Schedule) -> Result<(), Error>; + async fn modify_backup_schedule(&self, backup_schedule: &Schedule) -> Result<(), Error>; async fn remove_backup_schedule(&self, uuid: Uuid) -> Result<(), Error>; - async fn get_backup_schedule(&self, uuid: Uuid) -> Result, Error>; - async fn get_all_backup_schedules(&self) -> Result, Error>; + async fn get_backup_schedule(&self, uuid: Uuid) -> Result, Error>; + async fn get_all_backup_schedules(&self) -> Result, Error>; } impl ScheduleRepository for DatabaseManager { @@ -43,7 +43,7 @@ impl ScheduleRepository for DatabaseManager { Ok(()) } - async fn create_backup_schedule(&self, backup_schedule: &BackupSchedule) -> Result<(), Error> { + async fn create_backup_schedule(&self, backup_schedule: &Schedule) -> Result<(), Error> { let pool = self.get_pool(); sqlx::query( r#" @@ -99,7 +99,7 @@ impl ScheduleRepository for DatabaseManager { Ok(()) } - async fn modify_backup_schedule(&self, backup_schedule: &BackupSchedule) -> Result<(), Error> { + async fn modify_backup_schedule(&self, backup_schedule: &Schedule) -> Result<(), Error> { let pool = self.get_pool(); sqlx::query( r#" @@ -164,7 +164,7 @@ impl ScheduleRepository for DatabaseManager { Ok(()) } - async fn get_backup_schedule(&self, uuid: Uuid) -> Result, Error> { + async fn get_backup_schedule(&self, uuid: Uuid) -> Result, Error> { let pool = self.get_pool(); let row = sqlx::query( r#" @@ -215,7 +215,7 @@ impl ScheduleRepository for DatabaseManager { let interval = serde_json::from_str(&interval_str) .map_err(MiscError::DeserializeError)?; - Ok(Some(BackupSchedule { + Ok(Some(Schedule { uuid, name: row.get("name"), state, @@ -235,7 +235,7 @@ impl ScheduleRepository for DatabaseManager { } } - async fn get_all_backup_schedules(&self) -> Result, Error> { + async fn get_all_backup_schedules(&self) -> Result, Error> { let pool = self.get_pool(); let rows = sqlx::query( r#" @@ -285,7 +285,7 @@ impl ScheduleRepository for DatabaseManager { let interval = serde_json::from_str(&interval_str) .map_err(MiscError::DeserializeError)?; - schedules.push(BackupSchedule { + schedules.push(Schedule { uuid, name: row.get("name"), state, diff --git a/src/main.rs b/src/main.rs index be8ddda..5cc4af7 100644 --- a/src/main.rs +++ b/src/main.rs @@ -13,6 +13,6 @@ mod utils; async fn main() -> Result<(), Box> { let system = System::new().await?; system.run().await?; - system.shutdown(); + system.shutdown().await; Ok(()) } diff --git a/src/model/config.rs b/src/model/config.rs index 15905b8..6acae16 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -13,4 +13,5 @@ pub struct Config { pub default_wakeup_time: i64, // second pub max_concurrency: u8, // number pub max_file_operations: usize, // number + pub channel_capacity: usize, } diff --git a/src/model/core/actor/actor_ref.rs b/src/model/core/actor/actor_ref.rs deleted file mode 100644 index 0545a9c..0000000 --- a/src/model/core/actor/actor_ref.rs +++ /dev/null @@ -1,44 +0,0 @@ -use crate::interface::actor::message::Message; -use crate::model::core::actor::envelope::Envelope; -use crate::model::error::actor::ActorError; -use crate::model::error::Error; -use tokio::sync::{mpsc, oneshot}; - -pub struct ActorRef { - tx: mpsc::UnboundedSender>, -} - -impl ActorRef { - pub fn new(tx: mpsc::UnboundedSender>) -> Self { - Self { tx } - } - - pub async fn tell(&self, message: M) -> Result<(), Error> { - let envelope = Envelope::Tell(message); - self.tx - .send(envelope) - .map_err(|_| ActorError::SendMessageError)?; - Ok(()) - } - - pub async fn ask(&self, message: M) -> Result { - let (reply_tx, reply_rx) = oneshot::channel::(); - let envelope = Envelope::Ask { - message, - reply_to: reply_tx, - }; - self.tx - .send(envelope) - .map_err(|_| ActorError::SendMessageError)?; - let reply = reply_rx.await.map_err(|_| ActorError::ActorNotResponding)?; - Ok(reply) - } -} - -impl Clone for ActorRef { - fn clone(&self) -> Self { - Self { - tx: self.tx.clone(), - } - } -} diff --git a/src/model/core/actor/actor_runtime.rs b/src/model/core/actor/actor_runtime.rs deleted file mode 100644 index 8874ab8..0000000 --- a/src/model/core/actor/actor_runtime.rs +++ /dev/null @@ -1,59 +0,0 @@ -use crate::interface::actor::actor::Actor; -use crate::model::core::actor::actor_ref::ActorRef; -use crate::model::core::actor::envelope::Envelope; -use crate::model::error::actor::ActorError; -use macros::log; -use tokio::select; -use tokio::sync::{mpsc, oneshot}; - -pub struct ActorRuntime { - actor: A, - rx: mpsc::UnboundedReceiver>, -} - -impl ActorRuntime { - pub fn new(actor: A) -> (Self, ActorRef) { - let (tx, rx) = mpsc::unbounded_channel(); - let actor_ref = ActorRef::new(tx); - let runtime = Self { actor, rx }; - (runtime, actor_ref) - } - - pub async fn run(mut self) -> oneshot::Sender<()> { - let (shutdown_tx, mut shutdown_rx) = oneshot::channel(); - tokio::spawn(async move { - self.actor.pre_start().await; - loop { - select! { - envelope = self.rx.recv() => { - match envelope { - Some(Envelope::Tell(message)) => { - if self.actor.receive(message).await.is_err() { - log!(ActorError::SendMessageError); - } - } - Some(Envelope::Ask { message, reply_to }) => { - match self.actor.receive(message).await { - Ok(response) => { - if reply_to.send(response).is_err() { - log!(ActorError::SendMessageError); - } - } - Err(_) => { - log!(ActorError::SendMessageError); - } - } - } - None => break, - } - } - _ = &mut shutdown_rx => { - break; - } - } - } - self.actor.post_stop().await; - }); - shutdown_tx - } -} diff --git a/src/model/core/actor/envelope.rs b/src/model/core/actor/envelope.rs deleted file mode 100644 index 16ac101..0000000 --- a/src/model/core/actor/envelope.rs +++ /dev/null @@ -1,10 +0,0 @@ -use crate::interface::actor::message::Message; -use tokio::sync::oneshot; - -pub enum Envelope { - Tell(M), - Ask { - message: M, - reply_to: oneshot::Sender, - } -} diff --git a/src/model/core/actor/mod.rs b/src/model/core/actor/mod.rs deleted file mode 100644 index 4a63d81..0000000 --- a/src/model/core/actor/mod.rs +++ /dev/null @@ -1,3 +0,0 @@ -pub mod actor_ref; -pub mod actor_runtime; -pub mod envelope; diff --git a/src/model/core/backup/communication.rs b/src/model/core/backup/communication.rs new file mode 100644 index 0000000..f743df4 --- /dev/null +++ b/src/model/core/backup/communication.rs @@ -0,0 +1,33 @@ +use crate::interface::communication::command::Command; +use crate::interface::communication::message::Message; +use crate::interface::communication::query::Query; +use crate::model::core::backup::execution::Execution; +use uuid::Uuid; + +pub enum BackupCommand { + AddExecution(Execution), + RemoveExecution(Uuid), + StartExecution(Uuid), + SuspendExecution(Uuid), + ResumeExecution(Uuid), +} + +impl Message for BackupCommand { + type Response = (); +} + +impl Command for BackupCommand {} + +pub enum BackupQuery { + GetExecutions, +} + +impl Message for BackupQuery { + type Response = BackupQueryResponse; +} + +impl Query for BackupQuery {} + +pub enum BackupQueryResponse { + GetExecutions(Vec<(Uuid, Execution)>), +} diff --git a/src/model/core/backup/backup_execution.rs b/src/model/core/backup/execution.rs similarity index 97% rename from src/model/core/backup/backup_execution.rs rename to src/model/core/backup/execution.rs index 025d68d..e73d6d8 100644 --- a/src/model/core/backup/backup_execution.rs +++ b/src/model/core/backup/execution.rs @@ -46,7 +46,7 @@ pub struct BackupOptions { } #[derive(Debug, Clone)] -pub struct BackupExecution { +pub struct Execution { pub uuid: Uuid, pub state: BackupState, pub source_path: PathBuf, diff --git a/src/model/core/backup/message.rs b/src/model/core/backup/message.rs deleted file mode 100644 index c1bfbc4..0000000 --- a/src/model/core/backup/message.rs +++ /dev/null @@ -1,29 +0,0 @@ -use crate::interface::actor::message::Message; -use crate::model::core::backup::backup_execution::BackupExecution; -use uuid::Uuid; - -pub enum BackupServiceMessage { - ServiceCall(ServiceCallMessage), -} - -pub enum BackupServiceResponse { - ServiceCall(ServiceCallResponse), - None, -} - -impl Message for BackupServiceMessage { - type Response = BackupServiceResponse; -} - -pub enum ServiceCallMessage { - AddExecution(BackupExecution), - RemoveExecution(Uuid), - StartExecution(Uuid), - SuspendExecution(Uuid), - ResumeExecution(Uuid), - GetExecutions, -} - -pub enum ServiceCallResponse { - GetExecutions(Vec<(Uuid, BackupExecution)>) -} diff --git a/src/model/core/backup/mod.rs b/src/model/core/backup/mod.rs index 0629796..bc9b0a5 100644 --- a/src/model/core/backup/mod.rs +++ b/src/model/core/backup/mod.rs @@ -1,3 +1,3 @@ -pub mod backup_execution; -pub mod message; +pub mod execution; pub mod progress_data; +pub mod communication; diff --git a/src/model/core/gui/communication.rs b/src/model/core/gui/communication.rs new file mode 100644 index 0000000..fb7c28a --- /dev/null +++ b/src/model/core/gui/communication.rs @@ -0,0 +1,29 @@ +use crate::interface::communication::event::Event; +use crate::model::error::Error; +use std::path::PathBuf; +use uuid::Uuid; + +#[derive(Clone)] +pub struct FolderProcess { + pub uuid: Uuid, + pub folder: PathBuf, +} + +impl Event for FolderProcess {} + +#[derive(Clone)] +pub struct ExecutionProgress { + pub uuid: Uuid, + pub processed_files: usize, + pub error_count: usize, +} + +impl Event for ExecutionProgress {} + +#[derive(Clone)] +pub struct ExecutionErrors { + pub uuid: Uuid, + pub errors: Vec, +} + +impl Event for ExecutionErrors {} diff --git a/src/model/core/gui/message.rs b/src/model/core/gui/message.rs deleted file mode 100644 index b21df89..0000000 --- a/src/model/core/gui/message.rs +++ /dev/null @@ -1,25 +0,0 @@ -use crate::interface::actor::message::Message; -use crate::model::error::Error; -use std::path::PathBuf; -use uuid::Uuid; - -#[derive(Clone)] -pub enum GuiMessage { - FolderProcess { - uuid: Uuid, - folder: PathBuf, - }, - ExecutionProgress { - uuid: Uuid, - processed_files: usize, - error_count: usize, - }, - ExecutionErrors { - uuid: Uuid, - errors: Vec, - }, -} - -impl Message for GuiMessage { - type Response = (); -} diff --git a/src/model/core/gui/mod.rs b/src/model/core/gui/mod.rs index e216a50..7575266 100644 --- a/src/model/core/gui/mod.rs +++ b/src/model/core/gui/mod.rs @@ -1 +1 @@ -pub mod message; +pub mod communication; \ No newline at end of file diff --git a/src/model/core/infrastructure/event_broadcaster.rs b/src/model/core/infrastructure/event_broadcaster.rs new file mode 100644 index 0000000..d029343 --- /dev/null +++ b/src/model/core/infrastructure/event_broadcaster.rs @@ -0,0 +1,22 @@ +use crate::interface::communication::event::Event; +use crate::interface::communication::event::EventBroadcaster; +use crate::model::error::misc::MiscError; +use crate::model::error::Error; +use std::any::Any; +use tokio::sync::broadcast; + +pub struct TypedEventBroadcaster { + pub sender: broadcast::Sender, +} + +impl EventBroadcaster for TypedEventBroadcaster { + fn subscribe_typed(&self) -> Box { + Box::new(self.sender.subscribe()) + } + + fn broadcast_event(&self, event: Box) -> Result<(), Error> { + let typed_event = *event.downcast::().map_err(|_| MiscError::TypeMismatch)?; + let _ = self.sender.send(typed_event); + Ok(()) + } +} diff --git a/src/model/core/infrastructure/mod.rs b/src/model/core/infrastructure/mod.rs new file mode 100644 index 0000000..dfd2e83 --- /dev/null +++ b/src/model/core/infrastructure/mod.rs @@ -0,0 +1 @@ +pub mod event_broadcaster; \ No newline at end of file diff --git a/src/model/core/mod.rs b/src/model/core/mod.rs index 3090ae6..392a123 100644 --- a/src/model/core/mod.rs +++ b/src/model/core/mod.rs @@ -1,4 +1,4 @@ -pub mod actor; pub mod backup; pub mod gui; +pub mod infrastructure; pub mod schedule; diff --git a/src/model/core/schedule/communication.rs b/src/model/core/schedule/communication.rs new file mode 100644 index 0000000..6b4aef0 --- /dev/null +++ b/src/model/core/schedule/communication.rs @@ -0,0 +1,45 @@ +use uuid::Uuid; +use crate::interface::communication::command::Command; +use crate::interface::communication::message::Message; +use crate::interface::communication::query::Query; +use crate::model::core::schedule::schedule::Schedule; + +pub enum ScheduleManagerCommand { + AddSchedule(Schedule), + ModifySchedule(Schedule), + RemoveSchedule(Uuid), + ActivateSchedule(Uuid), + PauseSchedule(Uuid), + DisableSchedule(Uuid), + ExecuteReadySchedules, +} + +impl Message for ScheduleManagerCommand { + type Response = (); +} + +impl Command for ScheduleManagerCommand {} + +pub enum ScheduleManagerQuery { + GetSchedules, +} + +impl Message for ScheduleManagerQuery { + type Response = ScheduleManagerQueryResponse; +} + +impl Query for ScheduleManagerQuery {} + +pub enum ScheduleManagerQueryResponse { + GetSchedules(Vec), +} + +pub enum ScheduleTimerCommand { + RefreshTimer +} + +impl Message for ScheduleTimerCommand { + type Response = (); +} + +impl Command for ScheduleTimerCommand {} diff --git a/src/model/core/schedule/message.rs b/src/model/core/schedule/message.rs deleted file mode 100644 index 02df6d5..0000000 --- a/src/model/core/schedule/message.rs +++ /dev/null @@ -1,35 +0,0 @@ -use uuid::Uuid; -use crate::interface::actor::message::Message; -use crate::model::core::schedule::backup_schedule::BackupSchedule; - -pub enum ScheduleServiceMessage { - UnitNotification(UnitNotificationMessage), - ServiceCall(ServiceCallMessage), -} - -pub enum ScheduleServiceResponse { - ServiceCall(ServiceCallResponse), - None, -} - -pub enum UnitNotificationMessage { - CheckSchedule, -} - -pub enum ServiceCallMessage { - AddSchedule(BackupSchedule), - ModifySchedule(BackupSchedule), - RemoveSchedule(Uuid), - ActivateSchedule(Uuid), - PauseSchedule(Uuid), - DisableSchedule(Uuid), - GetSchedules, -} - -pub enum ServiceCallResponse { - GetSchedules(Vec), -} - -impl Message for ScheduleServiceMessage { - type Response = ScheduleServiceResponse; -} diff --git a/src/model/core/schedule/mod.rs b/src/model/core/schedule/mod.rs index befe86e..e67367f 100644 --- a/src/model/core/schedule/mod.rs +++ b/src/model/core/schedule/mod.rs @@ -1,2 +1,2 @@ -pub mod backup_schedule; -pub mod message; +pub mod schedule; +pub mod communication; diff --git a/src/model/core/schedule/backup_schedule.rs b/src/model/core/schedule/schedule.rs similarity index 87% rename from src/model/core/schedule/backup_schedule.rs rename to src/model/core/schedule/schedule.rs index a0c0536..27a7a59 100644 --- a/src/model/core/schedule/backup_schedule.rs +++ b/src/model/core/schedule/schedule.rs @@ -1,4 +1,4 @@ -use crate::model::core::backup::backup_execution::*; +use crate::model::core::backup::execution::*; use chrono::NaiveDateTime; use serde::{Deserialize, Serialize}; use std::path::PathBuf; @@ -20,7 +20,7 @@ pub enum ScheduleInterval { } #[derive(Debug, Clone)] -pub struct BackupSchedule { +pub struct Schedule { pub uuid: Uuid, pub name: String, pub state: ScheduleState, @@ -36,9 +36,9 @@ pub struct BackupSchedule { pub updated_at: NaiveDateTime, } -impl BackupSchedule { - pub fn to_execution(&self) -> BackupExecution { - BackupExecution { +impl Schedule { + pub fn to_execution(&self) -> Execution { + Execution { uuid: self.uuid, state: BackupState::Pending, source_path: self.source_path.clone(), diff --git a/src/model/error/actor.rs b/src/model/error/actor.rs deleted file mode 100644 index 7f9371e..0000000 --- a/src/model/error/actor.rs +++ /dev/null @@ -1,15 +0,0 @@ -use macros::traceable; - -traceable! { - ActorError { - #[no_source] - #[error("Actor not found")] - ActorNotFound => tracing::Level::ERROR, - #[no_source] - #[error("Actor not responding")] - ActorNotResponding => tracing::Level::WARN, - #[no_source] - #[error("Failed to send message to actor")] - SendMessageError => tracing::Level::ERROR, - } -} diff --git a/src/model/error/misc.rs b/src/model/error/misc.rs index 2812190..9666a09 100644 --- a/src/model/error/misc.rs +++ b/src/model/error/misc.rs @@ -7,7 +7,7 @@ traceable! { #[error("Failed to serialize object")] SerializeError => tracing::Level::ERROR, - + #[error("Failed to deserialize object")] DeserializeError => tracing::Level::ERROR, @@ -17,5 +17,25 @@ traceable! { #[no_source] #[error("Assert file not found")] AssertFileNotFound => tracing::Level::ERROR, + + #[no_source] + #[error("Handler not found")] + HandlerNotFound => tracing::Level::ERROR, + + #[no_source] + #[error("Type mismatch")] + TypeMismatch => tracing::Level::ERROR, + + #[no_source] + #[error("Type not registered")] + TypeNotRegistered => tracing::Level::ERROR, + + #[no_source] + #[error("Channel closed")] + ChannelClosed => tracing::Level::ERROR, + + #[no_source] + #[error("Channel empty")] + ChannelEmpty => tracing::Level::INFO, } } diff --git a/src/model/error/mod.rs b/src/model/error/mod.rs index 0d50b92..fea2212 100644 --- a/src/model/error/mod.rs +++ b/src/model/error/mod.rs @@ -1,11 +1,9 @@ -pub mod actor; pub mod database; pub mod io; pub mod misc; pub mod system; pub mod task; -use crate::model::error::actor::ActorError; use crate::model::error::database::DatabaseError; use crate::model::error::io::IOError; use crate::model::error::misc::MiscError; @@ -15,8 +13,6 @@ use serde::{Deserialize, Serialize}; #[derive(Clone, Debug, thiserror::Error, Serialize, Deserialize)] pub enum Error { - #[error("{0}")] - Actor(ActorError), #[error("{0}")] Database(DatabaseError), #[error("{0}")] @@ -29,12 +25,6 @@ pub enum Error { Task(TaskError), } -impl From for Error { - fn from(error: ActorError) -> Self { - Self::Actor(error) - } -} - impl From for Error { fn from(error: DatabaseError) -> Self { Self::Database(error) diff --git a/src/platform/linux/file_system.rs b/src/platform/linux/file_system.rs index 07d319f..6d515e8 100644 --- a/src/platform/linux/file_system.rs +++ b/src/platform/linux/file_system.rs @@ -1,4 +1,4 @@ -use crate::interface::file_system::FileSystemTrait; +use crate::interface::core::file_system::FileSystemTrait; use crate::model::error::Error; use crate::model::error::io::IOError; use crate::model::error::system::SystemError; diff --git a/src/platform/windows/file_system.rs b/src/platform/windows/file_system.rs index 9df3ea2..84e3c91 100644 --- a/src/platform/windows/file_system.rs +++ b/src/platform/windows/file_system.rs @@ -1,4 +1,4 @@ -use crate::interface::file_system::FileSystemTrait; +use crate::interface::core::file_system::FileSystemTrait; use crate::model::error::Error; use crate::model::error::io::IOError; use crate::model::error::misc::MiscError; diff --git a/src/ui/common.rs b/src/ui/common.rs index cd2e831..074747c 100644 --- a/src/ui/common.rs +++ b/src/ui/common.rs @@ -1,4 +1,4 @@ -use crate::model::core::backup::backup_execution::BackupExecution; +use crate::model::core::backup::execution::Execution; #[derive(Debug, Clone, PartialEq)] pub enum PageType { @@ -21,14 +21,14 @@ pub enum ComparisonModeSelection { #[derive(Debug, Clone)] pub struct ExecutionDisplay { - pub execution: BackupExecution, + pub execution: Execution, pub current_folder: String, pub processed_files: usize, pub error_count: usize, } -impl From for ExecutionDisplay { - fn from(execution: BackupExecution) -> Self { +impl From for ExecutionDisplay { + fn from(execution: Execution) -> Self { Self { execution, current_folder: String::new(), diff --git a/src/ui/execution_page.rs b/src/ui/execution_page.rs index eefa020..e2e3aa8 100644 --- a/src/ui/execution_page.rs +++ b/src/ui/execution_page.rs @@ -1,12 +1,8 @@ -use crate::core::backup::backup_service::BackupService; -use crate::core::infrastructure::actor_system::ActorSystem; use crate::core::infrastructure::app_config::AppConfig; -use crate::model::core::backup::backup_execution::*; -use crate::model::core::backup::message::{ - BackupServiceMessage, BackupServiceResponse, ServiceCallMessage, ServiceCallResponse, -}; -use crate::model::core::gui::message::GuiMessage; -use crate::model::error::actor::ActorError; +use crate::core::infrastructure::communication_manager::CommunicationManager; +use crate::model::core::backup::communication::{BackupCommand, BackupQuery, BackupQueryResponse}; +use crate::model::core::backup::execution::*; +use crate::model::core::gui::communication::{ExecutionErrors, ExecutionProgress, FolderProcess}; use crate::model::error::Error; use crate::ui::common::{ComparisonModeSelection, ExecutionDisplay, FolderSelectionMode}; use dashmap::DashMap; @@ -15,16 +11,19 @@ use egui_file_dialog::FileDialog; use futures::executor::block_on; use std::collections::HashSet; use std::path::PathBuf; -use std::sync::{mpsc, Arc}; +use std::sync::Arc; use std::time::{Duration, Instant}; +use tokio::sync::broadcast; use tracing::error; use uuid::Uuid; pub struct ExecutionPage { app_config: Arc, - actor_system: Arc, + communication_manager: Arc, - message_rx: mpsc::Receiver, + folder_process: broadcast::Receiver, + execution_progress: broadcast::Receiver, + execution_errors: broadcast::Receiver, executions: DashMap, error_messages: DashMap>, @@ -50,13 +49,17 @@ pub struct ExecutionPage { impl ExecutionPage { pub fn new( app_config: Arc, - actor_system: Arc, - message_rx: mpsc::Receiver, - ) -> Self { - Self { + communication_manager: Arc, + ) -> Result { + let folder_process = communication_manager.subscribe_event::()?; + let execution_progress = communication_manager.subscribe_event::()?; + let execution_errors = communication_manager.subscribe_event::()?; + let execution_page = Self { app_config, - actor_system, - message_rx, + communication_manager, + folder_process, + execution_progress, + execution_errors, executions: DashMap::new(), error_messages: DashMap::new(), new_task_source: String::new(), @@ -73,34 +76,30 @@ impl ExecutionPage { show_completed_tasks: true, viewing_errors_for_task: None, last_refresh: None, - } + }; + Ok(execution_page) } fn process_events(&mut self) { - while let Ok(message) = self.message_rx.try_recv() { - match message { - GuiMessage::FolderProcess { uuid, folder } => { - if let Some(mut task_display) = self.executions.get_mut(&uuid) { - task_display.current_folder = folder.to_string_lossy().to_string(); - } - } - GuiMessage::ExecutionProgress { - uuid, - processed_files, - error_count, - } => { - if let Some(mut task_display) = self.executions.get_mut(&uuid) { - task_display.processed_files = processed_files; - task_display.error_count = error_count; - } - } - GuiMessage::ExecutionErrors { uuid, errors } => { - match self.error_messages.get_mut(&uuid) { - Some(mut errors_ref) => errors_ref.extend(errors), - None => { - self.error_messages.insert(uuid, errors); - } - } + while let Ok(event) = self.folder_process.try_recv() { + let FolderProcess { uuid, folder } = event; + if let Some(mut task_display) = self.executions.get_mut(&uuid) { + task_display.current_folder = folder.to_string_lossy().to_string(); + } + } + while let Ok(event) = self.execution_progress.try_recv() { + let ExecutionProgress { uuid, processed_files, error_count } = event; + if let Some(mut task_display) = self.executions.get_mut(&uuid) { + task_display.processed_files = processed_files; + task_display.error_count = error_count; + } + } + while let Ok(event) = self.execution_errors.try_recv() { + let ExecutionErrors { uuid, errors } = event; + match self.error_messages.get_mut(&uuid) { + Some(mut errors_ref) => errors_ref.extend(errors), + None => { + self.error_messages.insert(uuid, errors); } } } @@ -108,20 +107,11 @@ impl ExecutionPage { fn sync_all_execution_states(&mut self) { if let Ok(response) = block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - backup_actor_ref - .ask(BackupServiceMessage::ServiceCall( - ServiceCallMessage::GetExecutions, - )) + self.communication_manager + .send_query(BackupQuery::GetExecutions) .await }) { - let BackupServiceResponse::ServiceCall(service_call) = response else { - return; - }; - let ServiceCallResponse::GetExecutions(latest_executions) = service_call; + let BackupQueryResponse::GetExecutions(latest_executions) = response; let latest_ids: HashSet = latest_executions.iter().map(|(id, _)| *id).collect(); for (task_id, latest_execution) in latest_executions { @@ -143,82 +133,47 @@ impl ExecutionPage { } } - fn handle_add_execution(&mut self, execution: BackupExecution) -> Result<(), Error> { + fn handle_add_execution(&mut self, execution: Execution) -> Result<(), Error> { block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let message = - BackupServiceMessage::ServiceCall(ServiceCallMessage::AddExecution(execution)); - backup_actor_ref - .tell(message) - .await - .map_err(|_| ActorError::SendMessageError)?; + self.communication_manager + .send_command(BackupCommand::AddExecution(execution)) + .await?; Ok(()) }) } fn handle_start_execution(&mut self, uuid: Uuid) -> Result<(), Error> { block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let message = - BackupServiceMessage::ServiceCall(ServiceCallMessage::StartExecution(uuid)); - backup_actor_ref - .tell(message) - .await - .map_err(|_| ActorError::SendMessageError)?; + self.communication_manager + .send_command(BackupCommand::StartExecution(uuid)) + .await?; Ok(()) }) } fn handle_suspend_execution(&mut self, uuid: Uuid) -> Result<(), Error> { block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let message = - BackupServiceMessage::ServiceCall(ServiceCallMessage::SuspendExecution(uuid)); - backup_actor_ref - .tell(message) - .await - .map_err(|_| ActorError::SendMessageError)?; + self.communication_manager + .send_command(BackupCommand::SuspendExecution(uuid)) + .await?; Ok(()) }) } fn handle_resume_execution(&mut self, uuid: Uuid) -> Result<(), Error> { block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let message = - BackupServiceMessage::ServiceCall(ServiceCallMessage::ResumeExecution(uuid)); - backup_actor_ref - .tell(message) - .await - .map_err(|_| ActorError::SendMessageError)?; + self.communication_manager + .send_command(BackupCommand::ResumeExecution(uuid)) + .await?; Ok(()) }) } fn handle_remove_execution(&mut self, uuid: Uuid) -> Result<(), Error> { block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let message = - BackupServiceMessage::ServiceCall(ServiceCallMessage::RemoveExecution(uuid)); - backup_actor_ref - .tell(message) - .await - .map_err(|_| ActorError::SendMessageError)?; + self.communication_manager + .send_command(BackupCommand::RemoveExecution(uuid)) + .await?; Ok(()) }) } @@ -543,7 +498,7 @@ impl ExecutionPage { } }; - let execution = BackupExecution { + let execution = Execution { uuid: Uuid::new_v4(), state: BackupState::Pending, source_path: PathBuf::from(&self.new_task_source), @@ -559,7 +514,8 @@ impl ExecutionPage { match self.handle_add_execution(execution.clone()) { Ok(_) => { - let execution_display = ExecutionDisplay::from(execution.clone()); + let execution_display = + ExecutionDisplay::from(execution.clone()); self.executions.insert(execution.uuid, execution_display); self.reset_form(); } diff --git a/src/ui/schedule_page.rs b/src/ui/schedule_page.rs index c1381d2..ec100e8 100644 --- a/src/ui/schedule_page.rs +++ b/src/ui/schedule_page.rs @@ -1,11 +1,10 @@ -use crate::core::infrastructure::actor_system::ActorSystem; use crate::core::infrastructure::app_config::AppConfig; -use crate::core::schedule::schedule_service::ScheduleService; -use crate::model::core::actor::actor_ref::ActorRef; -use crate::model::core::backup::backup_execution::*; -use crate::model::core::schedule::backup_schedule::*; -use crate::model::core::schedule::message::*; -use crate::model::error::actor::ActorError; +use crate::core::infrastructure::communication_manager::CommunicationManager; +use crate::model::core::backup::execution::*; +use crate::model::core::schedule::communication::{ + ScheduleManagerCommand, ScheduleManagerQuery, ScheduleManagerQueryResponse, +}; +use crate::model::core::schedule::schedule::*; use crate::model::error::Error; use crate::ui::common::{ComparisonModeSelection, FolderSelectionMode}; use eframe::egui; @@ -19,11 +18,9 @@ use uuid::Uuid; pub struct SchedulePage { app_config: Arc, - actor_system: Arc, + communication_manager: Arc, - schedule_service_ref: ActorRef, - - schedules: Vec, + schedules: Vec, new_schedule_name: String, new_schedule_source: String, @@ -45,14 +42,13 @@ pub struct SchedulePage { } impl SchedulePage { - pub fn new(app_config: Arc, actor_system: Arc) -> Result { - let schedule_service_ref = actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; + pub fn new( + app_config: Arc, + communication_manager: Arc, + ) -> Result { let schedule_page = Self { app_config, - actor_system, - schedule_service_ref, + communication_manager, schedules: Vec::new(), new_schedule_name: String::new(), new_schedule_source: String::new(), @@ -75,98 +71,69 @@ impl SchedulePage { fn load_schedules(&mut self) { match block_on(async { - self.schedule_service_ref - .ask(ScheduleServiceMessage::ServiceCall( - ServiceCallMessage::GetSchedules, - )) + self.communication_manager + .send_query(ScheduleManagerQuery::GetSchedules) .await }) { - Ok(ScheduleServiceResponse::ServiceCall(ServiceCallResponse::GetSchedules( - schedules, - ))) => { + Ok(ScheduleManagerQueryResponse::GetSchedules(schedules)) => { self.schedules = schedules; } - Ok(ScheduleServiceResponse::None) => {} Err(err) => { error!("{}", err); } } } - fn handle_add_schedule(&self, schedule: BackupSchedule) -> Result<(), Error> { + fn handle_add_schedule(&self, schedule: Schedule) -> Result<(), Error> { block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let message = - ScheduleServiceMessage::ServiceCall(ServiceCallMessage::AddSchedule(schedule)); - backup_actor_ref.tell(message).await?; + self.communication_manager + .send_command(ScheduleManagerCommand::AddSchedule(schedule)) + .await?; Ok(()) }) } - fn handle_modify_schedule(&self, schedule: BackupSchedule) -> Result<(), Error> { + fn handle_modify_schedule(&self, schedule: Schedule) -> Result<(), Error> { block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let message = - ScheduleServiceMessage::ServiceCall(ServiceCallMessage::ModifySchedule(schedule)); - backup_actor_ref.tell(message).await?; + self.communication_manager + .send_command(ScheduleManagerCommand::ModifySchedule(schedule)) + .await?; Ok(()) }) } fn handle_remove_schedule(&self, uuid: Uuid) -> Result<(), Error> { block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let message = - ScheduleServiceMessage::ServiceCall(ServiceCallMessage::RemoveSchedule(uuid)); - backup_actor_ref.tell(message).await?; + self.communication_manager + .send_command(ScheduleManagerCommand::RemoveSchedule(uuid)) + .await?; Ok(()) }) } fn handle_active_schedule(&self, uuid: Uuid) -> Result<(), Error> { block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let message = - ScheduleServiceMessage::ServiceCall(ServiceCallMessage::ActivateSchedule(uuid)); - backup_actor_ref.tell(message).await?; + self.communication_manager + .send_command(ScheduleManagerCommand::ActivateSchedule(uuid)) + .await?; Ok(()) }) } fn handle_pause_schedule(&self, uuid: Uuid) -> Result<(), Error> { block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let message = - ScheduleServiceMessage::ServiceCall(ServiceCallMessage::PauseSchedule(uuid)); - backup_actor_ref.tell(message).await?; + self.communication_manager + .send_command(ScheduleManagerCommand::PauseSchedule(uuid)) + .await?; Ok(()) }) } fn handle_disable_schedule(&self, uuid: Uuid) -> Result<(), Error> { block_on(async { - let backup_actor_ref = self - .actor_system - .actor_of::() - .ok_or(ActorError::ActorNotFound)?; - let message = - ScheduleServiceMessage::ServiceCall(ServiceCallMessage::DisableSchedule(uuid)); - backup_actor_ref.tell(message).await?; + self.communication_manager + .send_command(ScheduleManagerCommand::DisableSchedule(uuid)) + .await?; Ok(()) }) } @@ -226,7 +193,7 @@ impl SchedulePage { egui::ScrollArea::vertical() .auto_shrink([false; 2]) .show(ui, |ui| { - let schedules_to_show: Vec = self + let schedules_to_show: Vec = self .schedules .iter() .filter(|schedule| { @@ -254,7 +221,7 @@ impl SchedulePage { self.draw_schedule_details_window(ctx); } - fn draw_schedule_item(&mut self, ui: &mut egui::Ui, schedule: &BackupSchedule) { + fn draw_schedule_item(&mut self, ui: &mut egui::Ui, schedule: &Schedule) { egui::Frame::new() .fill(ui.visuals().faint_bg_color) .inner_margin(8.0) @@ -511,7 +478,7 @@ impl SchedulePage { } }; - let schedule = BackupSchedule { + let schedule = Schedule { uuid: Uuid::new_v4(), name: self.new_schedule_name.clone(), state: ScheduleState::Active,