Coverage for open_webui/utils/anthropic.py: 7%

485 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 05:07 +0000

1import logging 

2 

3import aiohttp 

4from open_webui.env import ( 

5 AIOHTTP_CLIENT_SESSION_SSL, 

6 AIOHTTP_CLIENT_TIMEOUT_MODEL_LIST, 

7 ENABLE_FORWARD_USER_INFO_HEADERS, 

8) 

9from open_webui.models.users import UserModel 

10from open_webui.utils.headers import include_user_info_headers 

11from open_webui.utils.json_codec import JSONCodec 

12 

13log = logging.getLogger(__name__) 

14 

15ANTHROPIC_VERSION = '2023-06-01' 

16 

17ANTHROPIC_CONVERTED_REQUEST_PARAMS = { 

18 'model', 

19 'messages', 

20 'system', 

21 'max_tokens', 

22 'temperature', 

23 'top_p', 

24 'top_k', 

25 'stop_sequences', 

26 'stream', 

27 'metadata', 

28 'service_tier', 

29 'tools', 

30 'tool_choice', 

31 'reasoning_effort', 

32} 

33 

34 

35def is_anthropic_url(url: str) -> bool: 

36 """Check if the URL is an Anthropic API endpoint.""" 

37 return 'api.anthropic.com' in url 

38 

39 

40async def get_anthropic_models(url: str, key: str, user: UserModel = None) -> dict: 

41 """ 

42 Fetch models from Anthropic's /v1/models endpoint with pagination. 

43 Normalizes the response to OpenAI format. 

44 """ 

45 timeout = aiohttp.ClientTimeout(total=AIOHTTP_CLIENT_TIMEOUT_MODEL_LIST) 

46 all_models = [] 

47 after_id = None 

48 

49 try: 

50 async with aiohttp.ClientSession(timeout=timeout, trust_env=True) as session: 

51 headers = { 

52 'x-api-key': key, 

53 'anthropic-version': ANTHROPIC_VERSION, 

54 } 

55 

56 if ENABLE_FORWARD_USER_INFO_HEADERS and user: 

57 headers = include_user_info_headers(headers, user) 

58 

59 while True: 

60 params = {'limit': 1000} 

61 if after_id: 

62 params['after_id'] = after_id 

63 

64 async with session.get( 

65 f'{url}/models', 

66 headers=headers, 

67 params=params, 

68 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

69 ) as response: 

70 if response.status != 200: 

71 error_detail = f'HTTP Error: {response.status}' 

72 try: 

73 res = await response.json() 

74 if 'error' in res: 

75 error_detail = f'External Error: {res["error"]}' 

76 except Exception: 

77 pass 

78 return {'object': 'list', 'data': [], 'error': error_detail} 

79 

80 data = await response.json() 

81 

82 for model in data.get('data', []): 

83 all_models.append( 

84 { 

85 'id': model.get('id'), 

86 'object': 'model', 

87 'created': 0, 

88 'owned_by': 'anthropic', 

89 'name': model.get('display_name', model.get('id')), 

90 } 

91 ) 

92 

93 if not data.get('has_more', False): 

94 break 

95 after_id = data.get('last_id') 

96 

97 except Exception as e: 

98 log.error(f'Anthropic connection error: {e}') 

99 return None 

100 

101 return {'object': 'list', 'data': all_models} 

102 

103 

104############################## 

105# 

106# Anthropic Messages API Conversion Utilities 

107# 

108############################## 

109 

110 

111def _copy_cache_control(source: dict, target: dict) -> dict: 

112 if isinstance(source, dict) and 'cache_control' in source: 

113 target['cache_control'] = source['cache_control'] 

114 return target 

115 

116 

117def _has_cache_control(blocks: list) -> bool: 

118 return any(isinstance(block, dict) and 'cache_control' in block for block in blocks) 

119 

120 

121def _finalize_openai_content(blocks: list) -> str | list: 

122 if not blocks: 

123 return '' 

124 

125 if len(blocks) == 1 and blocks[0].get('type') == 'text' and not _has_cache_control(blocks): 

126 return blocks[0].get('text', '') 

127 

128 return blocks 

129 

130 

131def is_anthropic_messages_passthrough(url: str, api_config: dict | None = None) -> bool: 

132 api_config = api_config or {} 

133 provider = str(api_config.get('provider', '')).lower() 

134 

135 return is_anthropic_url(url or '') or provider == 'litellm' 

136 

137 

138def convert_anthropic_to_openai_payload( 

139 anthropic_payload: dict, passthrough_params: list[str] | str | None = None 

140) -> dict: 

141 """ 

142 Convert an Anthropic Messages API request to OpenAI Chat Completions format. 

143 

144 Anthropic format: 

145 {model, messages: [{role, content}], system, max_tokens, ...} 

146 OpenAI format: 

147 {model, messages: [{role, content}], max_tokens, ...} 

148 """ 

