use std::sync::atomic::{AtomicI32, Ordering};
use async_trait::async_trait;
use everruns_core::error::Result as CoreResult;
use everruns_core::events::{
self, Event, EventData, EventRequest, OutputMessageDeltaData, ToolCompletedData,
ToolStartedData, TurnFailedData,
};
use everruns_core::traits::EventEmitter;
use everruns_core::typed_id::EventId;
use everruns_runtime::EventBus;
use serde_json::Value;
use tokio::sync::broadcast;
const EVENT_CHANNEL_CAPACITY: usize = 4096;
#[derive(Clone, Debug)]
#[non_exhaustive]
pub struct SessionEvent {
pub event_id: String,
pub session_id: String,
pub turn_id: Option<String>,
pub kind: SessionEventKind,
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub enum SessionEventKind {
TurnStarted,
TurnCompleted,
TurnFailed {
error: String,
},
TurnCancelled,
TextDelta {
delta: String,
},
ToolStarted {
tool_call_id: String,
tool_name: String,
},
ToolCompleted {
tool_call_id: String,
tool_name: String,
success: bool,
},
Other {
event_type: String,
payload: Value,
},
}
impl SessionEvent {
pub fn event_type(&self) -> &str {
match &self.kind {
SessionEventKind::TurnStarted => events::TURN_STARTED,
SessionEventKind::TurnCompleted => events::TURN_COMPLETED,
SessionEventKind::TurnFailed { .. } => events::TURN_FAILED,
SessionEventKind::TurnCancelled => events::TURN_CANCELLED,
SessionEventKind::TextDelta { .. } => events::OUTPUT_MESSAGE_DELTA,
SessionEventKind::ToolStarted { .. } => events::TOOL_STARTED,
SessionEventKind::ToolCompleted { .. } => events::TOOL_COMPLETED,
SessionEventKind::Other { event_type, .. } => event_type,
}
}
fn from_core_event(event: &Event) -> Self {
let turn_id = event.context.turn_id.map(|id| id.to_string());
let kind = match event.event_type.as_str() {
events::TURN_STARTED => SessionEventKind::TurnStarted,
events::TURN_COMPLETED => SessionEventKind::TurnCompleted,
events::TURN_CANCELLED => SessionEventKind::TurnCancelled,
events::TURN_FAILED => {
let error = match &event.data {
EventData::TurnFailed(TurnFailedData { error, .. }) => error.clone(),
_ => String::new(),
};
SessionEventKind::TurnFailed { error }
}
events::OUTPUT_MESSAGE_DELTA => {
let delta = match &event.data {
EventData::OutputMessageDelta(OutputMessageDeltaData { delta, .. }) => {
delta.clone()
}
_ => String::new(),
};
SessionEventKind::TextDelta { delta }
}
events::TOOL_STARTED => match &event.data {
EventData::ToolStarted(ToolStartedData { tool_call, .. }) => {
SessionEventKind::ToolStarted {
tool_call_id: tool_call.id.clone(),
tool_name: tool_call.name.clone(),
}
}
_ => Self::other_kind(event),
},
events::TOOL_COMPLETED => match &event.data {
EventData::ToolCompleted(ToolCompletedData {
tool_call_id,
tool_name,
success,
..
}) => SessionEventKind::ToolCompleted {
tool_call_id: tool_call_id.clone(),
tool_name: tool_name.clone(),
success: *success,
},
_ => Self::other_kind(event),
},
_ => Self::other_kind(event),
};
Self {
event_id: event.id.to_string(),
session_id: event.session_id.to_string(),
turn_id,
kind,
}
}
fn other_kind(event: &Event) -> SessionEventKind {
SessionEventKind::Other {
event_type: event.event_type.clone(),
payload: serde_json::to_value(&event.data).unwrap_or(Value::Null),
}
}
}
pub struct EventStream {
rx: broadcast::Receiver<SessionEvent>,
}
impl EventStream {
fn new(rx: broadcast::Receiver<SessionEvent>) -> Self {
Self { rx }
}
pub async fn recv(&mut self) -> Option<SessionEvent> {
loop {
match self.rx.recv().await {
Ok(event) => return Some(event),
Err(broadcast::error::RecvError::Lagged(_)) => continue,
Err(broadcast::error::RecvError::Closed) => return None,
}
}
}
pub fn try_recv(&mut self) -> Option<SessionEvent> {
loop {
match self.rx.try_recv() {
Ok(event) => return Some(event),
Err(broadcast::error::TryRecvError::Lagged(_)) => continue,
Err(broadcast::error::TryRecvError::Empty)
| Err(broadcast::error::TryRecvError::Closed) => return None,
}
}
}
}
#[derive(Clone, Default)]
pub struct RunOptions {
pub(crate) cancel: Option<CancellationToken>,
}
impl RunOptions {
pub fn new() -> Self {
Self::default()
}
pub fn cancel_token(mut self, token: CancellationToken) -> Self {
self.cancel = Some(token);
self
}
}
#[derive(Clone, Default)]
pub struct CancellationToken {
inner: tokio_util::sync::CancellationToken,
}
impl CancellationToken {
pub fn new() -> Self {
Self::default()
}
pub fn cancel(&self) {
self.inner.cancel();
}
pub fn is_cancelled(&self) -> bool {
self.inner.is_cancelled()
}
pub(crate) async fn cancelled(&self) {
self.inner.cancelled().await;
}
}
pub(crate) struct FacadeEventBus {
sender: broadcast::Sender<SessionEvent>,
sequence: AtomicI32,
}
impl FacadeEventBus {
pub(crate) fn new() -> Self {
let (sender, _rx) = broadcast::channel(EVENT_CHANNEL_CAPACITY);
Self {
sender,
sequence: AtomicI32::new(0),
}
}
pub(crate) fn subscribe(&self) -> EventStream {
EventStream::new(self.sender.subscribe())
}
}
#[async_trait]
impl EventEmitter for FacadeEventBus {
async fn emit(&self, request: EventRequest) -> CoreResult<Event> {
let seq = self.sequence.fetch_add(1, Ordering::Relaxed) + 1;
let event = request.into_event(EventId::new(), seq);
let _ = self.sender.send(SessionEvent::from_core_event(&event));
Ok(event)
}
}
impl EventBus for FacadeEventBus {}