1use 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#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
37pub struct Chat {
38 pub provider: OpenAIConfig,
40 pub model: String,
42 pub strict_tools: bool,
46 pub tool_result_array_content: bool,
48 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 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 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 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 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 pub fn with_strict_tools(mut self) -> Self {
156 self.strict_tools = true;
157 self
158 }
159
160 pub fn with_tool_result_array_content(mut self) -> Self {
162 self.tool_result_array_content = true;
163 self
164 }
165
166 pub fn with_prompt_caching(mut self) -> Self {
168 self.prompt_caching = true;
169 self
170 }
171
172 fn prepare(&self, request: &mut unary::CompletionRequest) -> Result<(), EncodeError> {
174 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 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 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 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 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 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
262const 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
268fn 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
275fn 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
312fn 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
324fn 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 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
371fn 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
395fn 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 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 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
448fn finalize_mistral(
454 map: &mut serde_json::Map<String, serde_json::Value>,
455) -> Result<(), EncodeError> {
456 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 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 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 message
515 .entry("prefix")
516 .or_insert(serde_json::Value::Bool(false));
517 message.remove("reasoning_content");
520 }
521 }
522 Ok(())
523}
524
525const MISTRAL_TEXT: &str = "text";
527const MISTRAL_IMAGE: &str = "image_url";
529const MISTRAL_AUDIO: &str = "input_audio";
531const MISTRAL_DOCUMENT: &str = "document_url";
533const MISTRAL_FILE: &str = "file";
535const MISTRAL_REFUSAL: &str = "refusal";
538
539fn 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
549fn 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
568fn 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 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 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
603fn 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
625fn mistral_chunk(part: &serde_json::Value) -> Result<serde_json::Value, EncodeError> {
631 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 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 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 None if mistral_part_text(part).is_some() => text_chunk(part),
666 None => Err(mistral_unsupported("untyped message content")),
667 }
668}
669
670fn 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 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
693fn 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 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 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 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
777fn 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 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
834pub enum ChatEvent {
836 Chunk(ChatFrame),
838 Whole(ChatFrame),
840 Done,
842 Failure(ProviderError),
844 BareText(String),
847}
848
849pub struct ChatDecoder<'id> {
852 provider: &'static str,
854 quirks: super::Quirks,
855 text: Option<TextPart<'id>>,
858 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 logprobs: Option<crate::message::AdditionalParams>,
867 additional_params: Option<crate::message::AdditionalParams>,
869 saw_terminal: bool,
871 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 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 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 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 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 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 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 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 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 out.close_pending(index, IfMalformed::Fail)?;
1069 }
1070 }
1071 Ok(())
1072 }
1073
1074 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 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 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 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 let text = {
1166 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 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 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 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 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 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 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 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 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 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 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 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 if data == "[DONE]" {
1339 return WireEvent::Known(ChatEvent::Done);
1340 }
1341 if let Some(error) = provider_error_envelope(&data) {
1345 return WireEvent::Known(ChatEvent::Failure(error));
1346 }
1347 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 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 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 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 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 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#[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
1501fn 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
1518fn 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 _ => return None,
1544 };
1545 Some(crate::message::Reasoning {
1546 id: id.clone().filter(|id| !id.is_empty()),
1547 content: vec![content],
1548 })
1549}
1550
1551fn 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
1563fn 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;