feat(advisor): aggregate cost + usage iterations for advisor calls

- Emit per-iteration usage (executor + advisor sub-calls) under
  usage.iterations[] for both chat completions and /v1/messages advisor flows.
- Publish advisor cost split (Main Model initial + Advisor Model) to the
  logging obj so the proxy UI renders it the same way as Azure Router.
- For /v1/messages orchestration, share the outer litellm_logging_obj with
  inner sub-calls (so provider/api_base/model_id get populated) while
  suppressing duplicate log rows via _is_litellm_internal_call.
- Restore call_type to anthropic_messages at end of orchestration so the
  final response is transformed to ModelResponse and the UI shows output.
- Preserve pre-computed response_cost on dict results in
  _response_cost_calculator to prevent cost_breakdown overwrites.

Made-with: Cursor
This commit is contained in:
Sameer Kankute 2026-04-17 17:21:50 +05:30
parent 1f7773c805
commit 6f022ef951
No known key found for this signature in database
8 changed files with 823 additions and 17 deletions

View file

@ -341,6 +341,12 @@ class AdvisorInterceptionLogger(CustomLogger):
current_response = response
advisor_uses = 0
total_response_cost = self._safe_get_response_cost(current_response)
advisor_subcall_cost: float = 0.0
iterations: List[Dict[str, Any]] = [
self._build_iteration_entry(
response=current_response, iteration_type="message"
)
]
advisor_interactions: List[Dict[str, str]] = []
try:
@ -349,12 +355,27 @@ class AdvisorInterceptionLogger(CustomLogger):
current_response
)
if not advisor_calls:
final_executor_cost = self._safe_get_response_cost(current_response)
advisor_first_call_cost = max(
total_response_cost - advisor_subcall_cost - final_executor_cost,
0.0,
)
self._set_response_cost_if_possible(
response=current_response, response_cost=total_response_cost
)
self._store_advisor_cost_breakdown(
logging_obj=kwargs.get("litellm_logging_obj") or logging_obj,
final_executor_cost=final_executor_cost,
advisor_first_call_cost=advisor_first_call_cost,
advisor_subcall_cost=advisor_subcall_cost,
total_response_cost=total_response_cost,
)
self._inject_advisor_results_into_response(
current_response, advisor_interactions
)
self._inject_advisor_iterations_into_response(
current_response, iterations
)
return current_response
if len(advisor_calls) != len(raw_tool_calls):
verbose_logger.debug(
@ -395,7 +416,16 @@ class AdvisorInterceptionLogger(CustomLogger):
api_key=advisor_api_key,
api_base=advisor_api_base,
)
total_response_cost += self._safe_get_response_cost(advisor_response)
advisor_call_cost = self._safe_get_response_cost(advisor_response)
total_response_cost += advisor_call_cost
advisor_subcall_cost += advisor_call_cost
iterations.append(
self._build_iteration_entry(
response=advisor_response,
iteration_type="advisor_message",
model=advisor_model,
)
)
advisor_text = self._extract_text_content(advisor_response)
advisor_interactions.append({
"tool_use_id": advisor_call["id"],
@ -435,6 +465,11 @@ class AdvisorInterceptionLogger(CustomLogger):
kwargs_for_followup=kwargs_for_followup,
)
total_response_cost += self._safe_get_response_cost(current_response)
iterations.append(
self._build_iteration_entry(
response=current_response, iteration_type="message"
)
)
finally:
if isinstance(call_id, str):
self._advisor_config_by_call_id.pop(call_id, None)
@ -895,6 +930,146 @@ class AdvisorInterceptionLogger(CustomLogger):
return 0.0
return 0.0
@staticmethod
def _build_iteration_entry(
response: Any,
iteration_type: str,
model: Optional[str] = None,
) -> Dict[str, Any]:
"""
Build one entry for ``iterations[]`` on advisor-orchestrated responses.
Mirrors the Anthropic usage shape: ``input_tokens``, ``output_tokens``,
``cache_read_input_tokens``, ``cache_creation_input_tokens``. For advisor
sub-calls, also includes the advisor model name.
"""
input_tokens = 0
output_tokens = 0
cache_read_input_tokens = 0
cache_creation_input_tokens = 0
usage: Any = None
if isinstance(response, dict):
usage = response.get("usage")
else:
usage = getattr(response, "usage", None)
if usage is not None:
input_tokens = (
AdvisorInterceptionLogger._get_usage_value(usage, "prompt_tokens")
or 0
)
output_tokens = (
AdvisorInterceptionLogger._get_usage_value(usage, "completion_tokens")
or 0
)
cache_read_input_tokens = (
AdvisorInterceptionLogger._get_usage_value(
usage, "cache_read_input_tokens"
)
or 0
)
cache_creation_input_tokens = (
AdvisorInterceptionLogger._get_usage_value(
usage, "cache_creation_input_tokens"
)
or 0
)
if not cache_read_input_tokens:
prompt_tokens_details = AdvisorInterceptionLogger._get_usage_value(
usage, "prompt_tokens_details"
)
if prompt_tokens_details is not None:
cache_read_input_tokens = (
AdvisorInterceptionLogger._get_usage_value(
prompt_tokens_details, "cached_tokens"
)
or 0
)
entry: Dict[str, Any] = {
"type": iteration_type,
"input_tokens": int(input_tokens or 0),
"cache_read_input_tokens": int(cache_read_input_tokens or 0),
"cache_creation_input_tokens": int(cache_creation_input_tokens or 0),
"output_tokens": int(output_tokens or 0),
}
if iteration_type == "advisor_message" and model is not None:
entry["model"] = model
return entry
@staticmethod
def _get_usage_value(usage_obj: Any, key: str) -> Any:
if usage_obj is None:
return None
if isinstance(usage_obj, dict):
return usage_obj.get(key)
return getattr(usage_obj, key, None)
@staticmethod
def _store_advisor_cost_breakdown(
logging_obj: Any,
final_executor_cost: float,
advisor_first_call_cost: float,
advisor_subcall_cost: float,
total_response_cost: float,
) -> None:
"""
Populate ``cost_breakdown.additional_costs`` on the logging object so
the proxy UI can display the advisor cost split (mirrors Azure Router).
"""
if logging_obj is None:
return
if not hasattr(logging_obj, "set_cost_breakdown"):
return
additional_costs: Dict[str, float] = {}
if advisor_first_call_cost > 0:
additional_costs["Main Model (initial)"] = advisor_first_call_cost
if advisor_subcall_cost > 0:
additional_costs["Advisor Model"] = advisor_subcall_cost
try:
logging_obj.set_cost_breakdown(
input_cost=final_executor_cost,
output_cost=0.0,
total_cost=total_response_cost,
cost_for_built_in_tools_cost_usd_dollar=0.0,
additional_costs=additional_costs or None,
)
except Exception as breakdown_error:
verbose_logger.debug(
"AdvisorInterception: failed to store cost breakdown: %s",
str(breakdown_error),
)
@staticmethod
def _inject_advisor_iterations_into_response(
response: Any, iterations: List[Dict[str, Any]]
) -> None:
"""
Attach the per-iteration usage breakdown under
``message.provider_specific_fields["advisor_iterations"]`` so callers
can inspect the full orchestration without breaking the OpenAI-compatible
``usage`` shape.
"""
if not iterations:
return
message = AdvisorInterceptionLogger._extract_first_choice_message_obj(response)
if message is None:
return
existing_psf = getattr(message, "provider_specific_fields", None) or {}
existing_psf["advisor_iterations"] = iterations
try:
message.provider_specific_fields = existing_psf
except Exception:
try:
setattr(message, "provider_specific_fields", existing_psf)
except Exception:
pass
@staticmethod
def _set_response_cost_if_possible(response: Any, response_cost: float) -> None:
if response is None:

View file

@ -1489,6 +1489,20 @@ class Logging(LiteLLMLoggingBaseClass):
router_model_id is None and "model_id" in hidden_params
): # use model_id if not already set
router_model_id = hidden_params["model_id"]
elif isinstance(result, dict):
# Dict-shaped responses (e.g. AnthropicMessagesResponse returned by
# ``/v1/messages``) can carry a pre-computed ``response_cost`` on
# ``_hidden_params`` — most commonly when an orchestrator (advisor
# tool loop) has already aggregated cost across multiple sub-calls.
# Returning it early avoids ``litellm.response_cost_calculator``
# recomputing from per-request usage and overwriting the
# ``cost_breakdown`` the orchestrator set via ``set_cost_breakdown``.
dict_hidden_params = result.get("_hidden_params") or {}
if (
isinstance(dict_hidden_params, dict)
and dict_hidden_params.get("response_cost") is not None
):
return dict_hidden_params["response_cost"]
# Fallback: extract router_model_id from litellm_params when not available
# from the result object. ResponsesAPIResponse objects (used by /v1/responses
@ -3445,13 +3459,19 @@ class Logging(LiteLLMLoggingBaseClass):
litellm_params={},
)
else:
from litellm.types.llms.anthropic import AnthropicResponse
pydantic_result = AnthropicResponse.model_validate(result)
import httpx
# NOTE: we intentionally pass the raw dict directly to
# ``transform_parsed_response`` instead of going through
# ``AnthropicResponse.model_validate(...).model_dump()``. The
# strict schema only allows ``text``/``tool_use``/``thinking``/
# ``redacted_thinking`` content blocks, which breaks logging for
# any response that carries provider-specific blocks like
# ``server_tool_use``, ``*_tool_result`` (web search, code
# execution, advisor tool, etc.). ``transform_parsed_response``
# and its ``extract_response_content`` helper already handle
# these natively, so validation adds no value here.
completion_response = result if isinstance(result, dict) else {}
result = litellm.AnthropicConfig().transform_parsed_response(
completion_response=pydantic_result.model_dump(),
completion_response=completion_response,
raw_response=httpx.Response(
status_code=200,
headers={},

View file

@ -14,14 +14,17 @@ How it works:
6. Wraps in FakeAnthropicMessagesStreamIterator if the caller requested streaming.
"""
import asyncio
import uuid
from typing import Any, AsyncIterator, Dict, List, Optional, Union, cast
import litellm
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,
AnthropicUsageIteration,
)
from litellm.types.llms.openai import AllMessageValues
from litellm.types.llms.anthropic import ANTHROPIC_ADVISOR_TOOL_TYPE
@ -117,8 +120,24 @@ class AdvisorOrchestrationHandler(MessagesInterceptor):
kwargs.pop("litellm_call_id", None) or uuid.uuid4()
)
metadata_base: Dict = dict(kwargs.pop("metadata", None) or {})
# Hold the outer ``anthropic_messages`` logging obj — we attach the
# aggregated cost/breakdown to it in _finalize_orchestrated_response.
#
# We *keep* this object in kwargs so inner executor sub-calls share it:
# their inner @client wrappers (aresponses/acompletion) still populate
# ``custom_llm_provider``, ``api_base``, ``model_id`` on this shared
# model_call_details via ``update_environment_variables`` — fields the
# outer anthropic_messages path never sets on its own. We mark the
# sub-calls ``_is_litellm_internal_call=True`` so their @client skips
# emitting a duplicate log row; only the outer call emits a single
# aggregated entry.
litellm_logging_obj = kwargs.get("litellm_logging_obj", None)
kwargs["_is_litellm_internal_call"] = True
iteration = 0
advisor_interactions: List[Dict] = []
iterations: List[AnthropicUsageIteration] = []
advisor_first_call_cost: float = 0.0
advisor_subcall_cost: float = 0.0
while True:
# --- Executor call (always non-streaming) ---
@ -137,6 +156,13 @@ class AdvisorOrchestrationHandler(MessagesInterceptor):
**kwargs,
)
executor_cost = _get_response_cost(executor_response, model=model)
iterations.append(
_build_iteration_entry(
response=executor_response, iteration_type="message"
)
)
advisor_use_block = _find_advisor_tool_use(executor_response)
if advisor_use_block is None:
@ -145,10 +171,27 @@ class AdvisorOrchestrationHandler(MessagesInterceptor):
_inject_advisor_blocks_into_response(
executor_response, advisor_interactions
)
total_cost = (
advisor_first_call_cost + advisor_subcall_cost + executor_cost
)
_finalize_orchestrated_response(
response=executor_response,
iterations=iterations,
total_cost=total_cost,
final_executor_cost=executor_cost,
advisor_first_call_cost=advisor_first_call_cost,
advisor_subcall_cost=advisor_subcall_cost,
litellm_logging_obj=litellm_logging_obj,
)
if stream:
return FakeAnthropicMessagesStreamIterator(executor_response)
return executor_response
# Executor response triggered another advisor call → count it as a
# "first/intermediate" executor turn. Only the terminating turn is
# treated as the base response.
advisor_first_call_cost += executor_cost
iteration += 1
if iteration > max_uses:
raise AdvisorMaxIterationsError(
@ -175,6 +218,16 @@ class AdvisorOrchestrationHandler(MessagesInterceptor):
api_base=advisor_api_base,
)
advisor_call_cost = _get_response_cost(advisor_response, model=advisor_model)
advisor_subcall_cost += advisor_call_cost
iterations.append(
_build_iteration_entry(
response=advisor_response,
iteration_type="advisor_message",
model=advisor_model,
)
)
advisor_text = _extract_response_text(advisor_response)
# Record the interaction for later injection into the final response.
@ -231,6 +284,165 @@ def _make_synthetic_advisor_tool() -> Dict:
}
def _get_response_cost(response: Any, model: Optional[str] = None) -> float:
"""
Extract the cost of a single sub-call response.
Prefers a pre-set ``_hidden_params["response_cost"]`` (produced by the
``@client`` wrapper on BaseModel responses / preserved in
:func:`_openai_response_to_anthropic_dict`). Falls back to
``litellm.completion_cost`` when the hidden param is missing — this covers
executor responses from ``anthropic_messages`` which return a plain dict.
"""
if response is None:
return 0.0
hidden_params: Any = None
if isinstance(response, dict):
hidden_params = response.get("_hidden_params")
else:
hidden_params = getattr(response, "_hidden_params", None)
if isinstance(hidden_params, dict):
cost = hidden_params.get("response_cost")
if isinstance(cost, (int, float)):
return float(cost)
try:
cost = litellm.completion_cost(completion_response=response, model=model)
if isinstance(cost, (int, float)):
return float(cost)
except Exception as cost_error:
verbose_logger.debug(
"AdvisorOrchestration: completion_cost fallback failed for model '%s': %s",
model,
str(cost_error),
)
return 0.0
def _build_iteration_entry(
response: Any,
iteration_type: str,
model: Optional[str] = None,
) -> AnthropicUsageIteration:
"""
Build one entry for ``usage.iterations[]`` on advisor-orchestrated
``/v1/messages`` responses. Reads tokens from Anthropic-shaped usage on
the dict response.
"""
usage: Dict[str, Any] = {}
if isinstance(response, dict):
maybe_usage = response.get("usage")
if isinstance(maybe_usage, dict):
usage = maybe_usage
entry: AnthropicUsageIteration = {
"type": iteration_type, # type: ignore[typeddict-item]
"input_tokens": int(usage.get("input_tokens", 0) or 0),
"cache_read_input_tokens": int(usage.get("cache_read_input_tokens", 0) or 0),
"cache_creation_input_tokens": int(
usage.get("cache_creation_input_tokens", 0) or 0
),
"output_tokens": int(usage.get("output_tokens", 0) or 0),
}
if iteration_type == "advisor_message" and model is not None:
entry["model"] = model
return entry
def _finalize_orchestrated_response(
response: Any,
iterations: List[AnthropicUsageIteration],
total_cost: float,
final_executor_cost: float,
advisor_first_call_cost: float,
advisor_subcall_cost: float,
litellm_logging_obj: Any,
) -> None:
"""
Attach aggregated usage, per-iteration breakdown and total cost to the
terminating executor response, and publish the advisor cost split to the
parent ``litellm_logging_obj`` so the proxy UI can render it (same
plumbing Azure Router uses).
"""
if not isinstance(response, dict):
return
aggregated_usage: Dict[str, Any] = {
"input_tokens": sum(it.get("input_tokens", 0) for it in iterations),
"output_tokens": sum(it.get("output_tokens", 0) for it in iterations),
"cache_read_input_tokens": sum(
it.get("cache_read_input_tokens", 0) for it in iterations
),
"cache_creation_input_tokens": sum(
it.get("cache_creation_input_tokens", 0) for it in iterations
),
"iterations": list(iterations),
}
response["usage"] = aggregated_usage
# Give the orchestrated response its own Anthropic-style id so the outer
# ``anthropic_messages`` log row is distinct from any inner sub-call log
# row in the proxy UI (the inner final executor sub-call shares this
# ``litellm_logging_obj`` and would otherwise overwrite its request_id).
response["id"] = f"msg_{uuid.uuid4().hex}"
hidden_params = response.get("_hidden_params")
if not isinstance(hidden_params, dict):
hidden_params = {}
hidden_params["response_cost"] = total_cost
response["_hidden_params"] = hidden_params
if litellm_logging_obj is not None:
try:
# Ensure downstream logging emits the aggregated cost even when the
# transformed ModelResponse path (which drops our dict hidden_params)
# is taken.
litellm_logging_obj.model_call_details["response_cost"] = total_cost
except Exception:
pass
# Restore call_type to ``anthropic_messages`` so the outer success
# handler runs ``_handle_anthropic_messages_response_logging`` and
# converts the final Anthropic dict to a ``ModelResponse`` — without
# this, the proxy UI's log row has no renderable output (the raw
# Anthropic ``content`` blocks are not recognised by the OpenAI-shaped
# renderer). Inner executor sub-calls (e.g. aresponses for OpenAI
# models) flip call_type to ``acompletion`` to dodge a separate bug in
# their own success path; that leaves us with the wrong call_type on
# the shared logging obj by the time orchestration finishes.
try:
from litellm.types.utils import CallTypes
litellm_logging_obj.call_type = CallTypes.anthropic_messages.value
litellm_logging_obj.model_call_details[
"call_type"
] = CallTypes.anthropic_messages.value
except Exception:
pass
if hasattr(litellm_logging_obj, "set_cost_breakdown"):
additional_costs: Dict[str, float] = {}
if advisor_first_call_cost > 0:
additional_costs["Main Model (initial)"] = advisor_first_call_cost
if advisor_subcall_cost > 0:
additional_costs["Advisor Model"] = advisor_subcall_cost
try:
litellm_logging_obj.set_cost_breakdown(
input_cost=final_executor_cost,
output_cost=0.0,
total_cost=total_cost,
cost_for_built_in_tools_cost_usd_dollar=0.0,
additional_costs=additional_costs or None,
)
except Exception as breakdown_error:
verbose_logger.debug(
"AdvisorOrchestration: failed to store cost breakdown: %s",
str(breakdown_error),
)
def _find_advisor_tool_use(response: Any) -> Optional[Dict]:
"""Return the first tool_use block whose name matches our synthetic advisor."""
content = response.get("content") if isinstance(response, dict) else []
@ -249,7 +461,8 @@ def _find_advisor_tool_use(response: Any) -> Optional[Dict]:
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.
Preserves ``usage`` (mapped to Anthropic's shape) and ``_hidden_params`` so
the orchestration loop can aggregate cost and per-iteration token usage.
"""
try:
choices = response.choices if hasattr(response, "choices") else []
@ -259,17 +472,62 @@ def _openai_response_to_anthropic_dict(response: Any) -> Dict:
text = msg.content if hasattr(msg, "content") else msg.get("content", "")
if text:
content_blocks.append({"type": "text", "text": text})
return {
anthropic_dict: Dict[str, Any] = {
"id": getattr(response, "id", ""),
"type": "message",
"role": "assistant",
"content": content_blocks,
"stop_reason": "end_turn",
"usage": _openai_usage_to_anthropic_usage(response),
}
hidden_params = getattr(response, "_hidden_params", None)
if hidden_params is not None:
anthropic_dict["_hidden_params"] = (
dict(hidden_params) if isinstance(hidden_params, dict) else hidden_params
)
return anthropic_dict
except Exception:
return {"content": [], "stop_reason": "end_turn"}
def _openai_usage_to_anthropic_usage(response: Any) -> Dict[str, int]:
"""Translate an OpenAI usage block (if any) to Anthropic token shape."""
usage = getattr(response, "usage", None)
if usage is None and isinstance(response, dict):
usage = response.get("usage")
if usage is None:
return {
"input_tokens": 0,
"output_tokens": 0,
"cache_creation_input_tokens": 0,
"cache_read_input_tokens": 0,
}
def _get(obj: Any, key: str, default: int = 0) -> int:
if isinstance(obj, dict):
value = obj.get(key, default)
else:
value = getattr(obj, key, default)
return int(value or 0)
cache_read = _get(usage, "cache_read_input_tokens")
if not cache_read:
prompt_details = (
usage.get("prompt_tokens_details")
if isinstance(usage, dict)
else getattr(usage, "prompt_tokens_details", None)
)
if prompt_details is not None:
cache_read = _get(prompt_details, "cached_tokens")
return {
"input_tokens": _get(usage, "prompt_tokens"),
"output_tokens": _get(usage, "completion_tokens"),
"cache_creation_input_tokens": _get(usage, "cache_creation_input_tokens"),
"cache_read_input_tokens": cache_read,
}
def _inject_advisor_blocks_into_response(
response: Any, advisor_interactions: List[Dict]
) -> None:
@ -298,10 +556,17 @@ def _inject_advisor_blocks_into_response(
tool_use_id = interaction["tool_use_id"]
advisor_text = interaction["advisor_text"]
# ``input`` is required by
# ``AnthropicConfig.convert_tool_use_to_openai_format`` (it does
# ``json.dumps(block["input"])`` unconditionally when logging converts
# the Anthropic response to OpenAI format). We did not observe a real
# tool-call payload here — the advisor was invoked via a synthetic
# tool — so an empty object is the correct stub.
content.append({
"type": "server_tool_use",
"id": tool_use_id,
"name": "advisor",
"input": {},
})
content.append({
"type": "advisor_tool_result",
@ -431,17 +696,33 @@ async def _call_messages_handler(
**kwargs,
) -> Any:
"""
Call anthropic_messages() — the public async /messages entry point — for
orchestration sub-calls (executor or advisor).
Dispatch an orchestration sub-call (executor turn) by invoking the
inner ``anthropic_messages_handler`` directly, **bypassing** the @client
wrapper around ``anthropic_messages``.
Using the public function (decorated with @client) ensures logging, retries,
and provider resolution all work correctly, identical to a direct user call.
Why bypass @client here:
- ``anthropic_messages`` is @client-decorated, which ``kwargs.pop``s
``_is_litellm_internal_call`` before the body runs. Once popped, the
flag is gone from the kwargs forwarded to the inner
``litellm.aresponses()`` / ``litellm.acompletion()`` call, so *those*
inner @client wrappers emit their own log rows.
- For advisor orchestration we want exactly ONE log row — the outer
``anthropic_messages`` call the proxy received. Calling the inner
handler directly keeps ``_is_litellm_internal_call=True`` in kwargs
all the way down, so the inner aresponses/acompletion skip logging
and the outer call's @client emits the single aggregated entry.
The shared ``litellm_logging_obj`` is intentionally passed through in
kwargs so the inner aresponses/acompletion @client populates provider
metadata (``custom_llm_provider``, ``api_base``, ``model_id``) on the
outer log row.
"""
from litellm.llms.anthropic.experimental_pass_through.messages.handler import (
anthropic_messages,
anthropic_messages_handler,
)
return await anthropic_messages(
kwargs["is_async"] = True
result = anthropic_messages_handler(
model=model,
messages=messages,
tools=tools,
@ -450,6 +731,9 @@ async def _call_messages_handler(
custom_llm_provider=custom_llm_provider,
**kwargs,
)
if asyncio.iscoroutine(result):
return await result
return result
async def _call_advisor_with_router(
@ -477,7 +761,10 @@ async def _call_advisor_with_router(
except ImportError:
pass
kwargs: Dict[str, Any] = {}
# Mark as internal so the @client decorator on acompletion skips
# emitting a log row — only the outer anthropic_messages call should
# produce a single aggregated log entry for the advisor request.
kwargs: Dict[str, Any] = {"_is_litellm_internal_call": True}
if metadata is not None:
kwargs["metadata"] = metadata
if api_key is not None:

View file

@ -38567,6 +38567,7 @@
"supports_tool_choice": true,
"supports_vision": true,
"tool_use_system_prompt_tokens": 346,
"supports_native_advisor_tool": true,
"supports_native_structured_output": true,
"supports_pdf_input": true
},
@ -38590,6 +38591,7 @@
"supports_tool_choice": true,
"supports_vision": true,
"tool_use_system_prompt_tokens": 346,
"supports_native_advisor_tool": true,
"supports_native_structured_output": true,
"supports_pdf_input": true
}

View file

@ -56,6 +56,23 @@ AnthropicResponseContentBlock: TypeAlias = Union[
]
class AnthropicUsageIteration(TypedDict, total=False):
"""
Per-iteration token usage for advisor-orchestrated requests.
Emitted in ``AnthropicUsage.iterations`` when LiteLLM runs the advisor
orchestration loop so callers can see the breakdown of executor vs.
advisor sub-calls.
"""
type: Literal["message", "advisor_message"]
model: Optional[str]
input_tokens: int
output_tokens: int
cache_creation_input_tokens: int
cache_read_input_tokens: int
class AnthropicUsage(TypedDict, total=False):
"""
Input and output tokens used in the request
@ -70,6 +87,15 @@ class AnthropicUsage(TypedDict, total=False):
cache_creation_input_tokens: int
cache_read_input_tokens: int
"""
Per-iteration breakdown for advisor-orchestrated requests.
Populated by LiteLLM's advisor orchestration loop; absent for normal
(non-advisor) requests. Each entry describes a single executor or advisor
sub-call that contributed to the aggregated usage above.
"""
iterations: List[AnthropicUsageIteration]
class AnthropicMessagesResponse(TypedDict, total=False):
"""

View file

@ -38567,6 +38567,7 @@
"supports_tool_choice": true,
"supports_vision": true,
"tool_use_system_prompt_tokens": 346,
"supports_native_advisor_tool": true,
"supports_native_structured_output": true,
"supports_pdf_input": true
},
@ -38590,6 +38591,7 @@
"supports_tool_choice": true,
"supports_vision": true,
"tool_use_system_prompt_tokens": 346,
"supports_native_advisor_tool": true,
"supports_native_structured_output": true,
"supports_pdf_input": true
}

View file

@ -189,6 +189,11 @@ async def test_run_chat_completion_agentic_loop_aggregates_subcall_costs(monkeyp
created=123,
)
initial_response._hidden_params["response_cost"] = 1.0
initial_response.usage = { # type: ignore[attr-defined]
"prompt_tokens": 100,
"completion_tokens": 20,
"total_tokens": 120,
}
advisor_subcall_response = ModelResponse(
id="advisor-subcall",
@ -207,6 +212,11 @@ async def test_run_chat_completion_agentic_loop_aggregates_subcall_costs(monkeyp
created=124,
)
advisor_subcall_response._hidden_params["response_cost"] = 0.3
advisor_subcall_response.usage = { # type: ignore[attr-defined]
"prompt_tokens": 200,
"completion_tokens": 150,
"total_tokens": 350,
}
final_response = ModelResponse(
id="final",
@ -225,6 +235,11 @@ async def test_run_chat_completion_agentic_loop_aggregates_subcall_costs(monkeyp
created=125,
)
final_response._hidden_params["response_cost"] = 0.7
final_response.usage = { # type: ignore[attr-defined]
"prompt_tokens": 300,
"completion_tokens": 40,
"total_tokens": 340,
}
calls = {"count": 0}
@ -238,13 +253,16 @@ async def test_run_chat_completion_agentic_loop_aggregates_subcall_costs(monkeyp
monkeypatch.setattr(litellm, "acompletion", mock_acompletion)
fake_logging_obj = MagicMock()
fake_logging_obj.set_cost_breakdown = MagicMock()
response = await logger.async_run_chat_completion_agentic_loop(
tools={"advisor_config": {"advisor_model": "claude-opus-4-6", "max_uses": 3}},
model="gpt-4o-mini",
messages=[{"role": "user", "content": "Test"}],
response=initial_response,
optional_params={"tools": [get_litellm_advisor_tool_openai()], "max_tokens": 256},
logging_obj=None,
logging_obj=fake_logging_obj,
stream=False,
kwargs={"litellm_call_id": "cost-loop-1", "custom_llm_provider": "openai"},
)
@ -253,6 +271,29 @@ async def test_run_chat_completion_agentic_loop_aggregates_subcall_costs(monkeyp
assert response is final_response
assert response._hidden_params["response_cost"] == pytest.approx(2.0)
fake_logging_obj.set_cost_breakdown.assert_called_once()
breakdown_kwargs = fake_logging_obj.set_cost_breakdown.call_args.kwargs
assert breakdown_kwargs["total_cost"] == pytest.approx(2.0)
assert breakdown_kwargs["input_cost"] == pytest.approx(0.7)
assert breakdown_kwargs["additional_costs"] == {
"Main Model (initial)": pytest.approx(1.0),
"Advisor Model": pytest.approx(0.3),
}
message = response.choices[0].message
advisor_iterations = message.provider_specific_fields["advisor_iterations"]
assert len(advisor_iterations) == 3
assert advisor_iterations[0]["type"] == "message"
assert advisor_iterations[0]["input_tokens"] == 100
assert advisor_iterations[0]["output_tokens"] == 20
assert advisor_iterations[1]["type"] == "advisor_message"
assert advisor_iterations[1]["model"] == "claude-opus-4-6"
assert advisor_iterations[1]["input_tokens"] == 200
assert advisor_iterations[1]["output_tokens"] == 150
assert advisor_iterations[2]["type"] == "message"
assert advisor_iterations[2]["input_tokens"] == 300
assert advisor_iterations[2]["output_tokens"] == 40
@pytest.mark.asyncio
async def test_should_run_chat_completion_agentic_loop_cleans_up_config_on_no_tool_call():

View file

@ -0,0 +1,253 @@
"""
Tests for advisor orchestration cost aggregation and per-iteration usage
reporting on the Anthropic ``/v1/messages`` path.
Covers:
1. Cost aggregation across executor + advisor sub-calls → final
``_hidden_params["response_cost"]`` equals the sum.
2. ``usage.iterations[]`` reports per-call token breakdowns in order.
3. ``litellm_logging_obj.set_cost_breakdown`` is called with the
"Main Model (initial)" + "Advisor Model" additional costs.
"""
from typing import Any, Dict
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
ADVISOR_TOOL = {
"type": "advisor_20260301",
"name": "advisor",
"model": "claude-opus-4-6",
}
MESSAGES = [{"role": "user", "content": "Help me plan this feature."}]
def _make_executor_tool_use_response(
tool_id: str = "toolu_advisor_01",
input_tokens: int = 100,
output_tokens: int = 20,
response_cost: float = 1.0,
) -> Dict[str, Any]:
return {
"id": "msg_executor_toolcall",
"type": "message",
"role": "assistant",
"model": "openai/gpt-4o-mini",
"content": [
{
"type": "tool_use",
"id": tool_id,
"name": "consult_advisor",
"input": {"question": "What approach should I use?"},
}
],
"stop_reason": "tool_use",
"usage": {
"input_tokens": input_tokens,
"output_tokens": output_tokens,
"cache_read_input_tokens": 0,
"cache_creation_input_tokens": 0,
},
"_hidden_params": {"response_cost": response_cost},
}
def _make_advisor_text_response(
text: str = "Use trial division up to sqrt(n).",
input_tokens: int = 200,
output_tokens: int = 150,
response_cost: float = 0.3,
) -> Dict[str, Any]:
return {
"id": "msg_advisor",
"type": "message",
"role": "assistant",
"model": "claude-opus-4-6",
"content": [{"type": "text", "text": text}],
"stop_reason": "end_turn",
"usage": {
"input_tokens": input_tokens,
"output_tokens": output_tokens,
"cache_read_input_tokens": 0,
"cache_creation_input_tokens": 0,
},
"_hidden_params": {"response_cost": response_cost},
}
def _make_final_executor_response(
text: str = "Here's the implementation.",
input_tokens: int = 300,
output_tokens: int = 40,
response_cost: float = 0.7,
) -> Dict[str, Any]:
return {
"id": "msg_executor_final",
"type": "message",
"role": "assistant",
"model": "openai/gpt-4o-mini",
"content": [{"type": "text", "text": text}],
"stop_reason": "end_turn",
"usage": {
"input_tokens": input_tokens,
"output_tokens": output_tokens,
"cache_read_input_tokens": 0,
"cache_creation_input_tokens": 0,
},
"_hidden_params": {"response_cost": response_cost},
}
@pytest.mark.asyncio
async def test_advisor_orchestration_aggregates_cost_and_iterations():
"""
Executor calls advisor once then produces the final response.
- Total cost = first executor (1.0) + advisor (0.3) + final executor (0.7) = 2.0
- ``usage.iterations`` contains 3 entries in order.
- ``set_cost_breakdown`` is called with ``Main Model (initial)`` and
``Advisor Model`` entries.
"""
from litellm.llms.anthropic.experimental_pass_through.messages.interceptors.advisor import (
AdvisorOrchestrationHandler,
)
first_executor_resp = _make_executor_tool_use_response(response_cost=1.0)
advisor_resp = _make_advisor_text_response(response_cost=0.3)
final_executor_resp = _make_final_executor_response(response_cost=0.7)
executor_call_count = 0
async def mock_messages(model, messages, tools, stream, max_tokens, **kwargs):
nonlocal executor_call_count
executor_call_count += 1
if executor_call_count == 1:
return first_executor_resp
return final_executor_resp
fake_logging_obj = MagicMock()
fake_logging_obj.model_call_details = {}
fake_logging_obj.set_cost_breakdown = MagicMock()
with patch(
"litellm.llms.anthropic.experimental_pass_through.messages.interceptors.advisor._call_messages_handler",
side_effect=mock_messages,
), patch(
"litellm.llms.anthropic.experimental_pass_through.messages.interceptors.advisor._call_advisor_with_router",
new_callable=AsyncMock,
return_value=advisor_resp,
):
h = AdvisorOrchestrationHandler()
result = await h.handle(
model="openai/gpt-4o-mini",
messages=MESSAGES,
tools=[ADVISOR_TOOL],
stream=False,
max_tokens=512,
custom_llm_provider="openai",
litellm_logging_obj=fake_logging_obj,
)
assert executor_call_count == 2
# Final result is the terminating executor response
assert result is final_executor_resp
# Outer response gets its own msg_<uuid> id so the proxy UI logs the
# orchestrated request as a distinct trace (not colliding with any
# inner sub-call's request_id).
assert isinstance(result["id"], str) and result["id"].startswith("msg_")
assert result["id"] != "msg_executor_final"
# Aggregated cost on the response
assert result["_hidden_params"]["response_cost"] == pytest.approx(2.0)
# usage.iterations[] reports per-call breakdowns
usage = result.get("usage")
assert usage is not None
iterations = usage.get("iterations")
assert iterations is not None and len(iterations) == 3
assert iterations[0]["type"] == "message"
assert iterations[0]["input_tokens"] == 100
assert iterations[0]["output_tokens"] == 20
assert iterations[1]["type"] == "advisor_message"
assert iterations[1]["model"] == "claude-opus-4-6"
assert iterations[1]["input_tokens"] == 200
assert iterations[1]["output_tokens"] == 150
assert iterations[2]["type"] == "message"
assert iterations[2]["input_tokens"] == 300
assert iterations[2]["output_tokens"] == 40
# Aggregated usage totals
assert usage["input_tokens"] == 100 + 200 + 300
assert usage["output_tokens"] == 20 + 150 + 40
# cost_breakdown surfaced on the logging object
fake_logging_obj.set_cost_breakdown.assert_called_once()
breakdown_kwargs = fake_logging_obj.set_cost_breakdown.call_args.kwargs
assert breakdown_kwargs["total_cost"] == pytest.approx(2.0)
assert breakdown_kwargs["input_cost"] == pytest.approx(0.7)
assert breakdown_kwargs["additional_costs"] == {
"Main Model (initial)": pytest.approx(1.0),
"Advisor Model": pytest.approx(0.3),
}
# Aggregate cost also recorded on logging model_call_details for downstream
# transforms that drop the dict ``_hidden_params``.
assert fake_logging_obj.model_call_details["response_cost"] == pytest.approx(2.0)
@pytest.mark.asyncio
async def test_advisor_orchestration_no_advisor_call_no_additional_costs():
"""
Executor produces the final response on first try (no advisor call).
- ``Main Model (initial)`` and ``Advisor Model`` should not be present
in ``additional_costs`` (both zero → dict is empty/None).
- ``usage.iterations`` has exactly one entry.
- Total cost matches the single executor turn.
"""
from litellm.llms.anthropic.experimental_pass_through.messages.interceptors.advisor import (
AdvisorOrchestrationHandler,
)
final_resp = _make_final_executor_response(response_cost=0.5)
fake_logging_obj = MagicMock()
fake_logging_obj.model_call_details = {}
fake_logging_obj.set_cost_breakdown = MagicMock()
with patch(
"litellm.llms.anthropic.experimental_pass_through.messages.interceptors.advisor._call_messages_handler",
new_callable=AsyncMock,
return_value=final_resp,
):
h = AdvisorOrchestrationHandler()
result = await h.handle(
model="openai/gpt-4o-mini",
messages=MESSAGES,
tools=[ADVISOR_TOOL],
stream=False,
max_tokens=512,
custom_llm_provider="openai",
litellm_logging_obj=fake_logging_obj,
)
assert result["_hidden_params"]["response_cost"] == pytest.approx(0.5)
# Even when no advisor is called, the orchestrator replaces the inner
# sub-call id with a fresh msg_<uuid> so the UI log row is distinct.
assert isinstance(result["id"], str) and result["id"].startswith("msg_")
assert result["id"] != "msg_executor_final"
iterations = result["usage"]["iterations"]
assert len(iterations) == 1
assert iterations[0]["type"] == "message"
fake_logging_obj.set_cost_breakdown.assert_called_once()
breakdown_kwargs = fake_logging_obj.set_cost_breakdown.call_args.kwargs
assert breakdown_kwargs["total_cost"] == pytest.approx(0.5)
assert breakdown_kwargs["input_cost"] == pytest.approx(0.5)
assert breakdown_kwargs["additional_costs"] in (None, {})