Skip to main content

relux_runtime/observe/
progress.rs

1use std::io::Write;
2use std::time::Instant;
3
4use colored::Colorize;
5use tokio::sync::mpsc;
6
7#[derive(Debug, Clone)]
8pub enum ProgressEvent {
9    Send,
10    MatchStart,
11    MatchDone,
12    SleepStart,
13    SleepDone,
14    ShellSwitch(String),
15    FnEnter(String),
16    FnExit,
17    ShellSpawn,
18    ShellTerminate,
19    EffectSetup(String),
20    EffectTeardown,
21    Cleanup,
22    FailPattern,
23    Timeout,
24    Failure,
25    Cancellation,
26    Error(String),
27    Warning(String),
28    Annotation(String),
29}
30
31pub type ProgressTx = mpsc::UnboundedSender<ProgressEvent>;
32
33pub fn channel() -> (ProgressTx, mpsc::UnboundedReceiver<ProgressEvent>) {
34    mpsc::unbounded_channel()
35}
36
37enum TimedWait {
38    Match,
39    Sleep,
40}
41
42/// Spawns the progress printer task. Returns a JoinHandle that resolves
43/// to the collected progress string once all senders are dropped.
44pub fn spawn_printer(
45    mut rx: mpsc::UnboundedReceiver<ProgressEvent>,
46) -> tokio::task::JoinHandle<String> {
47    tokio::spawn(async move {
48        let mut collected = String::new();
49        let mut timed: Option<(TimedWait, Instant)> = None;
50        let mut timed_tick_count: usize = 0;
51
52        loop {
53            let event = if timed.is_some() {
54                match tokio::time::timeout(std::time::Duration::from_secs(1), rx.recv()).await {
55                    Ok(Some(ev)) => Some(ev),
56                    Ok(None) => None,
57                    Err(_) => {
58                        if let Some((kind, started)) = &timed {
59                            let ch = match kind {
60                                TimedWait::Match => '~',
61                                TimedWait::Sleep => 'z',
62                            };
63                            let elapsed_secs = started.elapsed().as_secs() as usize;
64                            while timed_tick_count < elapsed_secs {
65                                emit(&mut collected, ch);
66                                timed_tick_count += 1;
67                            }
68                        }
69                        continue;
70                    }
71                }
72            } else {
73                rx.recv().await
74            };
75
76            let Some(event) = event else {
77                break;
78            };
79
80            match event {
81                ProgressEvent::Send => {
82                    emit(&mut collected, '.');
83                }
84                ProgressEvent::MatchStart => {
85                    timed = Some((TimedWait::Match, Instant::now()));
86                    timed_tick_count = 0;
87                }
88                ProgressEvent::MatchDone => {
89                    timed = None;
90                    emit(&mut collected, '.');
91                }
92                ProgressEvent::SleepStart => {
93                    timed = Some((TimedWait::Sleep, Instant::now()));
94                    timed_tick_count = 0;
95                }
96                ProgressEvent::SleepDone => {
97                    timed = None;
98                }
99                ProgressEvent::ShellSwitch(_) => {
100                    emit(&mut collected, '|');
101                }
102                ProgressEvent::FnEnter(_) => {
103                    emit(&mut collected, '{');
104                }
105                ProgressEvent::FnExit => {
106                    emit(&mut collected, '}');
107                }
108                ProgressEvent::ShellSpawn => {
109                    emit(&mut collected, '+');
110                }
111                ProgressEvent::ShellTerminate => {
112                    emit(&mut collected, '-');
113                }
114                ProgressEvent::EffectSetup(_) => {
115                    emit(&mut collected, '+');
116                }
117                ProgressEvent::EffectTeardown => {
118                    emit(&mut collected, '-');
119                }
120                ProgressEvent::Cleanup => {
121                    emit(&mut collected, 'c');
122                }
123                ProgressEvent::FailPattern => {
124                    timed = None;
125                    emit(&mut collected, '!');
126                }
127                ProgressEvent::Timeout => {
128                    timed = None;
129                    emit(&mut collected, 'T');
130                }
131                ProgressEvent::Failure => {
132                    timed = None;
133                    emit(&mut collected, 'F');
134                }
135                ProgressEvent::Cancellation => {
136                    timed = None;
137                    emit(&mut collected, 'C');
138                }
139                ProgressEvent::Error(_) => {
140                    timed = None;
141                    emit(&mut collected, 'E');
142                }
143                ProgressEvent::Warning(_) => {
144                    emit(&mut collected, 'W');
145                }
146                ProgressEvent::Annotation(text) => {
147                    let s = format!("({text})");
148                    collected.push_str(&s);
149                    eprint!("{}", s.dimmed());
150                    let _ = std::io::stderr().flush();
151                }
152            }
153        }
154
155        collected
156    })
157}
158
159fn emit(collected: &mut String, ch: char) {
160    collected.push(ch);
161    let s = ch.to_string().dimmed();
162    eprint!("{s}");
163    let _ = std::io::stderr().flush();
164}