Coverage for src/lilbee/server/chat_completions_api/translate.py: 100%

170 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-09-28 17:20 +0000

1"""Translation between OpenAI chat-completions models and the canonical types.""" 

2 

3from __future__ import annotations 

4 

5import json 

6import time 

7from collections.abc import AsyncIterator 

8from typing import Literal, assert_never 

9 

10from lilbee.core.config.enums import ReasoningMode 

11from lilbee.retrieval.reasoning import StreamToken, TagParser, split_reasoning 

12from lilbee.server.chat_completions_api.models import ( 

13 CompletionsImageContent, 

14 CompletionsMessage, 

15 CompletionsNamedToolChoice, 

16 CompletionsRequest, 

17 CompletionsResponse, 

18 CompletionsResponseChoice, 

19 CompletionsResponseMessage, 

20 CompletionsResponseToolCall, 

21 CompletionsResponseToolCallFunction, 

22 CompletionsStreamChoice, 

23 CompletionsStreamChunk, 

24 CompletionsStreamDelta, 

25 CompletionsStreamToolCall, 

26 CompletionsStreamToolCallFunction, 

27 CompletionsTextContent, 

28 CompletionsTool, 

29 CompletionsUsage, 

30 FinishReason, 

31 PromptTokensDetails, 

32 ToolChoiceMode, 

33) 

34from lilbee.server.chat_dispatch.canonical import ( 

35 CanonicalChatRequest, 

36 CanonicalMessage, 

37 CanonicalResponse, 

38 CanonicalStreamEvent, 

39 CanonicalTool, 

40 CanonicalToolChoice, 

41 CanonicalUsage, 

42 ContentBlock, 

43 ContentBlockDelta, 

44 ContentBlockStart, 

45 ContentBlockStop, 

46 MessageDelta, 

47 MessageStart, 

48 MessageStop, 

49 StopReason, 

50 TextBlock, 

51 TextDelta, 

52 ToolResultBlock, 

53 ToolUseBlock, 

54 ToolUseDelta, 

55) 

56from lilbee.server.chat_dispatch.tool_args import parse_tool_arguments 

57 

58_TOOL_CHOICE_MODES: dict[ToolChoiceMode, Literal["auto", "any", "none"]] = { 

59 ToolChoiceMode.AUTO: "auto", 

60 ToolChoiceMode.NONE: "none", 

61 ToolChoiceMode.REQUIRED: "any", 

62} 

63 

64_STOP_REASON_TO_FINISH: dict[StopReason, FinishReason] = { 

65 StopReason.END_TURN: FinishReason.STOP, 

66 StopReason.MAX_TOKENS: FinishReason.LENGTH, 

67 StopReason.TOOL_USE: FinishReason.TOOL_CALLS, 

68} 

69 

70_TOOL_CALL_ID_REQUIRED = "A tool message requires tool_call_id naming the call it answers." 

71_IMAGE_CONTENT_UNSUPPORTED = ( 

72 "Image content is not supported by /v1/chat/completions yet. Send a text-only request." 

73) 

74 

75 

76def completions_to_canonical_request( 

77 request: CompletionsRequest, *, mode: ReasoningMode = ReasoningMode.SEPARATE 

78) -> CanonicalChatRequest: 

79 """Translate a validated ``CompletionsRequest`` to the canonical request.""" 

80 system_parts: list[str] = [] 

81 messages: list[CanonicalMessage] = [] 

82 for msg in request.messages: 

83 if msg.role == "system": 

84 system_parts.append(_system_text(msg)) 

85 continue 

86 messages.append(_message_from_request(msg)) 

87 

88 return CanonicalChatRequest( 

89 model=request.model, 

90 messages=messages, 

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

92 tools=_tools_from_request(request.tools), 

93 tool_choice=_tool_choice_from_request(request.tool_choice), 

94 temperature=request.temperature, 

95 top_p=request.top_p, 

96 top_k=request.top_k, 

97 max_tokens=request.max_tokens, 

98 seed=request.seed, 

99 frequency_penalty=request.frequency_penalty, 

100 presence_penalty=request.presence_penalty, 

101 stop=_stop_from_request(request.stop), 

102 stream=request.stream, 

103 # OFF asks the template to skip thinking; the other modes only change 

104 # presentation, so the template default stands. 

105 think=False if mode is ReasoningMode.OFF else None, 

106 ) 

