summaryrefslogtreecommitdiffstats
path: root/src/client/reply_handler.rs
diff options
context:
space:
mode:
Diffstat (limited to 'src/client/reply_handler.rs')
-rw-r--r--src/client/reply_handler.rs65
1 files changed, 65 insertions, 0 deletions
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<ReplyEvent>,
+ buffer: String,
+ abort: AbortSignal,
+}
+
+impl ReplyHandler {
+ pub fn new(sender: UnboundedSender<ReplyEvent>, 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,
+}