From d7471b9828f5aad6ae046f97094d679514ed8340 Mon Sep 17 00:00:00 2001 From: rezaho Date: Tue, 18 Aug 2026 23:55:59 +0200 Subject: [PATCH 1/4] feat(models): a caller can ask the provider how big a prompt is MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The model port had no way to learn a token count. Every consumer that needed one estimated it from characters, which is a property of the content rather than of the model — the same text can run 1.1 to 3.2 characters per token depending on what it is, so a bound computed that way is not an estimate with a wide error bar, it is a number with no relationship to the quantity it names. Anthropic's /v1/messages/count_tokens answers with the model's own tokenizer, is free of charge, is rate-limited independently of message creation, and does not participate in prompt caching, so a caller may count the exact payload it is about to send without perturbing anything. BaseAPIModel.acount_tokens returns None by default: a provider without a counting service must say so rather than be guessed for, and the caller needs to tell an absent count from a small one. The two Anthropic legs implement it on their async adapters, each through its OWN format_request_payload — the OAuth leg prepends a required system block and renames reserved tools, so a body assembled any other way would count a request nobody makes. Only the generation controls come off. Failure is None, never an exception: the caller is sizing something, not producing the turn's answer. A 404/405 says the route is absent and is remembered for the process; a credential-shaped refusal is deliberately not, because the OAuth token file has several writers and a refresh in flight looks exactly like a rejection for one request. --- src/marsys/models/adapters/anthropic.py | 125 ++++++ src/marsys/models/adapters/anthropic_oauth.py | 57 ++- src/marsys/models/models.py | 26 ++ tests/models/test_count_tokens.py | 359 ++++++++++++++++++ 4 files changed, 566 insertions(+), 1 deletion(-) create mode 100644 tests/models/test_count_tokens.py diff --git a/src/marsys/models/adapters/anthropic.py b/src/marsys/models/adapters/anthropic.py index d480a4eb..ab258f35 100644 --- a/src/marsys/models/adapters/anthropic.py +++ b/src/marsys/models/adapters/anthropic.py @@ -179,6 +179,76 @@ def mark_conversation_tail_for_cache( return +# --- token counting (POST /v1/messages/count_tokens) ------------------------- +# +# The endpoint takes the SAME body as message creation minus the generation +# controls, and answers ``{"input_tokens": N}`` with the model's own tokenizer. +# It is free of charge, rate-limited independently of message creation, and does +# not participate in prompt caching, so a caller may count the exact payload it is +# about to send without perturbing anything. +# +# Both Anthropic legs (api-key and OAuth) build their payload through their own +# ``format_request_payload`` — that is the point: the OAuth leg injects a required +# system block and renames tools, so a count assembled any other way would not be a +# count of what the request actually sends. Only the generation-side keys come off. +_COUNT_TOKENS_REJECTED_KEYS = ( + "max_tokens", + "stream", + "temperature", + "top_p", + "top_k", + "output_config", +) +# Endpoints (by URL) that answered "there is no such route". Structural and +# permanent for the process; a credential-shaped refusal (401/403) is deliberately +# NOT recorded here, because the OAuth token file has several writers and a refresh +# in flight looks exactly like a rejection for one request. +_COUNT_TOKENS_UNSUPPORTED: set = set() + + +def count_tokens_url_for(messages_url: str) -> str: + """The count endpoint beside a Messages endpoint, query string preserved (the + OAuth leg's URL carries ``?beta=true`` and the count has to keep it).""" + base, sep, query = messages_url.partition("?") + return f"{base.rstrip('/')}/count_tokens" + (f"{sep}{query}" if sep else "") + + +def strip_for_count_tokens(payload: Dict[str, Any]) -> Dict[str, Any]: + """A Messages payload reduced to what the count endpoint accepts.""" + return {k: v for k, v in payload.items() if k not in _COUNT_TOKENS_REJECTED_KEYS} + + +def count_tokens_supported(url: str) -> bool: + return url not in _COUNT_TOKENS_UNSUPPORTED + + +def read_count_tokens_response( + url: str, status: int, body: Any, *, provider: str +) -> Optional[int]: + """One count response → a token count, or ``None``. + + ``None`` is the whole error vocabulary: a count is an optimisation on the + caller's side, never the turn's business, so nothing here raises. A 404/405 + says the route does not exist on this endpoint and is remembered for the + process; every other failure is treated as this-request-only. + """ + if status == 200 and isinstance(body, dict): + count = body.get("input_tokens") + if isinstance(count, int) and count >= 0: + return count + logger.warning("%s count_tokens returned no input_tokens: %r", provider, body) + return None + if status in (404, 405): + _COUNT_TOKENS_UNSUPPORTED.add(url) + logger.warning( + "%s does not serve %s (HTTP %s); token counting is off for this process", + provider, url, status, + ) + return None + logger.warning("%s count_tokens failed with HTTP %s", provider, status) + return None + + class AnthropicAdapter(APIProviderAdapter): """Adapter for Anthropic Claude API""" @@ -582,6 +652,14 @@ def format_request_payload(self, messages: List[Dict], **kwargs) -> Dict[str, An def get_endpoint_url(self) -> str: return f"{self.base_url.rstrip('/')}/messages" + def get_count_tokens_url(self) -> str: + return count_tokens_url_for(self.get_endpoint_url()) + + def format_count_tokens_payload( + self, messages: List[Dict], **kwargs + ) -> Dict[str, Any]: + return strip_for_count_tokens(self.format_request_payload(messages, **kwargs)) + def handle_api_error(self, error: Exception, response=None) -> ErrorResponse: """Enhanced error handling using ModelAPIError classification.""" from marsys.agents.exceptions import ModelAPIError @@ -887,3 +965,50 @@ async def arun_streaming( provider=self._provider_name() or "anthropic", response=stream_error_payload({"type": "max_retries"}, 0), ) + + async def acount_tokens( + self, + messages: List[Dict], + *, + tools: Optional[List[Dict[str, Any]]] = None, + system: Optional[str] = None, + ) -> Optional[int]: + """How many input tokens a Messages call with this body would carry, by the + model's own tokenizer, or ``None`` when the endpoint cannot say. + + ``system`` rides in as the leading ``role:"system"`` message the payload + builder already knows how to hoist, so the counted body is assembled by the + same code the real request goes through — including this leg's tool + conversion and cache markers. + + No retry loop and no raise: the caller is sizing something, not producing + the turn's answer, and a count that fails must cost the turn nothing. + """ + import asyncio + + import aiohttp + + url = self.get_count_tokens_url() + if not count_tokens_supported(url): + return None + if system is not None: + messages = [{"role": "system", "content": system}, *messages] + payload = self.format_count_tokens_payload(messages, tools=tools) + headers = {**self.get_headers(), "accept": "application/json"} + session = await self._ensure_session() + try: + async with session.post( + url, + headers=headers, + json=payload, + timeout=aiohttp.ClientTimeout(total=30), + ) as response: + status = response.status + try: + body = await response.json(content_type=None) + except ValueError: + body = None + except (aiohttp.ClientError, asyncio.TimeoutError) as exc: + logger.warning("anthropic count_tokens transport failure: %s", exc) + return None + return read_count_tokens_response(url, status, body, provider="anthropic") diff --git a/src/marsys/models/adapters/anthropic_oauth.py b/src/marsys/models/adapters/anthropic_oauth.py index cab6add6..3f844a44 100644 --- a/src/marsys/models/adapters/anthropic_oauth.py +++ b/src/marsys/models/adapters/anthropic_oauth.py @@ -10,7 +10,11 @@ CACHE_EXEMPT_KEY, _anthropic_model_rejects_temperature, _anthropic_model_requires_adaptive_thinking, + count_tokens_supported, + count_tokens_url_for, mark_conversation_tail_for_cache, + read_count_tokens_response, + strip_for_count_tokens, ) from marsys.models.adapters.base import APIProviderAdapter, AsyncBaseAPIAdapter from marsys.models.response_models import ( @@ -296,6 +300,14 @@ def get_endpoint_url(self) -> str: """Return Claude API endpoint URL.""" return self.API_URL + def get_count_tokens_url(self) -> str: + return count_tokens_url_for(self.get_endpoint_url()) + + def format_count_tokens_payload( + self, messages: List[Dict], **kwargs + ) -> Dict[str, Any]: + return strip_for_count_tokens(self.format_request_payload(messages, **kwargs)) + def _build_system_array(self, system_message: Optional[str] = None) -> List[Dict]: """ Build system prompt array with required prefix. @@ -1148,4 +1160,47 @@ async def arun_streaming( exception=e ) - + async def acount_tokens( + self, + messages: List[Dict], + *, + tools: Optional[List[Dict[str, Any]]] = None, + system: Optional[str] = None, + ) -> Optional[int]: + """How many input tokens a Messages call with this body would carry, by the + model's own tokenizer, or ``None`` when the endpoint cannot say. + + Built through this leg's own payload builder, which is why the method exists + here rather than once for both Anthropic legs: this one prepends the required + Claude-Code system block and renames tools, so a body assembled any other way + would count something the request never sends. The ``accept`` header is the + one deliberate difference from the messages call — that call is + streaming-always, and a count is a single JSON object. + + No retry loop and no raise: the caller is sizing something, not producing the + turn's answer, and a count that fails must cost the turn nothing. + """ + import httpx + + url = self.get_count_tokens_url() + if not count_tokens_supported(url): + return None + if system is not None: + messages = [{"role": "system", "content": system}, *messages] + await asyncio.to_thread(self._ensure_fresh_token) + payload = self.format_count_tokens_payload(messages, tools=tools) + headers = {**self.get_headers(), "accept": "application/json"} + try: + async with httpx.AsyncClient(timeout=30.0) as client: + response = await client.post(url, headers=headers, json=payload) + status = response.status_code + try: + body = response.json() + except ValueError: + body = None + except httpx.HTTPError as exc: + logger.warning("anthropic-oauth count_tokens transport failure: %s", exc) + return None + return read_count_tokens_response( + url, status, body, provider="anthropic-oauth" + ) diff --git a/src/marsys/models/models.py b/src/marsys/models/models.py index c8ed73d6..fcf827e2 100644 --- a/src/marsys/models/models.py +++ b/src/marsys/models/models.py @@ -927,6 +927,32 @@ def sync_run(): response = await loop.run_in_executor(None, sync_run) return response + async def acount_tokens( + self, + messages: List[Dict[str, Any]], + *, + tools: Optional[List[Dict[str, Any]]] = None, + system: Optional[str] = None, + ) -> Optional[int]: + """How many input tokens this model would charge for that body, counted by + the provider rather than estimated, or ``None`` when the provider offers no + such service. + + ``None`` is the default and the only honest answer for a provider without a + counting endpoint: a caller that needs a number can fall back to whatever + estimate it already has, but it must be able to tell an estimate from a + count. Providers that CAN answer implement ``acount_tokens`` on their async + adapter, where the payload rendering and credentials live. + + This is not the per-message ``TokenCounter`` in ``marsys.utils.tokens``: + that protocol returns a count per message from a character heuristic, which + one endpoint call cannot produce. Two different capabilities. + """ + counter = getattr(self.async_adapter, "acount_tokens", None) + if counter is None: + return None + return await counter(messages, tools=tools, system=system) + async def cleanup(self): """Clean up async resources.""" if self.async_adapter and hasattr(self.async_adapter, 'cleanup'): diff --git a/tests/models/test_count_tokens.py b/tests/models/test_count_tokens.py new file mode 100644 index 00000000..4c866752 --- /dev/null +++ b/tests/models/test_count_tokens.py @@ -0,0 +1,359 @@ +"""Provider-counted input tokens on the two Anthropic legs. + +No network. What is pinned here: + +* **The counted body is the body that would be SENT.** Each leg builds the count + payload through its own ``format_request_payload``, so the OAuth leg's required + Claude-Code system block and its tool renaming are counted too. A count assembled + any other way measures a request nobody makes. +* **Only the generation controls come off.** ``max_tokens``/``stream``/sampling are + rejected by the count endpoint; ``system``/``messages``/``tools``/``thinking`` are + exactly what decides the number. +* **A failure is ``None``, never an exception.** The caller is sizing something, not + producing an answer. A missing route (404/405) is remembered for the process; a + credential-shaped refusal is not — the OAuth token file has several writers and a + refresh in flight looks exactly like a rejection for one request. +""" + +import json + +import httpx +import pytest + +from marsys.models.adapters import anthropic as anthropic_mod +from marsys.models.adapters.anthropic import ( + AsyncAnthropicAdapter, + count_tokens_url_for, + strip_for_count_tokens, +) +from marsys.models.adapters.anthropic_oauth import AsyncAnthropicOAuthAdapter +from marsys.models.models import BaseAPIModel + +MESSAGES = [{"role": "user", "content": "hi"}] +TOOLS = [ + { + "type": "function", + "function": { + "name": "read_file", + "description": "read a file", + "parameters": {"type": "object", "properties": {}}, + }, + } +] + + +@pytest.fixture(autouse=True) +def _forget_unsupported_endpoints(): + """The unavailable-endpoint memo is process-wide by design; tests must not + inherit each other's.""" + anthropic_mod._COUNT_TOKENS_UNSUPPORTED.clear() + yield + anthropic_mod._COUNT_TOKENS_UNSUPPORTED.clear() + + +class _FakeResponse: + def __init__(self, status: int, body): + self.status = status + self._body = body + + async def json(self, content_type=None): + if self._body is None: + raise ValueError("not json") + return self._body + + async def __aenter__(self): + return self + + async def __aexit__(self, *exc): + return False + + +_DEFAULT_BODY = object() # so ``body=None`` can mean "the response is not JSON" + + +class _FakeSession: + """The aiohttp seam ``AsyncBaseAPIAdapter._ensure_session`` hands out.""" + + # ``_ensure_session`` reuses a live session and rebuilds a closed one. + closed = False + + def __init__(self, status: int = 200, body=_DEFAULT_BODY, raises=None): + self.status = status + self.body = {"input_tokens": 1234} if body is _DEFAULT_BODY else body + self.raises = raises + self.calls: list[dict] = [] + + def post(self, url, headers=None, json=None, timeout=None): + self.calls.append({"url": url, "headers": headers, "payload": json}) + if self.raises is not None: + raise self.raises + return _FakeResponse(self.status, self.body) + + +def _api(session: _FakeSession) -> AsyncAnthropicAdapter: + adapter = AsyncAnthropicAdapter( + model_name="claude-opus-5", + api_key="not-a-real-key", + base_url="https://api.anthropic.com/v1", + max_tokens=8192, + ) + adapter._session = session + return adapter + + +def _oauth() -> AsyncAnthropicOAuthAdapter: + """An OAuth adapter that never touches the credentials file on disk.""" + adapter = object.__new__(AsyncAnthropicOAuthAdapter) + adapter.model_name = "claude-sonnet-5" + adapter.max_tokens = 8192 + adapter.temperature = 0.7 + adapter.enable_thinking = False + adapter.thinking_budget = 0 + adapter.auto_refresh = False + adapter.access_token = "not-a-real-token" + adapter.credentials = {"access_token": "not-a-real-token"} + adapter._credentials_path = "unused" + adapter._session = None + return adapter + + +def _oauth_transport(monkeypatch, handler): + real_async_client = httpx.AsyncClient + transport = httpx.MockTransport(handler) + monkeypatch.setattr( + httpx, "AsyncClient", lambda **kw: real_async_client(transport=transport) + ) + + +# ── the URL and the payload ─────────────────────────────────────────────────── + + +def test_the_count_url_sits_beside_the_messages_url_and_keeps_its_query(): + assert ( + count_tokens_url_for("https://api.anthropic.com/v1/messages") + == "https://api.anthropic.com/v1/messages/count_tokens" + ) + # The OAuth leg's endpoint carries a query string; dropping it would send the + # count somewhere the messages call never goes. + assert ( + count_tokens_url_for("https://api.anthropic.com/v1/messages?beta=true") + == "https://api.anthropic.com/v1/messages/count_tokens?beta=true" + ) + + +def test_only_the_generation_controls_are_stripped(): + payload = { + "model": "claude-opus-5", + "max_tokens": 8192, + "stream": True, + "temperature": 0.7, + "top_p": 0.9, + "output_config": {"effort": "high"}, + "system": [{"type": "text", "text": "SYS"}], + "messages": [{"role": "user", "content": "hi"}], + "tools": [{"name": "t", "input_schema": {}}], + "thinking": {"type": "adaptive"}, + } + stripped = strip_for_count_tokens(payload) + assert set(stripped) == {"model", "system", "messages", "tools", "thinking"} + # the survivors are untouched, not rebuilt + assert stripped["system"] == payload["system"] + assert stripped["tools"] == payload["tools"] + + +# ── the api-key leg ─────────────────────────────────────────────────────────── + + +async def test_the_api_key_leg_counts_the_body_it_would_send(): + session = _FakeSession(body={"input_tokens": 4_242}) + adapter = _api(session) + + count = await adapter.acount_tokens(MESSAGES, tools=TOOLS, system="SYS") + + assert count == 4_242 + call = session.calls[0] + assert call["url"] == "https://api.anthropic.com/v1/messages/count_tokens" + # The system string rides as the payload builder's own rendered system field — + # the block shape the messages call sends, not the raw string. + assert call["payload"]["system"][0]["text"] == "SYS" + # …the tools are the converted Anthropic shape the request would carry… + assert call["payload"]["tools"][0]["name"] == "read_file" + assert "input_schema" in call["payload"]["tools"][0] + # …and the generation controls are gone. + assert "max_tokens" not in call["payload"] and "stream" not in call["payload"] + assert call["headers"]["accept"] == "application/json" + assert call["headers"]["x-api-key"] == "not-a-real-key" + + +async def test_a_body_with_no_system_is_counted_as_is(): + session = _FakeSession(body={"input_tokens": 11}) + adapter = _api(session) + assert await adapter.acount_tokens(MESSAGES) == 11 + assert session.calls[0]["payload"]["messages"][0]["role"] == "user" + + +async def test_a_transport_failure_answers_none_rather_than_raising(): + import aiohttp + + session = _FakeSession(raises=aiohttp.ClientError("connection reset")) + adapter = _api(session) + assert await adapter.acount_tokens(MESSAGES) is None + + +async def test_a_body_that_is_not_json_answers_none(): + session = _FakeSession(status=200, body=None) + adapter = _api(session) + assert await adapter.acount_tokens(MESSAGES) is None + + +async def test_a_200_without_input_tokens_answers_none(): + session = _FakeSession(status=200, body={"unexpected": True}) + adapter = _api(session) + assert await adapter.acount_tokens(MESSAGES) is None + + +# ── availability: structural vs credential-shaped ───────────────────────────── + + +@pytest.mark.parametrize("status", [404, 405]) +async def test_a_missing_route_is_remembered_for_the_process(status): + """The endpoint is not there; asking again every turn buys nothing.""" + session = _FakeSession(status=status, body={"error": "not found"}) + adapter = _api(session) + + assert await adapter.acount_tokens(MESSAGES) is None + assert await adapter.acount_tokens(MESSAGES) is None + assert len(session.calls) == 1 # the second call never left the process + + +@pytest.mark.parametrize("status", [401, 403, 429, 500]) +async def test_a_credential_or_transient_refusal_is_never_remembered(status): + """The OAuth token file has several writers, so a refresh in flight looks like a + rejection for one request. Memoising that would disable counting for the life of + a daemon that never restarts.""" + session = _FakeSession(status=status, body={"error": "nope"}) + adapter = _api(session) + + assert await adapter.acount_tokens(MESSAGES) is None + assert await adapter.acount_tokens(MESSAGES) is None + assert len(session.calls) == 2 # asked again, as it must be + + +async def test_the_memo_is_per_endpoint_not_global(): + unavailable = _api(_FakeSession(status=404, body={})) + assert await unavailable.acount_tokens(MESSAGES) is None + + other_session = _FakeSession(body={"input_tokens": 7}) + other = AsyncAnthropicAdapter( + model_name="claude-opus-5", + api_key="k", + base_url="https://other.example.invalid/v1", + max_tokens=1024, + ) + other._session = other_session + assert await other.acount_tokens(MESSAGES) == 7 + + +# ── the OAuth leg ───────────────────────────────────────────────────────────── + + +async def test_the_oauth_leg_counts_what_its_own_payload_builder_renders(monkeypatch): + seen: dict = {} + + def handler(request: httpx.Request) -> httpx.Response: + seen["url"] = str(request.url) + seen["headers"] = dict(request.headers) + seen["payload"] = json.loads(request.content) + return httpx.Response(200, json={"input_tokens": 9_001}) + + _oauth_transport(monkeypatch, handler) + adapter = _oauth() + monkeypatch.setattr(adapter, "_ensure_fresh_token", lambda: None, raising=False) + + count = await adapter.acount_tokens(MESSAGES, tools=TOOLS, system="SYS") + + assert count == 9_001 + assert seen["url"] == "https://api.anthropic.com/v1/messages/count_tokens?beta=true" + # The required Claude-Code prefix block is COUNTED, because the messages call + # always sends it — a count without it undercounts every request. + assert seen["payload"]["system"][0]["text"] == adapter.CLAUDE_CODE_PREFIX + assert seen["payload"]["system"][1]["text"] == "SYS" + # …as is the reserved-name transform this leg applies to tools. + assert seen["payload"]["tools"][0]["name"] == "Read" + assert "max_tokens" not in seen["payload"] and "stream" not in seen["payload"] + # JSON, not the messages call's SSE. + assert seen["headers"]["accept"] == "application/json" + assert seen["headers"]["authorization"].startswith("Bearer ") + assert "oauth-2025-04-20" in seen["headers"]["anthropic-beta"] + + +async def test_the_oauth_leg_refreshes_its_token_before_counting(monkeypatch): + """Same discipline as the messages call: the cached token is a per-request + snapshot of a file other processes rewrite.""" + refreshed: list[bool] = [] + _oauth_transport( + monkeypatch, lambda request: httpx.Response(200, json={"input_tokens": 5}) + ) + adapter = _oauth() + monkeypatch.setattr( + adapter, "_ensure_fresh_token", lambda: refreshed.append(True), raising=False + ) + + assert await adapter.acount_tokens(MESSAGES) == 5 + assert refreshed == [True] + + +async def test_an_unauthorised_oauth_leg_answers_none_without_raising(monkeypatch): + """The state AC-10's probe exists to settle: if the subscription credential is + not accepted here, compaction reports "not measured" and the turn is untouched.""" + _oauth_transport( + monkeypatch, + lambda request: httpx.Response(403, json={"error": {"type": "forbidden"}}), + ) + adapter = _oauth() + monkeypatch.setattr(adapter, "_ensure_fresh_token", lambda: None, raising=False) + + assert await adapter.acount_tokens(MESSAGES) is None + + +async def test_an_oauth_transport_failure_answers_none(monkeypatch): + def handler(request: httpx.Request) -> httpx.Response: + raise httpx.ConnectError("no route to host") + + _oauth_transport(monkeypatch, handler) + adapter = _oauth() + monkeypatch.setattr(adapter, "_ensure_fresh_token", lambda: None, raising=False) + + assert await adapter.acount_tokens(MESSAGES) is None + + +# ── the model port ──────────────────────────────────────────────────────────── + + +class _AdapterWithCount: + def __init__(self): + self.seen: dict = {} + + async def acount_tokens(self, messages, *, tools=None, system=None): + self.seen = {"messages": messages, "tools": tools, "system": system} + return 321 + + +async def test_the_model_delegates_the_count_to_its_async_adapter(): + model = object.__new__(BaseAPIModel) + adapter = _AdapterWithCount() + model.async_adapter = adapter + + assert await model.acount_tokens(MESSAGES, tools=TOOLS, system="SYS") == 321 + assert adapter.seen == {"messages": MESSAGES, "tools": TOOLS, "system": "SYS"} + + +async def test_a_provider_that_cannot_count_answers_none_rather_than_guessing(): + """Every non-Anthropic leg takes this path. The caller must be able to tell an + absent count from a small one.""" + model = object.__new__(BaseAPIModel) + model.async_adapter = object() + assert await model.acount_tokens(MESSAGES) is None + + model.async_adapter = None + assert await model.acount_tokens(MESSAGES) is None From 58cf0442be9ceaac58cc57e1e98ca851eda077f7 Mon Sep 17 00:00:00 2001 From: rezaho Date: Wed, 19 Aug 2026 00:47:59 +0200 Subject: [PATCH 2/4] fix(models): a count of a finished conversation is a legal request MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Measured against the live API: counting a message list whose last row is an assistant reply ending in whitespace is refused — invalid_request_error: messages: final assistant content cannot end with trailing whitespace The messages call almost never meets this rule, because something is appended after the last assistant turn on every real request. A count meets it constantly: counting a conversation means presenting its last row as final, a settled conversation ends with an assistant reply by definition, and models end replies with a newline. So the first caller of this endpoint would have found its post-fold counts failing on exactly the conversations most likely to need them, with nothing on the wire to say why. The count payload now right-trims the final assistant message's last text block, copying rather than mutating (those blocks can be the caller's own durable rows), and drops the message when nothing is left of it. And a non-200 logs the provider's own sentence rather than a bare status — this one answers in a line, and a status alone cost a round trip to diagnose. --- src/marsys/models/adapters/anthropic.py | 55 +++++++++++++++++++++++-- tests/models/test_count_tokens.py | 49 ++++++++++++++++++++++ 2 files changed, 101 insertions(+), 3 deletions(-) diff --git a/src/marsys/models/adapters/anthropic.py b/src/marsys/models/adapters/anthropic.py index ab258f35..b5bf9fc0 100644 --- a/src/marsys/models/adapters/anthropic.py +++ b/src/marsys/models/adapters/anthropic.py @@ -214,8 +214,45 @@ def count_tokens_url_for(messages_url: str) -> str: def strip_for_count_tokens(payload: Dict[str, Any]) -> Dict[str, Any]: - """A Messages payload reduced to what the count endpoint accepts.""" - return {k: v for k, v in payload.items() if k not in _COUNT_TOKENS_REJECTED_KEYS} + """A Messages payload reduced to a LEGAL count request. + + Two things happen here, and the second is not cosmetic. The generation controls + come off because the endpoint rejects them. And the final assistant message is + right-trimmed because the API rejects one that ends in whitespace — a rule the + messages call rarely meets (something is almost always appended after the last + assistant turn) and a count meets constantly, since counting a conversation means + presenting its last row as final. A settled conversation ends with an assistant + reply by definition, and models end replies with a newline, so without this a + fold's post-fold count fails on exactly the conversations most likely to fold. + """ + reduced = {k: v for k, v in payload.items() if k not in _COUNT_TOKENS_REJECTED_KEYS} + messages = reduced.get("messages") + if isinstance(messages, list) and messages: + trimmed = _trim_final_assistant(messages[-1]) + reduced["messages"] = ( + messages[:-1] if trimmed is None else [*messages[:-1], trimmed] + ) + return reduced + + +def _trim_final_assistant(message: Any) -> Optional[Dict[str, Any]]: + """The message with trailing whitespace off its last text block, or ``None`` when + nothing is left of it. Copies rather than mutating: the payload's blocks may be + the caller's own rows.""" + if not isinstance(message, dict) or message.get("role") != "assistant": + return message + content = message.get("content") + if isinstance(content, str): + trimmed = content.rstrip() + return {**message, "content": trimmed} if trimmed else None + if not isinstance(content, list) or not content: + return message + last = content[-1] + if not isinstance(last, dict) or last.get("type") != "text": + return message + text = str(last.get("text", "")).rstrip() + blocks = list(content[:-1]) if not text else [*content[:-1], {**last, "text": text}] + return {**message, "content": blocks} if blocks else None def count_tokens_supported(url: str) -> bool: @@ -245,10 +282,22 @@ def read_count_tokens_response( provider, url, status, ) return None - logger.warning("%s count_tokens failed with HTTP %s", provider, status) + # The provider's own words, not just the status: a bare "HTTP 400" on a request + # nobody sees is undiagnosable, and this one answers in a sentence. + logger.warning( + "%s count_tokens failed with HTTP %s: %s", provider, status, _error_text(body) + ) return None +def _error_text(body: Any) -> str: + if isinstance(body, dict): + error = body.get("error") + if isinstance(error, dict) and error.get("message"): + return str(error["message"]) + return "no error body" + + class AnthropicAdapter(APIProviderAdapter): """Adapter for Anthropic Claude API""" diff --git a/tests/models/test_count_tokens.py b/tests/models/test_count_tokens.py index 4c866752..8e649896 100644 --- a/tests/models/test_count_tokens.py +++ b/tests/models/test_count_tokens.py @@ -141,6 +141,55 @@ def test_the_count_url_sits_beside_the_messages_url_and_keeps_its_query(): ) +def test_the_final_assistant_message_is_trimmed_of_trailing_whitespace(): + """The API rejects a final assistant message that ends in whitespace, and a count + is the one request that routinely presents one: a settled conversation ends with + an assistant reply, and models end replies with a newline. Measured live — the + provider answers "final assistant content cannot end with trailing whitespace".""" + payload = { + "model": "m", + "messages": [ + {"role": "user", "content": "ask"}, + {"role": "assistant", "content": [{"type": "text", "text": "answered.\n\n"}]}, + ], + } + stripped = strip_for_count_tokens(payload) + assert stripped["messages"][-1]["content"][-1]["text"] == "answered." + # …and the caller's own rows are untouched: the payload's blocks may be the + # durable conversation's dicts. + assert payload["messages"][-1]["content"][-1]["text"] == "answered.\n\n" + + +def test_a_plain_string_final_assistant_message_is_trimmed_too(): + stripped = strip_for_count_tokens( + {"model": "m", "messages": [{"role": "assistant", "content": "done "}]} + ) + assert stripped["messages"][-1]["content"] == "done" + + +def test_a_whitespace_only_final_assistant_message_is_dropped(): + """Nothing is left of it to count, and an empty text block is itself rejected.""" + stripped = strip_for_count_tokens( + { + "model": "m", + "messages": [ + {"role": "user", "content": "ask"}, + {"role": "assistant", "content": [{"type": "text", "text": " "}]}, + ], + } + ) + assert stripped["messages"] == [{"role": "user", "content": "ask"}] + + +def test_a_final_user_message_is_left_exactly_as_it_is(): + """The rule is about assistant content. Trimming a user row would change what is + being counted for no reason.""" + stripped = strip_for_count_tokens( + {"model": "m", "messages": [{"role": "user", "content": "ask \n"}]} + ) + assert stripped["messages"] == [{"role": "user", "content": "ask \n"}] + + def test_only_the_generation_controls_are_stripped(): payload = { "model": "claude-opus-5", From b4d60bb6013732d0e56d6b9648a7f1a36f8b0c55 Mon Sep 17 00:00:00 2001 From: rezaho Date: Wed, 19 Aug 2026 00:55:12 +0200 Subject: [PATCH 3/4] test(models): a credential-shaped refusal is asked again, not remembered --- tests/models/test_count_tokens.py | 22 +++++++++++++++------- 1 file changed, 15 insertions(+), 7 deletions(-) diff --git a/tests/models/test_count_tokens.py b/tests/models/test_count_tokens.py index 8e649896..24588403 100644 --- a/tests/models/test_count_tokens.py +++ b/tests/models/test_count_tokens.py @@ -352,17 +352,25 @@ async def test_the_oauth_leg_refreshes_its_token_before_counting(monkeypatch): assert refreshed == [True] -async def test_an_unauthorised_oauth_leg_answers_none_without_raising(monkeypatch): - """The state AC-10's probe exists to settle: if the subscription credential is - not accepted here, compaction reports "not measured" and the turn is untouched.""" - _oauth_transport( - monkeypatch, - lambda request: httpx.Response(403, json={"error": {"type": "forbidden"}}), - ) +async def test_an_unauthorised_oauth_leg_answers_none_and_is_asked_again(monkeypatch): + """If the subscription credential is not accepted here, compaction reports "not + measured" and the turn is untouched — and the next turn asks again. This leg's + token file has several writers, so a refusal is as likely to be a refresh in + flight as a verdict, and a daemon that never restarts would otherwise stop + measuring for days on one blip.""" + calls: list[int] = [] + + def handler(request: httpx.Request) -> httpx.Response: + calls.append(1) + return httpx.Response(403, json={"error": {"type": "forbidden"}}) + + _oauth_transport(monkeypatch, handler) adapter = _oauth() monkeypatch.setattr(adapter, "_ensure_fresh_token", lambda: None, raising=False) assert await adapter.acount_tokens(MESSAGES) is None + assert await adapter.acount_tokens(MESSAGES) is None + assert len(calls) == 2 async def test_an_oauth_transport_failure_answers_none(monkeypatch): From ebdda1697541b965427c68e406d7b714d52b795b Mon Sep 17 00:00:00 2001 From: rezaho Date: Wed, 19 Aug 2026 01:18:01 +0200 Subject: [PATCH 4/4] style(models): type the unsupported-endpoint memo --- src/marsys/models/adapters/anthropic.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/marsys/models/adapters/anthropic.py b/src/marsys/models/adapters/anthropic.py index b5bf9fc0..eed52c06 100644 --- a/src/marsys/models/adapters/anthropic.py +++ b/src/marsys/models/adapters/anthropic.py @@ -203,7 +203,7 @@ def mark_conversation_tail_for_cache( # permanent for the process; a credential-shaped refusal (401/403) is deliberately # NOT recorded here, because the OAuth token file has several writers and a refresh # in flight looks exactly like a rejection for one request. -_COUNT_TOKENS_UNSUPPORTED: set = set() +_COUNT_TOKENS_UNSUPPORTED: "set[str]" = set() def count_tokens_url_for(messages_url: str) -> str: