Skip to main content

rskit_process/runner/
observer.rs

1use std::sync::Arc;
2
3/// Callback invoked for line-oriented process output.
4pub type OutputLineCallback = Arc<dyn Fn(&str) + Send + Sync + 'static>;
5
6/// Callback invoked for raw process output bytes.
7pub type OutputBytesCallback = Arc<dyn Fn(&[u8]) + Send + Sync + 'static>;
8
9/// Optional callbacks for line-oriented process output.
10#[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    /// Create an observer without callbacks.
32    #[must_use]
33    pub fn new() -> Self {
34        Self::default()
35    }
36
37    /// Observe each stdout line.
38    #[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    /// Observe each stderr line.
45    #[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    /// Observe raw stdout bytes.
52    #[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    /// Observe raw stderr bytes.
59    #[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}