1use serde::{Deserialize, Serialize};
10use serde_json::{Map, Value, json};
11
12use crate::completion::options::{BaseInput, RawAt, Rewrite, request_params};
13use crate::completion::{CompletionRequest, FinishReason, ProviderCapabilities, Replay};
14use crate::error::{EncodeError, ProviderError};
15use crate::json_utils::Lenient;
16use crate::message::{
17 AssistantContent, AssistantMessage, DocumentMediaType, DocumentSourceKind as Source, Message,
18 MimeType, ToolCall, ToolResult, ToolResultContent, UserContent,
19};
20use crate::observe::ObservedError;
21use crate::operation::{Block, CallFragment, Completion, Finish};
22use crate::providers::internal::openai_chat_completions_compatible::{
23 finish_reason, native_finish_reason, provider_error_envelope,
24};
25use crate::providers::internal::wire::classify_chat_completions_frame;
26use crate::providers::internal::wire_ids::WireIds;
27use crate::wire::{
28 AdapterEvent, AdapterUsage, AdapterVerdict, Capabilities, Decoder, Descriptor, Encoded, Flow,
29 Framing, Mode, ObservationSink, Out, Wire, WireEvent, WireFrame,
30};
31
32use super::dto::merge_fields;
33use super::{BodyRewrite, OpenAIConfig, OutputCap, Quirks};
34
35#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
38pub struct Chat {
39 pub provider: OpenAIConfig,
41 pub model: String,
43 pub strict_tools: bool,
47 pub tool_result_array_content: bool,
49}
50
51fn unsendable(what: &str) -> EncodeError {
55 EncodeError::request(format!("Chat Completions cannot carry {what} in this form"))
56}
57
58fn media_url(data: &Source, mime: Option<&str>) -> Option<String> {
61 match (data, mime) {
62 (Source::Url(url), _) => Some(url.clone()),
63 (Source::Base64(data), Some(mime)) => Some(format!("data:{mime};base64,{data}")),
64 _ => None,
65 }
66}
67
68fn image_url(image: &crate::message::Image, user: bool) -> Result<Value, EncodeError> {
71 let mime = image.media_type.as_ref().map(MimeType::to_mime_type);
72 let url = media_url(&image.data, mime).ok_or_else(|| unsendable("an image"))?;
73 let mut part = Map::from_iter([("url".to_owned(), Value::String(url))]);
74 let detail = match &image.detail {
75 Some(detail) => Some(detail.clone()),
76 None => user.then(Default::default),
77 };
78 if let Some(detail) = detail {
79 part.insert("detail".to_owned(), serde_json::to_value(detail)?);
80 }
81 Ok(json!({"type": "image_url", "image_url": part}))
82}
83
84fn user_part(part: &UserContent) -> Result<Value, EncodeError> {
86 let file = |file: Value| json!({"type": "file", "file": file});
87 Ok(match part {
88 UserContent::Text(text) => json!({"type": "text", "text": text.text}),
89 UserContent::Image(image) => image_url(image, true)?,
90 UserContent::Document(document) => {
91 let pdf = document.media_type == Some(DocumentMediaType::PDF);
92 match &document.data {
93 Source::FileId(id) => file(json!({ "file_id": id })),
94 Source::Base64(data) if pdf => file(json!({
95 "file_data": format!("data:application/pdf;base64,{data}"),
96 "filename": "document.pdf",
97 })),
98 Source::Url(url) if pdf => {
100 file(json!({"file_data": url, "filename": "document.pdf"}))
101 }
102 Source::String(text) if !pdf => json!({"type": "text", "text": text}),
103 _ => return Err(unsendable("a document")),
104 }
105 }
106 UserContent::Audio(audio) => match &audio.data {
107 Source::Base64(data) => json!({"type": "input_audio", "input_audio": {
108 "data": data,
109 "format": audio.media_type.clone().unwrap_or(crate::message::AudioMediaType::MP3),
110 }}),
111 _ => return Err(unsendable("audio")),
112 },
113 UserContent::Video(video) => {
114 let mime = video.media_type.as_ref().map(MimeType::to_mime_type);
115 let url = media_url(&video.data, mime).ok_or_else(|| unsendable("a video"))?;
116 json!({"type": "video_url", "video_url": {"url": url}})
117 }
118 UserContent::ToolResult(_) => return Err(unsendable("a tool result as a content part")),
119 })
120}
121
122fn tool_message(result: &ToolResult, ids: &WireIds, array: bool) -> Result<Value, EncodeError> {
126 let parts = result
127 .content
128 .iter()
129 .map(|part| match part {
130 ToolResultContent::Text(text) => Ok(json!({"type": "text", "text": text.text})),
131 ToolResultContent::Json { value } => {
132 Ok(json!({"type": "text", "text": value.to_string()}))
133 }
134 ToolResultContent::Image(image) => image_url(image, false),
135 })
136 .collect::<Result<Vec<_>, EncodeError>>()?;
137 let content = if array
138 || parts
139 .iter()
140 .any(|part| part.str("type") == Some("image_url"))
141 {
142 Value::Array(parts)
143 } else {
144 let texts: Vec<&str> = parts.iter().filter_map(|part| part.str("text")).collect();
145 Value::String(texts.join("\n"))
146 };
147 let id = ids.spell(&result.call);
148 Ok(json!({"role": "tool", "tool_call_id": id, "content": content}))
149}
150
151fn user_message(messages: &mut Vec<Value>, parts: &mut Vec<Value>) {
154 let content = match parts.as_slice() {
155 [] => return,
156 [part] if part.str("type") == Some("text") => part.get("text").cloned().unwrap_or_default(),
157 _ => Value::Array(std::mem::take(parts)),
158 };
159 parts.clear();
160 messages.push(json!({"role": "user", "content": content}));
161}
162
163impl Chat {
164 pub(crate) fn encode_with_headers(
165 &self,
166 request: CompletionRequest,
167 mode: Mode,
168 headers: impl FnOnce(
169 &OpenAIConfig,
170 &CompletionRequest,
171 http::request::Builder,
172 ) -> http::request::Builder,
173 ) -> Result<Encoded, EncodeError> {
174 let quirks = &self.provider.dialect.quirks;
175 let uri = self.provider.uri(
177 quirks.completion_path,
178 self.provider.deployment(&self.model),
179 );
180 let builder = headers(
181 &self.provider,
182 &request,
183 http::Request::post(uri).header("Content-Type", "application/json"),
184 );
185 let mut rewrites = Vec::new();
186 if quirks.output_cap == OutputCap::OpenAiReasoningFamilies {
187 rewrites.push(Rewrite::OutputCapRename);
188 }
189 if mode == Mode::Streaming {
190 if quirks.stream_include_usage {
192 rewrites.push(Rewrite::StreamUsage);
193 }
194 rewrites.push(Rewrite::Stream(true));
195 }
196 rewrites.push(Rewrite::ChatDialect(quirks.rewrite));
197 let raw_at = match quirks.rewrite {
198 BodyRewrite::Mira => RawAt::Ignored(
199 "Additional parameters are not supported by Mira and will be ignored",
200 ),
201 _ => RawAt::Top,
202 };
203 let body = request_params(
204 self,
205 &request,
206 |input| self.base(&request, input),
207 raw_at,
208 &rewrites,
209 )?;
210 let body = crate::providers::openai::options::check_body(
211 self,
212 &request,
213 body,
214 crate::providers::openai::options::Endpoint::ChatCompletions,
215 )?;
216 crate::providers::internal::trace_json(
217 crate::providers::internal::LogTarget::Completions,
218 "OpenAI Chat Completions request",
219 &body,
220 );
221 let request = builder.body(body.into_body())?;
222 let framing = match mode {
223 Mode::Streaming => Framing::Sse,
224 Mode::Unary => Framing::Whole,
225 };
226 Ok(Encoded::new(request, framing)
227 .with_request_id_header(self.provider.dialect.request_id_header)
228 .with_projection(ChatDecoder::project)
229 .with_route(Some(quirks.completion_path)))
230 }
231
232 fn reasoning_field(&self, model: &str) -> Option<&'static str> {
242 use crate::providers::openai::wire::{MOONSHOT, OPENROUTER};
243 let deepseek = self
244 .provider
245 .base_url
246 .to_ascii_lowercase()
247 .contains("deepseek.com");
248 let name = model.rsplit('/').next().unwrap_or(model);
249 let field = |vendor: &str, id: &str| {
250 crate::catalog::lookup(vendor, id)
251 .and_then(|spec| spec.compat.reasoning_field.as_deref())
252 };
253 let listed = field(self.provider.dialect.name, model)
254 .or_else(|| field(MOONSHOT.name, name))
255 .or_else(|| field(OPENROUTER.name, model));
256 self.provider
257 .dialect
258 .quirks
259 .reasoning_field
260 .or(deepseek.then_some("reasoning_content"))
261 .or(listed)
262 }
263
264 pub fn new(provider: OpenAIConfig, model: impl Into<String>) -> Self {
266 Self {
267 provider,
268 model: model.into(),
269 strict_tools: false,
270 tool_result_array_content: false,
271 }
272 }
273
274 pub fn with_strict_tools(mut self) -> Self {
277 self.strict_tools = true;
278 self
279 }
280
281 pub fn with_tool_result_array_content(mut self) -> Self {
283 self.tool_result_array_content = true;
284 self
285 }
286
287 fn base(
290 &self,
291 request: &CompletionRequest,
292 input: &mut BaseInput<'_>,
293 ) -> Result<Map<String, Value>, EncodeError> {
294 let quirks = &self.provider.dialect.quirks;
295 let mut model = request.model.clone().unwrap_or_else(|| self.model.clone());
296 if quirks.rewrite == BodyRewrite::HuggingFaceRouter {
297 model = self.provider.route().model_identifier(&model);
299 }
300 let passthrough = input.raw_tools()?;
301 let custom: Vec<String> = passthrough
303 .iter()
304 .filter(|tool| tool.str("type") == Some("custom"))
305 .filter_map(|tool| tool.at("/custom/name").and_then(Value::as_str))
306 .map(str::to_owned)
307 .collect();
308 let messages = self.messages(&request.chat_history, &model, &custom)?;
309
310 let mut tools: Vec<Value> = Vec::new();
311 let mut tool_choice = None;
312 if quirks.supports_tools {
313 for tool in &request.tools {
314 let mut parameters = tool.parameters.clone();
315 let mut function = json!({"name": tool.name, "description": tool.description});
316 if self.strict_tools {
317 crate::providers::openai::sanitize_schema(&mut parameters);
318 }
319 if let Some(function) = function.as_object_mut() {
320 function.insert("parameters".to_owned(), parameters);
321 if self.strict_tools {
322 function.insert("strict".to_owned(), Value::Bool(true));
323 }
324 }
325 tools.push(json!({"type": "function", "function": function}));
326 }
327 tools.extend(passthrough);
330 tool_choice = request
331 .tool_choice
332 .clone()
333 .map(tool_choice_value)
334 .transpose()?
335 .filter(|_| !tools.is_empty());
336 } else {
337 if !request.tools.is_empty() {
338 tracing::warn!("Tool use is not supported by this provider; tools will be ignored");
339 }
340 if request.tool_choice.is_some() {
341 tracing::warn!("Tool choice is not supported by this provider and will be ignored");
342 }
343 tools.extend(passthrough);
346 }
347
348 if request.output_schema.is_some() && !quirks.supports_response_format {
349 tracing::warn!(
350 "Structured outputs are not supported by this provider; ignoring output_schema"
351 );
352 }
353 let answered = messages
355 .iter()
356 .any(|message| message.str("role") == Some("tool"));
357 let response_format = match request.output_schema.clone() {
358 Some(schema)
359 if quirks.supports_response_format
360 && (quirks.response_format_with_tools || tools.is_empty() || answered) =>
361 {
362 let (name, schema) = crate::providers::openai::structured_output_schema(schema);
363 Some(json!({"type": "json_schema",
364 "json_schema": {"name": name, "strict": true, "schema": schema}}))
365 }
366 _ => None,
367 };
368
369 let fields = [
370 ("model", Some(Value::String(model))),
371 ("messages", Some(Value::Array(messages))),
372 ("tools", (!tools.is_empty()).then_some(Value::Array(tools))),
373 ("tool_choice", tool_choice),
374 ("temperature", request.temperature.map(Value::from)),
375 ("max_tokens", request.max_tokens.map(Value::from)),
376 ("response_format", response_format),
377 ];
378 Ok(fields
379 .into_iter()
380 .filter_map(|(key, value)| Some((key.to_owned(), value?)))
381 .collect())
382 }
383
384 fn messages(
387 &self,
388 history: &[Message],
389 model: &str,
390 custom: &[String],
391 ) -> Result<Vec<Value>, EncodeError> {
392 let ids = WireIds::for_target(history, self, model);
393 let mut messages = Vec::new();
394 for message in history {
395 match message {
396 Message::System { content } => messages.push(json!({"role": "system",
397 "content": [{"type": "text", "text": content}]})),
398 Message::User { content } => {
399 let mut parts = Vec::new();
400 for part in content {
401 if let UserContent::ToolResult(result) = part {
402 user_message(&mut messages, &mut parts);
403 let array = self.tool_result_array_content;
404 messages.push(tool_message(result, &ids, array)?);
405 } else {
406 parts.push(user_part(part)?);
407 }
408 }
409 user_message(&mut messages, &mut parts);
410 }
411 Message::Assistant(turn) => {
412 messages.extend(self.assistant(turn, &ids, custom, model));
413 }
414 }
415 }
416 if messages.is_empty() {
417 return Err(EncodeError::request(
418 "OpenAI Chat Completions request has no provider-compatible messages after \
419 conversion",
420 ));
421 }
422 Ok(messages)
423 }
424
425 fn assistant(
432 &self,
433 turn: &AssistantMessage,
434 ids: &WireIds,
435 custom: &[String],
436 model: &str,
437 ) -> Option<Value> {
438 let reasoning_field = self.reasoning_field(model);
439 let (mut text, mut parts, mut has_parts) = (String::new(), Vec::new(), false);
440 let mut reasoning: Vec<(String, String)> = Vec::new();
441 let (mut fields, mut calls) = (Map::new(), Vec::new());
442 for block in &turn.content {
443 let replay = block.replay(self, ids);
444 match block {
445 AssistantContent::Text(block) => {
446 text.push_str(&block.text);
447 parts.push(json!({"type": "text", "text": block.text}));
448 if let Replay::Item(item) = replay
450 && let Value::Object(item) = item.into_owned()
451 {
452 fields.extend(item);
453 }
454 }
455 AssistantContent::Reasoning(block) => {
456 let field = match replay {
457 Replay::Item(item) if item.get("type").is_some() => {
458 has_parts = true;
459 parts.push(item.into_owned());
460 None
461 }
462 Replay::Item(item) => {
463 let mut field = None;
464 for (key, value) in item.as_object().into_iter().flatten() {
465 if value.is_string() {
466 field.get_or_insert_with(|| key.clone());
467 } else {
468 let detail = Map::from_iter([(key.clone(), value.clone())]);
469 merge_fields(&mut fields, &detail);
470 }
471 }
472 field
473 }
474 Replay::Identity(identity)
475 if identity.get("type").and_then(Value::as_str) == Some("thinking") =>
476 {
477 has_parts = true;
478 parts.push(json!({"type": "thinking",
479 "thinking": [{"type": "text", "text": block.text}]}));
480 None
481 }
482 Replay::Identity(identity) => identity.keys().next().cloned(),
483 Replay::Rebuild => reasoning_field.map(str::to_owned),
484 };
485 if let Some(field) = field {
486 match reasoning.iter_mut().find(|(name, _)| *name == field) {
487 Some((_, joined)) => {
488 joined.push('\n');
489 joined.push_str(&block.text);
490 }
491 None => reasoning.push((field, block.text.clone())),
492 }
493 }
494 }
495 AssistantContent::ToolCall(call) => {
496 calls.push(call_item(call, replay, ids, custom));
497 }
498 AssistantContent::Opaque(opaque) if opaque.replay => {
501 if opaque.item.get("type").is_some() {
502 has_parts = true;
503 parts.push(opaque.item.clone());
504 }
505 }
506 AssistantContent::Opaque(_) | AssistantContent::Image(_) => {}
507 }
508 }
509 let mut message = Map::from_iter([("role".to_owned(), Value::from("assistant"))]);
510 if has_parts {
511 message.insert("content".to_owned(), Value::Array(parts));
512 } else if !text.is_empty() {
513 message.insert("content".to_owned(), Value::String(text));
514 }
515 for (field, text) in reasoning {
516 message.insert(field, Value::String(text));
517 }
518 message.extend(fields);
519 if !calls.is_empty() {
520 message.insert("tool_calls".to_owned(), Value::Array(calls));
521 }
522 let has_content = message.contains_key("audio")
523 || message.contains_key("tool_calls")
524 || match message.get("content") {
525 Some(Value::String(text)) => !text.is_empty(),
526 Some(Value::Array(parts)) => !parts.is_empty(),
527 _ => false,
528 };
529 if let Some(field) = reasoning_field {
530 message
531 .entry(field)
532 .or_insert_with(|| Value::String(String::new()));
533 }
534 has_content.then_some(Value::Object(message))
535 }
536}
537
538pub(crate) fn rewrite_body(
543 kind: BodyRewrite,
544 map: &mut Map<String, Value>,
545) -> Result<(), EncodeError> {
546 let forced = map
547 .get("tool_choice")
548 .and_then(|choice| choice.at("/function/name"))
549 .and_then(Value::as_str)
550 .map(str::to_owned);
551 match kind {
552 BodyRewrite::LlamaCpp => {
553 if let Some(name) = forced {
554 return Err(EncodeError::request(format!(
555 "llama.cpp cannot force a specific tool: `llama-server` accepts only \
556 `auto`, `none` or `required` for tool_choice and silently treats \
557 anything else as `auto`, so requesting `{name}` would return whichever \
558 tool the model picked. Use `ToolChoice::Required` to force a call, or \
559 advertise only `{name}` in `tools`."
560 )));
561 }
562 }
563 BodyRewrite::Moonshot => {
564 if forced.is_some() {
565 return Err(EncodeError::request(
566 "Moonshot does not support forcing a specific tool",
567 ));
568 }
569 if map.get("tool_choice").and_then(Value::as_str) == Some("required") {
570 tracing::warn!(
571 "Moonshot does not support tool_choice=required; coercing to auto with an \
572 additional steering message"
573 );
574 map.insert("tool_choice".to_owned(), Value::from("auto"));
575 if let Some(Value::Array(messages)) = map.get_mut("messages") {
576 messages.push(json!({"role": "user",
577 "content": "Please select a tool to handle the current issue."}));
578 }
579 }
580 }
581 BodyRewrite::Perplexity | BodyRewrite::Mira => {
587 let all = kind == BodyRewrite::Mira;
588 for content in messages_mut(map).filter_map(|message| message.get_mut("content")) {
589 if let Value::Array(parts) = content
590 && (all || parts.iter().all(|part| part.str("type") == Some("text")))
591 {
592 let texts: Vec<&str> =
593 parts.iter().filter_map(|part| part.str("text")).collect();
594 *content = Value::String(texts.join("\n"));
595 }
596 }
597 }
598 BodyRewrite::DeepSeek => finalize_deepseek(map),
599 BodyRewrite::Mistral => finalize_mistral(map),
600 BodyRewrite::Ollama => finalize_ollama(map)?,
601 BodyRewrite::None | BodyRewrite::OpenRouter | BodyRewrite::HuggingFaceRouter => {}
602 }
603 Ok(())
604}
605
606fn call_item(call: &ToolCall, replay: Replay<'_>, ids: &WireIds, custom: &[String]) -> Value {
612 let mut item = match replay {
613 Replay::Item(item) => match item.into_owned() {
614 Value::Object(item) => item,
615 _ => Map::new(),
616 },
617 Replay::Identity(identity) => identity,
618 Replay::Rebuild => Map::new(),
619 };
620 let name = call.function.name.as_str();
621 let custom = match item.get("type").and_then(Value::as_str) {
622 Some(kind) => kind == "custom",
623 None => {
624 let declared = custom.iter().any(|tool| tool == name);
625 let kind = if declared { "custom" } else { "function" };
626 item.insert("type".to_owned(), Value::from(kind));
627 declared
628 }
629 };
630 let id = ids.spell(&call.id);
631 item.insert("id".to_owned(), Value::String(id));
632 let (slot, key, value) = if custom {
633 let input = match call.function.arguments.get("input") {
634 Some(Value::String(input)) => input.clone(),
635 Some(input) => input.to_string(),
636 None => String::new(),
637 };
638 ("custom", "input", input)
639 } else {
640 let arguments = Value::Object(call.function.arguments.clone()).to_string();
641 ("function", "arguments", arguments)
642 };
643 let slot = item
644 .entry(slot)
645 .or_insert_with(|| Value::Object(Map::new()));
646 if !slot.is_object() {
647 *slot = Value::Object(Map::new());
648 }
649 if let Value::Object(fields) = slot {
650 fields.insert("name".to_owned(), Value::from(name));
651 fields.insert(key.to_owned(), Value::String(value));
652 }
653 Value::Object(item)
654}
655
656fn tool_choice_value(choice: crate::message::ToolChoice) -> Result<Value, EncodeError> {
658 use crate::message::ToolChoice;
659 Ok(match choice {
660 ToolChoice::Auto => Value::from("auto"),
661 ToolChoice::None => Value::from("none"),
662 ToolChoice::Required => Value::from("required"),
663 ToolChoice::Specific { function_names } => {
664 let [name] = function_names.as_slice() else {
665 return Err(EncodeError::request(
666 "Provider only supports forcing exactly one specific tool",
667 ));
668 };
669 json!({"type": "function", "function": {"name": name}})
670 }
671 })
672}
673
674fn finalize_ollama(map: &mut Map<String, Value>) -> Result<(), EncodeError> {
677 if let Some(key) = ["num_ctx", "options"]
678 .into_iter()
679 .find(|key| map.contains_key(*key))
680 {
681 return Err(EncodeError::request(format!(
682 "Ollama's OpenAI-compatible API ignores `{key}`; send it through the native \
683 route (`Ollama::native_completion`)"
684 )));
685 }
686 Ok(())
687}
688
689fn messages_mut(map: &mut Map<String, Value>) -> impl Iterator<Item = &mut Map<String, Value>> {
691 let messages = map.get_mut("messages").and_then(Value::as_array_mut);
692 messages
693 .into_iter()
694 .flatten()
695 .filter_map(Value::as_object_mut)
696}
697
698fn finalize_deepseek(map: &mut Map<String, Value>) {
702 let thinking = match map
703 .get("thinking")
704 .and_then(|thinking| thinking.str("type"))
705 {
706 Some(mode) => !mode.eq_ignore_ascii_case("disabled"),
707 None => map.get("model").and_then(Value::as_str) != Some("deepseek-chat"),
708 };
709 let forced = map
710 .get("tool_choice")
711 .is_some_and(|choice| choice.is_object() || choice.as_str() == Some("required"));
712 if thinking && forced {
713 tracing::debug!(
714 "dropping tool_choice: DeepSeek rejects a forced tool choice while thinking"
715 );
716 map.shift_remove("tool_choice");
717 }
718}
719
720fn finalize_mistral(map: &mut Map<String, Value>) {
723 let forces_a_tool_call = map
726 .get("tool_choice")
727 .is_some_and(|choice| !matches!(choice.as_str(), Some("auto" | "none")));
728 let has_tools = map
729 .get("tools")
730 .and_then(Value::as_array)
731 .is_some_and(|tools| !tools.is_empty());
732 let structured = map
733 .get("response_format")
734 .and_then(|format| format.str("type"))
735 .is_some_and(|kind| matches!(kind, "json_schema" | "json_object"));
736 if forces_a_tool_call && has_tools && structured {
737 tracing::debug!(
738 "relaxing tool_choice to `auto`: Mistral rejects a forced tool choice \
739 alongside a response format"
740 );
741 map.insert("tool_choice".to_owned(), Value::from("auto"));
742 }
743 for content in messages_mut(map).filter_map(|message| message.get_mut("content")) {
744 for part in content.as_array_mut().into_iter().flatten() {
745 *part = mistral_chunk(part);
746 }
747 }
748}
749
750fn mistral_chunk(part: &Value) -> Value {
756 let field = |pointer: &str| part.at(pointer).and_then(Value::as_str);
757 match part.str("type") {
758 Some("image_url") => {
759 json!({"type": "image_url", "image_url": part.get("image_url")})
760 }
761 Some("input_audio") => json!({"type": "input_audio",
762 "input_audio": field("/input_audio/data")}),
763 Some("file") => match field("/file/file_data") {
764 Some(data) => json!({"type": "document_url", "document_url": data,
765 "document_name": field("/file/filename")}),
766 None => json!({"type": "file", "file_id": field("/file/file_id")}),
767 },
768 _ => part.clone(),
769 }
770}
771
772impl Wire for Chat {
773 type Op = crate::operation::Completion;
774 type Payload = crate::wire::Encoded;
775 type Frame = crate::wire::WireFrame;
776 type Decoder<'id> = ChatDecoder;
777 type Reassembler = document::ChatCompletion;
778
779 fn describe(&self) -> Descriptor<'_> {
782 Descriptor::new(self.provider.dialect.name)
783 .model(self.model.as_str())
784 .capabilities(Capabilities::completion(
785 ProviderCapabilities::default().with_native_output_tool_composition(
786 self.provider.dialect.quirks.supports_response_format,
787 ),
788 ))
789 .replay(self)
790 }
791
792 fn encode(&self, request: CompletionRequest, mode: Mode) -> Result<Encoded, EncodeError> {
793 self.encode_with_headers(request, mode, OpenAIConfig::completion_headers)
794 }
795
796 fn decoder<'id>(&self) -> Self::Decoder<'id> {
797 ChatDecoder::new(self.provider.dialect.quirks)
798 }
799
800 fn reassembler(&self) -> Self::Reassembler {
801 document::ChatCompletion::new(self.provider.dialect.quirks)
802 }
803}
804
805impl crate::completion::ReplayTarget for Chat {
806 fn map_options(
808 &self,
809 request: &CompletionRequest,
810 fields: crate::completion::options::OptionFields<'_>,
811 ) -> crate::completion::options::OptionMap {
812 crate::providers::openai::options::chat_options(self, request, fields)
813 }
814
815 fn api(&self) -> crate::message::Api {
816 crate::message::Api::from_static("openai.chat")
817 }
818
819 fn provider(&self) -> &str {
820 self.provider.dialect.name
821 }
822
823 fn model(&self) -> &str {
824 &self.model
825 }
826
827 fn accepts(&self, model: &str) -> crate::completion::Accepts {
832 let quirks = &self.provider.dialect.quirks;
833 let user_images = reads_images(&self.provider.dialect, model);
834 crate::completion::Accepts {
835 user_images,
836 assistant_images: false,
837 tool_result_images: user_images && quirks.supports_image_tool_results,
838 tools: quirks.supports_tools,
839 }
840 }
841
842 fn encodes(&self, _model: &str, media: crate::completion::Media<'_>) -> bool {
849 use crate::completion::Media;
850 let dialect = &self.provider.dialect;
851 let rewrite = dialect.quirks.rewrite;
852 let parts = !matches!(rewrite, BodyRewrite::DeepSeek | BodyRewrite::Mira);
853 let files = parts
855 && rewrite != BodyRewrite::Perplexity
856 && ![super::dialects::COHERE.name, super::dialects::OLLAMA.name]
857 .contains(&dialect.name);
858 let linked = |source: &Source, typed: bool| match source {
859 Source::Url(_) => true,
860 Source::Base64(_) => typed,
861 Source::Raw(_) | Source::FileId(_) | Source::String(_) | Source::Unknown => false,
862 };
863 match media {
864 Media::Image(image, place) => {
865 parts
866 && place != crate::completion::Place::Assistant
867 && linked(&image.data, image.media_type.is_some())
868 && !(dialect.name == super::dialects::OLLAMA.name
870 && matches!(image.data, Source::Url(_)))
871 }
872 Media::Audio(audio) => files && matches!(audio.data, Source::Base64(_)),
873 Media::Video(video) => {
874 files
875 && rewrite != BodyRewrite::Mistral
876 && ![super::dialects::OPENAI.name, super::dialects::AZURE.name]
877 .contains(&dialect.name)
878 && linked(&video.data, video.media_type.is_some())
879 }
880 Media::Document(document) => {
881 let pdf = document.media_type == Some(DocumentMediaType::PDF);
882 match &document.data {
883 Source::String(_) => !pdf,
884 Source::FileId(_) => files && dialect.quirks.accepts_file_ids,
885 Source::Base64(_) => files && pdf,
886 Source::Url(_) => {
887 pdf && matches!(rewrite, BodyRewrite::OpenRouter | BodyRewrite::Mistral)
888 }
889 Source::Raw(_) | Source::Unknown => false,
890 }
891 }
892 }
893 }
894
895 fn normalize_tool_call_id(
901 &self,
902 id: &str,
903 _model: &str,
904 _: Option<&crate::message::Origin>,
905 ) -> String {
906 let sanitized =
907 |part: &str| crate::providers::internal::wire_ids::legal_call_id(part, usize::MAX);
908 if let Some((call, item)) = id.split_once('|') {
909 let call = sanitized(call);
910 let item = sanitized(item);
911 let combined = if item.is_empty() {
912 call.clone()
913 } else {
914 format!("{call}_{item}")
915 };
916 if combined.len() <= 40 {
917 return combined;
918 }
919 let hash: String = crate::providers::internal::wire_ids::short_hash(id)
920 .chars()
921 .take(8)
922 .collect();
923 let prefix: String = call.chars().take((40 - hash.len() - 1).max(1)).collect();
924 return format!("{prefix}_{hash}");
925 }
926 if self.provider.dialect.name == super::dialects::OPENAI.name {
927 return id.chars().take(40).collect();
928 }
929 id.to_owned()
930 }
931
932 fn identity(&self, item: &Value) -> Map<String, Value> {
937 if let Some(kind @ ("custom" | "thinking")) = item.str("type") {
938 return Map::from_iter([("type".to_owned(), Value::from(kind))]);
939 }
940 REASONING_TEXT_KEYS
941 .iter()
942 .find(|key| item.get(**key).is_some_and(Value::is_string))
943 .map(|key| Map::from_iter([((*key).to_owned(), Value::String(String::new()))]))
944 .unwrap_or_default()
945 }
946
947 fn call_id_slot(&self) -> Option<&'static str> {
948 Some("/id")
949 }
950
951 fn later_system(&self, _model: &str) -> crate::completion::LaterSystem {
952 self.provider.dialect.quirks.later_system
953 }
954
955 fn result_parts(&self, _model: &str) -> bool {
957 self.tool_result_array_content
958 }
959
960 fn sends_alone(&self, block: &AssistantContent) -> bool {
965 let ids = crate::providers::internal::wire_ids::WireIds::default();
966 match (block, block.replay(self, &ids)) {
967 (AssistantContent::ToolCall(_), _) => true,
968 (AssistantContent::Text(text), Replay::Item(item)) => {
969 !text.text.is_empty() || item.get("audio").is_some()
970 }
971 (AssistantContent::Text(text), _) => !text.text.is_empty(),
972 (AssistantContent::Reasoning(_), Replay::Item(item)) => item.get("type").is_some(),
973 (AssistantContent::Reasoning(_), Replay::Identity(identity)) => {
974 identity.get("type").and_then(Value::as_str) == Some("thinking")
975 }
976 (AssistantContent::Opaque(opaque), _) => opaque.item.get("type").is_some(),
977 (AssistantContent::Reasoning(_), Replay::Rebuild) | (AssistantContent::Image(_), _) => {
978 false
979 }
980 }
981 }
982}
983
984fn image_rule(vendor: &str) -> Option<(&str, fn(&str) -> bool)> {
994 use crate::providers::{minimax, moonshot, xiaomimimo, zai};
995 Some(match vendor {
996 "deepseek" | "mira" => return None,
997 "groq" => ("groq", |model| model.contains("llama-4")),
998 "mistral" | "mistralai" => ("mistral", |model| {
999 !(model.starts_with("codestral") || model.starts_with("devstral"))
1000 }),
1001 "openai" | "azure.openai" => ("openai", crate::providers::openai::reads_images),
1002 "xai" | "x-ai" => ("xai", crate::providers::xai::reads_images),
1003 "cohere" => ("cohere", crate::providers::cohere::reads_images),
1004 "zai" | "z-ai" => ("zai", zai::reads_images),
1005 "moonshot" | "moonshotai" => ("moonshot", moonshot::reads_images),
1006 "minimax" => ("minimax", minimax::reads_images),
1007 "xiaomimimo" | "xiaomi" => ("xiaomimimo", xiaomimimo::reads_images),
1008 vendor => (vendor, |_| true),
1009 })
1010}
1011
1012fn reads_images(dialect: &super::Dialect, model: &str) -> bool {
1017 match model.split_once('/') {
1018 Some((vendor, upstream)) if dialect.quirks.rewrite == BodyRewrite::OpenRouter => {
1019 crate::catalog::reads_images_or(dialect.name, model, |_| {
1020 image_rule(vendor).is_some_and(|(_, rule)| rule(upstream))
1021 })
1022 }
1023 _ => image_rule(dialect.name)
1024 .is_some_and(|(vendor, rule)| crate::catalog::reads_images_or(vendor, model, rule)),
1025 }
1026}
1027
1028pub enum ChatEvent {
1030 Chunk(Value),
1032 Whole(Value),
1034 Done,
1036 Failure(ProviderError),
1038 BareText(String),
1041}
1042
1043const REASONING_TEXT_KEYS: [&str; 3] = ["reasoning_content", "reasoning", "reasoning_text"];
1047
1048const REASONING_DETAILS: &str = "reasoning_details";
1050
1051#[derive(Debug, PartialEq)]
1053pub(crate) enum Part {
1054 Text(String),
1056 Thinking(String),
1058 Image,
1060 Unknown,
1062}
1063
1064impl Part {
1065 pub(crate) fn of(part: &Value) -> Self {
1068 let text = |key: &str| part.str(key).unwrap_or_default().to_owned();
1069 match part.str("type") {
1070 Some("text") => Self::Text(text("text")),
1071 None if part.str("text").is_some() => Self::Text(text("text")),
1072 Some("refusal") => Self::Text(text("refusal")),
1073 Some("thinking") => Self::Thinking(match part.get("thinking") {
1074 Some(Value::Array(chunks)) => chunks
1075 .iter()
1076 .filter_map(|chunk| chunk.str("text"))
1077 .collect(),
1078 _ => text("thinking"),
1079 }),
1080 Some("image_url") => Self::Image,
1081 _ => Self::Unknown,
1082 }
1083 }
1084}
1085
1086#[derive(Debug, PartialEq)]
1088pub(crate) enum CallKind {
1089 Function,
1091 Custom,
1093 Unknown,
1095}
1096
1097impl CallKind {
1098 pub(crate) fn of(call: &Value) -> Self {
1101 match call.str("type") {
1102 None | Some("function") => Self::Function,
1103 Some("custom") => Self::Custom,
1104 Some(_) => Self::Unknown,
1105 }
1106 }
1107}
1108
1109#[derive(Clone, Copy, PartialEq)]
1112enum Writing {
1113 Text,
1114 Thinking,
1115}
1116
1117struct OpenCall {
1121 at: usize,
1122 id: Option<String>,
1123 opaque: bool,
1124 arguments: String,
1125}
1126
1127impl OpenCall {
1128 fn complete(&self) -> bool {
1130 matches!(
1131 crate::json_utils::parse_tool_arguments(&self.arguments),
1132 Ok(Value::Object(_))
1133 )
1134 }
1135}
1136
1137#[derive(Default)]
1148pub struct ChatDecoder {
1149 quirks: Quirks,
1150 writing: Option<(Writing, usize)>,
1152 thinking_text: String,
1153 reasoning: Option<usize>,
1157 reasoning_field: Option<&'static str>,
1158 reasoning_text: String,
1159 reasoning_details: Map<String, Value>,
1160 audio_id: Option<Value>,
1162 calls: Vec<OpenCall>,
1164 usage: Option<Value>,
1165 reply_fields: Map<String, Value>,
1167 first_text: Option<usize>,
1170 last_text: Option<usize>,
1173 turn_citations: Vec<crate::wire::WireCitation>,
1176 finish: Option<FinishReason>,
1177 response_id: Option<String>,
1178 response_model: Option<String>,
1179 ended: bool,
1181 chunked: bool,
1183}
1184
1185impl ChatDecoder {
1186 fn new(quirks: Quirks) -> Self {
1187 Self {
1188 quirks,
1189 ..Self::default()
1190 }
1191 }
1192
1193 fn absorb(&mut self, frame: &Value) -> Option<Value> {
1196 if let Some(id) = frame.str("id") {
1197 self.response_id = Some(id.to_owned());
1198 }
1199 if let Some(model) = frame.str("model") {
1200 self.response_model = Some(model.to_owned());
1201 }
1202 let choice = frame
1205 .arr("choices")
1206 .iter()
1207 .find(|choice| {
1208 let index = choice.get("index").and_then(Value::as_u64);
1209 index.is_none_or(|index| index == 0)
1210 })
1211 .cloned();
1212 if let Some(usage) = frame
1213 .at("/usage")
1214 .or_else(|| choice.as_ref().and_then(|choice| choice.at("/usage")))
1215 {
1216 self.usage = Some(usage.clone());
1217 }
1218 for key in reported::REPLY_FIELDS {
1219 if let Some(value) = frame.get(key).filter(|value| !value.is_null()) {
1220 self.reply_fields.insert(key.to_owned(), value.clone());
1221 }
1222 }
1223 let choice = choice?;
1224 let reason = match choice
1227 .str("finish_reason")
1228 .filter(|reason| !reason.is_empty())
1229 {
1230 Some(reason) => Some(finish_reason(reason, &self.quirks)),
1231 None => choice
1232 .str("native_finish_reason")
1233 .filter(|reason| self.quirks.native_finish_reason && !reason.is_empty())
1234 .map(native_finish_reason),
1235 };
1236 if let Some(reason) = reason {
1237 self.finish = Some(reason);
1238 self.ended = true;
1239 }
1240 Some(choice)
1241 }
1242
1243 fn delta(
1247 &mut self,
1248 delta: &Map<String, Value>,
1249 out: &mut Out<'_, Completion>,
1250 ) -> Result<(), ProviderError> {
1251 let reasoning = REASONING_TEXT_KEYS.iter().find_map(|key| {
1252 delta
1253 .get(*key)
1254 .and_then(Value::as_str)
1255 .filter(|text| !text.is_empty())
1256 .map(|text| (*key, text))
1257 });
1258 if let Some((key, text)) = reasoning {
1259 self.reason(text, out)?;
1260 self.reasoning_field.get_or_insert(key);
1261 }
1262 if let Some(Value::Array(details)) = delta.get(REASONING_DETAILS)
1263 && !details.is_empty()
1264 {
1265 self.reason("", out)?;
1266 merge_details(&mut self.reasoning_details, details.clone());
1267 }
1268 let audio = delta.get("audio");
1269 if let Some(id) = audio.and_then(|audio| audio.at("/id")) {
1270 self.audio_id = Some(id.clone());
1271 }
1272 match delta.get("content") {
1273 Some(Value::Array(parts)) => {
1274 for part in parts {
1275 self.part(part, out)?;
1276 }
1277 }
1278 Some(part @ Value::Object(_)) => self.part(part, out)?,
1279 _ => {
1280 let text = ["content", "refusal"]
1281 .iter()
1282 .find_map(|key| {
1283 delta
1284 .get(*key)
1285 .and_then(Value::as_str)
1286 .filter(|text| !text.is_empty())
1287 })
1288 .or_else(|| audio.and_then(|audio| audio.str("transcript")))
1289 .filter(|text| !text.is_empty());
1290 if let Some(text) = text {
1291 self.write(Writing::Text, text, out)?;
1292 }
1293 }
1294 }
1295 for image in delta
1296 .get("images")
1297 .and_then(Value::as_array)
1298 .into_iter()
1299 .flatten()
1300 {
1301 self.image(image, out)?;
1302 }
1303 for call in delta
1304 .get("tool_calls")
1305 .and_then(Value::as_array)
1306 .into_iter()
1307 .flatten()
1308 {
1309 self.call(call, out)?;
1310 }
1311 self.turn_citations.extend(
1312 delta
1313 .get("annotations")
1314 .and_then(Value::as_array)
1315 .into_iter()
1316 .flatten()
1317 .filter_map(reported::annotation),
1318 );
1319 self.cite_turn(out);
1320 Ok(())
1321 }
1322
1323 fn cite_turn(&mut self, out: &mut Out<'_, Completion>) {
1326 if let Some(text) = self.first_text {
1327 for citation in self.turn_citations.drain(..) {
1328 out.cite(text, citation);
1329 }
1330 }
1331 }
1332
1333 fn reason(&mut self, text: &str, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
1335 let at = match self.reasoning {
1336 Some(at) => at,
1337 None => {
1338 let at = out.fresh_index();
1339 out.open(at, Block::Reasoning { redacted: false }, Value::Null)?;
1340 out.lead(at)?;
1341 self.reasoning = Some(at);
1342 at
1343 }
1344 };
1345 self.reasoning_text.push_str(text);
1346 out.push(at, text)
1347 }
1348
1349 fn close_reasoning(&mut self, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
1352 let Some(at) = self.reasoning.take() else {
1353 return Ok(());
1354 };
1355 let text = std::mem::take(&mut self.reasoning_text);
1356 let mut item = Map::new();
1357 if let Some(field) = self.reasoning_field.take() {
1358 item.insert(field.to_owned(), text.into());
1359 }
1360 item.append(&mut self.reasoning_details);
1361 out.edit(at, |slot| *slot = Value::Object(item))?;
1362 out.finish(at)
1363 }
1364
1365 fn write(
1368 &mut self,
1369 writing: Writing,
1370 text: &str,
1371 out: &mut Out<'_, Completion>,
1372 ) -> Result<(), ProviderError> {
1373 let index = match self.writing {
1374 Some((current, index)) if current == writing => index,
1375 _ => {
1376 self.close_writing(out)?;
1377 let index = out.fresh_index();
1378 let block = match writing {
1379 Writing::Thinking => Block::Reasoning { redacted: false },
1380 Writing::Text => {
1381 self.first_text.get_or_insert(index);
1382 self.last_text = Some(index);
1383 Block::Text
1384 }
1385 };
1386 out.open(index, block, Value::Null)?;
1387 self.writing = Some((writing, index));
1388 index
1389 }
1390 };
1391 if writing == Writing::Thinking {
1392 self.thinking_text.push_str(text);
1393 }
1394 out.push(index, text)
1395 }
1396
1397 fn close_writing(&mut self, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
1400 let (index, item) = match self.writing.take() {
1401 None => return Ok(()),
1402 Some((Writing::Text, index)) => (
1403 index,
1404 self.audio_id
1405 .take()
1406 .map(|id| json!({ "audio": { "id": id } })),
1407 ),
1408 Some((Writing::Thinking, index)) => {
1409 let text = std::mem::take(&mut self.thinking_text);
1410 (
1411 index,
1412 Some(json!({"type": "thinking", "thinking": [{"type": "text", "text": text}]})),
1413 )
1414 }
1415 };
1416 if let Some(item) = item {
1417 out.edit(index, |slot| *slot = item)?;
1418 }
1419 out.finish(index)
1420 }
1421
1422 #[deny(clippy::wildcard_enum_match_arm)]
1424 fn part(&mut self, part: &Value, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
1425 match Part::of(part) {
1426 Part::Text(text) => self.write(Writing::Text, &text, out),
1427 Part::Thinking(text) => self.write(Writing::Thinking, &text, out),
1428 Part::Image => self.image(part, out),
1429 Part::Unknown => {
1430 self.close_writing(out)?;
1431 if let Some(citation) = reported::reference(part) {
1432 match self.last_text {
1433 Some(text) => out.cite(text, citation),
1434 None => self.turn_citations.push(citation),
1435 }
1436 }
1437 let index = out.fresh_index();
1438 out.whole(index, Block::Opaque { replay: true }, part.clone(), "")
1439 }
1440 }
1441 }
1442
1443 fn image(&mut self, part: &Value, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
1446 use crate::message::{Image, ImageMediaType};
1447 let url = part
1448 .at("/image_url/url")
1449 .or_else(|| part.get("image_url"))
1450 .and_then(Value::as_str)
1451 .unwrap_or_default();
1452 let inline = url
1453 .strip_prefix("data:")
1454 .and_then(|rest| rest.split_once(";base64,"));
1455 let mut image = Image::default();
1456 let data = match inline {
1457 Some((mime, data)) => {
1458 image.data = Source::Base64(String::new());
1459 image.media_type = ImageMediaType::from_mime_type(mime);
1460 data
1461 }
1462 None => {
1463 image.data = Source::Url(url.to_owned());
1464 ""
1465 }
1466 };
1467 self.close_writing(out)?;
1468 let index = out.fresh_index();
1469 out.whole(index, Block::Image(image), part.clone(), data)
1470 }
1471
1472 fn call(&mut self, call: &Value, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
1480 self.close_writing(out)?;
1481 let id = match call.get("id") {
1482 Some(Value::String(id)) if !id.is_empty() && id != "null" => Some(id.clone()),
1483 Some(id @ Value::Number(_)) => Some(id.to_string()),
1484 _ => None,
1485 };
1486 let index = call
1487 .get("index")
1488 .and_then(Value::as_u64)
1489 .and_then(|index| usize::try_from(index).ok());
1490 let arguments = call
1491 .at("/function/arguments")
1492 .map(crate::json_utils::value_to_json_string);
1493 let name = call
1494 .at("/function/name")
1495 .or_else(|| call.at("/custom/name"))
1496 .and_then(Value::as_str);
1497 let starts = name.is_some_and(|name| !name.is_empty())
1498 || arguments.as_deref().is_some_and(|text| !text.is_empty());
1499 let at = match index {
1500 Some(index) => index,
1501 None => match id.as_deref() {
1502 Some(id) => self
1503 .calls
1504 .iter()
1505 .find(|open| open.id.as_deref() == Some(id)),
1506 None => self
1507 .calls
1508 .last()
1509 .filter(|open| open.opaque || !open.complete() || !starts),
1510 }
1511 .map_or_else(|| out.fresh_index(), |open| open.at),
1512 };
1513 let opaque = match self.calls.iter_mut().find(|open| open.at == at) {
1514 Some(open) => {
1515 if id.is_some() {
1516 open.id.clone_from(&id);
1517 }
1518 open.arguments
1519 .push_str(arguments.as_deref().unwrap_or_default());
1520 open.opaque
1521 }
1522 None => {
1523 let opaque = CallKind::of(call) == CallKind::Unknown;
1524 if opaque {
1525 out.open(at, Block::Opaque { replay: false }, Value::Null)?;
1526 }
1527 self.calls.push(OpenCall {
1528 at,
1529 id: id.clone(),
1530 opaque,
1531 arguments: arguments.clone().unwrap_or_default(),
1532 });
1533 opaque
1534 }
1535 };
1536 if !opaque {
1537 out.fragment(
1538 Some(at),
1539 CallFragment {
1540 id: id.as_deref(),
1541 name,
1542 arguments: arguments.as_deref(),
1543 },
1544 )?;
1545 }
1546 let mut fields = call.as_object().cloned().unwrap_or_default();
1547 fields.shift_remove("index");
1549 out.edit(at, |item| {
1550 if !item.is_object() {
1551 *item = Value::Object(Map::new());
1552 }
1553 if let Value::Object(item) = item {
1554 merge_fields(item, &fields);
1555 }
1556 })
1557 }
1558
1559 fn close_calls(
1566 &mut self,
1567 out: &mut Out<'_, Completion>,
1568 all: bool,
1569 ) -> Result<(), ProviderError> {
1570 let cut = self.finish == Some(FinishReason::Length);
1571 let mut waiting = Vec::new();
1572 for call in std::mem::take(&mut self.calls) {
1573 if call.opaque {
1574 out.finish(call.at)?;
1575 continue;
1576 }
1577 let (mut input, mut complete) = (None, false);
1578 out.edit(call.at, |item| {
1579 if CallKind::of(item) == CallKind::Custom {
1580 let custom = item.at("/custom/input").cloned();
1581 input = Some(custom.unwrap_or_else(|| Value::from("")));
1582 }
1583 complete = match item.at("/function/arguments") {
1584 Some(Value::String(text)) => matches!(
1585 crate::json_utils::parse_tool_arguments(text),
1586 Ok(Value::Object(_))
1587 ),
1588 Some(arguments) => arguments.is_object(),
1589 None => input.is_some(),
1590 };
1591 })?;
1592 if !all && !complete {
1593 waiting.push(call);
1594 continue;
1595 }
1596 if let Some(input) = input {
1597 out.announce(call.at, json!({ "input": input }))?;
1598 }
1599 if cut && !complete {
1600 out.close(call.at)?;
1601 } else {
1602 out.finish(call.at)?;
1603 }
1604 }
1605 self.calls = waiting;
1606 Ok(())
1607 }
1608
1609 fn chunk(&mut self, frame: &Value, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
1611 let Some(choice) = self.absorb(frame) else {
1612 return Ok(());
1613 };
1614 if let Some(delta) = choice.obj("delta") {
1615 self.delta(delta, out)?;
1616 }
1617 if self.finish == Some(FinishReason::ToolCalls) {
1621 self.close_calls(out, false)?;
1622 }
1623 Ok(())
1624 }
1625
1626 fn whole(
1629 &mut self,
1630 frame: &Value,
1631 mut out: Out<'_, Completion>,
1632 ) -> Result<Flow, ProviderError> {
1633 let Some(choice) = self.absorb(frame) else {
1634 return Err(ProviderError::Response(
1635 "Response contained no choices".to_owned(),
1636 ));
1637 };
1638 let Some(mut message) = choice.obj("message").cloned() else {
1639 return Err(ProviderError::Response(
1640 "Response did not contain a valid message or tool call".to_owned(),
1641 ));
1642 };
1643 self.ended = true;
1644 if let Some(Value::Array(calls)) = message.get_mut("tool_calls") {
1645 for (index, call) in calls.iter_mut().enumerate() {
1646 if let Some(call) = call.as_object_mut() {
1647 call.insert("index".to_owned(), index.into());
1648 }
1649 }
1650 }
1651 self.delta(&message, &mut out)?;
1652 self.close_calls(&mut out, true)?;
1653 self.end(out)
1654 }
1655
1656 fn end(&mut self, mut out: Out<'_, Completion>) -> Result<Flow, ProviderError> {
1659 self.close_writing(&mut out)?;
1660 self.close_reasoning(&mut out)?;
1661 self.cite_turn(&mut out);
1662 if let Some(text) = self.first_text {
1663 for citation in reported::listed(&self.reply_fields) {
1664 out.cite(text, citation);
1665 }
1666 }
1667 let cost = reported::cost(self.usage.as_ref(), &self.reply_fields);
1668 let usage = self
1669 .usage
1670 .as_ref()
1671 .map(|usage| normalized_usage(usage, &self.quirks))
1672 .unwrap_or_default()
1673 .cost(cost);
1674 Ok(out.end(Finish {
1675 usage,
1676 reason: self.finish.take(),
1677 response_id: self.response_id.take(),
1678 model: self.response_model.take(),
1679 ..Finish::default()
1680 }))
1681 }
1682
1683 fn finish(&mut self, mut out: Out<'_, Completion>) -> Result<Flow, ProviderError> {
1686 if !self.ended {
1687 return Err(ProviderError::Truncated);
1688 }
1689 self.close_calls(&mut out, true)?;
1690 self.end(out)
1691 }
1692
1693 fn done(&mut self, out: Out<'_, Completion>) -> Result<Flow, ProviderError> {
1697 if !self.ended && self.chunked && self.quirks.done_without_finish_reason {
1698 self.finish = Some(if self.calls.is_empty() {
1699 FinishReason::Stop
1700 } else {
1701 FinishReason::ToolCalls
1702 });
1703 self.ended = true;
1704 }
1705 self.finish(out)
1706 }
1707}
1708
1709fn merge_details(block: &mut Map<String, Value>, details: Vec<Value>) {
1714 let Value::Array(merged) = block
1715 .entry(REASONING_DETAILS)
1716 .or_insert_with(|| Value::Array(Vec::new()))
1717 else {
1718 return;
1719 };
1720 for detail in details {
1721 let kind = detail.str("type");
1722 let continues = matches!(kind, Some("reasoning.text" | "reasoning.summary"))
1723 && merged.last().is_some_and(|last| {
1724 last.str("type") == kind && last.get("index") == detail.get("index")
1725 });
1726 match (continues, merged.last_mut(), detail) {
1727 (true, Some(Value::Object(last)), Value::Object(fields)) => {
1728 for (key, value) in fields {
1729 let missing = last
1730 .get(&key)
1731 .is_none_or(|existing| existing.is_null() || existing.as_str() == Some(""));
1732 match (last.get_mut(&key), value) {
1733 (Some(Value::String(text)), Value::String(more))
1734 if matches!(key.as_str(), "text" | "summary") =>
1735 {
1736 text.push_str(&more);
1737 }
1738 (_, value) if missing => {
1739 last.insert(key, value);
1740 }
1741 _ => {}
1742 }
1743 }
1744 }
1745 (_, _, detail) => merged.push(detail),
1746 }
1747 }
1748}
1749
1750#[deny(clippy::wildcard_enum_match_arm)]
1751impl<'id> Decoder<'id, Completion> for ChatDecoder {
1752 type Event = ChatEvent;
1753
1754 fn classify(&self, frame: WireFrame) -> WireEvent<ChatEvent> {
1755 let data = frame.as_str();
1756 if data == "[DONE]" {
1758 return WireEvent::Known(ChatEvent::Done);
1759 }
1760 if let Some(error) = provider_error_envelope(&data) {
1763 return WireEvent::Known(ChatEvent::Failure(error));
1764 }
1765 if self.quirks.accepts_bare_string_reply
1766 && let Ok(Value::String(text)) = serde_json::from_str::<serde_json::Value>(&data)
1767 {
1768 return WireEvent::Known(ChatEvent::BareText(text));
1769 }
1770 classify_chat_completions_frame::<Value>(&data).map(|frame| {
1771 let whole = match frame.str("object") {
1774 Some(object) => object == "chat.completion",
1775 None => frame
1776 .arr("choices")
1777 .iter()
1778 .any(|choice| choice.obj("message").is_some()),
1779 };
1780 if whole {
1781 ChatEvent::Whole(frame)
1782 } else {
1783 ChatEvent::Chunk(frame)
1784 }
1785 })
1786 }
1787
1788 fn decode(
1789 &mut self,
1790 event: ChatEvent,
1791 mut out: Out<'id, Completion>,
1792 ) -> Result<Flow, ProviderError> {
1793 match event {
1794 ChatEvent::Chunk(frame) => {
1795 self.chunked = true;
1796 self.chunk(&frame, &mut out)?;
1797 Ok(Flow::More)
1798 }
1799 ChatEvent::Whole(frame) => self.whole(&frame, out),
1800 ChatEvent::Done => self.done(out),
1801 ChatEvent::BareText(text) => {
1804 self.ended = true;
1805 self.finish = Some(FinishReason::Stop);
1806 self.write(Writing::Text, &text, &mut out)?;
1807 self.end(out)
1808 }
1809 ChatEvent::Failure(error) => Err(error),
1810 }
1811 }
1812
1813 fn eof(&mut self, out: Out<'id, Completion>) -> Result<Flow, ProviderError> {
1816 self.finish(out)
1817 }
1818}
1819
1820fn normalized_usage(usage: &Value, quirks: &Quirks) -> crate::completion::Usage {
1830 let count = |pointer: &str| usage.at(pointer).and_then(Value::as_u64);
1831 let detail = |object: &str, key: &str| {
1832 count(&format!("/{object}/{key}")).or(usage.obj(object).map(|_| 0))
1833 };
1834 let (prompt, completion, total) = (
1835 count("/prompt_tokens"),
1836 count("/completion_tokens"),
1837 count("/total_tokens"),
1838 );
1839 let audio = count("/prompt_tokens_details/audio_tokens").unwrap_or(0);
1840 let input = prompt.map(|prompt| {
1841 let beside = prompt.saturating_add(audio);
1842 let accounted = beside.saturating_add(completion.unwrap_or(0));
1843 if audio != 0 && Some(accounted) == total {
1844 beside
1845 } else {
1846 prompt
1847 }
1848 });
1849 crate::completion::Usage {
1850 input_tokens: input,
1851 output_tokens: completion
1852 .or_else(|| total.map(|total| total.saturating_sub(input.unwrap_or(0)))),
1853 total_tokens: total,
1854 cached_input_tokens: detail("prompt_tokens_details", "cached_tokens")
1855 .or(count("/num_cached_tokens"))
1856 .or(count("/prompt_cache_hit_tokens")),
1857 cache_creation_input_tokens: count("/prompt_tokens_details/cache_write_tokens"),
1858 reasoning_tokens: detail("completion_tokens_details", "reasoning_tokens")
1859 .filter(|_| quirks.reliable_reasoning_count),
1860 ..Default::default()
1861 }
1862}
1863
1864impl ChatDecoder {
1865 pub(crate) fn project(payload: &[u8], sink: &mut ObservationSink<'_>) {
1869 let Ok(payload) = serde_json::from_slice::<Value>(payload) else {
1870 return;
1871 };
1872 if let Some(usage) = payload.at("/usage") {
1873 let count = |pointer: &str| usage.at(pointer).and_then(Value::as_u64);
1874 sink.emit(AdapterEvent::Usage {
1875 usage: AdapterUsage {
1876 input_tokens: count("/prompt_tokens"),
1877 output_tokens: count("/completion_tokens"),
1878 total_tokens: count("/total_tokens"),
1879 cached_input_tokens: count("/prompt_tokens_details/cached_tokens"),
1880 reasoning_tokens: count("/completion_tokens_details/reasoning_tokens"),
1881 tool_input_tokens: None,
1882 },
1883 });
1884 }
1885 let verdict = match payload
1888 .at("/choices/0/finish_reason")
1889 .and_then(Value::as_str)
1890 {
1891 Some(reason) => AdapterVerdict {
1892 finish_reason: Some(sink.scrub(reason)),
1893 block_reason: None,
1894 detail: None,
1895 model: payload.str("model").map(|value| sink.scrub(value)),
1896 },
1897 None => AdapterVerdict::default(),
1898 };
1899 let response_id = payload.str("id").map(|value| sink.scrub(value));
1900 sink.provider(verdict, response_id);
1901 if let Some(error) = payload
1902 .get("error")
1903 .and_then(|error| ObservedError::deserialize(error).ok())
1904 {
1905 error.emit(sink);
1906 }
1907 }
1908}
1909
1910mod document;
1911mod reported;
1912
1913#[cfg(test)]
1914mod tests;
1915
1916#[cfg(test)]
1917mod hard_case_tests;
1918
1919#[cfg(test)]
1920mod history_tests;
1921
1922#[cfg(test)]
1923mod request_tests;