moirai_iter/execution/
hybrid.rs1#![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#[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#[derive(Debug, Clone)]
33pub struct HybridConfig {
34 pub parallel_threshold: usize,
36 pub async_threshold: usize,
38 pub adaptation_factor: f64,
40 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#[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#[derive(Debug, Clone, Copy, PartialEq)]
66pub enum ExecutionStrategy {
67 Parallel,
69 Async,
71}
72
73impl Default for PerformanceHistory {
74 fn default() -> Self {
75 Self::new()
76 }
77}
78
79impl PerformanceHistory {
80 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 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 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 pub fn recommend_strategy(
106 &self,
107 item_count: usize,
108 config: &HybridConfig,
109 ) -> ExecutionStrategy {
110 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 if self.parallel_times.is_empty() && self.async_times.is_empty() {
121 if item_count < (config.parallel_threshold + config.async_threshold) / 2 {
123 ExecutionStrategy::Async
124 } else {
125 ExecutionStrategy::Parallel
126 }
127 } else {
128 let parallel_avg = if self.parallel_times.is_empty() {
130 Duration::from_secs(1) } 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) } else {
139 let sum: Duration = self.async_times.iter().sum();
140 sum / self.async_times.len() as u32
141 };
142
143 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 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 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 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 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 if let Ok(mut history) = self.performance_history.lock() {
203 match strategy {
204 ExecutionStrategy::Parallel => {
205 history.record_parallel_time(duration);
206 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 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 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 let strategy = self.choose_strategy(1); 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
245pub 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}