From 1e72e22ebfea823c850c816280bf0efab4c9db77 Mon Sep 17 00:00:00 2001 From: Sameer Kankute Date: Wed, 15 Apr 2026 17:37:23 +0530 Subject: [PATCH] feat(messages-api): cross-provider advisor orchestration for /v1/messages MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Run MessagesInterceptor checks before pre-request hooks so the synthetic advisor tool is registered before any tool conversion pass - AdvisorOrchestrationHandler resolves default_advisor_model from litellm.advisor_interception_params when the tool definition omits it - Sub-calls route through llm_router.acompletion() for proper credential resolution across any provider - Native Anthropic path preserved: if executor is Anthropic and advisor is claude-opus-4-6 the request passes through unchanged - Any other combination (Anthropic executor + non-Opus advisor, or non-Anthropic executor) uses the LiteLLM orchestration loop - Final response always includes server_tool_use and advisor_tool_result content blocks matching Anthropic native format - FakeAnthropicMessagesStreamIterator emits correct SSE events for those new block types in streaming responses - Normalize proxy aliases → canonical Anthropic model IDs before native passthrough so tools.0.model is never a proxy alias Made-with: Cursor --- .../messages/fake_stream_iterator.py | 39 +++ .../messages/handler.py | 44 +-- .../messages/interceptors/advisor.py | 260 +++++++++++++++++- .../messages/transformation.py | 53 ++++ .../responses_adapters/transformation.py | 16 +- 5 files changed, 379 insertions(+), 33 deletions(-) diff --git a/litellm/llms/anthropic/experimental_pass_through/messages/fake_stream_iterator.py b/litellm/llms/anthropic/experimental_pass_through/messages/fake_stream_iterator.py index 7fc9b00f2c7..cc481201685 100644 --- a/litellm/llms/anthropic/experimental_pass_through/messages/fake_stream_iterator.py +++ b/litellm/llms/anthropic/experimental_pass_through/messages/fake_stream_iterator.py @@ -128,6 +128,45 @@ class FakeAnthropicMessagesStreamIterator: f"event: content_block_delta\ndata: {json.dumps(content_block_delta)}\n\n".encode() ) + elif block_type == "server_tool_use": + content_block_start = { + "type": "content_block_start", + "index": index, + "content_block": { + "type": "server_tool_use", + "id": block_dict.get("id"), + "name": block_dict.get("name"), + }, + } + chunks.append( + f"event: content_block_start\ndata: {json.dumps(content_block_start)}\n\n".encode() + ) + + elif block_type == "advisor_tool_result": + content_block_start = { + "type": "content_block_start", + "index": index, + "content_block": { + "type": "advisor_tool_result", + "tool_use_id": block_dict.get("tool_use_id"), + "content": {"type": "advisor_result", "text": ""}, + }, + } + chunks.append( + f"event: content_block_start\ndata: {json.dumps(content_block_start)}\n\n".encode() + ) + advisor_content = block_dict.get("content") or {} + advisor_text = advisor_content.get("text", "") if isinstance(advisor_content, dict) else "" + if advisor_text: + content_block_delta = { + "type": "content_block_delta", + "index": index, + "delta": {"type": "advisor_result_delta", "text": advisor_text}, + } + chunks.append( + f"event: content_block_delta\ndata: {json.dumps(content_block_delta)}\n\n".encode() + ) + content_block_stop = {"type": "content_block_stop", "index": index} chunks.append( f"event: content_block_stop\ndata: {json.dumps(content_block_stop)}\n\n".encode() diff --git a/litellm/llms/anthropic/experimental_pass_through/messages/handler.py b/litellm/llms/anthropic/experimental_pass_through/messages/handler.py index c400d82b7cf..6ede735e392 100644 --- a/litellm/llms/anthropic/experimental_pass_through/messages/handler.py +++ b/litellm/llms/anthropic/experimental_pass_through/messages/handler.py @@ -192,6 +192,31 @@ async def anthropic_messages( "_websearch_interception_converted_stream", False ) + # Resolve custom_llm_provider early so interceptors can use it. + if not custom_llm_provider: + try: + _, custom_llm_provider, _, _ = litellm.get_llm_provider(model=model) + except Exception: + pass + + # Run registered MessagesInterceptors (e.g. advisor orchestration loop) + # BEFORE pre-request hooks, because hooks like AdvisorInterceptionLogger + # would convert the advisor_20260301 tool into an OpenAI function tool, + # making it invisible to the Messages API interceptor. + for interceptor in get_messages_interceptors(): + if interceptor.can_handle(tools, custom_llm_provider): + return await interceptor.handle( + model=model, + messages=messages, + tools=tools, + stream=original_stream, + max_tokens=max_tokens, + custom_llm_provider=custom_llm_provider, + api_key=api_key, + api_base=api_base, + **kwargs, + ) + # Execute pre-request hooks to allow CustomLoggers to modify request request_kwargs = await _execute_pre_request_hooks( model=model, @@ -206,8 +231,6 @@ async def anthropic_messages( tools = request_kwargs.pop("tools", tools) stream = request_kwargs.pop("stream", stream) # Propagate the provider derived inside pre-request hooks, if not already set. - # The litellm_params dict may have been overwritten by **kwargs in - # _execute_pre_request_hooks, so fall back to get_llm_provider() if needed. if not custom_llm_provider: custom_llm_provider = request_kwargs.get("litellm_params", {}).get( "custom_llm_provider" @@ -237,23 +260,6 @@ async def anthropic_messages( if short_circuit_response is not None: return short_circuit_response - # Run registered MessagesInterceptors (e.g. advisor orchestration loop). - # api_key and api_base are explicit params (not in **kwargs) so pass them - # explicitly so interceptor sub-calls can route to the same backend. - for interceptor in get_messages_interceptors(): - if interceptor.can_handle(tools, custom_llm_provider): - return await interceptor.handle( - model=model, - messages=messages, - tools=tools, - stream=original_stream, - max_tokens=max_tokens, - custom_llm_provider=custom_llm_provider, - api_key=api_key, - api_base=api_base, - **kwargs, - ) - loop = asyncio.get_event_loop() kwargs["is_async"] = True diff --git a/litellm/llms/anthropic/experimental_pass_through/messages/interceptors/advisor.py b/litellm/llms/anthropic/experimental_pass_through/messages/interceptors/advisor.py index 02437c9b63f..968e598a7fa 100644 --- a/litellm/llms/anthropic/experimental_pass_through/messages/interceptors/advisor.py +++ b/litellm/llms/anthropic/experimental_pass_through/messages/interceptors/advisor.py @@ -18,6 +18,7 @@ import uuid from typing import Any, AsyncIterator, Dict, List, Optional, Union import litellm.constants as _c +from litellm._logging import verbose_logger from litellm.llms.anthropic.common_utils import strip_advisor_blocks_from_messages from litellm.types.llms.anthropic_messages.anthropic_response import ( AnthropicMessagesResponse, @@ -25,7 +26,6 @@ from litellm.types.llms.anthropic_messages.anthropic_response import ( from litellm.types.llms.anthropic import ANTHROPIC_ADVISOR_TOOL_TYPE ADVISOR_MAX_USES: int = _c.ADVISOR_MAX_USES -ADVISOR_NATIVE_PROVIDERS: frozenset = _c.ADVISOR_NATIVE_PROVIDERS ADVISOR_TOOL_DESCRIPTION: str = _c.ADVISOR_TOOL_DESCRIPTION from .base import MessagesInterceptor @@ -36,7 +36,7 @@ class AdvisorMaxIterationsError(Exception): class AdvisorOrchestrationHandler(MessagesInterceptor): - """Orchestrates the advisor tool loop for non-native providers.""" + """Orchestrates the advisor tool loop for /v1/messages requests.""" def can_handle( self, @@ -46,8 +46,12 @@ class AdvisorOrchestrationHandler(MessagesInterceptor): if not tools: return False has_advisor = any(t.get("type") == ANTHROPIC_ADVISOR_TOOL_TYPE for t in tools) - is_non_native = custom_llm_provider not in ADVISOR_NATIVE_PROVIDERS - return has_advisor and is_non_native + if not has_advisor: + return False + # Keep Anthropic-native advisor behavior for Claude Opus 4.6. + if _should_use_native_anthropic_advisor(tools, custom_llm_provider): + return False + return True async def handle( self, @@ -74,9 +78,13 @@ class AdvisorOrchestrationHandler(MessagesInterceptor): f"handle() called but no {ANTHROPIC_ADVISOR_TOOL_TYPE} tool found in tools list" ) advisor_model: str = advisor_tool.get("model") or "" + if not advisor_model: + advisor_model = _resolve_default_advisor_model() if not advisor_model: raise ValueError( - "advisor tool definition must include a 'model' field specifying the advisor model" + "No advisor model specified. Either:\n" + " 1. Set 'default_advisor_model' in advisor_interception_params in your proxy config YAML, or\n" + " 2. Include a 'model' field in the advisor tool definition." ) _raw_max_uses = advisor_tool.get("max_uses") max_uses: int = ADVISOR_MAX_USES if _raw_max_uses is None else int(_raw_max_uses) @@ -108,6 +116,7 @@ class AdvisorOrchestrationHandler(MessagesInterceptor): ) metadata_base: Dict = dict(kwargs.pop("metadata", None) or {}) iteration = 0 + advisor_interactions: List[Dict] = [] while True: # --- Executor call (always non-streaming) --- @@ -130,6 +139,10 @@ class AdvisorOrchestrationHandler(MessagesInterceptor): if advisor_use_block is None: # No more advisor calls — this is the final response. + # Inject advisor_tool_result blocks to match Anthropic native format. + _inject_advisor_blocks_into_response( + executor_response, advisor_interactions + ) if stream: return FakeAnthropicMessagesStreamIterator(executor_response) return executor_response @@ -147,13 +160,10 @@ class AdvisorOrchestrationHandler(MessagesInterceptor): ) # --- Advisor sub-call (always non-streaming, no tools) --- - advisor_response: AnthropicMessagesResponse = await _call_messages_handler( + advisor_response: AnthropicMessagesResponse = await _call_advisor_with_router( model=advisor_model, messages=advisor_messages, - tools=None, - stream=False, max_tokens=max_tokens, - custom_llm_provider=None, # let litellm resolve from model name metadata={ **metadata_base, "advisor_sub_call": True, @@ -165,6 +175,12 @@ class AdvisorOrchestrationHandler(MessagesInterceptor): advisor_text = _extract_response_text(advisor_response) + # Record the interaction for later injection into the final response. + advisor_interactions.append({ + "tool_use_id": advisor_use_block.get("id", f"srvtoolu_{uuid.uuid4().hex[:24]}"), + "advisor_text": advisor_text, + }) + # --- Inject advisor result and continue loop --- current_messages = _inject_advisor_turn( current_messages, @@ -179,10 +195,99 @@ class AdvisorOrchestrationHandler(MessagesInterceptor): # --------------------------------------------------------------------------- +def _resolve_default_advisor_model() -> str: + """Resolve the default advisor model from proxy config / litellm settings.""" + import litellm + + params = getattr(litellm, "advisor_interception_params", None) or {} + return params.get("default_advisor_model", "") or "" + + +def _is_anthropic_opus_46_model(model: str) -> bool: + """Return True for Anthropic Claude Opus 4.6 model identifiers.""" + normalized = model.lower().replace("_", "-") + return "anthropic/" in normalized and "claude-opus-4-6" in normalized + + +def _resolve_proxy_model_alias_to_litellm_model(model: str) -> str: + """ + Resolve a proxy ``model_name`` alias to its configured ``litellm_params.model``. + + Example: ``claude_opus`` -> ``anthropic/claude-opus-4-6``. + """ + try: + from litellm.proxy.proxy_server import llm_router + except Exception: + return "" + + model_list = getattr(llm_router, "model_list", None) or [] + for deployment in model_list: + if not isinstance(deployment, dict): + continue + if deployment.get("model_name") != model: + continue + litellm_params = deployment.get("litellm_params") or {} + configured_model = litellm_params.get("model") + if isinstance(configured_model, str): + return configured_model + return "" + + +def _should_use_native_anthropic_advisor( + tools: List[Dict], custom_llm_provider: Optional[str] +) -> bool: + """ + Use Anthropic's native advisor path only when: + - executor provider is Anthropic, and + - advisor model resolves to Anthropic Claude Opus 4.6. + """ + if custom_llm_provider != "anthropic": + return False + + advisor_tool = next( + (t for t in tools if t.get("type") == ANTHROPIC_ADVISOR_TOOL_TYPE), + None, + ) + if advisor_tool is None: + return False + + advisor_model = (advisor_tool.get("model") or _resolve_default_advisor_model() or "").strip() + if not advisor_model: + return False + + if _is_anthropic_opus_46_model(advisor_model): + return True + + # Proxy requests commonly pass advisor model as a model_name alias. + resolved_proxy_model = _resolve_proxy_model_alias_to_litellm_model(advisor_model) + if _is_anthropic_opus_46_model(resolved_proxy_model): + return True + + try: + import litellm + + resolved_model, advisor_provider, _, _ = litellm.get_llm_provider( + model=advisor_model + ) + return advisor_provider == "anthropic" and _is_anthropic_opus_46_model( + resolved_model + ) + except Exception: + return False + + +_SYNTHETIC_ADVISOR_TOOL_NAME = "consult_advisor" + + def _make_synthetic_advisor_tool() -> Dict: - """Build a regular tool definition the executor provider can understand.""" + """Build a regular tool definition the executor provider can understand. + + Uses a name that does NOT collide with ``_ADVISOR_TOOL_NAMES`` in the + chat-completions interception handler so pre-request hooks won't + double-convert it. + """ return { - "name": "advisor", + "name": _SYNTHETIC_ADVISOR_TOOL_NAME, "description": ADVISOR_TOOL_DESCRIPTION, "input_schema": { "type": "object", @@ -198,7 +303,7 @@ def _make_synthetic_advisor_tool() -> Dict: def _find_advisor_tool_use(response: Any) -> Optional[Dict]: - """Return the first tool_use block with name='advisor', or None.""" + """Return the first tool_use block whose name matches our synthetic advisor.""" content = response.get("content") if isinstance(response, dict) else [] if not isinstance(content, list): return None @@ -206,12 +311,79 @@ def _find_advisor_tool_use(response: Any) -> Optional[Dict]: if ( isinstance(block, dict) and block.get("type") == "tool_use" - and block.get("name") == "advisor" + and block.get("name") == _SYNTHETIC_ADVISOR_TOOL_NAME ): return block return None +def _openai_response_to_anthropic_dict(response: Any) -> Dict: + """Convert an OpenAI ChatCompletion response to a minimal Anthropic Messages dict. + + Only the fields used by ``_extract_response_text`` are needed. + """ + try: + choices = response.choices if hasattr(response, "choices") else [] + content_blocks: List[Dict] = [] + for choice in choices: + msg = choice.message if hasattr(choice, "message") else choice.get("message", {}) + text = msg.content if hasattr(msg, "content") else msg.get("content", "") + if text: + content_blocks.append({"type": "text", "text": text}) + return { + "id": getattr(response, "id", ""), + "type": "message", + "role": "assistant", + "content": content_blocks, + "stop_reason": "end_turn", + } + except Exception: + return {"content": [], "stop_reason": "end_turn"} + + +def _inject_advisor_blocks_into_response( + response: Any, advisor_interactions: List[Dict] +) -> None: + """ + Mutate *response* in place so its ``content`` array includes + ``advisor_tool_result`` blocks that mirror Anthropic's native advisor + response format. + + Each advisor interaction produces two blocks appended after existing + content: + + * ``server_tool_use`` – records the executor's call to the advisor + * ``advisor_tool_result`` – carries the advisor's answer + + This ensures callers see an identical structure regardless of whether the + advisor ran natively or via LiteLLM interception. + """ + if not advisor_interactions: + return + + content = response.get("content") if isinstance(response, dict) else None + if not isinstance(content, list): + return + + for interaction in advisor_interactions: + tool_use_id = interaction["tool_use_id"] + advisor_text = interaction["advisor_text"] + + content.append({ + "type": "server_tool_use", + "id": tool_use_id, + "name": "advisor", + }) + content.append({ + "type": "advisor_tool_result", + "tool_use_id": tool_use_id, + "content": { + "type": "advisor_result", + "text": advisor_text, + }, + }) + + def _extract_response_text(response: Any) -> str: """Extract concatenated text from all text blocks in a response.""" content = response.get("content") if isinstance(response, dict) else [] @@ -349,3 +521,65 @@ async def _call_messages_handler( custom_llm_provider=custom_llm_provider, **kwargs, ) + + +async def _call_advisor_with_router( + model: str, + messages: List[Dict], + max_tokens: int, + metadata: Optional[Dict] = None, + api_key: Optional[str] = None, + api_base: Optional[str] = None, +) -> Any: + """ + Call the advisor model via ``llm_router.acompletion()`` (proxy) or + ``litellm.acompletion()`` (SDK-only). + + Returns a dict in Anthropic Messages format so the orchestration loop + can process it uniformly. + """ + import litellm as _litellm + + llm_router = None + try: + from litellm.proxy.proxy_server import llm_router as _router + + llm_router = _router + except ImportError: + pass + + kwargs: Dict[str, Any] = {} + if metadata is not None: + kwargs["metadata"] = metadata + if api_key is not None: + kwargs["api_key"] = api_key + if api_base is not None: + kwargs["api_base"] = api_base + + openai_response = None + if llm_router is not None: + try: + openai_response = await llm_router.acompletion( + model=model, + messages=messages, + tools=None, + max_tokens=max_tokens, + **kwargs, + ) + except Exception: + verbose_logger.debug( + "AdvisorOrchestration: Router call for advisor model '%s' failed, " + "falling back to direct litellm.acompletion()", + model, + ) + + if openai_response is None: + openai_response = await _litellm.acompletion( + model=model, + messages=messages, + tools=None, + max_tokens=max_tokens, + **kwargs, + ) + + return _openai_response_to_anthropic_dict(openai_response) diff --git a/litellm/llms/anthropic/experimental_pass_through/messages/transformation.py b/litellm/llms/anthropic/experimental_pass_through/messages/transformation.py index 46af1f7fbd1..06d1cab0750 100644 --- a/litellm/llms/anthropic/experimental_pass_through/messages/transformation.py +++ b/litellm/llms/anthropic/experimental_pass_through/messages/transformation.py @@ -28,6 +28,55 @@ from ...common_utils import ( DEFAULT_ANTHROPIC_API_VERSION = "2023-06-01" +def _resolve_proxy_model_alias_to_litellm_model(model: str) -> str: + """Resolve proxy model_name alias to configured litellm model string.""" + try: + from litellm.proxy.proxy_server import llm_router + except Exception: + return "" + + model_list = getattr(llm_router, "model_list", None) or [] + for deployment in model_list: + if not isinstance(deployment, dict): + continue + if deployment.get("model_name") != model: + continue + litellm_params = deployment.get("litellm_params") or {} + configured_model = litellm_params.get("model") + if isinstance(configured_model, str): + return configured_model + return "" + + +def _normalize_anthropic_advisor_tool_models(tools: List[Dict]) -> List[Dict]: + """ + Normalize advisor tool model names for Anthropic native /v1/messages calls. + + Anthropic expects advisor tool model values like ``claude-opus-4-6``. + Proxy alias names (e.g. ``claude_opus``) and provider-prefixed values + (e.g. ``anthropic/claude-opus-4-6``) are converted. + """ + normalized_tools: List[Dict] = [] + for tool in tools: + if not isinstance(tool, dict): + normalized_tools.append(tool) + continue + if tool.get("type") != ANTHROPIC_ADVISOR_TOOL_TYPE: + normalized_tools.append(tool) + continue + + updated_tool = dict(tool) + advisor_model = updated_tool.get("model") + if isinstance(advisor_model, str) and advisor_model.strip(): + resolved = _resolve_proxy_model_alias_to_litellm_model(advisor_model.strip()) + canonical_model = resolved or advisor_model.strip() + if canonical_model.startswith("anthropic/"): + canonical_model = canonical_model.split("/", 1)[1] + updated_tool["model"] = canonical_model + normalized_tools.append(updated_tool) + return normalized_tools + + class AnthropicMessagesConfig(BaseAnthropicMessagesConfig): def get_supported_anthropic_messages_params(self, model: str) -> list: return [ @@ -220,6 +269,10 @@ class AnthropicMessagesConfig(BaseAnthropicMessagesConfig): # Auto-strip advisor blocks from history if advisor tool is absent. # Prevents Anthropic 400: advisor_tool_result in history requires advisor tool. _tools = anthropic_messages_optional_request_params.get("tools") or [] + if _tools: + normalized_tools = _normalize_anthropic_advisor_tool_models(_tools) + anthropic_messages_optional_request_params["tools"] = normalized_tools + _tools = normalized_tools _has_advisor = any( isinstance(t, dict) and t.get("type") == ANTHROPIC_ADVISOR_TOOL_TYPE for t in _tools diff --git a/litellm/llms/anthropic/experimental_pass_through/responses_adapters/transformation.py b/litellm/llms/anthropic/experimental_pass_through/responses_adapters/transformation.py index dae7044a5bc..02e6064fd7f 100644 --- a/litellm/llms/anthropic/experimental_pass_through/responses_adapters/transformation.py +++ b/litellm/llms/anthropic/experimental_pass_through/responses_adapters/transformation.py @@ -198,7 +198,21 @@ class LiteLLMAnthropicToResponsesAPIAdapter: ) or tool_name == "web_search": result.append({"type": "web_search_preview"}) continue - func_tool: Dict[str, Any] = {"type": "function", "name": tool_name} + # Handle OpenAI chat completions format tools that may arrive here + # (e.g. {"type": "function", "function": {"name": ..., "parameters": ...}}) + if tool_type == "function" and "function" in tool_dict: + func_def = tool_dict["function"] + func_tool: Dict[str, Any] = { + "type": "function", + "name": func_def.get("name", ""), + } + if "description" in func_def: + func_tool["description"] = func_def["description"] + if "parameters" in func_def: + func_tool["parameters"] = func_def["parameters"] + result.append(func_tool) + continue + func_tool = {"type": "function", "name": tool_name} if "description" in tool_dict: func_tool["description"] = tool_dict["description"] if "input_schema" in tool_dict: