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
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
1import logging
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
13log = logging.getLogger(__name__)
15ANTHROPIC_VERSION = '2023-06-01'
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}
35def is_anthropic_url(url: str) -> bool:
36 """Check if the URL is an Anthropic API endpoint."""
37 return 'api.anthropic.com' in url
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
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 }
56 if ENABLE_FORWARD_USER_INFO_HEADERS and user:
57 headers = include_user_info_headers(headers, user)
59 while True:
60 params = {'limit': 1000}
61 if after_id:
62 params['after_id'] = after_id
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}
80 data = await response.json()
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 )
93 if not data.get('has_more', False):
94 break
95 after_id = data.get('last_id')
97 except Exception as e:
98 log.error(f'Anthropic connection error: {e}')
99 return None
101 return {'object': 'list', 'data': all_models}
104##############################
105#
106# Anthropic Messages API Conversion Utilities
107#
108##############################
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
117def _has_cache_control(blocks: list) -> bool:
118 return any(isinstance(block, dict) and 'cache_control' in block for block in blocks)
121def _finalize_openai_content(blocks: list) -> str | list:
122 if not blocks:
123 return ''
125 if len(blocks) == 1 and blocks[0].get('type') == 'text' and not _has_cache_control(blocks):
126 return blocks[0].get('text', '')
128 return blocks
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()
135 return is_anthropic_url(url or '') or provider == 'litellm'
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.
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 = {}
151 # Model
152 openai_payload['model'] = anthropic_payload.get('model', '')
154 # Build messages list
155 messages = []
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)})
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')
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 = []
191 for block in content:
192 block_type = block.get('type', 'text')
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 = ''
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')
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})
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 = ''
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 )
356 messages.append(
357 {
358 'role': 'tool',
359 'tool_call_id': block.get('tool_use_id', ''),
360 'content': tool_content,
361 }
362 )
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 ''})
379 openai_payload['messages'] = messages
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']
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]
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']
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}
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']
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]
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
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 }
469 return openai_payload
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
480 choice = {}
481 if openai_response.get('choices'):
482 choice = openai_response['choices'][0]
484 message = choice.get('message', {})
485 finish_reason = choice.get('finish_reason', 'stop')
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')
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 []
503 has_thinking = False
504 for block in thinking_blocks:
505 if not isinstance(block, dict):
506 continue
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
513 thinking = block.get('thinking') or block.get('content') or block.get('text')
514 if not thinking:
515 continue
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
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})
529 message_content = message.get('content')
530 if message_content:
531 content.append({'type': 'text', 'text': message_content})
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 )
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')
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)
564 usage_output = openai_usage.get('output_tokens')
565 if usage_output is None:
566 usage_output = openai_usage.get('completion_tokens')
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']
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 }
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.
599 OpenAI sends: data: {"choices": [{"delta": {"content": "..."}}]}
600 Anthropic sends: event: content_block_delta\\ndata: {"type": "content_block_delta", ...}
602 Handles text content, tool calls, and mixed content with proper
603 multi-block indexing as required by Anthropic's streaming protocol.
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
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'
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
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
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()
652 try:
653 async for chunk in openai_stream_generator:
654 if isinstance(chunk, bytes):
655 chunk = chunk.decode('utf-8', errors='ignore')
657 for line in chunk.strip().split('\n'):
658 line = line.strip()
660 if not line or not line.startswith('data:'):
661 continue
663 data_string = line[5:].strip()
664 if data_string == '[DONE]':
665 continue
666 if data_string == '{}':
667 continue
669 try:
670 data = JSONCodec.loads(data_string)
671 except (JSONCodec.JSONDecodeError, TypeError):
672 continue
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')
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)
688 usage_output = usage_data.get('output_tokens')
689 if usage_output is None:
690 usage_output = usage_data.get('completion_tokens')
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']
707 choices = data.get('choices', [])
708 if not choices:
709 continue
711 delta = choices[0].get('delta', {})
712 finish_reason = choices[0].get('finish_reason')
713 message = choices[0].get('message') or {}
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
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
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()
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
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
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()
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']
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
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
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', '')
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]
843 # Update name if provided on a later delta
844 if tool_call_name and not tool['name']:
845 tool['name'] = tool_call_name
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
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
866 # Buffer arguments and emit as input_json_delta
867 if arguments_chunk:
868 tool['arguments'] += arguments_chunk
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()
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
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')
907 except Exception as e:
908 log.error(f'Error in Anthropic stream conversion: {e}')
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
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
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
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()
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()
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()
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
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()
982 # Emit message_stop
983 yield f'event: message_stop\ndata: {JSONCodec.dumps({"type": "message_stop"})}\n\n'.encode()