mirror of
https://github.com/headroomlabs-ai/headroom.git
synced 2026-08-27 14:17:10 -04:00
## Description Large native Claude Code requests on the Anthropic path can still forward with zero request compression after CCR tool injection is deferred on a frozen prefix. The stale skip branch treats deferred injection as a reason to bypass request compression entirely, even though the later sticky CCR path already knows when new markers actually require the tool. This removes that stale bypass so compression still runs while the reversible CCR path stays intact. Closes #2291. ## Type of Change - [x] Bug fix (non-breaking change that fixes an issue) - [ ] New feature (non-breaking change that adds functionality) - [ ] Breaking change (fix or feature that would cause existing functionality to change) - [ ] Documentation update - [ ] Performance improvement - [ ] Code refactoring (no functional changes) ## Changes Made - Remove the stale `should_skip_ccr_request_compression` branch from `headroom/proxy/handlers/anthropic.py`, so deferred CCR tool injection no longer bypasses request compression in token, non-cache, or cache mode. - Keep the existing sticky CCR injection path as the only place that decides whether historical markers need the retrieval tool reintroduced. - Update `tests/test_proxy/test_anthropic_ccr_deferred_injection.py` to cover the two broken zero-compression cases and preserve the already-reversible frozen-prefix path. ## Testing - [x] Unit tests pass (`uv run pytest tests/test_proxy/test_anthropic_ccr_deferred_injection.py::test_cache_mode_compresses_delta_but_replays_cached_prefix_when_markers_are_historical tests/test_proxy/test_anthropic_ccr_deferred_injection.py::test_non_token_non_cache_mode_keeps_compression_and_injects_tool_for_new_markers tests/test_proxy/test_anthropic_ccr_deferred_injection.py::test_existing_retrieve_tool_keeps_reversible_ccr_path_when_prefix_is_frozen tests/test_openai_tool_search_deferral.py tests/test_openai_responses_compression_units.py -q`) - [x] Linting passes (`uv run ruff check .`) - [ ] Type checking passes (`uv run mypy headroom`) - [x] New tests added for new functionality when applicable - [ ] Manual testing performed ### Test Output ```text ======================= 53 passed, 1 warning in 12.71s ======================== All checks passed! 1310 files already formatted ``` ## Real Behavior Proof - Environment: synced branch and base worktrees on a local Windows proxy test host - Exact command / steps: run the updated Anthropic deferred-injection regressions directly against the base package tree and the branch package tree, then run the focused branch pytest suite above - Observed result: the base package tree fails the two updated zero-compression regressions (`FAIL test_cache_mode_compresses_delta_but_replays_cached_prefix_when_markers_are_historical`, `FAIL test_non_token_non_cache_mode_keeps_compression_and_injects_tool_for_new_markers`) while preserving the already-reversible path; the branch package tree prints three `PASS` lines for the same trio and keeps the neighboring OpenAI suites green in the 53-test focused run - Not tested: a live upstream Claude Code request with the reporter's exact provider/model credentials ## Review Readiness - [x] I have performed a self-review - [x] This PR is ready for human review ## Checklist - [x] My code follows the project's style guidelines - [x] I have performed a self-review of my code - [x] I have commented my code, particularly in hard-to-understand areas - [ ] I have made corresponding changes to the documentation - [x] My changes generate no new warnings - [x] I have added tests that prove my fix is effective or that my feature works - [x] New and existing unit tests pass locally with my changes - [ ] I have updated the CHANGELOG.md if applicable ## Additional Notes `CHANGELOG.md` is not applicable here because Headroom's release pipeline derives it from conventional commits. Scope is limited to the Anthropic CCR request-compression seam that current #2291 evidence exercises. OpenAI tool-search deferral is untouched because the current live issue is a native Claude Code path and the concrete stale skip branch on `origin/main` is in `anthropic.py`. Co-authored-by: Tejas Chopra <chopratejas@gmail.com>
This commit is contained in:
parent
1d79e70f95
commit
26b43f64d6
2 changed files with 223 additions and 276 deletions
|
|
@ -1286,26 +1286,6 @@ class AnthropicHandlerMixin:
|
|||
from headroom.transforms.compression_policy import resolve_policy
|
||||
|
||||
compression_policy = resolve_policy(getattr(request.state, "auth_mode", None))
|
||||
from headroom.ccr.tool_injection import CCR_TOOL_NAME
|
||||
|
||||
existing_tool_names = {
|
||||
tool.get("name") or tool.get("function", {}).get("name")
|
||||
for tool in (body.get("tools") or [])
|
||||
if isinstance(tool, dict)
|
||||
}
|
||||
|
||||
def should_skip_ccr_request_compression(
|
||||
current_frozen_message_count: int,
|
||||
) -> bool:
|
||||
if is_token_mode(self.config.mode):
|
||||
return False
|
||||
# If the tool is already present, CCR stays reversible even on frozen turns.
|
||||
return (
|
||||
self.config.ccr_inject_tool
|
||||
and current_frozen_message_count > 0
|
||||
and CCR_TOOL_NAME not in existing_tool_names
|
||||
)
|
||||
|
||||
if is_token_mode(self.config.mode):
|
||||
comp_cache = self._get_compression_cache(session_id)
|
||||
|
||||
|
|
@ -1337,29 +1317,109 @@ class AnthropicHandlerMixin:
|
|||
# Record all tool_results in the verified frozen prefix as stable
|
||||
comp_cache.mark_stable_from_messages(messages, frozen_message_count)
|
||||
|
||||
skip_ccr_request_compression = should_skip_ccr_request_compression(
|
||||
frozen_message_count
|
||||
)
|
||||
if skip_ccr_request_compression:
|
||||
logger.info(
|
||||
f"[{request_id}] CCR: skipping request-side compression "
|
||||
f"(frozen prefix={frozen_message_count}) because tool injection is deferred"
|
||||
# Zone 1: Swap cached compressed versions into working copy
|
||||
working_messages = comp_cache.apply_cached(messages)
|
||||
if (
|
||||
getattr(self, "_background_compression_enabled", False)
|
||||
and frozen_message_count == 0
|
||||
and original_tokens >= self._background_compression_min_tokens
|
||||
):
|
||||
accepted = self._background_compressor.enqueue(
|
||||
session_id,
|
||||
lambda: self.anthropic_pipeline.apply(
|
||||
messages=working_messages,
|
||||
model=model,
|
||||
model_limit=context_limit,
|
||||
context=extract_user_query(working_messages),
|
||||
frozen_message_count=frozen_message_count,
|
||||
idle_seconds=idle_seconds,
|
||||
biases=biases,
|
||||
request_id=request_id,
|
||||
compression_policy=compression_policy,
|
||||
**proxy_pipeline_kwargs(self.config),
|
||||
),
|
||||
lambda bg_result: comp_cache.update_from_result(
|
||||
messages, bg_result.messages
|
||||
),
|
||||
)
|
||||
if skip_ccr_request_compression:
|
||||
optimized_messages = messages
|
||||
_, optimized_tokens = await self._count_tokens_offloaded(
|
||||
model, optimized_messages
|
||||
|
||||
# Cold-start fast pass: run everything EXCEPT the
|
||||
# Kompress ML stage synchronously before forwarding.
|
||||
# The byte-identical freeze (#1850) locks a session
|
||||
# to whatever form its cold start put in the provider
|
||||
# cache; deferring the WHOLE pipeline locks in the raw
|
||||
# transcript and forfeits the session's savings for
|
||||
# its lifetime — including sub-second wins like
|
||||
# read_lifecycle stale-read drops. Only Kompress can
|
||||
# blow the request budget (#1171), so only Kompress
|
||||
# stays deferred. Fail-open: on timeout/error the
|
||||
# request forwards exactly as before this pass
|
||||
# existed. On timeout the worker can't be cancelled
|
||||
# (Python can't preempt a running thread), but
|
||||
# _run_compression_in_executor runs it on the bounded
|
||||
# compression pool and tracks it via the leaked-thread
|
||||
# metric, so stragglers are capped and observable
|
||||
# rather than unbounded. The pass is also bounded by
|
||||
# routing + statistical crushers (observed seconds even
|
||||
# on multi-M-token counts).
|
||||
from headroom.proxy.helpers import (
|
||||
COLD_START_FAST_PASS_TIMEOUT_SECONDS,
|
||||
)
|
||||
|
||||
_fast_pass = None
|
||||
try:
|
||||
async with stage_timer.measure("compression_first_stage"):
|
||||
_fast_pass = await self._run_compression_in_executor(
|
||||
lambda: self.anthropic_pipeline.apply(
|
||||
messages=working_messages,
|
||||
model=model,
|
||||
model_limit=context_limit,
|
||||
context=extract_user_query(working_messages),
|
||||
frozen_message_count=frozen_message_count,
|
||||
idle_seconds=idle_seconds,
|
||||
biases=biases,
|
||||
request_id=request_id,
|
||||
compression_policy=compression_policy,
|
||||
skip_kompress=True,
|
||||
**proxy_pipeline_kwargs(self.config),
|
||||
),
|
||||
timeout=COLD_START_FAST_PASS_TIMEOUT_SECONDS,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.info(
|
||||
"[%s] Cold-start fast pass skipped (%s: %s); "
|
||||
"deferring full pipeline to background",
|
||||
request_id,
|
||||
type(e).__name__,
|
||||
e,
|
||||
)
|
||||
|
||||
if _fast_pass is not None:
|
||||
comp_cache.update_from_result(messages, _fast_pass.messages)
|
||||
_fast_pass.transforms_applied = list(
|
||||
_fast_pass.transforms_applied
|
||||
) + [
|
||||
"deferred:kompress_background"
|
||||
if accepted
|
||||
else "deferred:dropped"
|
||||
]
|
||||
result = _fast_pass
|
||||
else:
|
||||
|
||||
class _DeferredCompressionResult:
|
||||
messages = working_messages
|
||||
transforms_applied = [
|
||||
"deferred:background_compression"
|
||||
if accepted
|
||||
else "deferred:dropped"
|
||||
]
|
||||
timing = {}
|
||||
waste_signals = None
|
||||
|
||||
result = _DeferredCompressionResult()
|
||||
else:
|
||||
# Zone 1: Swap cached compressed versions into working copy
|
||||
working_messages = comp_cache.apply_cached(messages)
|
||||
if (
|
||||
getattr(self, "_background_compression_enabled", False)
|
||||
and frozen_message_count == 0
|
||||
and original_tokens >= self._background_compression_min_tokens
|
||||
):
|
||||
accepted = self._background_compressor.enqueue(
|
||||
session_id,
|
||||
async with stage_timer.measure("compression_first_stage"):
|
||||
result = await self._run_compression_in_executor(
|
||||
lambda: self.anthropic_pipeline.apply(
|
||||
messages=working_messages,
|
||||
model=model,
|
||||
|
|
@ -1372,168 +1432,57 @@ class AnthropicHandlerMixin:
|
|||
compression_policy=compression_policy,
|
||||
**proxy_pipeline_kwargs(self.config),
|
||||
),
|
||||
lambda bg_result: comp_cache.update_from_result(
|
||||
messages, bg_result.messages
|
||||
),
|
||||
)
|
||||
|
||||
# Cold-start fast pass: run everything EXCEPT the
|
||||
# Kompress ML stage synchronously before forwarding.
|
||||
# The byte-identical freeze (#1850) locks a session
|
||||
# to whatever form its cold start put in the provider
|
||||
# cache; deferring the WHOLE pipeline locks in the raw
|
||||
# transcript and forfeits the session's savings for
|
||||
# its lifetime — including sub-second wins like
|
||||
# read_lifecycle stale-read drops. Only Kompress can
|
||||
# blow the request budget (#1171), so only Kompress
|
||||
# stays deferred. Fail-open: on timeout/error the
|
||||
# request forwards exactly as before this pass
|
||||
# existed. On timeout the worker can't be cancelled
|
||||
# (Python can't preempt a running thread), but
|
||||
# _run_compression_in_executor runs it on the bounded
|
||||
# compression pool and tracks it via the leaked-thread
|
||||
# metric, so stragglers are capped and observable
|
||||
# rather than unbounded. The pass is also bounded by
|
||||
# routing + statistical crushers (observed seconds even
|
||||
# on multi-M-token counts).
|
||||
from headroom.proxy.helpers import (
|
||||
COLD_START_FAST_PASS_TIMEOUT_SECONDS,
|
||||
)
|
||||
|
||||
_fast_pass = None
|
||||
try:
|
||||
async with stage_timer.measure("compression_first_stage"):
|
||||
_fast_pass = await self._run_compression_in_executor(
|
||||
lambda: self.anthropic_pipeline.apply(
|
||||
messages=working_messages,
|
||||
model=model,
|
||||
model_limit=context_limit,
|
||||
context=extract_user_query(working_messages),
|
||||
frozen_message_count=frozen_message_count,
|
||||
idle_seconds=idle_seconds,
|
||||
biases=biases,
|
||||
request_id=request_id,
|
||||
compression_policy=compression_policy,
|
||||
skip_kompress=True,
|
||||
**proxy_pipeline_kwargs(self.config),
|
||||
),
|
||||
timeout=COLD_START_FAST_PASS_TIMEOUT_SECONDS,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.info(
|
||||
"[%s] Cold-start fast pass skipped (%s: %s); "
|
||||
"deferring full pipeline to background",
|
||||
request_id,
|
||||
type(e).__name__,
|
||||
e,
|
||||
)
|
||||
|
||||
if _fast_pass is not None:
|
||||
comp_cache.update_from_result(messages, _fast_pass.messages)
|
||||
_fast_pass.transforms_applied = list(
|
||||
_fast_pass.transforms_applied
|
||||
) + [
|
||||
"deferred:kompress_background"
|
||||
if accepted
|
||||
else "deferred:dropped"
|
||||
]
|
||||
result = _fast_pass
|
||||
else:
|
||||
|
||||
class _DeferredCompressionResult:
|
||||
messages = working_messages
|
||||
transforms_applied = [
|
||||
"deferred:background_compression"
|
||||
if accepted
|
||||
else "deferred:dropped"
|
||||
]
|
||||
timing = {}
|
||||
waste_signals = None
|
||||
|
||||
result = _DeferredCompressionResult()
|
||||
else:
|
||||
async with stage_timer.measure("compression_first_stage"):
|
||||
result = await self._run_compression_in_executor(
|
||||
lambda: self.anthropic_pipeline.apply(
|
||||
messages=working_messages,
|
||||
model=model,
|
||||
model_limit=context_limit,
|
||||
context=extract_user_query(working_messages),
|
||||
frozen_message_count=frozen_message_count,
|
||||
idle_seconds=idle_seconds,
|
||||
biases=biases,
|
||||
request_id=request_id,
|
||||
compression_policy=compression_policy,
|
||||
**proxy_pipeline_kwargs(self.config),
|
||||
),
|
||||
timeout=COMPRESSION_TIMEOUT_SECONDS,
|
||||
)
|
||||
|
||||
# Cache newly compressed messages (index-aligned diff)
|
||||
if result.messages != working_messages:
|
||||
comp_cache.update_from_result(messages, result.messages)
|
||||
|
||||
# Always use pipeline result — Zone 1 swaps are already applied
|
||||
optimized_messages = result.messages
|
||||
transforms_applied = result.transforms_applied
|
||||
pipeline_timing = result.timing
|
||||
# Issue #327 / Bug 3: pipeline.apply uses the provider-
|
||||
# side tokenizer (AnthropicProvider tiktoken estimator),
|
||||
# which counts ~25% higher than the proxy-side
|
||||
# EstimatingTokenCounter used to set `original_tokens`
|
||||
# at line 634. Reusing `result.tokens_after` here
|
||||
# produced an apples-vs-oranges comparison against
|
||||
# `original_tokens` in the inflation guard below
|
||||
# (line ~901): even after a real 12% compression the
|
||||
# provider-tokenizer figure was higher than the proxy-
|
||||
# tokenizer baseline, triggering a spurious revert.
|
||||
# Recount optimized_messages with the proxy tokenizer
|
||||
# so original_tokens vs optimized_tokens is self-
|
||||
# consistent. The recount cost (~ms on a 50K-token
|
||||
# request) is paid once per request and is dwarfed by
|
||||
# the upstream call latency.
|
||||
optimized_tokens = tokenizer.count_messages(optimized_messages)
|
||||
elif not is_cache_mode(self.config.mode):
|
||||
skip_ccr_request_compression = should_skip_ccr_request_compression(
|
||||
frozen_message_count
|
||||
)
|
||||
if skip_ccr_request_compression:
|
||||
logger.info(
|
||||
f"[{request_id}] CCR: skipping request-side compression "
|
||||
f"(frozen prefix={frozen_message_count}) because tool injection is deferred"
|
||||
)
|
||||
if not skip_ccr_request_compression:
|
||||
async with stage_timer.measure("compression_first_stage"):
|
||||
result = await self._run_compression_in_executor(
|
||||
lambda: self.anthropic_pipeline.apply(
|
||||
messages=messages,
|
||||
model=model,
|
||||
model_limit=context_limit,
|
||||
context=extract_user_query(messages),
|
||||
frozen_message_count=frozen_message_count,
|
||||
biases=biases,
|
||||
request_id=request_id,
|
||||
compression_policy=compression_policy,
|
||||
**proxy_pipeline_kwargs(self.config),
|
||||
),
|
||||
timeout=COMPRESSION_TIMEOUT_SECONDS,
|
||||
)
|
||||
|
||||
if result.messages != messages:
|
||||
optimized_messages = result.messages
|
||||
transforms_applied = result.transforms_applied
|
||||
pipeline_timing = result.timing
|
||||
original_tokens = result.tokens_before
|
||||
optimized_tokens = result.tokens_after
|
||||
else:
|
||||
skip_ccr_request_compression = should_skip_ccr_request_compression(
|
||||
frozen_message_count
|
||||
)
|
||||
if skip_ccr_request_compression:
|
||||
logger.info(
|
||||
f"[{request_id}] CCR: skipping request-side compression "
|
||||
f"(frozen prefix={frozen_message_count}) because tool injection is deferred"
|
||||
# Cache newly compressed messages (index-aligned diff)
|
||||
if result.messages != working_messages:
|
||||
comp_cache.update_from_result(messages, result.messages)
|
||||
|
||||
# Always use pipeline result — Zone 1 swaps are already applied
|
||||
optimized_messages = result.messages
|
||||
transforms_applied = result.transforms_applied
|
||||
pipeline_timing = result.timing
|
||||
# Issue #327 / Bug 3: pipeline.apply uses the provider-
|
||||
# side tokenizer (AnthropicProvider tiktoken estimator),
|
||||
# which counts ~25% higher than the proxy-side
|
||||
# EstimatingTokenCounter used to set `original_tokens`
|
||||
# at line 634. Reusing `result.tokens_after` here
|
||||
# produced an apples-vs-oranges comparison against
|
||||
# `original_tokens` in the inflation guard below
|
||||
# (line ~901): even after a real 12% compression the
|
||||
# provider-tokenizer figure was higher than the proxy-
|
||||
# tokenizer baseline, triggering a spurious revert.
|
||||
# Recount optimized_messages with the proxy tokenizer
|
||||
# so original_tokens vs optimized_tokens is self-
|
||||
# consistent. The recount cost (~ms on a 50K-token
|
||||
# request) is paid once per request and is dwarfed by
|
||||
# the upstream call latency.
|
||||
optimized_tokens = tokenizer.count_messages(optimized_messages)
|
||||
elif not is_cache_mode(self.config.mode):
|
||||
async with stage_timer.measure("compression_first_stage"):
|
||||
result = await self._run_compression_in_executor(
|
||||
lambda: self.anthropic_pipeline.apply(
|
||||
messages=messages,
|
||||
model=model,
|
||||
model_limit=context_limit,
|
||||
context=extract_user_query(messages),
|
||||
frozen_message_count=frozen_message_count,
|
||||
biases=biases,
|
||||
request_id=request_id,
|
||||
compression_policy=compression_policy,
|
||||
**proxy_pipeline_kwargs(self.config),
|
||||
),
|
||||
timeout=COMPRESSION_TIMEOUT_SECONDS,
|
||||
)
|
||||
|
||||
if result.messages != messages:
|
||||
optimized_messages = result.messages
|
||||
transforms_applied = result.transforms_applied
|
||||
pipeline_timing = result.timing
|
||||
original_tokens = result.tokens_before
|
||||
optimized_tokens = result.tokens_after
|
||||
else:
|
||||
previous_original_messages = prefix_tracker.get_last_original_messages()
|
||||
previous_forwarded_messages = prefix_tracker.get_last_forwarded_messages()
|
||||
delta = self._extract_cache_stable_delta(
|
||||
|
|
@ -1544,78 +1493,70 @@ class AnthropicHandlerMixin:
|
|||
if delta is not None:
|
||||
stable_forwarded_prefix, delta_messages = delta
|
||||
if delta_messages:
|
||||
if skip_ccr_request_compression:
|
||||
optimized_messages = messages
|
||||
optimized_tokens = tokenizer.count_messages(optimized_messages)
|
||||
else:
|
||||
# Compress the delta, with two cache-mode adjustments:
|
||||
#
|
||||
# fix-5: strip the client's transient cache_control marker so
|
||||
# the router's per-block "never compress an explicit cache
|
||||
# key" guard (content_router.py:4006) doesn't skip the ONLY
|
||||
# compressible content every turn (route_counts had
|
||||
# cache_control_protected == the whole delta -> 0%). In cache
|
||||
# mode that marker is NOT the real forwarded breakpoint: the
|
||||
# compressed delta is frozen + replayed verbatim next turn and
|
||||
# normalize_message_cache_control (AFTER compression, below)
|
||||
# owns the single forwarded breakpoint. Cache-safety is
|
||||
# enforced post-compression, not by protecting the delta.
|
||||
#
|
||||
# fix-6: the delta is a lone tool_result whose tool_use (tool
|
||||
# NAME + call args) lives in the frozen prefix. Passing only
|
||||
# the delta to the router leaves tool_name="" so
|
||||
# _bash_search_fold (lossless grep/rg folding, no size floor),
|
||||
# per-tool bias, and relevance-query enrichment all degrade.
|
||||
# Pass the FULL current messages with frozen_message_count =
|
||||
# prefix length: _build_tool_name_map scans ALL messages (the
|
||||
# delta resolves its tool_name from the prefix's tool_use) but
|
||||
# the compression loop only touches indices >= frozen count,
|
||||
# so ONLY the delta is compressed. Splice the compressed delta
|
||||
# onto the byte-stable forwarded prefix.
|
||||
from headroom.cache.prefix_tracker import _strip_cache_control
|
||||
# Compress the delta, with two cache-mode adjustments:
|
||||
#
|
||||
# fix-5: strip the client's transient cache_control marker so
|
||||
# the router's per-block "never compress an explicit cache
|
||||
# key" guard (content_router.py:4006) doesn't skip the ONLY
|
||||
# compressible content every turn (route_counts had
|
||||
# cache_control_protected == the whole delta -> 0%). In cache
|
||||
# mode that marker is NOT the real forwarded breakpoint: the
|
||||
# compressed delta is frozen + replayed verbatim next turn and
|
||||
# normalize_message_cache_control (AFTER compression, below)
|
||||
# owns the single forwarded breakpoint. Cache-safety is
|
||||
# enforced post-compression, not by protecting the delta.
|
||||
#
|
||||
# fix-6: the delta is a lone tool_result whose tool_use (tool
|
||||
# NAME + call args) lives in the frozen prefix. Passing only
|
||||
# the delta to the router leaves tool_name="" so
|
||||
# _bash_search_fold (lossless grep/rg folding, no size floor),
|
||||
# per-tool bias, and relevance-query enrichment all degrade.
|
||||
# Pass the FULL current messages with frozen_message_count =
|
||||
# prefix length: _build_tool_name_map scans ALL messages (the
|
||||
# delta resolves its tool_name from the prefix's tool_use) but
|
||||
# the compression loop only touches indices >= frozen count,
|
||||
# so ONLY the delta is compressed. Splice the compressed delta
|
||||
# onto the byte-stable forwarded prefix.
|
||||
from headroom.cache.prefix_tracker import _strip_cache_control
|
||||
|
||||
# Compression context = the EXACT forwarded (cached) prefix
|
||||
# + the stripped delta, with the prefix frozen. Using the
|
||||
# forwarded prefix (not the original) keeps _build_tool_name_map
|
||||
# AND cross-turn dedup consistent with what is actually cached:
|
||||
# dedup can only reference bytes that are truly present in the
|
||||
# forwarded context, so no pointer can dangle. The prefix is
|
||||
# frozen (never compressed) and we discard the router's copy of
|
||||
# it below, so the forwarded prefix stays byte-identical to last
|
||||
# turn -> append-only -> no bust.
|
||||
prefix_n = len(stable_forwarded_prefix)
|
||||
compression_input = list(stable_forwarded_prefix) + list(
|
||||
_strip_cache_control(delta_messages)
|
||||
)
|
||||
result = await self._run_compression_in_executor(
|
||||
lambda: self.anthropic_pipeline.apply(
|
||||
messages=compression_input,
|
||||
model=model,
|
||||
model_limit=context_limit,
|
||||
context=extract_user_query(compression_input),
|
||||
frozen_message_count=prefix_n,
|
||||
idle_seconds=idle_seconds,
|
||||
biases=biases,
|
||||
request_id=request_id,
|
||||
compression_policy=compression_policy,
|
||||
**proxy_pipeline_kwargs(self.config),
|
||||
),
|
||||
timeout=COMPRESSION_TIMEOUT_SECONDS,
|
||||
)
|
||||
# Only the delta was eligible for compression (prefix frozen);
|
||||
# forward the byte-identical cached prefix + the compressed delta.
|
||||
compressed_delta = result.messages[prefix_n:]
|
||||
optimized_messages = stable_forwarded_prefix + compressed_delta
|
||||
transforms_applied = result.transforms_applied
|
||||
pipeline_timing = result.timing
|
||||
optimized_tokens = tokenizer.count_messages(optimized_messages)
|
||||
# Compression context = the EXACT forwarded (cached) prefix
|
||||
# + the stripped delta, with the prefix frozen. Using the
|
||||
# forwarded prefix (not the original) keeps _build_tool_name_map
|
||||
# AND cross-turn dedup consistent with what is actually cached:
|
||||
# dedup can only reference bytes that are truly present in the
|
||||
# forwarded context, so no pointer can dangle. The prefix is
|
||||
# frozen (never compressed) and we discard the router's copy of
|
||||
# it below, so the forwarded prefix stays byte-identical to last
|
||||
# turn -> append-only -> no bust.
|
||||
prefix_n = len(stable_forwarded_prefix)
|
||||
compression_input = list(stable_forwarded_prefix) + list(
|
||||
_strip_cache_control(delta_messages)
|
||||
)
|
||||
result = await self._run_compression_in_executor(
|
||||
lambda: self.anthropic_pipeline.apply(
|
||||
messages=compression_input,
|
||||
model=model,
|
||||
model_limit=context_limit,
|
||||
context=extract_user_query(compression_input),
|
||||
frozen_message_count=prefix_n,
|
||||
idle_seconds=idle_seconds,
|
||||
biases=biases,
|
||||
request_id=request_id,
|
||||
compression_policy=compression_policy,
|
||||
**proxy_pipeline_kwargs(self.config),
|
||||
),
|
||||
timeout=COMPRESSION_TIMEOUT_SECONDS,
|
||||
)
|
||||
# Only the delta was eligible for compression (prefix frozen);
|
||||
# forward the byte-identical cached prefix + the compressed delta.
|
||||
compressed_delta = result.messages[prefix_n:]
|
||||
optimized_messages = stable_forwarded_prefix + compressed_delta
|
||||
transforms_applied = result.transforms_applied
|
||||
pipeline_timing = result.timing
|
||||
optimized_tokens = tokenizer.count_messages(optimized_messages)
|
||||
else:
|
||||
if skip_ccr_request_compression:
|
||||
optimized_messages = messages
|
||||
optimized_tokens = tokenizer.count_messages(optimized_messages)
|
||||
else:
|
||||
optimized_messages = stable_forwarded_prefix
|
||||
optimized_tokens = tokenizer.count_messages(optimized_messages)
|
||||
optimized_messages = stable_forwarded_prefix
|
||||
optimized_tokens = tokenizer.count_messages(optimized_messages)
|
||||
else:
|
||||
# Conservative rule for cache mode:
|
||||
# only replay exact stable message-prefix extensions.
|
||||
|
|
|
|||
|
|
@ -490,7 +490,7 @@ def test_existing_retrieve_tool_keeps_reversible_ccr_path_when_prefix_is_frozen(
|
|||
assert [tool["name"] for tool in forwarded["tools"]] == ["headroom_retrieve"]
|
||||
|
||||
|
||||
def test_cache_mode_skip_replays_cached_compressed_prefix_when_tool_injection_is_deferred(
|
||||
def test_cache_mode_compresses_delta_but_replays_cached_prefix_when_markers_are_historical(
|
||||
monkeypatch,
|
||||
) -> None:
|
||||
captured: dict[str, object] = {}
|
||||
|
|
@ -574,13 +574,14 @@ def test_cache_mode_skip_replays_cached_compressed_prefix_when_tool_injection_is
|
|||
)
|
||||
|
||||
assert response.status_code == 200
|
||||
assert captured.get("compression_calls", []) == []
|
||||
assert len(captured.get("compression_calls", [])) == 1
|
||||
forwarded = captured["body"]
|
||||
# Tool injection is deferred (no CCR tool this turn), but the frozen
|
||||
# prefix was cached COMPRESSED last turn. Replay it byte-identical so the
|
||||
# prompt cache still hits instead of busting on original bytes (#1850);
|
||||
# the mutable tail stays original. Tool absent AND cache intact.
|
||||
assert forwarded["messages"] == previous_forwarded_messages + original_messages[1:]
|
||||
# the historical marker does not force tool injection back on. Tool absent
|
||||
# AND cache intact.
|
||||
assert forwarded["messages"] == previous_forwarded_messages
|
||||
assert "tools" not in forwarded
|
||||
|
||||
|
||||
|
|
@ -755,7 +756,7 @@ def test_token_mode_cached_messages_skip_cache_update_when_pipeline_result_is_un
|
|||
assert forwarded["messages"] == marker_messages
|
||||
|
||||
|
||||
def test_non_token_non_cache_mode_still_skips_marker_emission_when_tool_is_unavailable(
|
||||
def test_non_token_non_cache_mode_keeps_compression_and_injects_tool_for_new_markers(
|
||||
monkeypatch,
|
||||
) -> None:
|
||||
captured: dict[str, object] = {}
|
||||
|
|
@ -828,11 +829,16 @@ def test_non_token_non_cache_mode_still_skips_marker_emission_when_tool_is_unava
|
|||
},
|
||||
)
|
||||
|
||||
marker_message = {
|
||||
"role": "user",
|
||||
"content": "[100 items compressed to 10. Retrieve more: hash=abc123def456abc123def456]",
|
||||
}
|
||||
|
||||
assert response.status_code == 200
|
||||
assert captured.get("compression_calls", []) == []
|
||||
assert len(captured.get("compression_calls", [])) == 1
|
||||
forwarded = captured["body"]
|
||||
assert forwarded["messages"] == original_messages
|
||||
assert "tools" not in forwarded
|
||||
assert forwarded["messages"] == [marker_message]
|
||||
assert [tool["name"] for tool in forwarded["tools"]] == ["headroom_retrieve"]
|
||||
|
||||
|
||||
def test_non_token_non_cache_mode_keeps_reversible_path_and_records_waste_signals(
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue