rskit_process/runner/
observer.rs1use std::sync::Arc;
2
3pub type OutputLineCallback = Arc<dyn Fn(&str) + Send + Sync + 'static>;
5
6pub type OutputBytesCallback = Arc<dyn Fn(&[u8]) + Send + Sync + 'static>;
8
9#[derive(Clone, Default)]
11pub struct OutputObserver {
12 pub(in crate::runner) stdout_line: Option<OutputLineCallback>,
13 pub(in crate::runner) stderr_line: Option<OutputLineCallback>,
14 pub(in crate::runner) stdout_bytes: Option<OutputBytesCallback>,
15 pub(in crate::runner) stderr_bytes: Option<OutputBytesCallback>,
16}
17
18impl std::fmt::Debug for OutputObserver {
19 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
20 formatter
21 .debug_struct("OutputObserver")
22 .field("stdout_line", &self.stdout_line.is_some())
23 .field("stderr_line", &self.stderr_line.is_some())
24 .field("stdout_bytes", &self.stdout_bytes.is_some())
25 .field("stderr_bytes", &self.stderr_bytes.is_some())
26 .finish()
27 }
28}
29
30impl OutputObserver {
31 #[must_use]
33 pub fn new() -> Self {
34 Self::default()
35 }
36
37 #[must_use]
39 pub fn with_stdout_line(mut self, callback: impl Fn(&str) + Send + Sync + 'static) -> Self {
40 self.stdout_line = Some(Arc::new(callback));
41 self
42 }
43
44 #[must_use]
46 pub fn with_stderr_line(mut self, callback: impl Fn(&str) + Send + Sync + 'static) -> Self {
47 self.stderr_line = Some(Arc::new(callback));
48 self
49 }
50
51 #[must_use]
53 pub fn with_stdout_bytes(mut self, callback: impl Fn(&[u8]) + Send + Sync + 'static) -> Self {
54 self.stdout_bytes = Some(Arc::new(callback));
55 self
56 }
57
58 #[must_use]
60 pub fn with_stderr_bytes(mut self, callback: impl Fn(&[u8]) + Send + Sync + 'static) -> Self {
61 self.stderr_bytes = Some(Arc::new(callback));
62 self
63 }
64}
65
66#[cfg(test)]
67mod tests {
68 use std::sync::{
69 Arc,
70 atomic::{AtomicUsize, Ordering},
71 };
72
73 use super::*;
74
75 #[test]
76 fn output_observer_builders_store_callbacks_and_debug_flags() {
77 let stdout_lines = Arc::new(AtomicUsize::new(0));
78 let stderr_lines = Arc::new(AtomicUsize::new(0));
79 let stdout_bytes = Arc::new(AtomicUsize::new(0));
80 let stderr_bytes = Arc::new(AtomicUsize::new(0));
81
82 let observer = OutputObserver::new()
83 .with_stdout_line({
84 let calls = Arc::clone(&stdout_lines);
85 move |_| {
86 calls.fetch_add(1, Ordering::SeqCst);
87 }
88 })
89 .with_stderr_line({
90 let calls = Arc::clone(&stderr_lines);
91 move |_| {
92 calls.fetch_add(1, Ordering::SeqCst);
93 }
94 })
95 .with_stdout_bytes({
96 let bytes = Arc::clone(&stdout_bytes);
97 move |chunk| {
98 bytes.fetch_add(chunk.len(), Ordering::SeqCst);
99 }
100 })
101 .with_stderr_bytes({
102 let bytes = Arc::clone(&stderr_bytes);
103 move |chunk| {
104 bytes.fetch_add(chunk.len(), Ordering::SeqCst);
105 }
106 });
107
108 (observer.stdout_line.as_ref().unwrap())("out");
109 (observer.stderr_line.as_ref().unwrap())("err");
110 (observer.stdout_bytes.as_ref().unwrap())(b"stdout");
111 (observer.stderr_bytes.as_ref().unwrap())(b"stderr");
112
113 assert_eq!(stdout_lines.load(Ordering::SeqCst), 1);
114 assert_eq!(stderr_lines.load(Ordering::SeqCst), 1);
115 assert_eq!(stdout_bytes.load(Ordering::SeqCst), 6);
116 assert_eq!(stderr_bytes.load(Ordering::SeqCst), 6);
117
118 let debug = format!("{observer:?}");
119 assert!(debug.contains("stdout_line: true"));
120 assert!(debug.contains("stderr_line: true"));
121 assert!(debug.contains("stdout_bytes: true"));
122 assert!(debug.contains("stderr_bytes: true"));
123 }
124}