107 

108 

109def canonical_to_completions_response( 

110 resp: CanonicalResponse, *, response_id: str, mode: ReasoningMode = ReasoningMode.SEPARATE 

111) -> CompletionsResponse: 

112 """Translate a canonical chat response to the OpenAI ``chat.completion`` model.""" 

113 text_parts = [b.text for b in resp.content if isinstance(b, TextBlock)] 

114 tool_calls = [_response_tool_call(b) for b in resp.content if isinstance(b, ToolUseBlock)] 

115 # lilbee carries a reasoning model's thinking inline as <think>...</think>; the 

116 # OpenAI surface reports it in its own field so agents render a clean answer. 

117 # INLINE streams the thinking as ordinary content for clients that never 

118 # render reasoning_content -- with the tag markers stripped, because a 

119 # client that ignores reasoning_content renders raw <think> text literally. 

120 # OFF keeps the split, because a template that ignores enable_thinking 

121 # still thinks. 

122 if mode is ReasoningMode.INLINE: 

123 reasoning, answer = split_reasoning("".join(text_parts)) 

124 if reasoning: 

125 answer = f"{reasoning}\n\n{answer}" if answer else reasoning 

126 reasoning = "" 

127 else: 

128 reasoning, answer = split_reasoning("".join(text_parts)) 

129 content: str | None = answer if answer or not tool_calls else None 

130 

131 total = resp.usage.input_tokens + resp.usage.output_tokens 

132 return CompletionsResponse( 

133 id=response_id, 

134 created=int(time.time()), 

135 model=resp.model, 

136 choices=[ 

137 CompletionsResponseChoice( 

138 index=0, 

139 message=CompletionsResponseMessage( 

140 content=content, 

141 reasoning_content=reasoning or None, 

142 tool_calls=tool_calls or None, 

143 ), 

144 finish_reason=_STOP_REASON_TO_FINISH[resp.stop_reason], 

145 ) 

146 ], 

147 usage=CompletionsUsage( 

148 prompt_tokens=resp.usage.input_tokens, 

149 completion_tokens=resp.usage.output_tokens, 

150 total_tokens=total, 

151 prompt_tokens_details=PromptTokensDetails(cached_tokens=resp.usage.cached_input_tokens), 

152 ), 

153 ) 

154 

155 

156class _StreamMapper: 

157 """Per-stream state for the canonical-to-OpenAI chunk converter.""" 

158 

159 def __init__(self, *, mode: ReasoningMode = ReasoningMode.SEPARATE) -> None: 

160 self._role_emitted = False 

161 self._tool_index_for_block: dict[int, int] = {} 

162 self._next_tool_index = 0 

163 # INLINE routes thinking into content instead of reasoning_content, with 

164 # the <think> markers stripped: a client that ignores reasoning_content 

165 # renders raw tags literally, which is exactly what INLINE exists to 

166 # avoid. The same stateful parse handles a tag split across deltas. 

167 self._inline = mode is ReasoningMode.INLINE 

168 # Splits lilbee's inline <think> text into its own delta field. Stateful 

169 # because a tag can arrive split across deltas. 

170 self._reasoning = TagParser(show=True) 

171 

172 def block_start(self, event: ContentBlockStart) -> CompletionsStreamDelta | None: 

173 # The first delta must carry role:assistant for OpenAI-SDK accumulation, 

174 # whether the response opens with text or a tool call. 

175 role: Literal["assistant"] | None = None 

176 if not self._role_emitted: 

177 self._role_emitted = True 

178 role = "assistant" 

179 if isinstance(event.block, TextBlock): 

180 return CompletionsStreamDelta(role=role) if role is not None else None 

181 if isinstance(event.block, ToolUseBlock): 

182 tool_index = self._next_tool_index 