149 openai_payload = {} 

150 

151 # Model 

152 openai_payload['model'] = anthropic_payload.get('model', '') 

153 

154 # Build messages list 

155 messages = [] 

156 

157 # System prompt (Anthropic has it as top-level, OpenAI as a system message) 

158 system = anthropic_payload.get('system') 

159 if system: 159 ↛ 160line 159 didn't jump to line 160 because the condition on line 159 was never true

160 if isinstance(system, str): 

161 messages.append({'role': 'system', 'content': system}) 

162 elif isinstance(system, list): 

163 openai_content = [] 

164 for block in system: 

165 if isinstance(block, dict) and block.get('type') == 'text': 

166 openai_content.append( 

167 _copy_cache_control( 

168 block, 

169 { 

170 'type': 'text', 

171 'text': block.get('text', ''), 

172 }, 

173 ) 

174 ) 

175 elif isinstance(block, str): 

176 openai_content.append({'type': 'text', 'text': block}) 

177 messages.append({'role': 'system', 'content': _finalize_openai_content(openai_content)}) 

178 

179 # Convert messages 

180 for msg in anthropic_payload.get('messages', []): 180 ↛ 181line 180 didn't jump to line 181 because the loop on line 180 never started

181 role = msg.get('role', 'user') 

182 content = msg.get('content') 

183 

184 if isinstance(content, str): 

185 messages.append({'role': role, 'content': content}) 

186 elif isinstance(content, list): 

187 # Convert Anthropic content blocks to OpenAI format 

188 openai_content = [] 

189 tool_calls = [] 

190 

191 for block in content: 

192 block_type = block.get('type', 'text') 

193 

194 if block_type == 'text': 

195 openai_content.append( 

196 _copy_cache_control( 

197 block, 

198 { 

199 'type': 'text', 

200 'text': block.get('text', ''), 

201 }, 

202 ) 

203 ) 

204 elif block_type in ('thinking', 'redacted_thinking'): 

205 # Unsigned thinking cannot be replayed upstream 

206 if block_type == 'redacted_thinking' or block.get('signature'): 

207 openai_content.append(_copy_cache_control(block, dict(block))) 

208 elif block_type == 'image': 

209 source = block.get('source', {}) 

210 if source.get('type') == 'base64': 

211 media_type = source.get('media_type', 'image/png') 

212 data = source.get('data', '') 

213 openai_content.append( 

214 _copy_cache_control( 

215 block, 

216 { 

217 'type': 'image_url', 

218 'image_url': { 

219 'url': f'data:{media_type};base64,{data}', 

220 }, 

221 }, 

222 ) 

223 ) 

224 elif source.get('type') == 'url': 

225 openai_content.append( 

226 _copy_cache_control( 

227 block, 

228 { 

229 'type': 'image_url', 

230 'image_url': {'url': source.get('url', '')}, 

231 }, 

232 ) 

233 ) 

234 elif block_type == 'tool_use': 

235 tool_calls.append( 

236 { 

237 'id': block.get('id', ''), 

238 'type': 'function', 

239 'function': { 

240 'name': block.get('name', ''), 

241 'arguments': ( 

242 JSONCodec.dumps(block.get('input', {})) 

243 if isinstance(block.get('input'), dict) 

244 else str(block.get('input', '{}')) 

245 ), 

246 }, 

247 } 

248 ) 

249 elif block_type == 'tool_result': 

250 # Tool results become separate tool messages in OpenAI format 

251 tool_result_content = block.get('content', '') 

252 tool_content: str | list = '' 

253 

254 if isinstance(tool_result_content, str): 

255 tool_content = tool_result_content 

256 elif isinstance(tool_result_content, list): 

257 # Build a multimodal content array to preserve 

258 # images and other non-text content types. 

259 converted_parts = [] 

260 for content_block in tool_result_content: 

261 if not isinstance(content_block, dict): 

262 continue 

263 content_type = content_block.get('type', 'text') 

264 

265 if content_type == 'text': 

266 converted_parts.append( 

267 _copy_cache_control( 

268 content_block, 

269 { 

270 'type': 'text', 

271 'text': content_block.get('text', ''), 

272 }, 

273 ) 

274 ) 

275 elif content_type == 'image': 

276 source = content_block.get('source', {}) 

277 if source.get('type') == 'base64': 

278 media_type = source.get('media_type', 'image/png') 

279 data = source.get('data', '') 

280 converted_parts.append( 

281 _copy_cache_control( 

282 content_block, 

283 { 

284 'type': 'image_url', 

285 'image_url': { 

286 'url': f'data:{media_type};base64,{data}', 

287 }, 

288 }, 

289 ) 

290 ) 

291 elif source.get('type') == 'url': 

