mirror of
https://github.com/headroomlabs-ai/headroom.git
synced 2026-08-27 14:17:10 -04:00
fix: A6 — anthropic-beta and openai-beta deterministic merge + session-sticky
PR-A6 of the Phase A cache-safety lockdown. Eliminates P5-50 and preps
P0-6 (memory tool injection toggling).
Two cache-killer patterns the merge + tracker defeat:
1. Mid-session mutation: when memory was enabled the proxy did an
ad-hoc concat of `context-management-2025-06-27` onto the client
value (anthropic.py:1244-1248). The order varied with the client
value, breaking byte-stable headers across turns.
2. Token drop-out across turns: clients (Claude Code, Codex CLI) MAY
drop a beta token between turn N and turn N+1 even when the proxy
mutated turn N to add it. The cache hot zone is positional, so the
next turn's prefix bytes hash differently and the prefix-cache
read misses.
Changes
-------
`headroom/proxy/helpers.py`
* `merge_anthropic_beta` / `merge_openai_beta`: pure, deterministic,
order-preserving merge. Client tokens first (in their original
order), then Headroom-required tokens (in the order passed). Dedupe
is case-insensitive but preserves the original casing of the first
occurrence. No regex.
* `SessionBetaTracker`: bounded LRU keyed by (provider, session_id),
unioning client tokens with previously-seen tokens. OrderedDict
LRU; threading.RLock for thread safety (mirrors the
CompressionCache pattern from compression_cache.py).
* `get_session_beta_tracker` / `_reset_session_beta_tracker_for_test`
process-wide singleton with test reset.
* `log_beta_header_merge`: structured log per cache-affecting merge.
* Env-var knobs (NO HARDCODES):
- HEADROOM_BETA_HEADER_STICKY=enabled|disabled (default enabled).
- HEADROOM_BETA_TRACKER_MAX_SESSIONS (default 1000).
`headroom/proxy/handlers/anthropic.py`
* After `compute_session_id` (line ~744): record client
`anthropic-beta` against the session tracker, write the sticky
value back into `headers` if changed. Order matters: sticky-merge
FIRST so memory-injection has the canonical baseline.
* Memory-injection site (line ~1244): replace the ad-hoc concat with
`merge_anthropic_beta(headers["anthropic-beta"], required_tokens)`.
`headroom/proxy/handlers/openai.py`
* Chat-completions (line ~360): record/merge `openai-beta`.
* /v1/responses HTTP (line ~1213): compute `_responses_session_id`
and record/merge `openai-beta`.
* /v1/responses WS (line ~1711): replace the ad-hoc absent-only
inject with `merge_openai_beta(sticky, ["responses_websockets=
2026-02-06"])`. Replaces any case-variants of the existing key.
Tests
-----
`tests/test_anthropic_beta_session_sticky.py` (26 tests):
* Pure helper: empty inputs, only-client, only-headroom, ordering,
dedupe casing, deterministic memory-injection order, no-double-
inject when token already present.
* Tracker: sticky-on across turns even when client drops, casing
preservation, provider namespace independence, LRU eviction at
max_sessions, env-var validation (loud failures), thread safety
under 16-thread concurrent access, blank-input rejection.
`tests/test_openai_beta_session_sticky.py` (17 tests):
* Mirror of the anthropic suite for `OpenAI-Beta`.
* Plus WS-specific coverage: sticky-then-merge of
`responses_websockets=2026-02-06` against client baseline.
`tests/test_openai_codex_routing.py`
* Add `session_tracker_store` stub to `_DummyOpenAIHandler` so the
routing tests still exercise the responses HTTP handler now that
it computes a session_id for beta-merge.
Notes
-----
Build constraints honored:
* Configurable: HEADROOM_BETA_HEADER_STICKY,
HEADROOM_BETA_TRACKER_MAX_SESSIONS.
* No regex, no hardcodes (env-var bounds), no fallbacks (disabled
mode is operator opt-in for diagnostics, loud failures on invalid
values).
* Structured tracing log via `log_beta_header_merge`.
Acceptance:
* 43 new tests pass.
* `cargo test --workspace` green (no Rust changes).
* `make ci-precheck` green.
This commit is contained in:
parent
2e874c5e3e
commit
aec5ba3253
10 changed files with 1128 additions and 16 deletions
|
|
@ -5,14 +5,14 @@
|
|||
},
|
||||
"metadata": {
|
||||
"description": "Headroom marketplace for Claude Code and GitHub Copilot CLI plugins.",
|
||||
"version": "0.20.11"
|
||||
"version": "0.20.12"
|
||||
},
|
||||
"plugins": [
|
||||
{
|
||||
"name": "headroom",
|
||||
"source": "./plugins/headroom-agent-hooks",
|
||||
"description": "Headroom startup hooks for Claude Code and GitHub Copilot CLI.",
|
||||
"version": "0.20.11",
|
||||
"version": "0.20.12",
|
||||
"author": {
|
||||
"name": "Headroom Contributors",
|
||||
"url": "https://github.com/chopratejas/headroom"
|
||||
|
|
|
|||
4
.github/plugin/marketplace.json
vendored
4
.github/plugin/marketplace.json
vendored
|
|
@ -5,14 +5,14 @@
|
|||
},
|
||||
"metadata": {
|
||||
"description": "Headroom marketplace for Claude Code and GitHub Copilot CLI plugins.",
|
||||
"version": "0.20.11"
|
||||
"version": "0.20.12"
|
||||
},
|
||||
"plugins": [
|
||||
{
|
||||
"name": "headroom",
|
||||
"source": "./plugins/headroom-agent-hooks",
|
||||
"description": "Headroom startup hooks for Claude Code and GitHub Copilot CLI.",
|
||||
"version": "0.20.11",
|
||||
"version": "0.20.12",
|
||||
"author": {
|
||||
"name": "Headroom Contributors",
|
||||
"url": "https://github.com/chopratejas/headroom"
|
||||
|
|
|
|||
|
|
@ -750,6 +750,51 @@ class AnthropicHandlerMixin:
|
|||
frozen_message_count,
|
||||
)
|
||||
|
||||
# PR-A6 (P5-50, preps P0-6): session-sticky `anthropic-beta` merge.
|
||||
# Read the client's beta value (note: anthropic-beta is NOT
|
||||
# an x-headroom-* header so it survived the A5 strip), union
|
||||
# with previously-seen tokens for this session, and update
|
||||
# the tracker. Memory-injection (below at line ~1244) uses
|
||||
# `merge_anthropic_beta` to add `context-management-2025-06-27`
|
||||
# on top of the sticky baseline. Order matters: session-sticky
|
||||
# FIRST so we have the canonical baseline; memory injection
|
||||
# adds Headroom-required tokens AFTER.
|
||||
from headroom.proxy.helpers import (
|
||||
get_session_beta_tracker,
|
||||
log_beta_header_merge,
|
||||
)
|
||||
|
||||
_client_beta_value = headers.get("anthropic-beta")
|
||||
_client_beta_count = (
|
||||
len([t for t in (_client_beta_value or "").split(",") if t.strip()])
|
||||
if _client_beta_value
|
||||
else 0
|
||||
)
|
||||
_sticky_beta_value = get_session_beta_tracker().record_and_get_sticky_betas(
|
||||
provider="anthropic",
|
||||
session_id=session_id,
|
||||
client_value=_client_beta_value,
|
||||
)
|
||||
_sticky_beta_count = (
|
||||
len([t for t in _sticky_beta_value.split(",") if t.strip()])
|
||||
if _sticky_beta_value
|
||||
else 0
|
||||
)
|
||||
if _sticky_beta_value and _sticky_beta_value != (_client_beta_value or ""):
|
||||
headers["anthropic-beta"] = _sticky_beta_value
|
||||
elif not _sticky_beta_value and "anthropic-beta" in headers:
|
||||
# Sticky value can only equal "" when both client and
|
||||
# session are empty; preserve the (absent) client state.
|
||||
pass
|
||||
log_beta_header_merge(
|
||||
provider="anthropic",
|
||||
session_id=session_id,
|
||||
client_betas_count=_client_beta_count,
|
||||
sticky_betas_count=_sticky_beta_count,
|
||||
headroom_added=[],
|
||||
request_id=request_id,
|
||||
)
|
||||
|
||||
# In cache mode, avoid rewriting any message body bytes. The latest user
|
||||
# turn becomes historical on the next request, so even "latest turn only"
|
||||
# rewrites can invalidate the next cache read when the client resends the
|
||||
|
|
@ -1236,18 +1281,55 @@ class AnthropicHandlerMixin:
|
|||
]
|
||||
logger.info(f"[{request_id}] Memory: Injected tools: {tool_names}")
|
||||
|
||||
# Add beta headers for native memory tool
|
||||
# Add beta headers for native memory tool. PR-A6
|
||||
# (P5-50): use the deterministic `merge_anthropic_beta`
|
||||
# helper instead of ad-hoc string concat. Order:
|
||||
# client tokens first (preserved from session-sticky
|
||||
# baseline above), then Headroom-required tokens.
|
||||
# The session tracker already recorded the client
|
||||
# value; we append Headroom-required tokens here so
|
||||
# the next turn re-applies them deterministically.
|
||||
beta_headers = self.memory_handler.get_beta_headers()
|
||||
if beta_headers:
|
||||
from headroom.proxy.helpers import (
|
||||
log_beta_header_merge as _log_beta_header_merge_mem,
|
||||
)
|
||||
from headroom.proxy.helpers import (
|
||||
merge_anthropic_beta,
|
||||
)
|
||||
|
||||
for key, value in beta_headers.items():
|
||||
# Merge with existing beta header if present
|
||||
existing = headers.get(key, "")
|
||||
if existing and value not in existing:
|
||||
headers[key] = f"{existing},{value}"
|
||||
else:
|
||||
if key.lower() != "anthropic-beta":
|
||||
# Defensive: memory handler currently
|
||||
# only emits anthropic-beta. Any future
|
||||
# provider-specific beta header would
|
||||
# need its own merge helper.
|
||||
headers[key] = value
|
||||
continue
|
||||
existing_value = headers.get(key, "")
|
||||
required_tokens = [t.strip() for t in value.split(",") if t.strip()]
|
||||
merged = merge_anthropic_beta(existing_value, required_tokens)
|
||||
_existing_count = (
|
||||
len([t for t in existing_value.split(",") if t.strip()])
|
||||
if existing_value
|
||||
else 0
|
||||
)
|
||||
_merged_count = (
|
||||
len([t for t in merged.split(",") if t.strip()])
|
||||
if merged
|
||||
else 0
|
||||
)
|
||||
headers[key] = merged
|
||||
_log_beta_header_merge_mem(
|
||||
provider="anthropic",
|
||||
session_id=session_id,
|
||||
client_betas_count=_existing_count,
|
||||
sticky_betas_count=_merged_count,
|
||||
headroom_added=required_tokens,
|
||||
request_id=request_id,
|
||||
)
|
||||
logger.info(
|
||||
f"[{request_id}] Memory: Added beta header: {key}={headers[key]}"
|
||||
f"[{request_id}] Memory: Added beta header: {key}={merged}"
|
||||
)
|
||||
|
||||
if memory_context_injected or memory_tools_injected:
|
||||
|
|
|
|||
|
|
@ -357,6 +357,48 @@ class OpenAIHandlerMixin:
|
|||
openai_prefix_tracker = self.session_tracker_store.get_or_create(
|
||||
openai_session_id, "openai"
|
||||
)
|
||||
|
||||
# PR-A6 (P5-50, preps P0-6): session-sticky `OpenAI-Beta` merge.
|
||||
# Same pattern as anthropic.py — read client value, union with
|
||||
# session-seen tokens, update tracker. WS auto-injection of
|
||||
# `responses_websockets=2026-02-06` lives on the WS handler;
|
||||
# chat-completions has no Headroom-required tokens today, so the
|
||||
# merge effectively just makes the client value byte-stable
|
||||
# across turns.
|
||||
from headroom.proxy.helpers import (
|
||||
get_session_beta_tracker as _get_session_beta_tracker_chat,
|
||||
)
|
||||
from headroom.proxy.helpers import (
|
||||
log_beta_header_merge as _log_beta_header_merge_chat,
|
||||
)
|
||||
|
||||
_client_openai_beta = headers.get("openai-beta")
|
||||
_client_openai_beta_count = (
|
||||
len([t for t in (_client_openai_beta or "").split(",") if t.strip()])
|
||||
if _client_openai_beta
|
||||
else 0
|
||||
)
|
||||
_sticky_openai_beta = _get_session_beta_tracker_chat().record_and_get_sticky_betas(
|
||||
provider="openai",
|
||||
session_id=openai_session_id,
|
||||
client_value=_client_openai_beta,
|
||||
)
|
||||
_sticky_openai_beta_count = (
|
||||
len([t for t in _sticky_openai_beta.split(",") if t.strip()])
|
||||
if _sticky_openai_beta
|
||||
else 0
|
||||
)
|
||||
if _sticky_openai_beta and _sticky_openai_beta != (_client_openai_beta or ""):
|
||||
headers["openai-beta"] = _sticky_openai_beta
|
||||
_log_beta_header_merge_chat(
|
||||
provider="openai",
|
||||
session_id=openai_session_id,
|
||||
client_betas_count=_client_openai_beta_count,
|
||||
sticky_betas_count=_sticky_openai_beta_count,
|
||||
headroom_added=[],
|
||||
request_id=request_id,
|
||||
)
|
||||
|
||||
openai_frozen_count = openai_prefix_tracker.get_frozen_message_count()
|
||||
if is_cache_mode(self.config.mode):
|
||||
openai_frozen_count = self._strict_previous_turn_frozen_count(
|
||||
|
|
@ -1166,6 +1208,45 @@ class OpenAIHandlerMixin:
|
|||
request_id=request_id,
|
||||
)
|
||||
|
||||
# PR-A6 (P5-50, preps P0-6): session-sticky `OpenAI-Beta` merge
|
||||
# for /v1/responses. Compute a session_id off the same store the
|
||||
# chat handler uses so multi-endpoint clients within one
|
||||
# conversation share the sticky-token set.
|
||||
_responses_session_id = self.session_tracker_store.compute_session_id(
|
||||
request, model, messages
|
||||
)
|
||||
from headroom.proxy.helpers import (
|
||||
get_session_beta_tracker as _get_session_beta_tracker_resp,
|
||||
)
|
||||
from headroom.proxy.helpers import (
|
||||
log_beta_header_merge as _log_beta_header_merge_resp,
|
||||
)
|
||||
|
||||
_client_resp_beta = headers.get("openai-beta")
|
||||
_client_resp_beta_count = (
|
||||
len([t for t in (_client_resp_beta or "").split(",") if t.strip()])
|
||||
if _client_resp_beta
|
||||
else 0
|
||||
)
|
||||
_sticky_resp_beta = _get_session_beta_tracker_resp().record_and_get_sticky_betas(
|
||||
provider="openai",
|
||||
session_id=_responses_session_id,
|
||||
client_value=_client_resp_beta,
|
||||
)
|
||||
_sticky_resp_beta_count = (
|
||||
len([t for t in _sticky_resp_beta.split(",") if t.strip()]) if _sticky_resp_beta else 0
|
||||
)
|
||||
if _sticky_resp_beta and _sticky_resp_beta != (_client_resp_beta or ""):
|
||||
headers["openai-beta"] = _sticky_resp_beta
|
||||
_log_beta_header_merge_resp(
|
||||
provider="openai",
|
||||
session_id=_responses_session_id,
|
||||
client_betas_count=_client_resp_beta_count,
|
||||
sticky_betas_count=_sticky_resp_beta_count,
|
||||
headroom_added=[],
|
||||
request_id=request_id,
|
||||
)
|
||||
|
||||
# Memory: Get user ID when memory is enabled. Reads `request.headers`
|
||||
# directly because `headers` was stripped of `x-headroom-*` (PR-A5).
|
||||
memory_user_id: str | None = None
|
||||
|
|
@ -1707,9 +1788,59 @@ class OpenAIHandlerMixin:
|
|||
upstream_headers = await apply_copilot_api_auth(upstream_headers, url=upstream_url)
|
||||
|
||||
# Ensure the required beta header is present — OpenAI returns 500 without it.
|
||||
# Codex sends `responses_websockets=2026-02-06`; only inject if missing entirely.
|
||||
if "openai-beta" not in _lower_headers:
|
||||
upstream_headers["OpenAI-Beta"] = "responses_websockets=2026-02-06"
|
||||
# PR-A6 (P5-50): use the deterministic `merge_openai_beta` helper
|
||||
# so the auto-injected `responses_websockets=2026-02-06` is
|
||||
# appended to the client's value (preserving order, deduping
|
||||
# case-insensitively) rather than overwriting it. The
|
||||
# SessionBetaTracker also records the merge so a future cross-
|
||||
# connection sticky model can replay tokens by session_id.
|
||||
from headroom.proxy.helpers import (
|
||||
get_session_beta_tracker as _get_session_beta_tracker_ws,
|
||||
)
|
||||
from headroom.proxy.helpers import (
|
||||
log_beta_header_merge as _log_beta_header_merge_ws,
|
||||
)
|
||||
from headroom.proxy.helpers import merge_openai_beta as _merge_openai_beta_ws
|
||||
|
||||
_ws_required_tokens = ["responses_websockets=2026-02-06"]
|
||||
# Read the original (pre-merge) client value from the WS headers
|
||||
# to preserve casing and ordering.
|
||||
_ws_client_beta_value: str | None = None
|
||||
for _k, _v in upstream_headers.items():
|
||||
if _k.lower() == "openai-beta":
|
||||
_ws_client_beta_value = _v
|
||||
break
|
||||
# Record session-stickiness BEFORE adding required tokens so the
|
||||
# tracker stores the canonical client baseline.
|
||||
_ws_sticky_beta = _get_session_beta_tracker_ws().record_and_get_sticky_betas(
|
||||
provider="openai",
|
||||
session_id=session_id,
|
||||
client_value=_ws_client_beta_value,
|
||||
)
|
||||
_ws_merged_beta = _merge_openai_beta_ws(_ws_sticky_beta, _ws_required_tokens)
|
||||
# Replace any existing case-variants of openai-beta with the
|
||||
# canonical "OpenAI-Beta" key carrying the merged value.
|
||||
_ws_existing_keys = [_k for _k in upstream_headers if _k.lower() == "openai-beta"]
|
||||
for _k in _ws_existing_keys:
|
||||
del upstream_headers[_k]
|
||||
if _ws_merged_beta:
|
||||
upstream_headers["OpenAI-Beta"] = _ws_merged_beta
|
||||
_ws_client_beta_count = (
|
||||
len([t for t in (_ws_client_beta_value or "").split(",") if t.strip()])
|
||||
if _ws_client_beta_value
|
||||
else 0
|
||||
)
|
||||
_ws_merged_beta_count = (
|
||||
len([t for t in _ws_merged_beta.split(",") if t.strip()]) if _ws_merged_beta else 0
|
||||
)
|
||||
_log_beta_header_merge_ws(
|
||||
provider="openai",
|
||||
session_id=session_id,
|
||||
client_betas_count=_ws_client_beta_count,
|
||||
sticky_betas_count=_ws_merged_beta_count,
|
||||
headroom_added=_ws_required_tokens,
|
||||
request_id=request_id,
|
||||
)
|
||||
|
||||
logger.debug(
|
||||
f"[{request_id}] WS upstream headers: "
|
||||
|
|
|
|||
|
|
@ -15,6 +15,7 @@ import os
|
|||
import random
|
||||
import threading
|
||||
import time
|
||||
from collections import OrderedDict
|
||||
from pathlib import Path
|
||||
from typing import TYPE_CHECKING, Any, Literal, cast
|
||||
|
||||
|
|
@ -655,6 +656,325 @@ def log_outbound_headers(
|
|||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Beta-header merge + per-session stickiness (PR-A6 — fixes P5-50; preps P0-6).
|
||||
# ---------------------------------------------------------------------------
|
||||
#
|
||||
# Anthropic's `anthropic-beta` and OpenAI's `OpenAI-Beta` request headers
|
||||
# carry a comma-separated list of opt-in beta tokens. Two cache-killer
|
||||
# patterns motivated PR-A6:
|
||||
#
|
||||
# 1. Mid-session mutation: when memory is enabled the proxy historically
|
||||
# did an ad-hoc concat of `context-management-2025-06-27` onto the
|
||||
# client value (anthropic.py:1244-1248) — every variant produced a
|
||||
# different byte sequence and the order was undefined when the same
|
||||
# client value already contained a Headroom-required token.
|
||||
#
|
||||
# 2. Token drop-out across turns: clients (Claude Code, Codex CLI) MAY
|
||||
# drop a beta token between turn N and turn N+1 even when the proxy
|
||||
# mutated turn N to add it. The cache hot zone is positional, so the
|
||||
# next turn's prefix bytes hash differently and the prefix-cache
|
||||
# read misses.
|
||||
#
|
||||
# PR-A6 introduces:
|
||||
# * `merge_anthropic_beta` / `merge_openai_beta`: deterministic, pure,
|
||||
# order-preserving merge. Client tokens first (in their original order),
|
||||
# then Headroom-required tokens (in the order passed). Dedupe is
|
||||
# case-insensitive but preserves original casing of first occurrence.
|
||||
# Per Anthropic guide §6.3 #6: sticky-on means we add but never reorder.
|
||||
#
|
||||
# * `SessionBetaTracker`: bounded LRU cache keyed by `(provider,
|
||||
# session_id)` tracking every beta token observed for that session.
|
||||
# On every request we union the client value with previously-seen
|
||||
# tokens and update the seen set — so a beta seen in turn N is
|
||||
# present in turn N+1 even if the client drops it. LRU bound (default
|
||||
# 1000 sessions) prevents unbounded growth. Reentrant lock so future
|
||||
# callers from inside another locked method don't self-deadlock.
|
||||
#
|
||||
# Operator opt-in `HEADROOM_BETA_HEADER_STICKY=disabled` short-circuits
|
||||
# the tracker (returns the client value verbatim). That mode is loud and
|
||||
# explicit per realignment build constraint #4 — NOT a silent fallback.
|
||||
|
||||
_BETA_HEADER_STICKY_ENV = "HEADROOM_BETA_HEADER_STICKY"
|
||||
BetaHeaderStickyMode = Literal["enabled", "disabled"]
|
||||
_BETA_HEADER_STICKY_DEFAULT: BetaHeaderStickyMode = "enabled"
|
||||
|
||||
_BETA_TRACKER_MAX_SESSIONS_ENV = "HEADROOM_BETA_TRACKER_MAX_SESSIONS"
|
||||
_BETA_TRACKER_MAX_SESSIONS_DEFAULT = 1000
|
||||
|
||||
|
||||
def get_beta_header_sticky_mode() -> BetaHeaderStickyMode:
|
||||
"""Return the active beta-header stickiness mode.
|
||||
|
||||
Read at request time so operators can flip behaviour without a
|
||||
restart. Unknown values raise loudly per the no-silent-fallback
|
||||
build constraint.
|
||||
"""
|
||||
raw = os.environ.get(_BETA_HEADER_STICKY_ENV, "").strip().lower()
|
||||
if not raw:
|
||||
return _BETA_HEADER_STICKY_DEFAULT
|
||||
if raw in ("enabled", "disabled"):
|
||||
return cast(BetaHeaderStickyMode, raw)
|
||||
raise ValueError(f"Invalid {_BETA_HEADER_STICKY_ENV}={raw!r}; expected 'enabled' or 'disabled'")
|
||||
|
||||
|
||||
def get_beta_tracker_max_sessions() -> int:
|
||||
"""Return the LRU bound for `SessionBetaTracker` (sessions cap)."""
|
||||
raw = os.environ.get(_BETA_TRACKER_MAX_SESSIONS_ENV, "").strip()
|
||||
if not raw:
|
||||
return _BETA_TRACKER_MAX_SESSIONS_DEFAULT
|
||||
try:
|
||||
value = int(raw)
|
||||
except ValueError as exc:
|
||||
raise ValueError(
|
||||
f"Invalid {_BETA_TRACKER_MAX_SESSIONS_ENV}={raw!r}; expected positive int"
|
||||
) from exc
|
||||
if value <= 0:
|
||||
raise ValueError(f"Invalid {_BETA_TRACKER_MAX_SESSIONS_ENV}={raw!r}; expected positive int")
|
||||
return value
|
||||
|
||||
|
||||
def _split_beta_tokens(value: str | None) -> list[str]:
|
||||
"""Split a comma-separated beta-header value into trimmed tokens.
|
||||
|
||||
Empty/whitespace-only entries are dropped. Pure function, no regex.
|
||||
"""
|
||||
if not value:
|
||||
return []
|
||||
out: list[str] = []
|
||||
for raw in value.split(","):
|
||||
token = raw.strip()
|
||||
if token:
|
||||
out.append(token)
|
||||
return out
|
||||
|
||||
|
||||
def _merge_beta_tokens(client_value: str | None, headroom_required: list[str]) -> str:
|
||||
"""Shared deterministic merge for `anthropic-beta` / `OpenAI-Beta` tokens.
|
||||
|
||||
Rules (per Anthropic guide §6.3 #6 "sticky-on; add but never reorder"):
|
||||
|
||||
* Client tokens come first, in their original order.
|
||||
* Headroom-required tokens append in the order given, skipping any
|
||||
token already present (case-insensitive).
|
||||
* Dedupe is case-insensitive but the FIRST occurrence's casing wins
|
||||
(prevents drift when client uses one casing across turns).
|
||||
* Returns ``""`` when both inputs are empty.
|
||||
|
||||
Pure function. No regex. No global state.
|
||||
"""
|
||||
seen_lower: set[str] = set()
|
||||
out: list[str] = []
|
||||
for token in _split_beta_tokens(client_value):
|
||||
lower = token.lower()
|
||||
if lower in seen_lower:
|
||||
continue
|
||||
seen_lower.add(lower)
|
||||
out.append(token)
|
||||
for token in headroom_required:
|
||||
if not token:
|
||||
continue
|
||||
token = token.strip()
|
||||
if not token:
|
||||
continue
|
||||
lower = token.lower()
|
||||
if lower in seen_lower:
|
||||
continue
|
||||
seen_lower.add(lower)
|
||||
out.append(token)
|
||||
return ",".join(out)
|
||||
|
||||
|
||||
def merge_anthropic_beta(client_value: str | None, headroom_required: list[str]) -> str:
|
||||
"""Merge client `anthropic-beta` value with Headroom-required tokens.
|
||||
|
||||
See `_merge_beta_tokens` for full semantics. Order is deterministic:
|
||||
client tokens first (in their original order), then headroom tokens
|
||||
(in the order passed). No sorting — sticky-on per Anthropic guide
|
||||
§6.3 #6 means we add but never reorder. Dedupe is case-insensitive
|
||||
but preserves the original casing of the first occurrence.
|
||||
|
||||
Returns ``""`` when both inputs are empty.
|
||||
"""
|
||||
return _merge_beta_tokens(client_value, headroom_required)
|
||||
|
||||
|
||||
def merge_openai_beta(client_value: str | None, headroom_required: list[str]) -> str:
|
||||
"""Merge client `OpenAI-Beta` value with Headroom-required tokens.
|
||||
|
||||
Mirror of `merge_anthropic_beta`. Same semantics — the OpenAI header
|
||||
follows the same comma-separated convention and the same cache-stable
|
||||
rules apply.
|
||||
"""
|
||||
return _merge_beta_tokens(client_value, headroom_required)
|
||||
|
||||
|
||||
class SessionBetaTracker:
|
||||
"""Bounded LRU tracker of beta-header tokens observed per (provider, session).
|
||||
|
||||
On every request:
|
||||
* Read the client's beta-header value.
|
||||
* Union with previously-seen tokens for this session (sticky-on).
|
||||
* Update the session's seen set.
|
||||
* Return the union (preserving first-seen order).
|
||||
|
||||
Bounded by `max_sessions` (default 1000) via `OrderedDict` LRU
|
||||
eviction: hits move-to-end; overflow pops oldest. Reentrant lock so
|
||||
future callers from inside another locked method don't self-deadlock
|
||||
(mirrors `CompressionCache` pattern).
|
||||
|
||||
The tracker is provider-aware: the same `session_id` for Anthropic
|
||||
and OpenAI keeps independent token sets (clients/upstreams differ on
|
||||
which tokens are valid).
|
||||
"""
|
||||
|
||||
def __init__(self, max_sessions: int | None = None) -> None:
|
||||
if max_sessions is None:
|
||||
max_sessions = get_beta_tracker_max_sessions()
|
||||
if max_sessions <= 0:
|
||||
raise ValueError("max_sessions must be > 0")
|
||||
self._max_sessions: int = max_sessions
|
||||
# OrderedDict per `compression_cache.py` LRU pattern. Entries
|
||||
# store the per-session ordered token list (preserving first-seen
|
||||
# order). RLock allows future callers from inside another locked
|
||||
# method to enter without self-deadlock.
|
||||
self._lock = threading.RLock()
|
||||
self._sessions: OrderedDict[tuple[str, str], list[str]] = OrderedDict()
|
||||
|
||||
@property
|
||||
def active_sessions(self) -> int:
|
||||
with self._lock:
|
||||
return len(self._sessions)
|
||||
|
||||
def _key(self, provider: str, session_id: str) -> tuple[str, str]:
|
||||
return (provider, session_id)
|
||||
|
||||
def record_and_get_sticky_betas(
|
||||
self,
|
||||
provider: str,
|
||||
session_id: str,
|
||||
client_value: str | None,
|
||||
) -> str:
|
||||
"""Union client tokens with session-seen tokens; update; return.
|
||||
|
||||
``provider`` is the upstream identifier (``anthropic`` /
|
||||
``openai``). ``session_id`` is the proxy's per-conversation ID
|
||||
(e.g. `SessionTrackerStore.compute_session_id` output for the
|
||||
HTTP path; the WS handler's per-connection UUID for the WS
|
||||
path — note WS sessions are short-lived and won't accumulate
|
||||
cross-turn).
|
||||
|
||||
When `HEADROOM_BETA_HEADER_STICKY=disabled` returns the client
|
||||
value verbatim (operator diagnostic opt-in; documented as a
|
||||
per-deploy choice, NOT a silent fallback).
|
||||
|
||||
Returns the merged comma-separated value (possibly empty).
|
||||
"""
|
||||
if not provider:
|
||||
raise ValueError("provider must be non-empty")
|
||||
if not session_id:
|
||||
raise ValueError("session_id must be non-empty")
|
||||
|
||||
if get_beta_header_sticky_mode() == "disabled":
|
||||
# Diagnostic mode — return the client value verbatim, do not
|
||||
# touch tracker state. This is loud (operators read the env
|
||||
# var) and per-deploy.
|
||||
return (client_value or "").strip()
|
||||
|
||||
client_tokens = _split_beta_tokens(client_value)
|
||||
key = self._key(provider, session_id)
|
||||
|
||||
with self._lock:
|
||||
previous = self._sessions.get(key)
|
||||
if previous is None:
|
||||
merged_list: list[str] = []
|
||||
seen_lower: set[str] = set()
|
||||
else:
|
||||
# Move-to-end on hit (LRU touch).
|
||||
self._sessions.move_to_end(key)
|
||||
merged_list = list(previous)
|
||||
seen_lower = {t.lower() for t in merged_list}
|
||||
|
||||
# Append client tokens preserving order; first-seen casing wins.
|
||||
for token in client_tokens:
|
||||
lower = token.lower()
|
||||
if lower in seen_lower:
|
||||
continue
|
||||
seen_lower.add(lower)
|
||||
merged_list.append(token)
|
||||
|
||||
self._sessions[key] = merged_list
|
||||
self._sessions.move_to_end(key)
|
||||
|
||||
# Bound: evict oldest until at-or-below cap.
|
||||
while len(self._sessions) > self._max_sessions:
|
||||
self._sessions.popitem(last=False)
|
||||
|
||||
return ",".join(merged_list)
|
||||
|
||||
def reset(self) -> None:
|
||||
"""Clear all session state (test helper)."""
|
||||
with self._lock:
|
||||
self._sessions.clear()
|
||||
|
||||
|
||||
# Process-wide singleton. Lazily replaced by tests via `reset` /
|
||||
# `_reset_session_beta_tracker_for_test`. One tracker for both providers
|
||||
# — the (provider, session_id) key keeps namespaces independent.
|
||||
_session_beta_tracker_lock = threading.Lock()
|
||||
_session_beta_tracker: SessionBetaTracker | None = None
|
||||
|
||||
|
||||
def get_session_beta_tracker() -> SessionBetaTracker:
|
||||
"""Return the process-wide `SessionBetaTracker` singleton.
|
||||
|
||||
Lazily constructed so the env-var bound (`HEADROOM_BETA_TRACKER_MAX_SESSIONS`)
|
||||
is honored at first use. Tests use `_reset_session_beta_tracker_for_test`.
|
||||
"""
|
||||
global _session_beta_tracker
|
||||
with _session_beta_tracker_lock:
|
||||
if _session_beta_tracker is None:
|
||||
_session_beta_tracker = SessionBetaTracker()
|
||||
return _session_beta_tracker
|
||||
|
||||
|
||||
def _reset_session_beta_tracker_for_test() -> None:
|
||||
"""Clear the process-wide tracker (test-only)."""
|
||||
global _session_beta_tracker
|
||||
with _session_beta_tracker_lock:
|
||||
_session_beta_tracker = None
|
||||
|
||||
|
||||
def log_beta_header_merge(
|
||||
*,
|
||||
provider: str,
|
||||
session_id: str | None,
|
||||
client_betas_count: int,
|
||||
sticky_betas_count: int,
|
||||
headroom_added: list[str],
|
||||
request_id: str | None,
|
||||
) -> None:
|
||||
"""Structured log for every cache-affecting beta-header merge.
|
||||
|
||||
`headroom_added` is a list of public, documented beta tokens
|
||||
(e.g. ``context-management-2025-06-27``,
|
||||
``responses_websockets=2026-02-06``) — safe to log. We intentionally
|
||||
do NOT log the raw client value because beta tokens, while public,
|
||||
can carry experiment IDs the user has not opted to share with
|
||||
Headroom logs. Emitting counts only makes the decision auditable.
|
||||
"""
|
||||
logger.info(
|
||||
"event=beta_header_merge provider=%s session_id=%s "
|
||||
"client_betas=%d sticky_betas=%d headroom_added=%s request_id=%s",
|
||||
provider,
|
||||
session_id or "",
|
||||
client_betas_count,
|
||||
sticky_betas_count,
|
||||
",".join(headroom_added) if headroom_added else "",
|
||||
request_id or "",
|
||||
)
|
||||
|
||||
|
||||
async def _read_request_body_bytes(request: Request) -> bytes:
|
||||
"""Read and (if needed) decompress the request body, returning raw UTF-8 bytes.
|
||||
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
{
|
||||
"name": "headroom",
|
||||
"version": "0.20.11",
|
||||
"version": "0.20.12",
|
||||
"description": "Headroom startup hooks for Claude Code and GitHub Copilot CLI.",
|
||||
"author": {
|
||||
"name": "Headroom Contributors",
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
{
|
||||
"name": "headroom",
|
||||
"version": "0.20.11",
|
||||
"version": "0.20.12",
|
||||
"description": "Headroom startup hooks for Claude Code and GitHub Copilot CLI.",
|
||||
"author": {
|
||||
"name": "Headroom Contributors",
|
||||
|
|
|
|||
346
tests/test_anthropic_beta_session_sticky.py
Normal file
346
tests/test_anthropic_beta_session_sticky.py
Normal file
|
|
@ -0,0 +1,346 @@
|
|||
"""Session-sticky `anthropic-beta` tests for PR-A6 (P5-50, preps P0-6).
|
||||
|
||||
Two cache-killer patterns the merge + tracker must defeat:
|
||||
|
||||
1. Mid-session mutation: when memory is enabled the proxy historically
|
||||
did an ad-hoc concat of `context-management-2025-06-27` onto the
|
||||
client value (anthropic.py:1244-1248). The order varied with the
|
||||
client value, breaking byte-stable headers across turns.
|
||||
|
||||
2. Token drop-out across turns: clients (Claude Code, Codex CLI) MAY
|
||||
drop a beta token between turn N and turn N+1 even when the proxy
|
||||
mutated turn N to add it. The cache hot zone is positional, so the
|
||||
next turn's prefix bytes hash differently and the prefix-cache
|
||||
read misses.
|
||||
|
||||
The fix:
|
||||
|
||||
* `merge_anthropic_beta`: deterministic, pure, order-preserving.
|
||||
* `SessionBetaTracker`: bounded LRU keyed by (provider, session_id),
|
||||
unioning client tokens with previously-seen tokens.
|
||||
|
||||
Operator opt-in `HEADROOM_BETA_HEADER_STICKY=disabled` short-circuits
|
||||
the tracker (returns the client value verbatim). That mode is loud and
|
||||
explicit per realignment build constraint #4 — NOT a silent fallback.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import threading
|
||||
|
||||
import pytest
|
||||
|
||||
from headroom.proxy.helpers import (
|
||||
SessionBetaTracker,
|
||||
_reset_session_beta_tracker_for_test,
|
||||
get_beta_header_sticky_mode,
|
||||
get_beta_tracker_max_sessions,
|
||||
get_session_beta_tracker,
|
||||
merge_anthropic_beta,
|
||||
)
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Pure helper: `merge_anthropic_beta`
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def test_merge_helper_pure_function() -> None:
|
||||
"""Same inputs → same output, no global state."""
|
||||
result1 = merge_anthropic_beta("a,b", ["c"])
|
||||
result2 = merge_anthropic_beta("a,b", ["c"])
|
||||
assert result1 == result2 == "a,b,c"
|
||||
|
||||
|
||||
def test_merge_helper_empty_inputs_returns_empty_string() -> None:
|
||||
assert merge_anthropic_beta(None, []) == ""
|
||||
assert merge_anthropic_beta("", []) == ""
|
||||
assert merge_anthropic_beta(" ", []) == ""
|
||||
|
||||
|
||||
def test_merge_helper_only_client() -> None:
|
||||
assert merge_anthropic_beta("a,b", []) == "a,b"
|
||||
|
||||
|
||||
def test_merge_helper_only_headroom() -> None:
|
||||
assert merge_anthropic_beta(None, ["a", "b"]) == "a,b"
|
||||
|
||||
|
||||
def test_merge_helper_preserves_client_order_and_appends_headroom() -> None:
|
||||
# Client tokens FIRST in their original order, headroom AFTER in passed order.
|
||||
assert (
|
||||
merge_anthropic_beta("client-1,client-2", ["headroom-a", "headroom-b"])
|
||||
== "client-1,client-2,headroom-a,headroom-b"
|
||||
)
|
||||
|
||||
|
||||
def test_dedupe_case_insensitive_preserves_first_casing() -> None:
|
||||
# Client casing wins over headroom casing.
|
||||
assert (
|
||||
merge_anthropic_beta("Context-Management-2025-06-27", ["context-management-2025-06-27"])
|
||||
== "Context-Management-2025-06-27"
|
||||
)
|
||||
# Within client list: first occurrence's casing wins.
|
||||
assert merge_anthropic_beta("Foo,foo", []) == "Foo"
|
||||
|
||||
|
||||
def test_test_memory_injection_appends_deterministic_order() -> None:
|
||||
"""Memory injection appends `context-management-2025-06-27` AFTER client tokens."""
|
||||
merged = merge_anthropic_beta(
|
||||
"interleaved-thinking-2025-05-14",
|
||||
["context-management-2025-06-27"],
|
||||
)
|
||||
assert merged == "interleaved-thinking-2025-05-14,context-management-2025-06-27"
|
||||
|
||||
|
||||
def test_merge_helper_skips_empty_tokens() -> None:
|
||||
# Whitespace-only entries are dropped.
|
||||
assert merge_anthropic_beta("a, ,b", []) == "a,b"
|
||||
assert merge_anthropic_beta("a", ["", "b", " "]) == "a,b"
|
||||
|
||||
|
||||
def test_merge_helper_no_double_inject_when_already_present() -> None:
|
||||
# Headroom token already in client value → not re-appended.
|
||||
assert (
|
||||
merge_anthropic_beta("context-management-2025-06-27", ["context-management-2025-06-27"])
|
||||
== "context-management-2025-06-27"
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# `SessionBetaTracker`
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _isolate_tracker(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Reset the process-wide tracker singleton + env flags between tests."""
|
||||
monkeypatch.delenv("HEADROOM_BETA_HEADER_STICKY", raising=False)
|
||||
monkeypatch.delenv("HEADROOM_BETA_TRACKER_MAX_SESSIONS", raising=False)
|
||||
_reset_session_beta_tracker_for_test()
|
||||
yield
|
||||
_reset_session_beta_tracker_for_test()
|
||||
|
||||
|
||||
def test_beta_seen_turn_1_present_in_turn_2_even_if_client_drops() -> None:
|
||||
"""The core sticky-on guarantee — token observed in turn 1 stays in turn 2."""
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
|
||||
# Turn 1: client sends two tokens.
|
||||
out1 = tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic",
|
||||
session_id="s-1",
|
||||
client_value="prompt-caching-2024-07-31,interleaved-thinking-2025-05-14",
|
||||
)
|
||||
assert out1 == "prompt-caching-2024-07-31,interleaved-thinking-2025-05-14"
|
||||
|
||||
# Turn 2: client drops `interleaved-thinking-2025-05-14` entirely.
|
||||
out2 = tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic",
|
||||
session_id="s-1",
|
||||
client_value="prompt-caching-2024-07-31",
|
||||
)
|
||||
# Sticky-on: token survives.
|
||||
assert out2 == "prompt-caching-2024-07-31,interleaved-thinking-2025-05-14"
|
||||
|
||||
|
||||
def test_client_value_preserved_when_no_injection() -> None:
|
||||
"""First-turn client value flows through unchanged when no headroom adds."""
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
out = tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic",
|
||||
session_id="s-1",
|
||||
client_value="alpha,beta",
|
||||
)
|
||||
assert out == "alpha,beta"
|
||||
|
||||
|
||||
def test_empty_client_value_returns_empty_string() -> None:
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
assert (
|
||||
tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic", session_id="s-1", client_value=None
|
||||
)
|
||||
== ""
|
||||
)
|
||||
assert (
|
||||
tracker.record_and_get_sticky_betas(provider="anthropic", session_id="s-1", client_value="")
|
||||
== ""
|
||||
)
|
||||
|
||||
|
||||
def test_provider_namespaces_are_independent() -> None:
|
||||
"""Same session_id under two providers keeps independent token sets."""
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic", session_id="shared", client_value="a-only"
|
||||
)
|
||||
tracker.record_and_get_sticky_betas(
|
||||
provider="openai", session_id="shared", client_value="o-only"
|
||||
)
|
||||
a_out = tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic", session_id="shared", client_value=None
|
||||
)
|
||||
o_out = tracker.record_and_get_sticky_betas(
|
||||
provider="openai", session_id="shared", client_value=None
|
||||
)
|
||||
assert a_out == "a-only"
|
||||
assert o_out == "o-only"
|
||||
|
||||
|
||||
def test_first_seen_casing_preserved_across_turns() -> None:
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic", session_id="s-1", client_value="Context-Management-2025-06-27"
|
||||
)
|
||||
out = tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic", session_id="s-1", client_value="context-management-2025-06-27"
|
||||
)
|
||||
# Original casing wins.
|
||||
assert out == "Context-Management-2025-06-27"
|
||||
|
||||
|
||||
def test_lru_eviction_at_max_sessions() -> None:
|
||||
"""Bounded LRU pops the oldest session when overflowing."""
|
||||
tracker = SessionBetaTracker(max_sessions=2)
|
||||
tracker.record_and_get_sticky_betas(provider="anthropic", session_id="s-1", client_value="a")
|
||||
tracker.record_and_get_sticky_betas(provider="anthropic", session_id="s-2", client_value="b")
|
||||
assert tracker.active_sessions == 2
|
||||
|
||||
# Touch s-1 so s-2 becomes the LRU.
|
||||
tracker.record_and_get_sticky_betas(provider="anthropic", session_id="s-1", client_value=None)
|
||||
|
||||
# Add s-3: pops s-2 (oldest by recent access).
|
||||
tracker.record_and_get_sticky_betas(provider="anthropic", session_id="s-3", client_value="c")
|
||||
assert tracker.active_sessions == 2
|
||||
|
||||
# s-2 was evicted: a fresh request shows no carry-over.
|
||||
out = tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic", session_id="s-2", client_value=None
|
||||
)
|
||||
assert out == ""
|
||||
|
||||
|
||||
def test_max_sessions_invalid_raises() -> None:
|
||||
with pytest.raises(ValueError):
|
||||
SessionBetaTracker(max_sessions=0)
|
||||
with pytest.raises(ValueError):
|
||||
SessionBetaTracker(max_sessions=-1)
|
||||
|
||||
|
||||
def test_disabled_mode_passes_through(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""`HEADROOM_BETA_HEADER_STICKY=disabled` returns client value verbatim."""
|
||||
monkeypatch.setenv("HEADROOM_BETA_HEADER_STICKY", "disabled")
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
|
||||
# Turn 1: record a token.
|
||||
out1 = tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic", session_id="s-1", client_value="alpha"
|
||||
)
|
||||
# Disabled mode: client value flows through verbatim, no tracker update.
|
||||
assert out1 == "alpha"
|
||||
|
||||
# Turn 2: client drops the token.
|
||||
out2 = tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic", session_id="s-1", client_value=None
|
||||
)
|
||||
# Disabled mode: empty client value → empty result; NOT sticky.
|
||||
assert out2 == ""
|
||||
|
||||
|
||||
def test_disabled_mode_invalid_value_raises(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Unknown values raise loudly — no silent fallback."""
|
||||
monkeypatch.setenv("HEADROOM_BETA_HEADER_STICKY", "yolo")
|
||||
with pytest.raises(ValueError, match="HEADROOM_BETA_HEADER_STICKY"):
|
||||
get_beta_header_sticky_mode()
|
||||
|
||||
|
||||
def test_max_sessions_env_var_default() -> None:
|
||||
assert get_beta_tracker_max_sessions() == 1000
|
||||
|
||||
|
||||
def test_max_sessions_env_var_custom(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setenv("HEADROOM_BETA_TRACKER_MAX_SESSIONS", "42")
|
||||
assert get_beta_tracker_max_sessions() == 42
|
||||
|
||||
|
||||
def test_max_sessions_env_var_invalid_raises(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setenv("HEADROOM_BETA_TRACKER_MAX_SESSIONS", "0")
|
||||
with pytest.raises(ValueError):
|
||||
get_beta_tracker_max_sessions()
|
||||
monkeypatch.setenv("HEADROOM_BETA_TRACKER_MAX_SESSIONS", "-3")
|
||||
with pytest.raises(ValueError):
|
||||
get_beta_tracker_max_sessions()
|
||||
monkeypatch.setenv("HEADROOM_BETA_TRACKER_MAX_SESSIONS", "not-int")
|
||||
with pytest.raises(ValueError):
|
||||
get_beta_tracker_max_sessions()
|
||||
|
||||
|
||||
def test_thread_safe_concurrent_access() -> None:
|
||||
"""Spawn N threads on same session_id; assert no exceptions and final state correct."""
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
n_threads = 16
|
||||
iterations = 50
|
||||
errors: list[BaseException] = []
|
||||
|
||||
def worker(thread_idx: int) -> None:
|
||||
try:
|
||||
for i in range(iterations):
|
||||
tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic",
|
||||
session_id="shared",
|
||||
client_value=f"t{thread_idx}-i{i}",
|
||||
)
|
||||
except BaseException as e: # noqa: BLE001
|
||||
errors.append(e)
|
||||
|
||||
threads = [threading.Thread(target=worker, args=(idx,)) for idx in range(n_threads)]
|
||||
for t in threads:
|
||||
t.start()
|
||||
for t in threads:
|
||||
t.join()
|
||||
|
||||
assert errors == []
|
||||
# Final state contains every thread's contributed token.
|
||||
final = tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic", session_id="shared", client_value=None
|
||||
)
|
||||
final_tokens = set(final.split(","))
|
||||
expected = {f"t{idx}-i{i}" for idx in range(n_threads) for i in range(iterations)}
|
||||
assert expected.issubset(final_tokens)
|
||||
|
||||
|
||||
def test_blank_provider_or_session_raises() -> None:
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
with pytest.raises(ValueError):
|
||||
tracker.record_and_get_sticky_betas(provider="", session_id="s", client_value="a")
|
||||
with pytest.raises(ValueError):
|
||||
tracker.record_and_get_sticky_betas(provider="anthropic", session_id="", client_value="a")
|
||||
|
||||
|
||||
def test_singleton_get_session_beta_tracker_returns_same_instance() -> None:
|
||||
a = get_session_beta_tracker()
|
||||
b = get_session_beta_tracker()
|
||||
assert a is b
|
||||
|
||||
|
||||
def test_singleton_reset_replaces_instance() -> None:
|
||||
a = get_session_beta_tracker()
|
||||
_reset_session_beta_tracker_for_test()
|
||||
b = get_session_beta_tracker()
|
||||
assert a is not b
|
||||
|
||||
|
||||
def test_memory_injection_appends_deterministic_order() -> None:
|
||||
"""End-to-end: client value + memory beta token → deterministic merged value.
|
||||
|
||||
Mirrors the ad-hoc concat that the handler used to do but via the
|
||||
new merge helper. Order is client first, headroom token after.
|
||||
"""
|
||||
client = "interleaved-thinking-2025-05-14"
|
||||
merged = merge_anthropic_beta(client, ["context-management-2025-06-27"])
|
||||
assert merged == "interleaved-thinking-2025-05-14,context-management-2025-06-27"
|
||||
|
||||
# When client already had the token, no double-inject.
|
||||
client2 = "context-management-2025-06-27,interleaved-thinking-2025-05-14"
|
||||
merged2 = merge_anthropic_beta(client2, ["context-management-2025-06-27"])
|
||||
assert merged2 == client2
|
||||
226
tests/test_openai_beta_session_sticky.py
Normal file
226
tests/test_openai_beta_session_sticky.py
Normal file
|
|
@ -0,0 +1,226 @@
|
|||
"""Session-sticky `OpenAI-Beta` tests for PR-A6 (P5-50, preps P0-6).
|
||||
|
||||
OpenAI's `OpenAI-Beta` follows the same comma-separated convention as
|
||||
Anthropic's `anthropic-beta`. The proxy auto-injects
|
||||
`responses_websockets=2026-02-06` on the WS path when absent
|
||||
(handlers/openai.py:~1711). PR-A6 routes that injection through
|
||||
`merge_openai_beta` so the client's tokens are preserved and the
|
||||
auto-injected token is appended deterministically.
|
||||
|
||||
The same `SessionBetaTracker` (provider-aware, keyed by
|
||||
``(provider, session_id)``) backs both providers — one tracker, two
|
||||
namespaces.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import threading
|
||||
|
||||
import pytest
|
||||
|
||||
from headroom.proxy.helpers import (
|
||||
SessionBetaTracker,
|
||||
_reset_session_beta_tracker_for_test,
|
||||
merge_openai_beta,
|
||||
)
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Pure helper: `merge_openai_beta`
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def test_merge_helper_pure_function() -> None:
|
||||
"""Same inputs → same output, no global state."""
|
||||
a = merge_openai_beta("a,b", ["c"])
|
||||
b = merge_openai_beta("a,b", ["c"])
|
||||
assert a == b == "a,b,c"
|
||||
|
||||
|
||||
def test_merge_helper_empty_inputs_returns_empty_string() -> None:
|
||||
assert merge_openai_beta(None, []) == ""
|
||||
assert merge_openai_beta("", []) == ""
|
||||
|
||||
|
||||
def test_merge_helper_only_client() -> None:
|
||||
assert merge_openai_beta("alpha=1,beta=2", []) == "alpha=1,beta=2"
|
||||
|
||||
|
||||
def test_merge_helper_only_headroom() -> None:
|
||||
assert merge_openai_beta(None, ["responses_websockets=2026-02-06"]) == (
|
||||
"responses_websockets=2026-02-06"
|
||||
)
|
||||
|
||||
|
||||
def test_merge_helper_preserves_client_order_appends_headroom() -> None:
|
||||
out = merge_openai_beta(
|
||||
"alpha=1,beta=2",
|
||||
["responses_websockets=2026-02-06"],
|
||||
)
|
||||
assert out == "alpha=1,beta=2,responses_websockets=2026-02-06"
|
||||
|
||||
|
||||
def test_merge_helper_no_double_inject_when_already_present() -> None:
|
||||
out = merge_openai_beta(
|
||||
"responses_websockets=2026-02-06,alpha=1",
|
||||
["responses_websockets=2026-02-06"],
|
||||
)
|
||||
assert out == "responses_websockets=2026-02-06,alpha=1"
|
||||
|
||||
|
||||
def test_dedupe_case_insensitive_preserves_first_casing() -> None:
|
||||
# OpenAI tokens are usually `kebab=value` so casing rarely differs in
|
||||
# practice, but the helper's contract is provider-agnostic.
|
||||
assert (
|
||||
merge_openai_beta("Responses_Websockets=2026-02-06", ["responses_websockets=2026-02-06"])
|
||||
== "Responses_Websockets=2026-02-06"
|
||||
)
|
||||
|
||||
|
||||
def test_merge_helper_skips_empty_tokens() -> None:
|
||||
assert merge_openai_beta("a, ,b", []) == "a,b"
|
||||
assert merge_openai_beta("a", ["", "b", " "]) == "a,b"
|
||||
|
||||
|
||||
def test_test_memory_injection_appends_deterministic_order() -> None:
|
||||
"""Auto-injected `responses_websockets` appends AFTER client tokens."""
|
||||
out = merge_openai_beta("client-token", ["responses_websockets=2026-02-06"])
|
||||
assert out == "client-token,responses_websockets=2026-02-06"
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# `SessionBetaTracker` — provider="openai"
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _isolate_tracker(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.delenv("HEADROOM_BETA_HEADER_STICKY", raising=False)
|
||||
monkeypatch.delenv("HEADROOM_BETA_TRACKER_MAX_SESSIONS", raising=False)
|
||||
_reset_session_beta_tracker_for_test()
|
||||
yield
|
||||
_reset_session_beta_tracker_for_test()
|
||||
|
||||
|
||||
def test_beta_seen_turn_1_present_in_turn_2_even_if_client_drops() -> None:
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
|
||||
out1 = tracker.record_and_get_sticky_betas(
|
||||
provider="openai",
|
||||
session_id="s-1",
|
||||
client_value="responses_websockets=2026-02-06,extra-beta=1",
|
||||
)
|
||||
assert out1 == "responses_websockets=2026-02-06,extra-beta=1"
|
||||
|
||||
out2 = tracker.record_and_get_sticky_betas(
|
||||
provider="openai",
|
||||
session_id="s-1",
|
||||
client_value="responses_websockets=2026-02-06",
|
||||
)
|
||||
assert out2 == "responses_websockets=2026-02-06,extra-beta=1"
|
||||
|
||||
|
||||
def test_client_value_preserved_when_no_injection() -> None:
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
out = tracker.record_and_get_sticky_betas(
|
||||
provider="openai", session_id="s-1", client_value="alpha,beta"
|
||||
)
|
||||
assert out == "alpha,beta"
|
||||
|
||||
|
||||
def test_lru_eviction_at_max_sessions() -> None:
|
||||
tracker = SessionBetaTracker(max_sessions=2)
|
||||
tracker.record_and_get_sticky_betas(provider="openai", session_id="s-1", client_value="a")
|
||||
tracker.record_and_get_sticky_betas(provider="openai", session_id="s-2", client_value="b")
|
||||
tracker.record_and_get_sticky_betas(provider="openai", session_id="s-1", client_value=None)
|
||||
tracker.record_and_get_sticky_betas(provider="openai", session_id="s-3", client_value="c")
|
||||
# s-2 evicted.
|
||||
out = tracker.record_and_get_sticky_betas(
|
||||
provider="openai", session_id="s-2", client_value=None
|
||||
)
|
||||
assert out == ""
|
||||
|
||||
|
||||
def test_disabled_mode_passes_through(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setenv("HEADROOM_BETA_HEADER_STICKY", "disabled")
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
|
||||
out1 = tracker.record_and_get_sticky_betas(
|
||||
provider="openai", session_id="s-1", client_value="alpha"
|
||||
)
|
||||
assert out1 == "alpha"
|
||||
out2 = tracker.record_and_get_sticky_betas(
|
||||
provider="openai", session_id="s-1", client_value=None
|
||||
)
|
||||
assert out2 == ""
|
||||
|
||||
|
||||
def test_thread_safe_concurrent_access() -> None:
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
n_threads = 16
|
||||
iterations = 50
|
||||
errors: list[BaseException] = []
|
||||
|
||||
def worker(idx: int) -> None:
|
||||
try:
|
||||
for i in range(iterations):
|
||||
tracker.record_and_get_sticky_betas(
|
||||
provider="openai",
|
||||
session_id="shared",
|
||||
client_value=f"o-t{idx}-i{i}",
|
||||
)
|
||||
except BaseException as e: # noqa: BLE001
|
||||
errors.append(e)
|
||||
|
||||
threads = [threading.Thread(target=worker, args=(idx,)) for idx in range(n_threads)]
|
||||
for t in threads:
|
||||
t.start()
|
||||
for t in threads:
|
||||
t.join()
|
||||
assert errors == []
|
||||
final = tracker.record_and_get_sticky_betas(
|
||||
provider="openai", session_id="shared", client_value=None
|
||||
)
|
||||
final_tokens = set(final.split(","))
|
||||
expected = {f"o-t{idx}-i{i}" for idx in range(n_threads) for i in range(iterations)}
|
||||
assert expected.issubset(final_tokens)
|
||||
|
||||
|
||||
def test_provider_namespaces_are_independent() -> None:
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
tracker.record_and_get_sticky_betas(
|
||||
provider="openai", session_id="shared", client_value="o-token"
|
||||
)
|
||||
tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic", session_id="shared", client_value="a-token"
|
||||
)
|
||||
out_openai = tracker.record_and_get_sticky_betas(
|
||||
provider="openai", session_id="shared", client_value=None
|
||||
)
|
||||
out_anth = tracker.record_and_get_sticky_betas(
|
||||
provider="anthropic", session_id="shared", client_value=None
|
||||
)
|
||||
assert out_openai == "o-token"
|
||||
assert out_anth == "a-token"
|
||||
|
||||
|
||||
def test_ws_required_token_appended_deterministically() -> None:
|
||||
"""Mirrors the WS handler logic — record client value, then merge required."""
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
sticky = tracker.record_and_get_sticky_betas(
|
||||
provider="openai",
|
||||
session_id="ws-1",
|
||||
client_value="custom-beta=1",
|
||||
)
|
||||
merged = merge_openai_beta(sticky, ["responses_websockets=2026-02-06"])
|
||||
assert merged == "custom-beta=1,responses_websockets=2026-02-06"
|
||||
|
||||
|
||||
def test_ws_required_token_no_double_when_client_already_has_it() -> None:
|
||||
tracker = SessionBetaTracker(max_sessions=10)
|
||||
sticky = tracker.record_and_get_sticky_betas(
|
||||
provider="openai",
|
||||
session_id="ws-2",
|
||||
client_value="responses_websockets=2026-02-06",
|
||||
)
|
||||
merged = merge_openai_beta(sticky, ["responses_websockets=2026-02-06"])
|
||||
assert merged == "responses_websockets=2026-02-06"
|
||||
|
|
@ -141,6 +141,13 @@ class _DummyOpenAIHandler(OpenAIHandlerMixin):
|
|||
self.anthropic_backend = None
|
||||
self.cost_tracker = None
|
||||
self.memory_handler = None
|
||||
# PR-A6 wires session-sticky `OpenAI-Beta` merging into the
|
||||
# responses HTTP handler — it reads `compute_session_id` to key
|
||||
# the SessionBetaTracker. The routing tests don't exercise the
|
||||
# tracker semantics themselves, so a fixed-id stub is enough.
|
||||
self.session_tracker_store = SimpleNamespace(
|
||||
compute_session_id=lambda *a, **k: "sess-openai-1",
|
||||
)
|
||||
self.captured_request: tuple[str, str, dict, dict] | None = None
|
||||
self.captured_stream_request: tuple[str, dict, dict] | None = None
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue