Skip to main content

moirai_iter/execution/
hybrid.rs

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