relux_runtime/observe/
progress.rs1use 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
42pub 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}