Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_datime/src/validation/parallel.rs

10.8 KiB, 37 runs

created by r1870400018:8624, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! Validation of large batches across several threads.
2//!
3//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
4//! Anthropic Claude
5
6use crate::{
7 calendar::CalendarDate,
8 clock::ClockTime,
9 time::CalClock,
10 validation::{CalClockValidator, ValidationError, ValidationResult},
11};
12
13use oxedyne_fe2o3_core::prelude::*;
14
15use std::{
16 sync::{Arc, Mutex},
17 thread,
18 time::{Duration, Instant},
19};
20
21/// # Examples
22///
23/// ```ignore
24/// use oxedyne_fe2o3_datime::validation::{ParallelValidator, CalClockValidator};
25///
26/// let validator = CalClockValidator::new();
27/// let parallel_validator = ParallelValidator::new(validator, 4); // 4 threads
28///
29/// let calclocks = vec![/* large collection */];
30/// let results = parallel_validator.validate_batch(&calclocks);
31///
32/// println!("Validated {} items in parallel", results.total_items);
33/// ```
34#[derive(Debug)]
35pub struct ParallelValidator {
36 validator: Arc<CalClockValidator>,
37 thread_count: usize,
38 chunk_size: usize, // items handed to a thread at a time
39}
40
41impl ParallelValidator {
42 pub fn new(validator: CalClockValidator, thread_count: usize) -> Self {
43 Self {
44 validator: Arc::new(validator),
45 thread_count: std::cmp::max(1, thread_count),
46 chunk_size: 100, // Default chunk size
47 }
48 }
49
50 pub fn with_chunk_size(mut self, chunk_size: usize) -> Self {
51 self.chunk_size = std::cmp::max(1, chunk_size);
52 self
53 }
54
55 pub fn validate_batch(&self, calclocks: &[CalClock]) -> BatchValidationResult {
56 let _start_time = Instant::now();
57
58 if calclocks.is_empty() {
59 return BatchValidationResult {
60 total_items: 0,
61 valid_items: 0,
62 invalid_items: 0,
63 validation_errors: Vec::new(),
64 execution_time: Duration::new(0, 0),
65 thread_count: self.thread_count,
66 chunk_size: self.chunk_size,
67 };
68 }
69
70 // Clone the data to avoid lifetime issues
71 let calclocks_owned: Vec<CalClock> = calclocks.to_vec();
72
73 // Distribute work across threads
74 let results = Arc::new(Mutex::new(Vec::new()));
75 let mut handles = Vec::new();
76
77 let items_per_thread = (calclocks_owned.len() + self.thread_count - 1) / self.thread_count;
78
79 for thread_id in 0..self.thread_count {
80 let start_idx = thread_id * items_per_thread;
81 let end_idx = std::cmp::min(start_idx + items_per_thread, calclocks_owned.len());
82
83 if start_idx >= calclocks_owned.len() {
84 break;
85 }
86
87 let validator = Arc::clone(&self.validator);
88 let results = Arc::clone(&results);
89 let thread_items: Vec<CalClock> = calclocks_owned[start_idx..end_idx].to_vec();
90
91 let handle = thread::spawn(move || {
92 let mut local_results = Vec::new();
93
94 for (local_index, calclock) in thread_items.iter().enumerate() {
95 let global_index = start_idx + local_index;
96 let validation_result = validator.validate_calclock(calclock);
97 let item_result = ValidationItemResult {
98 index: global_index,
99 calclock: calclock.clone(),
100 result: validation_result,
101 };
102 local_results.push(item_result);
103 }
104
105 // Add to shared results
106 if let Ok(mut results) = results.lock() {
107 results.extend(local_results);
108 }
109 });
110
111 handles.push(handle);
112 }
113
114 // Wait for all threads to complete
115 for handle in handles {
116 if let Err(_) = handle.join() {
117 // Handle thread panic - in production you'd want better error handling
118 }
119 }
120
121 let execution_time = _start_time.elapsed();
122
123 // Collect results
124 let all_results = if let Ok(results) = results.lock() {
125 results.clone()
126 } else {
127 Vec::new()
128 };
129
130 // Aggregate statistics
131 let total_items = all_results.len();
132 let valid_items = all_results.iter().filter(|r| r.result.is_ok()).count();
133 let invalid_items = total_items - valid_items;
134
135 let validation_errors: Vec<ValidationItemError> = all_results
136 .iter()
137 .filter_map(|r| {
138 if let Err(errors) = &r.result {
139 Some(ValidationItemError {
140 index: r.index,
141 calclock: r.calclock.clone(),
142 errors: errors.clone(),
143 })
144 } else {
145 None
146 }
147 })
148 .collect();
149
150 BatchValidationResult {
151 total_items,
152 valid_items,
153 invalid_items,
154 validation_errors,
155 execution_time,
156 thread_count: self.thread_count,
157 chunk_size: self.chunk_size,
158 }
159 }
160
161 pub fn validate_dates_batch(&self, dates: &[CalendarDate]) -> BatchValidationResult {
162 let _start_time = Instant::now();
163
164 if dates.is_empty() {
165 return BatchValidationResult {
166 total_items: 0,
167 valid_items: 0,
168 invalid_items: 0,
169 validation_errors: Vec::new(),
170 execution_time: Duration::new(0, 0),
171 thread_count: self.thread_count,
172 chunk_size: self.chunk_size,
173 };
174 }
175
176 // Convert dates to minimal CalClocks for validation
177 let calclocks: Vec<CalClock> = dates
178 .iter()
179 .filter_map(|date| {
180 let zone = date.zone().clone();
181 if let Ok(time) = crate::clock::ClockTime::new(0, 0, 0, 0, zone) {
182 crate::time::CalClock::from_date_time(date.clone(), time).ok()
183 } else {
184 None
185 }
186 })
187 .collect();
188
189 self.validate_batch(&calclocks)
190 }
191
192 pub fn validate_times_batch(&self, times: &[ClockTime]) -> BatchValidationResult {
193 let _start_time = Instant::now();
194
195 if times.is_empty() {
196 return BatchValidationResult {
197 total_items: 0,
198 valid_items: 0,
199 invalid_items: 0,
200 validation_errors: Vec::new(),
201 execution_time: Duration::new(0, 0),
202 thread_count: self.thread_count,
203 chunk_size: self.chunk_size,
204 };
205 }
206
207 // Convert times to minimal CalClocks for validation
208 let calclocks: Vec<CalClock> = times
209 .iter()
210 .filter_map(|time| {
211 let zone = time.zone().clone();
212 if let Ok(date) = crate::calendar::CalendarDate::new(2024, 1, 1, zone) {
213 crate::time::CalClock::from_date_time(date, time.clone()).ok()
214 } else {
215 None
216 }
217 })
218 .collect();
219
220 self.validate_batch(&calclocks)
221 }
222
223 pub fn filter_valid_parallel(&self, calclocks: Vec<CalClock>) -> Vec<CalClock> {
224 let results = self.validate_batch(&calclocks);
225
226 calclocks
227 .into_iter()
228 .enumerate()
229 .filter(|(index, _)| {
230 !results.validation_errors.iter().any(|err| err.index == *index)
231 })
232 .map(|(_, calclock)| calclock)
233 .collect()
234 }
235
236 pub fn thread_count(&self) -> usize {
237 self.thread_count
238 }
239
240 pub fn chunk_size(&self) -> usize {
241 self.chunk_size
242 }
243}
244
245#[derive(Debug, Clone)]
246pub struct BatchValidationResult {
247 pub total_items: usize,
248 pub valid_items: usize,
249 pub invalid_items: usize,
250 pub validation_errors: Vec<ValidationItemError>,
251 pub execution_time: Duration,
252 pub thread_count: usize,
253 pub chunk_size: usize,
254}
255
256impl BatchValidationResult {
257 /// A fraction between zero and one, not a percentage.
258 pub fn success_rate(&self) -> f64 {
259 if self.total_items == 0 {
260 0.0
261 } else {
262 self.valid_items as f64 / self.total_items as f64
263 }
264 }
265
266 pub fn throughput(&self) -> f64 {
267 if self.execution_time.as_secs_f64() == 0.0 {
268 0.0
269 } else {
270 self.total_items as f64 / self.execution_time.as_secs_f64()
271 }
272 }
273
274 pub fn average_time_per_item(&self) -> Duration {
275 if self.total_items == 0 {
276 Duration::new(0, 0)
277 } else {
278 self.execution_time / self.total_items as u32
279 }
280 }
281
282 pub fn all_valid(&self) -> bool {
283 self.invalid_items == 0
284 }
285
286 pub fn error_summary(&self) -> std::collections::HashMap<String, usize> {
287 let mut summary = std::collections::HashMap::new();
288
289 for item_error in &self.validation_errors {
290 for error in &item_error.errors {
291 *summary.entry(error.rule.clone()).or_insert(0) += 1;
292 }
293 }
294
295 summary
296 }
297
298 pub fn format(&self) -> String {
299 format!(
300 "Batch Validation Result:\n\
301 - Total items: {}\n\
302 - Valid: {} ({:.1}%)\n\
303 - Invalid: {} ({:.1}%)\n\
304 - Execution time: {:?}\n\
305 - Throughput: {:.0} items/sec\n\
306 - Threads used: {}\n\
307 - Chunk size: {}",
308 self.total_items,
309 self.valid_items,
310 self.success_rate() * 100.0,
311 self.invalid_items,
312 (self.invalid_items as f64 / self.total_items as f64) * 100.0,
313 self.execution_time,
314 self.throughput(),
315 self.thread_count,
316 self.chunk_size
317 )
318 }
319}
320
321#[derive(Debug, Clone)]
322pub struct ValidationItemError {
323 pub index: usize, // position in the original batch
324 pub calclock: CalClock,
325 pub errors: Vec<ValidationError>,
326}
327
328#[derive(Debug, Clone)]
329struct ValidationItemResult {
330 index: usize,
331 calclock: CalClock,
332 result: ValidationResult,
333}
334
335pub fn optimal_thread_count() -> usize {
336 // Use number of logical CPUs, but cap at reasonable limits
337 let cpu_count = std::thread::available_parallelism()
338 .map(|n| n.get())
339 .unwrap_or(4);
340
341 // Cap between 2 and 16 threads for validation workloads
342 std::cmp::min(16, std::cmp::max(2, cpu_count))
343}
344
345pub fn create_optimal_parallel_validator(validator: CalClockValidator) -> ParallelValidator {
346 let thread_count = optimal_thread_count();
347 ParallelValidator::new(validator, thread_count)
348 .with_chunk_size(50) // Balanced chunk size for most workloads
349}