Coverage for src/lilbee/server/anthropic_api/translate.py: 100%
199 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 Anthropic Messages models and the canonical types."""
3from __future__ import annotations
5import json
6from collections.abc import AsyncIterator
7from dataclasses import replace
8from enum import StrEnum
9from typing import Any, Literal
11from lilbee.core.config.enums import ReasoningMode
12from lilbee.retrieval.reasoning import (
13 PseudoThinkingNormalizer,
14 StreamToken,
15 TagParser,
16 normalize_pseudo_thinking,
17 split_reasoning,
18)
19from lilbee.server.anthropic_api.models import (
20 _THINKING_DISABLED,
21 AnthropicEventType,
22 AnthropicMessage,
23 AnthropicThinking,
24 AnthropicTool,
25 AnthropicToolChoice,
26 AnthropicUsage,
27 ContentBlockParam,
28 CountTokensRequest,
29 ImageBlockParam,
30 MessagesRequest,
31 MessagesResponse,
32 SystemTextBlock,
33 TextBlockParam,
34 ToolResultBlockParam,
35 ToolUseBlockParam,
36 UnknownBlockParam,
37 _PromptBody,
38)
39from lilbee.server.chat_dispatch.canonical import (
40 CanonicalChatRequest,
41 CanonicalMessage,
42 CanonicalResponse,
43 CanonicalStreamEvent,
44 CanonicalTool,
45 CanonicalToolChoice,
46 CanonicalUsage,
47 ContentBlock,
48 ContentBlockDelta,
49 ContentBlockStart,
50 ContentBlockStop,
51 MessageDelta,
52 MessageStart,
53 MessageStop,
54 StopReason,
55 TextBlock,
56 TextDelta,
57 ToolResultBlock,
58 ToolUseBlock,
59)
61_IMAGE_CONTENT_UNSUPPORTED = (
62 "Image content is not supported by /v1/messages yet. Send a text-only request."
63)
64_TOOL_CHOICE_NAME_REQUIRED = 'tool_choice type "tool" requires a name.'
67class _BlockKind(StrEnum):
68 """Kind of the mapper's open output block."""
70 THINKING = "thinking"
71 TEXT = "text"
72 TOOL = "tool"
75_ANTHROPIC_CHOICE_MODES: dict[str, str] = {
76 "auto": "auto",
77 "any": "any",
78 "none": "none",
79}
82def resolve_reasoning_mode(
83 thinking: AnthropicThinking | None, *, default: ReasoningMode
84) -> ReasoningMode:
85 """Pick the reasoning mode for one call from the request and the setting.
87 Thinking is opt-in per request, as on the Anthropic API: a body with no
88 ``thinking`` gets none, whatever the setting presents it as. The setting
89 says how to present thinking a request asked for, and ``off`` refuses it
90 outright, so a request can only tighten.
91 """
92 if default is ReasoningMode.OFF:
93 return ReasoningMode.OFF
94 if thinking is None or thinking.type == _THINKING_DISABLED:
95 return ReasoningMode.OFF
96 return default
99def count_tokens_to_canonical_request(
100 request: CountTokensRequest, *, mode: ReasoningMode = ReasoningMode.SEPARATE
101) -> CanonicalChatRequest:
102 """Translate a validated ``CountTokensRequest`` to the canonical request.
104 The count_tokens contract carries no sampling or streaming fields, so the
105 result is the prompt translation alone.
106 """
107 return _canonical_prompt(request, mode=mode)
110def messages_to_canonical_request(
111 request: MessagesRequest, *, mode: ReasoningMode = ReasoningMode.SEPARATE
112) -> CanonicalChatRequest:
113 """Translate a validated ``MessagesRequest`` to the canonical request."""
114 return replace(
115 _canonical_prompt(request, mode=mode),
116 temperature=request.temperature,
117 top_p=request.top_p,
118 top_k=request.top_k,
119 max_tokens=request.max_tokens,
120 stop=list(request.stop_sequences) if request.stop_sequences else None,
121 stream=request.stream,
122 )
125def _canonical_prompt(request: _PromptBody, *, mode: ReasoningMode) -> CanonicalChatRequest:
126 """Translate the fields that decide the rendered prompt.
128 Both request bodies carry these fields and both routes translate them here,
129 so a field added to the prompt reaches the chat call and the count together.
130 """
131 return CanonicalChatRequest(
132 model=request.model,
133 messages=_canonical_messages(request.messages),
134 system=_system_text(request.system),
135 tools=_tools_from_request(request.tools),
136 tool_choice=_tool_choice_from_request(request.tool_choice),
137 # OFF asks the template to skip thinking; the other modes only change
138 # presentation, so the template default stands.
139 think=False if mode is ReasoningMode.OFF else None,
140 )
143def _canonical_messages(messages: list[AnthropicMessage]) -> list[CanonicalMessage]:
144 """Fan a request's messages out to the canonical message list."""
145 out: list[CanonicalMessage] = []
146 for msg in messages:
147 out.extend(_canonical_messages_for(msg))
148 return out
151def _system_text(system: str | list[SystemTextBlock] | None) -> str | None:
152 if system is None:
153 return None
154 if isinstance(system, str):
155 return system or None
156 joined = "\n\n".join(block.text for block in system)
157 return joined or None
160def _canonical_messages_for(msg: AnthropicMessage) -> list[CanonicalMessage]:
161 """Fan one Anthropic message out to canonical messages.
163 Tool results become their own ``role: "tool"`` messages, emitted before
164 the user's text so the provider sees results adjacent to the calls they
165 answer. Unknown blocks (replayed thinking) are dropped. A mid-conversation
166 ``system`` message becomes a system-reminder user turn -- Anthropic's own
167 documented degradation for models without the operator channel, and it
168 keeps the canonical layer's role set unchanged.
169 """
170 if msg.role == "system":
171 text = _message_text(msg)
172 if not text:
173 return []
174 return [
175 CanonicalMessage.from_string(
176 role="user", text=f"<system-reminder>\n{text}\n</system-reminder>"
177 )
178 ]
179 if isinstance(msg.content, str):
180 if not msg.content:
181 return []
182 return [CanonicalMessage.from_string(role=msg.role, text=msg.content)]
183 return _block_messages(msg.role, msg.content)
186def _block_messages(
187 role: Literal["user", "assistant"], content: list[ContentBlockParam]
188) -> list[CanonicalMessage]:
189 """Canonical messages for a block-form user or assistant message."""
190 tool_messages: list[CanonicalMessage] = []
191 blocks: list[ContentBlock] = []
192 for block in content:
193 if isinstance(block, TextBlockParam):
194 blocks.append(TextBlock(text=block.text))
195 elif isinstance(block, ToolUseBlockParam):
196 blocks.append(ToolUseBlock(id=block.id, name=block.name, input=block.input))
197 elif isinstance(block, ToolResultBlockParam):
198 tool_messages.append(_tool_result_message(block))
199 elif isinstance(block, ImageBlockParam):
200 raise ValueError(_IMAGE_CONTENT_UNSUPPORTED)
201 elif isinstance(block, UnknownBlockParam):
202 continue
204 out = tool_messages
205 if blocks:
206 out = [*tool_messages, CanonicalMessage(role=role, content=blocks)]
207 return out
210def _tool_result_message(block: ToolResultBlockParam) -> CanonicalMessage:
211 return CanonicalMessage(
212 role="tool",
213 content=[
214 ToolResultBlock(
215 tool_use_id=block.tool_use_id,
216 content=_tool_result_content(block),
217 is_error=block.is_error,
218 )
219 ],
220 )
223def _message_text(msg: AnthropicMessage) -> str:
224 """The concatenated text of a message, ignoring non-text blocks."""
225 if isinstance(msg.content, str):
226 return msg.content
227 return "".join(b.text for b in msg.content if isinstance(b, TextBlockParam))
230def _tool_result_content(block: ToolResultBlockParam) -> list[ContentBlock]:
231 if block.content is None:
232 return []
233 if isinstance(block.content, str):
234 return [TextBlock(text=block.content)]
235 parts: list[ContentBlock] = []
236 for part in block.content:
237 if isinstance(part, TextBlockParam):
238 parts.append(TextBlock(text=part.text))
239 elif isinstance(part, ImageBlockParam):
240 raise ValueError(_IMAGE_CONTENT_UNSUPPORTED)
241 # UnknownBlockParam: dropped
242 return parts
245def _tools_from_request(tools: list[AnthropicTool] | None) -> list[CanonicalTool] | None:
246 if not tools:
247 return None
248 return [
249 CanonicalTool(
250 name=tool.name,
251 description=tool.description or "",
252 input_schema=tool.input_schema,
253 )
254 for tool in tools
255 ]
258def _tool_choice_from_request(
259 choice: AnthropicToolChoice | None,
260) -> CanonicalToolChoice | None:
261 if choice is None:
262 return None
263 if choice.type == "tool":
264 if not choice.name:
265 raise ValueError(_TOOL_CHOICE_NAME_REQUIRED)
266 return CanonicalToolChoice(mode="tool", tool_name=choice.name)
267 mode = _ANTHROPIC_CHOICE_MODES[choice.type]
268 return CanonicalToolChoice(mode=mode) # type: ignore[arg-type]
271def canonical_to_messages_response(
272 resp: CanonicalResponse, *, response_id: str, mode: ReasoningMode = ReasoningMode.SEPARATE
273) -> MessagesResponse:
274 """Translate a canonical chat response to the Anthropic message shape.
276 lilbee carries a reasoning model's thinking inline as ``<think>...</think>``;
277 SEPARATE reports it as a leading ``thinking`` block so clients render a
278 clean answer. INLINE folds it into the answer text with the markers
279 stripped, for clients that never render thinking blocks. OFF drops it: the
280 caller asked for no thinking, and a template that ignores the request still
281 thinks, so the block would contradict the answer the caller asked for. OFF
282 also drops a reply-initial pseudo-thinking block a model emits as plain text.
283 """
284 text = "".join(b.text for b in resp.content if isinstance(b, TextBlock))
285 if mode is ReasoningMode.OFF:
286 text = normalize_pseudo_thinking(text)
287 reasoning, answer = split_reasoning(text)
288 if mode is ReasoningMode.INLINE and reasoning:
289 answer = f"{reasoning}\n\n{answer}" if answer else reasoning
290 if mode is not ReasoningMode.SEPARATE:
291 reasoning = ""
292 content: list[dict[str, Any]] = []
293 if reasoning:
294 content.append({"type": "thinking", "thinking": reasoning})
295 tool_uses = [b for b in resp.content if isinstance(b, ToolUseBlock)]
296 if answer or not (reasoning or tool_uses):
297 content.append({"type": "text", "text": answer})
298 content.extend(
299 {"type": "tool_use", "id": b.id, "name": b.name, "input": b.input} for b in tool_uses
300 )
301 return MessagesResponse(
302 id=response_id,
303 model=resp.model,
304 content=content,
305 stop_reason=str(resp.stop_reason),
306 usage=AnthropicUsage(
307 input_tokens=resp.usage.uncached_input_tokens,
308 output_tokens=resp.usage.output_tokens,
309 cache_read_input_tokens=resp.usage.cached_input_tokens,
310 ),
311 )
314class _AnthropicStreamMapper:
315 """Per-stream state for the canonical-to-Anthropic event converter.
317 lilbee streams reasoning inline as ``<think>`` text, and Anthropic's wire
318 format wants thinking and answer text in separate indexed blocks; the
319 mapper re-blocks the stream, closing the open block whenever the token
320 kind (thinking / text / tool_use) changes.
322 INLINE routes reasoning into the text block instead, and OFF drops it: a
323 parser built with ``show=False`` reports reasoning tokens with empty
324 content, which never opens a block or emits a delta. OFF also rewrites a
325 reply-initial pseudo-thinking tag to the ``<think>`` tags before parsing,
326 so a planning block a model emits as plain text is dropped too.
327 """
329 def __init__(self, *, mode: ReasoningMode = ReasoningMode.SEPARATE) -> None:
330 self._reasoning = TagParser(show=mode is not ReasoningMode.OFF)
331 self._pseudo: PseudoThinkingNormalizer | None = (
332 PseudoThinkingNormalizer() if mode is ReasoningMode.OFF else None
333 )
334 self._inline = mode is ReasoningMode.INLINE
335 self._next_index = 0
336 self._open: _BlockKind | None = None
338 def _close_open(self) -> list[tuple[AnthropicEventType, dict[str, Any]]]:
339 if self._open is None:
340 return []
341 index = self._next_index - 1
342 self._open = None
343 return [
344 (AnthropicEventType.CONTENT_BLOCK_STOP, {"type": "content_block_stop", "index": index})
345 ]
347 def _ensure_block(self, kind: _BlockKind) -> list[tuple[AnthropicEventType, dict[str, Any]]]:
348 if self._open == kind:
349 return []
350 events = self._close_open()
351 shell = (
352 {"type": "thinking", "thinking": ""}
353 if kind is _BlockKind.THINKING
354 else {"type": "text", "text": ""}
355 )
356 events.append(
357 (
358 AnthropicEventType.CONTENT_BLOCK_START,
359 {
360 "type": "content_block_start",
361 "index": self._next_index,
362 "content_block": shell,
363 },
364 )
365 )
366 self._open = kind
367 self._next_index += 1
368 return events
370 def _text_events(
371 self, tokens: list[StreamToken]
372 ) -> list[tuple[AnthropicEventType, dict[str, Any]]]:
373 events: list[tuple[AnthropicEventType, dict[str, Any]]] = []
374 for token in tokens:
375 if not token.content:
376 continue
377 thinking = token.is_reasoning and not self._inline
378 kind = _BlockKind.THINKING if thinking else _BlockKind.TEXT
379 events.extend(self._ensure_block(kind))
380 index = self._next_index - 1
381 delta = (
382 {"type": "thinking_delta", "thinking": token.content}
383 if thinking
384 else {"type": "text_delta", "text": token.content}
385 )
386 events.append(
387 (
388 AnthropicEventType.CONTENT_BLOCK_DELTA,
389 {"type": "content_block_delta", "index": index, "delta": delta},
390 )
391 )
392 return events
394 def block_start(
395 self, event: ContentBlockStart
396 ) -> list[tuple[AnthropicEventType, dict[str, Any]]]:
397 if isinstance(event.block, ToolUseBlock):
398 events = self._close_open()
399 index = self._next_index
400 self._next_index += 1
401 self._open = _BlockKind.TOOL
402 events.append(
403 (
404 AnthropicEventType.CONTENT_BLOCK_START,
405 {
406 "type": "content_block_start",
407 "index": index,
408 "content_block": {
409 "type": "tool_use",
410 "id": event.block.id,
411 "name": event.block.name,
412 "input": {},
413 },
414 },
415 )
416 )
417 if event.block.input:
418 # A provider that announces a whole call up front carries the
419 # parsed input on the start block; forward it as one delta so
420 # SDK accumulation still sees arguments.
421 events.append(
422 (
423 AnthropicEventType.CONTENT_BLOCK_DELTA,
424 {
425 "type": "content_block_delta",
426 "index": index,
427 "delta": {
428 "type": "input_json_delta",
429 "partial_json": json.dumps(event.block.input),
430 },
431 },
432 )
433 )
434 return events
435 # Text blocks open lazily on the first delta so an empty block never
436 # emits a start/stop pair.
437 return []
439 def block_delta(
440 self, event: ContentBlockDelta
441 ) -> list[tuple[AnthropicEventType, dict[str, Any]]]:
442 if isinstance(event.delta, TextDelta):
443 text = event.delta.text
444 if self._pseudo is not None:
445 text = self._pseudo.feed(text)
446 return self._text_events(self._reasoning.feed(text))
447 if self._open is not _BlockKind.TOOL:
448 # A tool delta for a block that never started is a provider quirk,
449 # not a stream error; dropping beats crashing the stream.
450 return []
451 return [
452 (
453 AnthropicEventType.CONTENT_BLOCK_DELTA,
454 {
455 "type": "content_block_delta",
456 "index": self._next_index - 1,
457 "delta": {
458 "type": "input_json_delta",
459 "partial_json": event.delta.partial_json,
460 },
461 },
462 )
463 ]
465 def block_stop(self) -> list[tuple[AnthropicEventType, dict[str, Any]]]:
466 tokens = self._reasoning.feed(self._pseudo.flush()) if self._pseudo is not None else []
467 remaining = self._reasoning.flush()
468 if remaining is not None:
469 tokens.append(remaining)
470 events = self._text_events(tokens)
471 events.extend(self._close_open())
472 return events
475async def canonical_stream_to_anthropic_events(
476 events: AsyncIterator[CanonicalStreamEvent],
477 *,
478 model: str,
479 response_id: str,
480 mode: ReasoningMode = ReasoningMode.SEPARATE,
481) -> AsyncIterator[tuple[AnthropicEventType, dict[str, Any]]]:
482 """Turn canonical stream events into Anthropic SSE ``(type, payload)`` pairs."""
483 mapper = _AnthropicStreamMapper(mode=mode)
484 yield (
485 AnthropicEventType.MESSAGE_START,
486 {
487 "type": "message_start",
488 "message": {
489 "id": response_id,
490 "type": "message",
491 "role": "assistant",
492 "model": model,
493 "content": [],
494 "stop_reason": None,
495 "stop_sequence": None,
496 # Anthropic puts the prompt-side counts here, but the engine
497 # reports usage only in its closing chunk, so every count is
498 # zeroed and the real numbers arrive with the message_delta.
499 "usage": {
500 "input_tokens": 0,
501 "output_tokens": 0,
502 "cache_creation_input_tokens": 0,
503 "cache_read_input_tokens": 0,
504 },
505 },
506 },
507 )
508 async for event in events:
509 if isinstance(event, ContentBlockStart):
510 for out in mapper.block_start(event):
511 yield out
512 elif isinstance(event, ContentBlockDelta):
513 for out in mapper.block_delta(event):
514 yield out
515 elif isinstance(event, ContentBlockStop):
516 for out in mapper.block_stop():
517 yield out
518 elif isinstance(event, MessageDelta):
519 usage = event.usage or CanonicalUsage(input_tokens=0, output_tokens=0)
520 yield (
521 AnthropicEventType.MESSAGE_DELTA,
522 {
523 "type": "message_delta",
524 "delta": {
525 "stop_reason": str(event.stop_reason or StopReason.END_TURN),
526 "stop_sequence": None,
527 },
528 "usage": {
529 "input_tokens": usage.uncached_input_tokens,
530 "output_tokens": usage.output_tokens,
531 "cache_creation_input_tokens": 0,
532 "cache_read_input_tokens": usage.cached_input_tokens,
533 },
534 },
535 )
536 elif isinstance(event, MessageStart | MessageStop):
537 # The Anthropic message_start is emitted eagerly above (it needs no
538 # canonical data), and message_stop follows the loop so it stays
539 # last even if the provider never sends one.
540 continue
541 yield (AnthropicEventType.MESSAGE_STOP, {"type": "message_stop"})