292 converted_parts.append( 

293 _copy_cache_control( 

294 content_block, 

295 { 

296 'type': 'image_url', 

297 'image_url': { 

298 'url': source.get('url', ''), 

299 }, 

300 }, 

301 ) 

302 ) 

303 elif content_type == 'document': 

304 # Documents have no direct OpenAI equivalent; 

305 # convert to a text representation. 

306 document_source = content_block.get('source', {}) 

307 document_title = content_block.get('title', 'Document') 

308 document_context = content_block.get('context', '') 

309 document_text = f'[Document: {document_title}]' 

310 if document_context: 

311 document_text += f'\n{document_context}' 

312 if document_source.get('type') == 'text' and document_source.get('data'): 

313 document_text += f'\n{document_source["data"]}' 

314 converted_parts.append({'type': 'text', 'text': document_text}) 

315 elif content_type == 'search_result': 

316 # Convert search results to a text 

317 # representation with source attribution. 

318 search_title = content_block.get('title', '') 

319 search_url = content_block.get('source', '') 

320 search_content_blocks = content_block.get('content', []) 

321 search_texts = [] 

322 for search_block in search_content_blocks: 

323 if isinstance(search_block, dict) and search_block.get('type') == 'text': 

324 search_texts.append(search_block.get('text', '')) 

325 search_body = '\n'.join(search_texts) 

326 search_text = f'[Search Result: {search_title}]' 

327 if search_url: 

328 search_text += f'\nSource: {search_url}' 

329 if search_body: 

330 search_text += f'\n{search_body}' 

331 converted_parts.append({'type': 'text', 'text': search_text}) 

332 

333 # Flatten to string when only text parts are present 

334 if all(part.get('type') == 'text' for part in converted_parts) and not _has_cache_control( 

335 converted_parts 

336 ): 

337 tool_content = '\n'.join(part.get('text', '') for part in converted_parts) 

338 elif converted_parts: 

339 tool_content = converted_parts 

340 else: 

341 tool_content = '' 

342 

343 # Propagate error status if present 

344 if block.get('is_error'): 

345 if isinstance(tool_content, str): 

346 tool_content = f'Error: {tool_content}' 

347 elif isinstance(tool_content, list): 

348 tool_content.insert( 

349 0, 

350 { 

351 'type': 'text', 

352 'text': 'Error: ', 

353 }, 

354 ) 

355 

356 messages.append( 

357 { 

358 'role': 'tool', 

359 'tool_call_id': block.get('tool_use_id', ''), 

360 'content': tool_content, 

361 } 

362 ) 

363 

364 # Build the message 

365 if tool_calls: 

366 # Assistant message with tool calls 

367 msg_dict = {'role': role} 

368 if openai_content: 

369 msg_dict['content'] = _finalize_openai_content(openai_content) 

370 else: 

371 msg_dict['content'] = '' 

372 msg_dict['tool_calls'] = tool_calls 

373 messages.append(msg_dict) 

374 elif openai_content or role == 'assistant': 

375 messages.append({'role': role, 'content': _finalize_openai_content(openai_content)}) 

376 else: 

377 messages.append({'role': role, 'content': str(content) if content else ''}) 

378 

379 openai_payload['messages'] = messages 

380 

381 # max_tokens 

382 if 'max_tokens' in anthropic_payload: 382 ↛ 383line 382 didn't jump to line 383 because the condition on line 382 was never true

383 openai_payload['max_tokens'] = anthropic_payload['max_tokens'] 

384 

385 captured_passthrough_params = { 

386 param: value for param, value in anthropic_payload.items() if param not in ANTHROPIC_CONVERTED_REQUEST_PARAMS 

387 } 

388 if isinstance(passthrough_params, str): 388 ↛ 389line 388 didn't jump to line 389 because the condition on line 388 was never true

389 passthrough_params = passthrough_params.split(',') 

390 elif not isinstance(passthrough_params, (list, tuple, set)): 390 ↛ 391line 390 didn't jump to line 391 because the condition on line 390 was never true

391 passthrough_params = [] 

392 passthrough_param_names = {str(item).strip() for item in passthrough_params if str(item).strip()} 

393 if '*' in passthrough_param_names: 393 ↛ 394line 393 didn't jump to line 394 because the condition on line 393 was never true

394 openai_payload.update(captured_passthrough_params) 

395 else: 

396 for param in passthrough_param_names: 396 ↛ 397line 396 didn't jump to line 397 because the loop on line 396 never started

397 if param in captured_passthrough_params: 

398 openai_payload[param] = captured_passthrough_params[param] 

399 

400 output_config = anthropic_payload.get('output_config') 

401 if isinstance(output_config, dict): 401 ↛ 402line 401 didn't jump to line 402 because the condition on line 401 was never true

