moirai_iter/execution/
hybrid.rs1use 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#[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#[derive(Debug, Clone)]
28pub struct HybridConfig {
29 pub parallel_threshold: usize,
31 pub async_threshold: usize,
33 pub adaptation_factor: f64,
35 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#[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#[derive(Debug, Clone, Copy, PartialEq)]
61pub enum ExecutionStrategy {
62 Parallel,
64 Async,
66}
67
68impl Default for PerformanceHistory {
69 fn default() -> Self {
70 Self::new()
71 }
72}
73
74impl PerformanceHistory {
75 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 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 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 pub fn recommend_strategy(
101 &self,
102 item_count: usize,
103 config: &HybridConfig,
104 ) -> ExecutionStrategy {
105 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 if self.parallel_times.is_empty() && self.async_times.is_empty() {
116 if item_count < (config.parallel_threshold + config.async_threshold) / 2 {
118 ExecutionStrategy::Async
119 } else {
120 ExecutionStrategy::Parallel
121 }
122 } else {
123 let parallel_avg = if self.parallel_times.is_empty() {
125 Duration::from_secs(1) } 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) } else {
134 let sum: Duration = self.async_times.iter().sum();
135 sum / self.async_times.len() as u32
136 };
137
138 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 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 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 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 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 if let Ok(mut history) = self.performance_history.lock() {
198 match strategy {
199 ExecutionStrategy::Parallel => {
200 history.record_parallel_time(duration);
201 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 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 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 let strategy = self.choose_strategy(1); 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
240pub 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}