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

1"""Translation between Anthropic Messages models and the canonical types.""" 

2 

3from __future__ import annotations 

4 

5import json 

6from collections.abc import AsyncIterator 

7from dataclasses import replace 

8from enum import StrEnum 

9from typing import Any, Literal 

10 

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) 

60 

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.' 

65 

66 

67class _BlockKind(StrEnum): 

68 """Kind of the mapper's open output block.""" 

69 

70 THINKING = "thinking" 

71 TEXT = "text" 

72 TOOL = "tool" 

73 

74 

75_ANTHROPIC_CHOICE_MODES: dict[str, str] = { 

76 "auto": "auto", 

77 "any": "any", 

78 "none": "none", 

79} 

80 

81 

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. 

86 

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 

97 

98 

99def count_tokens_to_canonical_request( 

100 request: CountTokensRequest, *, mode: ReasoningMode = ReasoningMode.SEPARATE 

101) -> CanonicalChatRequest: 

102 """Translate a validated ``CountTokensRequest`` to the canonical request. 

103 

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) 

108 

109 

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 ) 

123 

124 

125def _canonical_prompt(request: _PromptBody, *, mode: ReasoningMode) -> CanonicalChatRequest: 

126 """Translate the fields that decide the rendered prompt. 

127 

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 ) 

141 

142 

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 

149 

150 

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 

158 

159 

160def _canonical_messages_for(msg: AnthropicMessage) -> list[CanonicalMessage]: 

161 """Fan one Anthropic message out to canonical messages. 

162 

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) 

184 

185 

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 

203 

204 out = tool_messages 

205 if blocks: 

206 out = [*tool_messages, CanonicalMessage(role=role, content=blocks)] 

207 return out 

208 

209 

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 ) 

221 

222 

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)) 

228 

229 

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 

243 

244 

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 ] 

256 

257 

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] 

269 

270 

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. 

275 

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 ) 

312 

313 

314class _AnthropicStreamMapper: 

315 """Per-stream state for the canonical-to-Anthropic event converter. 

316 

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. 

321 

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 """ 

328 

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 

337 

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 ] 

346 

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 

369 

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 

393 

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 [] 

438 

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 ] 

464 

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 

473 

474 

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"})