402 if 'effort' in output_config and 'reasoning_effort' not in anthropic_payload: 

403 openai_payload['reasoning_effort'] = output_config['effort'] 

404 

405 format_config = output_config.get('format') 

406 if isinstance(format_config, dict): 

407 format_type = format_config.get('type') 

408 if format_type == 'json_schema': 

409 json_schema = { 

410 'name': format_config.get('name', 'response_schema'), 

411 'schema': format_config.get('schema', {}), 

412 } 

413 if 'description' in format_config: 

414 json_schema['description'] = format_config['description'] 

415 if 'strict' in format_config: 

416 json_schema['strict'] = format_config['strict'] 

417 openai_payload['response_format'] = { 

418 'type': 'json_schema', 

419 'json_schema': json_schema, 

420 } 

421 elif format_type == 'json_object': 

422 openai_payload['response_format'] = {'type': format_type} 

423 

424 if 'reasoning_effort' in anthropic_payload: 424 ↛ 425line 424 didn't jump to line 425 because the condition on line 424 was never true

425 openai_payload['reasoning_effort'] = anthropic_payload['reasoning_effort'] 

426 

427 # Common parameters 

428 for param in ('temperature', 'top_p', 'top_k', 'stop_sequences', 'stream', 'metadata', 'service_tier'): 

429 if param in anthropic_payload: 429 ↛ 430line 429 didn't jump to line 430 because the condition on line 429 was never true

430 if param == 'stop_sequences': 

431 openai_payload['stop'] = anthropic_payload[param] 

432 else: 

433 openai_payload[param] = anthropic_payload[param] 

434 

435 # Tools conversion: Anthropic → OpenAI 

436 if 'tools' in anthropic_payload: 436 ↛ 437line 436 didn't jump to line 437 because the condition on line 436 was never true

437 openai_tools = [] 

438 for tool in anthropic_payload['tools']: 

439 openai_tools.append( 

440 _copy_cache_control( 

441 tool, 

442 { 

443 'type': 'function', 

444 'function': { 

445 'name': tool.get('name', ''), 

446 'description': tool.get('description', ''), 

447 'parameters': tool.get('input_schema', {}), 

448 }, 

449 }, 

450 ) 

451 ) 

452 openai_payload['tools'] = openai_tools 

453 

454 # tool_choice 

455 if 'tool_choice' in anthropic_payload: 455 ↛ 456line 455 didn't jump to line 456 because the condition on line 455 was never true

456 tool_choice = anthropic_payload['tool_choice'] 

457 if isinstance(tool_choice, dict): 

458 tool_choice_type = tool_choice.get('type', 'auto') 

459 if tool_choice_type == 'auto': 

460 openai_payload['tool_choice'] = 'auto' 

461 elif tool_choice_type == 'any': 

462 openai_payload['tool_choice'] = 'required' 

463 elif tool_choice_type == 'tool': 

464 openai_payload['tool_choice'] = { 

465 'type': 'function', 

466 'function': {'name': tool_choice.get('name', '')}, 

467 } 

468 

469 return openai_payload 

470 

471 

472def convert_openai_to_anthropic_response( 

473 openai_response: dict, model: str = '', input_tokens: int | None = None 

474) -> dict: 

475 """ 

476 Convert a non-streaming OpenAI Chat Completions response to Anthropic Messages format. 

477 """ 

478 import uuid as _uuid 

479 

480 choice = {} 

481 if openai_response.get('choices'): 

482 choice = openai_response['choices'][0] 

483 

484 message = choice.get('message', {}) 

485 finish_reason = choice.get('finish_reason', 'stop') 

486 

487 # Map finish_reason to stop_reason 

488 stop_reason_map = { 

489 'stop': 'end_turn', 

490 'length': 'max_tokens', 

491 'tool_calls': 'tool_use', 

492 'content_filter': 'end_turn', 

493 } 

494 stop_reason = stop_reason_map.get(finish_reason, 'end_turn') 

495 

496 # Build content blocks 

497 content = [] 

498 message_thinking = message.get('thinking') 

499 thinking_blocks = message.get('thinking_blocks') or [] 

500 if not thinking_blocks and isinstance(message_thinking, dict): 

501 thinking_blocks = message_thinking.get('blocks') or [] 

502 

503 has_thinking = False 

504 for block in thinking_blocks: 

505 if not isinstance(block, dict): 

506 continue 

507 

508 if block.get('type') == 'redacted_thinking': 

509 content.append({k: v for k, v in block.items() if k in {'type', 'data'}}) 

510 has_thinking = True 

511 continue 

512 

513 thinking = block.get('thinking') or block.get('content') or block.get('text') 

514 if not thinking: 

515 continue 

516 

517 thinking_block = {'type': 'thinking', 'thinking': thinking} 

518 if block.get('signature'): 

519 thinking_block['signature'] = block['signature'] 

520 content.append(thinking_block) 

521 has_thinking = True 

522 

523 reasoning_content = message.get('reasoning_content') or message.get('reasoning') 

524 if not reasoning_content and isinstance(message_thinking, str): 

525 reasoning_content = message_thinking 

526 if reasoning_content and not has_thinking: 

527 content.append({'type': 'thinking', 'thinking': reasoning_content}) 

528 

529 message_content = message.get('content') 

530 if message_content: 

531 content.append({'type': 'text', 'text': message_content}) 

532 

533 # Tool calls -> tool_use blocks 

534 tool_calls = message.get('tool_calls') or [] 

535 for tool_call in tool_calls: 

536 function = tool_call.get('function', {}) 

537 try: 

538 tool_input = JSONCodec.loads(function.get('arguments', '{}')) 

539 except (JSONCodec.JSONDecodeError, TypeError): 

540 tool_input = {} 

541 content.append( 

542 { 

543 'type': 'tool_use', 

544 'id': tool_call.get('id', f'toolu_{_uuid.uuid4().hex[:24]}'), 

545 'name': function.get('name', ''), 

546 'input': tool_input, 

547 } 

548 ) 

549 

550 # Usage 

551 openai_usage = openai_response.get('usage') or {} 

552 cache_creation = openai_usage.get('cache_creation_input_tokens') 

553 cache_read = openai_usage.get('cache_read_input_tokens') 

554 prompt_details = openai_usage.get('prompt_tokens_details') 

555 if cache_read is None and isinstance(prompt_details, dict): 

556 cache_read = prompt_details.get('cached_tokens') 

557 

558 usage_input = openai_usage.get('input_tokens') 

559 if usage_input is None: 

560 prompt_tokens = openai_usage.get('prompt_tokens') 

561 if prompt_tokens is not None: 

562 usage_input = max(prompt_tokens - (cache_creation or 0) - (cache_read or 0), 0) 

563 

564 usage_output = openai_usage.get('output_tokens') 

565 if usage_output is None: 

566 usage_output = openai_usage.get('completion_tokens') 

567 

568 usage = { 

569 'input_tokens': usage_input if usage_input is not None else (input_tokens if input_tokens is not None else 0), 

570 'output_tokens': usage_output if usage_output is not None else 0, 

571 } 

572 if cache_creation is not None: 

573 usage['cache_creation_input_tokens'] = cache_creation 

574 if cache_read is not None: 

575 usage['cache_read_input_tokens'] = cache_read 

576 if isinstance(openai_usage.get('output_tokens_details'), dict): 

577 usage['output_tokens_details'] = openai_usage['output_tokens_details'] 

578 if isinstance(openai_usage.get('server_tool_use'), dict): 

579 usage['server_tool_use'] = openai_usage['server_tool_use'] 

580 if openai_usage.get('service_tier') is not None: 

581 usage['service_tier'] = openai_usage['service_tier'] 

582 

583 return { 

584 'id': openai_response.get('id', f'msg_{_uuid.uuid4().hex[:24]}'), 

585 'type': 'message', 

586 'role': 'assistant', 

587 'content': content, 

588 'model': model or openai_response.get('model', ''), 

589 'stop_reason': stop_reason, 

590 'stop_sequence': None, 

591 'usage': usage, 

592 } 

593 

594 

595async def openai_stream_to_anthropic_stream(openai_stream_generator, model: str = '', input_tokens: int | None = None): 

596 """ 

597 Convert an OpenAI SSE streaming response to Anthropic Messages SSE format. 

598 

599 OpenAI sends: data: {"choices": [{"delta": {"content": "..."}}]} 

600 Anthropic sends: event: content_block_delta\\ndata: {"type": "content_block_delta", ...} 

601 

602 Handles text content, tool calls, and mixed content with proper 

603 multi-block indexing as required by Anthropic's streaming protocol. 

604 

605 Tool calls are tracked by their unique id (not OpenAI index) so that 

606 parallel calls sharing the same index get distinct Anthropic tool_use 

607 blocks. Each block follows the Anthropic lifecycle: start -> delta -> stop. 

608 """ 

609 import uuid as _uuid 

610 

611 message_id = f'msg_{_uuid.uuid4().hex[:24]}' 

612 output_tokens = 0 

613 cache_creation_input_tokens = None 

614 cache_read_input_tokens = None 

615 output_tokens_details = None 

616 server_tool_use = None 

617 service_tier = None 

618 stop_reason = 'end_turn' 

619 

620 # Track content blocks with a running index. 

621 # Each text block or tool_use block gets its own index. 

622 current_block_index = 0 

623 thinking_block_open = False 

624 text_block_open = False 

625 

626 # Accumulated state for each tool call, keyed by tool call id. 

627 # Parallel calls that share the same OpenAI index get distinct entries. 

628 # Each entry: {id, name, arguments, block_index, started, stopped} 

629 tracked_tool_calls = {} 

630 # Map OpenAI tool call index -> tool call id for routing 

631 # argument-only deltas (deltas that carry arguments but no id). 

632 index_to_tool_id = {} 

633 # Whether any tool call block has been emitted (suppresses further text) 

634 has_tool_calls = False 

635 

636 # Emit message_start 

637 message_start = { 

638 'type': 'message_start', 

639 'message': { 

640 'id': message_id, 

641 'type': 'message', 

642 'role': 'assistant', 

643 'content': [], 

644 'model': model, 

645 'stop_reason': None, 

646 'stop_sequence': None, 

647 'usage': {'input_tokens': input_tokens or 0, 'output_tokens': 0}, 

648 }, 

649 } 

650 yield f'event: message_start\ndata: {JSONCodec.dumps(message_start)}\n\n'.encode() 

651 

652 try: 

653 async for chunk in openai_stream_generator: 

654 if isinstance(chunk, bytes): 

655 chunk = chunk.decode('utf-8', errors='ignore') 

656 

657 for line in chunk.strip().split('\n'): 

658 line = line.strip() 

659 

660 if not line or not line.startswith('data:'): 

661 continue 

662 

663 data_string = line[5:].strip() 

664 if data_string == '[DONE]': 

665 continue 

666 if data_string == '{}': 

667 continue 

668 

669 try: 

670 data = JSONCodec.loads(data_string) 

671 except (JSONCodec.JSONDecodeError, TypeError): 

672 continue 

673 

674 usage_data = data.get('usage') 

675 if isinstance(usage_data, dict): 

676 cache_creation = usage_data.get('cache_creation_input_tokens') 

677 cache_read = usage_data.get('cache_read_input_tokens') 

678 prompt_details = usage_data.get('prompt_tokens_details') 

679 if cache_read is None and isinstance(prompt_details, dict): 

680 cache_read = prompt_details.get('cached_tokens') 

681 

682 usage_input = usage_data.get('input_tokens') 

683 if usage_input is None: 

684 prompt_tokens = usage_data.get('prompt_tokens') 

685 if prompt_tokens is not None: 

686 usage_input = max(prompt_tokens - (cache_creation or 0) - (cache_read or 0), 0) 

687 

688 usage_output = usage_data.get('output_tokens') 

689 if usage_output is None: 

690 usage_output = usage_data.get('completion_tokens') 

691 

692 if usage_input is not None: 

693 input_tokens = usage_input 

694 if usage_output is not None: 

695 output_tokens = usage_output 

696 if cache_creation is not None: 

697 cache_creation_input_tokens = cache_creation 

698 if cache_read is not None: 

699 cache_read_input_tokens = cache_read 

700 if isinstance(usage_data.get('output_tokens_details'), dict): 

701 output_tokens_details = usage_data['output_tokens_details'] 

702 if isinstance(usage_data.get('server_tool_use'), dict): 

703 server_tool_use = usage_data['server_tool_use'] 

704 if usage_data.get('service_tier') is not None: 

705 service_tier = usage_data['service_tier'] 

706 

707 choices = data.get('choices', []) 

708 if not choices: 

709 continue 

710 

711 delta = choices[0].get('delta', {}) 

712 finish_reason = choices[0].get('finish_reason') 

713 message = choices[0].get('message') or {} 

714 

715 reasoning_content = ( 

716 delta.get('reasoning_content') 

717 or delta.get('reasoning') 

718 or delta.get('thinking') 

719 or message.get('reasoning_content') 

720 or message.get('reasoning') 

721 ) 

722 if not reasoning_content: 

723 thinking_blocks = delta.get('thinking_blocks') or message.get('thinking_blocks') or [] 

724 for block in thinking_blocks: 

725 if isinstance(block, dict): 

726 reasoning_content = block.get('thinking') or block.get('content') or block.get('text') 

727 if reasoning_content: 

728 break 

729 

730 if reasoning_content and not text_block_open and not has_tool_calls: 

731 if not thinking_block_open: 

732 block_start = { 

733 'type': 'content_block_start', 

734 'index': current_block_index, 

735 'content_block': {'type': 'thinking', 'thinking': ''}, 

736 } 

737 yield f'event: content_block_start\ndata: {JSONCodec.dumps(block_start)}\n\n'.encode() 

738 thinking_block_open = True 

739 

740 block_delta = { 

741 'type': 'content_block_delta', 

742 'index': current_block_index, 

743 'delta': {'type': 'thinking_delta', 'thinking': reasoning_content}, 

744 } 

745 yield f'event: content_block_delta\ndata: {JSONCodec.dumps(block_delta)}\n\n'.encode() 

746 

747 # --- Handle text content --- 

