1use std::sync::Arc;
32use std::sync::Mutex;
33use std::sync::atomic::{AtomicBool, Ordering};
34use std::time::Instant;
35
36use crate::activity::ActivityMetrics;
37use crate::adapter::{AdapterError, ExecutionError, OpDispenser, OpResult, WrappingDispenser};
38use crate::error_policy::ErrorPolicy;
39use crate::phase_outcome::PhaseErrorDetail;
40use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
41
42pub const NAME: WrapperName = WrapperName::new("errors");
43
44const PHASE_ERROR_CAPTURE_CAP: usize = 64;
48
49fn triggers(s: WrapperSubject) -> bool {
52 s.op().is_some()
53}
54
55fn describe_assignment(s: WrapperSubject) -> Option<String> {
58 let op = s.op()?;
59 op.params
60 .get("errors")
61 .and_then(|v| v.as_str())
62 .map(|spec| format!("errors: {spec}"))
63}
64
65inventory::submit! {
66 WrapperRegistration {
67 name: NAME,
68 owned_fields: &["errors"],
69 triggers,
70 requires_inner: &[],
71 forbids_outer: &[],
72 mutually_exclusive_with: &[],
73 describe_assignment,
74 levels: &[crate::wrapper_registry::WrapperLevel::Op],
75 }
76}
77
78pub struct ErrorHandlerDispenser {
84 inner: Arc<dyn OpDispenser>,
85 policy: Arc<ErrorPolicy>,
90 metrics: Arc<ActivityMetrics>,
93 phase_errors: Arc<Mutex<Vec<PhaseErrorDetail>>>,
96 stop_flag: Arc<AtomicBool>,
98 stop_reason: Arc<Mutex<Option<String>>>,
100 op_name: String,
102 records_attempts: bool,
109}
110
111impl ErrorHandlerDispenser {
112 #[allow(clippy::too_many_arguments)]
113 pub fn wrap(
114 inner: Arc<dyn OpDispenser>,
115 policy: Arc<ErrorPolicy>,
116 metrics: Arc<ActivityMetrics>,
117 phase_errors: Arc<Mutex<Vec<PhaseErrorDetail>>>,
118 stop_flag: Arc<AtomicBool>,
119 stop_reason: Arc<Mutex<Option<String>>>,
120 op_name: String,
121 records_attempts: bool,
122 ) -> Arc<dyn OpDispenser> {
123 Arc::new(Self {
124 inner,
125 policy,
126 metrics,
127 phase_errors,
128 stop_flag,
129 stop_reason,
130 op_name,
131 records_attempts,
132 })
133 }
134
135 fn route_terminal_error(
139 &self,
140 e: &ExecutionError,
141 cycle: u64,
142 wires: &dyn crate::wires::WireSource,
143 service_nanos: u64,
144 ) {
145 let inner_err = e.error();
146 let detail = self.policy.router.handle_error(
147 &inner_err.error_name,
148 &inner_err.message,
149 cycle,
150 service_nanos,
151 );
152 self.metrics.errors_total.inc();
153 self.metrics.count_error_type(&detail.name);
154
155 if let Ok(mut errs) = self.phase_errors.lock() {
159 if errs.len() < PHASE_ERROR_CAPTURE_CAP {
160 let op_template = self.inner.describe();
161 let op_resolved = self.inner.describe_resolved(wires);
162 errs.push(PhaseErrorDetail {
163 class: inner_err.error_name.clone(),
164 message: inner_err.message.clone(),
165 op_name: Some(self.op_name.clone()),
166 cycle: Some(cycle),
167 op_template,
168 op_resolved,
169 at_nanos: std::time::SystemTime::now()
170 .duration_since(std::time::UNIX_EPOCH)
171 .map(|d| d.as_nanos() as u64)
172 .unwrap_or(0),
173 retryable: detail.is_retryable(),
174 });
175 }
176 }
177
178 if detail.should_stop {
179 self.stop_flag.store(true, Ordering::Relaxed);
180 if let Ok(mut slot) = self.stop_reason.lock()
186 && slot.is_none()
187 {
188 let op_shape = self
189 .inner
190 .describe()
191 .map(|d| format!("\n op-template: {d}"))
192 .unwrap_or_default();
193 let op_resolved = self
194 .inner
195 .describe_resolved(wires)
196 .map(|d| format!("\n op-resolved: {d}"))
197 .unwrap_or_default();
198 let first = inner_err
204 .message
205 .lines()
206 .next()
207 .unwrap_or(&inner_err.message);
208 *slot = Some(format!(
209 "[{}] op '{}' at cycle {}: {first}{op_shape}{op_resolved}",
210 inner_err.error_name, self.op_name, cycle,
211 ));
212 }
213 }
214 }
215}
216
217impl WrappingDispenser for ErrorHandlerDispenser {}
218
219impl OpDispenser for ErrorHandlerDispenser {
220 fn execute<'a>(
221 &'a self,
222 cycle: u64,
223 ctx: &'a crate::fixture::ExecCtx<'a>,
224 ) -> std::pin::Pin<
225 Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
226 > {
227 Box::pin(async move {
228 let service_start = Instant::now();
229 let outcome: Result<OpResult, ExecutionError> = {
235 use futures::FutureExt as _;
236 match std::panic::AssertUnwindSafe(self.inner.execute(cycle, ctx))
237 .catch_unwind()
238 .await
239 {
240 Ok(r) => r,
241 Err(payload) => {
242 let msg = payload
243 .downcast_ref::<&'static str>()
244 .map(|s| (*s).to_string())
245 .or_else(|| payload.downcast_ref::<String>().cloned())
246 .unwrap_or_else(|| "<non-string panic payload>".into());
247 Err(ExecutionError::Op(AdapterError {
248 error_name: "panic".into(),
249 message: msg,
250 retryable: false,
251 }))
252 }
253 }
254 };
255 match outcome {
256 Ok(result) => {
262 if self.records_attempts && !result.skipped {
263 let dt = service_start.elapsed().as_nanos() as u64;
264 self.metrics.attempt_total.inc();
265 self.metrics.attempt_success.observe(dt);
266 self.metrics.tries_histogram.record(1);
267 }
268 Ok(result)
269 }
270 Err(e) => {
271 let service_nanos = service_start.elapsed().as_nanos() as u64;
272 if self.records_attempts {
273 self.metrics.attempt_total.inc();
274 self.metrics.attempt_failure.observe(service_nanos);
275 self.metrics.tries_histogram.record(1);
276 }
277 self.route_terminal_error(&e, cycle, ctx.wires, service_nanos);
278 Err(e)
279 }
280 }
281 })
282 }
283
284 fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
285 Some(self.inner.as_ref())
286 }
287}
288
289#[cfg(test)]
290mod tests {
291 use super::*;
292 use crate::adapter::ResultBody;
293 use crate::fixture::{ExecCtx, ResolvedPulls};
294 use nmbrs_metrics::labels::Labels;
295
296 struct FakeInner {
298 error: Option<(String, String)>,
299 panics: bool,
300 }
301
302 impl OpDispenser for FakeInner {
303 fn execute<'a>(
304 &'a self,
305 _cycle: u64,
306 _ctx: &'a ExecCtx<'a>,
307 ) -> std::pin::Pin<
308 Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
309 > {
310 Box::pin(async move {
311 if self.panics {
312 panic!("inner blew up");
313 }
314 match &self.error {
315 Some((name, msg)) => Err(ExecutionError::Op(AdapterError {
316 error_name: name.clone(),
317 message: msg.clone(),
318 retryable: false,
319 })),
320 None => Ok(OpResult {
321 body: None::<Box<dyn ResultBody>>,
322 skipped: false,
323 }),
324 }
325 })
326 }
327 }
328
329 struct Harness {
330 metrics: Arc<ActivityMetrics>,
331 phase_errors: Arc<Mutex<Vec<PhaseErrorDetail>>>,
332 stop_flag: Arc<AtomicBool>,
333 stop_reason: Arc<Mutex<Option<String>>>,
334 }
335
336 fn wrap_with(spec: &str, inner: FakeInner) -> (Arc<dyn OpDispenser>, Harness) {
337 let policy = ErrorPolicy::standalone(crate::error_policy::PolicyConfig::new(spec, None));
338 let h = Harness {
339 metrics: Arc::new(ActivityMetrics::new(&Labels::empty())),
340 phase_errors: Arc::new(Mutex::new(Vec::new())),
341 stop_flag: Arc::new(AtomicBool::new(false)),
342 stop_reason: Arc::new(Mutex::new(None)),
343 };
344 let d = ErrorHandlerDispenser::wrap(
345 Arc::new(inner),
346 policy,
347 h.metrics.clone(),
348 h.phase_errors.clone(),
349 h.stop_flag.clone(),
350 h.stop_reason.clone(),
351 "test_op".into(),
352 true,
353 );
354 (d, h)
355 }
356
357 fn empty_ctx() -> (crate::adapter::ResolvedFields, ResolvedPulls) {
358 let fields = crate::adapter::ResolvedFields::new(vec![], vec![]);
359 let pulls = ResolvedPulls::empty();
360 (fields, pulls)
361 }
362
363 #[tokio::test]
366 async fn ok_passes_through_untouched() {
367 let (d, h) = wrap_with(
368 ".*:warn,stop",
369 FakeInner {
370 error: None,
371 panics: false,
372 },
373 );
374 let (fields, pulls) = empty_ctx();
375 let ctx = ExecCtx::new(&fields, &pulls);
376 d.execute(0, &ctx).await.expect("ok");
377 assert_eq!(h.metrics.errors_total.get(), 0);
378 assert!(h.phase_errors.lock().unwrap().is_empty());
379 assert!(!h.stop_flag.load(Ordering::Relaxed));
380 }
381
382 #[tokio::test]
385 async fn stop_policy_sets_flag_and_captures() {
386 let (d, h) = wrap_with(
387 ".*:warn,stop",
388 FakeInner {
389 error: Some(("ModelError".into(), "boom".into())),
390 panics: false,
391 },
392 );
393 let (fields, pulls) = empty_ctx();
394 let ctx = ExecCtx::new(&fields, &pulls);
395 let err = d.execute(7, &ctx).await.expect_err("must propagate");
396 assert_eq!(err.error().error_name, "ModelError");
397 assert_eq!(h.metrics.errors_total.get(), 1);
398 assert!(
399 h.stop_flag.load(Ordering::Relaxed),
400 "stop verb must set the flag"
401 );
402 let reason = h
403 .stop_reason
404 .lock()
405 .unwrap()
406 .clone()
407 .expect("reason captured");
408 assert!(
409 reason.contains("test_op") && reason.contains("cycle 7"),
410 "diagnostic names op + cycle: {reason}"
411 );
412 let errs = h.phase_errors.lock().unwrap();
413 assert_eq!(errs.len(), 1);
414 assert_eq!(errs[0].class, "ModelError");
415 }
416
417 #[tokio::test]
420 async fn lenient_policy_counts_without_stopping() {
421 let (d, h) = wrap_with(
422 ".*:warn,counter",
423 FakeInner {
424 error: Some(("Timeout".into(), "slow".into())),
425 panics: false,
426 },
427 );
428 let (fields, pulls) = empty_ctx();
429 let ctx = ExecCtx::new(&fields, &pulls);
430 let _ = d.execute(0, &ctx).await.expect_err("must propagate");
431 assert_eq!(h.metrics.errors_total.get(), 1);
432 assert!(
433 !h.stop_flag.load(Ordering::Relaxed),
434 "counter/warn must not stop"
435 );
436 assert!(h.stop_reason.lock().unwrap().is_none());
437 }
438
439 #[tokio::test]
443 async fn panic_below_is_routed_not_unwound() {
444 let (d, h) = wrap_with(
445 ".*:warn,stop",
446 FakeInner {
447 error: None,
448 panics: true,
449 },
450 );
451 let (fields, pulls) = empty_ctx();
452 let ctx = ExecCtx::new(&fields, &pulls);
453 let err = d.execute(0, &ctx).await.expect_err("panic becomes Err");
454 assert_eq!(err.error().error_name, "panic");
455 assert!(h.stop_flag.load(Ordering::Relaxed));
456 assert_eq!(h.phase_errors.lock().unwrap()[0].class, "panic");
457 }
458
459 #[tokio::test]
462 async fn capture_buffer_caps_but_totals_keep_counting() {
463 let (d, h) = wrap_with(
464 ".*:counter",
465 FakeInner {
466 error: Some(("E".into(), "m".into())),
467 panics: false,
468 },
469 );
470 let (fields, pulls) = empty_ctx();
471 let ctx = ExecCtx::new(&fields, &pulls);
472 for c in 0..(PHASE_ERROR_CAPTURE_CAP as u64 + 10) {
473 let _ = d.execute(c, &ctx).await;
474 }
475 assert_eq!(
476 h.phase_errors.lock().unwrap().len(),
477 PHASE_ERROR_CAPTURE_CAP
478 );
479 assert_eq!(
480 h.metrics.errors_total.get(),
481 PHASE_ERROR_CAPTURE_CAP as u64 + 10
482 );
483 }
484}