Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-relay/src/lexigram/ai/relay/mappers/openai_responses/request.py: 17%

230 statements  

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

1"""Request-direction conversion for the OpenAI Responses mapper.""" 

2 

3from __future__ import annotations 

4 

5from typing import TYPE_CHECKING, Any, cast 

6 

7from lexigram.ai.relay.context import ConversionContext 

8from lexigram.ai.relay.errors import translate, unsupported_format 

9from lexigram.ai.relay.mappers.base import record_loss 

10from lexigram.ai.relay.mappers.openai_responses.utils import ( 

11 _TARGET, 

12 _parse_arguments, 

13) 

14from lexigram.contracts.ai.agents import ToolDefinition 

15from lexigram.contracts.ai.exceptions import RelayError 

16from lexigram.contracts.ai.llm import ChatMessage, FunctionCall, ToolCall 

17from lexigram.contracts.ai.multimodal import ImageUrlPart, TextPart 

18from lexigram.contracts.ai.relay.dto import ResponsesItem, ResponsesRequest 

19from lexigram.contracts.ai.relay.ir import RelayRequest 

20from lexigram.contracts.ai.thinking import ThinkingConfig 

21from lexigram.contracts.core.result import Err, Ok, Result 

22 

23if TYPE_CHECKING: 

24 from lexigram.ai.relay.mappers.openai_responses import OpenAIResponsesMapper 

25 

26 

27class RequestMixin: 

28 """Request conversion: wire ``ResponsesRequest`` to IR and back.""" 

29 

30 def request_to_ir( 

31 self: OpenAIResponsesMapper, 

32 payload: Any, 

33 *, 

34 context: ConversionContext, 

35 ) -> Result[RelayRequest, RelayError]: 

36 """Convert a ``ResponsesRequest`` into canonical ``RelayRequest``. 

37 

38 Args: 

39 payload: A wire request DTO. 

40 context: Per-conversion context with loss sink. 

41 

42 Returns: 

43 Ok(request) on success, Err(relay_error) on malformed payload. 

44 """ 

45 if not isinstance(payload, ResponsesRequest): 

46 return Err( 

47 unsupported_format( 

48 f"expected ResponsesRequest, got {type(payload).__name__}" 

49 ) 

50 ) 

51 system_parts: list[str] = [] 

52 if payload.instructions: 

53 system_parts.append(payload.instructions) 

54 messages: list[ChatMessage] = [] 

55 pending_tools: list[ToolCall] = [] 

56 pending_ids: list[str | None] = [] 

57 web_search_calls: list[dict[str, Any]] = [] 

58 

59 def flush_tools() -> None: 

60 """Emit accumulated function calls as one assistant turn.""" 

61 if not pending_tools: 

62 return 

63 metadata: dict[str, Any] = {} 

64 if any(pending_ids): 

65 metadata["function_call_item_ids"] = list(pending_ids) 

66 messages.append( 

67 ChatMessage( 

68 role="assistant", 

69 content="", 

70 tool_calls=list(pending_tools), 

71 metadata=metadata or None, 

72 ) 

73 ) 

74 pending_tools.clear() 

75 pending_ids.clear() 

76 

77 if isinstance(payload.input, str): 

78 messages.append(ChatMessage(role="user", content=payload.input)) 

79 else: 

80 for index, wire_item in enumerate(payload.input): 

81 item_type = wire_item.type 

82 if item_type == "message": 

83 flush_tools() 

84 if wire_item.role == "system": 

85 system_parts.append(self._message_text_to_ir(wire_item)) 

86 if index > 0: 

87 record_loss( 

88 context, 

89 field="system_message", 

90 target=_TARGET, 

91 reason="system_message_reordered", 

92 ) 

93 continue 

94 messages.append(self._message_from_item(wire_item, context)) 

95 elif item_type == "function_call": 

96 pending_tools.append(self._tool_from_item(wire_item)) 

97 pending_ids.append(wire_item.id or wire_item.call_id) 

98 elif item_type == "function_call_output": 

99 flush_tools() 

100 messages.append(self._tool_result_from_item(wire_item)) 

101 elif item_type == "reasoning": 

102 flush_tools() 

103 messages.append(self._reasoning_from_item(wire_item)) 

104 elif item_type == "web_search_call": 

105 web_search_calls.append(wire_item.to_dict()) 

