use std::sync::{Mutex, OnceLock};
use everruns_core::events::{
self, Event, EventContext, EventData, EventRequest, InputMessageData,
OutputMessageCompletedData, OutputMessageDeltaData, OutputMessageReplacedData,
OutputMessageStartedData, ReasonCompletedData, ToolCompletedData, ToolOutputDeltaData,
ToolProgressData, ToolStartedData, TurnCancelledData, TurnFailedData,
};
use everruns_host::{EventSink, EventSinkError};
use everruns_provider::typed_id::{SessionId, TurnId};
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,
event_type: String,
timestamp: String,
data: Value,
raw: Value,
canonical: 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>,
},
ReasoningDelta {
delta: String,
accumulated: String,
},
ReasoningCompleted {
text: String,
},
ReasoningItem {
provider: String,
item_id: Option<String>,
summary: Vec<String>,
},
ModelGeneration {
model: String,
provider: Option<String>,
input_tokens: Option<u32>,
output_tokens: Option<u32>,
cost_usd: Option<f64>,
duration_ms: Option<u64>,
success: bool,
},
Other {
event_type: String,
},
}
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.event_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.timestamp
}
pub fn as_json(&self) -> &Value {
&self.raw
}
pub fn into_json(self) -> Value {
self.raw
}
pub fn canonical_json(&self) -> &Value {
&self.canonical
}
pub fn data(&self) -> &Value {
&self.data
}
pub fn narration(&self) -> Option<&str> {
self.data().get("narration").and_then(Value::as_str)
}
fn from_core_event(event: &Event) -> Self {
let mut raw = serde_json::to_value(event).expect("canonical events are JSON serializable");
let mut data = serde_json::to_value(&event.data)
.expect("canonical event payloads are JSON serializable");
if event.event_type == events::OUTPUT_MESSAGE_DELTA {
data.as_object_mut()
.and_then(|data| data.remove("accumulated"));
raw.get_mut("data")
.and_then(Value::as_object_mut)
.and_then(|data| data.remove("accumulated"));
}
let canonical = raw.clone();
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::REASON_THINKING_DELTA => match &event.data {
EventData::ReasonThinkingDelta(data) => SessionEventKind::ReasoningDelta {
delta: data.delta.clone(),
accumulated: data.accumulated.clone(),
},
_ => Self::other_kind(event),
},
events::REASON_THINKING_COMPLETED => match &event.data {
EventData::ReasonThinkingCompleted(data) => SessionEventKind::ReasoningCompleted {
text: data.thinking.clone(),
},
_ => Self::other_kind(event),
},
events::REASON_ITEM => match &event.data {
EventData::ReasonItem(data) => SessionEventKind::ReasoningItem {
provider: data.provider.clone(),
item_id: (!data.item_id.is_empty()).then(|| data.item_id.clone()),
summary: data.summary.clone(),
},
_ => Self::other_kind(event),
},
events::LLM_GENERATION => match &event.data {
EventData::LlmGeneration(data) => {
let usage = data.metadata.usage.as_ref();
SessionEventKind::ModelGeneration {
model: data.metadata.model.clone(),
provider: data.metadata.provider.clone(),
input_tokens: usage.map(|usage| usage.input_tokens),
output_tokens: usage.map(|usage| usage.output_tokens),
cost_usd: usage.and_then(|usage| usage.effective_cost_usd()),
duration_ms: data.metadata.duration_ms,
success: data.metadata.success,
}
}
_ => Self::other_kind(event),
},
_ => Self::other_kind(event),
};
let data = Self::reviewed_data(&kind, &data);
raw["data"] = data.clone();
Self {
event_id: event.id.to_string(),
session_id: event.session_id.to_string(),
turn_id,
kind,
event_type: event.event_type.clone(),
timestamp: event.ts.to_rfc3339(),
data,
raw,
canonical,
}
}
fn reviewed_data(kind: &SessionEventKind, data: &Value) -> Value {
const DISPLAY_FIELDS: [&str; 2] = ["narration", "display_name"];
let mut reviewed = match kind {
SessionEventKind::InputMessage { message_id }
| SessionEventKind::OutputStarted { message_id } => {
serde_json::json!({ "message_id": message_id })
}
SessionEventKind::OutputCompleted { message_id } => {
serde_json::json!({ "message_id": message_id })
}
SessionEventKind::OutputReplaced {
message_id,
replacement,
} => serde_json::json!({ "message_id": message_id, "replacement": replacement }),
SessionEventKind::TurnFailed { error } => serde_json::json!({ "error": error }),
SessionEventKind::TextDelta { delta } => serde_json::json!({ "delta": delta }),
SessionEventKind::ToolStarted {
tool_call_id,
tool_name,
} => serde_json::json!({ "tool_call_id": tool_call_id, "tool_name": tool_name }),
SessionEventKind::ToolCompleted {
tool_call_id,
tool_name,
success,
} => serde_json::json!({
"tool_call_id": tool_call_id,
"tool_name": tool_name,
"success": success,
}),
SessionEventKind::ToolProgress {
tool_call_id,
tool_name,
message,
} => serde_json::json!({
"tool_call_id": tool_call_id,
"tool_name": tool_name,
"message": message,
}),
SessionEventKind::ToolOutputDelta {
tool_call_id,
tool_name,
stream,
delta,
} => serde_json::json!({
"tool_call_id": tool_call_id,
"tool_name": tool_name,
"stream": stream,
"delta": delta,
}),
SessionEventKind::ReasoningDelta { delta, accumulated } => {
serde_json::json!({ "delta": delta, "accumulated": accumulated })
}
SessionEventKind::ReasonCompleted { success, error } => {
serde_json::json!({ "success": success, "error": error })
}
SessionEventKind::ReasoningCompleted { text } => serde_json::json!({ "text": text }),
SessionEventKind::ReasoningItem {
provider,
item_id,
summary,
} => serde_json::json!({
"provider": provider,
"item_id": item_id,
"summary": summary,
}),
SessionEventKind::ModelGeneration {
model,
provider,
input_tokens,
output_tokens,
cost_usd,
duration_ms,
success,
} => serde_json::json!({
"model": model,
"provider": provider,
"input_tokens": input_tokens,
"output_tokens": output_tokens,
"cost_usd": cost_usd,
"duration_ms": duration_ms,
"success": success,
}),
SessionEventKind::TurnCancelled => {
let mut cancelled = serde_json::Map::new();
for field in ["turn_id", "reason", "usage"] {
if let Some(value) = data.get(field) {
cancelled.insert(field.to_string(), value.clone());
}
}
Value::Object(cancelled)
}
SessionEventKind::Other { .. }
| SessionEventKind::TurnStarted
| SessionEventKind::TurnCompleted
| SessionEventKind::ReasonStarted => serde_json::json!({}),
};
if let Some(object) = reviewed.as_object_mut() {
for field in DISPLAY_FIELDS {
if let Some(value) = data.get(field) {
object.insert(field.to_string(), value.clone());
}
}
}
reviewed
}
fn other_kind(event: &Event) -> SessionEventKind {
SessionEventKind::Other {
event_type: event.event_type.clone(),
}
}
}
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>,
pub(crate) timeout: Option<std::time::Duration>,
}
impl RunOptions {
pub fn new() -> Self {
Self::default()
}
pub fn cancel_token(mut self, token: CancellationToken) -> Self {
self.cancel = Some(token);
self
}
pub fn timeout(mut self, timeout: std::time::Duration) -> Self {
self.timeout = Some(timeout);
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: OnceLock<broadcast::Sender<SessionEvent>>,
capacity: usize,
active_turn: Mutex<Option<EventContext>>,
}
impl FacadeEventBus {
pub(crate) fn new() -> Self {
Self::with_capacity(EVENT_STREAM_CAPACITY)
}
fn with_capacity(capacity: usize) -> Self {
Self {
sender: OnceLock::new(),
capacity,
active_turn: Mutex::new(None),
}
}
pub(crate) fn subscribe(&self) -> EventStream {
let sender = self
.sender
.get_or_init(|| broadcast::channel(self.capacity).0);
EventStream::new(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();
}
_ => {}
}
if let Some(sender) = self.sender.get() {
let _ = 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::event_emitter::EventEmitter;
use everruns_core::events::{
ActStartedData, OutputMessageDeltaData, OutputMessageReplacedData, ToolStartedData,
TurnCancelledData, TurnStartedData,
};
use everruns_host::{HostEventEmitter, InMemoryEventLog};
use everruns_provider::tool_types::ToolCall;
use everruns_provider::typed_id::{MessageId, SessionId, TurnId};
use serde_json::json;
use super::{EventStreamError, FacadeEventBus, SessionEvent, SessionEventKind};
use crate::{Agent, InMemoryEngine, 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()),
agent_id: None,
agent_name: None,
agent_description: None,
},
)
}
#[test]
fn reviewed_data_never_exceeds_the_canonical_payload() {
let canonical = json!({
"tool_call_id": "call_1",
"tool_name": "lookup",
"success": true,
"status": "success",
"result": ["secret"],
"narration": "Looking it up",
"display_name": "Knowledge lookup",
});
let kind = SessionEventKind::ToolCompleted {
tool_call_id: "call_1".to_string(),
tool_name: "lookup".to_string(),
success: true,
};
let reviewed = SessionEvent::reviewed_data(&kind, &canonical);
let reviewed = reviewed.as_object().expect("reviewed data is an object");
for (key, value) in reviewed {
assert_eq!(
Some(value),
canonical.get(key),
"reviewed key {key} is absent from or differs in the canonical payload"
);
}
assert_eq!(reviewed["tool_call_id"], "call_1");
assert_eq!(reviewed["narration"], "Looking it up");
assert!(!reviewed.contains_key("result"));
}
#[test]
fn a_cancellation_field_nobody_promoted_stays_off_the_reviewed_surface() {
let canonical = json!({
"turn_id": "turn_1",
"reason": "user cancelled",
"usage": { "input_tokens": 12, "output_tokens": 3 },
"partial_output": "the model had written this far",
});
let reviewed = SessionEvent::reviewed_data(&SessionEventKind::TurnCancelled, &canonical);
assert_eq!(reviewed["turn_id"], "turn_1");
assert_eq!(reviewed["reason"], "user cancelled");
assert_eq!(reviewed["usage"]["input_tokens"], 12);
assert!(
reviewed.get("partial_output").is_none(),
"an unpromoted cancellation field must not reach the reviewed surface"
);
}
#[tokio::test]
async fn envelope_is_complete_while_data_stays_reviewed() {
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.data(), &json!({}));
assert!(!observed.as_json().to_string().contains("hello"));
assert_eq!(
observed.canonical_json(),
&serde_json::to_value(canonical).expect("canonical event serializes")
);
assert_eq!(observed.canonical_json()["data"]["input_content"], "hello");
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.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 subscriber_after_earlier_events_still_observes_later_ones() {
let session_id = SessionId::new();
let bus = Arc::new(FacadeEventBus::new());
let emitter = host(bus.clone());
emitter
.emit(turn_started(session_id, TurnId::new()))
.await
.expect("emitting without a subscriber succeeds");
let mut stream = bus.subscribe();
let turn_id = TurnId::new();
emitter
.emit(turn_started(session_id, turn_id))
.await
.expect("emitting to a live subscriber succeeds");
let observed = stream
.recv()
.await
.expect("stream stays open")
.expect("the event emitted after subscribing arrives");
assert_eq!(
observed.turn_id.as_deref(),
Some(turn_id.to_string().as_str())
);
}
#[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!(delta.data()["delta"], "hi");
assert!(delta.data().get("accumulated").is_none());
assert!(delta.as_json()["data"].get("accumulated").is_none());
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_is_identified_but_not_projected() {
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 } if event_type == "act.started"
));
assert_eq!(observed.data(), &json!({}));
assert!(!observed.as_json().to_string().contains("running tools"));
assert_eq!(
observed.canonical_json(),
&serde_json::to_value(canonical).expect("canonical event serializes")
);
assert_eq!(
observed.canonical_json()["data"]["headline"],
"running tools"
);
}
#[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 = InMemoryEngine::new().create(agent.clone());
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_keeps_order_and_reaches_arguments_canonically() {
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 = InMemoryEngine::new().create(agent.clone());
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!(completed.data()["tool_call_id"], "call_lookup_1");
assert_eq!(completed.data()["success"], true);
assert!(started.data()["tool_call"].is_null());
assert!(completed.data()["result"].is_null());
assert_eq!(
started.canonical_json()["data"]["tool_call"]["arguments"]["key"],
"answer"
);
assert_eq!(completed.canonical_json()["data"]["status"], "success");
assert!(completed.canonical_json()["data"]["result"].is_array());
}
}