748 # Anthropic expects text blocks before tool blocks, so skip 

749 # text deltas once any tool call has started. 

750 content = delta.get('content') 

751 if content and not has_tool_calls: 

752 if thinking_block_open: 

753 block_stop = { 

754 'type': 'content_block_stop', 

755 'index': current_block_index, 

756 } 

757 yield f'event: content_block_stop\ndata: {JSONCodec.dumps(block_stop)}\n\n'.encode() 

758 thinking_block_open = False 

759 current_block_index += 1 

760 

761 if not text_block_open: 

762 block_start = { 

763 'type': 'content_block_start', 

764 'index': current_block_index, 

765 'content_block': {'type': 'text', 'text': ''}, 

766 } 

767 yield f'event: content_block_start\ndata: {JSONCodec.dumps(block_start)}\n\n'.encode() 

768 text_block_open = True 

769 

770 block_delta = { 

771 'type': 'content_block_delta', 

772 'index': current_block_index, 

773 'delta': {'type': 'text_delta', 'text': content}, 

774 } 

775 yield f'event: content_block_delta\ndata: {JSONCodec.dumps(block_delta)}\n\n'.encode() 

776 

777 # --- Handle tool calls --- 

778 # Some providers put tool_calls on the final message object 

779 # instead of the delta; fall back to that when needed. 

780 tool_calls = delta.get('tool_calls') or [] 

781 if not tool_calls and message.get('tool_calls'): 

782 tool_calls = message['tool_calls'] 

783 

784 if tool_calls: 

785 # Close text block if one is open (text comes before tools) 

786 if thinking_block_open: 

787 block_stop = { 

788 'type': 'content_block_stop', 

789 'index': current_block_index, 

790 } 

791 yield f'event: content_block_stop\ndata: {JSONCodec.dumps(block_stop)}\n\n'.encode() 

792 thinking_block_open = False 

793 current_block_index += 1 

794 

795 if text_block_open: 

796 block_stop = { 

797 'type': 'content_block_stop', 

798 'index': current_block_index, 

799 } 

800 yield f'event: content_block_stop\ndata: {JSONCodec.dumps(block_stop)}\n\n'.encode() 

801 text_block_open = False 

802 current_block_index += 1 

803 

804 for tool_call in tool_calls: 

805 tool_call_index = tool_call.get('index', 0) 

806 tool_call_id = tool_call.get('id', '') 

807 tool_call_name = (tool_call.get('function') or {}).get('name', '') 

808 arguments_chunk = (tool_call.get('function') or {}).get('arguments', '') 

809 

810 # Resolve which tracked tool call this delta belongs to. 

811 # A delta with an id starts or identifies a specific tool. 

812 # A delta without an id carries arguments for the most 

813 # recent tool at this OpenAI index. 

814 if tool_call_id: 

815 if tool_call_id not in tracked_tool_calls: 

816 tracked_tool_calls[tool_call_id] = { 

817 'id': tool_call_id, 

818 'name': tool_call_name, 

819 'arguments': '', 

820 'block_index': -1, 

821 'started': False, 

822 'stopped': False, 

823 } 

824 index_to_tool_id[tool_call_index] = tool_call_id 

825 tool = tracked_tool_calls[tool_call_id] 

826 elif tool_call_index in index_to_tool_id: 

827 tool = tracked_tool_calls[index_to_tool_id[tool_call_index]] 

828 else: 

829 # First delta for this index with no id; create a 

830 # provisional entry with a generated fallback id. 

831 fallback_id = f'toolu_{_uuid.uuid4().hex[:24]}' 

832 tracked_tool_calls[fallback_id] = { 

833 'id': fallback_id, 

834 'name': tool_call_name, 

835 'arguments': '', 

836 'block_index': -1, 

837 'started': False, 

838 'stopped': False, 

839 } 

840 index_to_tool_id[tool_call_index] = fallback_id 

841 tool = tracked_tool_calls[fallback_id] 

842 

843 # Update name if provided on a later delta 

844 if tool_call_name and not tool['name']: 

845 tool['name'] = tool_call_name 

846 

847 # Emit content_block_start once we have a name 

848 if not tool['started'] and tool['name']: 

849 tool['block_index'] = current_block_index 

850 tool['started'] = True 

851 has_tool_calls = True 

852 

853 block_start = { 

854 'type': 'content_block_start', 

855 'index': current_block_index, 

856 'content_block': { 

857 'type': 'tool_use', 

858 'id': tool['id'], 

859 'name': tool['name'], 

860 'input': {}, 

861 }, 

862 } 

863 yield f'event: content_block_start\ndata: {JSONCodec.dumps(block_start)}\n\n'.encode() 

864 current_block_index += 1 

865 

866 # Buffer arguments and emit as input_json_delta 

867 if arguments_chunk: 

