Skip to main content

moirai_iter/execution/
hybrid.rs

1//! Hybrid execution context.
2
3use std::collections::VecDeque;
4use std::sync::{Arc, Mutex};
5use std::time::{Duration, Instant};
6
7use super::async_ctx::AsyncContext;
8use super::base::ExecutionBase;
9use super::parallel::ParallelContext;
10
11/// Hybrid context that adapts between parallel and async execution
12#[derive(Clone)]
13pub struct HybridContext {
14    pub(super) parallel_context: ParallelContext,
15    pub(super) async_context: AsyncContext,
16    pub(super) performance_history: Arc<Mutex<PerformanceHistory>>,
17    pub(super) config: HybridConfig,
18}
19
20impl Default for HybridContext {
21    fn default() -> Self {
22        Self::new()
23    }
24}
25
26/// Configuration for hybrid execution strategy
27#[derive(Debug, Clone)]
28pub struct HybridConfig {
29    pub parallel_threshold: usize,
30    pub async_threshold: usize,
31    pub adaptation_factor: f64,
32    pub history_window: usize,
33}
34
35impl Default for HybridConfig {
36    fn default() -> Self {
37        Self {
38            parallel_threshold: 1000,
39            async_threshold: 10000,
40            adaptation_factor: 0.1,
41            history_window: 10,
42        }
43    }
44}
45
46/// Performance history for adaptive execution decisions
47#[derive(Debug)]
48pub struct PerformanceHistory {
49    pub(super) parallel_times: VecDeque<Duration>,
50    pub(super) async_times: VecDeque<Duration>,
51    pub(super) last_decision: Option<ExecutionStrategy>,
52    pub(super) decision_count: usize,
53}
54
55#[derive(Debug, Clone, Copy, PartialEq)]
56pub enum ExecutionStrategy {
57    Parallel,
58    Async,
59}
60
61impl Default for PerformanceHistory {
62    fn default() -> Self {
63        Self::new()
64    }
65}
66
67impl PerformanceHistory {
68    /// Create a new PerformanceHistory instance
69    pub fn new() -> Self {
70        Self {
71            parallel_times: VecDeque::new(),
72            async_times: VecDeque::new(),
73            last_decision: None,
74            decision_count: 0,
75        }
76    }
77
78    /// Record performance of parallel strategy
79    pub fn record_parallel_time(&mut self, duration: Duration) {
80        self.parallel_times.push_back(duration);
81        self.last_decision = Some(ExecutionStrategy::Parallel);
82        self.decision_count += 1;
83    }
84
85    /// Record performance of async strategy
86    pub fn record_async_time(&mut self, duration: Duration) {
87        self.async_times.push_back(duration);
88        self.last_decision = Some(ExecutionStrategy::Async);
89        self.decision_count += 1;
90    }
91
92    /// Recommend strategy based on history and item count
93    pub fn recommend_strategy(
94        &self,
95        item_count: usize,
96        config: &HybridConfig,
97    ) -> ExecutionStrategy {
98        // Use config thresholds and performance history for decision
99        if item_count < config.parallel_threshold {
100            return ExecutionStrategy::Async;
101        }
102
103        if item_count > config.async_threshold {
104            return ExecutionStrategy::Parallel;
105        }
106
107        // In the middle range, use performance history to decide
108        if self.parallel_times.is_empty() && self.async_times.is_empty() {
109            // No history, use simple heuristic
110            if item_count < (config.parallel_threshold + config.async_threshold) / 2 {
111                ExecutionStrategy::Async
112            } else {
113                ExecutionStrategy::Parallel
114            }
115        } else {
116            // Compare average performance
117            let parallel_avg = if self.parallel_times.is_empty() {
118                Duration::from_secs(1) // Assume high if no data
119            } else {
120                let sum: Duration = self.parallel_times.iter().sum();
121                sum / self.parallel_times.len() as u32
122            };
123
124            let async_avg = if self.async_times.is_empty() {
125                Duration::from_secs(1) // Assume high if no data
126            } else {
127                let sum: Duration = self.async_times.iter().sum();
128                sum / self.async_times.len() as u32
129            };
130
131            // Choose the faster strategy, with adaptation factor
132            let factor = config.adaptation_factor;
133            if parallel_avg.as_secs_f64() * factor < async_avg.as_secs_f64() {
134                ExecutionStrategy::Parallel
135            } else {
136                ExecutionStrategy::Async
137            }
138        }
139    }
140}
141
142impl HybridContext {
143    /// Create a new hybrid context with default configuration
144    pub fn new() -> Self {
145        Self {
146            parallel_context: ParallelContext::new(),
147            async_context: AsyncContext::new(),
148            performance_history: Arc::new(Mutex::new(PerformanceHistory::new())),
149            config: HybridConfig::default(),
150        }
151    }
152
153    /// Create with custom configuration
154    pub fn with_config(config: HybridConfig) -> Self {
155        Self {
156            parallel_context: ParallelContext::new(),
157            async_context: AsyncContext::new(),
158            performance_history: Arc::new(Mutex::new(PerformanceHistory::new())),
159            config,
160        }
161    }
162
163    /// Choose execution strategy based on workload characteristics
164    pub fn choose_strategy(&self, item_count: usize) -> ExecutionStrategy {
165        let history = self.performance_history.lock().unwrap();
166        history.recommend_strategy(item_count, &self.config)
167    }
168
169    /// Execute an iterator operation with hybrid processing
170    pub fn execute_iter<T, F, R>(
171        &self,
172        items: Vec<T>,
173        func: F,
174    ) -> Result<Vec<R>, Box<dyn std::error::Error + Send + Sync>>
175    where
176        T: Send + 'static,
177        F: Fn(T) -> R + Send + Sync + 'static,
178        R: Send + 'static,
179    {
180        let strategy = self.choose_strategy(items.len());
181
182        let start = Instant::now();
183        let result = match strategy {
184            ExecutionStrategy::Parallel => self.parallel_context.execute_iter(items, func),
185            ExecutionStrategy::Async => self.async_context.execute_iter(items, func),
186        };
187        let duration = start.elapsed();
188
189        // Record performance for future decisions
190        if let Ok(mut history) = self.performance_history.lock() {
191            match strategy {
192                ExecutionStrategy::Parallel => {
193                    history.record_parallel_time(duration);
194                    // Apply config window size
195                    while history.parallel_times.len() > self.config.history_window {
196                        history.parallel_times.pop_front();
197                    }
198                }
199                ExecutionStrategy::Async => {
200                    history.record_async_time(duration);
201                    // Apply config window size
202                    while history.async_times.len() > self.config.history_window {
203                        history.async_times.pop_front();
204                    }
205                }
206            }
207        }
208
209        result
210    }
211
212    /// Execute a closure with the context
213    pub fn execute<F, R>(&self, func: F) -> Result<R, Box<dyn std::error::Error + Send + Sync>>
214    where
215        F: FnOnce() -> R + Send,
216        R: Send,
217    {
218        // Choose strategy and delegate
219        let strategy = self.choose_strategy(1); // Single item
220        match strategy {
221            ExecutionStrategy::Parallel => self.parallel_context.execute(func),
222            ExecutionStrategy::Async => self.async_context.execute(func),
223        }
224    }
225}
226
227impl ExecutionBase for HybridContext {
228    fn context_type(&self) -> &'static str {
229        "Hybrid"
230    }
231}
232
233/// Helper function to chunk items
234pub fn owned_chunks<T>(items: Vec<T>, chunk_size: usize) -> Vec<Vec<T>> {
235    let chunk_size = chunk_size.max(1);
236    let item_count = items.len();
237    let mut remaining = item_count;
238    let mut iter = items.into_iter();
239    let mut chunks = Vec::with_capacity(item_count.div_ceil(chunk_size));
240
241    while remaining > 0 {
242        let take = remaining.min(chunk_size);
243        chunks.push(iter.by_ref().take(take).collect());
244        remaining -= take;
245    }
246
247    chunks
248}