Skip to main content

rig_core/providers/openai/wire/
chat.rs

1//! Chat Completions request encoding, and one decoder for a whole reply and
2//! a stream of chunks.
3//!
4//! ```
5//! use rig_core::providers::openai::OpenAI;
6//! let wire = OpenAI::new("key").chat("gpt-5.2");
7//! ```
8
9use serde::{Deserialize, Serialize};
10
11use crate::completion::{CompletionRequest, FinishReason, ProviderCapabilities};
12use crate::error::EncodeError;
13use crate::error::ProviderError;
14use crate::observe::ObservedError;
15use crate::providers::internal::openai_chat_completions_compatible::{
16    drop_tool_calls_cut_by_budget, map_native_finish_reason, map_openai_finish_reason,
17    provider_error_envelope,
18};
19use crate::providers::internal::wire::classify_chat_completions_frame;
20use crate::providers::openai::completion::{
21    self as unary, AssistantContent, Message, ToolChoice, assistant_refusal_fallback,
22    is_openai_reasoning_model, request_body,
23};
24use crate::wire::{
25    AdapterEvent, AdapterUsage, AdapterVerdict, Body, Capabilities, Decoder, Descriptor, Encoded,
26    Framing, Mode, ObservationSink, Out, Wire, WireEvent, WireFrame,
27};
28
29use super::dto::{
30    ChatChoice, ChatFrame, ChatUsage, StreamingCompletionResponse, StreamingDelta, delta_text,
31};
32use super::{BodyRewrite, OpenAIConfig, OutputCap};
33
34/// The chat-completions wire: a provider configuration, a model, and the
35/// per-turn options the endpoint takes.
36#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
37pub struct Chat {
38    /// Which provider, and how to reach it.
39    pub provider: OpenAIConfig,
40    /// The model this wire addresses.
41    pub model: String,
42    /// Whether tool schemas are sanitized for OpenAI's strict mode:
43    /// `additionalProperties: false` on every object, every property
44    /// required, and `strict: true` on each function definition.
45    pub strict_tools: bool,
46    /// Whether tool-result messages serialize their content as arrays.
47    pub tool_result_array_content: bool,
48    /// Whether the request asks for provider-side prompt caching
49    /// (OpenRouter's ephemeral `cache_control` on the system prompt).
50    pub prompt_caching: bool,
51}
52
53impl Chat {
54    pub(crate) fn encode_with_headers(
55        &self,
56        request: CompletionRequest,
57        mode: Mode,
58        headers: impl FnOnce(
59            &OpenAIConfig,
60            &CompletionRequest,
61            http::request::Builder,
62        ) -> http::request::Builder,
63    ) -> Result<Encoded, EncodeError> {
64        let (request, issuers) =
65            super::scope_reasoning(&self.provider.dialect, &self.model, request)?;
66        let quirks = &self.provider.dialect.quirks;
67        // Azure's deployment URL remains pinned to the handle, not a request override.
68        let uri = self.provider.uri(
69            quirks.completion_path,
70            self.provider.deployment(&self.model),
71        );
72        let builder = headers(
73            &self.provider,
74            &request,
75            http::Request::post(uri).header("Content-Type", "application/json"),
76        );
77        if !quirks.accepts_file_ids {
78            refuse_file_ids(&request)?;
79        }
80        let mut typed = unary::CompletionRequest::try_from(unary::OpenAIRequestParams {
81            model: self.model.clone(),
82            request,
83            strict_tools: self.strict_tools,
84            tool_result_array_content: self.tool_result_array_content,
85            supports_response_format: quirks.supports_response_format,
86            response_format_with_tools: quirks.response_format_with_tools,
87            supports_tools: quirks.supports_tools,
88            supports_image_tool_results: quirks.supports_image_tool_results,
89            reasoning_details: quirks.reasoning_details,
90            issuers,
91        })?;
92        self.prepare(&mut typed)?;
93
94        // The resolved model, not the handle's: a per-request override
95        // changes which endpoint answers, so it decides the spelling too.
96        let modern_output_cap = match quirks.output_cap {
97            OutputCap::Legacy => false,
98            OutputCap::OpenAiReasoningFamilies => is_openai_reasoning_model(&typed.model),
99        };
100        let mut body = request_body(&typed, modern_output_cap)?;
101
102        if mode == Mode::Streaming {
103            if quirks.stream_include_usage {
104                // Preserve caller stream options, including an explicit include_usage value.
105                match body.get_mut("stream_options") {
106                    Some(serde_json::Value::Object(options)) => {
107                        options
108                            .entry("include_usage")
109                            .or_insert(serde_json::Value::Bool(true));
110                    }
111                    Some(_) => {}
112                    None => {
113                        body = crate::json_utils::merge(
114                            body,
115                            serde_json::json!({"stream_options": {"include_usage": true}}),
116                        );
117                    }
118                }
119            }
120            body = crate::json_utils::merge(body, serde_json::json!({"stream": true}));
121        }
122        self.finalize(&mut body)?;
123
124        crate::providers::internal::trace_json(
125            crate::providers::internal::LogTarget::Completions,
126            "OpenAI Chat Completions request",
127            &body,
128        );
129
130        let request = builder.body(Body::Bytes(serde_json::to_vec(&body)?))?;
131
132        let framing = match mode {
133            Mode::Streaming => Framing::Sse,
134            Mode::Unary => Framing::Whole,
135        };
136        Ok(Encoded::new(request, framing)
137            .with_request_id_header(self.provider.dialect.request_id_header)
138            .with_projection(ChatDecoder::project)
139            .with_route(Some(self.provider.dialect.quirks.completion_path)))
140    }
141
142    /// The wire for `model` on `provider`, with every option off.
143    pub fn new(provider: OpenAIConfig, model: impl Into<String>) -> Self {
144        Self {
145            provider,
146            model: model.into(),
147            strict_tools: false,
148            tool_result_array_content: false,
149            prompt_caching: false,
150        }
151    }
152
153    /// Sanitize tool schemas for OpenAI's strict mode, so the provider can
154    /// guarantee a tool call matches its schema exactly.
155    pub fn with_strict_tools(mut self) -> Self {
156        self.strict_tools = true;
157        self
158    }
159
160    /// Serialize tool-result content as arrays.
161    pub fn with_tool_result_array_content(mut self) -> Self {
162        self.tool_result_array_content = true;
163        self
164    }
165
166    /// Ask the provider to cache the prompt.
167    pub fn with_prompt_caching(mut self) -> Self {
168        self.prompt_caching = true;
169        self
170    }
171
172    /// Apply typed dialect rewrites, rejecting unsupported tool choices or parameters.
173    fn prepare(&self, request: &mut unary::CompletionRequest) -> Result<(), EncodeError> {
174        // Only OpenAI's own endpoint serves its reasoning families.
175        if matches!(
176            self.provider.dialect.quirks.output_cap,
177            OutputCap::OpenAiReasoningFamilies
178        ) {
179            refuse_tools_while_reasoning(request)?;
180        }
181        match self.provider.dialect.quirks.rewrite {
182            BodyRewrite::GroqCompoundTools => {
183                fold_groq_native_tools(request)?;
184                strip_assistant_reasoning(request);
185            }
186            BodyRewrite::LlamaCpp => {
187                if let Some(ToolChoice::Function { name }) = &request.tool_choice {
188                    return Err(EncodeError::request(format!(
189                        "llama.cpp cannot force a specific tool: `llama-server` accepts only \
190                         `auto`, `none` or `required` for tool_choice and silently treats \
191                         anything else as `auto`, so requesting `{name}` would return whichever \
192                         tool the model picked. Use `ToolChoice::Required` to force a call, or \
193                         advertise only `{name}` in `tools`."
194                    )));
195                }
196            }
197            BodyRewrite::Moonshot => steer_moonshot_tool_choice(request)?,
198            BodyRewrite::Mira => {
199                // The gateway rejects pass-through parameters.
200                if request.additional_params.take().is_some() {
201                    tracing::warn!(
202                        "Additional parameters are not supported by Mira and will be ignored"
203                    );
204                }
205            }
206            BodyRewrite::HuggingFaceRouter => {
207                // Some sub-providers (Fireworks) address models through a
208                // qualified identifier in the request body.
209                request.model = self.provider.route().model_identifier(&request.model);
210            }
211            BodyRewrite::None
212            | BodyRewrite::DeepSeek
213            | BodyRewrite::Perplexity
214            | BodyRewrite::Hyperbolic
215            | BodyRewrite::Mistral
216            | BodyRewrite::OpenRouter => {}
217        }
218        Ok(())
219    }
220
221    /// Apply dialect body rewrites after merging streaming parameters.
222    /// Return conversion errors for unsupported content.
223    fn finalize(&self, body: &mut serde_json::Value) -> Result<(), EncodeError> {
224        let Some(map) = body.as_object_mut() else {
225            return Ok(());
226        };
227        match self.provider.dialect.quirks.rewrite {
228            BodyRewrite::Perplexity => {
229                // Perplexity accepts only system/user/assistant roles with
230                // strict user/assistant alternation. Text-only content-part
231                // arrays flatten; arrays with non-text parts are left for
232                // its own multimodal handling on the sonar models.
233                if let Some(messages) = map.get_mut("messages").and_then(as_array_mut) {
234                    unary::sanitize_plain_text_history(messages, Some(("\n", true)), false, true);
235                }
236            }
237            BodyRewrite::Hyperbolic => {
238                // Strip tool-exchange remnants a shared history may carry;
239                // content-part arrays stay as-is for the vision models.
240                if let Some(messages) = map.get_mut("messages").and_then(as_array_mut) {
241                    unary::sanitize_plain_text_history(messages, None, false, false);
242                }
243            }
244            BodyRewrite::Mira => {
245                if let Some(messages) = map.get_mut("messages").and_then(as_array_mut) {
246                    unary::sanitize_plain_text_history(messages, Some(("\n", false)), true, false);
247                }
248            }
249            BodyRewrite::DeepSeek => finalize_deepseek(map),
250            BodyRewrite::Mistral => finalize_mistral(map)?,
251            BodyRewrite::OpenRouter => finalize_openrouter(map, self.prompt_caching),
252            BodyRewrite::None
253            | BodyRewrite::HuggingFaceRouter
254            | BodyRewrite::GroqCompoundTools
255            | BodyRewrite::LlamaCpp
256            | BodyRewrite::Moonshot => {}
257        }
258        Ok(())
259    }
260}
261
262/// GPT-6 models that call function tools on Chat Completions only at
263/// `reasoning_effort: "none"` (their model pages), and those that do not
264/// support `"none"` at all, so never call tools there.
265const TOOLS_ONLY_WITHOUT_REASONING: [&str; 2] = [unary::GPT_6_SOL, unary::GPT_6_LUNA];
266const NO_TOOLS_ON_CHAT: [&str; 2] = [unary::GPT_6_ASTRA, unary::GPT_6_1_SOL];
267
268/// Whether `model` is `id` or one of its dated snapshots (`<id>-YYYY-MM-DD`).
269fn is_model(model: &str, id: &str) -> bool {
270    model
271        .strip_prefix(id)
272        .is_some_and(|rest| rest.is_empty() || rest.starts_with("-20"))
273}
274
275/// Refuse a Chat Completions request with function tools on a GPT-6 model
276/// that would answer it with a 400, naming the fix. A caller who already
277/// sends `reasoning_effort: "none"` where the model takes it is not refused.
278fn refuse_tools_while_reasoning(request: &unary::CompletionRequest) -> Result<(), EncodeError> {
279    if request.tools.is_empty() {
280        return Ok(());
281    }
282    let model = request.model.as_str();
283    if NO_TOOLS_ON_CHAT.iter().any(|id| is_model(model, id)) {
284        return Err(EncodeError::request(format!(
285            "{model} cannot call function tools on Chat Completions: it takes them there only \
286             at reasoning_effort \"none\", which it does not support. Use the Responses wire."
287        )));
288    }
289    let effort_none = request
290        .additional_params
291        .as_ref()
292        .and_then(|params| params.get("reasoning_effort"))
293        .and_then(serde_json::Value::as_str)
294        == Some("none");
295    if TOOLS_ONLY_WITHOUT_REASONING
296        .iter()
297        .any(|id| is_model(model, id))
298        && !effort_none
299    {
300        return Err(EncodeError::request(format!(
301            "{model} calls function tools on Chat Completions only at reasoning_effort \"none\": \
302             send `\"reasoning_effort\": \"none\"` in additional_params, or use the Responses wire."
303        )));
304    }
305    Ok(())
306}
307
308fn as_array_mut(value: &mut serde_json::Value) -> Option<&mut Vec<serde_json::Value>> {
309    value.as_array_mut()
310}
311
312/// Groq delivers hidden reasoning under `reasoning` and rejects an assistant
313/// message that carries `reasoning_content` with a 400, so a reasoning turn
314/// replays without it. The reasoning is dropped from the replay rather than
315/// respelled: whether Groq accepts it back under `reasoning` is unverified.
316fn strip_assistant_reasoning(request: &mut unary::CompletionRequest) {
317    for message in &mut request.messages {
318        if let Message::Assistant { reasoning, .. } = message {
319            *reasoning = None;
320        }
321    }
322}
323
324/// Groq's compound-system native tools (`browser_search`, `code_interpreter`,
325/// …) arrive through `additional_params.tools`. Left there they would clobber
326/// the function-tool array on serialization, because `additional_params` is
327/// flattened into the body and the flattened key wins.
328fn fold_groq_native_tools(request: &mut unary::CompletionRequest) -> Result<(), EncodeError> {
329    let Some(map) = request
330        .additional_params
331        .as_mut()
332        .and_then(serde_json::Value::as_object_mut)
333    else {
334        return Ok(());
335    };
336    let Some(raw_tools) = map.remove("tools") else {
337        return Ok(());
338    };
339    let serde_json::Value::Array(native_tools) = raw_tools else {
340        return Err(EncodeError::request(
341            "Groq `additional_params.tools` must be an array of native tool objects",
342        ));
343    };
344
345    // `compound_custom.enabled_tools` is a set keyed by tool type, so a
346    // caller who names the same tool twice enables it once.
347    let enabled = map
348        .entry("compound_custom")
349        .or_insert_with(|| serde_json::json!({}))
350        .as_object_mut()
351        .map(|custom| {
352            custom
353                .entry("enabled_tools")
354                .or_insert_with(|| serde_json::Value::Array(Vec::new()))
355        });
356    let Some(serde_json::Value::Array(enabled)) = enabled else {
357        return Ok(());
358    };
359    for tool in native_tools {
360        let kind = tool.get("type").and_then(serde_json::Value::as_str);
361        let already_enabled = enabled
362            .iter()
363            .any(|existing| existing.get("type").and_then(serde_json::Value::as_str) == kind);
364        if !already_enabled {
365            enabled.push(tool);
366        }
367    }
368    Ok(())
369}
370
371/// Moonshot supports only `auto`/`none`: forcing one specific tool has no
372/// workaround, and `required` is steered with an extra user message.
373fn steer_moonshot_tool_choice(request: &mut unary::CompletionRequest) -> Result<(), EncodeError> {
374    if matches!(request.tool_choice, Some(ToolChoice::Function { .. })) {
375        return Err(EncodeError::request(
376            "Moonshot does not support forcing a specific tool".to_owned(),
377        ));
378    }
379    if matches!(request.tool_choice, Some(ToolChoice::Required)) {
380        tracing::warn!(
381            "Moonshot does not support tool_choice=required; coercing to auto with an \
382             additional steering message"
383        );
384        request.tool_choice = Some(ToolChoice::Auto);
385        request.messages.push(Message::User {
386            content: vec![unary::UserContent::Text {
387                text: "Please select a tool to handle the current issue.".to_owned(),
388            }],
389            name: None,
390        });
391    }
392    Ok(())
393}
394
395/// DeepSeek takes message `content` as a plain string, echoes tool calls back
396/// with an `index`, and needs an explicit empty `content` on a tool-call-only
397/// assistant turn.
398fn finalize_deepseek(map: &mut serde_json::Map<String, serde_json::Value>) {
399    if let Some(messages) = map.get_mut("messages").and_then(as_array_mut) {
400        for message in messages {
401            let Some(message) = message.as_object_mut() else {
402                continue;
403            };
404            let is_assistant =
405                message.get("role").and_then(serde_json::Value::as_str) == Some("assistant");
406
407            if let Some(content) = message.get_mut("content") {
408                let separator = if is_assistant { "" } else { "\n" };
409                // Preserve nontext parts so unsupported attachments are rejected
410                // rather than silently omitted from the prompt.
411                unary::flatten_text_content_parts(content, separator, true);
412            } else if is_assistant {
413                message.insert(
414                    "content".to_owned(),
415                    serde_json::Value::String(String::new()),
416                );
417            }
418
419            if is_assistant
420                && let Some(tool_calls) = message.get_mut("tool_calls").and_then(as_array_mut)
421            {
422                for tool_call in tool_calls {
423                    if let Some(tool_call) = tool_call.as_object_mut() {
424                        tool_call
425                            .entry("index")
426                            .or_insert_with(|| serde_json::json!(0));
427                    }
428                }
429            }
430        }
431    }
432
433    // DeepSeek rejects forced tool choices unless thinking is explicitly
434    // disabled; suppress them to an explicit `null` otherwise.
435    let thinking_disabled = map
436        .get("thinking")
437        .and_then(|thinking| thinking.get("type"))
438        .and_then(serde_json::Value::as_str)
439        .is_some_and(|mode| mode.eq_ignore_ascii_case("disabled"));
440    if !thinking_disabled
441        && let Some(tool_choice) = map.get_mut("tool_choice")
442        && (tool_choice.is_object() || tool_choice.as_str() == Some("required"))
443    {
444        *tool_choice = serde_json::Value::Null;
445    }
446}
447
448/// Mistral's wire-level differences: its own spelling for a forced tool
449/// choice, the relaxation that lets a structured format ride beside tools,
450/// and its assistant-message schema.
451///
452/// Its multimodal content mapping is [`mistral_content`].
453fn finalize_mistral(
454    map: &mut serde_json::Map<String, serde_json::Value>,
455) -> Result<(), EncodeError> {
456    // Mistral spells the "must call some tool" mode `any`, not `required`.
457    if let Some(tool_choice) = map.get_mut("tool_choice")
458        && tool_choice.as_str() == Some("required")
459    {
460        *tool_choice = serde_json::Value::String("any".to_owned());
461    }
462
463    // Mistral rejects forced tool calls beside JSON response formats.
464    // Relax the choice rather than discard the requested schema; text formats are exempt.
465    let forces_a_tool_call = map
466        .get("tool_choice")
467        .is_some_and(|choice| !matches!(choice.as_str(), Some("auto" | "none")));
468    let has_tools = map
469        .get("tools")
470        .and_then(serde_json::Value::as_array)
471        .is_some_and(|tools| !tools.is_empty());
472    let has_structured_format = map
473        .get("response_format")
474        .and_then(|format| format.get("type"))
475        .and_then(serde_json::Value::as_str)
476        .is_some_and(|kind| matches!(kind, "json_schema" | "json_object"));
477    if forces_a_tool_call && has_tools && has_structured_format {
478        tracing::debug!(
479            "relaxing tool_choice to `auto`: Mistral rejects a forced tool choice \
480             alongside a response format"
481        );
482        map.insert(
483            "tool_choice".to_owned(),
484            serde_json::Value::String("auto".to_owned()),
485        );
486    }
487
488    let Some(messages) = map.get_mut("messages").and_then(as_array_mut) else {
489        return Ok(());
490    };
491    for message in messages {
492        let Some(message) = message.as_object_mut() else {
493            continue;
494        };
495        let is_assistant =
496            message.get("role").and_then(serde_json::Value::as_str) == Some("assistant");
497
498        // Mistral takes text-only message `content` as a plain string and
499        // carries images, audio and documents as its own chunk array.
500        // Content it has no chunk for fails here rather than reaching the API
501        // with the part removed.
502        if let Some(content) = message.get_mut("content") {
503            mistral_content(content)?;
504        }
505
506        if is_assistant {
507            if !message.contains_key("content") {
508                message.insert(
509                    "content".to_owned(),
510                    serde_json::Value::String(String::new()),
511                );
512            }
513            // `prefix` is part of Mistral's assistant message schema.
514            message
515                .entry("prefix")
516                .or_insert(serde_json::Value::Bool(false));
517            // Mistral rejects unknown assistant fields; hidden reasoning
518            // cannot be echoed back.
519            message.remove("reasoning_content");
520        }
521    }
522    Ok(())
523}
524
525/// Mistral's text chunk tag.
526const MISTRAL_TEXT: &str = "text";
527/// Mistral's image chunk tag.
528const MISTRAL_IMAGE: &str = "image_url";
529/// Mistral's audio chunk tag.
530const MISTRAL_AUDIO: &str = "input_audio";
531/// Mistral's document chunk tag.
532const MISTRAL_DOCUMENT: &str = "document_url";
533/// Mistral's uploaded-file chunk tag.
534const MISTRAL_FILE: &str = "file";
535/// OpenAI's refusal part: textual, but under a key Mistral's chunk schema
536/// has no field for, so it is re-tagged rather than forwarded.
537const MISTRAL_REFUSAL: &str = "refusal";
538
539/// The text a part carries, under either key the shared conversion uses.
540fn mistral_part_text(part: &serde_json::Value) -> Option<&str> {
541    part.get(MISTRAL_TEXT)
542        .and_then(serde_json::Value::as_str)
543        .or_else(|| {
544            part.get(MISTRAL_REFUSAL)
545                .and_then(serde_json::Value::as_str)
546        })
547}
548
549/// Identify text and refusal parts by tag, falling back to payload keys without a tag.
550/// An explicit nontext tag prevents flattening even when a text key is present.
551fn is_mistral_text_part(part: &serde_json::Value) -> bool {
552    match part.get("type").and_then(serde_json::Value::as_str) {
553        Some(MISTRAL_TEXT | MISTRAL_REFUSAL) => true,
554        Some(_) => false,
555        None => mistral_part_text(part).is_some(),
556    }
557}
558
559fn mistral_unsupported(what: &str) -> EncodeError {
560    crate::message::MessageError::ConversionError(format!(
561        "Mistral cannot carry {what}. Mistral messages accept text, `{MISTRAL_IMAGE}`, \
562         `{MISTRAL_AUDIO}`, `{MISTRAL_DOCUMENT}` and `{MISTRAL_FILE}` content; convert the \
563         content to one of those before sending it."
564    ))
565    .into()
566}
567
568/// Convert file data to a document URL or a file reference to a top-level file ID.
569/// Preserve optional filenames for inline documents. Return a conversion error
570/// when neither file data nor a file ID is present.
571fn mistral_file_chunk(part: &serde_json::Value) -> Result<serde_json::Value, EncodeError> {
572    let file = part.get(MISTRAL_FILE);
573    let field = |name: &str| {
574        file.and_then(|file| file.get(name))
575            .and_then(serde_json::Value::as_str)
576    };
577
578    // Accept already-converted file references to keep finalization idempotent.
579    if let Some(file_id) = part.get("file_id").and_then(serde_json::Value::as_str) {
580        return Ok(serde_json::json!({"type": MISTRAL_FILE, "file_id": file_id}));
581    }
582
583    if let Some(data) = field("file_data") {
584        // `document_name` is left out entirely rather than sent as null when
585        // the part has no filename.
586        Ok(match field("filename") {
587            Some(filename) => serde_json::json!({
588                "type": MISTRAL_DOCUMENT,
589                MISTRAL_DOCUMENT: data,
590                "document_name": filename,
591            }),
592            None => serde_json::json!({"type": MISTRAL_DOCUMENT, MISTRAL_DOCUMENT: data}),
593        })
594    } else if let Some(file_id) = field("file_id") {
595        Ok(serde_json::json!({"type": MISTRAL_FILE, "file_id": file_id}))
596    } else {
597        Err(mistral_unsupported(
598            "a file content part carrying neither `file_data` nor `file_id`",
599        ))
600    }
601}
602
603/// Convert string or object audio payloads to Mistral's base64-string form.
604/// Return a conversion error for missing or nonstring audio data.
605fn mistral_audio_chunk(part: &serde_json::Value) -> Result<serde_json::Value, EncodeError> {
606    let payload = part.get(MISTRAL_AUDIO).ok_or_else(|| {
607        mistral_unsupported("an audio content part carrying no `input_audio` payload")
608    })?;
609
610    let data = match payload {
611        serde_json::Value::String(data) => data.as_str(),
612        payload => payload
613            .get("data")
614            .and_then(serde_json::Value::as_str)
615            .ok_or_else(|| {
616                mistral_unsupported(
617                    "an audio content part whose `input_audio` payload is not base64 data",
618                )
619            })?,
620    };
621
622    Ok(serde_json::json!({"type": MISTRAL_AUDIO, MISTRAL_AUDIO: data}))
623}
624
625/// One content part as the Mistral chunk that carries it.
626///
627/// Dispatched on the `type` tag, which the shared conversion always emits, so
628/// a part naming a chunk kind is converted as that kind regardless of what
629/// other keys it carries.
630fn mistral_chunk(part: &serde_json::Value) -> Result<serde_json::Value, EncodeError> {
631    /// Text and refusal parts are both re-tagged `text`: Mistral's schema has
632    /// no `refusal` field and every chunk forbids unknown keys.
633    fn text_chunk(part: &serde_json::Value) -> Result<serde_json::Value, EncodeError> {
634        let text = mistral_part_text(part)
635            .ok_or_else(|| mistral_unsupported("a text content part carrying no text"))?;
636        Ok(serde_json::json!({"type": MISTRAL_TEXT, MISTRAL_TEXT: text}))
637    }
638
639    match part.get("type").and_then(serde_json::Value::as_str) {
640        Some(MISTRAL_TEXT | MISTRAL_REFUSAL) => text_chunk(part),
641        // Rebuild the envelope because Mistral rejects unknown sibling fields.
642        Some(MISTRAL_IMAGE) => {
643            let image = part.get(MISTRAL_IMAGE).ok_or_else(|| {
644                mistral_unsupported("an image content part carrying no `image_url` payload")
645            })?;
646            Ok(serde_json::json!({"type": MISTRAL_IMAGE, MISTRAL_IMAGE: image}))
647        }
648        Some(MISTRAL_AUDIO) => mistral_audio_chunk(part),
649        Some(MISTRAL_FILE) => mistral_file_chunk(part),
650        // Accept already-converted document parts to keep finalization idempotent.
651        Some(MISTRAL_DOCUMENT) => {
652            let url = part.get(MISTRAL_DOCUMENT).ok_or_else(|| {
653                mistral_unsupported("a document content part carrying no `document_url`")
654            })?;
655            Ok(match part.get("document_name") {
656                Some(name) => serde_json::json!({
657                    "type": MISTRAL_DOCUMENT, MISTRAL_DOCUMENT: url, "document_name": name,
658                }),
659                None => serde_json::json!({"type": MISTRAL_DOCUMENT, MISTRAL_DOCUMENT: url}),
660            })
661        }
662        Some(kind) => Err(mistral_unsupported(&format!("`{kind}` message content"))),
663        // Untagged, but textual: the flattening would have taken it, so it
664        // converts rather than failing.
665        None if mistral_part_text(part).is_some() => text_chunk(part),
666        None => Err(mistral_unsupported("untyped message content")),
667    }
668}
669
670/// Flatten text-only arrays and convert mixed arrays to Mistral content chunks.
671/// Nonarrays remain unchanged. Unsupported mixed content returns a conversion
672/// error; text-only parts without string payloads are omitted.
673fn mistral_content(content: &mut serde_json::Value) -> Result<(), EncodeError> {
674    let Some(parts) = content.as_array() else {
675        return Ok(());
676    };
677
678    if parts.iter().all(is_mistral_text_part) {
679        // The tag-based guard is authoritative even for text parts missing their payload.
680        unary::flatten_text_content_parts(content, "", false);
681        return Ok(());
682    }
683
684    if let Some(parts) = content.as_array_mut() {
685        for part in parts {
686            *part = mistral_chunk(part)?;
687        }
688    }
689
690    Ok(())
691}
692
693/// OpenRouter's body rewrites.
694///
695/// OpenRouter's routing preferences (`ProviderPreferences`) need no rewrite:
696/// they reach the body through the request's `additional_params` as
697/// `{"provider": …}`.
698fn finalize_openrouter(map: &mut serde_json::Map<String, serde_json::Value>, prompt_caching: bool) {
699    if prompt_caching {
700        apply_openrouter_prompt_caching(map);
701    }
702
703    let Some(messages) = map.get_mut("messages").and_then(as_array_mut) else {
704        return;
705    };
706    for message in messages {
707        let Some(message) = message.as_object_mut() else {
708            continue;
709        };
710        // The shared assistant message serializes hidden reasoning under the
711        // llama.cpp/DeepSeek key `reasoning_content`; OpenRouter's documented
712        // assistant field is `reasoning`.
713        if message.get("role").and_then(serde_json::Value::as_str) == Some("assistant")
714            && let Some(reasoning) = message.remove("reasoning_content")
715        {
716            message.insert("reasoning".to_owned(), reasoning);
717        }
718
719        // OpenRouter image parts omit the shared fidelity hint.
720        for part in message
721            .get_mut("content")
722            .and_then(as_array_mut)
723            .into_iter()
724            .flatten()
725        {
726            if let Some(image) = part
727                .get_mut("image_url")
728                .and_then(serde_json::Value::as_object_mut)
729            {
730                image.remove("detail");
731            }
732        }
733    }
734}
735
736fn apply_openrouter_prompt_caching(map: &mut serde_json::Map<String, serde_json::Value>) {
737    let Some(messages) = map.get_mut("messages").and_then(as_array_mut) else {
738        return;
739    };
740    let Some(system) = messages
741        .iter_mut()
742        .find(|message| message.get("role").and_then(serde_json::Value::as_str) == Some("system"))
743    else {
744        return;
745    };
746    match system.get("content").cloned() {
747        Some(serde_json::Value::String(text)) => {
748            if let Some(object) = system.as_object_mut() {
749                object.insert(
750                    "content".to_owned(),
751                    serde_json::json!([{
752                        "type": "text",
753                        "text": text,
754                        "cache_control": { "type": "ephemeral" }
755                    }]),
756                );
757            }
758        }
759        Some(serde_json::Value::Array(mut parts)) => {
760            // The last block marks the cache boundary without altering earlier content.
761            if let Some(last) = parts.last_mut()
762                && let Some(object) = last.as_object_mut()
763            {
764                object.insert(
765                    "cache_control".to_owned(),
766                    serde_json::json!({ "type": "ephemeral" }),
767                );
768            }
769            if let Some(object) = system.as_object_mut() {
770                object.insert("content".to_owned(), serde_json::Value::Array(parts));
771            }
772        }
773        _ => {}
774    }
775}
776
777/// Return a request error for document or image inputs using provider file IDs.
778fn refuse_file_ids(request: &CompletionRequest) -> Result<(), EncodeError> {
779    use crate::message::{DocumentSourceKind, Message, UserContent};
780
781    let refusal = || {
782        EncodeError::request("Provider file IDs are not supported for OpenRouter document inputs")
783    };
784    for message in &request.chat_history {
785        let Message::User { content, .. } = message else {
786            continue;
787        };
788        for part in content {
789            match part {
790                UserContent::Document(document) => {
791                    if matches!(document.data, DocumentSourceKind::FileId(_)) {
792                        return Err(refusal());
793                    }
794                }
795                UserContent::Image(image) => {
796                    if matches!(image.data, DocumentSourceKind::FileId(_)) {
797                        return Err(refusal());
798                    }
799                }
800                _ => {}
801            }
802        }
803    }
804    Ok(())
805}
806
807impl Wire for Chat {
808    type Op = crate::operation::Completion;
809    type Payload = crate::wire::Encoded;
810    type Frame = crate::wire::WireFrame;
811    type Decoder<'id> = ChatDecoder<'id>;
812
813    /// Format deferral permits tool composition; dialects without schema
814    /// support require the agent's tool-mode enforcement instead.
815    fn describe(&self) -> Descriptor<'_> {
816        Descriptor::new(self.provider.dialect.name)
817            .model(self.model.as_str())
818            .capabilities(Capabilities::completion(
819                ProviderCapabilities::default().with_native_output_tool_composition(
820                    self.provider.dialect.quirks.supports_response_format,
821                ),
822            ))
823    }
824
825    fn encode(&self, request: CompletionRequest, mode: Mode) -> Result<Encoded, EncodeError> {
826        self.encode_with_headers(request, mode, OpenAIConfig::completion_headers)
827    }
828
829    fn decoder<'id>(&self) -> ChatDecoder<'id> {
830        ChatDecoder::new(self.provider.dialect.name, self.provider.dialect.quirks)
831    }
832}
833
834/// Classified Chat Completions frame, including whole replies and terminal signals.
835pub enum ChatEvent {
836    /// A `chat.completion.chunk`: one step of a streamed turn.
837    Chunk(ChatFrame),
838    /// A `chat.completion`: the whole turn in one frame.
839    Whole(ChatFrame),
840    /// The `[DONE]` sentinel: the provider ended the stream.
841    Done,
842    /// The wire's in-band error envelope, delivered with a 200 status.
843    Failure(ProviderError),
844    /// A bare JSON string where an envelope belongs: the whole answer, with
845    /// no metadata and no terminal reason. Mira's gateway sends this.
846    BareText(String),
847}
848
849/// The chat-completions decoder: one state machine for a whole reply and a
850/// stream of chunks.
851pub struct ChatDecoder<'id> {
852    /// Descriptor name the reply is attributed to.
853    provider: &'static str,
854    quirks: super::Quirks,
855    /// The text part bare text extends; a reasoning fragment or a tool call
856    /// closes it.
857    text: Option<TextPart<'id>>,
858    /// `reasoning_content` carries no ids or boundaries.
859    thoughts: Thoughts<'id>,
860    final_usage: Option<ChatUsage>,
861    final_finish_reason: Option<FinishReason>,
862    response_id: Option<String>,
863    response_model: Option<String>,
864    /// Accumulated primary-choice token metadata. `AdditionalParams::merge`
865    /// concatenates nested arrays, which is the wire's token order.
866    logprobs: Option<crate::message::AdditionalParams>,
867    /// Accumulated provider-specific top-level chunk metadata.
868    additional_params: Option<crate::message::AdditionalParams>,
869    /// Whether a finish reason established the turn complete.
870    saw_terminal: bool,
871    /// Whether any frame decoded successfully. A bare `[DONE]` after only
872    /// parse failures must not dress the failure up as a default-usage
873    /// success.
874    saw_any_valid_frame: bool,
875}
876
877impl<'id> ChatDecoder<'id> {
878    fn new(provider: &'static str, quirks: super::Quirks) -> Self {
879        Self {
880            provider,
881            quirks,
882            text: None,
883            thoughts: Thoughts::new(),
884            final_usage: None,
885            final_finish_reason: None,
886            response_id: None,
887            response_model: None,
888            logprobs: None,
889            additional_params: None,
890            saw_terminal: false,
891            saw_any_valid_frame: false,
892        }
893    }
894
895    /// The normalized finish reason a choice reported, `None` when the
896    /// chunk carried no `finish_reason`.
897    ///
898    /// A gateway's upstream-native reason is consulted only when the
899    /// normalized field is absent or empty, which is OpenRouter's documented
900    /// precedence; a direct provider has no native field to consult.
901    fn finish_reason(&self, choice: &ChatChoice) -> Option<FinishReason> {
902        if let Some(reason) = choice
903            .finish_reason
904            .as_ref()
905            .map(super::dto::FinishReason::as_wire)
906            .filter(|reason| !reason.is_empty())
907        {
908            return Some(map_openai_finish_reason(reason));
909        }
910        if self.quirks.native_finish_reason
911            && let Some(native) = choice
912                .native_finish_reason
913                .as_deref()
914                .filter(|reason| !reason.is_empty())
915        {
916            return Some(map_native_finish_reason(native));
917        }
918        None
919    }
920
921    /// Absorb the metadata every frame carries, whichever shape it is.
922    fn absorb_metadata(&mut self, frame: &mut ChatFrame) {
923        if let Some(id) = frame.id.take() {
924            self.response_id = Some(id);
925        }
926        if let Some(model) = frame.model.take() {
927            self.response_model = Some(model);
928        }
929        if let Some(usage) = frame.usage.take() {
930            self.final_usage = Some(usage);
931        }
932        if let Some(additional_params) =
933            crate::message::AdditionalParams::new(std::mem::take(&mut frame.additional_params))
934        {
935            match self.additional_params.as_mut() {
936                Some(accumulated) => accumulated.merge(additional_params),
937                None => self.additional_params = Some(additional_params),
938            }
939        }
940    }
941
942    fn close_text(&mut self, out: &mut Out<'id, Completion>) {
943        if let Some(part) = self.text.take() {
944            out.close_text(part);
945        }
946    }
947
948    /// One chunk's parts, in the order the wire implies: reasoning, its
949    /// signature or the boundary that stops it, text, then tool calls.
950    fn emit_parts(
951        &mut self,
952        out: &mut Out<'id, Completion>,
953        reasoning: Option<String>,
954        signature: Option<String>,
955        text: Option<String>,
956        calls: bool,
957    ) {
958        if let Some(reasoning) = reasoning.filter(|reasoning| !reasoning.is_empty()) {
959            self.close_text(out);
960            self.thoughts.fragment(out, &reasoning);
961        }
962        if let Some(signature) = signature {
963            self.thoughts.signature(out, signature);
964        }
965        let text = text.filter(|text| !text.is_empty());
966        if text.is_some() || calls {
967            // These wires omit the boundary before interleaving output.
968            self.thoughts.boundary();
969        }
970        if let Some(text) = text {
971            let part = self.text.get_or_insert_with(|| out.text());
972            out.push_text(part, &text);
973        }
974    }
975
976    /// One `chat.completion.chunk`.
977    fn interpret_chunk(
978        &mut self,
979        mut frame: ChatFrame,
980        out: &mut Out<'id, Completion>,
981    ) -> Result<(), ProviderError> {
982        self.saw_any_valid_frame = true;
983        self.absorb_metadata(&mut frame);
984        let Some(choice) = frame.into_primary() else {
985            return Ok(());
986        };
987        let finish_reason = self.finish_reason(&choice);
988        let text = delta_text(&choice.delta);
989        let StreamingDelta {
990            reasoning_content,
991            reasoning,
992            tool_calls,
993            reasoning_details,
994            ..
995        } = choice.delta;
996        let reasoning = reasoning_content.or(reasoning);
997        let details: Vec<unary::ReasoningDetails> =
998            reasoning_details.iter().filter_map(typed_detail).collect();
999
1000        if let Some(reason) = &finish_reason {
1001            self.final_finish_reason = Some(reason.clone());
1002            self.saw_terminal = true;
1003        }
1004
1005        if let Some(logprobs) = choice.logprobs {
1006            match self.logprobs.as_mut() {
1007                Some(accumulated) => accumulated.merge(logprobs),
1008                None => self.logprobs = Some(logprobs),
1009            }
1010        }
1011
1012        // Replayable reasoning must precede the tool calls it accompanies.
1013        if self.quirks.reasoning_details {
1014            for detail in &details {
1015                if let Some(reasoning) = detail_reasoning(detail) {
1016                    self.close_text(out);
1017                    out.reasoning_block(reasoning);
1018                }
1019            }
1020        }
1021
1022        let reasoning_signature = self
1023            .quirks
1024            .reasoning_details
1025            .then(|| details.iter().find_map(reasoning_signature))
1026            .flatten();
1027
1028        self.emit_parts(
1029            out,
1030            reasoning,
1031            reasoning_signature,
1032            text,
1033            !tool_calls.is_empty(),
1034        );
1035
1036        for incoming in tool_calls {
1037            self.close_text(out);
1038            if let Some(existing) = out.pending_id(incoming.index)
1039                && incoming.evicts(&existing, &out.pending_name(incoming.index))
1040            {
1041                // The wire reused this call's index: the call it held is
1042                // delivered even when its arguments never parse.
1043                out.close_pending(incoming.index, IfMalformed::EmptyObject)?;
1044            }
1045            out.call_fragment(
1046                incoming.index,
1047                CallFragment {
1048                    id: incoming.id.as_deref(),
1049                    name: incoming.function.name.as_deref(),
1050                    arguments: incoming.function.arguments.as_deref(),
1051                    ..CallFragment::default()
1052                },
1053            )?;
1054            if self.quirks.emits_complete_single_chunk_tool_calls
1055                && incoming.is_complete_single_chunk()
1056            {
1057                // A probe: the call closes if its input parses, and stays
1058                // open for more fragments otherwise.
1059                out.close_pending(incoming.index, IfMalformed::KeepOpen)?;
1060            }
1061        }
1062
1063        if matches!(finish_reason, Some(FinishReason::ToolCalls)) {
1064            for index in out.pending_calls() {
1065                // Completed calls with malformed arguments must fail, not
1066                // disappear. Empty arguments remain valid for zero-argument
1067                // tools.
1068                out.close_pending(index, IfMalformed::Fail)?;
1069            }
1070        }
1071        Ok(())
1072    }
1073
1074    /// Whether a length-truncated unary choice contains tool calls needing raw inspection.
1075    /// Empty argument strings normalize to `{}`, so typed arguments alone cannot
1076    /// distinguish truncation before the first token from a zero-argument call.
1077    fn is_budget_cut_tool_turn(&self, frame: &ChatFrame) -> bool {
1078        let Some(choice) = frame.primary() else {
1079            return false;
1080        };
1081        if !matches!(self.finish_reason(choice), Some(FinishReason::Length)) {
1082            return false;
1083        }
1084        matches!(
1085            &choice.message,
1086            Some(Message::Assistant { tool_calls, .. }) if !tool_calls.is_empty()
1087        )
1088    }
1089
1090    /// Whether a raw choice blames the output-token budget, under the same
1091    /// precedence [`Self::finish_reason`] applies to a decoded one.
1092    fn reports_output_length(&self, choice: &serde_json::Value) -> bool {
1093        let reason = |key: &str| {
1094            choice
1095                .get(key)
1096                .and_then(serde_json::Value::as_str)
1097                .filter(|reason| !reason.is_empty())
1098        };
1099        if let Some(normalized) = reason("finish_reason") {
1100            return matches!(map_openai_finish_reason(normalized), FinishReason::Length);
1101        }
1102        self.quirks.native_finish_reason
1103            && reason("native_finish_reason").is_some_and(|native| {
1104                matches!(map_native_finish_reason(native), FinishReason::Length)
1105            })
1106    }
1107
1108    /// Decode a unary body after dropping incomplete calls from length-truncated choices.
1109    /// Preserve valid arguments and require the shared compound-defect check before
1110    /// dropping calls. Return `None` if nothing is dropped or decoding still fails.
1111    fn body_without_calls_cut_by_the_budget(&self, data: &str) -> Option<ChatFrame> {
1112        let mut body = serde_json::from_str::<serde_json::Value>(data).ok()?;
1113        let mut dropped = 0;
1114        for choice in body.get_mut("choices").and_then(as_array_mut)? {
1115            if self.reports_output_length(choice) {
1116                dropped += drop_tool_calls_cut_by_budget::<ChatChoice>(choice);
1117            }
1118        }
1119        if dropped == 0 {
1120            return None;
1121        }
1122        let frame = serde_json::from_value::<ChatFrame>(body).ok()?;
1123        tracing::debug!(
1124            provider = self.provider,
1125            dropped,
1126            "dropping unary tool calls whose arguments the output-token budget cut short"
1127        );
1128        Some(frame)
1129    }
1130
1131    /// The `chat.completion` body: the parts a stream of the same turn
1132    /// would have written, then the end.
1133    fn interpret_whole(
1134        &mut self,
1135        mut frame: ChatFrame,
1136        mut out: Out<'id, Completion>,
1137    ) -> Result<Flow, ProviderError> {
1138        self.saw_any_valid_frame = true;
1139        let Some(choice) = frame.primary() else {
1140            return Err(ProviderError::Response(
1141                "Response contained no choices".to_owned(),
1142            ));
1143        };
1144        let finish_reason = self.finish_reason(choice);
1145        let Some(Message::Assistant {
1146            content,
1147            reasoning,
1148            refusal,
1149            tool_calls,
1150            reasoning_details,
1151            ..
1152        }) = choice.message.clone()
1153        else {
1154            return Err(ProviderError::Response(
1155                "Response did not contain a valid message or tool call".to_owned(),
1156            ));
1157        };
1158        let logprobs = choice.logprobs.clone();
1159        self.absorb_metadata(&mut frame);
1160        self.logprobs = logprobs;
1161        self.final_finish_reason = finish_reason;
1162        self.saw_terminal = true;
1163
1164        // Response IDs are not replayable message IDs; retain them only as terminal metadata.
1165        let text = {
1166            // The streamed path concatenates a turn's text fragments into one
1167            // part, so the unary body's parts join the same way rather than
1168            // producing a different number of parts for the same turn.
1169            let mut text = String::new();
1170            for part in &content {
1171                let part = match part {
1172                    AssistantContent::Text { text } => text,
1173                    AssistantContent::Refusal { refusal } => refusal,
1174                };
1175                text.push_str(part);
1176            }
1177            // This wire spells a refusal as a *sibling* of `content`
1178            // (`{"content": null, "refusal": "…"}`), so a path reading
1179            // `content` alone would drop it entirely.
1180            if let Some(refusal) = assistant_refusal_fallback(&content, refusal.as_deref()) {
1181                text.push_str(refusal);
1182            }
1183            text
1184        };
1185
1186        let reasoning = reasoning.filter(|reasoning| !reasoning.is_empty());
1187        // Structured details retain signatures needed to replay reasoning.
1188        let details: Vec<&unary::ReasoningDetails> = if self.quirks.reasoning_details {
1189            reasoning_details.iter().collect()
1190        } else {
1191            Vec::new()
1192        };
1193        let blocks: Vec<_> = details
1194            .iter()
1195            .copied()
1196            .filter_map(whole_detail_reasoning)
1197            .collect();
1198        // Prefer replayable structured blocks to avoid duplicating their plaintext display.
1199        // Without those blocks, preserve plaintext and any signature-only detail.
1200        let (reasoning, reasoning_signature) = if blocks.is_empty() {
1201            (
1202                reasoning,
1203                details.iter().copied().find_map(reasoning_signature),
1204            )
1205        } else {
1206            (None, None)
1207        };
1208        // Truncation or filtering can leave no visible content; retain its reason and usage.
1209        // Empty replies without such a reason are response errors.
1210        let cut_short = self
1211            .final_finish_reason
1212            .as_ref()
1213            .is_some_and(FinishReason::truncated_output);
1214        if text.is_empty()
1215            && tool_calls.is_empty()
1216            && reasoning.is_none()
1217            && blocks.is_empty()
1218            && !cut_short
1219        {
1220            return Err(ProviderError::Response(
1221                crate::message::EMPTY_RESPONSE_ERROR.to_owned(),
1222            ));
1223        }
1224
1225        // Reasoning details are the turn's own output, so they are written
1226        // before the text and tool calls, exactly as the streamed path orders
1227        // them.
1228        for reasoning in blocks {
1229            out.reasoning_block(reasoning);
1230        }
1231
1232        self.emit_parts(
1233            &mut out,
1234            reasoning,
1235            reasoning_signature,
1236            (!text.is_empty()).then_some(text),
1237            !tool_calls.is_empty(),
1238        );
1239        self.close_text(&mut out);
1240
1241        // Each call is buffered at its own position, so separate id-less
1242        // calls stay distinct.
1243        for (index, call) in tool_calls.iter().enumerate() {
1244            out.call_fragment(
1245                index,
1246                CallFragment {
1247                    id: Some(call.id.as_str()),
1248                    name: Some(call.function.name.as_str()),
1249                    ..CallFragment::default()
1250                },
1251            )?;
1252            out.announce_pending(index, call.function.arguments.clone());
1253            out.close_pending(index, IfMalformed::Fail)?;
1254        }
1255
1256        self.end(out, false)
1257    }
1258
1259    /// Write the provider's end of the reply. A stream's `raw` is the
1260    /// native terminal record the chunks built; a whole body's is the body
1261    /// itself, which the transport keeps.
1262    fn end(
1263        &mut self,
1264        mut out: Out<'id, Completion>,
1265        streamed: bool,
1266    ) -> Result<Flow, ProviderError> {
1267        self.close_text(&mut out);
1268        self.thoughts.close(&mut out, None);
1269        // A gateway's reasoning belongs to the upstream model that produced it.
1270        if self.quirks.upstream_reasoning_issuer
1271            && let Some(model) = self.response_model.as_deref()
1272        {
1273            out.issued_by(super::upstream_reasoning_issuer(self.provider, model));
1274        }
1275        let usage = self
1276            .final_usage
1277            .as_ref()
1278            .map(|usage| usage.to_normalized_for(&self.quirks));
1279        let native = StreamingCompletionResponse {
1280            usage: self.final_usage.take(),
1281            finish_reason: self.final_finish_reason.take(),
1282            response_id: self.response_id.take(),
1283            model: self.response_model.take(),
1284            logprobs: self.logprobs.take().map(Into::into),
1285            additional_params: self.additional_params.take(),
1286        };
1287        if streamed {
1288            out.raw(serde_json::to_value(&native)?);
1289        }
1290        let mut finish = native.into_finish();
1291        finish.usage = usage.unwrap_or_default();
1292        Ok(out.end(finish))
1293    }
1294
1295    /// The stream ended: flush the calls the provider delivered, then end
1296    /// the reply. Tool calls the provider fully delivered are content, so
1297    /// a length cut still flushes them.
1298    fn finish(&mut self, mut out: Out<'id, Completion>) -> Result<Flow, ProviderError> {
1299        let output_length_truncation = matches!(
1300            self.final_finish_reason.as_ref(),
1301            Some(FinishReason::Length)
1302        );
1303        for index in out.pending_calls() {
1304            if output_length_truncation && !out.pending_has_arguments(index) {
1305                tracing::debug!(
1306                    "dropping streamed tool call cut off before its first argument token"
1307                );
1308                out.drop_pending(index);
1309                continue;
1310            }
1311            // Only an explicit length finish permits dropping malformed arguments.
1312            let if_malformed = if output_length_truncation {
1313                IfMalformed::Drop
1314            } else {
1315                IfMalformed::Fail
1316            };
1317            out.close_pending(index, if_malformed)?;
1318        }
1319        // A bare terminator without valid content cannot establish success.
1320        if !self.saw_any_valid_frame {
1321            return Err(ProviderError::Truncated);
1322        }
1323        self.end(out, true)
1324    }
1325}
1326
1327use crate::operation::{CallFragment, Completion, IfMalformed, TextPart};
1328use crate::providers::internal::thoughts::Thoughts;
1329use crate::wire::Flow;
1330
1331impl<'id> Decoder<'id, Completion> for ChatDecoder<'id> {
1332    type Event = ChatEvent;
1333
1334    fn classify(&self, frame: WireFrame) -> WireEvent<ChatEvent> {
1335        let data = frame.as_str();
1336        // `[DONE]` is the wire's terminal sentinel, not JSON; it is Known by
1337        // definition and its decode ends the reply.
1338        if data == "[DONE]" {
1339            return WireEvent::Known(ChatEvent::Done);
1340        }
1341        // The wire's in-band error envelope arrives with a 200 status and is
1342        // this wire's own terminal failure, so it is a modeled event rather
1343        // than a transport-level filter.
1344        if let Some(error) = provider_error_envelope(&data) {
1345            return WireEvent::Known(ChatEvent::Failure(error));
1346        }
1347        // Supported bare-string replies must be recognized before object classification.
1348        if self.quirks.accepts_bare_string_reply
1349            && let Ok(serde_json::Value::String(text)) =
1350                serde_json::from_str::<serde_json::Value>(&data)
1351        {
1352            return WireEvent::Known(ChatEvent::BareText(text));
1353        }
1354        let classified = classify_chat_completions_frame::<ChatFrame>(&data);
1355        // Inspect raw arguments for length cuts, including empty strings normalized to {}.
1356        let may_be_budget_cut = match &classified {
1357            WireEvent::Corrupt(_) => true,
1358            WireEvent::Known(frame) => self.is_budget_cut_tool_turn(frame),
1359            WireEvent::Unknown { .. } => false,
1360        };
1361        if may_be_budget_cut && let Some(frame) = self.body_without_calls_cut_by_the_budget(&data) {
1362            return WireEvent::Known(ChatEvent::Whole(frame));
1363        }
1364        classified.map(|frame| {
1365            if frame.is_whole() {
1366                ChatEvent::Whole(frame)
1367            } else {
1368                ChatEvent::Chunk(frame)
1369            }
1370        })
1371    }
1372
1373    fn decode(
1374        &mut self,
1375        event: ChatEvent,
1376        mut out: Out<'id, Completion>,
1377    ) -> Result<Flow, ProviderError> {
1378        match event {
1379            ChatEvent::Chunk(frame) => {
1380                self.interpret_chunk(frame, &mut out)?;
1381                Ok(Flow::More)
1382            }
1383            ChatEvent::Whole(frame) => self.interpret_whole(frame, out),
1384            ChatEvent::Done => {
1385                if !self.saw_terminal {
1386                    // `[DONE]` without a finish reason still ends the turn.
1387                    self.saw_terminal = true;
1388                }
1389                self.finish(out)
1390            }
1391            ChatEvent::BareText(text) => {
1392                self.saw_any_valid_frame = true;
1393                self.saw_terminal = true;
1394                if !text.is_empty() {
1395                    let part = out.text();
1396                    out.push_text(&part, &text);
1397                    out.close_text(part);
1398                }
1399                self.end(out, false)
1400            }
1401            ChatEvent::Failure(error) => Err(error),
1402        }
1403    }
1404
1405    /// A stream that stops after a finish reason without `[DONE]` still
1406    /// ended: some dialects (Perplexity) never send the sentinel.
1407    fn eof(&mut self, out: Out<'id, Completion>) -> Result<Flow, ProviderError> {
1408        if !self.saw_terminal {
1409            return Err(ProviderError::Truncated);
1410        }
1411        self.finish(out)
1412    }
1413}
1414
1415impl ChatDecoder<'_> {
1416    /// Verdict, model, response id, usage and error envelope, read off a raw
1417    /// payload before normalization discards them. The driver calls it for
1418    /// the unary reply and for every stream frame without anyone having to
1419    /// attach it.
1420    pub(crate) fn project(payload: &[u8], sink: &mut ObservationSink<'_>) {
1421        let Ok(payload) = serde_json::from_slice::<ObservedPayload>(payload) else {
1422            return;
1423        };
1424        if let Some(usage) = payload.usage {
1425            sink.emit(AdapterEvent::Usage {
1426                usage: AdapterUsage {
1427                    input_tokens: usage.prompt_tokens,
1428                    output_tokens: usage.completion_tokens,
1429                    total_tokens: usage.total_tokens,
1430                    cached_input_tokens: usage
1431                        .prompt_tokens_details
1432                        .and_then(|details| details.cached_tokens),
1433                    reasoning_tokens: usage
1434                        .completion_tokens_details
1435                        .and_then(|details| details.reasoning_tokens),
1436                    tool_input_tokens: None,
1437                },
1438            });
1439        }
1440        // Every chunk names the model; only the chunk that carries the finish
1441        // reason is a verdict, so the model rides with it rather than on each
1442        // delta. The id still lands on the terminal verdict or the closure.
1443        let choice = payload.choices.into_iter().next().unwrap_or_default();
1444        let verdict = match choice.finish_reason {
1445            Some(reason) => AdapterVerdict {
1446                finish_reason: Some(sink.scrub(&reason)),
1447                block_reason: None,
1448                detail: None,
1449                model: payload.model.map(|value| sink.scrub(&value)),
1450            },
1451            None => AdapterVerdict::default(),
1452        };
1453        let response_id = payload.id.map(|value| sink.scrub(&value));
1454        sink.provider(verdict, response_id);
1455        if let Some(error) = payload.error {
1456            error.emit(sink);
1457        }
1458    }
1459}
1460
1461/// One object for the unary reply and each stream chunk `project` above
1462/// sees. Every field is optional: a chunk carries a delta, the last chunk
1463/// (or the reply) carries the usage, and `[DONE]` is not JSON at all.
1464#[derive(Deserialize)]
1465struct ObservedPayload {
1466    id: Option<String>,
1467    model: Option<String>,
1468    usage: Option<ObservedUsage>,
1469    #[serde(default)]
1470    choices: Vec<ObservedChoice>,
1471    error: Option<ObservedError>,
1472}
1473
1474#[derive(Deserialize)]
1475struct ObservedUsage {
1476    #[serde(default, deserialize_with = "crate::observe::lenient_count")]
1477    prompt_tokens: Option<u64>,
1478    #[serde(default, deserialize_with = "crate::observe::lenient_count")]
1479    completion_tokens: Option<u64>,
1480    #[serde(default, deserialize_with = "crate::observe::lenient_count")]
1481    total_tokens: Option<u64>,
1482    #[serde(default)]
1483    prompt_tokens_details: Option<ObservedTokenDetails>,
1484    #[serde(default)]
1485    completion_tokens_details: Option<ObservedTokenDetails>,
1486}
1487
1488#[derive(Default, Deserialize)]
1489struct ObservedTokenDetails {
1490    #[serde(default, deserialize_with = "crate::observe::lenient_count")]
1491    cached_tokens: Option<u64>,
1492    #[serde(default, deserialize_with = "crate::observe::lenient_count")]
1493    reasoning_tokens: Option<u64>,
1494}
1495
1496#[derive(Default, Deserialize)]
1497struct ObservedChoice {
1498    finish_reason: Option<String>,
1499}
1500
1501/// A gateway's encrypted-reasoning detail as a whole reasoning block.
1502///
1503/// Encrypted reasoning (`{"type":"reasoning.encrypted"}`) is the turn's own
1504/// output, not tool-call metadata: it arrives with `reasoning: null` and an
1505/// `rs_*` id of its own, which never matches a `call_*` tool-call id, and it
1506/// arrives before any tool call opens. Writing it as a reasoning part is
1507/// what lets the blob reach the aggregated choice and be replayed next turn.
1508fn detail_reasoning(detail: &unary::ReasoningDetails) -> Option<crate::message::Reasoning> {
1509    let unary::ReasoningDetails::Encrypted { id, data, .. } = detail else {
1510        return None;
1511    };
1512    Some(crate::message::Reasoning {
1513        id: id.clone().filter(|id| !id.is_empty()),
1514        content: vec![crate::message::ReasoningContent::Encrypted(data.clone())],
1515    })
1516}
1517
1518/// A nonempty unary reasoning detail as one replayable part. Empty and
1519/// signature-only entries return `None`.
1520fn whole_detail_reasoning(detail: &unary::ReasoningDetails) -> Option<crate::message::Reasoning> {
1521    let (id, content) = match detail {
1522        unary::ReasoningDetails::Summary { id, summary, .. } if !summary.is_empty() => (
1523            id,
1524            crate::message::ReasoningContent::Summary(summary.clone()),
1525        ),
1526        unary::ReasoningDetails::Encrypted { id, data, .. } if !data.is_empty() => (
1527            id,
1528            crate::message::ReasoningContent::Encrypted(data.clone()),
1529        ),
1530        unary::ReasoningDetails::Text {
1531            id,
1532            text: Some(text),
1533            signature,
1534            ..
1535        } if !text.is_empty() => (
1536            id,
1537            crate::message::ReasoningContent::Text {
1538                text: text.clone(),
1539                signature: signature.clone().filter(|signature| !signature.is_empty()),
1540            },
1541        ),
1542        // Signature-only details attach to separately supplied plaintext.
1543        _ => return None,
1544    };
1545    Some(crate::message::Reasoning {
1546        id: id.clone().filter(|id| !id.is_empty()),
1547        content: vec![content],
1548    })
1549}
1550
1551/// Return a nonempty signature from a text reasoning detail.
1552fn reasoning_signature(detail: &unary::ReasoningDetails) -> Option<String> {
1553    let unary::ReasoningDetails::Text {
1554        signature: Some(signature),
1555        ..
1556    } = detail
1557    else {
1558        return None;
1559    };
1560    (!signature.is_empty()).then(|| signature.clone())
1561}
1562
1563/// Decode a modeled reasoning detail, returning `None` for unrecognized or invalid shapes.
1564fn typed_detail(detail: &serde_json::Value) -> Option<unary::ReasoningDetails> {
1565    serde_json::from_value(detail.clone()).ok()
1566}
1567
1568#[cfg(test)]
1569mod tests;
1570
1571#[cfg(test)]
1572mod hard_case_tests;