868 tool['arguments'] += arguments_chunk 

869 

870 if tool['started'] and not tool['stopped']: 

871 block_delta = { 

872 'type': 'content_block_delta', 

873 'index': tool['block_index'], 

874 'delta': { 

875 'type': 'input_json_delta', 

876 'partial_json': arguments_chunk, 

877 }, 

878 } 

879 yield f'event: content_block_delta\ndata: {JSONCodec.dumps(block_delta)}\n\n'.encode() 

880 

881 # Close the block once arguments form complete JSON 

882 if ( 

883 tool['started'] 

884 and not tool['stopped'] 

885 and (tool['arguments'].rstrip()[-1:] == '}' or tool['arguments'].lstrip()[:1] != '{') 

886 ): 

887 try: 

888 JSONCodec.loads(tool['arguments']) 

889 tool['stopped'] = True 

890 block_stop = { 

891 'type': 'content_block_stop', 

892 'index': tool['block_index'], 

893 } 

894 yield f'event: content_block_stop\ndata: {JSONCodec.dumps(block_stop)}\n\n'.encode() 

895 except (JSONCodec.JSONDecodeError, ValueError): 

896 pass 

897 

898 # --- Handle finish reason --- 

899 if finish_reason is not None: 

900 stop_reason_map = { 

901 'stop': 'end_turn', 

902 'length': 'max_tokens', 

903 'tool_calls': 'tool_use', 

904 } 

905 stop_reason = stop_reason_map.get(finish_reason, 'end_turn') 

906 

907 except Exception as e: 

908 log.error(f'Error in Anthropic stream conversion: {e}') 

909 

910 # Close any open thinking block 

911 if thinking_block_open: 

912 block_stop = {'type': 'content_block_stop', 'index': current_block_index} 

913 yield f'event: content_block_stop\ndata: {JSONCodec.dumps(block_stop)}\n\n'.encode() 

914 current_block_index += 1 

915 

916 # Flush any tools that buffered arguments but never emitted a block 

917 for tool in tracked_tool_calls.values(): 

918 if not tool['started'] and tool['name']: 

919 tool['block_index'] = current_block_index 

920 tool['started'] = True 

921 

922 block_start = { 

923 'type': 'content_block_start', 

924 'index': current_block_index, 

925 'content_block': { 

926 'type': 'tool_use', 

927 'id': tool['id'], 

928 'name': tool['name'], 

929 'input': {}, 

930 }, 

931 } 

932 yield f'event: content_block_start\ndata: {JSONCodec.dumps(block_start)}\n\n'.encode() 

933 current_block_index += 1 

934 

935 if tool['arguments']: 

936 block_delta = { 

937 'type': 'content_block_delta', 

938 'index': tool['block_index'], 

939 'delta': { 

940 'type': 'input_json_delta', 

941 'partial_json': tool['arguments'], 

942 }, 

943 } 

944 yield f'event: content_block_delta\ndata: {JSONCodec.dumps(block_delta)}\n\n'.encode() 

945 

946 # Close any open text block 

947 if text_block_open: 

948 block_stop = {'type': 'content_block_stop', 'index': current_block_index} 

949 yield f'event: content_block_stop\ndata: {JSONCodec.dumps(block_stop)}\n\n'.encode() 

950 

951 # Close any tool call blocks that are still open 

952 for tool in tracked_tool_calls.values(): 

953 if tool['started'] and not tool['stopped']: 

954 block_stop = {'type': 'content_block_stop', 'index': tool['block_index']} 

955 yield f'event: content_block_stop\ndata: {JSONCodec.dumps(block_stop)}\n\n'.encode() 

956 

957 # Emit message_delta with stop reason 

958 usage = {'output_tokens': output_tokens} 

959 if input_tokens is not None: 

960 usage['input_tokens'] = input_tokens 

961 if cache_creation_input_tokens is not None: 

962 usage['cache_creation_input_tokens'] = cache_creation_input_tokens 

963 if cache_read_input_tokens is not None: 

964 usage['cache_read_input_tokens'] = cache_read_input_tokens 

965 if output_tokens_details is not None: 

966 usage['output_tokens_details'] = output_tokens_details 

967 if server_tool_use is not None: 

968 usage['server_tool_use'] = server_tool_use 

969 if service_tier is not None: 

970 usage['service_tier'] = service_tier 

971 

972 message_delta = { 

973 'type': 'message_delta', 

974 'delta': { 

975 'stop_reason': stop_reason, 

976 'stop_sequence': None, 

977 }, 

978 'usage': usage, 

979 } 

980 yield f'event: message_delta\ndata: {JSONCodec.dumps(message_delta)}\n\n'.encode() 

981 

982 # Emit message_stop 

983 yield f'event: message_stop\ndata: {JSONCodec.dumps({"type": "message_stop"})}\n\n'.encode()