183 self._tool_index_for_block[event.index] = tool_index 

184 self._next_tool_index += 1 

185 return CompletionsStreamDelta( 

186 role=role, 

187 tool_calls=[_tool_call_open(tool_index, event.block.id, event.block.name)], 

188 ) 

189 return None 

190 

191 def block_delta(self, event: ContentBlockDelta) -> CompletionsStreamDelta | None: 

192 if isinstance(event.delta, TextDelta): 

193 return self._text_delta(self._reasoning.feed(event.delta.text)) 

194 if isinstance(event.delta, ToolUseDelta): 

195 # A delta for a block we never saw start is a provider quirk, not a 

196 # server fault; a bare subscript turned it into a KeyError that 

197 # surfaced as a stream-level internal error. 

198 tool_index = self._tool_index_for_block.get(event.index) 

199 if tool_index is None: 

200 return None 

201 return CompletionsStreamDelta( 

202 tool_calls=[_tool_call_args(tool_index, event.delta.partial_json)], 

203 ) 

204 return None 

205 

206 def block_stop(self) -> CompletionsStreamDelta | None: 

207 """Emit whatever the reasoning splitter still holds (a partial or unclosed tag).""" 

208 remaining = self._reasoning.flush() 

209 return self._text_delta([remaining] if remaining else []) 

210 

211 def _text_delta(self, tokens: list[StreamToken]) -> CompletionsStreamDelta | None: 

212 """One delta carrying the parsed text: split fields, or tag-free content.""" 

213 if self._inline: 

214 # Arrival order preserved; the parser already dropped the markers. 

215 text = "".join(t.content for t in tokens) 

216 return CompletionsStreamDelta(content=text) if text else None 

217 reasoning = "".join(t.content for t in tokens if t.is_reasoning) 

218 answer = "".join(t.content for t in tokens if not t.is_reasoning) 

219 if not reasoning and not answer: 

220 return None 

221 return CompletionsStreamDelta( 

222 content=answer or None, 

223 reasoning_content=reasoning or None, 

224 ) 

225 

226 

227def _tool_call_open(index: int, call_id: str, name: str) -> CompletionsStreamToolCall: 

228 return CompletionsStreamToolCall( 

229 index=index, 

230 id=call_id, 

231 type="function", 

232 function=CompletionsStreamToolCallFunction(name=name, arguments=""), 

233 ) 

234 

235 

236def _tool_call_args(index: int, partial_json: str) -> CompletionsStreamToolCall: 

237 return CompletionsStreamToolCall( 

238 index=index, 

239 function=CompletionsStreamToolCallFunction(arguments=partial_json), 

240 ) 

241 

242 

243def _finish_reason_for(event: MessageDelta) -> FinishReason: 

244 if event.stop_reason is None: 

245 return FinishReason.STOP 

246 return _STOP_REASON_TO_FINISH[event.stop_reason] 

247 

248 

249async def canonical_stream_to_completions_chunks( 

250 events: AsyncIterator[CanonicalStreamEvent], 

251 *, 

252 model: str, 

253 response_id: str, 

254 include_usage: bool = False, 

255 mode: ReasoningMode = ReasoningMode.SEPARATE, 

256) -> AsyncIterator[CompletionsStreamChunk]: 

257 """Turn canonical stream events into ``CompletionsStreamChunk`` instances. 

258 

259 The trailing usage-only chunk is emitted only when *include_usage* is set, 

260 matching OpenAI's ``stream_options.include_usage`` contract. 

261 """ 

262 mapper = _StreamMapper(mode=mode) 

263 # One timestamp for the whole completion. OpenAI holds created constant 

264 # across a stream; recomputing it per chunk reported several creation times 

265 # for one completion, and clients order or dedupe on it. 

266 created = int(time.time()) 

267 async for event in events: 

268 for chunk in _chunks_for_event( 

269 event, 

270 mapper, 

271 model=model, 

272 response_id=response_id, 

273 include_usage=include_usage, 

274 created=created, 

275 ): 

276 yield chunk 

277 

278 

