1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
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,
}
|