use std::sync::Mutex;
use everruns_core::events::{
self, Event, EventContext, EventData, EventRequest, InputMessageData,
OutputMessageCompletedData, OutputMessageDeltaData, OutputMessageReplacedData,
OutputMessageStartedData, ReasonCompletedData, ToolCompletedData, ToolOutputDeltaData,
ToolProgressData, ToolStartedData, TurnCancelledData, TurnFailedData,
};
use everruns_core::typed_id::{SessionId, TurnId};
use everruns_host::{EventSink, EventSinkError};
use serde_json::Value;
use tokio::sync::broadcast;
pub const EVENT_STREAM_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,
raw: Value,
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub enum SessionEventKind {
InputMessage {
message_id: String,
},
OutputStarted {
message_id: String,
},
TurnStarted,
TurnCompleted,
TurnFailed {
error: String,
},
TurnCancelled,
TextDelta {
delta: String,
},
OutputReplaced {
message_id: String,
replacement: String,
},
OutputCompleted {
message_id: String,
},
ToolStarted {
tool_call_id: String,
tool_name: String,
},
ToolCompleted {
tool_call_id: String,
tool_name: String,
success: bool,
},
ToolProgress {
tool_call_id: String,
tool_name: String,
message: String,
},
ToolOutputDelta {
tool_call_id: String,
tool_name: String,
stream: String,
delta: String,
},
ReasonStarted,
ReasonCompleted {
success: bool,
error: Option<String>,
},
ModelGeneration,
Other {
event_type: String,
payload: Value,
},
}
impl SessionEventKind {
pub fn is_terminal(&self) -> bool {
matches!(
self,
Self::TurnCompleted | Self::TurnFailed { .. } | Self::TurnCancelled
)
}
}
impl SessionEvent {
pub fn event_type(&self) -> &str {
self.raw
.get("type")
.and_then(Value::as_str)
.expect("canonical events always carry a string type")
}
pub fn sequence(&self) -> Option<i32> {
self.raw
.get("sequence")
.and_then(Value::as_i64)
.and_then(|value| i32::try_from(value).ok())
}
pub fn timestamp(&self) -> &str {
self.raw
.get("ts")
.and_then(Value::as_str)
.expect("canonical events always carry a string timestamp")
}
pub fn as_json(&self) -> &Value {
&self.raw
}
pub fn into_json(self) -> Value {
self.raw
}
pub fn data(&self) -> &Value {
self.raw
.get("data")
.expect("canonical events always carry data")
}
pub fn narration(&self) -> Option<&str> {
self.data().get("narration").and_then(Value::as_str)
}
fn from_core_event(event: &Event) -> Self {
let raw = serde_json::to_value(event).expect("canonical events are JSON serializable");
let turn_id = event.context.turn_id.map(|id| id.to_string());
let kind = match event.event_type.as_str() {
events::INPUT_MESSAGE => match &event.data {
EventData::InputMessage(InputMessageData { message }) => {
SessionEventKind::InputMessage {
message_id: message.id.to_string(),
}
}
_ => Self::other_kind(event),
},
events::OUTPUT_MESSAGE_STARTED => match &event.data {
EventData::OutputMessageStarted(OutputMessageStartedData {
message_id, ..
}) => SessionEventKind::OutputStarted {
message_id: message_id.to_string(),
},
_ => Self::other_kind(event),
},
events::TURN_STARTED => SessionEventKind::TurnStarted,
events::TURN_COMPLETED => SessionEventKind::TurnCompleted,
events::TURN_CANCELLED => SessionEventKind::TurnCancelled,
events::TURN_FAILED => match &event.data {
EventData::TurnFailed(TurnFailedData { error, .. }) => {
SessionEventKind::TurnFailed {
error: error.clone(),
}
}
_ => Self::other_kind(event),
},
events::OUTPUT_MESSAGE_DELTA => match &event.data {
EventData::OutputMessageDelta(OutputMessageDeltaData { delta, .. }) => {
SessionEventKind::TextDelta {
delta: delta.clone(),
}
}
_ => Self::other_kind(event),
},
events::OUTPUT_MESSAGE_REPLACED => match &event.data {
EventData::OutputMessageReplaced(OutputMessageReplacedData {
message_id,
replacement,
..
}) => SessionEventKind::OutputReplaced {
message_id: message_id.to_string(),
replacement: replacement.clone(),
},
_ => Self::other_kind(event),
},
events::OUTPUT_MESSAGE_COMPLETED => match &event.data {
EventData::OutputMessageCompleted(OutputMessageCompletedData {
message, ..
}) => SessionEventKind::OutputCompleted {
message_id: message.id.to_string(),
},
_ => Self::other_kind(event),
},
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),
},
events::TOOL_PROGRESS => match &event.data {
EventData::ToolProgress(ToolProgressData {
tool_call_id,
tool_name,
message,
..
}) => SessionEventKind::ToolProgress {
tool_call_id: tool_call_id.clone(),
tool_name: tool_name.clone(),
message: message.clone(),
},
_ => Self::other_kind(event),
},
events::TOOL_OUTPUT_DELTA => match &event.data {
EventData::ToolOutputDelta(ToolOutputDeltaData {
tool_call_id,
tool_name,
stream,
delta,
}) => SessionEventKind::ToolOutputDelta {
tool_call_id: tool_call_id.clone(),
tool_name: tool_name.clone(),
stream: stream.clone(),
delta: delta.clone(),
},
_ => Self::other_kind(event),
},
events::REASON_STARTED => SessionEventKind::ReasonStarted,
events::REASON_COMPLETED => match &event.data {
EventData::ReasonCompleted(ReasonCompletedData { success, error, .. }) => {
SessionEventKind::ReasonCompleted {
success: *success,
error: error.clone(),
}
}
_ => Self::other_kind(event),
},
events::LLM_GENERATION => SessionEventKind::ModelGeneration,
_ => Self::other_kind(event),
};
Self {
event_id: event.id.to_string(),
session_id: event.session_id.to_string(),
turn_id,
kind,
raw,
}
}
fn other_kind(event: &Event) -> SessionEventKind {
SessionEventKind::Other {
event_type: event.event_type.clone(),
payload: serde_json::to_value(&event.data)
.expect("canonical event data is JSON serializable"),
}
}
}
pub struct EventStream {
rx: broadcast::Receiver<SessionEvent>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum EventStreamError {
Lagged {
missed: u64,
},
}
impl std::fmt::Display for EventStreamError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Lagged { missed } => write!(f, "event stream lagged by {missed} events"),
}
}
}
impl std::error::Error for EventStreamError {}
impl EventStream {
fn new(rx: broadcast::Receiver<SessionEvent>) -> Self {
Self { rx }
}
pub async fn recv(&mut self) -> Result<Option<SessionEvent>, EventStreamError> {
match self.rx.recv().await {
Ok(event) => Ok(Some(event)),
Err(broadcast::error::RecvError::Lagged(missed)) => {
Err(EventStreamError::Lagged { missed })
}
Err(broadcast::error::RecvError::Closed) => Ok(None),
}
}
pub fn try_recv(&mut self) -> Result<Option<SessionEvent>, EventStreamError> {
match self.rx.try_recv() {
Ok(event) => Ok(Some(event)),
Err(broadcast::error::TryRecvError::Lagged(missed)) => {
Err(EventStreamError::Lagged { missed })
}
Err(broadcast::error::TryRecvError::Empty)
| Err(broadcast::error::TryRecvError::Closed) => Ok(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>,
active_turn: Mutex<Option<EventContext>>,
}
impl FacadeEventBus {
pub(crate) fn new() -> Self {
Self::with_capacity(EVENT_STREAM_CAPACITY)
}
fn with_capacity(capacity: usize) -> Self {
let (sender, _rx) = broadcast::channel(capacity);
Self {
sender,
active_turn: Mutex::new(None),
}
}
pub(crate) fn subscribe(&self) -> EventStream {
EventStream::new(self.sender.subscribe())
}
#[cfg(test)]
pub(crate) fn cancellation_request(&self, session_id: SessionId) -> (TurnId, EventRequest) {
self.cancellation_request_for_turn(session_id, TurnId::new())
}
pub(crate) fn cancellation_request_for_turn(
&self,
session_id: SessionId,
fallback_turn_id: TurnId,
) -> (TurnId, EventRequest) {
let context = self
.active_turn
.lock()
.expect("active-turn lock poisoned")
.take()
.unwrap_or_else(|| EventContext {
turn_id: Some(fallback_turn_id),
..EventContext::default()
});
let turn_id = context.turn_id.unwrap_or(fallback_turn_id);
let request = EventRequest::new(
session_id,
EventContext {
turn_id: Some(turn_id),
..context
},
TurnCancelledData {
turn_id,
reason: Some("cancelled by application".to_string()),
usage: None,
},
);
(turn_id, request)
}
fn observe(&self, event: &Event) -> Result<(), EventSinkError> {
match event.event_type.as_str() {
events::TURN_STARTED => {
*self.active_turn.lock().expect("active-turn lock poisoned") =
Some(event.context.clone());
}
events::TURN_COMPLETED
| events::TURN_FAILED
| events::TURN_CANCELLED
| events::TURN_SEALED => {
self.active_turn
.lock()
.expect("active-turn lock poisoned")
.take();
}
_ => {}
}
let _ = self.sender.send(SessionEvent::from_core_event(event));
Ok(())
}
}
impl EventSink for FacadeEventBus {
fn try_send(&self, event: Event) -> Result<(), EventSinkError> {
self.observe(&event)
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use everruns_core::events::{
ActStartedData, OutputMessageDeltaData, OutputMessageReplacedData, ToolStartedData,
TurnCancelledData, TurnStartedData,
};
use everruns_core::traits::EventEmitter;
use everruns_core::{MessageId, SessionId, ToolCall, TurnId};
use everruns_host::{HostEventEmitter, InMemoryEventLog};
use serde_json::json;
use super::{EventStreamError, FacadeEventBus, SessionEventKind};
use crate::{Agent, Model};
fn host(bus: Arc<FacadeEventBus>) -> HostEventEmitter {
HostEventEmitter::new(Arc::new(InMemoryEventLog::new()), bus)
}
fn turn_started(session_id: SessionId, turn_id: TurnId) -> everruns_core::EventRequest {
let input_message_id = MessageId::new();
everruns_core::EventRequest::new(
session_id,
everruns_core::EventContext::turn(turn_id, input_message_id),
TurnStartedData {
turn_id,
input_message_id,
input_content: Some("hello".to_string()),
},
)
}
#[tokio::test]
async fn public_event_json_is_the_exact_canonical_envelope() {
let session_id = SessionId::new();
let turn_id = TurnId::new();
let bus = Arc::new(FacadeEventBus::new());
let emitter = host(bus.clone());
let mut stream = bus.subscribe();
let request = turn_started(session_id, turn_id)
.with_metadata(json!({"provider": "simulated"}))
.with_tags(vec!["terminal".to_string()]);
let canonical = emitter.emit(request).await.expect("event emits");
let observed = stream
.recv()
.await
.expect("stream remains lossless")
.expect("event is delivered");
assert!(matches!(observed.kind, SessionEventKind::TurnStarted));
assert_eq!(
observed.as_json(),
&serde_json::to_value(canonical).expect("canonical event serializes")
);
assert_eq!(observed.event_type(), "turn.started");
assert_eq!(
observed.turn_id.as_deref(),
Some(turn_id.to_string().as_str())
);
assert_eq!(observed.as_json()["sequence"], 1);
assert_eq!(observed.sequence(), Some(1));
assert!(!observed.timestamp().is_empty());
assert_eq!(observed.data()["input_content"], "hello");
assert_eq!(observed.as_json()["metadata"]["provider"], "simulated");
assert_eq!(observed.as_json()["tags"], json!(["terminal"]));
}
#[tokio::test]
async fn bounded_stream_reports_lag_instead_of_hiding_loss() {
let session_id = SessionId::new();
let bus = Arc::new(FacadeEventBus::with_capacity(2));
let emitter = host(bus.clone());
let mut stream = bus.subscribe();
for _ in 0..3 {
emitter
.emit(turn_started(session_id, TurnId::new()))
.await
.expect("event emits without observer backpressure");
}
assert!(matches!(
stream.recv().await,
Err(EventStreamError::Lagged { missed: 1 })
));
assert!(stream.recv().await.expect("gap reported").is_some());
}
#[tokio::test]
async fn no_subscriber_is_a_noop_not_a_closed_sink_failure() {
let bus = Arc::new(FacadeEventBus::new());
let emitter = host(bus);
emitter
.emit(turn_started(SessionId::new(), TurnId::new()))
.await
.expect("observation absence cannot reverse the append");
assert_eq!(emitter.delivery_stats().closed, 0);
}
#[tokio::test]
async fn live_arrival_interleaves_durable_and_sequence_less_ephemeral_events() {
let session_id = SessionId::new();
let turn_id = TurnId::new();
let message_id = MessageId::new();
let bus = Arc::new(FacadeEventBus::new());
let emitter = host(bus.clone());
let mut stream = bus.subscribe();
emitter
.emit(turn_started(session_id, turn_id))
.await
.unwrap();
emitter
.emit(everruns_core::EventRequest::new(
session_id,
everruns_core::EventContext::turn(turn_id, message_id),
OutputMessageDeltaData {
turn_id,
message_id,
delta: "hi".to_string(),
accumulated: "hi".to_string(),
phase: None,
},
))
.await
.unwrap();
emitter
.emit(everruns_core::EventRequest::new(
session_id,
everruns_core::EventContext::turn(turn_id, message_id),
TurnCancelledData {
turn_id,
reason: Some("test".to_string()),
usage: None,
},
))
.await
.unwrap();
let started = stream.recv().await.unwrap().unwrap();
let delta = stream.recv().await.unwrap().unwrap();
let cancelled = stream.recv().await.unwrap().unwrap();
assert_eq!(started.event_type(), "turn.started");
assert_eq!(delta.event_type(), "output.message.delta");
assert_eq!(cancelled.event_type(), "turn.cancelled");
assert_eq!(started.sequence(), Some(1));
assert_eq!(delta.sequence(), None);
assert!(delta.as_json().get("sequence").is_none());
assert_eq!(cancelled.sequence(), Some(2));
}
#[tokio::test]
async fn cancellation_uses_the_active_turn_and_canonical_sequence() {
let session_id = SessionId::new();
let turn_id = TurnId::new();
let bus = Arc::new(FacadeEventBus::new());
let emitter = host(bus.clone());
let mut stream = bus.subscribe();
emitter
.emit(turn_started(session_id, turn_id))
.await
.expect("turn starts");
let (cancelled_turn_id, request) = bus.cancellation_request(session_id);
emitter.emit(request).await.expect("cancellation commits");
assert_eq!(cancelled_turn_id, turn_id);
let started = stream.recv().await.expect("no lag").expect("start event");
let cancelled = stream.recv().await.expect("no lag").expect("cancel event");
assert!(matches!(started.kind, SessionEventKind::TurnStarted));
assert!(matches!(cancelled.kind, SessionEventKind::TurnCancelled));
assert_eq!(
cancelled.turn_id.as_deref(),
Some(turn_id.to_string().as_str())
);
assert_eq!(cancelled.as_json()["sequence"], 2);
assert_eq!(
cancelled.data(),
&serde_json::to_value(TurnCancelledData {
turn_id,
reason: Some("cancelled by application".to_string()),
usage: None,
})
.expect("cancel data serializes")
);
}
#[tokio::test]
async fn output_replacement_retains_rebuildable_message_identity_and_text() {
let session_id = SessionId::new();
let turn_id = TurnId::new();
let message_id = MessageId::new();
let bus = Arc::new(FacadeEventBus::new());
let emitter = host(bus.clone());
let mut stream = bus.subscribe();
emitter
.emit(everruns_core::EventRequest::new(
session_id,
everruns_core::EventContext {
turn_id: Some(turn_id),
..everruns_core::EventContext::default()
},
OutputMessageReplacedData {
turn_id,
message_id,
guardrail_capability_id: "guardrails".to_string(),
guardrail_id: "output-policy".to_string(),
reason_code: "blocked".to_string(),
replacement: "Response withheld.".to_string(),
},
))
.await
.expect("replacement emits");
let replacement = stream
.recv()
.await
.expect("no lag")
.expect("replacement delivered");
assert!(matches!(
&replacement.kind,
SessionEventKind::OutputReplaced {
message_id: observed_id,
replacement,
} if observed_id == &message_id.to_string() && replacement == "Response withheld."
));
assert_eq!(replacement.data()["message_id"], message_id.to_string());
assert_eq!(replacement.data()["replacement"], "Response withheld.");
}
#[tokio::test]
async fn tool_narration_is_preserved_for_renderers() {
let session_id = SessionId::new();
let turn_id = TurnId::new();
let bus = Arc::new(FacadeEventBus::new());
let emitter = host(bus.clone());
let mut stream = bus.subscribe();
emitter
.emit(everruns_core::EventRequest::new(
session_id,
everruns_core::EventContext {
turn_id: Some(turn_id),
..everruns_core::EventContext::default()
},
ToolStartedData {
tool_call: ToolCall {
id: "call_1".to_string(),
name: "lookup".to_string(),
arguments: json!({"key": "answer"}),
},
tool_call_fingerprint: None,
display_name: Some("Knowledge lookup".to_string()),
narration: Some("Looking up the answer".to_string()),
},
))
.await
.expect("tool start emits");
let observed = stream.recv().await.unwrap().unwrap();
assert_eq!(observed.narration(), Some("Looking up the answer"));
assert_eq!(observed.data()["display_name"], "Knowledge lookup");
}
#[tokio::test]
async fn unpromoted_event_kind_retains_its_complete_canonical_payload() {
let session_id = SessionId::new();
let turn_id = TurnId::new();
let bus = Arc::new(FacadeEventBus::new());
let emitter = host(bus.clone());
let mut stream = bus.subscribe();
let canonical = emitter
.emit(everruns_core::EventRequest::new(
session_id,
everruns_core::EventContext {
turn_id: Some(turn_id),
..everruns_core::EventContext::default()
},
ActStartedData {
tool_calls: Vec::new(),
headline: Some("running tools".to_string()),
},
))
.await
.expect("event emits");
let observed = stream.recv().await.unwrap().unwrap();
assert!(matches!(
&observed.kind,
SessionEventKind::Other { event_type, payload }
if event_type == "act.started" && payload["headline"] == "running tools"
));
assert_eq!(
observed.as_json(),
&serde_json::to_value(canonical).expect("canonical event serializes")
);
}
#[tokio::test]
async fn provider_failure_retains_reason_and_turn_terminal_payloads() {
let agent = Agent::builder()
.instructions("Answer concisely.")
.model(Model::simulated_error("provider unavailable"))
.build()
.expect("valid agent");
let session = agent.session();
let mut stream = session.events();
let result = session
.run("hello")
.await
.expect("provider failure resolves to a failed turn");
assert!(!result.success);
drop(session);
let mut observed = Vec::new();
while let Some(event) = stream.recv().await.expect("failure stream does not lag") {
observed.push(event);
}
let reason_failure = observed
.iter()
.find(|event| {
matches!(
event.kind,
SessionEventKind::ReasonCompleted { success: false, .. }
)
})
.expect("reason.completed preserves the provider failure");
assert!(
reason_failure.data()["error"]
.as_str()
.is_some_and(|error| error.contains("provider unavailable"))
);
let turn_failure = observed
.iter()
.find(|event| matches!(event.kind, SessionEventKind::TurnFailed { .. }))
.expect("turn.failed is the terminal event");
assert_eq!(turn_failure.event_type(), "turn.failed");
assert_eq!(
turn_failure.turn_id.as_deref(),
Some(result.turn_id.as_str())
);
assert!(turn_failure.data()["error"].as_str().is_some());
}
#[tokio::test]
async fn tool_lifecycle_retains_arguments_result_and_order() {
let tool = crate::FunctionTool::new(
"lookup",
"Look up a value.",
json!({
"type": "object",
"properties": { "key": { "type": "string" } },
"required": ["key"]
}),
|arguments: serde_json::Value| async move {
Ok::<_, String>(json!({ "value": arguments["key"] }))
},
);
let agent = Agent::builder()
.instructions("Use the lookup tool.")
.model(Model::simulated_scripted(
"done",
vec![
vec![ToolCall {
id: "call_lookup_1".to_string(),
name: "lookup".to_string(),
arguments: json!({ "key": "answer" }),
}],
vec![],
],
))
.tool(tool)
.build()
.expect("valid agent");
let session = agent.session();
let mut stream = session.events();
let result = session.run("look it up").await.expect("tool turn runs");
assert!(result.success);
drop(session);
let mut observed = Vec::new();
while let Some(event) = stream.recv().await.expect("tool stream does not lag") {
observed.push(event);
}
let started = observed
.iter()
.find(|event| matches!(event.kind, SessionEventKind::ToolStarted { .. }))
.expect("tool.started");
let completed = observed
.iter()
.find(|event| matches!(event.kind, SessionEventKind::ToolCompleted { .. }))
.expect("tool.completed");
assert!(
started.sequence().expect("tool start is durable")
< completed.sequence().expect("tool completion is durable")
);
assert_eq!(started.data()["tool_call"]["id"], "call_lookup_1");
assert_eq!(started.data()["tool_call"]["arguments"]["key"], "answer");
assert_eq!(completed.data()["tool_call_id"], "call_lookup_1");
assert_eq!(completed.data()["status"], "success");
assert!(completed.data()["result"].is_array());
}
}