use std::sync::Arc;
use serde_json::Value as JsonValue;
use tokio::sync::{broadcast, mpsc, oneshot};
use uuid::Uuid;
use crate::{
chat::message::FuneraMessage,
event_bus::env_state_bus::EnvStateEvent,
middleware::{ErrorsEnabled, EventSenderFn, MiddlewareChain, MiddlewareEvent},
re_act::{ReActLoop, ReActLoopConfig},
};
pub enum SessionCmd {
PushMessages { msgs: Vec<FuneraMessage> },
FetchContext {
respond: oneshot::Sender<Vec<JsonValue>>,
},
GetMessages {
respond: oneshot::Sender<Vec<FuneraMessage>>,
},
Clear,
}
pub fn spawn_session_actor() -> mpsc::UnboundedSender<SessionCmd> {
let (tx, mut rx) = mpsc::unbounded_channel();
tokio::spawn(async move {
let mut msgs: Vec<FuneraMessage> = Vec::new();
while let Some(cmd) = rx.recv().await {
match cmd {
SessionCmd::PushMessages { msgs: new } => msgs.extend(new),
SessionCmd::FetchContext { respond } => {
let ctx: Vec<JsonValue> = msgs.iter().map(|m| m.format_json()).collect();
let _ = respond.send(ctx);
}
SessionCmd::GetMessages { respond } => {
let _ = respond.send(msgs.clone());
}
SessionCmd::Clear => msgs.clear(),
}
}
});
tx
}
pub struct FuneraSession {
id: Uuid,
session_tx: mpsc::UnboundedSender<SessionCmd>,
}
impl FuneraSession {
pub fn new(session_tx: mpsc::UnboundedSender<SessionCmd>) -> Self {
Self {
id: Uuid::new_v4(),
session_tx,
}
}
pub fn id(&self) -> Uuid {
self.id
}
pub fn session_tx(&self) -> mpsc::UnboundedSender<SessionCmd> {
self.session_tx.clone()
}
pub fn push_message(&self, msg: FuneraMessage) {
let _ = self
.session_tx
.send(SessionCmd::PushMessages { msgs: vec![msg] });
}
pub async fn session_context(&self) -> Vec<JsonValue> {
let (respond, rx) = oneshot::channel();
let _ = self.session_tx.send(SessionCmd::FetchContext { respond });
rx.await.unwrap_or_default()
}
pub async fn get_messages(&self) -> Vec<FuneraMessage> {
let (respond, rx) = oneshot::channel();
let _ = self.session_tx.send(SessionCmd::GetMessages { respond });
rx.await.unwrap_or_default()
}
pub async fn react_loop<P: crate::provider::ChatProvider, E: MiddlewareEvent>(
&self,
init_msg: FuneraMessage,
mut config: ReActLoopConfig,
env_state_tx: broadcast::Sender<EnvStateEvent>,
middleware: Option<Arc<MiddlewareChain<E, ErrorsEnabled>>>,
event_sender: Option<EventSenderFn<E>>,
) -> anyhow::Result<()> {
let _ = env_state_tx.send(EnvStateEvent::SessionStart);
self.push_message(init_msg);
config.session_tx = Some(self.session_tx.clone());
let react_loop = ReActLoop::<P>::from_config(config);
let loop_handle = react_loop.run::<E>(middleware, event_sender);
loop_handle.task.await??;
let _ = env_state_tx.send(EnvStateEvent::SessionClosed);
Ok(())
}
}