use chrono::{DateTime, Utc};
use serde::{Deserialize, Deserializer, Serialize};
use uuid::Uuid;
#[cfg(feature = "openapi")]
use utoipa::ToSchema;
use crate::typed_id::{EventId, ExecId, MessageId, SessionId, TurnId};
mod compaction_data;
mod file_voice_data;
mod llm_data;
mod message_data;
mod reason_data;
mod tool_data;
mod turn_data;
mod usage;
pub use compaction_data::*;
pub use file_voice_data::*;
pub use llm_data::*;
pub use message_data::*;
pub use reason_data::*;
pub use tool_data::*;
pub use turn_data::*;
pub use usage::*;
pub const INPUT_MESSAGE: &str = "input.message";
pub const OUTPUT_MESSAGE_STARTED: &str = "output.message.started";
pub const OUTPUT_MESSAGE_DELTA: &str = "output.message.delta";
pub const OUTPUT_MESSAGE_COMPLETED: &str = "output.message.completed";
pub const OUTPUT_MESSAGE_REPLACED: &str = "output.message.replaced";
pub const TURN_STARTED: &str = "turn.started";
pub const TURN_COMPLETED: &str = "turn.completed";
pub const TURN_FAILED: &str = "turn.failed";
pub const TURN_SEALED: &str = "turn.sealed";
pub const TURN_CANCELLED: &str = "turn.cancelled";
pub const REASON_STARTED: &str = "reason.started";
pub const REASON_COMPLETED: &str = "reason.completed";
pub const REASON_RECOVERED: &str = "reason.recovered";
pub const CAPABILITY_USAGE: &str = "capability.usage";
pub const ACT_STARTED: &str = "act.started";
pub const ACT_COMPLETED: &str = "act.completed";
pub const TOOL_STARTED: &str = "tool.started";
pub const TOOL_COMPLETED: &str = "tool.completed";
pub const TOOL_PROGRESS: &str = "tool.progress";
pub const TOOL_OUTPUT_DELTA: &str = "tool.output.delta";
pub const TOOL_CALL_REQUESTED: &str = "tool.call_requested";
pub const TRANSCRIPT_REPAIRED: &str = "transcript.repaired";
pub const TOOL_CALL_REPAIRED: &str = "tool.call_repaired";
pub const LLM_GENERATION: &str = "llm.generation";
fn is_ephemeral_event_type(event_type: &str) -> bool {
matches!(
event_type,
OUTPUT_MESSAGE_DELTA
| REASON_THINKING_DELTA
| TOOL_OUTPUT_DELTA
| VOICE_INPUT_TRANSCRIPT_DELTA
| VOICE_OUTPUT_TRANSCRIPT_DELTA
)
}
pub const REASON_THINKING_STARTED: &str = "reason.thinking.started";
pub const REASON_THINKING_DELTA: &str = "reason.thinking.delta";
pub const REASON_THINKING_COMPLETED: &str = "reason.thinking.completed";
pub const REASON_ITEM: &str = "reason.item";
pub const SESSION_STARTED: &str = "session.started";
pub const SESSION_ACTIVATED: &str = "session.activated";
pub const SESSION_IDLED: &str = "session.idled";
pub const SESSION_TITLE_UPDATED: &str = "session.title.updated";
pub const SESSION_MODEL_CHANGED: &str = "session.model.changed";
pub const SCHEDULE_TRIGGERED: &str = "schedule.triggered";
pub const TASK_CREATED: &str = "task.created";
pub const TASK_UPDATED: &str = "task.updated";
pub const TASK_MESSAGE_SENT: &str = "task.message.sent";
pub const TASK_MESSAGE_RECEIVED: &str = "task.message.received";
pub const CONTEXT_COMPACTING: &str = "context.compacting";
pub const CONTEXT_COMPACTED: &str = "context.compacted";
pub const CONTEXT_COMPACTION_SKIPPED: &str = "context.compaction.skipped";
pub const CONTEXT_COMPACTION_FAILED: &str = "context.compaction.failed";
pub const FILE_WRITTEN: &str = "file.written";
pub const BUDGET_WARNING: &str = "budget.warning";
pub const BUDGET_PAUSED: &str = "budget.paused";
pub const BUDGET_EXHAUSTED: &str = "budget.exhausted";
pub const BUDGET_RESUMED: &str = "budget.resumed";
pub const VOICE_SESSION_STARTED: &str = "voice.session.started";
pub const VOICE_INPUT_TRANSCRIPT_DELTA: &str = "voice.input_transcript.delta";
pub const VOICE_INPUT_TRANSCRIPT_COMPLETED: &str = "voice.input_transcript.completed";
pub const VOICE_OUTPUT_TRANSCRIPT_DELTA: &str = "voice.output_transcript.delta";
pub const VOICE_OUTPUT_TRANSCRIPT_COMPLETED: &str = "voice.output_transcript.completed";
pub const VOICE_SESSION_ENDED: &str = "voice.session.ended";
pub const VOICE_SESSION_FAILED: &str = "voice.session.failed";
pub const VALID_EVENT_TYPES: &[&str] = &[
INPUT_MESSAGE,
OUTPUT_MESSAGE_STARTED,
OUTPUT_MESSAGE_DELTA,
OUTPUT_MESSAGE_COMPLETED,
OUTPUT_MESSAGE_REPLACED,
TURN_STARTED,
TURN_COMPLETED,
TURN_FAILED,
TURN_SEALED,
TURN_CANCELLED,
REASON_STARTED,
REASON_COMPLETED,
REASON_RECOVERED,
ACT_STARTED,
ACT_COMPLETED,
TOOL_STARTED,
TOOL_COMPLETED,
TOOL_PROGRESS,
TOOL_OUTPUT_DELTA,
TOOL_CALL_REQUESTED,
TRANSCRIPT_REPAIRED,
TOOL_CALL_REPAIRED,
LLM_GENERATION,
REASON_THINKING_STARTED,
REASON_THINKING_DELTA,
REASON_THINKING_COMPLETED,
REASON_ITEM,
SESSION_STARTED,
SESSION_ACTIVATED,
SESSION_IDLED,
SESSION_TITLE_UPDATED,
SESSION_MODEL_CHANGED,
SCHEDULE_TRIGGERED,
CONTEXT_COMPACTING,
CONTEXT_COMPACTED,
CONTEXT_COMPACTION_SKIPPED,
CONTEXT_COMPACTION_FAILED,
BUDGET_WARNING,
BUDGET_PAUSED,
BUDGET_EXHAUSTED,
BUDGET_RESUMED,
VOICE_SESSION_STARTED,
VOICE_INPUT_TRANSCRIPT_DELTA,
VOICE_INPUT_TRANSCRIPT_COMPLETED,
VOICE_OUTPUT_TRANSCRIPT_DELTA,
VOICE_OUTPUT_TRANSCRIPT_COMPLETED,
VOICE_SESSION_ENDED,
VOICE_SESSION_FAILED,
FILE_WRITTEN,
CAPABILITY_USAGE,
];
use crate::execution_context::ExecutionContext;
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
#[cfg_attr(feature = "openapi", derive(ToSchema))]
pub struct EventContext {
#[serde(skip_serializing_if = "Option::is_none")]
#[cfg_attr(feature = "openapi", schema(value_type = Option<String>, example = "turn_01933b5a00007000800000000000001"))]
pub turn_id: Option<TurnId>,
#[serde(skip_serializing_if = "Option::is_none")]
#[cfg_attr(feature = "openapi", schema(value_type = Option<String>, example = "message_01933b5a00007000800000000000001"))]
pub input_message_id: Option<MessageId>,
#[serde(skip_serializing_if = "Option::is_none")]
#[cfg_attr(feature = "openapi", schema(value_type = Option<String>, example = "exec_01933b5a00007000800000000000001"))]
pub exec_id: Option<ExecId>,
#[serde(skip_serializing_if = "Option::is_none")]
pub trace_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub span_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub parent_span_id: Option<String>,
}
impl EventContext {
pub fn empty() -> Self {
Self::default()
}
pub fn from_execution_context(ctx: &ExecutionContext) -> Self {
Self {
turn_id: Some(ctx.turn_id),
input_message_id: Some(ctx.input_message_id),
exec_id: Some(ctx.exec_id),
trace_id: None,
span_id: None,
parent_span_id: None,
}
}
pub fn turn(turn_id: TurnId, input_message_id: MessageId) -> Self {
Self {
turn_id: Some(turn_id),
input_message_id: Some(input_message_id),
exec_id: None,
trace_id: None,
span_id: None,
parent_span_id: None,
}
}
pub fn with_span(
mut self,
trace_id: String,
span_id: String,
parent_span_id: Option<String>,
) -> Self {
self.trace_id = Some(trace_id);
self.span_id = Some(span_id);
self.parent_span_id = parent_span_id;
self
}
}
#[derive(Debug, Clone, Serialize)]
#[cfg_attr(feature = "openapi", derive(ToSchema))]
pub struct Event {
#[cfg_attr(feature = "openapi", schema(value_type = String, example = "event_01933b5a00007000800000000000001"))]
pub id: EventId,
#[serde(rename = "type")]
pub event_type: String,
pub ts: DateTime<Utc>,
#[cfg_attr(feature = "openapi", schema(value_type = String, example = "session_01933b5a00007000800000000000001"))]
pub session_id: SessionId,
pub context: EventContext,
pub data: EventData,
#[serde(skip_serializing_if = "Option::is_none")]
pub metadata: Option<serde_json::Value>,
#[serde(skip_serializing_if = "Option::is_none")]
pub tags: Option<Vec<String>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub sequence: Option<i32>,
}
#[derive(Debug, Deserialize)]
struct RawEvent {
id: EventId,
#[serde(rename = "type")]
event_type: String,
ts: DateTime<Utc>,
session_id: SessionId,
context: EventContext,
data: serde_json::Value,
metadata: Option<serde_json::Value>,
tags: Option<Vec<String>>,
sequence: Option<i32>,
}
impl<'de> Deserialize<'de> for Event {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let raw = RawEvent::deserialize(deserializer)?;
let data = deserialize_event_data(&raw.event_type, raw.data);
Ok(Self {
id: raw.id,
event_type: raw.event_type,
ts: raw.ts,
session_id: raw.session_id,
context: raw.context,
data,
metadata: raw.metadata,
tags: raw.tags,
sequence: raw.sequence,
})
}
}
impl Event {
pub fn into_public(mut self) -> Self {
self.data = self.data.into_public();
self
}
pub fn new(session_id: SessionId, context: EventContext, data: impl Into<EventData>) -> Self {
let data = data.into();
let event_type = data.event_type().to_string();
Self {
id: EventId::new(),
event_type,
ts: Utc::now(),
session_id,
context,
data,
metadata: None,
tags: None,
sequence: None,
}
}
pub fn with_id(
id: EventId,
session_id: SessionId,
context: EventContext,
data: impl Into<EventData>,
) -> Self {
let data = data.into();
let event_type = data.event_type().to_string();
Self {
id,
event_type,
ts: Utc::now(),
session_id,
context,
data,
metadata: None,
tags: None,
sequence: None,
}
}
pub fn with_sequence(mut self, sequence: i32) -> Self {
self.sequence = Some(sequence);
self
}
pub fn with_metadata(mut self, metadata: serde_json::Value) -> Self {
self.metadata = Some(metadata);
self
}
pub fn with_tags(mut self, tags: Vec<String>) -> Self {
self.tags = Some(tags);
self
}
pub fn session_uuid(&self) -> Uuid {
self.session_id.uuid()
}
pub fn is_message_event(&self) -> bool {
self.event_type == INPUT_MESSAGE || self.event_type == OUTPUT_MESSAGE_COMPLETED
}
pub fn is_ephemeral(&self) -> bool {
is_ephemeral_event_type(&self.event_type)
}
pub fn is_input_event(&self) -> bool {
self.event_type.starts_with("input.")
}
pub fn is_output_event(&self) -> bool {
self.event_type.starts_with("output.")
}
pub fn is_atom_event(&self) -> bool {
matches!(
self.event_type.as_str(),
REASON_STARTED
| REASON_COMPLETED
| REASON_RECOVERED
| ACT_STARTED
| ACT_COMPLETED
| TOOL_STARTED
| TOOL_COMPLETED
| TOOL_PROGRESS
| TOOL_CALL_REQUESTED
| TRANSCRIPT_REPAIRED
)
}
pub fn is_turn_event(&self) -> bool {
self.event_type.starts_with("turn.")
}
pub fn is_session_event(&self) -> bool {
self.event_type.starts_with("session.")
}
pub fn is_unsupported(&self) -> bool {
self.data.is_unsupported()
}
}
use crate::message::{ContentPart, Message};
use crate::tool_narration::{
ToolNarrationPhase, render_group_headline_with_locale, render_tool_narration_with_locale,
};
use crate::tool_types::ToolCall;
use everruns_provider::execution_phase::ExecutionPhase;
pub const FILE_OP_CREATE: &str = "create";
#[derive(Debug, Clone, Serialize)]
#[serde(untagged)]
#[cfg_attr(feature = "openapi", derive(ToSchema))]
#[cfg_attr(feature = "openapi", schema(
title = "EventData",
description = "Event-specific payload. The schema depends on the event type field.",
example = json!({"message": {"id": "...", "role": "user", "content": []}})
))]
pub enum EventData {
InputMessage(InputMessageData),
OutputMessageDelta(OutputMessageDeltaData),
OutputMessageStarted(OutputMessageStartedData),
OutputMessageReplaced(OutputMessageReplacedData),
OutputMessageCompleted(OutputMessageCompletedData),
TurnStarted(TurnStartedData),
TurnCompleted(TurnCompletedData),
TurnFailed(TurnFailedData),
ReasonStarted(ReasonStartedData),
ReasonCompleted(ReasonCompletedData),
ReasonRecovered(ReasonRecoveredData),
CapabilityUsage(CapabilityUsageData),
ActStarted(ActStartedData),
ActCompleted(ActCompletedData),
ToolStarted(ToolStartedData),
ToolCompleted(ToolCompletedData),
ToolProgress(ToolProgressData),
ToolOutputDelta(ToolOutputDeltaData),
ToolCallRequested(ToolCallRequestedData),
TranscriptRepaired(TranscriptRepairedData),
ToolCallRepaired(ToolCallRepairedData),
LlmGeneration(LlmGenerationData),
ReasonThinkingDelta(ReasonThinkingDeltaData),
ReasonItem(ReasonItemData),
ReasonThinkingStarted(ReasonThinkingStartedData),
ReasonThinkingCompleted(ReasonThinkingCompletedData),
TurnSealed(TurnSealedData),
TurnCancelled(TurnCancelledData),
SessionStarted(SessionStartedData),
SessionActivated(SessionActivatedData),
SessionIdled(SessionIdledData),
SessionTitleUpdated(SessionTitleUpdatedData),
SessionModelChanged(SessionModelChangedData),
TaskCreated(SessionTaskEventData),
TaskUpdated(SessionTaskEventData),
TaskMessageSent(TaskMessageEventData),
TaskMessageReceived(TaskMessageEventData),
ContextCompacting(ContextCompactingData),
ContextCompacted(ContextCompactedData),
ContextCompactionSkipped(ContextCompactionSkippedData),
ContextCompactionFailed(ContextCompactionFailedData),
FileWritten(FileWrittenData),
BudgetWarning(BudgetEventData),
BudgetPaused(BudgetEventData),
BudgetExhausted(BudgetEventData),
BudgetResumed(BudgetEventData),
VoiceSessionStarted(VoiceSessionStartedData),
VoiceInputTranscriptDelta(VoiceTranscriptData),
VoiceInputTranscriptCompleted(VoiceTranscriptData),
VoiceOutputTranscriptDelta(VoiceTranscriptData),
VoiceOutputTranscriptCompleted(VoiceTranscriptData),
VoiceSessionEnded(VoiceSessionEndedData),
VoiceSessionFailed(VoiceSessionFailedData),
#[serde(skip)]
Unsupported {
event_type: String,
data: serde_json::Value,
},
}
impl EventData {
pub fn needs_public_projection(&self) -> bool {
matches!(
self,
EventData::InputMessage(_)
| EventData::OutputMessageCompleted(_)
| EventData::LlmGeneration(_)
)
}
pub fn into_public(self) -> Self {
match self {
EventData::InputMessage(mut data) => {
data.message = data.message.into_public();
EventData::InputMessage(data)
}
EventData::OutputMessageCompleted(mut data) => {
data.message = data.message.into_public();
EventData::OutputMessageCompleted(data)
}
EventData::LlmGeneration(mut data) => {
data.messages = data
.messages
.into_iter()
.map(Message::into_public)
.collect();
EventData::LlmGeneration(data)
}
other => other,
}
}
pub fn is_unsupported(&self) -> bool {
matches!(self, EventData::Unsupported { .. })
}
pub fn unsupported(event_type: String, data: serde_json::Value) -> Self {
tracing::warn!(
event_type = %event_type,
"Encountered unsupported event type - will be filtered from API responses"
);
EventData::Unsupported { event_type, data }
}
}
macro_rules! event_data_kinds {
($( $variant:ident($data:ty) = $type_const:path ),+ $(,)?) => {
impl EventData {
pub fn event_type(&self) -> &'static str {
match self {
$( EventData::$variant(_) => $type_const, )+
EventData::Unsupported { .. } => "unsupported",
}
}
}
pub fn deserialize_event_data(event_type: &str, data: serde_json::Value) -> EventData {
let result = match event_type {
$(
$type_const => serde_json::from_value::<$data>(data.clone())
.map(EventData::$variant),
)+
_ => return EventData::unsupported(event_type.to_string(), data),
};
result.unwrap_or_else(|e| {
tracing::warn!(
event_type = %event_type,
error = %e,
"Failed to deserialize known event type - treating as unsupported"
);
EventData::Unsupported {
event_type: event_type.to_string(),
data,
}
})
}
};
}
event_data_kinds! {
InputMessage(InputMessageData) = INPUT_MESSAGE,
OutputMessageStarted(OutputMessageStartedData) = OUTPUT_MESSAGE_STARTED,
OutputMessageDelta(OutputMessageDeltaData) = OUTPUT_MESSAGE_DELTA,
OutputMessageReplaced(OutputMessageReplacedData) = OUTPUT_MESSAGE_REPLACED,
OutputMessageCompleted(OutputMessageCompletedData) = OUTPUT_MESSAGE_COMPLETED,
TurnStarted(TurnStartedData) = TURN_STARTED,
TurnCompleted(TurnCompletedData) = TURN_COMPLETED,
TurnFailed(TurnFailedData) = TURN_FAILED,
TurnSealed(TurnSealedData) = TURN_SEALED,
TurnCancelled(TurnCancelledData) = TURN_CANCELLED,
ReasonStarted(ReasonStartedData) = REASON_STARTED,
ReasonCompleted(ReasonCompletedData) = REASON_COMPLETED,
ReasonRecovered(ReasonRecoveredData) = REASON_RECOVERED,
CapabilityUsage(CapabilityUsageData) = CAPABILITY_USAGE,
ActStarted(ActStartedData) = ACT_STARTED,
ActCompleted(ActCompletedData) = ACT_COMPLETED,
ToolStarted(ToolStartedData) = TOOL_STARTED,
ToolCompleted(ToolCompletedData) = TOOL_COMPLETED,
ToolProgress(ToolProgressData) = TOOL_PROGRESS,
ToolOutputDelta(ToolOutputDeltaData) = TOOL_OUTPUT_DELTA,
ToolCallRequested(ToolCallRequestedData) = TOOL_CALL_REQUESTED,
TranscriptRepaired(TranscriptRepairedData) = TRANSCRIPT_REPAIRED,
ToolCallRepaired(ToolCallRepairedData) = TOOL_CALL_REPAIRED,
LlmGeneration(LlmGenerationData) = LLM_GENERATION,
ReasonThinkingStarted(ReasonThinkingStartedData) = REASON_THINKING_STARTED,
ReasonThinkingDelta(ReasonThinkingDeltaData) = REASON_THINKING_DELTA,
ReasonThinkingCompleted(ReasonThinkingCompletedData) = REASON_THINKING_COMPLETED,
ReasonItem(ReasonItemData) = REASON_ITEM,
SessionStarted(SessionStartedData) = SESSION_STARTED,
SessionActivated(SessionActivatedData) = SESSION_ACTIVATED,
SessionIdled(SessionIdledData) = SESSION_IDLED,
SessionTitleUpdated(SessionTitleUpdatedData) = SESSION_TITLE_UPDATED,
SessionModelChanged(SessionModelChangedData) = SESSION_MODEL_CHANGED,
ContextCompacting(ContextCompactingData) = CONTEXT_COMPACTING,
ContextCompacted(ContextCompactedData) = CONTEXT_COMPACTED,
ContextCompactionSkipped(ContextCompactionSkippedData) = CONTEXT_COMPACTION_SKIPPED,
ContextCompactionFailed(ContextCompactionFailedData) = CONTEXT_COMPACTION_FAILED,
FileWritten(FileWrittenData) = FILE_WRITTEN,
BudgetWarning(BudgetEventData) = BUDGET_WARNING,
BudgetPaused(BudgetEventData) = BUDGET_PAUSED,
BudgetExhausted(BudgetEventData) = BUDGET_EXHAUSTED,
BudgetResumed(BudgetEventData) = BUDGET_RESUMED,
VoiceSessionStarted(VoiceSessionStartedData) = VOICE_SESSION_STARTED,
VoiceInputTranscriptDelta(VoiceTranscriptData) = VOICE_INPUT_TRANSCRIPT_DELTA,
VoiceInputTranscriptCompleted(VoiceTranscriptData) = VOICE_INPUT_TRANSCRIPT_COMPLETED,
VoiceOutputTranscriptDelta(VoiceTranscriptData) = VOICE_OUTPUT_TRANSCRIPT_DELTA,
VoiceOutputTranscriptCompleted(VoiceTranscriptData) = VOICE_OUTPUT_TRANSCRIPT_COMPLETED,
VoiceSessionEnded(VoiceSessionEndedData) = VOICE_SESSION_ENDED,
VoiceSessionFailed(VoiceSessionFailedData) = VOICE_SESSION_FAILED,
TaskCreated(SessionTaskEventData) = TASK_CREATED,
TaskUpdated(SessionTaskEventData) = TASK_UPDATED,
TaskMessageSent(TaskMessageEventData) = TASK_MESSAGE_SENT,
TaskMessageReceived(TaskMessageEventData) = TASK_MESSAGE_RECEIVED,
}
macro_rules! impl_from_event_data {
($($data_type:ty => $variant:ident),* $(,)?) => {
$(
impl From<$data_type> for EventData {
fn from(data: $data_type) -> Self {
EventData::$variant(data)
}
}
)*
};
}
impl_from_event_data! {
InputMessageData => InputMessage,
OutputMessageStartedData => OutputMessageStarted,
OutputMessageDeltaData => OutputMessageDelta,
OutputMessageReplacedData => OutputMessageReplaced,
OutputMessageCompletedData => OutputMessageCompleted,
TurnStartedData => TurnStarted,
TurnCompletedData => TurnCompleted,
TurnFailedData => TurnFailed,
TurnSealedData => TurnSealed,
TurnCancelledData => TurnCancelled,
ReasonStartedData => ReasonStarted,
ReasonCompletedData => ReasonCompleted,
ReasonRecoveredData => ReasonRecovered,
CapabilityUsageData => CapabilityUsage,
ActStartedData => ActStarted,
ActCompletedData => ActCompleted,
ToolStartedData => ToolStarted,
ToolCompletedData => ToolCompleted,
ToolProgressData => ToolProgress,
ToolOutputDeltaData => ToolOutputDelta,
ToolCallRequestedData => ToolCallRequested,
TranscriptRepairedData => TranscriptRepaired,
ToolCallRepairedData => ToolCallRepaired,
LlmGenerationData => LlmGeneration,
ReasonThinkingStartedData => ReasonThinkingStarted,
ReasonThinkingDeltaData => ReasonThinkingDelta,
ReasonThinkingCompletedData => ReasonThinkingCompleted,
ReasonItemData => ReasonItem,
SessionStartedData => SessionStarted,
SessionActivatedData => SessionActivated,
SessionIdledData => SessionIdled,
SessionTitleUpdatedData => SessionTitleUpdated,
SessionModelChangedData => SessionModelChanged,
ContextCompactingData => ContextCompacting,
ContextCompactedData => ContextCompacted,
ContextCompactionSkippedData => ContextCompactionSkipped,
ContextCompactionFailedData => ContextCompactionFailed,
FileWrittenData => FileWritten,
VoiceSessionStartedData => VoiceSessionStarted,
VoiceSessionEndedData => VoiceSessionEnded,
VoiceSessionFailedData => VoiceSessionFailed,
}
impl EventData {
pub fn voice_transcript_event(data: VoiceTranscriptData, event_type: &str) -> Self {
match event_type {
VOICE_INPUT_TRANSCRIPT_DELTA => EventData::VoiceInputTranscriptDelta(data),
VOICE_INPUT_TRANSCRIPT_COMPLETED => EventData::VoiceInputTranscriptCompleted(data),
VOICE_OUTPUT_TRANSCRIPT_DELTA => EventData::VoiceOutputTranscriptDelta(data),
VOICE_OUTPUT_TRANSCRIPT_COMPLETED => EventData::VoiceOutputTranscriptCompleted(data),
_ => EventData::unsupported(
event_type.to_string(),
serde_json::to_value(&data).unwrap_or(serde_json::Value::Null),
),
}
}
}
impl EventData {
pub fn budget_event(data: BudgetEventData, event_type: &str) -> Self {
match event_type {
BUDGET_WARNING => EventData::BudgetWarning(data),
BUDGET_PAUSED => EventData::BudgetPaused(data),
BUDGET_EXHAUSTED => EventData::BudgetExhausted(data),
BUDGET_RESUMED => EventData::BudgetResumed(data),
_ => EventData::unsupported(
event_type.to_string(),
serde_json::to_value(&data).unwrap_or(serde_json::Value::Null),
),
}
}
}
#[derive(Debug, Clone, Serialize)]
#[cfg_attr(feature = "openapi", derive(ToSchema))]
pub struct EventRequest {
#[serde(rename = "type")]
pub event_type: String,
pub ts: DateTime<Utc>,
pub session_id: SessionId,
pub context: EventContext,
pub data: EventData,
#[serde(skip_serializing_if = "Option::is_none")]
pub metadata: Option<serde_json::Value>,
#[serde(skip_serializing_if = "Option::is_none")]
pub tags: Option<Vec<String>>,
}
#[derive(Debug, Deserialize)]
struct RawEventRequest {
#[serde(rename = "type")]
event_type: String,
ts: DateTime<Utc>,
session_id: SessionId,
context: EventContext,
data: serde_json::Value,
metadata: Option<serde_json::Value>,
tags: Option<Vec<String>>,
}
impl<'de> Deserialize<'de> for EventRequest {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let raw = RawEventRequest::deserialize(deserializer)?;
let data = deserialize_event_data(&raw.event_type, raw.data);
Ok(Self {
event_type: raw.event_type,
ts: raw.ts,
session_id: raw.session_id,
context: raw.context,
data,
metadata: raw.metadata,
tags: raw.tags,
})
}
}
impl EventRequest {
pub fn new(session_id: SessionId, context: EventContext, data: impl Into<EventData>) -> Self {
let data = data.into();
let event_type = data.event_type().to_string();
Self {
event_type,
ts: Utc::now(),
session_id,
context,
data,
metadata: None,
tags: None,
}
}
pub fn with_metadata(mut self, metadata: serde_json::Value) -> Self {
self.metadata = Some(metadata);
self
}
pub fn with_tags(mut self, tags: Vec<String>) -> Self {
self.tags = Some(tags);
self
}
pub fn is_ephemeral(&self) -> bool {
is_ephemeral_event_type(&self.event_type)
}
pub fn into_event(self, id: EventId, sequence: i32) -> Event {
Event {
id,
event_type: self.event_type,
ts: self.ts,
session_id: self.session_id,
context: self.context,
data: self.data,
metadata: self.metadata,
tags: self.tags,
sequence: Some(sequence),
}
}
}
pub struct EventBuilder {
session_id: SessionId,
context: EventContext,
}
impl EventBuilder {
pub fn new(session_id: SessionId) -> Self {
Self {
session_id,
context: EventContext::empty(),
}
}
pub fn with_turn(mut self, turn_id: TurnId, input_message_id: MessageId) -> Self {
self.context.turn_id = Some(turn_id);
self.context.input_message_id = Some(input_message_id);
self
}
pub fn with_exec(mut self, exec_id: ExecId) -> Self {
self.context.exec_id = Some(exec_id);
self
}
pub fn build(self, data: impl Into<EventData>) -> Event {
Event::new(self.session_id, self.context, data)
}
}
#[cfg(test)]
mod contract_tests;
#[cfg(test)]
mod tests;
#[cfg(test)]
mod token_usage_tests;