use std::sync::Arc;
use serde::Deserialize;
use serde_json::value::RawValue;
use crate::CostStatus;
use crate::{AgentEventKind, EstimatedUsdCost, MessagePhase, ToolOutputBody, Usage};
#[derive(Clone, Debug)]
#[non_exhaustive]
pub enum AgentEventData {
OpenAi(OpenAiEvent),
Assistant(AssistantEvent),
Reasoning(ReasoningEvent),
Run(RunEvent),
Tool(ToolEvent),
Model(ModelEvent),
Context(ContextEvent),
Transport(TransportEvent),
}
#[derive(Clone, Debug, Deserialize)]
pub struct OpenAiEvent {
pub direction: String,
pub transport: String,
pub phase: String,
pub model_call_index: Option<u32>,
pub event: Box<RawValue>,
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub enum AssistantEvent {
Delta(AssistantDelta),
Message(AssistantMessage),
}
#[derive(Clone, Debug, Deserialize)]
pub struct AssistantDelta {
pub model_call_index: u32,
pub item_id: Option<String>,
pub phase: Option<MessagePhase>,
pub text: String,
}
#[derive(Clone, Debug, Deserialize)]
pub struct AssistantMessage {
pub model_call_index: u32,
pub item_id: Option<String>,
pub phase: Option<MessagePhase>,
pub text: String,
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub enum ReasoningEvent {
SummaryDelta(ReasoningSummaryDelta),
}
#[derive(Clone, Debug, Deserialize)]
pub struct ReasoningSummaryDelta {
pub model_call_index: u32,
pub text: String,
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub enum RunEvent {
Started(RunStarted),
Steered(RunSteered),
Error(RunError),
Completed(Box<RunTerminal>),
Failed(Box<RunTerminal>),
}
#[derive(Clone, Debug, Deserialize)]
pub struct RunStarted {
pub mode: String,
pub model: String,
pub reasoning_mode: String,
pub effort: String,
pub transport: String,
pub orchestration: String,
pub websocket_url: String,
pub workspace: Option<String>,
pub instruction_bytes: usize,
}
#[derive(Clone, Copy, Debug, Deserialize)]
pub struct RunSteered {
pub steer_index: u32,
pub instruction_bytes: usize,
}
#[derive(Clone, Debug, Deserialize)]
pub struct RunError {
pub message: String,
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum RunStatus {
Completed,
Cancelled,
Failed,
}
#[derive(Clone, Copy, Debug, Default, Deserialize)]
pub struct EventUsage {
pub input_tokens: u64,
pub cached_input_tokens: u64,
pub cache_write_input_tokens: u64,
pub output_tokens: u64,
pub reasoning_output_tokens: u64,
pub total_tokens: u64,
}
#[derive(Clone, Debug, Default, Deserialize)]
pub struct RunMetrics {
pub model_calls: u32,
pub steers: u32,
pub compactions: u32,
pub tool_calls: u32,
pub connection_attempts: u32,
pub websocket_reconnects: u32,
pub response_attempts: u32,
pub response_retries: u32,
pub connection_duration_ns: u64,
pub retry_backoff_duration_ns: u64,
pub model_duration_ns: u64,
pub compaction_duration_ns: u64,
pub warmup_duration_ns: u64,
pub tool_work_duration_ns: u64,
pub tool_wall_duration_ns: u64,
pub usage: EventUsage,
pub warmup_usage: EventUsage,
}
#[derive(Clone, Debug, Deserialize)]
pub struct RunTerminal {
pub status: RunStatus,
pub model: String,
pub reasoning_mode: String,
pub effort: String,
pub transport: String,
pub orchestration: String,
pub duration_ms: u64,
pub duration_ns: u64,
#[serde(flatten)]
pub metrics: RunMetrics,
pub estimated_cost: Option<EstimatedUsdCost>,
pub cost_usd: Option<f64>,
#[serde(default)]
pub cost_status: CostStatus,
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub enum ToolEvent {
Call(ToolCall),
Result(ToolResultEvent),
}
#[derive(Clone, Debug, Deserialize)]
pub struct ToolCall {
pub call_id: String,
pub tool: String,
pub arguments: Box<RawValue>,
pub model_call_index: u32,
}
impl ToolCall {
pub fn decode_arguments<T: serde::de::DeserializeOwned>(&self) -> Result<T, serde_json::Error> {
serde_json::from_str(self.arguments.get())
}
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum ToolStatus {
Completed,
Failed,
Cancelled,
}
#[derive(Clone, Debug, Deserialize)]
pub struct ToolResultEvent {
pub call_id: String,
pub tool: String,
pub status: ToolStatus,
pub duration_ns: u64,
pub started_after_ns: Option<u64>,
pub result: ToolOutputBody,
pub metadata: Option<Box<RawValue>>,
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub enum ModelEvent {
WarmupStarted(ModelWarmupStarted),
WarmupCompleted(ModelWarmupCompleted),
WarmupFailed(ModelWarmupFailed),
CallStarted(ModelCallStarted),
CallCompleted(ModelCallCompleted),
CallFailed(ModelCallFailed),
}
#[derive(Clone, Debug, Deserialize)]
pub struct ModelWarmupStarted {
pub model: String,
pub prompt_cache_key: String,
}
#[derive(Clone, Debug, Deserialize)]
pub struct ModelWarmupCompleted {
pub source: String,
pub attempt: Option<u32>,
pub connection_generation: Option<u32>,
pub duration_ns: u64,
pub usage: Option<Usage>,
}
#[derive(Clone, Debug, Deserialize)]
pub struct ModelWarmupFailed {
pub duration_ns: u64,
pub error: String,
}
#[derive(Clone, Debug, Deserialize)]
pub struct ModelCallStarted {
pub call_index: u32,
pub model: String,
pub reasoning_mode: String,
pub effort: String,
}
#[derive(Clone, Debug, Deserialize)]
pub struct ModelCallCompleted {
pub call_index: u32,
pub model: String,
pub attempt: u32,
pub connection_generation: u32,
pub status: String,
pub duration_ns: u64,
pub time_to_first_event_ns: u64,
pub time_to_first_output_ns: Option<u64>,
pub tool_calls: usize,
pub usage: Option<Usage>,
}
#[derive(Clone, Debug, Deserialize)]
pub struct ModelCallFailed {
pub call_index: u32,
pub model: String,
pub duration_ns: u64,
pub error: String,
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub enum ContextEvent {
CompactionStarted(CompactionStarted),
CompactionCompleted(CompactionCompleted),
CompactionFailed(CompactionFailed),
}
#[derive(Clone, Debug, Deserialize)]
pub struct CompactionStarted {
pub after_model_call_index: u32,
pub active_context_tokens: u64,
pub auto_compact_token_limit: u64,
}
#[derive(Clone, Debug, Deserialize)]
pub struct CompactionCompleted {
pub after_model_call_index: u32,
pub attempt: u32,
pub connection_generation: u32,
pub status: String,
pub duration_ns: u64,
pub time_to_first_event_ns: u64,
pub time_to_first_output_ns: Option<u64>,
pub usage: Option<Usage>,
}
#[derive(Clone, Debug, Deserialize)]
pub struct CompactionFailed {
pub after_model_call_index: u32,
pub duration_ns: u64,
pub error: String,
}
#[derive(Clone, Debug)]
pub struct TransportEvent {
kind: AgentEventKind,
payload: Arc<RawValue>,
}
impl TransportEvent {
pub(crate) const fn new(kind: AgentEventKind, payload: Arc<RawValue>) -> Self {
Self { kind, payload }
}
#[must_use]
pub const fn kind(&self) -> AgentEventKind {
self.kind
}
#[must_use]
pub fn raw_payload(&self) -> &RawValue {
&self.payload
}
pub fn decode_payload<T: serde::de::DeserializeOwned>(&self) -> Result<T, serde_json::Error> {
serde_json::from_str(self.payload.get())
}
}