use std::{sync::Arc, time::Duration};
use tokio::{sync::broadcast, time::interval};
use crate::{contains_ansi_escape_sequence,
get_terminal_width,
is_fully_uninteractive_terminal,
is_stdout_piped,
ok,
spinner_print,
spinner_render,
InlineString,
LineStateControlSignal,
OutputDevice,
SafeBool,
SharedWriter,
SpinnerStyle,
StdMutex,
StdoutIsPipedResult,
TTYResult};
pub struct Spinner {
pub tick_delay: Duration,
pub interval_message: InlineString,
pub final_message: InlineString,
pub style: SpinnerStyle,
pub output_device: OutputDevice,
pub maybe_shared_writer: Option<SharedWriter>,
pub shutdown_sender: broadcast::Sender<()>,
safe_is_shutdown: SafeBool,
maybe_shutdown_complete_rx: Option<tokio::sync::oneshot::Receiver<()>>,
}
impl Spinner {
pub async fn try_start(
arg_interval_msg: impl AsRef<str>,
arg_final_msg: impl AsRef<str>,
tick_delay: Duration,
style: SpinnerStyle,
output_device: OutputDevice,
maybe_shared_writer: Option<SharedWriter>,
) -> miette::Result<Option<Spinner>> {
if let StdoutIsPipedResult::StdoutIsPiped = is_stdout_piped() {
return Ok(None);
}
if let TTYResult::IsNotInteractive = is_fully_uninteractive_terminal() {
return Ok(None);
}
let interval_msg = {
let msg = arg_interval_msg.as_ref();
if contains_ansi_escape_sequence(msg) {
strip_ansi_escapes::strip_str(msg)
} else {
msg.to_string()
}
};
let final_msg = {
let msg = arg_final_msg.as_ref();
if contains_ansi_escape_sequence(msg) {
strip_ansi_escapes::strip_str(msg)
} else {
msg.to_string()
}
};
let (shutdown_sender, _) = broadcast::channel::<()>(1);
let mut spinner = Spinner {
interval_message: interval_msg.into(),
final_message: final_msg.into(),
tick_delay,
style,
output_device,
maybe_shared_writer,
shutdown_sender,
safe_is_shutdown: Arc::new(StdMutex::new(false)),
maybe_shutdown_complete_rx: None,
};
spinner.try_start_task().await?;
Ok(Some(spinner))
}
pub fn is_shutdown(&self) -> bool { *self.safe_is_shutdown.lock().unwrap() }
pub async fn try_start_task(&mut self) -> miette::Result<()> {
if let Some(shared_writer) = self.maybe_shared_writer.as_ref() {
_ = shared_writer
.line_state_control_channel_sender
.send(LineStateControlSignal::SpinnerActive(
self.shutdown_sender.clone(),
))
.await;
_ = shared_writer
.line_state_control_channel_sender
.send(LineStateControlSignal::Pause)
.await;
};
let mut shutdown_receiver = self.shutdown_sender.subscribe();
let self_safe_is_shutdown = self.safe_is_shutdown.clone();
spinner_print::print_start_if_standalone(
self.output_device.clone(),
self.maybe_shared_writer.clone(),
)?;
let (shutdown_complete_sender, shutdown_complete_receiver) =
tokio::sync::oneshot::channel::<()>();
self.maybe_shutdown_complete_rx = Some(shutdown_complete_receiver);
let output_device_clone = self.output_device.clone();
let interval_message_clone = self.interval_message.clone();
let final_message_clone = self.final_message.clone();
let maybe_shared_writer_clone = self.maybe_shared_writer.clone();
let mut style_clone = self.style.clone();
let tick_delay_clone = self.tick_delay;
tokio::spawn(async move {
let mut interval = interval(tick_delay_clone);
let mut count = 0;
loop {
tokio::select! {
_ = shutdown_receiver.recv() => {
drop(interval);
if let Some(shared_writer) = maybe_shared_writer_clone.as_ref() {
_ = shared_writer
.line_state_control_channel_sender
.send(LineStateControlSignal::SpinnerInactive)
.await;
}
let final_output = spinner_render::render_final_tick(
&style_clone,
&final_message_clone,
get_terminal_width(),
);
_ = spinner_print::print_tick_final_msg(
&style_clone,
&final_output,
output_device_clone.clone(),
maybe_shared_writer_clone.clone(),
);
if let Some(shared_writer) = maybe_shared_writer_clone.as_ref() {
let _ = shared_writer
.line_state_control_channel_sender
.send(LineStateControlSignal::Resume)
.await;
}
*self_safe_is_shutdown.lock().unwrap() = true;
let _ = shutdown_complete_sender.send(());
break;
}
_ = interval.tick() => {
if *self_safe_is_shutdown.lock().unwrap() {
break;
}
let output = spinner_render::render_tick(
&mut style_clone,
&interval_message_clone,
count,
get_terminal_width(),
);
_ = spinner_print::print_tick_interval_msg(
&style_clone,
&output,
output_device_clone.clone()
);
count += 1;
},
}
}
});
ok!()
}
pub async fn request_shutdown(&mut self) { _ = self.shutdown_sender.send(()); }
pub async fn await_shutdown(mut self) {
if let Some(receiver) = self.maybe_shutdown_complete_rx.take() {
let _ = receiver.await;
}
}
}
#[cfg(test)]
mod tests {
use smallvec::SmallVec;
use super::{Duration,
LineStateControlSignal,
SharedWriter,
Spinner,
SpinnerStyle,
TTYResult};
use crate::{return_if_not_interactive_terminal,
OutputDevice,
OutputDeviceExt,
SpinnerColor,
SpinnerTemplate};
type ArrayVec = SmallVec<[LineStateControlSignal; FACTOR as usize]>;
const FACTOR: u32 = 5;
const QUANTUM: Duration = Duration::from_millis(100);
#[serial_test::serial]
#[tokio::test]
#[allow(clippy::needless_return)]
async fn test_spinner_color() {
let (output_device_mock, stdout_mock) = OutputDevice::new_mock();
let (line_sender, mut line_receiver) = tokio::sync::mpsc::channel(1_000);
let shared_writer = SharedWriter::new(line_sender);
let res_maybe_spinner = Spinner::try_start(
"message",
"final message",
QUANTUM,
SpinnerStyle {
template: SpinnerTemplate::Braille,
color: SpinnerColor::None,
},
output_device_mock,
Some(shared_writer),
)
.await;
return_if_not_interactive_terminal!();
let mut spinner = res_maybe_spinner.unwrap().unwrap();
tokio::time::sleep(QUANTUM * FACTOR).await;
spinner.request_shutdown().await;
spinner.await_shutdown().await;
let output_buffer_data = stdout_mock.get_copy_of_buffer_as_string_strip_ansi();
println!("{output_buffer_data:?}");
assert!(output_buffer_data.contains("⠁ message\n"));
assert!(output_buffer_data.contains("final message\n"));
let line_control_signal_sink = {
let mut acc = ArrayVec::new();
loop {
let it = line_receiver.try_recv();
match it {
Ok(signal) => {
acc.push(signal);
}
Err(error) => match error {
tokio::sync::mpsc::error::TryRecvError::Empty => {
break;
}
tokio::sync::mpsc::error::TryRecvError::Disconnected => {
break;
}
},
}
}
acc
};
assert_eq!(line_control_signal_sink.len(), 4);
matches!(
line_control_signal_sink[0],
LineStateControlSignal::SpinnerActive(_)
);
matches!(line_control_signal_sink[1], LineStateControlSignal::Pause);
matches!(
line_control_signal_sink[2],
LineStateControlSignal::SpinnerInactive
);
matches!(line_control_signal_sink[3], LineStateControlSignal::Resume);
drop(line_receiver);
}
#[serial_test::serial]
#[tokio::test]
#[allow(clippy::needless_return)]
async fn test_spinner_no_color() {
let (output_device_mock, stdout_mock) = OutputDevice::new_mock();
let (line_sender, mut line_receiver) = tokio::sync::mpsc::channel(1_000);
let shared_writer = SharedWriter::new(line_sender);
let res_maybe_spinner = Spinner::try_start(
"message",
"final message",
QUANTUM,
SpinnerStyle::default(),
output_device_mock,
Some(shared_writer),
)
.await;
return_if_not_interactive_terminal!();
let mut spinner = res_maybe_spinner.unwrap().unwrap();
tokio::time::sleep(QUANTUM * FACTOR).await;
spinner.request_shutdown().await;
spinner.await_shutdown().await;
let output_buffer_data = stdout_mock.get_copy_of_buffer_as_string();
println!("{output_buffer_data:?}");
let stripped_output_buffer_data =
strip_ansi_escapes::strip_str(&output_buffer_data);
assert!(stripped_output_buffer_data.contains("⠁ message\n"));
assert!(stripped_output_buffer_data.contains("final message\n"));
let line_control_signal_sink = {
let mut acc = ArrayVec::new();
loop {
let it = line_receiver.try_recv();
match it {
Ok(signal) => {
acc.push(signal);
}
Err(error) => match error {
tokio::sync::mpsc::error::TryRecvError::Empty => {
break;
}
tokio::sync::mpsc::error::TryRecvError::Disconnected => {
break;
}
},
}
}
acc
};
assert_eq!(line_control_signal_sink.len(), 4);
matches!(
line_control_signal_sink[0],
LineStateControlSignal::SpinnerActive(_)
);
matches!(line_control_signal_sink[1], LineStateControlSignal::Pause);
matches!(
line_control_signal_sink[2],
LineStateControlSignal::SpinnerInactive
);
matches!(line_control_signal_sink[3], LineStateControlSignal::Resume);
drop(line_receiver);
}
}