106 record_loss( 

107 context, 

108 field=f"input[{index}]", 

109 target=_TARGET, 

110 reason="unsupported_item_preserved", 

111 severity="info", 

112 ) 

113 else: 

114 record_loss( 

115 context, 

116 field=f"input[{index}]", 

117 target=_TARGET, 

118 reason="unknown_item_dropped", 

119 ) 

120 flush_tools() 

121 metadata: dict[str, Any] = {} 

122 if payload.include is not None: 

123 metadata["include"] = list(payload.include) 

124 if payload.max_output_tokens is not None: 

125 metadata["max_tokens_kind"] = "max_completion_tokens" 

126 if web_search_calls: 

127 metadata["input_web_search_calls"] = web_search_calls 

128 reasoning = payload.reasoning 

129 if reasoning is not None: 

130 metadata["reasoning"] = reasoning 

131 if payload.text is not None: 

132 metadata["text"] = payload.text 

133 if payload.service_tier is not None: 

134 metadata["service_tier"] = payload.service_tier 

135 thinking: ThinkingConfig | None = None 

136 if isinstance(reasoning, dict): 

137 thinking = ThinkingConfig(effort=reasoning.get("effort")) 

138 return Ok( 

139 RelayRequest( 

140 model=context.normalize_model(payload.model), 

141 messages=messages, 

142 system="\n".join(system_parts) if system_parts else None, 

143 tools=self._tools_to_ir(payload.tools, context), 

144 temperature=payload.temperature, 

145 max_tokens=payload.max_output_tokens, 

146 stream=payload.stream, 

147 include_usage=bool(payload.include and "usage" in payload.include), 

148 parallel_tool_calls=payload.parallel_tool_calls, 

149 thinking=thinking, 

150 response_format=self._text_to_response_format(payload.text), 

151 metadata=metadata, 

152 passthrough=dict(payload.passthrough), 

153 ) 

154 ) 

155 

156 def ir_to_request( 

157 self: OpenAIResponsesMapper, 

158 request: RelayRequest, 

159 *, 

160 context: ConversionContext, 

161 ) -> Result[Any, RelayError]: 

162 """Convert canonical ``RelayRequest`` into a ``ResponsesRequest``. 

163 

164 Args: 

165 request: Canonical request IR. 

166 context: Per-conversion context with loss sink. 

167 

168 Returns: 

169 Ok(request) on success, Err(relay_error) on failure. 

170 """ 

171 try: 

172 items: list[ResponsesItem] = [] 

173 instructions = request.system 

174 system_parts = [request.system] if request.system else [] 

175 for message in request.messages: 

176 if message.role == "system": 

177 system_parts.append(self._system_text(message)) 

178 record_loss( 

179 context, 

180 field="system_message", 

181 target=_TARGET, 

182 reason="system_message_reordered", 

183 ) 

184 continue 

185 if message.role == "tool": 

186 items.append(self._tool_result_to_item(message)) 

187 continue 

188 if message.tool_calls: 

189 items.extend(self._tool_calls_to_items(message, context)) 

190 continue 

191 items.append(self._message_to_item(message, context)) 

192 items.extend(self._web_search_items(request)) 

193 if system_parts: 

194 instructions = "\n".join(system_parts) 

195 handled_metadata = { 

196 "include", 

197 "reasoning", 

198 "text", 

199 "service_tier", 

200 "input_web_search_calls", 

201 "generation_config", 

202 "safety_settings", 

203 "tool_config", 

204 } 

205 return Ok( 

206 ResponsesRequest( 

207 model=context.resolve_model(request.model), 

208 input=items, 

209 instructions=instructions, 

210 tools=( 

211 [self._tool_from_ir(tool) for tool in request.tools] 

212 if request.tools 

213 else None 

214 ), 

215 temperature=request.temperature, 

216 max_output_tokens=request.max_tokens, 

217 stream=request.stream, 

218 include=self._include_from_ir(request), 

219 parallel_tool_calls=request.parallel_tool_calls, 

220 reasoning=self._reasoning_from_ir(request, context), 

221 text=self._text_from_ir(request), 

222 service_tier=request.metadata.get("service_tier"), 

223 tool_choice=request.tool_choice, 

224 passthrough={ 

225 **request.passthrough, 

226 **{ 

227 key: value 

228 for key, value in request.metadata.items() 

229 if key not in handled_metadata 

230 }, 

231 }, 

232 ) 

233 ) 

234 except (RelayError, ValueError, TypeError, KeyError) as exc: 

