mocra_core/engine/task/
task_manager.rs1#![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 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 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 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 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 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 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 pub async fn module_names(&self) -> Vec<String> {
100 let assembler = self.module_assembler.read().await;
101 assembler.module_names()
102 }
103 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 pub async fn load_with_model(&self, task_model: &TaskEvent) -> Result<Task> {
110 self.factory.load_with_model(task_model).await
111 }
112
113 pub async fn load_parser(&self, parser_model: &TaskParserEvent) -> Result<Task> {
115 self.factory.load_parser_model(parser_model).await
116 }
117
118 pub async fn load_error(&self, error_model: &TaskErrorEvent) -> Result<Task> {
120 self.factory.load_error_model(error_model).await
121 }
122
123 pub async fn load_with_response(&self, response: &Response) -> Result<Task> {
125 self.factory.load_with_response(response).await
126 }
127
128 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 pub async fn clear_factory_cache(&self) {
138 self.factory.clear_cache().await;
139 }
140
141 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}