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, VerticalOverflow};
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 vertical_overflow(mut self, vertical_overflow: VerticalOverflow) -> Self {
62 self.live_render.set_vertical_overflow(vertical_overflow);
63 self
64 }
65
66 pub fn start(&mut self) {
68 if self.started {
69 return;
70 }
71 self.started = true;
72 if self.console.is_terminal() {
73 let _ = write!(self.writer, "{}", Control::show_cursor(false).as_str());
74 }
75 self.refresh();
76 }
77
78 pub fn update(&mut self, renderable: Box<dyn Renderable>) {
80 self.live_render.set_renderable(renderable);
81 self.refresh();
82 }
83
84 pub fn refresh(&mut self) {
87 if !self.console.is_terminal() {
88 return;
89 }
90 let position = self.live_render.position_cursor();
93 let content = self.console.render_to_string(&self.live_render);
94 let _ = write!(self.writer, "{}{}", position.as_str(), content);
95 }
96
97 pub fn stop(&mut self) {
99 if !self.started {
100 return;
101 }
102 self.live_render
104 .set_vertical_overflow(VerticalOverflow::Visible);
105 if !self.console.is_terminal() {
106 if !self.transient {
109 let content = self.console.render_to_string(&self.live_render);
110 let _ = write!(self.writer, "{content}");
111 }
112 self.started = false;
113 return;
114 }
115 let position = self.live_render.position_cursor();
116 let content = self.console.render_to_string(&self.live_render);
117 let newline = if self.live_render.last_render_height() > 0 {
120 "\n"
121 } else {
122 ""
123 };
124 let _ = write!(
125 self.writer,
126 "{}{}{newline}{}",
127 position.as_str(),
128 content,
129 Control::show_cursor(true).as_str()
130 );
131 if self.transient {
132 let _ = write!(
133 self.writer,
134 "{}",
135 self.live_render.restore_cursor().as_str()
136 );
137 }
138 self.started = false;
139 }
140
141 fn abort(&mut self) {
145 if !self.started {
146 return;
147 }
148 self.started = false;
149 if self.console.is_terminal() {
150 let _ = write!(self.writer, "\n{}", Control::show_cursor(true).as_str());
151 let _ = self.writer.flush();
152 }
153 }
154
155 pub fn writer(&self) -> &W {
157 &self.writer
158 }
159
160 pub fn into_writer(self) -> W {
162 self.writer
163 }
164}
165
166enum LiveMessage {
168 Update(Box<dyn Renderable + Send>),
170 Refresh,
172 RefreshAck(mpsc::Sender<()>),
174 Stop,
176}
177
178impl<W: Write + Send + 'static> Live<W> {
179 pub fn spawn(
185 renderable: Box<dyn Renderable + Send>,
186 console: Console,
187 writer: W,
188 refresh_per_second: f64,
189 ) -> AutoLive<W> {
190 Live::spawn_with(renderable, console, writer, refresh_per_second, false)
191 }
192
193 pub fn spawn_with(
196 renderable: Box<dyn Renderable + Send>,
197 console: Console,
198 writer: W,
199 refresh_per_second: f64,
200 transient: bool,
201 ) -> AutoLive<W> {
202 let (sender, receiver) = mpsc::channel::<LiveMessage>();
203 let (started, wait_started) = mpsc::channel::<()>();
204 let interval = Duration::from_secs_f64(1.0 / refresh_per_second.max(f64::MIN_POSITIVE));
205 let handle = thread::spawn(move || {
206 let mut live = Live::new(renderable, console, writer).transient(transient);
209 let outcome = panic::catch_unwind(AssertUnwindSafe(|| {
214 live.start();
215 let _ = started.send(());
216 loop {
217 match receiver.recv_timeout(interval) {
218 Ok(LiveMessage::Update(renderable)) => live.update(renderable),
219 Ok(LiveMessage::Refresh) | Err(RecvTimeoutError::Timeout) => live.refresh(),
220 Ok(LiveMessage::RefreshAck(done)) => {
221 live.refresh();
222 let _ = done.send(());
223 }
224 Ok(LiveMessage::Stop) | Err(RecvTimeoutError::Disconnected) => {
226 live.stop();
227 break;
228 }
229 }
230 }
231 }));
232 let failure = outcome.err().map(|payload| {
233 live.abort();
234 panic_message(payload.as_ref())
235 });
236 (live.into_writer(), failure)
237 });
238 let _ = wait_started.recv();
242 AutoLive {
243 sender,
244 handle: Some(handle),
245 }
246 }
247}
248
249pub struct AutoLive<W: Write + Send + 'static> {
252 sender: mpsc::Sender<LiveMessage>,
253 handle: Option<JoinHandle<(W, Option<String>)>>,
254}
255
256#[derive(Debug)]
259pub struct LivePanic<W> {
260 pub writer: W,
263 pub message: String,
265}
266
267fn panic_message(payload: &(dyn Any + Send)) -> String {
269 if let Some(message) = payload.downcast_ref::<&str>() {
270 (*message).to_string()
271 } else if let Some(message) = payload.downcast_ref::<String>() {
272 message.clone()
273 } else {
274 "<non-string panic payload>".to_string()
275 }
276}
277
278impl<W: Write + Send + 'static> AutoLive<W> {
279 pub fn update(&self, renderable: Box<dyn Renderable + Send>) {
281 let _ = self.sender.send(LiveMessage::Update(renderable));
282 }
283
284 pub fn refresh(&self) {
286 let _ = self.sender.send(LiveMessage::Refresh);
287 }
288
289 pub fn refresh_wait(&self) {
293 let (done, wait) = mpsc::channel();
294 if self.sender.send(LiveMessage::RefreshAck(done)).is_ok() {
295 let _ = wait.recv();
296 }
297 }
298
299 pub fn stop(self) -> W {
305 match self.try_stop() {
306 Ok(writer) => writer,
307 Err(failure) => failure.writer,
308 }
309 }
310
311 pub fn try_stop(mut self) -> Result<W, LivePanic<W>> {
314 let _ = self.sender.send(LiveMessage::Stop);
315 let handle = self
316 .handle
317 .take()
318 .expect("thread handle present until stop/drop");
319 let (writer, failure) = match handle.join() {
322 Ok(result) => result,
323 Err(payload) => panic::resume_unwind(payload),
324 };
325 match failure {
326 None => Ok(writer),
327 Some(message) => Err(LivePanic { writer, message }),
328 }
329 }
330}
331
332impl<W: Write + Send + 'static> Drop for AutoLive<W> {
333 fn drop(&mut self) {
334 if let Some(handle) = self.handle.take() {
336 let _ = self.sender.send(LiveMessage::Stop);
337 let _ = handle.join();
338 }
339 }
340}
341
342#[cfg(test)]
343mod tests {
344 use super::*;
345 use crate::color::ColorSystem;
346 use crate::text::Text;
347
348 fn console() -> Console {
349 Console::builder()
350 .force_terminal(true)
351 .color_system(Some(ColorSystem::Truecolor))
352 .width(20)
353 .no_color(false)
354 .build()
355 }
356
357 #[test]
358 fn redirected_live_emits_only_the_final_frame() {
359 let console = Console::builder().force_terminal(false).width(20).build();
360 let mut live = Live::new(Box::new(Text::new("first")), console, Vec::<u8>::new());
361 live.start();
362 live.update(Box::new(Text::new("last")));
363 live.refresh();
364 assert!(live.writer().is_empty());
365 live.stop();
366 assert_eq!(live.writer(), b"last");
367 }
368
369 #[test]
370 fn manual_refresh_stream_matches_upstream() {
371 let mut live = Live::new(
372 Box::new(Text::new("frame one")),
373 console(),
374 Vec::<u8>::new(),
375 );
376 live.start();
377 live.update(Box::new(Text::new("frame two")));
378 live.update(Box::new(Text::new("frame three")));
379 live.stop();
380
381 let expected = "\x1b[?25lframe one\r\x1b[2Kframe two\r\x1b[2Kframe three\r\x1b[2Kframe three\n\x1b[?25h";
384 assert_eq!(String::from_utf8(live.writer().clone()).unwrap(), expected);
385 }
386
387 #[test]
388 fn auto_refresh_thread_produces_the_same_stream() {
389 let auto = Live::spawn(
394 Box::new(Text::new("frame one")),
395 console(),
396 Vec::<u8>::new(),
397 0.1,
398 );
399 auto.update(Box::new(Text::new("frame two")));
400 auto.update(Box::new(Text::new("frame three")));
401 let output = auto.stop();
402
403 let expected = "\x1b[?25lframe one\r\x1b[2Kframe two\r\x1b[2Kframe three\r\x1b[2Kframe three\n\x1b[?25h";
404 assert_eq!(String::from_utf8(output).unwrap(), expected);
405 }
406
407 struct PanicsLater {
409 renders: std::sync::Arc<std::sync::atomic::AtomicUsize>,
410 fail_after: usize,
411 }
412
413 impl Renderable for PanicsLater {
414 fn rich_render(
415 &self,
416 console: &Console,
417 options: &crate::console::ConsoleOptions,
418 ) -> Vec<crate::segment::Segment> {
419 use std::sync::atomic::Ordering;
420 if self.renders.fetch_add(1, Ordering::SeqCst) >= self.fail_after {
421 panic!("render failed");
422 }
423 Text::new("ok").rich_render(console, options)
424 }
425 }
426
427 #[test]
428 fn a_panicking_render_still_restores_the_cursor() {
429 let renders = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
433 let auto = Live::spawn(
434 Box::new(PanicsLater {
435 renders: renders.clone(),
436 fail_after: 1,
437 }),
438 console(),
439 Vec::<u8>::new(),
440 0.1,
441 );
442 auto.refresh_wait(); let (output, error) = match auto.try_stop() {
444 Ok(output) => (output, None),
445 Err(failure) => (failure.writer, Some(failure.message)),
446 };
447 let output = String::from_utf8(output).unwrap();
448 assert!(output.starts_with("\x1b[?25lok"), "{output:?}");
449 assert!(
450 output.ends_with("\x1b[?25h"),
451 "cursor left hidden: {output:?}"
452 );
453 assert_eq!(error.as_deref(), Some("render failed"));
454
455 let auto = Live::spawn(
457 Box::new(PanicsLater {
458 renders: std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)),
459 fail_after: 1,
460 }),
461 console(),
462 Vec::<u8>::new(),
463 0.1,
464 );
465 auto.refresh();
466 let output = String::from_utf8(auto.stop()).unwrap();
467 assert!(
468 output.ends_with("\x1b[?25h"),
469 "cursor left hidden: {output:?}"
470 );
471 }
472
473 #[test]
474 fn dropping_the_handle_finalizes_the_display() {
475 let console = console();
478 let auto = Live::spawn(Box::new(Text::new("only")), console, Vec::<u8>::new(), 0.1);
480 drop(auto); }
485}