1use std::any::Any;
16use std::io::Write;
17use std::panic::{self, AssertUnwindSafe};
18use std::sync::mpsc::{self, RecvTimeoutError};
19use std::thread::{self, JoinHandle};
20use std::time::Duration;
21
22use crate::console::Console;
23use crate::control::Control;
24use crate::live_render::LiveRender;
25use crate::protocol::Renderable;
26
27pub struct Live<W: Write> {
31 live_render: LiveRender,
32 console: Console,
33 writer: W,
34 started: bool,
35 transient: bool,
36}
37
38impl<W: Write> Live<W> {
39 pub fn new(renderable: Box<dyn Renderable>, console: Console, writer: W) -> Self {
42 Live {
43 live_render: LiveRender::new(renderable),
44 console,
45 writer,
46 started: false,
47 transient: false,
48 }
49 }
50
51 pub fn transient(mut self, transient: bool) -> Self {
54 self.transient = transient;
55 self
56 }
57
58 pub fn start(&mut self) {
60 if self.started {
61 return;
62 }
63 self.started = true;
64 if self.console.is_terminal() {
65 let _ = write!(self.writer, "{}", Control::show_cursor(false).as_str());
66 }
67 self.refresh();
68 }
69
70 pub fn update(&mut self, renderable: Box<dyn Renderable>) {
72 self.live_render.set_renderable(renderable);
73 self.refresh();
74 }
75
76 pub fn refresh(&mut self) {
79 if !self.console.is_terminal() {
80 return;
81 }
82 let position = self.live_render.position_cursor();
85 let content = self.console.render_to_string(&self.live_render);
86 let _ = write!(self.writer, "{}{}", position.as_str(), content);
87 }
88
89 pub fn stop(&mut self) {
91 if !self.started {
92 return;
93 }
94 if !self.console.is_terminal() {
95 if !self.transient {
98 let content = self.console.render_to_string(&self.live_render);
99 let _ = write!(self.writer, "{content}");
100 }
101 self.started = false;
102 return;
103 }
104 let position = self.live_render.position_cursor();
105 let content = self.console.render_to_string(&self.live_render);
106 let newline = if self.live_render.last_render_height() > 0 {
109 "\n"
110 } else {
111 ""
112 };
113 let _ = write!(
114 self.writer,
115 "{}{}{newline}{}",
116 position.as_str(),
117 content,
118 Control::show_cursor(true).as_str()
119 );
120 if self.transient {
121 let _ = write!(
122 self.writer,
123 "{}",
124 self.live_render.restore_cursor().as_str()
125 );
126 }
127 self.started = false;
128 }
129
130 fn abort(&mut self) {
134 if !self.started {
135 return;
136 }
137 self.started = false;
138 if self.console.is_terminal() {
139 let _ = write!(self.writer, "\n{}", Control::show_cursor(true).as_str());
140 let _ = self.writer.flush();
141 }
142 }
143
144 pub fn writer(&self) -> &W {
146 &self.writer
147 }
148
149 pub fn into_writer(self) -> W {
151 self.writer
152 }
153}
154
155enum LiveMessage {
157 Update(Box<dyn Renderable + Send>),
159 Refresh,
161 RefreshAck(mpsc::Sender<()>),
163 Stop,
165}
166
167impl<W: Write + Send + 'static> Live<W> {
168 pub fn spawn(
174 renderable: Box<dyn Renderable + Send>,
175 console: Console,
176 writer: W,
177 refresh_per_second: f64,
178 ) -> AutoLive<W> {
179 Live::spawn_with(renderable, console, writer, refresh_per_second, false)
180 }
181
182 pub fn spawn_with(
185 renderable: Box<dyn Renderable + Send>,
186 console: Console,
187 writer: W,
188 refresh_per_second: f64,
189 transient: bool,
190 ) -> AutoLive<W> {
191 let (sender, receiver) = mpsc::channel::<LiveMessage>();
192 let (started, wait_started) = mpsc::channel::<()>();
193 let interval = Duration::from_secs_f64(1.0 / refresh_per_second.max(f64::MIN_POSITIVE));
194 let handle = thread::spawn(move || {
195 let mut live = Live::new(renderable, console, writer).transient(transient);
198 let outcome = panic::catch_unwind(AssertUnwindSafe(|| {
203 live.start();
204 let _ = started.send(());
205 loop {
206 match receiver.recv_timeout(interval) {
207 Ok(LiveMessage::Update(renderable)) => live.update(renderable),
208 Ok(LiveMessage::Refresh) | Err(RecvTimeoutError::Timeout) => live.refresh(),
209 Ok(LiveMessage::RefreshAck(done)) => {
210 live.refresh();
211 let _ = done.send(());
212 }
213 Ok(LiveMessage::Stop) | Err(RecvTimeoutError::Disconnected) => {
215 live.stop();
216 break;
217 }
218 }
219 }
220 }));
221 let failure = outcome.err().map(|payload| {
222 live.abort();
223 panic_message(payload.as_ref())
224 });
225 (live.into_writer(), failure)
226 });
227 let _ = wait_started.recv();
231 AutoLive {
232 sender,
233 handle: Some(handle),
234 }
235 }
236}
237
238pub struct AutoLive<W: Write + Send + 'static> {
241 sender: mpsc::Sender<LiveMessage>,
242 handle: Option<JoinHandle<(W, Option<String>)>>,
243}
244
245#[derive(Debug)]
248pub struct LivePanic<W> {
249 pub writer: W,
252 pub message: String,
254}
255
256fn panic_message(payload: &(dyn Any + Send)) -> String {
258 if let Some(message) = payload.downcast_ref::<&str>() {
259 (*message).to_string()
260 } else if let Some(message) = payload.downcast_ref::<String>() {
261 message.clone()
262 } else {
263 "<non-string panic payload>".to_string()
264 }
265}
266
267impl<W: Write + Send + 'static> AutoLive<W> {
268 pub fn update(&self, renderable: Box<dyn Renderable + Send>) {
270 let _ = self.sender.send(LiveMessage::Update(renderable));
271 }
272
273 pub fn refresh(&self) {
275 let _ = self.sender.send(LiveMessage::Refresh);
276 }
277
278 pub fn refresh_wait(&self) {
282 let (done, wait) = mpsc::channel();
283 if self.sender.send(LiveMessage::RefreshAck(done)).is_ok() {
284 let _ = wait.recv();
285 }
286 }
287
288 pub fn stop(self) -> W {
294 match self.try_stop() {
295 Ok(writer) => writer,
296 Err(failure) => failure.writer,
297 }
298 }
299
300 pub fn try_stop(mut self) -> Result<W, LivePanic<W>> {
303 let _ = self.sender.send(LiveMessage::Stop);
304 let handle = self
305 .handle
306 .take()
307 .expect("thread handle present until stop/drop");
308 let (writer, failure) = match handle.join() {
311 Ok(result) => result,
312 Err(payload) => panic::resume_unwind(payload),
313 };
314 match failure {
315 None => Ok(writer),
316 Some(message) => Err(LivePanic { writer, message }),
317 }
318 }
319}
320
321impl<W: Write + Send + 'static> Drop for AutoLive<W> {
322 fn drop(&mut self) {
323 if let Some(handle) = self.handle.take() {
325 let _ = self.sender.send(LiveMessage::Stop);
326 let _ = handle.join();
327 }
328 }
329}
330
331#[cfg(test)]
332mod tests {
333 use super::*;
334 use crate::color::ColorSystem;
335 use crate::text::Text;
336
337 fn console() -> Console {
338 Console::builder()
339 .force_terminal(true)
340 .color_system(Some(ColorSystem::Truecolor))
341 .width(20)
342 .no_color(false)
343 .build()
344 }
345
346 #[test]
347 fn redirected_live_emits_only_the_final_frame() {
348 let console = Console::builder().force_terminal(false).width(20).build();
349 let mut live = Live::new(Box::new(Text::new("first")), console, Vec::<u8>::new());
350 live.start();
351 live.update(Box::new(Text::new("last")));
352 live.refresh();
353 assert!(live.writer().is_empty());
354 live.stop();
355 assert_eq!(live.writer(), b"last");
356 }
357
358 #[test]
359 fn manual_refresh_stream_matches_upstream() {
360 let mut live = Live::new(
361 Box::new(Text::new("frame one")),
362 console(),
363 Vec::<u8>::new(),
364 );
365 live.start();
366 live.update(Box::new(Text::new("frame two")));
367 live.update(Box::new(Text::new("frame three")));
368 live.stop();
369
370 let expected = "\x1b[?25lframe one\r\x1b[2Kframe two\r\x1b[2Kframe three\r\x1b[2Kframe three\n\x1b[?25h";
373 assert_eq!(String::from_utf8(live.writer().clone()).unwrap(), expected);
374 }
375
376 #[test]
377 fn auto_refresh_thread_produces_the_same_stream() {
378 let auto = Live::spawn(
383 Box::new(Text::new("frame one")),
384 console(),
385 Vec::<u8>::new(),
386 0.1,
387 );
388 auto.update(Box::new(Text::new("frame two")));
389 auto.update(Box::new(Text::new("frame three")));
390 let output = auto.stop();
391
392 let expected = "\x1b[?25lframe one\r\x1b[2Kframe two\r\x1b[2Kframe three\r\x1b[2Kframe three\n\x1b[?25h";
393 assert_eq!(String::from_utf8(output).unwrap(), expected);
394 }
395
396 struct PanicsLater {
398 renders: std::sync::Arc<std::sync::atomic::AtomicUsize>,
399 fail_after: usize,
400 }
401
402 impl Renderable for PanicsLater {
403 fn rich_render(
404 &self,
405 console: &Console,
406 options: &crate::console::ConsoleOptions,
407 ) -> Vec<crate::segment::Segment> {
408 use std::sync::atomic::Ordering;
409 if self.renders.fetch_add(1, Ordering::SeqCst) >= self.fail_after {
410 panic!("render failed");
411 }
412 Text::new("ok").rich_render(console, options)
413 }
414 }
415
416 #[test]
417 fn a_panicking_render_still_restores_the_cursor() {
418 let renders = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
422 let auto = Live::spawn(
423 Box::new(PanicsLater {
424 renders: renders.clone(),
425 fail_after: 1,
426 }),
427 console(),
428 Vec::<u8>::new(),
429 0.1,
430 );
431 auto.refresh_wait(); let (output, error) = match auto.try_stop() {
433 Ok(output) => (output, None),
434 Err(failure) => (failure.writer, Some(failure.message)),
435 };
436 let output = String::from_utf8(output).unwrap();
437 assert!(output.starts_with("\x1b[?25lok"), "{output:?}");
438 assert!(
439 output.ends_with("\x1b[?25h"),
440 "cursor left hidden: {output:?}"
441 );
442 assert_eq!(error.as_deref(), Some("render failed"));
443
444 let auto = Live::spawn(
446 Box::new(PanicsLater {
447 renders: std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)),
448 fail_after: 1,
449 }),
450 console(),
451 Vec::<u8>::new(),
452 0.1,
453 );
454 auto.refresh();
455 let output = String::from_utf8(auto.stop()).unwrap();
456 assert!(
457 output.ends_with("\x1b[?25h"),
458 "cursor left hidden: {output:?}"
459 );
460 }
461
462 #[test]
463 fn dropping_the_handle_finalizes_the_display() {
464 let console = console();
467 let auto = Live::spawn(Box::new(Text::new("only")), console, Vec::<u8>::new(), 0.1);
469 drop(auto); }
474}