279def _chunks_for_event( 

280 event: CanonicalStreamEvent, 

281 mapper: _StreamMapper, 

282 *, 

283 model: str, 

284 response_id: str, 

285 include_usage: bool, 

286 created: int, 

287) -> list[CompletionsStreamChunk]: 

288 """The OpenAI chunks one canonical event translates to; empty for a no-op event.""" 

289 if isinstance(event, ContentBlockStart): 

290 return _maybe_chunk(model, response_id, mapper.block_start(event), created=created) 

291 if isinstance(event, ContentBlockDelta): 

292 return _maybe_chunk(model, response_id, mapper.block_delta(event), created=created) 

293 if isinstance(event, ContentBlockStop): 

294 # Closing a text block flushes any text the reasoning splitter still 

295 # buffers (an unclosed <think>, or a tag that never completed). 

296 return _maybe_chunk(model, response_id, mapper.block_stop(), created=created) 

297 if isinstance(event, MessageDelta): 

298 return _message_delta_chunks( 

299 event, 

300 model=model, 

301 response_id=response_id, 

302 include_usage=include_usage, 

303 created=created, 

304 ) 

305 if isinstance(event, MessageStart | MessageStop): 

306 # OpenAI's wire format has no equivalent: MessageStart carries metadata 

307 # already encoded in the chunk header, and MessageStop is replaced by the 

308 # final chunk's finish_reason. 

309 return [] 

310 # Unreachable: the branches above exhaust CanonicalStreamEvent. This makes a new 

311 # event type a type error rather than a silently dropped frame. 

312 assert_never(event) # pragma: no cover 

313 

314 

315def _message_delta_chunks( 

316 event: MessageDelta, *, model: str, response_id: str, include_usage: bool, created: int 

317) -> list[CompletionsStreamChunk]: 

318 """The finish chunk, plus the usage-only chunk when the client asked for it.""" 

319 chunks = [ 

320 _chunk( 

321 model, 

322 response_id, 

323 CompletionsStreamDelta(), 

324 finish_reason=_finish_reason_for(event), 

325 created=created, 

326 ) 

327 ] 

328 if include_usage: 

329 # OpenAI's contract sends the usage-only chunk unconditionally when 

330 # include_usage is set; a client blocking on it must not hang because the 

331 # provider streamed no usage frame. 

332 usage = event.usage or CanonicalUsage(input_tokens=0, output_tokens=0) 

333 chunks.append(_usage_chunk(model, response_id, usage, created=created)) 

334 return chunks 

335 

336 

337def _maybe_chunk( 

338 model: str, response_id: str, delta: CompletionsStreamDelta | None, *, created: int 

339) -> list[CompletionsStreamChunk]: 

340 """Wrap a delta in a chunk, or nothing when the event produced no delta.""" 

341 return [] if delta is None else [_chunk(model, response_id, delta, created=created)] 

342 

343 

344def _chunk( 

345 model: str, 

346 response_id: str, 

347 delta: CompletionsStreamDelta, 

348 *, 

349 finish_reason: FinishReason | None = None, 

350 created: int, 

351) -> CompletionsStreamChunk: 

352 return CompletionsStreamChunk( 

353 id=response_id, 

354 created=created, 

355 model=model, 

356 choices=[CompletionsStreamChoice(index=0, delta=delta, finish_reason=finish_reason)], 

357 ) 

358 

359 

360def _usage_chunk( 

361 model: str, response_id: str, usage: CanonicalUsage, *, created: int 

362) -> CompletionsStreamChunk: 

363 """Final include_usage chunk: empty choices, populated usage totals.""" 

364 total = usage.input_tokens + usage.output_tokens 

365 return CompletionsStreamChunk( 

366 id=response_id, 

367 created=created, 

368 model=model, 

369 choices=[], 

370 usage=CompletionsUsage( 

371 prompt_tokens=usage.input_tokens, 

372 completion_tokens=usage.output_tokens, 

373 total_tokens=total, 

374 prompt_tokens_details=PromptTokensDetails(cached_tokens=usage.cached_input_tokens), 

375 ), 

376 ) 

377 

378 

