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/client/reply_handler.rs | 65 +++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 65 insertions(+) create mode 100644 src/client/reply_handler.rs (limited to 'src/client/reply_handler.rs') diff --git a/src/client/reply_handler.rs b/src/client/reply_handler.rs new file mode 100644 index 0000000..e11ea1d --- /dev/null +++ b/src/client/reply_handler.rs @@ -0,0 +1,65 @@ +use crate::utils::AbortSignal; + +use anyhow::{Context, Result}; +use tokio::sync::mpsc::UnboundedSender; + +pub struct ReplyHandler { + sender: UnboundedSender, + buffer: String, + abort: AbortSignal, +} + +impl ReplyHandler { + pub fn new(sender: UnboundedSender, 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 + } +} + +#[derive(Debug)] +pub enum ReplyEvent { + Text(String), + Done, +} -- cgit v1.2.3