From c8660da1279e755dba6e3b89b331ef8496b9fa57 Mon Sep 17 00:00:00 2001 From: DaLaw2 Date: Sat, 9 Aug 2025 00:31:57 +0800 Subject: [PATCH 1/8] refactor: Introduce actor model framework and reorganize module hierarchy --- src/core/actor/actor_system.rs | 48 ++++++++++++++++++ src/core/actor/mod.rs | 1 + src/core/schedule/timer.rs | 93 ++++++++++++++++++++++++++++++++++ 3 files changed, 142 insertions(+) create mode 100644 src/core/actor/actor_system.rs create mode 100644 src/core/actor/mod.rs create mode 100644 src/core/schedule/timer.rs diff --git a/src/core/actor/actor_system.rs b/src/core/actor/actor_system.rs new file mode 100644 index 0000000..328af4b --- /dev/null +++ b/src/core/actor/actor_system.rs @@ -0,0 +1,48 @@ +use crate::interface::actor::actor::Actor; +use crate::model::core::actor::actor_ref::ActorRef; +use crate::model::core::actor::actor_runtime::ActorRuntime; +use crossbeam_queue::SegQueue; +use dashmap::DashMap; +use std::any::{Any, TypeId}; +use std::mem; +use tokio::sync::oneshot; + +pub struct ActorSystem { + actors: DashMap>, + shutdowns: SegQueue>, +} + +impl ActorSystem { + pub fn new() -> Self { + Self { + actors: DashMap::new(), + shutdowns: SegQueue::new(), + } + } + + pub async fn spawn(&mut 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.push(shutdown); + } + + pub fn shutdown(&mut self) { + let shutdowns = mem::take(&mut self.shutdowns); + for shutdown in shutdowns { + let _ = shutdown.send(()); + } + } + + pub fn actor_of(&self) -> Option> { + let type_id = TypeId::of::(); + self.actors + .get(&type_id)? + .downcast_ref::>() + .cloned() + } +} diff --git a/src/core/actor/mod.rs b/src/core/actor/mod.rs new file mode 100644 index 0000000..6e40f45 --- /dev/null +++ b/src/core/actor/mod.rs @@ -0,0 +1 @@ +pub mod actor_system; diff --git a/src/core/schedule/timer.rs b/src/core/schedule/timer.rs new file mode 100644 index 0000000..669201a --- /dev/null +++ b/src/core/schedule/timer.rs @@ -0,0 +1,93 @@ +use std::sync::Arc; +use tokio::sync::{mpsc, oneshot}; +use chrono::{Duration, Utc}; +use macros::log; +use tracing::error; +use tokio::select; +use tokio::time::sleep; +use crate::core::infrastructure::app_config::AppConfig; +use crate::core::schedule::schedule_manager::ScheduleManager; +use crate::model::core::schedule::backup_schedule::ScheduleState; +use crate::model::error::Error; +use crate::model::error::task::TaskError; + +pub struct ScheduleTimer { + app_config: Arc, + schedule_manager: Arc, + shutdown_rx: Option>, + refresh_rx: mpsc::UnboundedReceiver<()>, +} + +impl ScheduleTimer { + pub fn new( + app_config: Arc, + schedule_manager: Arc, + shutdown_rx: oneshot::Receiver<()>, + refresh_rx: mpsc::UnboundedReceiver<()>, + ) -> Self { + ScheduleTimer { + app_config, + schedule_manager, + shutdown_rx: Some(shutdown_rx), + refresh_rx, + } + } + + pub async fn run(mut self) { + let schedule_manager = self.schedule_manager.clone(); + match self.shutdown_rx.take() { + Some(mut shutdown_rx) => 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_rx.recv() => {} + _ = sleep(sleep_time.to_std().unwrap()) => {} + } + if let Err(err) = schedule_manager.execute_ready_schedule().await { + error!("{}", err); + } + }, + None => log!(TaskError::IllegalRunState), + } + } + + async fn calculate_sleep_duration(&self) -> Result, Error> { + let mut next_time = None; + + let schedules = self.schedule_manager.get_all_schedules().await?; + + for schedule in schedules { + if schedule.state != ScheduleState::Active { + continue; + } + if let Some(schedule_next_time) = schedule.next_run_time { + match next_time { + Some(current_time) => { + if schedule_next_time < current_time { + next_time = Some(schedule_next_time); + } + } + None => next_time = Some(schedule_next_time), + } + } + } + if let Some(schedule_next_time) = next_time { + let now = Utc::now().naive_utc(); + let duration = schedule_next_time.signed_duration_since(now); + Ok(Some(Duration::seconds(duration.num_seconds().max(0)))) + } else { + Ok(None) + } + } +} \ No newline at end of file From e36a3a205bfd343e5367f9109b8935d5248f73ba Mon Sep 17 00:00:00 2001 From: DaLaw2 Date: Sun, 10 Aug 2025 04:17:54 +0800 Subject: [PATCH 2/8] wip: Remove event system and restructure module hierarchy for actor-based architecture --- src/core/actor/actor_system.rs | 48 ------------------ src/core/actor/mod.rs | 1 - src/core/schedule/timer.rs | 93 ---------------------------------- src/model/error/message.rs | 11 ++++ src/model/error/mod.rs | 1 + 5 files changed, 12 insertions(+), 142 deletions(-) delete mode 100644 src/core/actor/actor_system.rs delete mode 100644 src/core/actor/mod.rs delete mode 100644 src/core/schedule/timer.rs create mode 100644 src/model/error/message.rs diff --git a/src/core/actor/actor_system.rs b/src/core/actor/actor_system.rs deleted file mode 100644 index 328af4b..0000000 --- a/src/core/actor/actor_system.rs +++ /dev/null @@ -1,48 +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 crossbeam_queue::SegQueue; -use dashmap::DashMap; -use std::any::{Any, TypeId}; -use std::mem; -use tokio::sync::oneshot; - -pub struct ActorSystem { - actors: DashMap>, - shutdowns: SegQueue>, -} - -impl ActorSystem { - pub fn new() -> Self { - Self { - actors: DashMap::new(), - shutdowns: SegQueue::new(), - } - } - - pub async fn spawn(&mut 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.push(shutdown); - } - - pub fn shutdown(&mut self) { - let shutdowns = mem::take(&mut self.shutdowns); - for shutdown in shutdowns { - let _ = shutdown.send(()); - } - } - - pub fn actor_of(&self) -> Option> { - let type_id = TypeId::of::(); - self.actors - .get(&type_id)? - .downcast_ref::>() - .cloned() - } -} diff --git a/src/core/actor/mod.rs b/src/core/actor/mod.rs deleted file mode 100644 index 6e40f45..0000000 --- a/src/core/actor/mod.rs +++ /dev/null @@ -1 +0,0 @@ -pub mod actor_system; diff --git a/src/core/schedule/timer.rs b/src/core/schedule/timer.rs deleted file mode 100644 index 669201a..0000000 --- a/src/core/schedule/timer.rs +++ /dev/null @@ -1,93 +0,0 @@ -use std::sync::Arc; -use tokio::sync::{mpsc, oneshot}; -use chrono::{Duration, Utc}; -use macros::log; -use tracing::error; -use tokio::select; -use tokio::time::sleep; -use crate::core::infrastructure::app_config::AppConfig; -use crate::core::schedule::schedule_manager::ScheduleManager; -use crate::model::core::schedule::backup_schedule::ScheduleState; -use crate::model::error::Error; -use crate::model::error::task::TaskError; - -pub struct ScheduleTimer { - app_config: Arc, - schedule_manager: Arc, - shutdown_rx: Option>, - refresh_rx: mpsc::UnboundedReceiver<()>, -} - -impl ScheduleTimer { - pub fn new( - app_config: Arc, - schedule_manager: Arc, - shutdown_rx: oneshot::Receiver<()>, - refresh_rx: mpsc::UnboundedReceiver<()>, - ) -> Self { - ScheduleTimer { - app_config, - schedule_manager, - shutdown_rx: Some(shutdown_rx), - refresh_rx, - } - } - - pub async fn run(mut self) { - let schedule_manager = self.schedule_manager.clone(); - match self.shutdown_rx.take() { - Some(mut shutdown_rx) => 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_rx.recv() => {} - _ = sleep(sleep_time.to_std().unwrap()) => {} - } - if let Err(err) = schedule_manager.execute_ready_schedule().await { - error!("{}", err); - } - }, - None => log!(TaskError::IllegalRunState), - } - } - - async fn calculate_sleep_duration(&self) -> Result, Error> { - let mut next_time = None; - - let schedules = self.schedule_manager.get_all_schedules().await?; - - for schedule in schedules { - if schedule.state != ScheduleState::Active { - continue; - } - if let Some(schedule_next_time) = schedule.next_run_time { - match next_time { - Some(current_time) => { - if schedule_next_time < current_time { - next_time = Some(schedule_next_time); - } - } - None => next_time = Some(schedule_next_time), - } - } - } - if let Some(schedule_next_time) = next_time { - let now = Utc::now().naive_utc(); - let duration = schedule_next_time.signed_duration_since(now); - Ok(Some(Duration::seconds(duration.num_seconds().max(0)))) - } else { - Ok(None) - } - } -} \ No newline at end of file diff --git a/src/model/error/message.rs b/src/model/error/message.rs new file mode 100644 index 0000000..b6d8532 --- /dev/null +++ b/src/model/error/message.rs @@ -0,0 +1,11 @@ +use uuid::Uuid; +use crate::interface::actor::message::Message; +use crate::model::error::Error; + +pub enum ErrorMessage { + BackupError(Uuid, Vec), +} + +impl Message for ErrorMessage { + type Response = (); +} diff --git a/src/model/error/mod.rs b/src/model/error/mod.rs index 0d50b92..83a44f8 100644 --- a/src/model/error/mod.rs +++ b/src/model/error/mod.rs @@ -1,6 +1,7 @@ pub mod actor; pub mod database; pub mod io; +pub mod message; pub mod misc; pub mod system; pub mod task; From 6933b76243e9959cbe3dd4afe52cf4c2220f9ab7 Mon Sep 17 00:00:00 2001 From: DaLaw2 Date: Sun, 10 Aug 2025 16:31:28 +0800 Subject: [PATCH 3/8] refactor: Remove unused traits and modules, enhance GUI handling, and refine actor-driven architecture --- src/model/error/message.rs | 11 ----------- src/model/error/mod.rs | 1 - 2 files changed, 12 deletions(-) delete mode 100644 src/model/error/message.rs diff --git a/src/model/error/message.rs b/src/model/error/message.rs deleted file mode 100644 index b6d8532..0000000 --- a/src/model/error/message.rs +++ /dev/null @@ -1,11 +0,0 @@ -use uuid::Uuid; -use crate::interface::actor::message::Message; -use crate::model::error::Error; - -pub enum ErrorMessage { - BackupError(Uuid, Vec), -} - -impl Message for ErrorMessage { - type Response = (); -} diff --git a/src/model/error/mod.rs b/src/model/error/mod.rs index 83a44f8..0d50b92 100644 --- a/src/model/error/mod.rs +++ b/src/model/error/mod.rs @@ -1,7 +1,6 @@ pub mod actor; pub mod database; pub mod io; -pub mod message; pub mod misc; pub mod system; pub mod task; From d5b4e91db6c021d2bc254bc2611d508784e84f54 Mon Sep 17 00:00:00 2001 From: DaLaw2 Date: Sun, 24 Aug 2025 15:21:53 +0800 Subject: [PATCH 4/8] wip: Remove actor-based architecture and add CommunicationManager for cross-service communication --- src/core/backup/backup_engine.rs | 35 +++- src/core/backup/backup_service.rs | 71 +------ src/core/gui/mod.rs | 2 - src/core/infrastructure/actor_system.rs | 56 ------ .../infrastructure/communication_manager.rs | 180 ++++++++++++++++++ src/core/infrastructure/mod.rs | 2 +- src/core/mod.rs | 1 - src/core/schedule/schedule_manager.rs | 2 - src/core/schedule/schedule_service.rs | 90 ++------- src/core/schedule/schedule_timer.rs | 3 - src/core/system.rs | 2 - src/interface/actor/actor.rs | 14 -- src/interface/actor/mod.rs | 2 - src/interface/communication/command.rs | 15 ++ src/interface/communication/event.rs | 9 + .../{actor => communication}/message.rs | 0 src/interface/communication/mod.rs | 4 + src/interface/communication/query.rs | 15 ++ src/interface/core/mod.rs | 2 + src/interface/core/service.rs | 16 ++ src/interface/core/unit.rs | 16 ++ src/interface/mod.rs | 3 +- src/model/config.rs | 1 + src/model/core/actor/actor_ref.rs | 44 ----- src/model/core/actor/actor_runtime.rs | 59 ------ src/model/core/actor/envelope.rs | 10 - src/model/core/actor/mod.rs | 3 - src/model/core/backup/message.rs | 29 --- src/model/core/backup/mod.rs | 1 - src/model/core/gui/message.rs | 25 --- src/model/core/gui/mod.rs | 1 - .../core/infrastructure/event_broadcaster.rs | 22 +++ src/model/core/infrastructure/mod.rs | 1 + src/model/core/mod.rs | 2 +- src/model/core/schedule/message.rs | 35 ---- src/model/core/schedule/mod.rs | 1 - src/model/error/actor.rs | 15 -- src/model/error/misc.rs | 22 ++- src/model/error/mod.rs | 10 - src/ui/execution_page.rs | 6 - src/ui/schedule_page.rs | 5 - 41 files changed, 358 insertions(+), 474 deletions(-) delete mode 100644 src/core/gui/mod.rs delete mode 100644 src/core/infrastructure/actor_system.rs create mode 100644 src/core/infrastructure/communication_manager.rs delete mode 100644 src/interface/actor/actor.rs delete mode 100644 src/interface/actor/mod.rs create mode 100644 src/interface/communication/command.rs create mode 100644 src/interface/communication/event.rs rename src/interface/{actor => communication}/message.rs (100%) create mode 100644 src/interface/communication/mod.rs create mode 100644 src/interface/communication/query.rs create mode 100644 src/interface/core/mod.rs create mode 100644 src/interface/core/service.rs create mode 100644 src/interface/core/unit.rs delete mode 100644 src/model/core/actor/actor_ref.rs delete mode 100644 src/model/core/actor/actor_runtime.rs delete mode 100644 src/model/core/actor/envelope.rs delete mode 100644 src/model/core/actor/mod.rs delete mode 100644 src/model/core/backup/message.rs delete mode 100644 src/model/core/gui/message.rs create mode 100644 src/model/core/infrastructure/event_broadcaster.rs create mode 100644 src/model/core/infrastructure/mod.rs delete mode 100644 src/model/core/schedule/message.rs delete mode 100644 src/model/error/actor.rs diff --git a/src/core/backup/backup_engine.rs b/src/core/backup/backup_engine.rs index c94dd35..6bf9616 100644 --- a/src/core/backup/backup_engine.rs +++ b/src/core/backup/backup_engine.rs @@ -1,14 +1,14 @@ 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::io_manager::IOManager; +use crate::interface::core::unit::Unit; use crate::interface::file_system::FileSystemTrait; use crate::model::core::backup::backup_execution::*; -use crate::model::core::gui::message::GuiMessage; 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; @@ -16,6 +16,7 @@ use macros::log; use std::collections::{HashSet, VecDeque}; use std::path::{Path, PathBuf}; use std::sync::Arc; +use tokio::sync::mpsc::UnboundedReceiver; use tokio::sync::oneshot; use tokio::task::JoinHandle; use tracing::error; @@ -24,7 +25,6 @@ use uuid::Uuid; pub struct BackupEngine { app_config: Arc, io_manager: Arc, - actor_system: Arc, progress_tracker: Arc, executions: Arc>, running_executions: Arc, JoinHandle<()>)>>, @@ -34,13 +34,11 @@ impl BackupEngine { pub fn new( app_config: Arc, io_manager: Arc, - actor_system: Arc, progress_tracker: Arc, ) -> Self { Self { app_config, io_manager, - actor_system, progress_tracker, executions: Arc::new(DashMap::new()), running_executions: Arc::new(DashMap::new()), @@ -391,7 +389,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); @@ -646,3 +645,27 @@ impl Worker { Ok(destination_root.join(relative_path)) } } + +#[async_trait] +impl Unit for BackupEngine { + type Command = (); + type InternalCommand = (); + type Query = (); + type QueryResponse = (); + + fn get_internal_channel(&self) -> UnboundedReceiver { + todo!() + } + + async fn handle_command(&self, command: Self::Command) -> Result<(), Error> { + todo!() + } + + async fn handle_internal_command(&self, command: Self::InternalCommand) -> Result<(), Error> { + todo!() + } + + async fn handle_query(&self, query: Self::Query) -> Result { + todo!() + } +} diff --git a/src/core/backup/backup_service.rs b/src/core/backup/backup_service.rs index 399dd27..32ede39 100644 --- a/src/core/backup/backup_service.rs +++ b/src/core/backup/backup_service.rs @@ -1,82 +1,27 @@ 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::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 crate::interface::core::service::Service; use async_trait::async_trait; use std::sync::Arc; +use tokio::sync::oneshot::Receiver; pub struct BackupService { backup_engine: Arc, } impl BackupService { - pub async fn init( - app_config: Arc, - io_manager: Arc, - actor_system: Arc, - ) { + pub async fn init(app_config: Arc, io_manager: Arc) { let progress_tracker = Arc::new(ProgressTracker::new(io_manager.clone())); - let backup_engine = Arc::new(BackupEngine::new( - app_config, - io_manager, - actor_system.clone(), - progress_tracker, - )); - let backup_service = Self { - backup_engine, - }; - actor_system.spawn(backup_service).await; + let backup_engine = Arc::new(BackupEngine::new(app_config, io_manager, progress_tracker)); + let backup_service = 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; - } - - 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), - )) - } - }, - } +impl Service for BackupService { + async fn run_message_loop(self: Arc, shutdown_rx: Receiver<()>) { + todo!() } } diff --git a/src/core/gui/mod.rs b/src/core/gui/mod.rs deleted file mode 100644 index 464dd49..0000000 --- a/src/core/gui/mod.rs +++ /dev/null @@ -1,2 +0,0 @@ -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..b1ef52e --- /dev/null +++ b/src/core/infrastructure/communication_manager.rs @@ -0,0 +1,180 @@ +use crate::core::infrastructure::app_config::AppConfig; +use crate::interface::communication::command::*; +use crate::interface::communication::event::Event; +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; +use crate::interface::communication::event::EventBroadcaster; + +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_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 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 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 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 has_command_handler(&self) -> bool { + self.command_handlers.contains_key(&TypeId::of::()) + } + + pub fn has_query_handler(&self) -> bool { + self.query_handlers.contains_key(&TypeId::of::()) + } + + pub fn has_event_type(&self) -> bool { + self.event_broadcasters.contains_key(&TypeId::of::()) + } + + pub fn clear_command_handlers(&self) { + self.command_handlers.clear(); + } + + pub fn clear_query_handlers(&self) { + self.query_handlers.clear(); + } + + pub fn clear_event_types(&self) { + self.event_broadcasters.clear(); + } +} + +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/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/mod.rs b/src/core/mod.rs index f8b01fe..3a57e46 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -1,5 +1,4 @@ pub mod backup; -pub mod gui; pub mod infrastructure; pub mod schedule; pub mod system; diff --git a/src/core/schedule/schedule_manager.rs b/src/core/schedule/schedule_manager.rs index bfef95e..ddd339e 100644 --- a/src/core/schedule/schedule_manager.rs +++ b/src/core/schedule/schedule_manager.rs @@ -1,8 +1,6 @@ use crate::core::backup::backup_service::BackupService; -use crate::core::infrastructure::actor_system::ActorSystem; use crate::core::infrastructure::database_manager::DatabaseManager; use crate::interface::repository::schedule::ScheduleRepository; -use crate::model::core::backup::message::{BackupServiceMessage, ServiceCallMessage}; use crate::model::core::schedule::backup_schedule::*; use crate::model::error::Error; use chrono::{Duration, Months, Utc}; diff --git a/src/core/schedule/schedule_service.rs b/src/core/schedule/schedule_service.rs index 8996786..4e76ae8 100644 --- a/src/core/schedule/schedule_service.rs +++ b/src/core/schedule/schedule_service.rs @@ -1,16 +1,14 @@ -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::service::Service; use crate::model::error::Error; use async_trait::async_trait; -use macros::log; +use std::any::Any; use std::sync::{Arc, OnceLock}; +use tokio::sync::oneshot::Receiver; use tokio::sync::{mpsc, oneshot}; pub struct ScheduleService { @@ -47,77 +45,23 @@ impl ScheduleService { } #[async_trait] -impl Actor for ScheduleService { - type Message = ScheduleServiceMessage; +impl Service for ScheduleService { + async fn run_message_loop(self: Arc, shutdown_rx: Receiver<()>) { + todo!() + } +} - async fn pre_start(&mut self) { - 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_trait] +impl CommunicationCapable for ScheduleService { + fn register_handlers(&self, comm: &CommunicationManager) { + todo!() } - 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 handle_command(&self, command: Box) -> Result<(), Error> { + todo!() } - 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), - )) - } - }, - } + async fn handle_query(&self, query: Box) -> Result, Error> { + todo!() } } diff --git a/src/core/schedule/schedule_timer.rs b/src/core/schedule/schedule_timer.rs index 2bfa16e..7145d5a 100644 --- a/src/core/schedule/schedule_timer.rs +++ b/src/core/schedule/schedule_timer.rs @@ -1,9 +1,6 @@ -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::model::error::Error; use chrono::{Duration, Utc}; use std::sync::Arc; diff --git a/src/core/system.rs b/src/core/system.rs index 0d37c5f..c1ece2f 100644 --- a/src/core/system.rs +++ b/src/core/system.rs @@ -1,6 +1,4 @@ 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::database_manager::DatabaseManager; use crate::core::infrastructure::io_manager::IOManager; 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/core/mod.rs b/src/interface/core/mod.rs new file mode 100644 index 0000000..eec556b --- /dev/null +++ b/src/interface/core/mod.rs @@ -0,0 +1,2 @@ +pub mod service; +pub mod unit; diff --git a/src/interface/core/service.rs b/src/interface/core/service.rs new file mode 100644 index 0000000..9857cf6 --- /dev/null +++ b/src/interface/core/service.rs @@ -0,0 +1,16 @@ +use async_trait::async_trait; +use std::sync::Arc; +use tokio::sync::oneshot; + +#[async_trait] +pub trait Service { + 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_message_loop(self: Arc, shutdown_rx: oneshot::Receiver<()>); +} diff --git a/src/interface/core/unit.rs b/src/interface/core/unit.rs new file mode 100644 index 0000000..caa30b7 --- /dev/null +++ b/src/interface/core/unit.rs @@ -0,0 +1,16 @@ +use crate::model::error::Error; +use async_trait::async_trait; +use tokio::sync::mpsc::UnboundedReceiver; + +#[async_trait] +pub trait Unit { + type Command: Send + 'static; + type InternalCommand: Send + 'static; + type Query: Send + 'static; + type QueryResponse: Send + 'static; + + fn get_internal_channel(&self) -> UnboundedReceiver; + async fn handle_command(&self, command: Self::Command) -> Result<(), Error>; + async fn handle_internal_command(&self, command: Self::InternalCommand) -> Result<(), Error>; + async fn handle_query(&self, query: Self::Query) -> Result; +} diff --git a/src/interface/mod.rs b/src/interface/mod.rs index e5e1678..128d71a 100644 --- a/src/interface/mod.rs +++ b/src/interface/mod.rs @@ -1,3 +1,4 @@ -pub mod actor; +pub mod communication; +pub mod core; pub mod file_system; pub mod repository; 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/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..ba27d88 100644 --- a/src/model/core/backup/mod.rs +++ b/src/model/core/backup/mod.rs @@ -1,3 +1,2 @@ pub mod backup_execution; -pub mod message; pub mod progress_data; 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..e69de29 100644 --- a/src/model/core/gui/mod.rs +++ b/src/model/core/gui/mod.rs @@ -1 +0,0 @@ -pub mod message; 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/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..7a0790a 100644 --- a/src/model/core/schedule/mod.rs +++ b/src/model/core/schedule/mod.rs @@ -1,2 +1 @@ pub mod backup_schedule; -pub mod message; 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..8e38645 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, @@ -15,7 +15,23 @@ traceable! { UIPlatformError => tracing::Level::ERROR, #[no_source] - #[error("Assert file not found")] - AssertFileNotFound => tracing::Level::ERROR, + #[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/ui/execution_page.rs b/src/ui/execution_page.rs index eefa020..7432892 100644 --- a/src/ui/execution_page.rs +++ b/src/ui/execution_page.rs @@ -1,12 +1,6 @@ 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::model::error::Error; use crate::ui::common::{ComparisonModeSelection, ExecutionDisplay, FolderSelectionMode}; use dashmap::DashMap; diff --git a/src/ui/schedule_page.rs b/src/ui/schedule_page.rs index c1381d2..3c89b04 100644 --- a/src/ui/schedule_page.rs +++ b/src/ui/schedule_page.rs @@ -1,12 +1,7 @@ -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::model::error::Error; use crate::ui::common::{ComparisonModeSelection, FolderSelectionMode}; use eframe::egui; use egui_file_dialog::FileDialog; From 1fd70f464fcc9384aa1ebf576fde090780797da3 Mon Sep 17 00:00:00 2001 From: DaLaw2 Date: Sun, 24 Aug 2025 23:56:40 +0800 Subject: [PATCH 5/8] wip: Introduce communication interfaces for services and replace actor system architecture --- src/core/backup/backup_engine.rs | 62 ++++++++++++++++-------- src/core/backup/backup_service.rs | 35 ++++++++++++- src/core/gui/mod.rs | 2 + src/core/mod.rs | 1 + src/core/schedule/schedule_manager.rs | 61 ++++++++++++++++++++++- src/core/schedule/schedule_service.rs | 24 +++++---- src/interface/core/service.rs | 2 +- src/interface/core/unit.rs | 12 +++-- src/model/core/backup/communication.rs | 41 ++++++++++++++++ src/model/core/backup/mod.rs | 1 + src/model/core/gui/communication.rs | 29 +++++++++++ src/model/core/gui/mod.rs | 1 + src/model/core/schedule/communication.rs | 45 +++++++++++++++++ src/model/core/schedule/mod.rs | 1 + 14 files changed, 275 insertions(+), 42 deletions(-) create mode 100644 src/core/gui/mod.rs create mode 100644 src/model/core/backup/communication.rs create mode 100644 src/model/core/gui/communication.rs create mode 100644 src/model/core/schedule/communication.rs diff --git a/src/core/backup/backup_engine.rs b/src/core/backup/backup_engine.rs index 6bf9616..b7248ed 100644 --- a/src/core/backup/backup_engine.rs +++ b/src/core/backup/backup_engine.rs @@ -5,6 +5,9 @@ use crate::core::infrastructure::io_manager::IOManager; use crate::interface::core::unit::Unit; use crate::interface::file_system::FileSystemTrait; use crate::model::core::backup::backup_execution::*; +use crate::model::core::backup::communication::{ + BackupCommand, BackupInternalCommand, BackupQuery, BackupQueryResponse, +}; use crate::model::error::system::SystemError; use crate::model::error::task::TaskError; use crate::model::error::Error; @@ -79,14 +82,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 { @@ -98,14 +101,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 { @@ -116,7 +119,7 @@ impl BackupEngine { let (_, (shutdown, handle)) = self .running_executions - .remove(&uuid) + .remove(uuid) .ok_or(TaskError::ExecutionNotFound)?; shutdown .send(()) @@ -125,14 +128,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 { @@ -144,7 +147,7 @@ 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(()) } @@ -648,24 +651,45 @@ impl Worker { #[async_trait] impl Unit for BackupEngine { - type Command = (); - type InternalCommand = (); - type Query = (); - type QueryResponse = (); + type Command = BackupCommand; + type InternalCommand = BackupInternalCommand; + type Query = BackupQuery; fn get_internal_channel(&self) -> UnboundedReceiver { - todo!() + unimplemented!() } async fn handle_command(&self, command: Self::Command) -> Result<(), Error> { - todo!() + 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 fn handle_internal_command(&self, command: Self::InternalCommand) -> Result<(), Error> { - todo!() + match command {} } - async fn handle_query(&self, query: Self::Query) -> Result { - todo!() + async fn handle_query(&self, query: Self::Query) -> 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 32ede39..4b9f99b 100644 --- a/src/core/backup/backup_service.rs +++ b/src/core/backup/backup_service.rs @@ -5,7 +5,14 @@ use crate::core::infrastructure::io_manager::IOManager; use crate::interface::core::service::Service; use async_trait::async_trait; use std::sync::Arc; +use tokio::io::AsyncWriteExt; +use tokio::select; use tokio::sync::oneshot::Receiver; +use crate::interface::communication::command::CommandHandler; +use crate::interface::communication::query::QueryHandler; +use crate::interface::core::unit::Unit; +use crate::model::core::backup::communication::{BackupCommand, BackupQuery, BackupQueryResponse}; +use crate::model::error::Error; pub struct BackupService { backup_engine: Arc, @@ -21,7 +28,31 @@ impl BackupService { #[async_trait] impl Service for BackupService { - async fn run_message_loop(self: Arc, shutdown_rx: Receiver<()>) { - todo!() + async fn process_internal_command(self: Arc, shutdown_rx: Receiver<()>) { + let mut backup_engine_rx = self.backup_engine.get_internal_channel(); + loop { + select! { + _ = &shutdown_rx => { + break + } + _ = backup_engine_rx.recv() => { + + } + } + } + } +} + +#[async_trait] +impl CommandHandler for BackupService { + async fn handle_command(&self, command: BackupCommand) -> Result<(), Error> { + self.backup_engine.handle_command(command) + } +} + +#[async_trait] +impl QueryHandler for BackupService { + async fn handle_query(&self, query: BackupQuery) -> Result { + self.backup_engine.handle_query(query) } } diff --git a/src/core/gui/mod.rs b/src/core/gui/mod.rs new file mode 100644 index 0000000..6fc230f --- /dev/null +++ b/src/core/gui/mod.rs @@ -0,0 +1,2 @@ +pub mod gui_manager; +pub mod gui_message_handler; diff --git a/src/core/mod.rs b/src/core/mod.rs index 3a57e46..f8b01fe 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -1,4 +1,5 @@ pub mod backup; +pub mod gui; pub mod infrastructure; pub mod schedule; pub mod system; diff --git a/src/core/schedule/schedule_manager.rs b/src/core/schedule/schedule_manager.rs index ddd339e..2b205d9 100644 --- a/src/core/schedule/schedule_manager.rs +++ b/src/core/schedule/schedule_manager.rs @@ -6,18 +6,21 @@ use crate::model::error::Error; use chrono::{Duration, Months, Utc}; use dashmap::DashMap; use std::sync::Arc; +use async_trait::async_trait; +use tokio::sync::mpsc::UnboundedReceiver; use uuid::Uuid; +use crate::interface::communication::message::Message; +use crate::interface::core::unit::Unit; +use crate::model::core::schedule::communication::{ScheduleCommand, ScheduleInternalCommand, ScheduleQuery, ScheduleQueryResponse}; pub struct ScheduleManager { database_manager: Arc, - actor_system: Arc, schedules: DashMap, } impl ScheduleManager { pub async fn new( database_manager: Arc, - actor_system: Arc, ) -> Result { let schedules = DashMap::new(); let database_schedules = database_manager.get_all_backup_schedules().await?; @@ -141,3 +144,57 @@ impl ScheduleManager { schedule.next_run_time = new_next_run_time; } } + +#[async_trait] +impl Unit for ScheduleManager { + type Command = ScheduleCommand; + type InternalCommand = ScheduleInternalCommand; + type Query = ScheduleQuery; + + fn get_internal_channel(&self) -> UnboundedReceiver { + todo!() + } + + async fn handle_command(&self, command: Self::Command) -> Result<(), Error> { + match command { + ScheduleCommand::AddSchedule(schedule) => { + self.create_schedule(schedule).await?; + } + ScheduleCommand::ModifySchedule(schedule) => { + self.modify_schedule(schedule).await?; + } + ScheduleCommand::RemoveSchedule(uuid) => { + self.remove_schedule(uuid).await?; + } + ScheduleCommand::ActivateSchedule(uuid) => { + self.active_schedule(uuid).await?; + } + ScheduleCommand::PauseSchedule(uuid) => { + self.pause_schedule(uuid).await?; + } + ScheduleCommand::DisableSchedule(uuid) => { + self.disable_schedule(uuid).await?; + } + } + Ok(()) + } + + async fn handle_internal_command(&self, command: Self::InternalCommand) -> Result<(), Error> { + match command { + ScheduleInternalCommand::TimerNotify => { + self.execute_ready_schedule().await?; + } + _ => {} + } + Ok(()) + } + + async fn handle_query(&self, query: Self::Query) -> Result<::Response, Error> { + match query { + ScheduleQuery::GetSchedules => { + let executions = self.get_all_schedules().await; + Ok(ScheduleQueryResponse::GetSchedules(executions)) + } + } + } +} diff --git a/src/core/schedule/schedule_service.rs b/src/core/schedule/schedule_service.rs index 4e76ae8..0936791 100644 --- a/src/core/schedule/schedule_service.rs +++ b/src/core/schedule/schedule_service.rs @@ -10,6 +10,9 @@ use std::any::Any; use std::sync::{Arc, OnceLock}; use tokio::sync::oneshot::Receiver; use tokio::sync::{mpsc, oneshot}; +use crate::interface::communication::command::CommandHandler; +use crate::interface::communication::query::QueryHandler; +use crate::model::core::schedule::communication::{ScheduleCommand, ScheduleQuery, ScheduleQueryResponse}; pub struct ScheduleService { schedule_manager: Arc, @@ -22,18 +25,16 @@ impl ScheduleService { pub async fn init( app_config: Arc, database_manager: Arc, - actor_system: Arc, ) -> Result<(), Error> { 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).await?); + let schedule_timer = Arc::new(ScheduleTimer::new(app_config)); let schedule_service = Self { schedule_manager, schedule_timer, timer_refresh: OnceLock::new(), shutdowns: Vec::new(), }; - actor_system.spawn(schedule_service).await; Ok(()) } @@ -46,22 +47,19 @@ impl ScheduleService { #[async_trait] impl Service for ScheduleService { - async fn run_message_loop(self: Arc, shutdown_rx: Receiver<()>) { + async fn process_internal_command(self: Arc, shutdown_rx: Receiver<()>) { todo!() } } -#[async_trait] -impl CommunicationCapable for ScheduleService { - fn register_handlers(&self, comm: &CommunicationManager) { - todo!() - } - - async fn handle_command(&self, command: Box) -> Result<(), Error> { +impl CommandHandler for ScheduleService { + async fn handle_command(&self, command: ScheduleCommand) -> Result<(), Error> { todo!() } +} - async fn handle_query(&self, query: Box) -> Result, Error> { +impl QueryHandler for ScheduleService { + async fn handle_query(&self, query: ScheduleQuery) -> Result { todo!() } } diff --git a/src/interface/core/service.rs b/src/interface/core/service.rs index 9857cf6..66af958 100644 --- a/src/interface/core/service.rs +++ b/src/interface/core/service.rs @@ -12,5 +12,5 @@ pub trait Service { shutdown_tx } - async fn run_message_loop(self: Arc, shutdown_rx: oneshot::Receiver<()>); + async fn process_internal_command(self: Arc, shutdown_rx: oneshot::Receiver<()>); } diff --git a/src/interface/core/unit.rs b/src/interface/core/unit.rs index caa30b7..d2b6425 100644 --- a/src/interface/core/unit.rs +++ b/src/interface/core/unit.rs @@ -1,16 +1,18 @@ +use crate::interface::communication::command::Command; +use crate::interface::communication::message::Message; +use crate::interface::communication::query::Query; use crate::model::error::Error; use async_trait::async_trait; use tokio::sync::mpsc::UnboundedReceiver; #[async_trait] pub trait Unit { - type Command: Send + 'static; - type InternalCommand: Send + 'static; - type Query: Send + 'static; - type QueryResponse: Send + 'static; + type Command: Command; + type InternalCommand: Command; + type Query: Query; fn get_internal_channel(&self) -> UnboundedReceiver; async fn handle_command(&self, command: Self::Command) -> Result<(), Error>; async fn handle_internal_command(&self, command: Self::InternalCommand) -> Result<(), Error>; - async fn handle_query(&self, query: Self::Query) -> Result; + async fn handle_query(&self, query: Self::Query) -> Result<::Response, Error>; } diff --git a/src/model/core/backup/communication.rs b/src/model/core/backup/communication.rs new file mode 100644 index 0000000..31b91c8 --- /dev/null +++ b/src/model/core/backup/communication.rs @@ -0,0 +1,41 @@ +use crate::interface::communication::command::Command; +use crate::interface::communication::message::Message; +use crate::interface::communication::query::Query; +use crate::model::core::backup::backup_execution::BackupExecution; +use uuid::Uuid; + +pub enum BackupCommand { + AddExecution(BackupExecution), + RemoveExecution(Uuid), + StartExecution(Uuid), + SuspendExecution(Uuid), + ResumeExecution(Uuid), +} + +impl Message for BackupCommand { + type Response = (); +} + +impl Command for BackupCommand {} + +pub enum BackupInternalCommand {} + +impl Message for BackupInternalCommand { + type Response = (); +} + +impl Command for BackupInternalCommand {} + +pub enum BackupQuery { + GetExecutions, +} + +impl Message for BackupQuery { + type Response = BackupQueryResponse; +} + +impl Query for BackupQuery {} + +pub enum BackupQueryResponse { + GetExecutions(Vec<(Uuid, BackupExecution)>), +} diff --git a/src/model/core/backup/mod.rs b/src/model/core/backup/mod.rs index ba27d88..f79d421 100644 --- a/src/model/core/backup/mod.rs +++ b/src/model/core/backup/mod.rs @@ -1,2 +1,3 @@ pub mod backup_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..e001c5b --- /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 { + uuid: Uuid, + folder: PathBuf, +} + +impl Event for FolderProcess {} + +#[derive(Clone)] +pub struct ExecutionProgress { + uuid: Uuid, + processed_files: usize, + error_count: usize, +} + +impl Event for ExecutionProgress {} + +#[derive(Clone)] +pub struct ExecutionErrors { + uuid: Uuid, + errors: Vec, +} + +impl Event for ExecutionErrors {} diff --git a/src/model/core/gui/mod.rs b/src/model/core/gui/mod.rs index e69de29..7575266 100644 --- a/src/model/core/gui/mod.rs +++ b/src/model/core/gui/mod.rs @@ -0,0 +1 @@ +pub mod communication; \ No newline at end of file diff --git a/src/model/core/schedule/communication.rs b/src/model/core/schedule/communication.rs new file mode 100644 index 0000000..edf5b11 --- /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::backup_schedule::BackupSchedule; + +pub enum ScheduleCommand { + AddSchedule(BackupSchedule), + ModifySchedule(BackupSchedule), + RemoveSchedule(Uuid), + ActivateSchedule(Uuid), + PauseSchedule(Uuid), + DisableSchedule(Uuid), +} + +impl Message for ScheduleCommand { + type Response = (); +} + +impl Command for ScheduleCommand {} + +pub enum ScheduleInternalCommand { + TimerNotify, + RefreshTimer, +} + +impl Message for ScheduleInternalCommand { + type Response = (); +} + +impl Command for ScheduleInternalCommand {} + +pub enum ScheduleQuery { + GetSchedules, +} + +impl Message for ScheduleQuery { + type Response = ScheduleQueryResponse; +} + +impl Query for ScheduleQuery {} + +pub enum ScheduleQueryResponse { + GetSchedules(Vec), +} diff --git a/src/model/core/schedule/mod.rs b/src/model/core/schedule/mod.rs index 7a0790a..ea316b6 100644 --- a/src/model/core/schedule/mod.rs +++ b/src/model/core/schedule/mod.rs @@ -1 +1,2 @@ pub mod backup_schedule; +pub mod communication; From f526a57f85b9d64173b33605226c1d5062252d6f Mon Sep 17 00:00:00 2001 From: DaLaw2 Date: Mon, 25 Aug 2025 11:44:34 +0800 Subject: [PATCH 6/8] wip: replace actor-based system with CommunicationManager --- src/core/backup/backup_engine.rs | 102 +++++++------- src/core/backup/backup_service.rs | 58 +++----- src/core/backup/progress_tracker.rs | 2 +- .../infrastructure/communication_manager.rs | 22 +-- src/core/infrastructure/io_manager.rs | 2 +- src/core/schedule/schedule_manager.rs | 91 ++++++------- src/core/schedule/schedule_service.rs | 54 +++----- src/core/schedule/schedule_timer.rs | 127 ++++++++++-------- src/core/system.rs | 52 ++++--- src/interface/{ => core}/file_system.rs | 2 +- src/interface/core/mod.rs | 4 +- .../core/{service.rs => runnable.rs} | 4 +- src/interface/core/unit.rs | 18 --- src/interface/mod.rs | 1 - src/interface/repository/schedule.rs | 22 +-- src/model/core/backup/communication.rs | 24 ++-- .../{backup_execution.rs => execution.rs} | 2 +- src/model/core/backup/mod.rs | 2 +- src/model/core/schedule/communication.rs | 40 +++--- src/model/core/schedule/mod.rs | 2 +- .../{backup_schedule.rs => schedule.rs} | 10 +- src/platform/linux/file_system.rs | 2 +- src/platform/windows/file_system.rs | 2 +- src/ui/common.rs | 8 +- src/ui/execution_page.rs | 6 +- src/ui/schedule_page.rs | 16 +-- 26 files changed, 321 insertions(+), 354 deletions(-) rename src/interface/{ => core}/file_system.rs (99%) rename src/interface/core/{service.rs => runnable.rs} (72%) delete mode 100644 src/interface/core/unit.rs rename src/model/core/backup/{backup_execution.rs => execution.rs} (97%) rename src/model/core/schedule/{backup_schedule.rs => schedule.rs} (87%) diff --git a/src/core/backup/backup_engine.rs b/src/core/backup/backup_engine.rs index b7248ed..99f1abe 100644 --- a/src/core/backup/backup_engine.rs +++ b/src/core/backup/backup_engine.rs @@ -1,13 +1,13 @@ use crate::core::backup::progress_tracker::ProgressTracker; use crate::core::gui::gui_message_handler::GuiMessageHandler; use crate::core::infrastructure::app_config::AppConfig; +use crate::core::infrastructure::communication_manager::CommunicationManager; use crate::core::infrastructure::io_manager::IOManager; -use crate::interface::core::unit::Unit; -use crate::interface::file_system::FileSystemTrait; -use crate::model::core::backup::backup_execution::*; -use crate::model::core::backup::communication::{ - BackupCommand, BackupInternalCommand, BackupQuery, BackupQueryResponse, -}; +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; @@ -28,8 +28,9 @@ use uuid::Uuid; pub struct BackupEngine { app_config: Arc, io_manager: Arc, + communication_manager: Arc, progress_tracker: Arc, - executions: Arc>, + executions: Arc>, running_executions: Arc, JoinHandle<()>)>>, } @@ -37,17 +38,29 @@ impl BackupEngine { pub fn new( app_config: Arc, io_manager: Arc, + communication_manager: Arc, progress_tracker: Arc, ) -> Self { Self { app_config, io_manager, + 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 @@ -67,14 +80,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); } @@ -172,9 +185,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<()>)>>, } @@ -182,27 +195,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; @@ -254,17 +262,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 = ExecutionErrorEvent { + 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); } } } @@ -318,7 +325,7 @@ impl Worker { async fn run( &self, - execution: BackupExecution, + execution: Execution, global_queue: Arc>, mut shutdown: oneshot::Receiver<()>, ) -> (Vec, Vec) { @@ -383,7 +390,7 @@ impl Worker { async fn process_entry( &self, - execution: &BackupExecution, + execution: &Execution, current_path: &Path, ) -> Result, Error> { let io_manager = &self.io_manager; @@ -415,7 +422,7 @@ impl Worker { async fn backup_directory( &self, - execution: &BackupExecution, + execution: &Execution, source_path: &Path, destination_path: &Path, ) -> Result, Error> { @@ -440,7 +447,7 @@ impl Worker { async fn backup_file( &self, - execution: &BackupExecution, + execution: &Execution, source_path: &Path, destination_path: &Path, ) -> Result, Error> { @@ -471,7 +478,7 @@ impl Worker { #[inline(always)] async fn process_symlink( &self, - execution: &BackupExecution, + execution: &Execution, source_path: &Path, destination_path: &Path, ) -> Result<(), Error> { @@ -486,7 +493,7 @@ impl Worker { async fn follow_symlink( &self, - execution: &BackupExecution, + execution: &Execution, source_path: &Path, destination_path: &Path, ) -> Result<(), Error> { @@ -543,7 +550,7 @@ impl Worker { async fn copy_symlink( &self, - execution: &BackupExecution, + execution: &Execution, source_path: &Path, destination_path: &Path, ) -> Result<(), Error> { @@ -650,16 +657,8 @@ impl Worker { } #[async_trait] -impl Unit for BackupEngine { - type Command = BackupCommand; - type InternalCommand = BackupInternalCommand; - type Query = BackupQuery; - - fn get_internal_channel(&self) -> UnboundedReceiver { - unimplemented!() - } - - async fn handle_command(&self, command: Self::Command) -> Result<(), Error> { +impl CommandHandler for BackupEngine { + async fn handle_command(&self, command: BackupCommand) -> Result<(), Error> { match command { BackupCommand::AddExecution(execution) => { self.add_execution(execution).await; @@ -679,12 +678,11 @@ impl Unit for BackupEngine { } Ok(()) } +} - async fn handle_internal_command(&self, command: Self::InternalCommand) -> Result<(), Error> { - match command {} - } - - async fn handle_query(&self, query: Self::Query) -> Result { +#[async_trait] +impl QueryHandler for BackupEngine { + async fn handle_query(&self, query: BackupQuery) -> Result { match query { BackupQuery::GetExecutions => { let executions = self.get_all_executions(); diff --git a/src/core/backup/backup_service.rs b/src/core/backup/backup_service.rs index 4b9f99b..0c87c25 100644 --- a/src/core/backup/backup_service.rs +++ b/src/core/backup/backup_service.rs @@ -1,58 +1,32 @@ use crate::core::backup::backup_engine::BackupEngine; use crate::core::backup::progress_tracker::ProgressTracker; use crate::core::infrastructure::app_config::AppConfig; +use crate::core::infrastructure::communication_manager::CommunicationManager; use crate::core::infrastructure::io_manager::IOManager; -use crate::interface::core::service::Service; -use async_trait::async_trait; use std::sync::Arc; -use tokio::io::AsyncWriteExt; -use tokio::select; -use tokio::sync::oneshot::Receiver; -use crate::interface::communication::command::CommandHandler; -use crate::interface::communication::query::QueryHandler; -use crate::interface::core::unit::Unit; -use crate::model::core::backup::communication::{BackupCommand, BackupQuery, BackupQueryResponse}; -use crate::model::error::Error; pub struct BackupService { backup_engine: Arc, } impl BackupService { - pub async fn init(app_config: Arc, io_manager: Arc) { + pub async fn new( + app_config: Arc, + io_manager: 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, progress_tracker)); - let backup_service = Self { backup_engine }; + let backup_engine = Arc::new(BackupEngine::new( + app_config, + io_manager, + communication_manager, + progress_tracker, + )); + Self { backup_engine } } -} - -#[async_trait] -impl Service for BackupService { - async fn process_internal_command(self: Arc, shutdown_rx: Receiver<()>) { - let mut backup_engine_rx = self.backup_engine.get_internal_channel(); - loop { - select! { - _ = &shutdown_rx => { - break - } - _ = backup_engine_rx.recv() => { - - } - } - } - } -} - -#[async_trait] -impl CommandHandler for BackupService { - async fn handle_command(&self, command: BackupCommand) -> Result<(), Error> { - self.backup_engine.handle_command(command) - } -} -#[async_trait] -impl QueryHandler for BackupService { - async fn handle_query(&self, query: BackupQuery) -> Result { - self.backup_engine.handle_query(query) + pub async fn register_services(&self) { + let backup_engine = self.backup_engine.clone(); + backup_engine.register_services().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/infrastructure/communication_manager.rs b/src/core/infrastructure/communication_manager.rs index b1ef52e..3203e4e 100644 --- a/src/core/infrastructure/communication_manager.rs +++ b/src/core/infrastructure/communication_manager.rs @@ -1,6 +1,7 @@ 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; @@ -9,7 +10,6 @@ use dashmap::DashMap; use std::any::{Any, TypeId}; use std::sync::Arc; use tokio::sync::broadcast; -use crate::interface::communication::event::EventBroadcaster; pub struct CommunicationManager { app_config: Arc, @@ -40,22 +40,28 @@ impl CommunicationManager { 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)); + self.event_broadcasters + .insert(type_id, Box::new(broadcaster)); } pub async fn publish_event(&self, event: E) -> Result<(), Error> { let type_id = TypeId::of::(); - let broadcaster = self.event_broadcasters.get(&type_id) + let broadcaster = self + .event_broadcasters + .get(&type_id) .ok_or(MiscError::TypeNotRegistered)?; broadcaster.broadcast_event(Box::new(event)) } pub fn subscribe_event(&self) -> Result, Error> { let type_id = TypeId::of::(); - let broadcaster = self.event_broadcasters.get(&type_id) + let broadcaster = self + .event_broadcasters + .get(&type_id) .ok_or(MiscError::TypeNotRegistered)?; let receiver_box = broadcaster.subscribe_typed(); - let receiver = *receiver_box.downcast::>() + let receiver = *receiver_box + .downcast::>() .map_err(|_| MiscError::TypeMismatch)?; Ok(receiver) } @@ -151,7 +157,7 @@ impl ServiceRegistrar { Self { service, comm } } - pub fn command(&self) -> &Self + pub fn command(self) -> Self where S: CommandHandler, { @@ -160,7 +166,7 @@ impl ServiceRegistrar { self } - pub fn query(&self) -> &Self + pub fn query(self) -> Self where S: QueryHandler, { @@ -169,7 +175,7 @@ impl ServiceRegistrar { self } - pub fn event(&self) -> &Self { + pub fn event(self) -> Self { self.comm.register_event_type::(); self } 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/schedule/schedule_manager.rs b/src/core/schedule/schedule_manager.rs index 2b205d9..8ef60eb 100644 --- a/src/core/schedule/schedule_manager.rs +++ b/src/core/schedule/schedule_manager.rs @@ -1,26 +1,28 @@ -use crate::core::backup::backup_service::BackupService; +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::schedule::backup_schedule::*; +use crate::model::core::backup::communication::BackupCommand; +use crate::model::core::schedule::schedule::*; +use crate::model::core::schedule::communication::*; use crate::model::error::Error; +use async_trait::async_trait; use chrono::{Duration, Months, Utc}; use dashmap::DashMap; use std::sync::Arc; -use async_trait::async_trait; -use tokio::sync::mpsc::UnboundedReceiver; use uuid::Uuid; -use crate::interface::communication::message::Message; -use crate::interface::core::unit::Unit; -use crate::model::core::schedule::communication::{ScheduleCommand, ScheduleInternalCommand, ScheduleQuery, ScheduleQueryResponse}; pub struct ScheduleManager { database_manager: Arc, - schedules: DashMap, + communication_manager: Arc, + schedules: DashMap, } impl ScheduleManager { pub async fn new( database_manager: Arc, + communication_manager: Arc, ) -> Result { let schedules = DashMap::new(); let database_schedules = database_manager.get_all_backup_schedules().await?; @@ -29,17 +31,26 @@ 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?; @@ -47,7 +58,7 @@ impl ScheduleManager { 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?; @@ -109,13 +120,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?; } @@ -124,7 +130,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; } @@ -146,54 +152,45 @@ impl ScheduleManager { } #[async_trait] -impl Unit for ScheduleManager { - type Command = ScheduleCommand; - type InternalCommand = ScheduleInternalCommand; - type Query = ScheduleQuery; - - fn get_internal_channel(&self) -> UnboundedReceiver { - todo!() - } - - async fn handle_command(&self, command: Self::Command) -> Result<(), Error> { +impl CommandHandler for ScheduleManager { + async fn handle_command(&self, command: ScheduleManagerCommand) -> Result<(), Error> { match command { - ScheduleCommand::AddSchedule(schedule) => { + ScheduleManagerCommand::AddSchedule(schedule) => { self.create_schedule(schedule).await?; } - ScheduleCommand::ModifySchedule(schedule) => { + ScheduleManagerCommand::ModifySchedule(schedule) => { self.modify_schedule(schedule).await?; } - ScheduleCommand::RemoveSchedule(uuid) => { + ScheduleManagerCommand::RemoveSchedule(uuid) => { self.remove_schedule(uuid).await?; } - ScheduleCommand::ActivateSchedule(uuid) => { + ScheduleManagerCommand::ActivateSchedule(uuid) => { self.active_schedule(uuid).await?; } - ScheduleCommand::PauseSchedule(uuid) => { + ScheduleManagerCommand::PauseSchedule(uuid) => { self.pause_schedule(uuid).await?; } - ScheduleCommand::DisableSchedule(uuid) => { + ScheduleManagerCommand::DisableSchedule(uuid) => { self.disable_schedule(uuid).await?; } - } - Ok(()) - } - - async fn handle_internal_command(&self, command: Self::InternalCommand) -> Result<(), Error> { - match command { - ScheduleInternalCommand::TimerNotify => { + ScheduleManagerCommand::ExecuteReadySchedules => { self.execute_ready_schedule().await?; } - _ => {} } Ok(()) } +} - async fn handle_query(&self, query: Self::Query) -> Result<::Response, Error> { +#[async_trait] +impl QueryHandler for ScheduleManager { + async fn handle_query( + &self, + query: ScheduleManagerQuery, + ) -> Result { match query { - ScheduleQuery::GetSchedules => { + ScheduleManagerQuery::GetSchedules => { let executions = self.get_all_schedules().await; - Ok(ScheduleQueryResponse::GetSchedules(executions)) + Ok(ScheduleManagerQueryResponse::GetSchedules(executions)) } } } diff --git a/src/core/schedule/schedule_service.rs b/src/core/schedule/schedule_service.rs index 0936791..7dd28d1 100644 --- a/src/core/schedule/schedule_service.rs +++ b/src/core/schedule/schedule_service.rs @@ -3,63 +3,47 @@ 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::core::service::Service; +use crate::interface::core::runnable::Runnable; use crate::model::error::Error; use async_trait::async_trait; -use std::any::Any; -use std::sync::{Arc, OnceLock}; +use std::sync::Arc; use tokio::sync::oneshot::Receiver; -use tokio::sync::{mpsc, oneshot}; -use crate::interface::communication::command::CommandHandler; -use crate::interface::communication::query::QueryHandler; -use crate::model::core::schedule::communication::{ScheduleCommand, ScheduleQuery, ScheduleQueryResponse}; 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, - ) -> Result<(), Error> { + communication_manager: Arc, + ) -> Result { let schedule_manager = - Arc::new(ScheduleManager::new(database_manager).await?); - let schedule_timer = Arc::new(ScheduleTimer::new(app_config)); + 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(), }; - 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 Service for ScheduleService { - async fn process_internal_command(self: Arc, shutdown_rx: Receiver<()>) { - todo!() - } -} - -impl CommandHandler for ScheduleService { - async fn handle_command(&self, command: ScheduleCommand) -> Result<(), Error> { - todo!() - } -} - -impl QueryHandler for ScheduleService { - async fn handle_query(&self, query: ScheduleQuery) -> Result { - todo!() +impl Runnable for ScheduleService { + async fn run_impl(self: Arc, shutdown_rx: Receiver<()>) { + let schedule_timer = self.schedule_timer.clone(); + 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 7145d5a..06f028a 100644 --- a/src/core/schedule/schedule_timer.rs +++ b/src/core/schedule/schedule_timer.rs @@ -1,83 +1,53 @@ use crate::core::infrastructure::app_config::AppConfig; -use crate::core::schedule::schedule_service::ScheduleService; -use crate::model::core::schedule::backup_schedule::ScheduleState; +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::oneshot::Receiver; +use tokio::sync::Notify; use tokio::sync::{mpsc, 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; @@ -102,3 +72,48 @@ impl ScheduleTimer { } } } + +#[async_trait] +impl Runnable for ScheduleTimer { + async fn run_impl(self: Arc, mut shutdown_rx: 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 c1ece2f..5e579e0 100644 --- a/src/core/system.rs +++ b/src/core/system.rs @@ -1,5 +1,6 @@ use crate::core::backup::backup_service::BackupService; 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; @@ -8,18 +9,20 @@ 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; +use crate::interface::core::runnable::Runnable; pub struct System { - app_config: Arc, - io_manager: Arc, - database_manager: Arc, - actor_system: Arc, + backup_service: Arc, + schedule_service: Arc, + shutdowns: SegQueue>, } impl System { @@ -27,12 +30,27 @@ 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 system = Self { - app_config, - io_manager, - database_manager, - actor_system, + backup_service, + schedule_service, + shutdowns: SegQueue::new(), }; Ok(system) } @@ -41,18 +59,10 @@ 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())); + self.backup_service.register_services().await; + self.schedule_service.register_services().await; + let schedule_service_shutdown = self.schedule_service.run().await; + self.shutdowns.push(schedule_service_shutdown); log!(SystemLog::InitializeComplete); gui_manager.start().await } 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 index eec556b..a758f05 100644 --- a/src/interface/core/mod.rs +++ b/src/interface/core/mod.rs @@ -1,2 +1,2 @@ -pub mod service; -pub mod unit; +pub mod file_system; +pub mod runnable; diff --git a/src/interface/core/service.rs b/src/interface/core/runnable.rs similarity index 72% rename from src/interface/core/service.rs rename to src/interface/core/runnable.rs index 66af958..99a956e 100644 --- a/src/interface/core/service.rs +++ b/src/interface/core/runnable.rs @@ -3,7 +3,7 @@ use std::sync::Arc; use tokio::sync::oneshot; #[async_trait] -pub trait Service { +pub trait Runnable { async fn run(self: Arc) -> oneshot::Sender<()> { let (shutdown_tx, shutdown_rx) = oneshot::channel(); @@ -12,5 +12,5 @@ pub trait Service { shutdown_tx } - async fn process_internal_command(self: Arc, shutdown_rx: oneshot::Receiver<()>); + async fn run_impl(self: Arc, shutdown_rx: oneshot::Receiver<()>); } diff --git a/src/interface/core/unit.rs b/src/interface/core/unit.rs deleted file mode 100644 index d2b6425..0000000 --- a/src/interface/core/unit.rs +++ /dev/null @@ -1,18 +0,0 @@ -use crate::interface::communication::command::Command; -use crate::interface::communication::message::Message; -use crate::interface::communication::query::Query; -use crate::model::error::Error; -use async_trait::async_trait; -use tokio::sync::mpsc::UnboundedReceiver; - -#[async_trait] -pub trait Unit { - type Command: Command; - type InternalCommand: Command; - type Query: Query; - - fn get_internal_channel(&self) -> UnboundedReceiver; - async fn handle_command(&self, command: Self::Command) -> Result<(), Error>; - async fn handle_internal_command(&self, command: Self::InternalCommand) -> Result<(), Error>; - async fn handle_query(&self, query: Self::Query) -> Result<::Response, Error>; -} diff --git a/src/interface/mod.rs b/src/interface/mod.rs index 128d71a..7a7d75d 100644 --- a/src/interface/mod.rs +++ b/src/interface/mod.rs @@ -1,4 +1,3 @@ pub mod communication; pub mod core; -pub mod file_system; 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/model/core/backup/communication.rs b/src/model/core/backup/communication.rs index 31b91c8..6752ec7 100644 --- a/src/model/core/backup/communication.rs +++ b/src/model/core/backup/communication.rs @@ -1,11 +1,13 @@ use crate::interface::communication::command::Command; +use crate::interface::communication::event::Event; use crate::interface::communication::message::Message; use crate::interface::communication::query::Query; -use crate::model::core::backup::backup_execution::BackupExecution; +use crate::model::core::backup::execution::Execution; +use crate::model::error::Error; use uuid::Uuid; pub enum BackupCommand { - AddExecution(BackupExecution), + AddExecution(Execution), RemoveExecution(Uuid), StartExecution(Uuid), SuspendExecution(Uuid), @@ -18,14 +20,6 @@ impl Message for BackupCommand { impl Command for BackupCommand {} -pub enum BackupInternalCommand {} - -impl Message for BackupInternalCommand { - type Response = (); -} - -impl Command for BackupInternalCommand {} - pub enum BackupQuery { GetExecutions, } @@ -37,5 +31,13 @@ impl Message for BackupQuery { impl Query for BackupQuery {} pub enum BackupQueryResponse { - GetExecutions(Vec<(Uuid, BackupExecution)>), + GetExecutions(Vec<(Uuid, Execution)>), +} + +#[derive(Clone)] +pub struct ExecutionErrorEvent { + pub uuid: Uuid, + pub errors: Vec, } + +impl Event for ExecutionErrorEvent {} 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/mod.rs b/src/model/core/backup/mod.rs index f79d421..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 execution; pub mod progress_data; pub mod communication; diff --git a/src/model/core/schedule/communication.rs b/src/model/core/schedule/communication.rs index edf5b11..6b4aef0 100644 --- a/src/model/core/schedule/communication.rs +++ b/src/model/core/schedule/communication.rs @@ -2,44 +2,44 @@ 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::backup_schedule::BackupSchedule; +use crate::model::core::schedule::schedule::Schedule; -pub enum ScheduleCommand { - AddSchedule(BackupSchedule), - ModifySchedule(BackupSchedule), +pub enum ScheduleManagerCommand { + AddSchedule(Schedule), + ModifySchedule(Schedule), RemoveSchedule(Uuid), ActivateSchedule(Uuid), PauseSchedule(Uuid), DisableSchedule(Uuid), + ExecuteReadySchedules, } -impl Message for ScheduleCommand { +impl Message for ScheduleManagerCommand { type Response = (); } -impl Command for ScheduleCommand {} +impl Command for ScheduleManagerCommand {} -pub enum ScheduleInternalCommand { - TimerNotify, - RefreshTimer, +pub enum ScheduleManagerQuery { + GetSchedules, } -impl Message for ScheduleInternalCommand { - type Response = (); +impl Message for ScheduleManagerQuery { + type Response = ScheduleManagerQueryResponse; } -impl Command for ScheduleInternalCommand {} +impl Query for ScheduleManagerQuery {} -pub enum ScheduleQuery { - GetSchedules, +pub enum ScheduleManagerQueryResponse { + GetSchedules(Vec), } -impl Message for ScheduleQuery { - type Response = ScheduleQueryResponse; +pub enum ScheduleTimerCommand { + RefreshTimer } -impl Query for ScheduleQuery {} - -pub enum ScheduleQueryResponse { - GetSchedules(Vec), +impl Message for ScheduleTimerCommand { + type Response = (); } + +impl Command for ScheduleTimerCommand {} diff --git a/src/model/core/schedule/mod.rs b/src/model/core/schedule/mod.rs index ea316b6..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 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/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 7432892..15eccfd 100644 --- a/src/ui/execution_page.rs +++ b/src/ui/execution_page.rs @@ -1,6 +1,6 @@ use crate::core::backup::backup_service::BackupService; use crate::core::infrastructure::app_config::AppConfig; -use crate::model::core::backup::backup_execution::*; +use crate::model::core::backup::execution::*; use crate::model::error::Error; use crate::ui::common::{ComparisonModeSelection, ExecutionDisplay, FolderSelectionMode}; use dashmap::DashMap; @@ -137,7 +137,7 @@ 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 @@ -537,7 +537,7 @@ impl ExecutionPage { } }; - let execution = BackupExecution { + let execution = Execution { uuid: Uuid::new_v4(), state: BackupState::Pending, source_path: PathBuf::from(&self.new_task_source), diff --git a/src/ui/schedule_page.rs b/src/ui/schedule_page.rs index 3c89b04..6f56475 100644 --- a/src/ui/schedule_page.rs +++ b/src/ui/schedule_page.rs @@ -1,7 +1,7 @@ use crate::core::infrastructure::app_config::AppConfig; use crate::core::schedule::schedule_service::ScheduleService; -use crate::model::core::backup::backup_execution::*; -use crate::model::core::schedule::backup_schedule::*; +use crate::model::core::backup::execution::*; +use crate::model::core::schedule::schedule::*; use crate::ui::common::{ComparisonModeSelection, FolderSelectionMode}; use eframe::egui; use egui_file_dialog::FileDialog; @@ -18,7 +18,7 @@ pub struct SchedulePage { schedule_service_ref: ActorRef, - schedules: Vec, + schedules: Vec, new_schedule_name: String, new_schedule_source: String, @@ -88,7 +88,7 @@ impl SchedulePage { } } - 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 @@ -101,7 +101,7 @@ impl SchedulePage { }) } - 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 @@ -221,7 +221,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| { @@ -249,7 +249,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) @@ -506,7 +506,7 @@ impl SchedulePage { } }; - let schedule = BackupSchedule { + let schedule = Schedule { uuid: Uuid::new_v4(), name: self.new_schedule_name.clone(), state: ScheduleState::Active, From cebcb6a2aacc9e9b7b4550076872b37510392db5 Mon Sep 17 00:00:00 2001 From: DaLaw2 Date: Mon, 25 Aug 2025 12:49:54 +0800 Subject: [PATCH 7/8] wip: Remove actor-based GUI components and integrate them into the CommunicationManager structure --- src/core/backup/backup_engine.rs | 5 +- src/core/gui/gui_manager.rs | 23 ++- src/core/gui/gui_message_handler.rs | 48 ------ src/core/gui/mod.rs | 1 - .../infrastructure/communication_manager.rs | 62 +++---- src/core/system.rs | 19 ++- src/interface/core/runnable.rs | 2 +- src/model/core/gui/communication.rs | 14 +- src/model/error/misc.rs | 4 + src/ui/execution_page.rs | 155 +++++++----------- src/ui/schedule_page.rs | 92 ++++------- 11 files changed, 158 insertions(+), 267 deletions(-) delete mode 100644 src/core/gui/gui_message_handler.rs diff --git a/src/core/backup/backup_engine.rs b/src/core/backup/backup_engine.rs index 99f1abe..1cf9434 100644 --- a/src/core/backup/backup_engine.rs +++ b/src/core/backup/backup_engine.rs @@ -1,5 +1,4 @@ use crate::core::backup::progress_tracker::ProgressTracker; -use crate::core::gui::gui_message_handler::GuiMessageHandler; use crate::core::infrastructure::app_config::AppConfig; use crate::core::infrastructure::communication_manager::CommunicationManager; use crate::core::infrastructure::io_manager::IOManager; @@ -167,14 +166,14 @@ impl BackupEngine { 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, 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 6fc230f..166deaf 100644 --- a/src/core/gui/mod.rs +++ b/src/core/gui/mod.rs @@ -1,2 +1 @@ pub mod gui_manager; -pub mod gui_message_handler; diff --git a/src/core/infrastructure/communication_manager.rs b/src/core/infrastructure/communication_manager.rs index 3203e4e..bc8afb5 100644 --- a/src/core/infrastructure/communication_manager.rs +++ b/src/core/infrastructure/communication_manager.rs @@ -35,37 +35,6 @@ impl CommunicationManager { ServiceRegistrar::new(service, self) } - 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 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 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 fn register_command_handler( &self, handler: Arc + Send + Sync>, @@ -122,6 +91,37 @@ impl CommunicationManager { } } + 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 fn has_command_handler(&self) -> bool { self.command_handlers.contains_key(&TypeId::of::()) } diff --git a/src/core/system.rs b/src/core/system.rs index 5e579e0..004f579 100644 --- a/src/core/system.rs +++ b/src/core/system.rs @@ -1,9 +1,11 @@ use crate::core::backup::backup_service::BackupService; +use crate::core::gui::gui_manager::GuiManager; 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)))] @@ -17,11 +19,11 @@ use privilege::user::privileged; use std::process; use std::sync::Arc; use tokio::sync::oneshot; -use crate::interface::core::runnable::Runnable; pub struct System { backup_service: Arc, schedule_service: Arc, + gui_manager: Arc, shutdowns: SegQueue>, } @@ -47,9 +49,11 @@ impl System { ) .await?, ); + let gui_manager = Arc::new(GuiManager::new(app_config, communication_manager)); let system = Self { backup_service, schedule_service, + gui_manager, shutdowns: SegQueue::new(), }; Ok(system) @@ -59,17 +63,18 @@ impl System { Logging::initialize().await; log!(SystemLog::Initializing); Self::elevate_privileges()?; - self.backup_service.register_services().await; - self.schedule_service.register_services().await; - let schedule_service_shutdown = self.schedule_service.run().await; + 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 fn shutdown(&self) {} fn elevate_privileges() -> Result<(), Error> { #[cfg(not(debug_assertions))] diff --git a/src/interface/core/runnable.rs b/src/interface/core/runnable.rs index 99a956e..7b2bc27 100644 --- a/src/interface/core/runnable.rs +++ b/src/interface/core/runnable.rs @@ -3,7 +3,7 @@ use std::sync::Arc; use tokio::sync::oneshot; #[async_trait] -pub trait Runnable { +pub trait Runnable: 'static { async fn run(self: Arc) -> oneshot::Sender<()> { let (shutdown_tx, shutdown_rx) = oneshot::channel(); diff --git a/src/model/core/gui/communication.rs b/src/model/core/gui/communication.rs index e001c5b..fb7c28a 100644 --- a/src/model/core/gui/communication.rs +++ b/src/model/core/gui/communication.rs @@ -5,25 +5,25 @@ use uuid::Uuid; #[derive(Clone)] pub struct FolderProcess { - uuid: Uuid, - folder: PathBuf, + pub uuid: Uuid, + pub folder: PathBuf, } impl Event for FolderProcess {} #[derive(Clone)] pub struct ExecutionProgress { - uuid: Uuid, - processed_files: usize, - error_count: usize, + pub uuid: Uuid, + pub processed_files: usize, + pub error_count: usize, } impl Event for ExecutionProgress {} #[derive(Clone)] pub struct ExecutionErrors { - uuid: Uuid, - errors: Vec, + pub uuid: Uuid, + pub errors: Vec, } impl Event for ExecutionErrors {} diff --git a/src/model/error/misc.rs b/src/model/error/misc.rs index 8e38645..9666a09 100644 --- a/src/model/error/misc.rs +++ b/src/model/error/misc.rs @@ -14,6 +14,10 @@ traceable! { #[error("Failed to initialize UI platform")] UIPlatformError => tracing::Level::ERROR, + #[no_source] + #[error("Assert file not found")] + AssertFileNotFound => tracing::Level::ERROR, + #[no_source] #[error("Handler not found")] HandlerNotFound => tracing::Level::ERROR, diff --git a/src/ui/execution_page.rs b/src/ui/execution_page.rs index 15eccfd..2cc7b62 100644 --- a/src/ui/execution_page.rs +++ b/src/ui/execution_page.rs @@ -1,6 +1,9 @@ use crate::core::backup::backup_service::BackupService; use crate::core::infrastructure::app_config::AppConfig; +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; @@ -11,14 +14,17 @@ use std::collections::HashSet; use std::path::PathBuf; use std::sync::{mpsc, 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>, @@ -44,13 +50,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(), @@ -67,34 +77,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); } } } @@ -102,20 +108,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 { @@ -139,80 +136,45 @@ impl ExecutionPage { 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(()) }) } @@ -553,7 +515,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 6f56475..ec100e8 100644 --- a/src/ui/schedule_page.rs +++ b/src/ui/schedule_page.rs @@ -1,7 +1,11 @@ use crate::core::infrastructure::app_config::AppConfig; -use crate::core::schedule::schedule_service::ScheduleService; +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; use egui_file_dialog::FileDialog; @@ -14,9 +18,7 @@ use uuid::Uuid; pub struct SchedulePage { app_config: Arc, - actor_system: Arc, - - schedule_service_ref: ActorRef, + communication_manager: Arc, schedules: Vec, @@ -40,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(), @@ -70,18 +71,13 @@ 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); } @@ -90,78 +86,54 @@ impl SchedulePage { 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: 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(()) }) } From 7ee6e4e10481b63024a39e02ff8af414b239ee53 Mon Sep 17 00:00:00 2001 From: DaLaw2 Date: Mon, 25 Aug 2025 13:15:17 +0800 Subject: [PATCH 8/8] feat: Complete replace actor-based architecture with CommunicationManager --- src/core/backup/backup_engine.rs | 8 +++---- src/core/backup/backup_service.rs | 4 ++++ .../infrastructure/communication_manager.rs | 24 ------------------- src/core/schedule/schedule_manager.rs | 20 +++++++++++++++- src/core/schedule/schedule_timer.rs | 5 ++-- src/core/system.rs | 7 +++++- src/main.rs | 2 +- src/model/core/backup/communication.rs | 10 -------- src/ui/execution_page.rs | 3 +-- 9 files changed, 37 insertions(+), 46 deletions(-) diff --git a/src/core/backup/backup_engine.rs b/src/core/backup/backup_engine.rs index 1cf9434..2492771 100644 --- a/src/core/backup/backup_engine.rs +++ b/src/core/backup/backup_engine.rs @@ -18,11 +18,11 @@ use macros::log; use std::collections::{HashSet, VecDeque}; use std::path::{Path, PathBuf}; use std::sync::Arc; -use tokio::sync::mpsc::UnboundedReceiver; 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, @@ -56,7 +56,7 @@ impl BackupEngine { .with_service(self) .command::() .query::() - .event::() + .event::() .build(); } @@ -261,13 +261,13 @@ impl ExecutionRunner { next_level.extend(worker_next_level); if !worker_errors.is_empty() { errors.extend(worker_errors.clone()); - let event = ExecutionErrorEvent { + let event = ExecutionErrors { uuid: execution.uuid, errors: worker_errors, }; if let Err(err) = self .communication_manager - .publish_event::(event) + .publish_event::(event) .await { error!("{}", err); diff --git a/src/core/backup/backup_service.rs b/src/core/backup/backup_service.rs index 0c87c25..1ed12b3 100644 --- a/src/core/backup/backup_service.rs +++ b/src/core/backup/backup_service.rs @@ -29,4 +29,8 @@ impl BackupService { let backup_engine = self.backup_engine.clone(); backup_engine.register_services().await; } + + pub async fn shutdown(&self) { + self.backup_engine.stop_all_executions().await; + } } diff --git a/src/core/infrastructure/communication_manager.rs b/src/core/infrastructure/communication_manager.rs index bc8afb5..8842751 100644 --- a/src/core/infrastructure/communication_manager.rs +++ b/src/core/infrastructure/communication_manager.rs @@ -121,30 +121,6 @@ impl CommunicationManager { .ok_or(MiscError::TypeNotRegistered)?; broadcaster.broadcast_event(Box::new(event)) } - - pub fn has_command_handler(&self) -> bool { - self.command_handlers.contains_key(&TypeId::of::()) - } - - pub fn has_query_handler(&self) -> bool { - self.query_handlers.contains_key(&TypeId::of::()) - } - - pub fn has_event_type(&self) -> bool { - self.event_broadcasters.contains_key(&TypeId::of::()) - } - - pub fn clear_command_handlers(&self) { - self.command_handlers.clear(); - } - - pub fn clear_query_handlers(&self) { - self.query_handlers.clear(); - } - - pub fn clear_event_types(&self) { - self.event_broadcasters.clear(); - } } pub struct ServiceRegistrar { diff --git a/src/core/schedule/schedule_manager.rs b/src/core/schedule/schedule_manager.rs index 8ef60eb..7c9200f 100644 --- a/src/core/schedule/schedule_manager.rs +++ b/src/core/schedule/schedule_manager.rs @@ -4,8 +4,8 @@ use crate::interface::communication::command::CommandHandler; use crate::interface::communication::query::QueryHandler; use crate::interface::repository::schedule::ScheduleRepository; use crate::model::core::backup::communication::BackupCommand; -use crate::model::core::schedule::schedule::*; 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}; @@ -55,6 +55,9 @@ impl ScheduleManager { .create_backup_schedule(&schedule) .await?; self.schedules.insert(schedule.uuid, schedule); + self.communication_manager + .send_command(ScheduleTimerCommand::RefreshTimer) + .await?; Ok(()) } @@ -63,12 +66,18 @@ impl ScheduleManager { .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(()) } @@ -79,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(()) } @@ -90,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(()) } @@ -101,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(()) } diff --git a/src/core/schedule/schedule_timer.rs b/src/core/schedule/schedule_timer.rs index 06f028a..1844799 100644 --- a/src/core/schedule/schedule_timer.rs +++ b/src/core/schedule/schedule_timer.rs @@ -9,9 +9,8 @@ use async_trait::async_trait; use chrono::{Duration, Utc}; use std::sync::Arc; use tokio::select; -use tokio::sync::oneshot::Receiver; use tokio::sync::Notify; -use tokio::sync::{mpsc, oneshot}; +use tokio::sync::oneshot; use tokio::time::sleep; use tracing::error; @@ -75,7 +74,7 @@ impl ScheduleTimer { #[async_trait] impl Runnable for ScheduleTimer { - async fn run_impl(self: Arc, mut shutdown_rx: Receiver<()>) { + async fn run_impl(self: Arc, mut shutdown_rx: oneshot::Receiver<()>) { let communication_manager = self.communication_manager.clone(); loop { diff --git a/src/core/system.rs b/src/core/system.rs index 004f579..c84f086 100644 --- a/src/core/system.rs +++ b/src/core/system.rs @@ -74,7 +74,12 @@ impl System { gui_manager.start().await } - pub fn shutdown(&self) {} + 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> { #[cfg(not(debug_assertions))] 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/core/backup/communication.rs b/src/model/core/backup/communication.rs index 6752ec7..f743df4 100644 --- a/src/model/core/backup/communication.rs +++ b/src/model/core/backup/communication.rs @@ -1,9 +1,7 @@ use crate::interface::communication::command::Command; -use crate::interface::communication::event::Event; use crate::interface::communication::message::Message; use crate::interface::communication::query::Query; use crate::model::core::backup::execution::Execution; -use crate::model::error::Error; use uuid::Uuid; pub enum BackupCommand { @@ -33,11 +31,3 @@ impl Query for BackupQuery {} pub enum BackupQueryResponse { GetExecutions(Vec<(Uuid, Execution)>), } - -#[derive(Clone)] -pub struct ExecutionErrorEvent { - pub uuid: Uuid, - pub errors: Vec, -} - -impl Event for ExecutionErrorEvent {} diff --git a/src/ui/execution_page.rs b/src/ui/execution_page.rs index 2cc7b62..e2e3aa8 100644 --- a/src/ui/execution_page.rs +++ b/src/ui/execution_page.rs @@ -1,4 +1,3 @@ -use crate::core::backup::backup_service::BackupService; use crate::core::infrastructure::app_config::AppConfig; use crate::core::infrastructure::communication_manager::CommunicationManager; use crate::model::core::backup::communication::{BackupCommand, BackupQuery, BackupQueryResponse}; @@ -12,7 +11,7 @@ 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;