Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-relay/src/lexigram/ai/relay/mappers/openai_chat/mapper.py: 20%

118 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-25 07:19 +0800

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 )