235 return Err(translate(exc, detail="ir_to_request")) 

236 

237 @staticmethod 

238 def _message_text_to_ir(wire_item: ResponsesItem) -> str: 

239 """Extract text from a wire message item.""" 

240 return "".join( 

241 str(part.get("text", "")) 

242 for part in wire_item.content or [] 

243 if isinstance(part, dict) and part.get("type") == "input_text" 

244 ) 

245 

246 @staticmethod 

247 def _input_parts_to_ir( 

248 content: list[dict[str, Any]] | None, context: ConversionContext 

249 ) -> tuple[list[Any], list[dict[str, Any]]]: 

250 """Convert wire content parts into canonical parts and files.""" 

251 converted: list[Any] = [] 

252 files: list[dict[str, Any]] = [] 

253 for part in content or []: 

254 if not isinstance(part, dict): 

255 converted.append(TextPart(text=str(part))) 

256 continue 

257 part_type = part.get("type") 

258 if part_type == "input_text": 

259 converted.append(TextPart(text=str(part.get("text", "")))) 

260 elif part_type == "input_image": 

261 image = part.get("image_url") 

262 if isinstance(image, dict): 

263 converted.append( 

264 ImageUrlPart( 

265 url=str(image.get("url", "")), 

266 detail=cast("Any", image.get("detail", "auto") or "auto"), 

267 ) 

268 ) 

269 else: 

270 converted.append( 

271 ImageUrlPart( 

272 url=str(image or ""), 

273 detail=cast("Any", part.get("detail", "auto") or "auto"), 

274 ) 

275 ) 

276 elif part_type == "input_file": 

277 files.append(part) 

278 record_loss( 

279 context, 

280 field="content", 

281 target=_TARGET, 

282 reason="unrepresentable_part_preserved", 

283 severity="info", 

284 ) 

285 else: 

286 record_loss( 

287 context, 

288 field=part_type or "part", 

289 target=_TARGET, 

290 reason="unknown_part_type", 

291 ) 

292 return converted, files 

293 

294 def _message_from_item( 

295 self: OpenAIResponsesMapper, 

296 wire_item: ResponsesItem, 

297 context: ConversionContext, 

298 ) -> ChatMessage: 

299 """Convert a wire message item into a canonical message.""" 

300 wire_content = wire_item.content 

301 if isinstance(wire_content, str): 

302 wire_content = [{"type": "input_text", "text": wire_content}] 

303 parts, files = self._input_parts_to_ir(wire_content, context) 

304 metadata: dict[str, Any] = {} 

305 if wire_item.id: 

306 metadata["item_id"] = wire_item.id 

307 if files: 

308 metadata["input_files"] = files 

309 content: str | list[Any] 

310 if len(parts) == 1 and isinstance(parts[0], TextPart): 

311 content = parts[0].text 

312 elif parts: 

313 content = parts 

314 else: 

315 content = "" 

316 return ChatMessage( 

317 role=wire_item.role or "user", 

318 content=content, 

319 metadata=metadata or None, 

320 ) 

321 

322 @staticmethod 

323 def _tool_from_item(wire_item: ResponsesItem) -> ToolCall: 

324 """Convert a wire function_call item into a canonical tool call.""" 

325 return ToolCall( 

326 id=wire_item.call_id or wire_item.id or "", 

327 type="function", 

328 function=FunctionCall( 

329 name=wire_item.name or "", 

330 arguments=_parse_arguments(wire_item.arguments or ""), 

331 ), 

332 ) 

333 

334 @staticmethod 

335 def _tool_result_from_item(wire_item: ResponsesItem) -> ChatMessage: 

336 """Convert a wire function_call_output item into a tool message.""" 

337 metadata: dict[str, Any] = {} 

338 if wire_item.id: 

339 metadata["item_id"] = wire_item.id 

340 return ChatMessage( 

341 role="tool", 

342 content=wire_item.output or "", 

343 tool_call_id=wire_item.call_id, 

344 metadata=metadata or None, 

345 ) 

346 

347 @staticmethod 

348 def _reasoning_from_item(wire_item: ResponsesItem) -> ChatMessage: 

349 """Convert a wire reasoning item into an assistant message.""" 

350 metadata: dict[str, Any] = {} 

351 if wire_item.id: 

352 metadata["item_id"] = wire_item.id 

353 return ChatMessage( 

354 role="assistant", 

355 content="", 

356 thinking_blocks=list(wire_item.summary or []), 

357 metadata=metadata or None, 

358 ) 

