refactor(anthropic/context_management): streaming iterator compaction fixes and compact polyfill improvements

- Extract usage-merge helper; guard empty slice-only compaction result
- Silently drop trailing chunks after usage; remove dead _polyfill_result key
- Fix bedrock None context_mgmt; stream per-instance queue; sync polyfill; trailing-chunk passthrough
- Apply context_management on sync path; clear held stop_reason chunk in async iterator
- Fix post-compaction question selection, system type, sync stream merge
- Skip tool_result-only user turns; bedrock: elif for context_management
- Add streaming iterator compaction test suite

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
Sameer Kankute 2026-05-27 17:25:10 +05:30
parent 6cf22cf9fd
commit 9196e3aa90
No known key found for this signature in database
8 changed files with 467 additions and 243 deletions

View file

@ -37,10 +37,51 @@ if TYPE_CHECKING:
# Anthropic-only keys already mapped by the translator; strip on extra_kwargs re-merge.
ANTHROPIC_ONLY_REQUEST_KEYS: frozenset[str] = frozenset(
{"output_config", "context_management"}
{"output_config"}
)
async def _prepare_context_managed_request(
*,
model: str,
messages: List[Dict],
tools: Optional[List[Dict]],
system: Optional[Any],
context_management_spec: Any,
metadata: Optional[Dict],
drop_params: Optional[bool],
llm_router: Any,
) -> Optional[PolyfillResult]:
"""Apply client compaction history, then optional context_management polyfill."""
from litellm.llms.anthropic.experimental_pass_through.context_management.editors.compact import (
apply_client_compaction_block_history,
)
history_result = apply_client_compaction_block_history(
messages=cast(List[Dict[str, Any]], messages),
system=system,
)
working_messages = (
history_result.messages if history_result is not None else messages
)
working_system = history_result.system if history_result is not None else system
polyfill_result = await _run_polyfill_if_enabled(
model=model,
messages=working_messages,
tools=tools,
system=working_system,
context_management_spec=context_management_spec,
metadata=metadata,
drop_params=drop_params,
llm_router=llm_router,
)
if polyfill_result is not None:
return polyfill_result
return history_result
async def _run_polyfill_if_enabled(
*,
model: str,
@ -372,7 +413,7 @@ class LiteLLMMessagesToCompletionTransformationHandler:
except Exception:
pass
polyfill_result = await _run_polyfill_if_enabled(
polyfill_result = await _prepare_context_managed_request(
model=model,
messages=messages,
tools=tools,
@ -495,7 +536,7 @@ class LiteLLMMessagesToCompletionTransformationHandler:
pass
polyfill_result = run_async_function(
_run_polyfill_if_enabled,
_prepare_context_managed_request,
model=model,
messages=messages,
tools=tools,

View file

@ -20,6 +20,7 @@ from litellm.types.llms.anthropic import (
AppliedEdit,
ContextManagementResponse,
UsageDelta,
UsageIteration,
)
from litellm.types.utils import AdapterCompletionStreamWrapper
@ -61,6 +62,8 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
model: str,
tool_name_mapping: Optional[Dict[str, str]] = None,
applied_edits: Optional[List[AppliedEdit]] = None,
compaction_block: Optional[Dict[str, Any]] = None,
iterations_usage: Optional[List[UsageIteration]] = None,
):
super().__init__(completion_stream)
self.model = model
@ -68,6 +71,10 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
self.tool_name_mapping = tool_name_mapping or {}
# Polyfill applied_edits on final message_delta.
self.applied_edits: List[AppliedEdit] = list(applied_edits or [])
# Synthesized compaction block from compact_20260112 polyfill (streaming).
self.compaction_block = compaction_block
self.iterations_usage = iterations_usage
self.sent_compaction_block: bool = False
# Per-instance queue for buffering multiple chunks. Must be initialized
# here (not at class level) so concurrent streams don't share the same
# deque and corrupt each other's SSE event order.
@ -120,7 +127,59 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
merged_chunk["context_management"] = ContextManagementResponse(
applied_edits=list(self.applied_edits)
)
return merged_chunk
return self._augment_message_delta_usage(merged_chunk)
def _augment_message_delta_usage(self, message_delta_chunk: Dict[str, Any]) -> Dict[str, Any]:
"""Attach polyfill compaction iteration usage to the final message_delta."""
if self.iterations_usage is None:
return message_delta_chunk
usage = message_delta_chunk.get("usage")
if not isinstance(usage, dict) or "iterations" in usage:
return message_delta_chunk
message_iteration: UsageIteration = {
"type": "message",
"input_tokens": usage.get("input_tokens", 0),
"output_tokens": usage.get("output_tokens", 0),
}
augmented = message_delta_chunk.copy()
augmented_usage = dict(usage)
augmented_usage["iterations"] = list(self.iterations_usage) + [message_iteration] # type: ignore[typeddict-unknown-key]
augmented["usage"] = augmented_usage
return augmented
def _queue_compaction_block_events(self) -> None:
"""Emit compaction content-block SSE events before the main response stream.
Anthropic delivers compaction as a single delta (no token-by-token streaming).
"""
if self.compaction_block is None:
return
compaction_index = self.current_content_block_index
summary_content = self.compaction_block.get("content") or ""
self.chunk_queue.append(
{
"type": "content_block_start",
"index": compaction_index,
"content_block": {"type": "compaction"},
}
)
self.chunk_queue.append(
{
"type": "content_block_delta",
"index": compaction_index,
"delta": {"type": "compaction_delta", "content": summary_content},
}
)
self.chunk_queue.append(
{
"type": "content_block_stop",
"index": compaction_index,
}
)
self._increment_content_block_index()
def _create_initial_usage_delta(self) -> UsageDelta:
"""
@ -171,6 +230,11 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
)
return self.chunk_queue.popleft()
if self.sent_compaction_block is False and self.compaction_block is not None:
self.sent_compaction_block = True
self._queue_compaction_block_events()
return self.chunk_queue.popleft()
if self.sent_content_block_start is False:
self.sent_content_block_start = True
self.chunk_queue.append(
@ -275,6 +339,10 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
if processed_chunk.get("delta", {}).get("stop_reason") is not None:
self.holding_stop_reason_chunk = processed_chunk
else:
if processed_chunk.get("type") == "message_delta":
processed_chunk = self._augment_message_delta_usage(
processed_chunk
)
self.chunk_queue.append(processed_chunk)
return self.chunk_queue.popleft()
elif self.holding_chunk is not None:
@ -283,13 +351,21 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
self.holding_chunk = None
return self.chunk_queue.popleft()
else:
if processed_chunk.get("type") == "message_delta":
processed_chunk = self._augment_message_delta_usage(
processed_chunk
)
self.chunk_queue.append(processed_chunk)
return self.chunk_queue.popleft()
# Handle any remaining held chunks after stream ends
if not self.queued_usage_chunk:
if self.holding_stop_reason_chunk is not None:
self.chunk_queue.append(self.holding_stop_reason_chunk)
self.chunk_queue.append(
self._augment_message_delta_usage(
self.holding_stop_reason_chunk
)
)
self.holding_stop_reason_chunk = None
if self.holding_chunk is not None:
@ -309,7 +385,9 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
return self.chunk_queue.popleft()
# Handle any held stop_reason chunk
if self.holding_stop_reason_chunk is not None:
held = self.holding_stop_reason_chunk
held = self._augment_message_delta_usage(
self.holding_stop_reason_chunk
)
self.holding_stop_reason_chunk = None
return held
if self.sent_last_message is False:
@ -350,6 +428,11 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
)
return self.chunk_queue.popleft()
if self.sent_compaction_block is False and self.compaction_block is not None:
self.sent_compaction_block = True
self._queue_compaction_block_events()
return self.chunk_queue.popleft()
if self.sent_content_block_start is False:
self.sent_content_block_start = True
self.chunk_queue.append(
@ -452,6 +535,10 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
):
self.holding_stop_reason_chunk = processed_chunk
else:
if processed_chunk.get("type") == "message_delta":
processed_chunk = self._augment_message_delta_usage(
processed_chunk
)
self.chunk_queue.append(processed_chunk)
return self.chunk_queue.popleft()
elif self.holding_chunk is not None:
@ -461,14 +548,21 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
self.holding_chunk = None
return self.chunk_queue.popleft()
else:
# Queue the current chunk
if processed_chunk.get("type") == "message_delta":
processed_chunk = self._augment_message_delta_usage(
processed_chunk
)
self.chunk_queue.append(processed_chunk)
return self.chunk_queue.popleft()
# Handle any remaining held chunks after stream ends
if not self.queued_usage_chunk:
if self.holding_stop_reason_chunk is not None:
self.chunk_queue.append(self.holding_stop_reason_chunk)
self.chunk_queue.append(
self._augment_message_delta_usage(
self.holding_stop_reason_chunk
)
)
self.holding_stop_reason_chunk = None
if self.holding_chunk is not None:
@ -493,7 +587,9 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
# subsequent ``__anext__`` call doesn't re-emit the same chunk
# (matches the sync ``__next__`` path).
if self.holding_stop_reason_chunk is not None:
held = self.holding_stop_reason_chunk
held = self._augment_message_delta_usage(
self.holding_stop_reason_chunk
)
self.holding_stop_reason_chunk = None
return held
if not self.sent_last_message:

View file

@ -236,12 +236,22 @@ class AnthropicAdapter:
tool_name_mapping: Optional mapping of truncated tool names to original names.
polyfill_result: PolyfillResult from context_management polyfill.
"""
applied_edits = polyfill_result.applied_edits if polyfill_result else None
applied_edits = (
polyfill_result.applied_edits_for_response() if polyfill_result else None
)
compaction_block = (
polyfill_result.compaction_block if polyfill_result is not None else None
)
iterations_usage = (
polyfill_result.iterations_usage if polyfill_result is not None else None
)
anthropic_wrapper = AnthropicStreamWrapper(
completion_stream=completion_stream,
model=model,
tool_name_mapping=tool_name_mapping,
applied_edits=applied_edits,
compaction_block=compaction_block,
iterations_usage=iterations_usage,
)
# Return the SSE-wrapped version for proper event formatting
return anthropic_wrapper.async_anthropic_sse_wrapper()
@ -1424,7 +1434,9 @@ class LiteLLMAnthropicMessagesAdapter:
stop_reason=anthropic_finish_reason,
)
applied_edits = polyfill_result.applied_edits if polyfill_result else None
applied_edits = (
polyfill_result.applied_edits_for_response() if polyfill_result else None
)
if applied_edits:
translated_obj["context_management"] = ContextManagementResponse(
applied_edits=list(applied_edits)

View file

@ -395,6 +395,47 @@ def _extract_usage(response: Any) -> Tuple[int, int]:
)
def apply_client_compaction_block_history(
*,
messages: List[Dict[str, Any]],
system: Optional[Union[str, List[Dict[str, Any]]]],
) -> Optional[PolyfillResult]:
"""Honor client-sent compaction blocks without a ``compact_20260112`` edit.
When the request omits ``context_management`` but the message history already
contains a ``compaction`` content block (e.g. Claude Code client-side compaction),
apply the same slice-only forwarding as the under-threshold path: summary on
``system``, latest user question only on the main call.
"""
effective_messages, prior_compaction_block = _slice_around_compaction_block(
messages
)
if prior_compaction_block is None:
return None
verbose_logger.info(
"compact_20260112: client compaction block in message history; "
"applying slice-only forwarding (no context_management edit)"
)
prior_summary_text = prior_compaction_block.get("content") or ""
augmented_system: Union[str, List[Dict[str, Any]], None] = system
if isinstance(prior_summary_text, str) and prior_summary_text:
augmented_system = _augment_system_with_summary(system, prior_summary_text)
verbose_logger.info(
"compact_20260112: compaction summary added to main call system prefix (%s chars)",
len(prior_summary_text),
)
downstream_messages = _select_last_user_question(effective_messages)
return PolyfillResult(
messages=downstream_messages,
system=augmented_system,
applied_edits=[],
)
async def apply_compact_20260112(
*,
model: str,
@ -415,6 +456,10 @@ async def apply_compact_20260112(
# Validation runs first. Raising AnthropicContextManagementError here is
# the only path on which the polyfill aborts the request.
trigger_tokens, warnings = _resolve_trigger_tokens(edit_spec)
verbose_logger.info(
"compact_20260112: request has compaction trigger (input_tokens threshold=%s)",
trigger_tokens,
)
if edit_spec.get("pause_after_compaction"):
warnings.append("pause_after_compaction_ignored")
@ -442,6 +487,10 @@ async def apply_compact_20260112(
augmented_system: Union[str, List[Dict[str, Any]], None] = system
if isinstance(prior_summary_text, str) and prior_summary_text:
augmented_system = _augment_system_with_summary(system, prior_summary_text)
verbose_logger.info(
"compact_20260112: compaction summary added to main call system prefix (%s chars)",
len(prior_summary_text),
)
downstream_messages = _strip_compaction_blocks(effective_messages)
@ -464,14 +513,14 @@ async def apply_compact_20260112(
)
if current_tokens <= trigger_tokens:
# Slice-only path. If Phase A fired we still slice + system-prefix.
# Guard against an empty downstream message list: this happens when the
# sliced result was a single assistant turn whose only block was the
# compaction block, which ``_strip_compaction_blocks`` then drops. The
# downstream API rejects an empty messages array, so fall back to
# picking a real user question (or a synthetic continuation prompt) —
# mirroring the full-summary path's behavior below.
if not downstream_messages:
# Slice-only path: prior context lives in ``augmented_system`` (the
# compaction summary prefix). The main model call must not re-send stale
# assistant turns from the post-compaction tail — only the latest user
# question, matching the full-summary path below.
if prior_compaction_block is not None:
downstream_messages = _select_last_user_question(effective_messages)
elif not downstream_messages:
# No compaction checkpoint: only substitute when strip left nothing.
downstream_messages = _select_last_user_question(effective_messages)
return PolyfillResult(
messages=downstream_messages,
@ -537,6 +586,10 @@ async def apply_compact_20260112(
# eligible turn exists, fall back to a synthetic continuation prompt so
# the downstream call still has a non-empty user message.
summarized_system = _augment_system_with_summary(system, summary_text)
verbose_logger.info(
"compact_20260112: compaction summary added to main call system prefix (%s chars)",
len(summary_text),
)
downstream_messages_after_summary = _select_last_user_question(effective_messages)
return PolyfillResult(

View file

@ -14,6 +14,8 @@ from litellm.types.llms.anthropic import (
UsageIteration,
)
from .constants import COMPACT_EDIT_TYPE
@dataclass
class PolyfillResult:
@ -22,3 +24,19 @@ class PolyfillResult:
applied_edits: List[AppliedEdit] = field(default_factory=list)
compaction_block: Optional[CompactionBlock] = None
iterations_usage: Optional[List[UsageIteration]] = None
def applied_edits_for_response(self) -> Optional[List[AppliedEdit]]:
"""``applied_edits`` to attach on the client-visible response.
``compact_20260112`` is included only when a new compaction block was
synthesized (slice-only / under-threshold paths omit it). Other edit
types are included when the editor returned an ``AppliedEdit``.
"""
visible: List[AppliedEdit] = []
for edit in self.applied_edits:
if edit.get("type") == COMPACT_EDIT_TYPE:
if self.compaction_block is not None:
visible.append(edit)
else:
visible.append(edit)
return visible or None

View file

@ -1,231 +1,19 @@
model_list:
- model_name: gpt-5-mini-end-user-test
- model_name: claude-sonnet-4-6
litellm_params:
model: gpt-5-mini
region_name: "eu"
model_info:
id: "1"
- model_name: gpt-5-mini-end-user-test
model: bedrock/converse/global.anthropic.claude-sonnet-4-6
# aws_region_name: "global"
- model_name: claude-sonnet-4-5
litellm_params:
model: openai/gpt-5-mini
api_key: os.environ/OPENAI_API_KEY # The `os.environ/` prefix tells litellm to read this from the env. See https://docs.litellm.ai/docs/simple_proxy#load-api-keys-from-vault
- model_name: gpt-3.5-turbo
model: bedrock/converse/global.anthropic.claude-sonnet-4-5-20250929-v1:0
# aws_region_name: "global"
- model_name: gpt-5.4-mini
litellm_params:
model: openai/gpt-4.1-mini
api_key: os.environ/OPENAI_API_KEY # The `os.environ/` prefix tells litellm to read this from the env. See https://docs.litellm.ai/docs/simple_proxy#load-api-keys-from-vault
- model_name: gpt-3.5-turbo-large
litellm_params:
model: "gpt-4.1"
api_key: os.environ/OPENAI_API_KEY
rpm: 480
timeout: 300
stream_timeout: 60
- model_name: gpt-4
litellm_params:
model: openai/gpt-4.1
api_key: os.environ/OPENAI_API_KEY # The `os.environ/` prefix tells litellm to read this from the env. See https://docs.litellm.ai/docs/simple_proxy#load-api-keys-from-vault
rpm: 480
timeout: 300
stream_timeout: 60
- model_name: sagemaker-completion-model
litellm_params:
model: sagemaker/berri-benchmarking-Llama-2-70b-chat-hf-4
input_cost_per_second: 0.000420
- model_name: text-embedding-ada-002
litellm_params:
model: openai/text-embedding-3-small
api_key: os.environ/OPENAI_API_KEY
model_info:
mode: embedding
base_model: text-embedding-3-small
- model_name: dall-e-2 # dall-e-2 and dall-e-3 were deprecated 2026-05-12; alias to gpt-image-1
litellm_params:
model: openai/gpt-image-1
- model_name: openai-dall-e-3 # dall-e-3 deprecated 2026-05-12; underlying now gpt-image-1
litellm_params:
model: gpt-image-1
- model_name: fake-openai-endpoint
litellm_params:
model: openai/gpt-5-mini
api_key: fake-key
api_base: https://exampleopenaiendpoint-production.up.railway.app/
- model_name: fake-openai-endpoint-2
litellm_params:
model: openai/my-fake-model
api_key: my-fake-key
api_base: https://exampleopenaiendpoint-production.up.railway.app/
stream_timeout: 0.001
rpm: 1
- model_name: fake-openai-endpoint-3
litellm_params:
model: openai/my-fake-model
api_key: my-fake-key
api_base: https://exampleopenaiendpoint-production.up.railway.app/
stream_timeout: 0.001
rpm: 1000
- model_name: fake-openai-endpoint-4
litellm_params:
model: openai/my-fake-model
api_key: my-fake-key
api_base: https://exampleopenaiendpoint-production.up.railway.app/
num_retries: 50
- model_name: fake-openai-endpoint-3
litellm_params:
model: openai/my-fake-model-2
api_key: my-fake-key
api_base: https://exampleopenaiendpoint-production.up.railway.app/
stream_timeout: 0.001
rpm: 1000
- model_name: bad-model
litellm_params:
model: openai/bad-model
api_key: os.environ/OPENAI_API_KEY
api_base: https://exampleopenaiendpoint-production.up.railway.app/
mock_timeout: True
timeout: 60
rpm: 1000
model_info:
health_check_timeout: 1
- model_name: good-model
litellm_params:
model: openai/bad-model
api_key: os.environ/OPENAI_API_KEY
api_base: https://exampleopenaiendpoint-production.up.railway.app/
rpm: 1000
model_info:
health_check_timeout: 1
- model_name: "*"
litellm_params:
model: openai/*
api_key: os.environ/OPENAI_API_KEY
- model_name: realtime-v1
litellm_params:
model: azure/gpt-realtime-20250828-standard
api_version: "2025-08-28"
realtime_protocol: GA # Possible values: "GA"/ "v1", "beta"
- model_name: realtime-beta
litellm_params:
model: azure/gpt-realtime-20250828-standard
api_version: 2025-04-01-preview
# provider specific wildcard routing
- model_name: "anthropic/*"
litellm_params:
model: "anthropic/*"
api_key: os.environ/ANTHROPIC_API_KEY
- model_name: "bedrock/*"
litellm_params:
model: "bedrock/*"
- model_name: "groq/*"
litellm_params:
model: "groq/*"
api_key: os.environ/GROQ_API_KEY
- model_name: mistral-embed
litellm_params:
model: mistral/mistral-embed
- model_name: gpt-instruct # [PROD TEST] - tests if `/health` automatically infers this to be a text completion model
litellm_params:
model: text-completion-openai/gpt-3.5-turbo-instruct
- model_name: fake-openai-endpoint-5
litellm_params:
model: openai/my-fake-model
api_key: my-fake-key
api_base: https://exampleopenaiendpoint-production.up.railway.app/
timeout: 1
- model_name: badly-configured-openai-endpoint
litellm_params:
model: openai/my-fake-model
api_key: my-fake-key
api_base: https://exampleopenaiendpoint-production.up.railway.appxxxx/
- model_name: gemini-2.5-flash
litellm_params:
model: gemini/gemini-2.5-flash
api_key: os.environ/GOOGLE_API_KEY
- model_name: gpt-5.5
litellm_params:
model: gpt-5.5
model: openai/gpt-5.4-mini
api_key: os.environ/OPENAI_API_KEY
general_settings:
context_management_summary_model: gpt-5.4-mini
litellm_settings:
# set_verbose: True # Uncomment this if you want to see verbose logs; not recommended in production
drop_params: True
success_callback: ["prometheus"]
# max_budget: 100
# budget_duration: 30d
num_retries: 5
request_timeout: 600
telemetry: False
context_window_fallbacks: [{"gpt-3.5-turbo": ["gpt-3.5-turbo-large"]}]
default_team_settings:
- team_id: team-1
success_callback: ["langfuse"]
failure_callback: ["langfuse"]
langfuse_public_key: os.environ/LANGFUSE_PROJECT1_PUBLIC # Project 1
langfuse_secret: os.environ/LANGFUSE_PROJECT1_SECRET # Project 1
- team_id: team-2
success_callback: ["langfuse"]
failure_callback: ["langfuse"]
langfuse_public_key: os.environ/LANGFUSE_PROJECT2_PUBLIC # Project 2
langfuse_secret: os.environ/LANGFUSE_PROJECT2_SECRET # Project 2
langfuse_host: https://us.cloud.langfuse.com
# cache: true # [OPTIONAL] use for caching responses
# enable_caching_on_provider_specific_optional_params: True # Include provider-specific params in cache keys
# cache_params: # And for shared health check
# type: redis
# host: localhost
# port: 6379
# For /fine_tuning/jobs endpoints
finetune_settings:
- custom_llm_provider: azure
api_base: os.environ/AZURE_API_BASE
api_key: os.environ/AZURE_API_KEY
api_version: "2023-03-15-preview"
- custom_llm_provider: openai
api_key: os.environ/OPENAI_API_KEY
# for /files endpoints
files_settings:
- custom_llm_provider: azure
api_base: os.environ/AZURE_API_BASE
api_key: os.environ/AZURE_API_KEY
api_version: "2023-03-15-preview"
- custom_llm_provider: openai
api_key: os.environ/OPENAI_API_KEY
router_settings:
routing_strategy: usage-based-routing-v2
redis_host: os.environ/REDIS_HOST
redis_password: os.environ/REDIS_PASSWORD
redis_port: os.environ/REDIS_PORT
enable_pre_call_checks: true
model_group_alias: {"my-special-fake-model-alias-name": "fake-openai-endpoint-3"}
general_settings:
master_key: sk-1234 # [OPTIONAL] Use to enforce auth on proxy. See - https://docs.litellm.ai/docs/proxy/virtual_keys
store_model_in_db: True
proxy_budget_rescheduler_min_time: 60
proxy_budget_rescheduler_max_time: 64
proxy_batch_write_at: 1
database_connection_pool_limit: 10
# background_health_checks: true
# use_shared_health_check: true
# health_check_interval: 30
# database_url: "postgresql://<user>:<password>@<host>:<port>/<dbname>" # [OPTIONAL] use for token-based auth to proxy
pass_through_endpoints:
- path: "/v1/rerank" # route you want to add to LiteLLM Proxy Server
target: "https://api.cohere.com/v1/rerank" # URL this route should forward requests to
headers: # headers to forward to this URL
content-type: application/json # (Optional) Extra Headers to pass to this endpoint
accept: application/json
forward_headers: True
# environment_variables:
# settings for using redis caching
# REDIS_HOST: redis-16337.c322.us-east-1-2.ec2.cloud.redislabs.com
# REDIS_PORT: "16337"
# REDIS_PASSWORD:
use_chat_completions_url_for_anthropic_messages: true

View file

@ -0,0 +1,151 @@
"""Compaction block SSE events from AnthropicStreamWrapper (compact_20260112 polyfill)."""
import os
import sys
from typing import List
from unittest.mock import MagicMock
import pytest
sys.path.insert(0, os.path.abspath("../../../../.."))
from litellm.llms.anthropic.experimental_pass_through.adapters.streaming_iterator import (
AnthropicStreamWrapper,
)
from litellm.types.utils import Delta, StreamingChoices
def _make_text_chunk(text: str, finish_reason: str = None) -> MagicMock:
chunk = MagicMock()
chunk.choices = [
StreamingChoices(
finish_reason=finish_reason,
index=0,
delta=Delta(content=text, role="assistant" if text else None, tool_calls=None),
logprobs=None,
)
]
chunk.usage = None
chunk._hidden_params = {}
return chunk
async def _collect_events_async(wrapper: AnthropicStreamWrapper) -> List[dict]:
events = []
async for event in wrapper:
events.append(event)
return events
@pytest.mark.asyncio
async def test_stream_emits_compaction_block_before_text():
"""Polyfill compaction_block must surface as compaction SSE events at index 0."""
async def mock_stream():
yield _make_text_chunk("Hi")
yield _make_text_chunk("", finish_reason="stop")
compaction_block = {
"type": "compaction",
"content": "Summary of prior conversation turns.",
}
iterations_usage = [
{"type": "compaction", "input_tokens": 100, "output_tokens": 50},
]
wrapper = AnthropicStreamWrapper(
completion_stream=mock_stream(),
model="claude-sonnet-4-6",
compaction_block=compaction_block,
iterations_usage=iterations_usage,
applied_edits=[{"type": "compact_20260112"}],
)
events = await _collect_events_async(wrapper)
compaction_start = next(
e
for e in events
if e.get("type") == "content_block_start"
and e.get("content_block", {}).get("type") == "compaction"
)
assert compaction_start["index"] == 0
compaction_delta = next(
e
for e in events
if e.get("type") == "content_block_delta"
and e.get("delta", {}).get("type") == "compaction_delta"
)
assert compaction_delta["index"] == 0
assert compaction_delta["delta"]["content"] == "Summary of prior conversation turns."
compaction_stop = next(
e
for e in events
if e.get("type") == "content_block_stop" and e.get("index") == 0
)
assert compaction_stop is not None
text_start = next(
e
for e in events
if e.get("type") == "content_block_start"
and e.get("content_block", {}).get("type") == "text"
)
assert text_start["index"] == 1
message_delta = next(e for e in events if e.get("type") == "message_delta")
iterations = message_delta.get("usage", {}).get("iterations")
assert iterations is not None
assert iterations[0]["type"] == "compaction"
assert iterations[1]["type"] == "message"
@pytest.mark.asyncio
async def test_stream_omits_context_management_when_no_compaction_applied():
"""applied_edits without a compaction block must not emit context_management."""
async def mock_stream():
yield _make_text_chunk("Hello")
yield _make_text_chunk("", finish_reason="stop")
wrapper = AnthropicStreamWrapper(
completion_stream=mock_stream(),
model="claude-sonnet-4-6",
applied_edits=None,
)
events = await _collect_events_async(wrapper)
message_deltas = [e for e in events if e.get("type") == "message_delta"]
assert message_deltas
assert "context_management" not in message_deltas[-1]
@pytest.mark.asyncio
async def test_stream_without_compaction_block_unchanged():
"""No compaction_block means no compaction SSE events."""
async def mock_stream():
yield _make_text_chunk("Hello")
yield _make_text_chunk("", finish_reason="stop")
wrapper = AnthropicStreamWrapper(
completion_stream=mock_stream(),
model="claude-sonnet-4-6",
)
events = await _collect_events_async(wrapper)
assert not any(
e.get("content_block", {}).get("type") == "compaction"
for e in events
if e.get("type") == "content_block_start"
)
text_start = next(
e
for e in events
if e.get("type") == "content_block_start"
and e.get("content_block", {}).get("type") == "text"
)
assert text_start["index"] == 0

View file

@ -26,8 +26,12 @@ from litellm.llms.anthropic.experimental_pass_through.context_management.editors
_extract_summary_text,
_slice_around_compaction_block,
_strip_compaction_blocks,
apply_client_compaction_block_history,
apply_compact_20260112,
)
from litellm.llms.anthropic.experimental_pass_through.context_management.result import (
PolyfillResult,
)
MODEL = "openai/gpt-4o"
@ -84,6 +88,36 @@ def _make_mock_response(
# ---------------------------------------------------------------------------
def test_applied_edits_for_response_omits_compact_without_compaction_block():
"""Slice-only / under-threshold: no client-visible context_management edit."""
result = PolyfillResult(
messages=[],
system="summary on system",
applied_edits=[{"type": "compact_20260112"}],
compaction_block=None,
)
assert result.applied_edits_for_response() is None
def test_applied_edits_for_response_includes_compact_when_block_present():
result = PolyfillResult(
messages=[],
system=None,
applied_edits=[
{
"type": "compact_20260112",
"summary_input_tokens": 10,
"summary_output_tokens": 5,
}
],
compaction_block={"type": "compaction", "content": "summary"},
)
visible = result.applied_edits_for_response()
assert visible is not None
assert visible[0]["type"] == "compact_20260112"
assert visible[0]["summary_input_tokens"] == 10
def test_slice_around_compaction_block_found():
messages = _messages_with_compaction("my summary")
sliced, block = _slice_around_compaction_block(messages)
@ -243,6 +277,34 @@ async def test_opt_in_gating_no_summary_model_configured():
assert result.iterations_usage is None
# ---------------------------------------------------------------------------
# Client compaction block without context_management
# ---------------------------------------------------------------------------
def test_client_compaction_block_history_without_context_management():
"""Compaction in messages alone triggers slice-only forwarding."""
messages = _messages_with_compaction("prior summary text")
result = apply_client_compaction_block_history(messages=messages, system=None)
assert result is not None
assert result.system is not None
assert "prior summary text" in str(result.system)
assert result.compaction_block is None
assert result.applied_edits == []
assert len(result.messages) == 1
assert result.messages[0]["role"] == "user"
assert result.messages[0]["content"] == "latest question"
def test_client_compaction_block_history_no_compaction_returns_none():
result = apply_client_compaction_block_history(
messages=_simple_messages(), system="base"
)
assert result is None
# ---------------------------------------------------------------------------
# Editor: slice-only path
# ---------------------------------------------------------------------------
@ -275,7 +337,10 @@ async def test_slice_only_path_with_existing_compaction_block():
assert result.compaction_block is None
assert result.iterations_usage is None
# No compaction block in downstream messages
# Main call: summary on system, latest user question only (no stale assistant).
assert len(result.messages) == 1
assert result.messages[0]["role"] == "user"
assert result.messages[0]["content"] == "latest question"
for msg in result.messages:
content = msg.get("content")
if isinstance(content, list):