summaryrefslogtreecommitdiffstats
path: root/src/client/sse_handler.rs
diff options
context:
space:
mode:
authorsigoden <sigoden@gmail.com>2024-05-23 19:28:56 +0800
committerGitHub <noreply@github.com>2024-05-23 19:28:56 +0800
commit5458150ed3203cf13b0371efa2c791ac696cee93 (patch)
tree702d17039d9677246a3b91deb8243e09cd768402 /src/client/sse_handler.rs
parent2ccbb0f06a4558e15642feb53ba7b2bd72804820 (diff)
downloadaichat-5458150ed3203cf13b0371efa2c791ac696cee93.tar.gz
fix: json stream parser and refine client modules (#538)
Diffstat (limited to 'src/client/sse_handler.rs')
-rw-r--r--src/client/sse_handler.rs78
1 files changed, 0 insertions, 78 deletions
diff --git a/src/client/sse_handler.rs b/src/client/sse_handler.rs
deleted file mode 100644
index ddbdbcd..0000000
--- a/src/client/sse_handler.rs
+++ /dev/null
@@ -1,78 +0,0 @@
-use crate::utils::AbortSignal;
-
-use anyhow::{Context, Result};
-use tokio::sync::mpsc::UnboundedSender;
-
-use super::ToolCall;
-
-pub struct SseHandler {
- sender: UnboundedSender<SseEvent>,
- abort: AbortSignal,
- buffer: String,
- tool_calls: Vec<ToolCall>,
-}
-
-impl SseHandler {
- pub fn new(sender: UnboundedSender<SseEvent>, abort: AbortSignal) -> Self {
- Self {
- sender,
- abort,
- buffer: String::new(),
- tool_calls: Vec::new(),
- }
- }
-
- pub fn text(&mut self, text: &str) -> Result<()> {
- // debug!("HandleText: {}", text);
- if text.is_empty() {
- return Ok(());
- }
- self.buffer.push_str(text);
- let ret = self
- .sender
- .send(SseEvent::Text(text.to_string()))
- .with_context(|| "Failed to send ReplyEvent:Text");
- self.safe_ret(ret)?;
- Ok(())
- }
-
- pub fn done(&mut self) -> Result<()> {
- // debug!("HandleDone");
- let ret = self
- .sender
- .send(SseEvent::Done)
- .with_context(|| "Failed to send ReplyEvent::Done");
- self.safe_ret(ret)?;
- Ok(())
- }
-
- pub fn tool_call(&mut self, call: ToolCall) -> Result<()> {
- // debug!("HandleCall: {:?}", call);
- self.tool_calls.push(call);
- Ok(())
- }
-
- pub fn get_abort(&self) -> AbortSignal {
- self.abort.clone()
- }
-
- pub fn take(self) -> (String, Vec<ToolCall>) {
- let Self {
- buffer, tool_calls, ..
- } = self;
- (buffer, tool_calls)
- }
-
- fn safe_ret(&self, ret: Result<()>) -> Result<()> {
- if ret.is_err() && self.abort.aborted() {
- return Ok(());
- }
- ret
- }
-}
-
-#[derive(Debug)]
-pub enum SseEvent {
- Text(String),
- Done,
-}