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,
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#[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 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 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 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 pub fn recommend_strategy(
94 &self,
95 item_count: usize,
96 config: &HybridConfig,
97 ) -> ExecutionStrategy {
98 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 if self.parallel_times.is_empty() && self.async_times.is_empty() {
109 if item_count < (config.parallel_threshold + config.async_threshold) / 2 {
111 ExecutionStrategy::Async
112 } else {
113 ExecutionStrategy::Parallel
114 }
115 } else {
116 let parallel_avg = if self.parallel_times.is_empty() {
118 Duration::from_secs(1) } 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) } else {
127 let sum: Duration = self.async_times.iter().sum();
128 sum / self.async_times.len() as u32
129 };
130
131 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 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 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 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 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 if let Ok(mut history) = self.performance_history.lock() {
191 match strategy {
192 ExecutionStrategy::Parallel => {
193 history.record_parallel_time(duration);
194 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 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 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 let strategy = self.choose_strategy(1); 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
233pub 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}