use rpi_ai::types::{AssistantMessageEvent, ToolResultMessage};
use std::sync::Arc;
use crate::message::AgentMessage;
use crate::types::AgentToolResult;
#[derive(Debug, Clone)]
pub enum AgentEvent {
AgentStart,
AgentEnd { messages: Vec<AgentMessage> },
TurnStart,
TurnEnd {
message: AgentMessage,
tool_results: Vec<ToolResultMessage>,
},
MessageStart { message: AgentMessage },
MessageUpdate {
message: AgentMessage,
assistant_message_event: AssistantMessageEvent,
},
MessageEnd { message: AgentMessage },
ToolExecutionStart {
tool_call_id: String,
tool_name: String,
args: serde_json::Value,
},
ToolExecutionUpdate {
tool_call_id: String,
tool_name: String,
args: serde_json::Value,
partial_result: Arc<AgentToolResult>,
},
ToolExecutionEnd {
tool_call_id: String,
tool_name: String,
result: AgentToolResult,
is_error: bool,
},
}
impl AgentEvent {
pub fn type_tag(&self) -> &'static str {
match self {
AgentEvent::AgentStart => "agent_start",
AgentEvent::AgentEnd { .. } => "agent_end",
AgentEvent::TurnStart => "turn_start",
AgentEvent::TurnEnd { .. } => "turn_end",
AgentEvent::MessageStart { .. } => "message_start",
AgentEvent::MessageUpdate { .. } => "message_update",
AgentEvent::MessageEnd { .. } => "message_end",
AgentEvent::ToolExecutionStart { .. } => "tool_execution_start",
AgentEvent::ToolExecutionUpdate { .. } => "tool_execution_update",
AgentEvent::ToolExecutionEnd { .. } => "tool_execution_end",
}
}
pub fn is_terminal(&self) -> bool {
matches!(self, AgentEvent::AgentEnd { .. })
}
}
pub trait AgentEmitter: Send + Sync {
fn emit(&self, event: AgentEvent) -> futures::future::BoxFuture<'static, ()>;
fn try_emit(&self, _event: AgentEvent) {}
}
pub struct BroadcastEmitter {
tx: tokio::sync::broadcast::Sender<AgentEvent>,
}
impl BroadcastEmitter {
pub fn new(buffer: usize) -> (Self, tokio::sync::broadcast::Receiver<AgentEvent>) {
let (tx, rx) = tokio::sync::broadcast::channel(buffer);
(Self { tx }, rx)
}
pub fn from_sender(tx: tokio::sync::broadcast::Sender<AgentEvent>) -> Self {
Self { tx }
}
pub fn subscribe(&self) -> tokio::sync::broadcast::Receiver<AgentEvent> {
self.tx.subscribe()
}
pub fn try_emit(&self, event: AgentEvent) {
let _ = self.tx.send(event);
}
}
impl AgentEmitter for BroadcastEmitter {
fn emit(&self, event: AgentEvent) -> futures::future::BoxFuture<'static, ()> {
let _ = self.tx.send(event);
Box::pin(async {})
}
fn try_emit(&self, event: AgentEvent) {
let _ = self.tx.send(event);
}
}
pub struct CollectorEmitter {
events: std::sync::Arc<std::sync::Mutex<Vec<AgentEvent>>>,
}
impl CollectorEmitter {
pub fn new() -> (Self, std::sync::Arc<std::sync::Mutex<Vec<AgentEvent>>>) {
let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
(Self { events: Arc::clone(&events) }, events)
}
}
impl Default for CollectorEmitter {
fn default() -> Self {
let (s, _) = Self::new();
s
}
}
impl AgentEmitter for CollectorEmitter {
fn emit(&self, event: AgentEvent) -> futures::future::BoxFuture<'static, ()> {
self.events.lock().expect("events lock").push(event);
Box::pin(async {})
}
fn try_emit(&self, event: AgentEvent) {
self.events.lock().expect("events lock").push(event);
}
}