From 1cc89eff514669d6459f62cc4eed8e0b34d8ef0c Mon Sep 17 00:00:00 2001 From: sigoden Date: Tue, 23 Apr 2024 14:32:06 +0800 Subject: refactor: more async code (#427) --- src/render/mod.rs | 119 ++++++------------------------------------------------ 1 file changed, 12 insertions(+), 107 deletions(-) (limited to 'src/render/mod.rs') diff --git a/src/render/mod.rs b/src/render/mod.rs index dc4c081..146577b 100644 --- a/src/render/mod.rs +++ b/src/render/mod.rs @@ -4,61 +4,26 @@ mod stream; pub use self::markdown::{MarkdownRender, RenderOptions}; use self::stream::{markdown_stream, raw_stream}; -use crate::client::Client; -use crate::config::{GlobalConfig, Input}; use crate::utils::AbortSignal; +use crate::{client::ReplyEvent, config::GlobalConfig}; -use anyhow::{Context, Result}; -use crossbeam::channel::{unbounded, Sender}; -use crossbeam::sync::WaitGroup; +use anyhow::Result; use is_terminal::IsTerminal; use nu_ansi_term::{Color, Style}; use std::io::stdout; -use std::thread::spawn; +use tokio::sync::mpsc::UnboundedReceiver; -pub fn render_stream( - input: &Input, - client: &dyn Client, +pub async fn render_stream( + rx: UnboundedReceiver, config: &GlobalConfig, abort: AbortSignal, -) -> Result { - let wg = WaitGroup::new(); - let wg_cloned = wg.clone(); - let render_options = config.read().get_render_options()?; - let mut stream_handler = { - let (tx, rx) = unbounded(); - let abort_clone = abort.clone(); - let highlight = config.read().highlight; - spawn(move || { - let run = move || { - if stdout().is_terminal() { - let mut render = MarkdownRender::init(render_options)?; - markdown_stream(&rx, &mut render, &abort) - } else { - raw_stream(&rx, &abort) - } - }; - if let Err(err) = run() { - render_error(err, highlight); - } - drop(wg_cloned); - }); - ReplyHandler::new(tx, abort_clone) - }; - let ret = client.send_message_streaming(input, &mut stream_handler); - wg.wait(); - let output = stream_handler.get_buffer().to_string(); - match ret { - Ok(_) => { - println!(); - Ok(output) - } - Err(err) => { - if !output.is_empty() { - println!(); - } - Err(err) - } +) -> Result<()> { + if stdout().is_terminal() { + let render_options = config.read().get_render_options()?; + let mut render = MarkdownRender::init(render_options)?; + markdown_stream(rx, &mut render, &abort).await + } else { + raw_stream(rx, &abort).await } } @@ -71,63 +36,3 @@ pub fn render_error(err: anyhow::Error, highlight: bool) { eprintln!("{err}"); } } - -pub struct ReplyHandler { - sender: Sender, - buffer: String, - abort: AbortSignal, -} - -impl ReplyHandler { - pub fn new(sender: Sender, abort: AbortSignal) -> Self { - Self { - sender, - abort, - buffer: String::new(), - } - } - - pub fn text(&mut self, text: &str) -> Result<()> { - debug!("ReplyText: {}", text); - if text.is_empty() { - return Ok(()); - } - self.buffer.push_str(text); - let ret = self - .sender - .send(ReplyEvent::Text(text.to_string())) - .with_context(|| "Failed to send ReplyEvent:Text"); - self.safe_ret(ret)?; - Ok(()) - } - - pub fn done(&mut self) -> Result<()> { - debug!("ReplyDone"); - let ret = self - .sender - .send(ReplyEvent::Done) - .with_context(|| "Failed to send ReplyEvent::Done"); - self.safe_ret(ret)?; - Ok(()) - } - - pub fn get_buffer(&self) -> &str { - &self.buffer - } - - pub fn get_abort(&self) -> AbortSignal { - self.abort.clone() - } - - fn safe_ret(&self, ret: Result<()>) -> Result<()> { - if ret.is_err() && self.abort.aborted() { - return Ok(()); - } - ret - } -} - -pub enum ReplyEvent { - Text(String), - Done, -} -- cgit v1.2.3