379def _system_text(msg: CompletionsMessage) -> str: 

380 """Flatten a system message to text, rejecting image parts. 

381 

382 The same image part is a 400 in a user message, so silently dropping it 

383 here answered as though the request had been honoured. 

384 """ 

385 if isinstance(msg.content, str): 

386 return msg.content 

387 if isinstance(msg.content, list): 

388 if any(isinstance(part, CompletionsImageContent) for part in msg.content): 

389 raise ValueError(_IMAGE_CONTENT_UNSUPPORTED) 

390 return "".join( 

391 part.text for part in msg.content if isinstance(part, CompletionsTextContent) 

392 ) 

393 return "" 

394 

395 

396def _message_from_request(msg: CompletionsMessage) -> CanonicalMessage: 

397 role = msg.role 

398 if role == "system": 

399 raise ValueError("system messages should be extracted by the caller") 

400 if role == "tool": 

401 if not msg.tool_call_id: 

402 # Substituting "" produced a tool result no tool call can pair with, 

403 # which the provider sees as a result for a call it never made. 

404 raise ValueError(_TOOL_CALL_ID_REQUIRED) 

405 return CanonicalMessage( 

406 role="tool", 

407 content=[ 

408 ToolResultBlock( 

409 tool_use_id=msg.tool_call_id, 

410 content=_tool_result_content(msg.content), 

411 ) 

412 ], 

413 ) 

414 

415 blocks: list[ContentBlock] = list(_content_blocks(msg.content)) 

416 for call in msg.tool_calls or []: 

417 blocks.append( 

418 ToolUseBlock( 

419 id=call.id, 

420 name=call.function.name, 

421 input=parse_tool_arguments(call.function.arguments), 

422 ) 

423 ) 

424 return CanonicalMessage(role=role, content=blocks) 

425 

426 

427def _content_blocks(content: str | list | None) -> list[ContentBlock]: 

428 if content is None or content == "": 

429 return [] 

430 if isinstance(content, str): 

431 return [TextBlock(text=content)] 

432 blocks: list[ContentBlock] = [] 

433 for part in content: 

434 if isinstance(part, CompletionsTextContent): 

435 blocks.append(TextBlock(text=part.text)) 

436 elif isinstance(part, CompletionsImageContent): 

437 raise ValueError(_IMAGE_CONTENT_UNSUPPORTED) 

438 return blocks 

439 

440 

441def _tool_result_content(content: str | list | None) -> list[ContentBlock]: 

442 if isinstance(content, str): 

443 return [TextBlock(text=content)] 

444 if isinstance(content, list): 

445 return _content_blocks(content) 

446 return [TextBlock(text="" if content is None else str(content))] 

447 

448 

449def _response_tool_call(block: ToolUseBlock) -> CompletionsResponseToolCall: 

450 return CompletionsResponseToolCall( 

451 id=block.id, 

452 function=CompletionsResponseToolCallFunction( 

453 name=block.name, arguments=json.dumps(block.input) 

454 ), 

455 ) 

456 

457 

458def _tools_from_request(tools: list[CompletionsTool] | None) -> list[CanonicalTool] | None: 

459 if not tools: 

460 return None 

461 return [ 

462 CanonicalTool( 

463 name=tool.function.name, 

464 description=tool.function.description or "", 

465 input_schema=tool.function.parameters, 

466 ) 

467 for tool in tools 

468 ] 

469 

470 

471def _tool_choice_from_request( 

472 choice: ToolChoiceMode | CompletionsNamedToolChoice | None, 

473) -> CanonicalToolChoice | None: 

474 if choice is None: 

475 return None 

476 if isinstance(choice, ToolChoiceMode): 

477 return CanonicalToolChoice(mode=_TOOL_CHOICE_MODES[choice]) 

478 return CanonicalToolChoice(mode="tool", tool_name=choice.function.name) 

479 

480 

481def _stop_from_request(stop: str | list[str] | None) -> list[str] | None: 

482 if stop is None: 

483 return None 

484 if isinstance(stop, str): 

485 return [stop] 

486 return list(stop)