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 | |
| 6 | use crate::{ |
| 7 | calendar::CalendarDate, |
| 8 | clock::ClockTime, |
| 9 | time::CalClock, |
| 10 | validation::{CalClockValidator, ValidationError, ValidationResult}, |
| 11 | }; |
| 12 | |
| 13 | use oxedyne_fe2o3_core::prelude::*; |
| 14 | |
| 15 | use 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)] |
| 35 | pub struct ParallelValidator { |
| 36 | validator: Arc<CalClockValidator>, |
| 37 | thread_count: usize, |
| 38 | chunk_size: usize, // items handed to a thread at a time |
| 39 | } |
| 40 | |
| 41 | impl 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)] |
| 246 | pub 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 | |
| 256 | impl 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)] |
| 322 | pub 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)] |
| 329 | struct ValidationItemResult { |
| 330 | index: usize, |
| 331 | calclock: CalClock, |
| 332 | result: ValidationResult, |
| 333 | } |
| 334 | |
| 335 | pub 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 | |
| 345 | pub 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 | } |