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
« 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."""
3from __future__ import annotations
5import json
6import time
7from collections.abc import AsyncIterator
8from typing import Literal, assert_never
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
58_TOOL_CHOICE_MODES: dict[ToolChoiceMode, Literal["auto", "any", "none"]] = {
59 ToolChoiceMode.AUTO: "auto",
60 ToolChoiceMode.NONE: "none",
61 ToolChoiceMode.REQUIRED: "any",
62}
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}
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)
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))
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 )
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
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 )
156class _StreamMapper:
157 """Per-stream state for the canonical-to-OpenAI chunk converter."""
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)
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
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
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 [])
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 )
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 )
236def _tool_call_args(index: int, partial_json: str) -> CompletionsStreamToolCall:
237 return CompletionsStreamToolCall(
238 index=index,
239 function=CompletionsStreamToolCallFunction(arguments=partial_json),
240 )
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]
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.
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
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
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
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)]
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 )
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 )
379def _system_text(msg: CompletionsMessage) -> str:
380 """Flatten a system message to text, rejecting image parts.
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 ""
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 )
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)
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
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))]
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 )
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 ]
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)
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)