use futures_util::FutureExt;
use miette::IntoDiagnostic;
use tokio::sync::broadcast;
use crate::{inline_string, is_fully_uninteractive_terminal, is_stdin_piped,
is_stdout_piped, ok, CommonResult, InputDevice, LineStateControlSignal,
OutputDevice, Readline, ReadlineEvent, SharedWriter, StdinIsPipedResult,
StdoutIsPipedResult, TTYResult,
READLINE_ASYNC_INITIAL_PROMPT_DISPLAY_CURSOR_SHOW_DELAY};
#[allow(missing_debug_implementations)]
pub struct ReadlineAsyncContext {
pub readline: Readline,
pub shared_writer: SharedWriter,
pub shutdown_complete_sender: broadcast::Sender<()>,
}
#[macro_export]
macro_rules! rla_println {
(
$rla:ident,
$($format:tt)*
) => {{
use std::io::Write;
writeln!($rla.shared_writer, $($format)*).ok();
}};
}
#[macro_export]
macro_rules! rla_print {
(
$rla:ident,
$($format:tt)*
) => {{
use std::io::Write;
write!($rla.shared_writer, $($format)*).ok();
}};
}
#[macro_export]
macro_rules! rla_println_prefixed {
(
$rla:ident,
$($format:tt)*
) => {{
use std::io::Write;
use $crate::fg_pink;
write!($rla.shared_writer, "{}", fg_pink(" > ").bold().bg_moonlight_blue()).ok();
writeln!($rla.shared_writer, $($format)*).ok();
}};
}
impl ReadlineAsyncContext {
pub async fn try_new(
read_line_prompt: Option<impl AsRef<str>>,
) -> miette::Result<Option<ReadlineAsyncContext>> {
if let StdinIsPipedResult::StdinIsPiped = is_stdin_piped() {
return Ok(None);
}
if let StdoutIsPipedResult::StdoutIsPiped = is_stdout_piped() {
return Ok(None);
}
if let TTYResult::IsNotInteractive = is_fully_uninteractive_terminal() {
return Ok(None);
}
let output_device = OutputDevice::new_stdout();
let input_device = InputDevice::new_event_stream();
let prompt =
read_line_prompt.map_or_else(|| "> ".to_owned(), |p| p.as_ref().to_string());
let shutdown_complete_channel = broadcast::channel::<()>(1);
let (shutdown_complete_sender, _) = shutdown_complete_channel;
let (readline, stdout) = Readline::try_new(
prompt.clone(),
output_device,
input_device,
shutdown_complete_sender.clone(),
)
.into_diagnostic()?;
tokio::time::sleep(READLINE_ASYNC_INITIAL_PROMPT_DISPLAY_CURSOR_SHOW_DELAY).await;
Ok(Some(ReadlineAsyncContext {
readline,
shared_writer: stdout,
shutdown_complete_sender,
}))
}
#[must_use]
pub fn clone_shared_writer(&self) -> SharedWriter { self.shared_writer.clone() }
pub fn mut_input_device(&mut self) -> &mut InputDevice {
&mut self.readline.input_device
}
pub fn clone_output_device(&mut self) -> OutputDevice {
self.readline.output_device.clone()
}
pub async fn read_line(&mut self) -> miette::Result<ReadlineEvent> {
self.readline.readline().fuse().await.into_diagnostic()
}
pub async fn flush(&mut self) {
self.shared_writer
.line_state_control_channel_sender
.send(LineStateControlSignal::Flush)
.await
.ok();
}
pub async fn pause(&mut self) {
self.shared_writer
.line_state_control_channel_sender
.send(LineStateControlSignal::Pause)
.await
.ok();
}
pub async fn resume(&mut self) {
self.shared_writer
.line_state_control_channel_sender
.send(LineStateControlSignal::Resume)
.await
.ok();
}
pub async fn request_shutdown(&self, message: Option<&str>) -> CommonResult<()> {
if let Some(message) = message {
let message = inline_string!("\r{message}");
self.shared_writer
.line_state_control_channel_sender
.send(LineStateControlSignal::Line(message))
.await
.map_err(std::io::Error::other)
.into_diagnostic()?;
self.shared_writer
.line_state_control_channel_sender
.send(LineStateControlSignal::Flush)
.await
.map_err(std::io::Error::other)
.into_diagnostic()?;
}
self.shared_writer
.line_state_control_channel_sender
.send(LineStateControlSignal::ExitReadlineLoop)
.await
.map_err(std::io::Error::other)
.into_diagnostic()?;
ok!()
}
pub async fn await_shutdown(self) {
let mut shutdown_complete_receiver = self.shutdown_complete_sender.subscribe();
shutdown_complete_receiver.recv().await.ok();
}
}