Skip to main content

ferrin_core/stream_text/
events.rs

1//! Events emitted by `stream_text`.
2
3use ferrin_spec::CustomKind;
4use ferrin_spec::FinishReason;
5use ferrin_spec::JsonValue;
6use ferrin_spec::PartId;
7use ferrin_spec::ProviderMetadata;
8use ferrin_spec::ToolCallId;
9use ferrin_spec::ToolName;
10use ferrin_spec::Usage;
11use ferrin_spec::Warning;
12use ferrin_spec::language_model::Source;
13use serde::Deserialize;
14use serde::Serialize;
15
16use crate::error::Error;
17use crate::error::ErrorKind;
18use crate::generate_text::GeneratedFile;
19use crate::generate_text::ParsedToolCall;
20use crate::generate_text::StepPerformance;
21use crate::generate_text::StepRequest;
22use crate::generate_text::StepResponse;
23use crate::generate_text::ToolApprovalRequestContent;
24use crate::generate_text::ToolApprovalResponseContent;
25use crate::generate_text::ToolExecutionError;
26use crate::generate_text::ToolOutputDenied;
27use crate::generate_text::ToolResult;
28use crate::telemetry::ModelIdentity;
29
30/// One event of a text stream. Serializable, so applications can forward
31/// events over SSE or WebSocket frames unchanged.
32#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
33#[serde(tag = "type", rename_all = "kebab-case")]
34#[non_exhaustive]
35#[allow(
36    clippy::large_enum_variant,
37    reason = "events flow through channels one at a time; flat variants keep pattern matching simple"
38)]
39pub enum StreamEvent {
40    /// The call started.
41    Start {
42        /// Call id.
43        call_id: String,
44    },
45    /// A step started.
46    StartStep {
47        /// Zero-based step index.
48        step_number: u32,
49        /// Application state used for this step.
50        #[serde(default, skip_serializing_if = "Option::is_none")]
51        runtime_context: Option<JsonValue>,
52        /// Shared tool context used for this step.
53        #[serde(default, skip_serializing_if = "Option::is_none")]
54        tools_context: Option<JsonValue>,
55        /// The model handling the step.
56        model: ModelIdentity,
57        /// Request metadata.
58        request: StepRequest,
59        /// Warnings from the adapter.
60        warnings: Vec<Warning>,
61    },
62    /// Start of a text part.
63    TextStart {
64        /// Part id.
65        id: PartId,
66        /// Provider-specific metadata.
67        #[serde(default, skip_serializing_if = "Option::is_none")]
68        provider_metadata: Option<ProviderMetadata>,
69    },
70    /// Text increment.
71    TextDelta {
72        /// Part id.
73        id: PartId,
74        /// Appended text.
75        text: String,
76        /// Provider-specific metadata.
77        #[serde(default, skip_serializing_if = "Option::is_none")]
78        provider_metadata: Option<ProviderMetadata>,
79    },
80    /// End of a text part.
81    TextEnd {
82        /// Part id.
83        id: PartId,
84        /// Provider-specific metadata.
85        #[serde(default, skip_serializing_if = "Option::is_none")]
86        provider_metadata: Option<ProviderMetadata>,
87    },
88    /// Start of a reasoning part.
89    ReasoningStart {
90        /// Part id.
91        id: PartId,
92        /// Provider-specific metadata.
93        #[serde(default, skip_serializing_if = "Option::is_none")]
94        provider_metadata: Option<ProviderMetadata>,
95    },
96    /// Reasoning increment.
97    ReasoningDelta {
98        /// Part id.
99        id: PartId,
100        /// Appended text.
101        text: String,
102        /// Provider-specific metadata.
103        #[serde(default, skip_serializing_if = "Option::is_none")]
104        provider_metadata: Option<ProviderMetadata>,
105    },
106    /// End of a reasoning part.
107    ReasoningEnd {
108        /// Part id.
109        id: PartId,
110        /// Provider-specific metadata.
111        #[serde(default, skip_serializing_if = "Option::is_none")]
112        provider_metadata: Option<ProviderMetadata>,
113    },
114    /// A reasoning artifact stored as a file.
115    ReasoningFile(GeneratedFile),
116    /// A generated file.
117    File(GeneratedFile),
118    /// A citation source.
119    Source(Source),
120    /// Provider-specific content.
121    Custom {
122        /// Kind of the content.
123        kind: CustomKind,
124        /// Provider-specific metadata.
125        #[serde(default, skip_serializing_if = "Option::is_none")]
126        provider_metadata: Option<ProviderMetadata>,
127    },
128    /// Start of streamed tool input.
129    ToolInputStart {
130        /// Tool call id.
131        id: ToolCallId,
132        /// Tool name.
133        tool_name: ToolName,
134        /// Whether the provider executes the tool.
135        #[serde(default, skip_serializing_if = "std::ops::Not::not")]
136        provider_executed: bool,
137        /// Whether the tool is dynamic.
138        #[serde(default, skip_serializing_if = "std::ops::Not::not")]
139        dynamic: bool,
140        /// Display title.
141        #[serde(default, skip_serializing_if = "Option::is_none")]
142        title: Option<String>,
143        /// Provider-specific metadata.
144        #[serde(default, skip_serializing_if = "Option::is_none")]
145        provider_metadata: Option<ProviderMetadata>,
146    },
147    /// Tool input increment (raw JSON text).
148    ToolInputDelta {
149        /// Tool call id.
150        id: ToolCallId,
151        /// Appended JSON text.
152        delta: String,
153        /// Provider-specific metadata.
154        #[serde(default, skip_serializing_if = "Option::is_none")]
155        provider_metadata: Option<ProviderMetadata>,
156    },
157    /// End of streamed tool input.
158    ToolInputEnd {
159        /// Tool call id.
160        id: ToolCallId,
161        /// Provider-specific metadata.
162        #[serde(default, skip_serializing_if = "Option::is_none")]
163        provider_metadata: Option<ProviderMetadata>,
164    },
165    /// A parsed tool call.
166    ToolCall(ParsedToolCall),
167    /// A tool result (may be preliminary).
168    ToolResult(ToolResult),
169    /// A tool error.
170    ToolError(ToolExecutionError),
171    /// A request to approve a tool call.
172    ToolApprovalRequest(ToolApprovalRequestContent),
173    /// An automatic approval decision.
174    ToolApprovalResponse(ToolApprovalResponseContent),
175    /// A denied tool call.
176    ToolOutputDenied(ToolOutputDenied),
177    /// A step finished.
178    FinishStep {
179        /// Zero-based step index.
180        step_number: u32,
181        /// Finish reason.
182        finish_reason: FinishReason,
183        /// Usage of the step.
184        usage: Usage,
185        /// Response metadata (without messages).
186        response: StepResponse,
187        /// Provider-specific metadata.
188        #[serde(default, skip_serializing_if = "Option::is_none")]
189        provider_metadata: Option<ProviderMetadata>,
190        /// Timing and throughput of the step.
191        performance: StepPerformance,
192    },
193    /// The call finished.
194    Finish {
195        /// Finish reason of the last step.
196        finish_reason: FinishReason,
197        /// Usage summed over all steps.
198        total_usage: Usage,
199    },
200    /// An error occurred (the stream may continue after a retry).
201    Error {
202        /// The error.
203        error: StreamErrorInfo,
204    },
205    /// The call was aborted by cancellation.
206    Abort,
207    /// The model call of the current step was retried after a stream error.
208    /// Content emitted by the previous attempt stays in the stream but is
209    /// excluded from the step result.
210    RetryAttempt {
211        /// Zero-based step index.
212        step_number: u32,
213        /// One-based retry attempt number.
214        attempt: u32,
215        /// Request metadata of the new attempt.
216        request: StepRequest,
217        /// Warnings from the new attempt's stream start.
218        warnings: Vec<Warning>,
219    },
220    /// A raw provider chunk (only with `include_raw_chunks`).
221    Raw {
222        /// The chunk.
223        raw_value: JsonValue,
224    },
225}
226
227impl StreamEvent {
228    /// Returns the wire name of the variant.
229    #[must_use]
230    pub fn kind_name(&self) -> &'static str {
231        match self {
232            Self::Start { .. } => "start",
233            Self::StartStep { .. } => "start-step",
234            Self::TextStart { .. } => "text-start",
235            Self::TextDelta { .. } => "text-delta",
236            Self::TextEnd { .. } => "text-end",
237            Self::ReasoningStart { .. } => "reasoning-start",
238            Self::ReasoningDelta { .. } => "reasoning-delta",
239            Self::ReasoningEnd { .. } => "reasoning-end",
240            Self::ReasoningFile(_) => "reasoning-file",
241            Self::File(_) => "file",
242            Self::Source(_) => "source",
243            Self::Custom { .. } => "custom",
244            Self::ToolInputStart { .. } => "tool-input-start",
245            Self::ToolInputDelta { .. } => "tool-input-delta",
246            Self::ToolInputEnd { .. } => "tool-input-end",
247            Self::ToolCall(_) => "tool-call",
248            Self::ToolResult(_) => "tool-result",
249            Self::ToolError(_) => "tool-error",
250            Self::ToolApprovalRequest(_) => "tool-approval-request",
251            Self::ToolApprovalResponse(_) => "tool-approval-response",
252            Self::ToolOutputDenied(_) => "tool-output-denied",
253            Self::FinishStep { .. } => "finish-step",
254            Self::Finish { .. } => "finish",
255            Self::Error { .. } => "error",
256            Self::Abort => "abort",
257            Self::RetryAttempt { .. } => "retry-attempt",
258            Self::Raw { .. } => "raw",
259        }
260    }
261
262    /// Returns the text of a text delta.
263    #[must_use]
264    pub fn as_text_delta(&self) -> Option<&str> {
265        match self {
266            Self::TextDelta { text, .. } => Some(text),
267            _ => None,
268        }
269    }
270}
271
272/// Serializable projection of an [`Error`] carried in a stream.
273#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
274pub struct StreamErrorInfo {
275    /// Coarse category.
276    pub kind: ErrorKind,
277    /// Display message.
278    pub message: String,
279    /// HTTP status when known.
280    #[serde(default, skip_serializing_if = "Option::is_none")]
281    pub status_code: Option<u16>,
282    /// Whether retrying may succeed.
283    pub is_retryable: bool,
284    /// Provider error payload when available.
285    #[serde(default, skip_serializing_if = "Option::is_none")]
286    pub provider_data: Option<JsonValue>,
287}
288
289impl StreamErrorInfo {
290    /// Projects an error.
291    #[must_use]
292    pub fn from_error(error: &Error) -> Self {
293        let provider_data = error
294            .as_provider()
295            .and_then(ferrin_spec::error::ProviderError::as_api_call)
296            .and_then(|api| api.data.clone());
297        Self {
298            kind: error.kind(),
299            message: error.to_string(),
300            status_code: error.status_code().map(|status| status.as_u16()),
301            is_retryable: error.is_retryable(),
302            provider_data,
303        }
304    }
305}
306
307impl From<&Error> for StreamErrorInfo {
308    fn from(error: &Error) -> Self {
309        Self::from_error(error)
310    }
311}