Skip to main content

scirs2_io/workflow/
functions_3.rs

1//! Auto-generated module
2//!
3//! 🤖 Generated with [SplitRS](https://github.com/cool-japan/splitrs)
4
5use 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
14/// Event-driven workflows
15pub mod events {
16    use super::*;
17    use crossbeam_channel::{Receiver, Sender};
18    /// Event types that can trigger workflows
19    #[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    /// Event-driven workflow executor
48    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        /// Register an event rule
94        pub fn register_rule(&mut self, rule: EventRule) {
95            self.rules.push(rule);
96        }
97        /// Get event sender for external systems
98        pub fn get_event_sender(&self) -> Sender<WorkflowEvent> {
99            self.event_tx.clone()
100        }
101        /// Process events and trigger workflows
102        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}
162/// Workflow versioning and history
163pub mod versioning {
164    use super::*;
165    /// Workflow version control
166    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        /// Create a new version
190        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        /// Get a specific version
212        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        /// Get latest version
219        pub fn get_latest(&self, workflowid: &str) -> Option<&WorkflowVersion> {
220            self.versions.get(workflowid)?.last()
221        }
222        /// Get version history
223        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        /// Diff two versions
230        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    /// Dependency change event
312    #[derive(Debug)]
313    pub enum DependencyChange {
314        /// Dependency was added
315        Added {
316            /// Task ID
317            task: String,
318            /// Dependency task ID
319            dependency: String,
320        },
321        /// Dependency was removed
322        Removed {
323            /// Task ID
324            task: String,
325            /// Dependency task ID
326            dependency: String,
327        },
328    }
329}
330/// Distributed execution support
331pub mod distributed {
332    use super::*;
333    /// Distributed workflow executor
334    pub struct DistributedExecutor {
335        coordinator_url: String,
336        worker_pool: WorkerPool,
337        task_queue: Arc<Mutex<Vec<DistributedTask>>>,
338    }
339    /// Task in a distributed workflow
340    #[derive(Debug, Clone)]
341    pub struct DistributedTask {
342        /// Task definition
343        pub task: Task,
344        /// Workflow identifier
345        pub workflowid: String,
346        /// Execution identifier
347        pub executionid: String,
348        /// Worker ID this task is assigned to
349        pub assigned_worker: Option<String>,
350        /// Current task status
351        pub status: TaskStatus,
352    }
353    /// Pool of worker nodes
354    pub struct WorkerPool {
355        workers: Vec<WorkerNode>,
356    }
357    /// Worker node in the distributed system
358    #[derive(Debug, Clone)]
359    pub struct WorkerNode {
360        /// Unique worker identifier
361        pub id: String,
362        /// Worker node URL
363        pub url: String,
364        /// Worker capabilities and resources
365        pub capabilities: WorkerCapabilities,
366        /// Current load (0.0-1.0)
367        pub current_load: f64,
368        /// Current worker status
369        pub status: WorkerStatus,
370    }
371    /// Worker node capabilities and resources
372    #[derive(Debug, Clone)]
373    pub struct WorkerCapabilities {
374        /// Number of CPU cores available
375        pub cpu_cores: usize,
376        /// Memory available in GB
377        pub memorygb: f64,
378        /// Whether GPU is available
379        pub gpu_available: bool,
380        /// Task types this worker can execute
381        pub supported_task_types: Vec<TaskType>,
382    }
383    /// Status of a worker node
384    #[derive(Debug, Clone, Copy, PartialEq)]
385    pub enum WorkerStatus {
386        /// Worker is available for new tasks
387        Available,
388        /// Worker is currently busy
389        Busy,
390        /// Worker is offline
391        Offline,
392    }
393    impl DistributedExecutor {
394        /// Create a new distributed executor
395        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        /// Register a worker node
405        pub fn register_worker(&mut self, worker: WorkerNode) {
406            self.worker_pool.workers.push(worker);
407        }
408        /// Schedule task to appropriate worker
409        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}