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