1use crate::error::{IoError, Result};
6use crate::metadata::{Metadata, MetadataValue};
7use chrono::{DateTime, Datelike, Duration, Utc};
8use serde::{Deserialize, Serialize};
9use std::collections::{HashMap, HashSet};
10use std::sync::{Arc, Mutex};
11
12use super::types::{ResourceRequirements, Task, TaskStatus, TaskType, Workflow, WorkflowExecutor};
13
14pub mod events {
16 use super::*;
17 use crossbeam_channel::{Receiver, Sender};
18 #[derive(Debug, Clone, Serialize, Deserialize)]
20 pub enum WorkflowEvent {
21 FileCreated {
22 path: String,
23 },
24 FileModified {
25 path: String,
26 },
27 DataAvailable {
28 source: String,
29 timestamp: DateTime<Utc>,
30 },
31 ScheduledTime {
32 workflowid: String,
33 },
34 ExternalTrigger {
35 source: String,
36 payload: serde_json::Value,
37 },
38 WorkflowCompleted {
39 workflowid: String,
40 executionid: String,
41 },
42 Custom {
43 event_type: String,
44 data: serde_json::Value,
45 },
46 }
47 pub struct EventDrivenExecutor {
49 event_rx: Receiver<WorkflowEvent>,
50 event_tx: Sender<WorkflowEvent>,
51 rules: Vec<EventRule>,
52 executor: Arc<WorkflowExecutor>,
53 }
54 #[derive(Debug, Clone)]
55 pub struct EventRule {
56 pub id: String,
57 pub event_pattern: EventPattern,
58 pub workflowid: String,
59 pub parameters: HashMap<String, serde_json::Value>,
60 }
61 #[derive(Debug, Clone)]
62 pub enum EventPattern {
63 FilePattern {
64 path_regex: String,
65 },
66 SourcePattern {
67 source: String,
68 },
69 EventTypePattern {
70 event_type: String,
71 },
72 CompositePattern {
73 patterns: Vec<EventPattern>,
74 operator: LogicalOperator,
75 },
76 }
77 #[derive(Debug, Clone)]
78 pub enum LogicalOperator {
79 And,
80 Or,
81 Not,
82 }
83 impl EventDrivenExecutor {
84 pub fn new(executor: Arc<WorkflowExecutor>) -> Self {
85 let (tx, rx) = crossbeam_channel::unbounded();
86 Self {
87 event_rx: rx,
88 event_tx: tx,
89 rules: Vec::new(),
90 executor,
91 }
92 }
93 pub fn register_rule(&mut self, rule: EventRule) {
95 self.rules.push(rule);
96 }
97 pub fn get_event_sender(&self) -> Sender<WorkflowEvent> {
99 self.event_tx.clone()
100 }
101 pub fn process_events(&self, workflows: &HashMap<String, Workflow>) -> Result<()> {
103 while let Ok(event) = self.event_rx.try_recv() {
104 for rule in &self.rules {
105 if self.matches_pattern(&event, &rule.event_pattern) {
106 if let Some(workflow) = workflows.get(&rule.workflowid) {
107 let mut workflow = workflow.clone();
108 workflow.metadata.set(
109 "trigger_event",
110 MetadataValue::String(
111 serde_json::to_string(&event).expect("Operation failed"),
112 ),
113 );
114 self.executor.execute(&workflow)?;
115 }
116 }
117 }
118 }
119 Ok(())
120 }
121 #[allow(clippy::only_used_in_recursion)]
122 fn matches_pattern(&self, event: &WorkflowEvent, pattern: &EventPattern) -> bool {
123 match pattern {
124 EventPattern::FilePattern { path_regex } => {
125 if let WorkflowEvent::FileCreated { path }
126 | WorkflowEvent::FileModified { path } = event
127 {
128 regex::Regex::new(path_regex)
129 .map(|re| re.is_match(path))
130 .unwrap_or(false)
131 } else {
132 false
133 }
134 }
135 EventPattern::SourcePattern { source } => match event {
136 WorkflowEvent::DataAvailable { source: s, .. } => s == source,
137 WorkflowEvent::ExternalTrigger { source: s, .. } => s == source,
138 WorkflowEvent::FileCreated { .. } => false,
139 WorkflowEvent::FileModified { .. } => false,
140 WorkflowEvent::ScheduledTime { .. } => false,
141 WorkflowEvent::WorkflowCompleted { .. } => false,
142 WorkflowEvent::Custom { .. } => false,
143 },
144 EventPattern::EventTypePattern { event_type } => {
145 if let WorkflowEvent::Custom { event_type: t, .. } = event {
146 t == event_type
147 } else {
148 false
149 }
150 }
151 EventPattern::CompositePattern { patterns, operator } => match operator {
152 LogicalOperator::And => patterns.iter().all(|p| self.matches_pattern(event, p)),
153 LogicalOperator::Or => patterns.iter().any(|p| self.matches_pattern(event, p)),
154 LogicalOperator::Not => {
155 !patterns.iter().any(|p| self.matches_pattern(event, p))
156 }
157 },
158 }
159 }
160 }
161}
162pub mod versioning {
164 use super::*;
165 pub struct WorkflowVersionControl {
167 versions: HashMap<String, Vec<WorkflowVersion>>,
168 }
169 #[derive(Debug, Clone)]
170 pub struct WorkflowVersion {
171 pub version: String,
172 pub workflow: Workflow,
173 pub created_at: DateTime<Utc>,
174 pub created_by: String,
175 pub change_description: String,
176 pub parent_version: Option<String>,
177 }
178 impl Default for WorkflowVersionControl {
179 fn default() -> Self {
180 Self::new()
181 }
182 }
183 impl WorkflowVersionControl {
184 pub fn new() -> Self {
185 Self {
186 versions: HashMap::new(),
187 }
188 }
189 pub fn create_version(
191 &mut self,
192 workflow: Workflow,
193 created_by: impl Into<String>,
194 description: impl Into<String>,
195 ) -> String {
196 let workflowid = workflow.id.clone();
197 let versions = self.versions.entry(workflowid.clone()).or_default();
198 let version_number = versions.len() + 1;
199 let version = format!("v{version_number}.0.0");
200 let parent_version = versions.last().map(|v| v.version.clone());
201 versions.push(WorkflowVersion {
202 version: version.clone(),
203 workflow,
204 created_at: Utc::now(),
205 created_by: created_by.into(),
206 change_description: description.into(),
207 parent_version,
208 });
209 version
210 }
211 pub fn get_version(&self, workflowid: &str, version: &str) -> Option<&WorkflowVersion> {
213 self.versions
214 .get(workflowid)?
215 .iter()
216 .find(|v| v.version == version)
217 }
218 pub fn get_latest(&self, workflowid: &str) -> Option<&WorkflowVersion> {
220 self.versions.get(workflowid)?.last()
221 }
222 pub fn get_history(&self, workflowid: &str) -> Vec<&WorkflowVersion> {
224 self.versions
225 .get(workflowid)
226 .map(|v| v.iter().collect())
227 .unwrap_or_default()
228 }
229 pub fn diff(
231 &self,
232 workflowid: &str,
233 version1: &str,
234 version2: &str,
235 ) -> Option<WorkflowDiff> {
236 let v1 = self.get_version(workflowid, version1)?;
237 let v2 = self.get_version(workflowid, version2)?;
238 Some(WorkflowDiff {
239 version1: version1.to_string(),
240 version2: version2.to_string(),
241 added_tasks: self.diff_tasks(&v1.workflow.tasks, &v2.workflow.tasks, true),
242 removed_tasks: self.diff_tasks(&v1.workflow.tasks, &v2.workflow.tasks, false),
243 modified_tasks: self.find_modified_tasks(&v1.workflow.tasks, &v2.workflow.tasks),
244 dependency_changes: self
245 .diff_dependencies(&v1.workflow.dependencies, &v2.workflow.dependencies),
246 })
247 }
248 fn diff_tasks(&self, tasks1: &[Task], tasks2: &[Task], added: bool) -> Vec<String> {
249 let set1: HashSet<_> = tasks1.iter().map(|t| &t.id).collect();
250 let set2: HashSet<_> = tasks2.iter().map(|t| &t.id).collect();
251 if added {
252 set2.difference(&set1).map(|id| (*id).clone()).collect()
253 } else {
254 set1.difference(&set2).map(|id| (*id).clone()).collect()
255 }
256 }
257 fn find_modified_tasks(&self, tasks1: &[Task], tasks2: &[Task]) -> Vec<String> {
258 let map1: HashMap<&String, &Task> = tasks1.iter().map(|t| (&t.id, t)).collect();
259 let map2: HashMap<&String, &Task> = tasks2.iter().map(|t| (&t.id, t)).collect();
260 let mut modified = Vec::new();
261 for (id, task1) in map1 {
262 if let Some(task2) = map2.get(id) {
263 if task1.name != task2.name || task1.config != task2.config {
264 modified.push(id.clone());
265 }
266 }
267 }
268 modified
269 }
270 fn diff_dependencies(
271 &self,
272 deps1: &HashMap<String, Vec<String>>,
273 deps2: &HashMap<String, Vec<String>>,
274 ) -> Vec<DependencyChange> {
275 let mut changes = Vec::new();
276 let all_tasks: HashSet<_> = deps1.keys().chain(deps2.keys()).collect();
277 for task in all_tasks {
278 let deps1_set: HashSet<_> = deps1
279 .get(task)
280 .map(|d| d.iter().collect())
281 .unwrap_or_default();
282 let deps2_set: HashSet<_> = deps2
283 .get(task)
284 .map(|d| d.iter().collect())
285 .unwrap_or_default();
286 for added in deps2_set.difference(&deps1_set) {
287 changes.push(DependencyChange::Added {
288 task: (*task).clone(),
289 dependency: (*added).clone(),
290 });
291 }
292 for removed in deps1_set.difference(&deps2_set) {
293 changes.push(DependencyChange::Removed {
294 task: (*task).clone(),
295 dependency: (*removed).clone(),
296 });
297 }
298 }
299 changes
300 }
301 }
302 #[derive(Debug)]
303 pub struct WorkflowDiff {
304 pub version1: String,
305 pub version2: String,
306 pub added_tasks: Vec<String>,
307 pub removed_tasks: Vec<String>,
308 pub modified_tasks: Vec<String>,
309 pub dependency_changes: Vec<DependencyChange>,
310 }
311 #[derive(Debug)]
313 pub enum DependencyChange {
314 Added {
316 task: String,
318 dependency: String,
320 },
321 Removed {
323 task: String,
325 dependency: String,
327 },
328 }
329}
330pub mod distributed {
332 use super::*;
333 pub struct DistributedExecutor {
335 coordinator_url: String,
336 worker_pool: WorkerPool,
337 task_queue: Arc<Mutex<Vec<DistributedTask>>>,
338 }
339 #[derive(Debug, Clone)]
341 pub struct DistributedTask {
342 pub task: Task,
344 pub workflowid: String,
346 pub executionid: String,
348 pub assigned_worker: Option<String>,
350 pub status: TaskStatus,
352 }
353 pub struct WorkerPool {
355 workers: Vec<WorkerNode>,
356 }
357 #[derive(Debug, Clone)]
359 pub struct WorkerNode {
360 pub id: String,
362 pub url: String,
364 pub capabilities: WorkerCapabilities,
366 pub current_load: f64,
368 pub status: WorkerStatus,
370 }
371 #[derive(Debug, Clone)]
373 pub struct WorkerCapabilities {
374 pub cpu_cores: usize,
376 pub memorygb: f64,
378 pub gpu_available: bool,
380 pub supported_task_types: Vec<TaskType>,
382 }
383 #[derive(Debug, Clone, Copy, PartialEq)]
385 pub enum WorkerStatus {
386 Available,
388 Busy,
390 Offline,
392 }
393 impl DistributedExecutor {
394 pub fn new(coordinator_url: impl Into<String>) -> Self {
396 Self {
397 coordinator_url: coordinator_url.into(),
398 worker_pool: WorkerPool {
399 workers: Vec::new(),
400 },
401 task_queue: Arc::new(Mutex::new(Vec::new())),
402 }
403 }
404 pub fn register_worker(&mut self, worker: WorkerNode) {
406 self.worker_pool.workers.push(worker);
407 }
408 pub fn schedule_task(&self, task: DistributedTask) -> Result<String> {
410 let worker = self.find_suitable_worker(&task)?;
411 let mut queue = self.task_queue.lock().expect("Operation failed");
412 let mut scheduled_task = task;
413 scheduled_task.assigned_worker = Some(worker.id.clone());
414 queue.push(scheduled_task);
415 Ok(worker.id.clone())
416 }
417 fn find_suitable_worker(&self, task: &DistributedTask) -> Result<&WorkerNode> {
418 let suitable_workers: Vec<_> = self
419 .worker_pool
420 .workers
421 .iter()
422 .filter(|w| {
423 w.status == WorkerStatus::Available
424 && w.capabilities
425 .supported_task_types
426 .contains(&task.task.task_type)
427 && self.meets_resource_requirements(w, &task.task.resources)
428 })
429 .collect();
430 suitable_workers
431 .into_iter()
432 .min_by(|a, b| {
433 a.current_load
434 .partial_cmp(&b.current_load)
435 .expect("Operation failed")
436 })
437 .ok_or_else(|| IoError::Other("No suitable worker available".to_string()))
438 }
439 fn meets_resource_requirements(
440 &self,
441 worker: &WorkerNode,
442 requirements: &ResourceRequirements,
443 ) -> bool {
444 if let Some(cpu) = requirements.cpu_cores {
445 if worker.capabilities.cpu_cores < cpu {
446 return false;
447 }
448 }
449 if let Some(memory) = requirements.memorygb {
450 if worker.capabilities.memorygb < memory {
451 return false;
452 }
453 }
454 if requirements.gpu.is_some() && !worker.capabilities.gpu_available {
455 return false;
456 }
457 true
458 }
459 }
460}