use crate::state::AgentState;
use crate::types::{Message, MessageId, Role, RunId, ThreadId, ToolCallId};
use crate::JsonValue;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub enum EventType {
TextMessageStart,
TextMessageContent,
TextMessageEnd,
TextMessageChunk,
ThinkingTextMessageStart,
ThinkingTextMessageContent,
ThinkingTextMessageEnd,
ToolCallStart,
ToolCallArgs,
ToolCallEnd,
ToolCallChunk,
ToolCallResult,
ThinkingStart,
ThinkingEnd,
StateSnapshot,
StateDelta,
MessagesSnapshot,
ActivitySnapshot,
ActivityDelta,
Raw,
Custom,
RunStarted,
RunFinished,
RunError,
StepStarted,
StepFinished,
}
impl EventType {
pub fn as_str(&self) -> &'static str {
match self {
EventType::TextMessageStart => "TEXT_MESSAGE_START",
EventType::TextMessageContent => "TEXT_MESSAGE_CONTENT",
EventType::TextMessageEnd => "TEXT_MESSAGE_END",
EventType::TextMessageChunk => "TEXT_MESSAGE_CHUNK",
EventType::ThinkingTextMessageStart => "THINKING_TEXT_MESSAGE_START",
EventType::ThinkingTextMessageContent => "THINKING_TEXT_MESSAGE_CONTENT",
EventType::ThinkingTextMessageEnd => "THINKING_TEXT_MESSAGE_END",
EventType::ToolCallStart => "TOOL_CALL_START",
EventType::ToolCallArgs => "TOOL_CALL_ARGS",
EventType::ToolCallEnd => "TOOL_CALL_END",
EventType::ToolCallChunk => "TOOL_CALL_CHUNK",
EventType::ToolCallResult => "TOOL_CALL_RESULT",
EventType::ThinkingStart => "THINKING_START",
EventType::ThinkingEnd => "THINKING_END",
EventType::StateSnapshot => "STATE_SNAPSHOT",
EventType::StateDelta => "STATE_DELTA",
EventType::MessagesSnapshot => "MESSAGES_SNAPSHOT",
EventType::ActivitySnapshot => "ACTIVITY_SNAPSHOT",
EventType::ActivityDelta => "ACTIVITY_DELTA",
EventType::Raw => "RAW",
EventType::Custom => "CUSTOM",
EventType::RunStarted => "RUN_STARTED",
EventType::RunFinished => "RUN_FINISHED",
EventType::RunError => "RUN_ERROR",
EventType::StepStarted => "STEP_STARTED",
EventType::StepFinished => "STEP_FINISHED",
}
}
}
impl std::fmt::Display for EventType {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.as_str())
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
pub struct BaseEvent {
#[serde(skip_serializing_if = "Option::is_none")]
pub timestamp: Option<f64>,
#[serde(rename = "rawEvent", skip_serializing_if = "Option::is_none")]
pub raw_event: Option<JsonValue>,
}
impl BaseEvent {
pub fn new() -> Self {
Self::default()
}
pub fn with_current_timestamp() -> Self {
Self {
timestamp: Some(
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as f64)
.unwrap_or(0.0),
),
raw_event: None,
}
}
pub fn timestamp(mut self, timestamp: f64) -> Self {
self.timestamp = Some(timestamp);
self
}
pub fn raw_event(mut self, raw_event: JsonValue) -> Self {
self.raw_event = Some(raw_event);
self
}
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum EventValidationError {
#[error("Delta must not be an empty string")]
EmptyDelta,
#[error("Invalid event format: {0}")]
InvalidFormat(String),
#[error("Missing required field: {0}")]
MissingField(String),
#[error("Event type mismatch: expected {expected}, got {actual}")]
TypeMismatch {
expected: String,
actual: String,
},
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TextMessageStartEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "messageId")]
pub message_id: MessageId,
pub role: Role,
}
impl TextMessageStartEvent {
pub fn new(message_id: impl Into<MessageId>) -> Self {
Self {
base: BaseEvent::default(),
message_id: message_id.into(),
role: Role::Assistant,
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
pub fn with_raw_event(mut self, raw_event: JsonValue) -> Self {
self.base.raw_event = Some(raw_event);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TextMessageContentEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "messageId")]
pub message_id: MessageId,
pub delta: String,
}
impl TextMessageContentEvent {
pub fn new(
message_id: impl Into<MessageId>,
delta: impl Into<String>,
) -> Result<Self, EventValidationError> {
let delta = delta.into();
if delta.is_empty() {
return Err(EventValidationError::EmptyDelta);
}
Ok(Self {
base: BaseEvent::default(),
message_id: message_id.into(),
delta,
})
}
pub fn new_unchecked(message_id: impl Into<MessageId>, delta: impl Into<String>) -> Self {
Self {
base: BaseEvent::default(),
message_id: message_id.into(),
delta: delta.into(),
}
}
pub fn validate(&self) -> Result<(), EventValidationError> {
if self.delta.is_empty() {
return Err(EventValidationError::EmptyDelta);
}
Ok(())
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TextMessageEndEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "messageId")]
pub message_id: MessageId,
}
impl TextMessageEndEvent {
pub fn new(message_id: impl Into<MessageId>) -> Self {
Self {
base: BaseEvent::default(),
message_id: message_id.into(),
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TextMessageChunkEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "messageId", skip_serializing_if = "Option::is_none")]
pub message_id: Option<MessageId>,
pub role: Role,
#[serde(skip_serializing_if = "Option::is_none")]
pub delta: Option<String>,
}
impl TextMessageChunkEvent {
pub fn new(role: Role) -> Self {
Self {
base: BaseEvent::default(),
message_id: None,
role,
delta: None,
}
}
pub fn with_message_id(mut self, message_id: impl Into<MessageId>) -> Self {
self.message_id = Some(message_id.into());
self
}
pub fn with_delta(mut self, delta: impl Into<String>) -> Self {
self.delta = Some(delta.into());
self
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ThinkingTextMessageStartEvent {
#[serde(flatten)]
pub base: BaseEvent,
}
impl ThinkingTextMessageStartEvent {
pub fn new() -> Self {
Self {
base: BaseEvent::default(),
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
impl Default for ThinkingTextMessageStartEvent {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ThinkingTextMessageContentEvent {
#[serde(flatten)]
pub base: BaseEvent,
pub delta: String,
}
impl ThinkingTextMessageContentEvent {
pub fn new(delta: impl Into<String>) -> Self {
Self {
base: BaseEvent::default(),
delta: delta.into(),
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ThinkingTextMessageEndEvent {
#[serde(flatten)]
pub base: BaseEvent,
}
impl ThinkingTextMessageEndEvent {
pub fn new() -> Self {
Self {
base: BaseEvent::default(),
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
impl Default for ThinkingTextMessageEndEvent {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ToolCallStartEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "toolCallId")]
pub tool_call_id: ToolCallId,
#[serde(rename = "toolCallName")]
pub tool_call_name: String,
#[serde(rename = "parentMessageId", skip_serializing_if = "Option::is_none")]
pub parent_message_id: Option<MessageId>,
}
impl ToolCallStartEvent {
pub fn new(tool_call_id: impl Into<ToolCallId>, tool_call_name: impl Into<String>) -> Self {
Self {
base: BaseEvent::default(),
tool_call_id: tool_call_id.into(),
tool_call_name: tool_call_name.into(),
parent_message_id: None,
}
}
pub fn with_parent_message_id(mut self, message_id: impl Into<MessageId>) -> Self {
self.parent_message_id = Some(message_id.into());
self
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ToolCallArgsEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "toolCallId")]
pub tool_call_id: ToolCallId,
pub delta: String,
}
impl ToolCallArgsEvent {
pub fn new(tool_call_id: impl Into<ToolCallId>, delta: impl Into<String>) -> Self {
Self {
base: BaseEvent::default(),
tool_call_id: tool_call_id.into(),
delta: delta.into(),
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ToolCallEndEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "toolCallId")]
pub tool_call_id: ToolCallId,
}
impl ToolCallEndEvent {
pub fn new(tool_call_id: impl Into<ToolCallId>) -> Self {
Self {
base: BaseEvent::default(),
tool_call_id: tool_call_id.into(),
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ToolCallChunkEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "toolCallId", skip_serializing_if = "Option::is_none")]
pub tool_call_id: Option<ToolCallId>,
#[serde(rename = "toolCallName", skip_serializing_if = "Option::is_none")]
pub tool_call_name: Option<String>,
#[serde(rename = "parentMessageId", skip_serializing_if = "Option::is_none")]
pub parent_message_id: Option<MessageId>,
#[serde(skip_serializing_if = "Option::is_none")]
pub delta: Option<String>,
}
impl ToolCallChunkEvent {
pub fn new() -> Self {
Self {
base: BaseEvent::default(),
tool_call_id: None,
tool_call_name: None,
parent_message_id: None,
delta: None,
}
}
pub fn with_tool_call_id(mut self, tool_call_id: impl Into<ToolCallId>) -> Self {
self.tool_call_id = Some(tool_call_id.into());
self
}
pub fn with_tool_call_name(mut self, name: impl Into<String>) -> Self {
self.tool_call_name = Some(name.into());
self
}
pub fn with_parent_message_id(mut self, message_id: impl Into<MessageId>) -> Self {
self.parent_message_id = Some(message_id.into());
self
}
pub fn with_delta(mut self, delta: impl Into<String>) -> Self {
self.delta = Some(delta.into());
self
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
impl Default for ToolCallChunkEvent {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ToolCallResultEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "messageId")]
pub message_id: MessageId,
#[serde(rename = "toolCallId")]
pub tool_call_id: ToolCallId,
pub content: String,
#[serde(default = "Role::tool")]
pub role: Role,
}
impl ToolCallResultEvent {
pub fn new(
message_id: impl Into<MessageId>,
tool_call_id: impl Into<ToolCallId>,
content: impl Into<String>,
) -> Self {
Self {
base: BaseEvent::default(),
message_id: message_id.into(),
tool_call_id: tool_call_id.into(),
content: content.into(),
role: Role::Tool,
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RunStartedEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "threadId")]
pub thread_id: ThreadId,
#[serde(rename = "runId")]
pub run_id: RunId,
}
impl RunStartedEvent {
pub fn new(thread_id: impl Into<ThreadId>, run_id: impl Into<RunId>) -> Self {
Self {
base: BaseEvent::default(),
thread_id: thread_id.into(),
run_id: run_id.into(),
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub enum RunFinishedOutcome {
Success,
Interrupt,
}
impl Default for RunFinishedOutcome {
fn default() -> Self {
Self::Success
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
pub struct InterruptInfo {
#[serde(skip_serializing_if = "Option::is_none")]
pub id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub reason: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub payload: Option<JsonValue>,
}
impl InterruptInfo {
pub fn new() -> Self {
Self::default()
}
pub fn with_id(mut self, id: impl Into<String>) -> Self {
self.id = Some(id.into());
self
}
pub fn with_reason(mut self, reason: impl Into<String>) -> Self {
self.reason = Some(reason.into());
self
}
pub fn with_payload(mut self, payload: JsonValue) -> Self {
self.payload = Some(payload);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RunFinishedEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "threadId")]
pub thread_id: ThreadId,
#[serde(rename = "runId")]
pub run_id: RunId,
#[serde(skip_serializing_if = "Option::is_none")]
pub outcome: Option<RunFinishedOutcome>,
#[serde(skip_serializing_if = "Option::is_none")]
pub result: Option<JsonValue>,
#[serde(skip_serializing_if = "Option::is_none")]
pub interrupt: Option<InterruptInfo>,
}
impl RunFinishedEvent {
pub fn new(thread_id: impl Into<ThreadId>, run_id: impl Into<RunId>) -> Self {
Self {
base: BaseEvent::default(),
thread_id: thread_id.into(),
run_id: run_id.into(),
outcome: None,
result: None,
interrupt: None,
}
}
pub fn with_outcome(mut self, outcome: RunFinishedOutcome) -> Self {
self.outcome = Some(outcome);
self
}
pub fn with_result(mut self, result: JsonValue) -> Self {
self.result = Some(result);
self
}
pub fn with_interrupt(mut self, interrupt: InterruptInfo) -> Self {
self.outcome = Some(RunFinishedOutcome::Interrupt);
self.interrupt = Some(interrupt);
self
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
pub fn effective_outcome(&self) -> RunFinishedOutcome {
self.outcome.unwrap_or_else(|| {
if self.interrupt.is_some() {
RunFinishedOutcome::Interrupt
} else {
RunFinishedOutcome::Success
}
})
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RunErrorEvent {
#[serde(flatten)]
pub base: BaseEvent,
pub message: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub code: Option<String>,
}
impl RunErrorEvent {
pub fn new(message: impl Into<String>) -> Self {
Self {
base: BaseEvent::default(),
message: message.into(),
code: None,
}
}
pub fn with_code(mut self, code: impl Into<String>) -> Self {
self.code = Some(code.into());
self
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct StepStartedEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "stepName")]
pub step_name: String,
}
impl StepStartedEvent {
pub fn new(step_name: impl Into<String>) -> Self {
Self {
base: BaseEvent::default(),
step_name: step_name.into(),
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct StepFinishedEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "stepName")]
pub step_name: String,
}
impl StepFinishedEvent {
pub fn new(step_name: impl Into<String>) -> Self {
Self {
base: BaseEvent::default(),
step_name: step_name.into(),
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(bound(deserialize = ""))]
pub struct StateSnapshotEvent<StateT: AgentState = JsonValue> {
#[serde(flatten)]
pub base: BaseEvent,
pub snapshot: StateT,
}
impl<StateT: AgentState> StateSnapshotEvent<StateT> {
pub fn new(snapshot: StateT) -> Self {
Self {
base: BaseEvent::default(),
snapshot,
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
impl<StateT: AgentState + Default> Default for StateSnapshotEvent<StateT> {
fn default() -> Self {
Self {
base: BaseEvent::default(),
snapshot: StateT::default(),
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct StateDeltaEvent {
#[serde(flatten)]
pub base: BaseEvent,
pub delta: Vec<JsonValue>,
}
impl StateDeltaEvent {
pub fn new(delta: Vec<JsonValue>) -> Self {
Self {
base: BaseEvent::default(),
delta,
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
impl Default for StateDeltaEvent {
fn default() -> Self {
Self {
base: BaseEvent::default(),
delta: Vec::new(),
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct MessagesSnapshotEvent {
#[serde(flatten)]
pub base: BaseEvent,
pub messages: Vec<Message>,
}
impl MessagesSnapshotEvent {
pub fn new(messages: Vec<Message>) -> Self {
Self {
base: BaseEvent::default(),
messages,
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
impl Default for MessagesSnapshotEvent {
fn default() -> Self {
Self {
base: BaseEvent::default(),
messages: Vec::new(),
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ActivitySnapshotEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "messageId")]
pub message_id: MessageId,
#[serde(rename = "activityType")]
pub activity_type: String,
pub content: JsonValue,
#[serde(skip_serializing_if = "Option::is_none")]
pub replace: Option<bool>,
}
impl ActivitySnapshotEvent {
pub fn new(
message_id: impl Into<MessageId>,
activity_type: impl Into<String>,
content: JsonValue,
) -> Self {
Self {
base: BaseEvent::default(),
message_id: message_id.into(),
activity_type: activity_type.into(),
content,
replace: None,
}
}
pub fn with_replace(mut self, replace: bool) -> Self {
self.replace = Some(replace);
self
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ActivityDeltaEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(rename = "messageId")]
pub message_id: MessageId,
#[serde(rename = "activityType")]
pub activity_type: String,
pub patch: Vec<JsonValue>,
}
impl ActivityDeltaEvent {
pub fn new(
message_id: impl Into<MessageId>,
activity_type: impl Into<String>,
patch: Vec<JsonValue>,
) -> Self {
Self {
base: BaseEvent::default(),
message_id: message_id.into(),
activity_type: activity_type.into(),
patch,
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ThinkingStartEvent {
#[serde(flatten)]
pub base: BaseEvent,
#[serde(skip_serializing_if = "Option::is_none")]
pub title: Option<String>,
}
impl ThinkingStartEvent {
pub fn new() -> Self {
Self {
base: BaseEvent::default(),
title: None,
}
}
pub fn with_title(mut self, title: impl Into<String>) -> Self {
self.title = Some(title.into());
self
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
impl Default for ThinkingStartEvent {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ThinkingEndEvent {
#[serde(flatten)]
pub base: BaseEvent,
}
impl ThinkingEndEvent {
pub fn new() -> Self {
Self {
base: BaseEvent::default(),
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
impl Default for ThinkingEndEvent {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RawEvent {
#[serde(flatten)]
pub base: BaseEvent,
pub event: JsonValue,
#[serde(skip_serializing_if = "Option::is_none")]
pub source: Option<String>,
}
impl RawEvent {
pub fn new(event: JsonValue) -> Self {
Self {
base: BaseEvent::default(),
event,
source: None,
}
}
pub fn with_source(mut self, source: impl Into<String>) -> Self {
self.source = Some(source.into());
self
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CustomEvent {
#[serde(flatten)]
pub base: BaseEvent,
pub name: String,
pub value: JsonValue,
}
impl CustomEvent {
pub fn new(name: impl Into<String>, value: JsonValue) -> Self {
Self {
base: BaseEvent::default(),
name: name.into(),
value,
}
}
pub fn with_timestamp(mut self, timestamp: f64) -> Self {
self.base.timestamp = Some(timestamp);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "SCREAMING_SNAKE_CASE", bound(deserialize = ""))]
pub enum Event<StateT: AgentState = JsonValue> {
TextMessageStart(TextMessageStartEvent),
TextMessageContent(TextMessageContentEvent),
TextMessageEnd(TextMessageEndEvent),
TextMessageChunk(TextMessageChunkEvent),
ThinkingTextMessageStart(ThinkingTextMessageStartEvent),
ThinkingTextMessageContent(ThinkingTextMessageContentEvent),
ThinkingTextMessageEnd(ThinkingTextMessageEndEvent),
ToolCallStart(ToolCallStartEvent),
ToolCallArgs(ToolCallArgsEvent),
ToolCallEnd(ToolCallEndEvent),
ToolCallChunk(ToolCallChunkEvent),
ToolCallResult(ToolCallResultEvent),
ThinkingStart(ThinkingStartEvent),
ThinkingEnd(ThinkingEndEvent),
StateSnapshot(StateSnapshotEvent<StateT>),
StateDelta(StateDeltaEvent),
MessagesSnapshot(MessagesSnapshotEvent),
ActivitySnapshot(ActivitySnapshotEvent),
ActivityDelta(ActivityDeltaEvent),
Raw(RawEvent),
Custom(CustomEvent),
RunStarted(RunStartedEvent),
RunFinished(RunFinishedEvent),
RunError(RunErrorEvent),
StepStarted(StepStartedEvent),
StepFinished(StepFinishedEvent),
}
impl<StateT: AgentState> Event<StateT> {
pub fn event_type(&self) -> EventType {
match self {
Event::TextMessageStart(_) => EventType::TextMessageStart,
Event::TextMessageContent(_) => EventType::TextMessageContent,
Event::TextMessageEnd(_) => EventType::TextMessageEnd,
Event::TextMessageChunk(_) => EventType::TextMessageChunk,
Event::ThinkingTextMessageStart(_) => EventType::ThinkingTextMessageStart,
Event::ThinkingTextMessageContent(_) => EventType::ThinkingTextMessageContent,
Event::ThinkingTextMessageEnd(_) => EventType::ThinkingTextMessageEnd,
Event::ToolCallStart(_) => EventType::ToolCallStart,
Event::ToolCallArgs(_) => EventType::ToolCallArgs,
Event::ToolCallEnd(_) => EventType::ToolCallEnd,
Event::ToolCallChunk(_) => EventType::ToolCallChunk,
Event::ToolCallResult(_) => EventType::ToolCallResult,
Event::ThinkingStart(_) => EventType::ThinkingStart,
Event::ThinkingEnd(_) => EventType::ThinkingEnd,
Event::StateSnapshot(_) => EventType::StateSnapshot,
Event::StateDelta(_) => EventType::StateDelta,
Event::MessagesSnapshot(_) => EventType::MessagesSnapshot,
Event::ActivitySnapshot(_) => EventType::ActivitySnapshot,
Event::ActivityDelta(_) => EventType::ActivityDelta,
Event::Raw(_) => EventType::Raw,
Event::Custom(_) => EventType::Custom,
Event::RunStarted(_) => EventType::RunStarted,
Event::RunFinished(_) => EventType::RunFinished,
Event::RunError(_) => EventType::RunError,
Event::StepStarted(_) => EventType::StepStarted,
Event::StepFinished(_) => EventType::StepFinished,
}
}
pub fn timestamp(&self) -> Option<f64> {
match self {
Event::TextMessageStart(e) => e.base.timestamp,
Event::TextMessageContent(e) => e.base.timestamp,
Event::TextMessageEnd(e) => e.base.timestamp,
Event::TextMessageChunk(e) => e.base.timestamp,
Event::ThinkingTextMessageStart(e) => e.base.timestamp,
Event::ThinkingTextMessageContent(e) => e.base.timestamp,
Event::ThinkingTextMessageEnd(e) => e.base.timestamp,
Event::ToolCallStart(e) => e.base.timestamp,
Event::ToolCallArgs(e) => e.base.timestamp,
Event::ToolCallEnd(e) => e.base.timestamp,
Event::ToolCallChunk(e) => e.base.timestamp,
Event::ToolCallResult(e) => e.base.timestamp,
Event::ThinkingStart(e) => e.base.timestamp,
Event::ThinkingEnd(e) => e.base.timestamp,
Event::StateSnapshot(e) => e.base.timestamp,
Event::StateDelta(e) => e.base.timestamp,
Event::MessagesSnapshot(e) => e.base.timestamp,
Event::ActivitySnapshot(e) => e.base.timestamp,
Event::ActivityDelta(e) => e.base.timestamp,
Event::Raw(e) => e.base.timestamp,
Event::Custom(e) => e.base.timestamp,
Event::RunStarted(e) => e.base.timestamp,
Event::RunFinished(e) => e.base.timestamp,
Event::RunError(e) => e.base.timestamp,
Event::StepStarted(e) => e.base.timestamp,
Event::StepFinished(e) => e.base.timestamp,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_event_type_serialization() {
let event_type = EventType::TextMessageStart;
let json = serde_json::to_string(&event_type).unwrap();
assert_eq!(json, "\"TEXT_MESSAGE_START\"");
let event_type = EventType::ToolCallArgs;
let json = serde_json::to_string(&event_type).unwrap();
assert_eq!(json, "\"TOOL_CALL_ARGS\"");
let event_type = EventType::StateSnapshot;
let json = serde_json::to_string(&event_type).unwrap();
assert_eq!(json, "\"STATE_SNAPSHOT\"");
}
#[test]
fn test_event_type_deserialization() {
let event_type: EventType = serde_json::from_str("\"RUN_STARTED\"").unwrap();
assert_eq!(event_type, EventType::RunStarted);
let event_type: EventType = serde_json::from_str("\"THINKING_TEXT_MESSAGE_CONTENT\"").unwrap();
assert_eq!(event_type, EventType::ThinkingTextMessageContent);
}
#[test]
fn test_event_type_as_str() {
assert_eq!(EventType::TextMessageStart.as_str(), "TEXT_MESSAGE_START");
assert_eq!(EventType::RunFinished.as_str(), "RUN_FINISHED");
assert_eq!(EventType::Custom.as_str(), "CUSTOM");
}
#[test]
fn test_event_type_display() {
assert_eq!(format!("{}", EventType::TextMessageStart), "TEXT_MESSAGE_START");
assert_eq!(format!("{}", EventType::StateDelta), "STATE_DELTA");
}
#[test]
fn test_base_event_serialization() {
let event = BaseEvent {
timestamp: Some(1706123456789.0),
raw_event: None,
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"timestamp\":1706123456789.0"));
assert!(!json.contains("rawEvent")); }
#[test]
fn test_base_event_with_raw_event() {
let event = BaseEvent {
timestamp: None,
raw_event: Some(serde_json::json!({"provider": "openai"})),
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"rawEvent\""));
assert!(json.contains("\"provider\":\"openai\""));
}
#[test]
fn test_base_event_builder() {
let event = BaseEvent::new()
.timestamp(1234567890.0)
.raw_event(serde_json::json!({"test": true}));
assert_eq!(event.timestamp, Some(1234567890.0));
assert!(event.raw_event.is_some());
}
#[test]
fn test_event_validation_error_display() {
let error = EventValidationError::EmptyDelta;
assert_eq!(error.to_string(), "Delta must not be an empty string");
let error = EventValidationError::InvalidFormat("bad json".to_string());
assert_eq!(error.to_string(), "Invalid event format: bad json");
let error = EventValidationError::MissingField("message_id".to_string());
assert_eq!(error.to_string(), "Missing required field: message_id");
let error = EventValidationError::TypeMismatch {
expected: "TEXT_MESSAGE_START".to_string(),
actual: "RUN_STARTED".to_string(),
};
assert_eq!(
error.to_string(),
"Event type mismatch: expected TEXT_MESSAGE_START, got RUN_STARTED"
);
}
#[test]
fn test_event_validation_error_is_std_error() {
fn requires_error<E: std::error::Error>(_: E) {}
requires_error(EventValidationError::EmptyDelta);
}
#[test]
fn test_all_event_types_roundtrip() {
let all_types = [
EventType::TextMessageStart,
EventType::TextMessageContent,
EventType::TextMessageEnd,
EventType::TextMessageChunk,
EventType::ThinkingTextMessageStart,
EventType::ThinkingTextMessageContent,
EventType::ThinkingTextMessageEnd,
EventType::ToolCallStart,
EventType::ToolCallArgs,
EventType::ToolCallEnd,
EventType::ToolCallChunk,
EventType::ToolCallResult,
EventType::ThinkingStart,
EventType::ThinkingEnd,
EventType::StateSnapshot,
EventType::StateDelta,
EventType::MessagesSnapshot,
EventType::ActivitySnapshot,
EventType::ActivityDelta,
EventType::Raw,
EventType::Custom,
EventType::RunStarted,
EventType::RunFinished,
EventType::RunError,
EventType::StepStarted,
EventType::StepFinished,
];
for event_type in all_types {
let json = serde_json::to_string(&event_type).unwrap();
let parsed: EventType = serde_json::from_str(&json).unwrap();
assert_eq!(event_type, parsed);
}
}
#[test]
fn test_text_message_start_event() {
use crate::types::{MessageId, Role};
let event = TextMessageStartEvent::new(MessageId::random());
assert_eq!(event.role, Role::Assistant);
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"messageId\""));
assert!(json.contains("\"role\":\"assistant\""));
}
#[test]
fn test_text_message_start_event_with_timestamp() {
use crate::types::MessageId;
let event = TextMessageStartEvent::new(MessageId::random()).with_timestamp(1234567890.0);
assert_eq!(event.base.timestamp, Some(1234567890.0));
}
#[test]
fn test_text_message_content_event_validation() {
use crate::types::MessageId;
let result = TextMessageContentEvent::new(MessageId::random(), "Hello");
assert!(result.is_ok());
let result = TextMessageContentEvent::new(MessageId::random(), "");
assert!(matches!(result, Err(EventValidationError::EmptyDelta)));
}
#[test]
fn test_text_message_content_event_validate_method() {
use crate::types::MessageId;
let event = TextMessageContentEvent::new_unchecked(MessageId::random(), "");
assert!(matches!(event.validate(), Err(EventValidationError::EmptyDelta)));
let event = TextMessageContentEvent::new_unchecked(MessageId::random(), "Hello");
assert!(event.validate().is_ok());
}
#[test]
fn test_text_message_content_event_serialization() {
use crate::types::MessageId;
let event = TextMessageContentEvent::new(MessageId::random(), "Hello, world!").unwrap();
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"messageId\""));
assert!(json.contains("\"delta\":\"Hello, world!\""));
}
#[test]
fn test_text_message_end_event() {
use crate::types::MessageId;
let msg_id = MessageId::random();
let event = TextMessageEndEvent::new(msg_id.clone());
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"messageId\""));
}
#[test]
fn test_text_message_chunk_event() {
use crate::types::{MessageId, Role};
let event = TextMessageChunkEvent::new(Role::Assistant)
.with_message_id(MessageId::random())
.with_delta("chunk content");
assert!(event.message_id.is_some());
assert_eq!(event.delta, Some("chunk content".to_string()));
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"messageId\""));
assert!(json.contains("\"delta\":\"chunk content\""));
}
#[test]
fn test_text_message_chunk_event_skips_none() {
use crate::types::Role;
let event = TextMessageChunkEvent::new(Role::Assistant);
let json = serde_json::to_string(&event).unwrap();
assert!(!json.contains("\"messageId\""));
assert!(!json.contains("\"delta\""));
assert!(json.contains("\"role\":\"assistant\""));
}
#[test]
fn test_thinking_text_message_start_event() {
let event = ThinkingTextMessageStartEvent::new();
let json = serde_json::to_string(&event).unwrap();
assert_eq!(json, "{}");
}
#[test]
fn test_thinking_text_message_start_event_with_timestamp() {
let event = ThinkingTextMessageStartEvent::new().with_timestamp(1234567890.0);
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"timestamp\":1234567890.0"));
}
#[test]
fn test_thinking_text_message_content_event() {
let event = ThinkingTextMessageContentEvent::new("Let me think about this...");
assert_eq!(event.delta, "Let me think about this...");
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"delta\":\"Let me think about this...\""));
}
#[test]
fn test_thinking_text_message_content_event_allows_empty() {
let event = ThinkingTextMessageContentEvent::new("");
assert_eq!(event.delta, "");
}
#[test]
fn test_thinking_text_message_end_event() {
let event = ThinkingTextMessageEndEvent::new();
let json = serde_json::to_string(&event).unwrap();
assert_eq!(json, "{}");
}
#[test]
fn test_thinking_text_message_events_default() {
let start = ThinkingTextMessageStartEvent::default();
let end = ThinkingTextMessageEndEvent::default();
assert!(start.base.timestamp.is_none());
assert!(end.base.timestamp.is_none());
}
#[test]
fn test_tool_call_start_event() {
use crate::types::ToolCallId;
let event = ToolCallStartEvent::new(ToolCallId::random(), "get_weather");
assert_eq!(event.tool_call_name, "get_weather");
assert!(event.parent_message_id.is_none());
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"toolCallId\""));
assert!(json.contains("\"toolCallName\":\"get_weather\""));
assert!(!json.contains("parentMessageId")); }
#[test]
fn test_tool_call_start_event_with_parent() {
use crate::types::{MessageId, ToolCallId};
let event = ToolCallStartEvent::new(ToolCallId::random(), "get_weather")
.with_parent_message_id(MessageId::random());
assert!(event.parent_message_id.is_some());
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"parentMessageId\""));
}
#[test]
fn test_tool_call_args_event() {
use crate::types::ToolCallId;
let event = ToolCallArgsEvent::new(ToolCallId::random(), r#"{"location":"#);
assert_eq!(event.delta, r#"{"location":"#);
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"toolCallId\""));
assert!(json.contains("\"delta\""));
}
#[test]
fn test_tool_call_end_event() {
use crate::types::ToolCallId;
let event = ToolCallEndEvent::new(ToolCallId::random());
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"toolCallId\""));
}
#[test]
fn test_tool_call_chunk_event() {
use crate::types::ToolCallId;
let event = ToolCallChunkEvent::new()
.with_tool_call_id(ToolCallId::random())
.with_tool_call_name("search")
.with_delta(r#"{"query": "rust"}"#);
assert!(event.tool_call_id.is_some());
assert_eq!(event.tool_call_name, Some("search".to_string()));
assert!(event.delta.is_some());
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"toolCallId\""));
assert!(json.contains("\"toolCallName\":\"search\""));
assert!(json.contains("\"delta\""));
}
#[test]
fn test_tool_call_chunk_event_skips_none() {
let event = ToolCallChunkEvent::new();
let json = serde_json::to_string(&event).unwrap();
assert_eq!(json, "{}");
}
#[test]
fn test_tool_call_result_event() {
use crate::types::{MessageId, Role, ToolCallId};
let event = ToolCallResultEvent::new(
MessageId::random(),
ToolCallId::random(),
r#"{"weather": "sunny", "temp": 72}"#,
);
assert_eq!(event.role, Role::Tool);
assert!(event.content.contains("sunny"));
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"messageId\""));
assert!(json.contains("\"toolCallId\""));
assert!(json.contains("\"content\""));
assert!(json.contains("\"role\":\"tool\""));
}
#[test]
fn test_tool_call_result_event_deserialize_default_role() {
let json = r#"{"messageId":"550e8400-e29b-41d4-a716-446655440000","toolCallId":"6ba7b810-9dad-11d1-80b4-00c04fd430c8","content":"result"}"#;
let event: ToolCallResultEvent = serde_json::from_str(json).unwrap();
assert_eq!(event.role, Role::Tool);
}
#[test]
fn test_run_started_event() {
use crate::types::{RunId, ThreadId};
let event = RunStartedEvent::new(ThreadId::random(), RunId::random());
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"threadId\""));
assert!(json.contains("\"runId\""));
}
#[test]
fn test_run_finished_event() {
use crate::types::{RunId, ThreadId};
let event = RunFinishedEvent::new(ThreadId::random(), RunId::random());
assert!(event.result.is_none());
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"threadId\""));
assert!(json.contains("\"runId\""));
assert!(!json.contains("\"result\"")); }
#[test]
fn test_run_finished_event_with_result() {
use crate::types::{RunId, ThreadId};
let event = RunFinishedEvent::new(ThreadId::random(), RunId::random())
.with_result(serde_json::json!({"success": true}));
assert!(event.result.is_some());
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"result\""));
assert!(json.contains("\"success\":true"));
}
#[test]
fn test_run_error_event() {
let event = RunErrorEvent::new("Connection timeout");
assert_eq!(event.message, "Connection timeout");
assert!(event.code.is_none());
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"message\":\"Connection timeout\""));
assert!(!json.contains("\"code\"")); }
#[test]
fn test_run_error_event_with_code() {
let event = RunErrorEvent::new("Rate limit exceeded").with_code("RATE_LIMITED");
assert_eq!(event.code, Some("RATE_LIMITED".to_string()));
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"code\":\"RATE_LIMITED\""));
}
#[test]
fn test_run_finished_outcome_serialization() {
let success = RunFinishedOutcome::Success;
let interrupt = RunFinishedOutcome::Interrupt;
let success_json = serde_json::to_string(&success).unwrap();
let interrupt_json = serde_json::to_string(&interrupt).unwrap();
assert_eq!(success_json, "\"SUCCESS\"");
assert_eq!(interrupt_json, "\"INTERRUPT\"");
let deserialized: RunFinishedOutcome = serde_json::from_str("\"SUCCESS\"").unwrap();
assert_eq!(deserialized, RunFinishedOutcome::Success);
let deserialized: RunFinishedOutcome = serde_json::from_str("\"INTERRUPT\"").unwrap();
assert_eq!(deserialized, RunFinishedOutcome::Interrupt);
}
#[test]
fn test_run_finished_outcome_default() {
let outcome = RunFinishedOutcome::default();
assert_eq!(outcome, RunFinishedOutcome::Success);
}
#[test]
fn test_interrupt_info_empty() {
let info = InterruptInfo::new();
assert!(info.id.is_none());
assert!(info.reason.is_none());
assert!(info.payload.is_none());
let json = serde_json::to_string(&info).unwrap();
assert_eq!(json, "{}");
}
#[test]
fn test_interrupt_info_with_all_fields() {
let info = InterruptInfo::new()
.with_id("approval-001")
.with_reason("human_approval")
.with_payload(serde_json::json!({"action": "delete", "rows": 42}));
assert_eq!(info.id, Some("approval-001".to_string()));
assert_eq!(info.reason, Some("human_approval".to_string()));
assert!(info.payload.is_some());
let json = serde_json::to_string(&info).unwrap();
assert!(json.contains("\"id\":\"approval-001\""));
assert!(json.contains("\"reason\":\"human_approval\""));
assert!(json.contains("\"action\":\"delete\""));
}
#[test]
fn test_run_finished_event_with_interrupt() {
use crate::types::{RunId, ThreadId};
let event = RunFinishedEvent::new(ThreadId::random(), RunId::random())
.with_interrupt(
InterruptInfo::new()
.with_reason("human_approval")
.with_payload(serde_json::json!({"proposal": "send email"}))
);
assert_eq!(event.outcome, Some(RunFinishedOutcome::Interrupt));
assert!(event.interrupt.is_some());
assert!(event.result.is_none());
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"outcome\":\"INTERRUPT\""));
assert!(json.contains("\"interrupt\""));
assert!(json.contains("\"reason\":\"human_approval\""));
}
#[test]
fn test_run_finished_event_backward_compatibility() {
use crate::types::{RunId, ThreadId};
let event = RunFinishedEvent::new(ThreadId::random(), RunId::random())
.with_result(serde_json::json!({"done": true}));
assert!(event.outcome.is_none());
assert!(event.interrupt.is_none());
assert_eq!(event.effective_outcome(), RunFinishedOutcome::Success);
let json = serde_json::to_string(&event).unwrap();
assert!(!json.contains("\"outcome\"")); }
#[test]
fn test_run_finished_event_effective_outcome() {
use crate::types::{RunId, ThreadId};
let event1 = RunFinishedEvent::new(ThreadId::random(), RunId::random());
assert_eq!(event1.effective_outcome(), RunFinishedOutcome::Success);
let mut event2 = RunFinishedEvent::new(ThreadId::random(), RunId::random());
event2.interrupt = Some(InterruptInfo::new());
assert_eq!(event2.effective_outcome(), RunFinishedOutcome::Interrupt);
let event3 = RunFinishedEvent::new(ThreadId::random(), RunId::random())
.with_outcome(RunFinishedOutcome::Interrupt);
assert_eq!(event3.effective_outcome(), RunFinishedOutcome::Interrupt);
}
#[test]
fn test_step_started_event() {
let event = StepStartedEvent::new("process_input");
assert_eq!(event.step_name, "process_input");
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"stepName\":\"process_input\""));
}
#[test]
fn test_step_finished_event() {
let event = StepFinishedEvent::new("generate_response");
assert_eq!(event.step_name, "generate_response");
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"stepName\":\"generate_response\""));
}
#[test]
fn test_step_events_with_timestamp() {
let start = StepStartedEvent::new("step1").with_timestamp(1234567890.0);
let end = StepFinishedEvent::new("step1").with_timestamp(1234567891.0);
assert_eq!(start.base.timestamp, Some(1234567890.0));
assert_eq!(end.base.timestamp, Some(1234567891.0));
}
#[test]
fn test_state_snapshot_event() {
let event = StateSnapshotEvent::new(serde_json::json!({"count": 42}));
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"snapshot\""));
assert!(json.contains("\"count\":42"));
}
#[test]
fn test_state_snapshot_event_default() {
let event: StateSnapshotEvent<()> = StateSnapshotEvent::default();
assert!(event.base.timestamp.is_none());
}
#[test]
fn test_state_delta_event() {
let patches = vec![
serde_json::json!({"op": "replace", "path": "/count", "value": 43}),
serde_json::json!({"op": "add", "path": "/new_field", "value": "hello"}),
];
let event = StateDeltaEvent::new(patches);
assert_eq!(event.delta.len(), 2);
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"delta\""));
assert!(json.contains("\"op\":\"replace\""));
}
#[test]
fn test_state_delta_event_default() {
let event = StateDeltaEvent::default();
assert!(event.delta.is_empty());
}
#[test]
fn test_messages_snapshot_event() {
use crate::types::{Message, MessageId};
let messages = vec![
Message::User {
id: MessageId::random(),
content: "Hello".to_string(),
name: None,
},
Message::Assistant {
id: MessageId::random(),
content: Some("Hi there!".to_string()),
name: None,
tool_calls: None,
},
];
let event = MessagesSnapshotEvent::new(messages);
assert_eq!(event.messages.len(), 2);
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"messages\""));
}
#[test]
fn test_messages_snapshot_event_default() {
let event = MessagesSnapshotEvent::default();
assert!(event.messages.is_empty());
}
#[test]
fn test_thinking_start_event() {
let event = ThinkingStartEvent::new();
assert!(event.title.is_none());
let json = serde_json::to_string(&event).unwrap();
assert!(!json.contains("\"title\"")); }
#[test]
fn test_thinking_start_event_with_title() {
let event = ThinkingStartEvent::new().with_title("Analyzing query");
assert_eq!(event.title, Some("Analyzing query".to_string()));
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"title\":\"Analyzing query\""));
}
#[test]
fn test_thinking_end_event() {
let event = ThinkingEndEvent::new();
let json = serde_json::to_string(&event).unwrap();
assert_eq!(json, "{}");
}
#[test]
fn test_thinking_step_events_default() {
let start = ThinkingStartEvent::default();
let end = ThinkingEndEvent::default();
assert!(start.title.is_none());
assert!(end.base.timestamp.is_none());
}
#[test]
fn test_raw_event() {
let event = RawEvent::new(serde_json::json!({"provider_data": "openai"}));
assert!(event.source.is_none());
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"event\""));
assert!(json.contains("\"provider_data\":\"openai\""));
assert!(!json.contains("\"source\"")); }
#[test]
fn test_raw_event_with_source() {
let event = RawEvent::new(serde_json::json!({})).with_source("anthropic");
assert_eq!(event.source, Some("anthropic".to_string()));
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"source\":\"anthropic\""));
}
#[test]
fn test_custom_event() {
let event = CustomEvent::new("user_action", serde_json::json!({"clicked": "button"}));
assert_eq!(event.name, "user_action");
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"name\":\"user_action\""));
assert!(json.contains("\"value\""));
assert!(json.contains("\"clicked\":\"button\""));
}
#[test]
fn test_event_enum_serialization() {
use crate::types::MessageId;
let event: Event = Event::TextMessageStart(TextMessageStartEvent::new(
MessageId::random(),
));
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"type\":\"TEXT_MESSAGE_START\""));
assert!(json.contains("\"messageId\""));
assert!(json.contains("\"role\":\"assistant\""));
}
#[test]
fn test_event_enum_deserialization() {
let json = r#"{"type":"RUN_ERROR","message":"Test error"}"#;
let event: Event = serde_json::from_str(json).unwrap();
match event {
Event::RunError(e) => assert_eq!(e.message, "Test error"),
_ => panic!("Expected RunError variant"),
}
}
#[test]
fn test_event_type_method() {
use crate::types::MessageId;
let event: Event = Event::TextMessageEnd(TextMessageEndEvent::new(MessageId::random()));
assert_eq!(event.event_type(), EventType::TextMessageEnd);
let event: Event = Event::RunStarted(RunStartedEvent::new(
crate::types::ThreadId::random(),
crate::types::RunId::random(),
));
assert_eq!(event.event_type(), EventType::RunStarted);
let event: Event = Event::Custom(CustomEvent::new("test", serde_json::json!({})));
assert_eq!(event.event_type(), EventType::Custom);
}
#[test]
fn test_event_timestamp_method() {
use crate::types::MessageId;
let event: Event = Event::TextMessageStart(
TextMessageStartEvent::new(MessageId::random())
.with_timestamp(1234567890.0),
);
assert_eq!(event.timestamp(), Some(1234567890.0));
let event: Event = Event::ThinkingEnd(ThinkingEndEvent::new());
assert_eq!(event.timestamp(), None);
}
#[test]
fn test_event_all_variants_serialize() {
use crate::types::{Message, MessageId, RunId, ThreadId, ToolCallId};
let events: Vec<Event> = vec![
Event::TextMessageStart(TextMessageStartEvent::new(MessageId::random())),
Event::TextMessageContent(TextMessageContentEvent::new_unchecked(MessageId::random(), "Hello")),
Event::TextMessageEnd(TextMessageEndEvent::new(MessageId::random())),
Event::TextMessageChunk(TextMessageChunkEvent::new(Role::Assistant).with_delta("Hi")),
Event::ThinkingTextMessageStart(ThinkingTextMessageStartEvent::new()),
Event::ThinkingTextMessageContent(ThinkingTextMessageContentEvent::new("thinking...")),
Event::ThinkingTextMessageEnd(ThinkingTextMessageEndEvent::new()),
Event::ToolCallStart(ToolCallStartEvent::new(ToolCallId::random(), "test_tool")),
Event::ToolCallArgs(ToolCallArgsEvent::new(ToolCallId::random(), "{}")),
Event::ToolCallEnd(ToolCallEndEvent::new(ToolCallId::random())),
Event::ToolCallChunk(ToolCallChunkEvent::new()),
Event::ToolCallResult(ToolCallResultEvent::new(MessageId::random(), ToolCallId::random(), "result")),
Event::ThinkingStart(ThinkingStartEvent::new()),
Event::ThinkingEnd(ThinkingEndEvent::new()),
Event::StateSnapshot(StateSnapshotEvent::new(serde_json::json!({}))),
Event::StateDelta(StateDeltaEvent::new(vec![])),
Event::MessagesSnapshot(MessagesSnapshotEvent::new(vec![Message::Assistant {
id: MessageId::random(),
content: Some("Hi".to_string()),
name: None,
tool_calls: None,
}])),
Event::ActivitySnapshot(ActivitySnapshotEvent::new(MessageId::random(), "PLAN", serde_json::json!({"steps": []}))),
Event::ActivityDelta(ActivityDeltaEvent::new(MessageId::random(), "PLAN", vec![serde_json::json!({"op": "add", "path": "/steps/-", "value": "test"})])),
Event::Raw(RawEvent::new(serde_json::json!({}))),
Event::Custom(CustomEvent::new("test", serde_json::json!({}))),
Event::RunStarted(RunStartedEvent::new(ThreadId::random(), RunId::random())),
Event::RunFinished(RunFinishedEvent::new(ThreadId::random(), RunId::random())),
Event::RunError(RunErrorEvent::new("error")),
Event::StepStarted(StepStartedEvent::new("step")),
Event::StepFinished(StepFinishedEvent::new("step")),
];
for event in events {
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"type\":"));
let deserialized: Event = serde_json::from_str(&json).unwrap();
assert_eq!(event.event_type(), deserialized.event_type());
}
}
}