Skip to main content

mocra_core/engine/task/
task_manager.rs

1#![allow(unused)]
2#[cfg(feature = "store")]
3use super::repository::{MetadataStore, TaskRepository};
4use super::{assembler::ModuleAssembler, factory::TaskFactory, module::Module, task::Task};
5use crate::cacheable::{CacheAble, CacheService};
6use crate::common::interface::ModuleTrait;
7use crate::common::model::config::Config;
8use crate::common::model::login_info::LoginInfo;
9use crate::common::model::message::{TaskErrorEvent, TaskEvent, TaskParserEvent};
10use crate::common::model::{Request, Response};
11use crate::common::state::DbHandle;
12use crate::errors::Result;
13use dashmap::DashMap;
14use std::sync::Arc;
15use std::sync::atomic::{AtomicU64, Ordering};
16use std::time::{SystemTime, UNIX_EPOCH};
17use tokio::sync::RwLock;
18
19pub struct TaskManager {
20    factory: TaskFactory,
21    pub cache_service: Arc<CacheService>,
22    module_assembler: Arc<RwLock<ModuleAssembler>>,
23}
24
25impl TaskManager {
26    /// Creates a task manager with repository, factory, and module assembler wiring.
27    /// Narrowed dependencies: only a DB handle + cache / cookies / config are needed (it used to
28    /// take the entire `Arc<State>` — refactor Phase 2).
29    /// With `db = None` (or without the `store` feature) it enters no-DB mode, synthesizing tasks
30    /// from the in-memory module registry.
31    pub fn new(
32        db: &DbHandle,
33        cache_service: Arc<CacheService>,
34        cookie_service: Option<Arc<CacheService>>,
35        config: Arc<RwLock<Config>>,
36    ) -> Self {
37        #[cfg(feature = "store")]
38        let repository = db
39            .as_ref()
40            .map(|db| Arc::new(TaskRepository::new((**db).clone())) as Arc<dyn MetadataStore>);
41        #[cfg(not(feature = "store"))]
42        let repository = ();
43
44        let module_assembler = Arc::new(RwLock::new(ModuleAssembler::new()));
45        let factory = TaskFactory::new(
46            repository,
47            Arc::clone(&cache_service),
48            cookie_service,
49            Arc::clone(&module_assembler),
50            Arc::clone(&config),
51        );
52
53        Self {
54            factory,
55            cache_service,
56            module_assembler,
57        }
58    }
59
60    /// Registers one module implementation.
61    pub async fn add_module(&self, work: Arc<dyn ModuleTrait>) {
62        let name = work.name();
63        {
64            let mut assembler = self.module_assembler.write().await;
65            assembler.register_module(work.clone());
66        }
67    }
68
69    /// Registers multiple module implementations.
70    pub async fn add_modules(&self, works: Vec<Arc<dyn ModuleTrait>>) {
71        {
72            let mut assembler = self.module_assembler.write().await;
73            for work in &works {
74                assembler.register_module(work.clone());
75            }
76        }
77        for work in works {
78            let name = work.name();
79        }
80    }
81
82    /// Returns true when module name is registered.
83    pub async fn exists_module(&self, name: &str) -> bool {
84        let assembler = self.module_assembler.read().await;
85        assembler.get_module(name).is_some()
86    }
87    /// Unregisters module by name.
88    pub async fn remove_work(&self, name: &str) {
89        let mut assembler = self.module_assembler.write().await;
90        assembler.remove_module(name);
91        drop(assembler);
92    }
93    /// Removes all modules loaded from a given origin path.
94    pub async fn remove_by_origin(&self, origin: &std::path::Path) {
95        let mut assembler = self.module_assembler.write().await;
96        assembler.remove_by_origin(origin);
97    }
98    /// Returns all registered module names.
99    pub async fn module_names(&self) -> Vec<String> {
100        let assembler = self.module_assembler.read().await;
101        assembler.module_names()
102    }
103    /// Tags module names with origin path for hot-reload bookkeeping.
104    pub async fn set_origin(&self, names: &[String], origin: &std::path::Path) {
105        let mut assembler = self.module_assembler.write().await;
106        assembler.set_origin(names, origin);
107    }
108    /// Loads task from TaskModel and synchronizes initial state.
109    pub async fn load_with_model(&self, task_model: &TaskEvent) -> Result<Task> {
110        self.factory.load_with_model(task_model).await
111    }
112
113    /// Loads task from ParserTaskModel with historical state restoration.
114    pub async fn load_parser(&self, parser_model: &TaskParserEvent) -> Result<Task> {
115        self.factory.load_parser_model(parser_model).await
116    }
117
118    /// Loads task from ErrorTaskModel and applies error accounting.
119    pub async fn load_error(&self, error_model: &TaskErrorEvent) -> Result<Task> {
120        self.factory.load_error_model(error_model).await
121    }
122
123    /// Loads task context from response metadata.
124    pub async fn load_with_response(&self, response: &Response) -> Result<Task> {
125        self.factory.load_with_response(response).await
126    }
127
128    /// Loads resolved module and optional login info for parser flow.
129    pub async fn load_module_with_response(
130        &self,
131        response: &Response,
132    ) -> Result<(Arc<Module>, Option<LoginInfo>)> {
133        self.factory.load_module_with_response(response).await
134    }
135
136    /// Clears internal factory caches.
137    pub async fn clear_factory_cache(&self) {
138        self.factory.clear_cache().await;
139    }
140
141    /// Returns all registered module implementations.
142    pub async fn get_all_modules(&self) -> Vec<Arc<dyn ModuleTrait>> {
143        let assembler = self.module_assembler.read().await;
144        assembler.get_all_modules()
145    }
146}