359 

360 @staticmethod 

361 def _tools_to_ir( 

362 tools: list[dict[str, Any]] | None, context: ConversionContext 

363 ) -> list[ToolDefinition]: 

364 """Convert wire tool dicts into canonical tool definitions.""" 

365 definitions: list[ToolDefinition] = [] 

366 if not tools: 

367 return definitions 

368 for index, tool in enumerate(tools): 

369 if not isinstance(tool, dict): 

370 record_loss( 

371 context, 

372 field=f"tools[{index}]", 

373 target=_TARGET, 

374 reason="non_dict_tool_dropped", 

375 ) 

376 continue 

377 if tool.get("type", "function") != "function": 

378 record_loss( 

379 context, 

380 field=f"tools[{index}]", 

381 target=_TARGET, 

382 reason="non_function_tool_dropped", 

383 ) 

384 continue 

385 if isinstance(tool.get("function"), dict): 

386 function = tool["function"] 

387 else: 

388 function = tool 

389 parameters = function.get("parameters", {}) 

390 definitions.append( 

391 ToolDefinition( 

392 name=str(function.get("name", "")), 

393 description=str(function.get("description", "")), 

394 parameters=parameters if isinstance(parameters, dict) else {}, 

395 ) 

396 ) 

397 return definitions 

398 

399 @staticmethod 

400 def _tool_from_ir(tool: ToolDefinition) -> dict[str, Any]: 

401 """Serialize a canonical tool definition as a wire tool dict.""" 

402 return { 

403 "type": "function", 

404 "name": tool.name, 

405 "description": tool.description, 

406 "parameters": tool.parameters, 

407 } 

408 

409 @staticmethod 

410 def _text_to_response_format( 

411 text: dict[str, Any] | None, 

412 ) -> dict[str, Any] | None: 

413 """Derive a canonical response format from the wire text config.""" 

414 if not isinstance(text, dict): 

415 return None 

416 fmt = text.get("format") 

417 if not isinstance(fmt, dict): 

418 return None 

419 fmt_type = fmt.get("type") 

420 if fmt_type == "json_object": 

421 return {"type": "json_object"} 

422 if fmt_type == "json_schema": 

423 result: dict[str, Any] = {"type": "json_schema"} 

424 for key in ("schema", "name", "strict"): 

425 if key in fmt: 

426 result[key] = fmt[key] 

427 return result 

428 return None 

429 

430 @classmethod 

431 def _text_from_ir(cls, request: RelayRequest) -> dict[str, Any] | None: 

432 """Rebuild the wire text config from canonical response format.""" 

433 raw = request.metadata.get("text") 

434 text: dict[str, Any] = dict(raw) if isinstance(raw, dict) else {} 

435 response_format = request.response_format 

436 if response_format is None: 

437 return text or None 

438 fmt_type = response_format.get("type") 

439 if fmt_type == "json_object": 

440 text["format"] = {"type": "json_object"} 

441 elif fmt_type == "json_schema": 

442 fmt: dict[str, Any] = {"type": "json_schema"} 

443 for key in ("schema", "name", "strict"): 

444 if key in response_format: 

445 fmt[key] = response_format[key] 

446 text["format"] = fmt 

447 else: 

448 text["format"] = {"type": "text"} 

449 return text or None 

450 

451 @staticmethod 

452 def _include_from_ir(request: RelayRequest) -> list[str] | None: 

453 """Rebuild the wire include list from canonical stream settings.""" 

454 raw = request.metadata.get("include") 

455 include = list(raw) if isinstance(raw, list) else [] 

456 if request.include_usage and "usage" not in include: 

457 include.append("usage") 

458 return include or None 

459 

460 def _reasoning_from_ir( 

461 self: OpenAIResponsesMapper, 

462 request: RelayRequest, 

463 context: ConversionContext, 

464 ) -> dict[str, Any] | None: 

465 """Rebuild the wire reasoning config from canonical thinking.""" 

466 thinking = request.thinking 

467 if thinking is not None: 

468 if thinking.effort is not None: 

469 return {"effort": thinking.effort} 

470 record_loss( 

471 context, 

472 field="thinking", 

473 target=_TARGET, 

474 reason="effort_only_supported", 

475 ) 

476 raw = request.metadata.get("reasoning") 

477 if isinstance(raw, dict): 

478 return dict(raw) 

479 return None