1"""OpenAI Chat Completions request and response mapper.
2
3Converts the OpenAI Chat Completions wire DTOs
4(:class:`OpenAIChatRequest` / :class:`OpenAIChatResponse`) into the
5canonical relay IR and back. Stream conversion is handled by the shared
6stream lifecycle task and reports ``unsupported_feature`` until then.
7
8The mapper class composes direction mixins: :class:`WireToIRMixin`
9(parsing) and :class:`IRToWireMixin` (building), with shared free
10helpers in :mod:`_helpers`.
11"""
12
13from __future__ import annotations
14
15from dataclasses import replace
16from typing import Any, cast
17
18from lexigram.ai.relay.context import ConversionContext
19from lexigram.ai.relay.errors import translate, unsupported_feature, unsupported_format
20from lexigram.ai.relay.mappers.base import new_uuid, record_loss
21from lexigram.ai.relay.mappers.openai_chat._helpers import (
22 _TARGET,
23 _extract_text,
24 _tool_call_to_wire,
25 _tool_calls_to_ir,
26)
27from lexigram.ai.relay.mappers.openai_chat.ir_to_wire import IRToWireMixin
28from lexigram.ai.relay.mappers.openai_chat.wire_to_ir import WireToIRMixin
29from lexigram.contracts.ai.exceptions import RelayError
30from lexigram.contracts.ai.llm import ChatMessage, ToolCall
31from lexigram.contracts.ai.relay.dto import (
32 OpenAIChatChoice,
33 OpenAIChatMessage,
34 OpenAIChatRequest,
35 OpenAIChatResponse,
36)
37from lexigram.contracts.ai.relay.ir import (
38 RelayRequest,
39 RelayResponse,
40 StreamDelta,
41 StreamState,
42 normalize_finish_reason,
43)
44from lexigram.contracts.ai.thinking import ThinkingConfig, ThinkingResult
45from lexigram.contracts.core.result import Err, Ok, Result
46
47__all__ = ["OpenAIChatMapper"]
48
49
50class OpenAIChatMapper(WireToIRMixin, IRToWireMixin):
51 """Bidirectional OpenAI Chat Completions converter.
52
53 Attributes:
54 format: The wire format this mapper handles.
55 """
56
57 format = _TARGET
58
59 def request_to_ir(
60 self, payload: Any, *, context: ConversionContext
61 ) -> Result[RelayRequest, RelayError]:
62 """Convert an ``OpenAIChatRequest`` into canonical ``RelayRequest``.
63
64 Args:
65 payload: A wire request DTO.
66 context: Per-conversion context with loss sink.
67
68 Returns:
69 Ok(request) on success, Err(relay_error) on malformed payload.
70 """
71 if not isinstance(payload, OpenAIChatRequest):
72 return Err(
73 unsupported_format(
74 f"expected OpenAIChatRequest, got {type(payload).__name__}"
75 )
76 )
77 system_parts: list[str] = []
78 messages: list[ChatMessage] = []
79 for position, message in enumerate(payload.messages):
80 if message.role == "system":
81 text = _extract_text(
82 message.content,
83 context,
84 field=f"system_message[{position}].content",
85 )
86 system_parts.append(text)
87 if position > 0:
88 record_loss(
89 context,
90 field="system_message",
91 target=_TARGET,
92 reason="system_message_reordered",
93 )
94 continue
95 content: str | list[Any]
96 if isinstance(message.content, list):
97 content = self._wire_parts_to_ir(message.content, context)
98 elif message.content is None:
99 content = ""
100 else:
101 content = message.content
102 tool_calls = _tool_calls_to_ir(message.tool_calls)
103 messages.append(
104 ChatMessage(
105 role=message.role,
106 content=cast("str | list[Any]", content),
107 name=message.name,
108 tool_call_id=message.tool_call_id,
109 tool_calls=tool_calls,
110 metadata=dict(message.passthrough) or None,
111 )
112 )
113 max_tokens = self._normalize_max_tokens(payload, context)
114 stop_sequences = (
115 [payload.stop]
116 if isinstance(payload.stop, str)
117 else (
118 [s for s in payload.stop if isinstance(s, str)] if payload.stop else []
119 )
120 )
121 include_usage = False
122 if isinstance(payload.stream_options, dict):
123 include_usage = bool(payload.stream_options.get("include_usage", False))
124 thinking: ThinkingConfig | None = None
125 reasoning = payload.reasoning
126 if isinstance(reasoning, dict):
127 thinking = ThinkingConfig(effort=reasoning.get("effort"))
128 metadata: dict[str, Any] = {}
129 if reasoning is not None:
130 metadata["reasoning"] = reasoning
131 if payload.stream_options is not None:
132 metadata["stream_options"] = payload.stream_options
133 if payload.service_tier is not None:
134 metadata["service_tier"] = payload.service_tier
135 return Ok(
136 RelayRequest(
137 model=context.normalize_model(payload.model),
138 messages=messages,
139 system="\n".join(system_parts) if system_parts else None,
140 tools=self._tools_to_ir(payload.tools, context),
141 tool_choice=payload.tool_choice,
142 temperature=payload.temperature,
143 top_p=payload.top_p,
144 max_tokens=max_tokens,
145 stop_sequences=stop_sequences,
146 response_format=payload.response_format,
147 stream=payload.stream,
148 include_usage=include_usage,
149 parallel_tool_calls=payload.parallel_tool_calls,
150 thinking=thinking,
151 metadata=metadata,
152 passthrough=dict(payload.passthrough),
153 )
154 )
155
156 def ir_to_request(
157 self, request: RelayRequest, *, context: ConversionContext
158 ) -> Result[Any, RelayError]:
159 """Convert canonical ``RelayRequest`` into an ``OpenAIChatRequest``.
160
161 Args:
162 request: Canonical request IR.
163 context: Per-conversion context with loss sink.
164
165 Returns:
166 Ok(request) on success, Err(relay_error) on failure.
167 """
168 try:
169 messages: list[OpenAIChatMessage] = []
170 if request.system:
171 messages.append(
172 OpenAIChatMessage(role="system", content=request.system)
173 )
174 for message in request.messages:
175 prepared = message
176 if message.role == "assistant" and message.tool_calls:
177 if any(not tool_call.id for tool_call in message.tool_calls):
178 prepared = replace(
179 message,
180 tool_calls=[
181 tool_call
182 if tool_call.id
183 else replace(tool_call, id=f"call_{index + 1}")
184 for index, tool_call in enumerate(message.tool_calls)
185 ],
186 )
187 elif message.role == "tool" and not message.tool_call_id:
188 prepared = replace(message, tool_call_id="call_0")
189 messages.append(self._message_from_ir(prepared, context))
190 stream_options = self._stream_options_from_ir(request)
191 reasoning = self._reasoning_from_ir(request, context)
192 if request.metadata.get("max_tokens_kind") == "max_completion_tokens":
193 max_completion_tokens: int | None = request.max_tokens
194 max_tokens: int | None = None
195 else:
196 max_completion_tokens = None
197 max_tokens = request.max_tokens
198 return Ok(
199 OpenAIChatRequest(
200 model=context.resolve_model(request.model),
201 messages=messages,
202 temperature=request.temperature,
203 top_p=request.top_p,
204 max_tokens=max_tokens,
205 max_completion_tokens=max_completion_tokens,
206 stream=request.stream,
207 stream_options=stream_options,
208 tools=(
209 [self._tool_from_ir(tool) for tool in request.tools]
210 if request.tools
211 else None
212 ),
213 tool_choice=request.tool_choice,
214 parallel_tool_calls=request.parallel_tool_calls,
215 stop=self._stop_from_ir(request.stop_sequences),
216 response_format=request.response_format,
217 reasoning=reasoning,
218 service_tier=request.metadata.get("service_tier"),
219 passthrough={
220 **request.passthrough,
221 **{
222 key: value
223 for key, value in request.metadata.items()
224 if key
225 not in {
226 "service_tier",
227 "reasoning",
228 "stream_options",
229 "generation_config",
230 "safety_settings",
231 "tool_config",
232 "max_tokens_kind",
233 }
234 },
235 },
236 )
237 )
238 except (RelayError, ValueError, TypeError, KeyError) as exc:
239 return Err(translate(exc, detail="ir_to_request"))
240
241 def response_to_ir(
242 self, payload: Any, *, context: ConversionContext
243 ) -> Result[RelayResponse, RelayError]:
244 """Convert an ``OpenAIChatResponse`` into canonical ``RelayResponse``.
245
246 Args:
247 payload: A wire response DTO.
248 context: Per-conversion context with loss sink.
249
250 Returns:
251 Ok(response) on success, Err(relay_error) on malformed payload.
252 """
253 if not isinstance(payload, OpenAIChatResponse):
254 return Err(
255 unsupported_format(
256 f"expected OpenAIChatResponse, got {type(payload).__name__}"
257 )
258 )
259 try:
260 passthrough = dict(payload.passthrough)
261 if payload.system_fingerprint is not None:
262 passthrough["system_fingerprint"] = payload.system_fingerprint
263 choice = payload.choices[0] if payload.choices else None
264 if len(payload.choices) > 1:
265 record_loss(
266 context,
267 field="choices",
268 target=_TARGET,
269 reason="multiple_choices_collapsed",
270 )
271 message = choice.message if choice is not None else None
272 content = ""
273 tool_calls: list[ToolCall] = []
274 thinking: ThinkingResult | None = None
275 if message is not None:
276 content = self._message_text_to_ir(message, context)
277 tool_calls = list(_tool_calls_to_ir(message.tool_calls) or [])
278 thinking = self._reasoning_from_message(message, payload.usage)
279 return Ok(
280 RelayResponse(
281 model=payload.model,
282 id=payload.id,
283 created=payload.created,
284 content=content,
285 thinking=thinking,
286 tool_calls=tool_calls,
287 finish_reason=normalize_finish_reason(
288 choice.finish_reason if choice is not None else None
289 ),
290 usage=self._usage_from_wire(payload.usage),
291 passthrough=passthrough,
292 )
293 )
294 except (RelayError, ValueError, TypeError, KeyError) as exc:
295 return Err(translate(exc, detail="response_to_ir"))
296
297 def ir_to_response(
298 self, response: RelayResponse, *, context: ConversionContext
299 ) -> Result[Any, RelayError]:
300 """Convert canonical ``RelayResponse`` into an ``OpenAIChatResponse``.
301
302 Args:
303 response: Canonical response IR.
304 context: Per-conversion context with loss sink.
305
306 Returns:
307 Ok(response) on success, Err(relay_error) on failure.
308 """
309 try:
310 passthrough = dict(response.passthrough)
311 system_fingerprint = passthrough.pop("system_fingerprint", None)
312 content: str | None = response.content or None
313 tool_calls: list[dict[str, Any]] = []
314 for tool in response.tool_calls:
315 wire = _tool_call_to_wire(tool)
316 if not wire["id"]:
317 wire["id"] = f"call_{new_uuid()}"
318 tool_calls.append(wire)
319 message = OpenAIChatMessage(
320 role="assistant",
321 content=content,
322 tool_calls=tool_calls or None,
323 )
324 finish_reason = (
325 "tool_calls" if response.tool_calls else response.finish_reason
326 )
327 return Ok(
328 OpenAIChatResponse(
329 id=response.id or f"chatcmpl-{new_uuid()}",
330 model=context.resolve_model(response.model),
331 created=response.created or 0,
332 choices=[
333 OpenAIChatChoice(
334 index=0,
335 message=message,
336 finish_reason=finish_reason,
337 )
338 ],
339 usage=self._usage_to_wire(response.usage),
340 system_fingerprint=(system_fingerprint),
341 passthrough=passthrough,
342 )
343 )
344 except (RelayError, ValueError, TypeError, KeyError) as exc:
345 return Err(translate(exc, detail="ir_to_response"))
346
347 def stream_to_delta(
348 self, event: Any, *, state: StreamState
349 ) -> Result[tuple[StreamDelta, ...], RelayError]:
350 """Stream conversion is deferred to the shared stream lifecycle task."""
351 return Err(
352 unsupported_feature("openai_chat stream conversion is not implemented yet")
353 )
354
355 def delta_to_stream(
356 self, delta: StreamDelta, *, state: StreamState
357 ) -> Result[tuple[Any, ...], RelayError]:
358 """Stream conversion is deferred to the shared stream lifecycle task."""
359 return Err(
360 unsupported_feature("openai_chat stream conversion is not implemented yet")
361 )