Coverage for src/lilbee/server/anthropic_api/routes.py: 100%
123 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
1"""HTTP routes for the Anthropic-compatible ``/v1/messages`` surface."""
3from __future__ import annotations
5import asyncio
6import logging
7import uuid
8from collections.abc import AsyncGenerator
10from litestar import Request, Router, post
11from litestar.background_tasks import BackgroundTask
12from litestar.exceptions import NotAuthorizedException, ValidationException
13from litestar.response import Response, Stream
15from lilbee.app.services import get_services
16from lilbee.core.config import cfg
17from lilbee.core.config.enums import ReasoningMode
18from lilbee.retrieval.reasoning import effective_reasoning_cap
19from lilbee.server.anthropic_api.errors import anthropic_error_body, anthropic_error_type
20from lilbee.server.anthropic_api.models import (
21 _THINKING_DISABLED,
22 AnthropicEventType,
23 CountTokensRequest,
24 CountTokensResponse,
25 MessagesRequest,
26 MessagesResponse,
27)
28from lilbee.server.anthropic_api.streaming import encode_anthropic_event, encode_anthropic_sse
29from lilbee.server.anthropic_api.translate import (
30 canonical_stream_to_anthropic_events,
31 canonical_to_messages_response,
32 count_tokens_to_canonical_request,
33 messages_to_canonical_request,
34 resolve_reasoning_mode,
35)
36from lilbee.server.auth import auth_checked_in_handler, session_manager
37from lilbee.server.chat_completions_api.errors import (
38 CompletionsErrorCode,
39 classify_provider_error,
40)
41from lilbee.server.chat_dispatch.canonical import CanonicalChatRequest
42from lilbee.server.chat_dispatch.concurrency import (
43 ChatBusyError,
44 ChatSlotGuard,
45 acquire_chat_slot_or_busy,
46)
47from lilbee.server.chat_dispatch.dispatch import (
48 count_request_tokens,
49 dispatch_chat_stream,
50 preflight_chat_request,
51 resolve_served_model,
52)
53from lilbee.server.chat_dispatch.reasoning_cap import (
54 budget_capped_chars,
55 cap_aware_chat,
56 cap_aware_chat_stream,
57)
58from lilbee.server.handlers.sse import SSE_MEDIA_TYPE
59from lilbee.server.validation_format import format_validation
61log = logging.getLogger(__name__)
63_INTERNAL_ERROR_MESSAGE = "Internal server error. Check the server logs for details."
64# Says whether the count_tokens answer was measured on the backend or estimated.
65COUNT_ACCURACY_HEADER = "X-Lilbee-Token-Count-Accuracy"
68async def _auth_before_request(request: Request) -> Response | None:
69 """Reject an unauthenticated caller before Litestar parses the body.
71 Same rationale as the completions surface: the auth answer must ride this
72 surface's own error envelope, and it must win over body validation.
73 """
74 return _auth_failure(request)
77@post("/v1/messages", status_code=200, before_request=_auth_before_request)
78@auth_checked_in_handler
79async def messages_endpoint(request: Request, data: MessagesRequest) -> Response | Stream:
80 """``/v1/messages`` (stream + non-stream + tools), Anthropic wire format.
82 ``stream: true`` switches the 200 response from JSON to the Anthropic SSE
83 event stream; the request body picks the arm, matching Anthropic's own
84 contract.
85 """
86 # Thinking is opt-in per request; the setting only presents it.
87 mode = resolve_reasoning_mode(data.thinking, default=cfg.messages_reasoning)
88 cap_chars = budget_capped_chars(effective_reasoning_cap(), _budget_tokens(data))
89 try:
90 req = messages_to_canonical_request(data, mode=mode)
91 except ValueError as exc:
92 # Wire-valid but untranslatable (image content, bare tool choice).
93 return _error_response(400, CompletionsErrorCode.INVALID_REQUEST, str(exc))
95 preflight = await _preflight_resolved_model(req)
96 if isinstance(preflight, Response):
97 return preflight
98 resolved_model = preflight
100 try:
101 await acquire_chat_slot_or_busy(get_services().provider.max_concurrent_chats())
102 except ChatBusyError:
103 return _error_response(
104 429,
105 CompletionsErrorCode.RATE_LIMIT_EXCEEDED,
106 "Backend is busy. Retry in a moment.",
107 headers={"Retry-After": "1"},
108 )
110 guard = ChatSlotGuard()
111 if req.stream:
112 # The after-send hook frees the slot when a disconnect lands before the
113 # generator's first iteration (its finally never runs in that case).
114 return Stream(
115 _gated_messages_stream(
116 req, guard, model=resolved_model, mode=mode, cap_chars=cap_chars
117 ),
118 media_type=SSE_MEDIA_TYPE,
119 background=BackgroundTask(guard.release),
120 )
121 return await _run_non_stream(
122 req, guard, canonical_model=resolved_model, mode=mode, cap_chars=cap_chars
123 )
126@post("/v1/messages/count_tokens", status_code=200, before_request=_auth_before_request)
127@auth_checked_in_handler
128async def count_tokens_endpoint(request: Request, data: CountTokensRequest) -> Response:
129 """``/v1/messages/count_tokens``: the served model's own count for a request.
131 Clients read ``input_tokens`` for context accounting, and
132 ``X-Lilbee-Token-Count-Accuracy`` for whether that number was measured on the
133 backend (``exact``) or estimated (``estimated``). Counting never runs a tool,
134 so the only thing a request has to be is translatable: the tool-capability
135 check ``/v1/messages`` runs is not applied. The count does not take a chat
136 slot, so it never answers 429; a count that finds the model cold loads it,
137 as a chat call would.
138 """
139 mode = resolve_reasoning_mode(data.thinking, default=cfg.messages_reasoning)
140 try:
141 req = count_tokens_to_canonical_request(data, mode=mode)
142 except ValueError as exc:
143 # Wire-valid but untranslatable (image content, bare tool choice).
144 return _error_response(400, CompletionsErrorCode.INVALID_REQUEST, str(exc))
146 try:
147 canonical_model = await asyncio.to_thread(resolve_served_model, req)
148 count = await asyncio.to_thread(count_request_tokens, req, canonical_model=canonical_model)
149 except Exception as exc:
150 return _classified_error_response(exc)
151 body = CountTokensResponse(input_tokens=count.tokens)
152 return Response(
153 body.model_dump(),
154 media_type="application/json",
155 headers={COUNT_ACCURACY_HEADER: count.accuracy},
156 )
159def _budget_tokens(data: MessagesRequest) -> int | None:
160 """The thinking budget this request asks for; ``disabled`` carries none."""
161 if data.thinking is None or data.thinking.type == _THINKING_DISABLED:
162 return None
163 return data.thinking.budget_tokens
166async def _preflight_resolved_model(req: CanonicalChatRequest) -> str | Response:
167 """Validate *req* before any streaming response starts (see completions)."""
168 try:
169 return await asyncio.to_thread(preflight_chat_request, req)
170 except Exception as exc:
171 return _classified_error_response(exc)
174def _classified_error_response(exc: Exception) -> Response:
175 """The envelope for *exc*, or the generic 500 when nothing classifies it."""
176 classified = classify_provider_error(exc)
177 if classified is None:
178 return _internal_error_response()
179 return _error_response(classified.http_status, classified.code, classified.message)
182async def _run_non_stream(
183 req: CanonicalChatRequest,
184 guard: ChatSlotGuard,
185 *,
186 canonical_model: str,
187 mode: ReasoningMode = ReasoningMode.SEPARATE,
188 cap_chars: int = 0,
189) -> Response:
190 """Dispatch a non-streaming chat call, translating errors to the envelope."""
191 try:
192 resp = await asyncio.to_thread(
193 cap_aware_chat, req, canonical_model=canonical_model, cap_chars=cap_chars
194 )
195 except Exception as exc:
196 return _classified_error_response(exc)
197 finally:
198 await guard.release()
199 body: MessagesResponse = canonical_to_messages_response(
200 resp, response_id=_response_id(), mode=mode
201 )
202 return Response(body.model_dump(), media_type="application/json")
205async def _gated_messages_stream(
206 req: CanonicalChatRequest,
207 guard: ChatSlotGuard,
208 *,
209 model: str,
210 mode: ReasoningMode = ReasoningMode.SEPARATE,
211 cap_chars: int = 0,
212) -> AsyncGenerator[bytes, None]:
213 """Drive dispatch -> translate -> SSE-encode, freeing the slot on exit.
215 A mid-stream failure surfaces as Anthropic's ``event: error`` frame; the
216 headers are already flushed at 200 by then, so the frame is the only
217 channel left.
218 """
219 response_id = _response_id()
220 try:
221 try:
222 events = cap_aware_chat_stream(
223 dispatch_chat_stream(req, canonical_model=model),
224 req,
225 canonical_model=model,
226 cap_chars=cap_chars,
227 )
228 pairs = canonical_stream_to_anthropic_events(
229 events, model=model, response_id=response_id, mode=mode
230 )
231 async for frame in encode_anthropic_sse(pairs):
232 yield frame
233 except Exception as exc:
234 classified = classify_provider_error(exc)
235 if classified is None:
236 log.exception("anthropic messages stream failed")
237 body = anthropic_error_body("api_error", _INTERNAL_ERROR_MESSAGE)
238 else:
239 body = anthropic_error_body(
240 anthropic_error_type(classified.code), classified.message
241 )
242 yield encode_anthropic_event(AnthropicEventType.ERROR, body)
243 finally:
244 await guard.release()
247def _internal_error_response() -> Response:
248 """Log and return the generic api_error 500 envelope."""
249 log.exception("anthropic messages surface failed")
250 return _error_response(500, CompletionsErrorCode.INTERNAL_ERROR, _INTERNAL_ERROR_MESSAGE)
253def _error_response(
254 status: int,
255 code: CompletionsErrorCode,
256 message: str,
257 *,
258 headers: dict[str, str] | None = None,
259) -> Response:
260 return Response(
261 anthropic_error_body(anthropic_error_type(code), message),
262 status_code=status,
263 headers=headers or {},
264 media_type="application/json",
265 )
268def _auth_failure(request: Request) -> Response | None:
269 """Return a 401 envelope if the bearer token is missing/wrong, else None.
271 Claude Code sends ``ANTHROPIC_AUTH_TOKEN`` as a bearer Authorization
272 header, which is exactly the session token check every /v1 route runs.
273 """
274 auth_header = request.headers.get("authorization", "")
275 try:
276 authorized = session_manager.validate(auth_header)
277 except NotAuthorizedException:
278 authorized = False
279 if authorized:
280 return None
281 return _error_response(401, CompletionsErrorCode.INVALID_API_KEY, "Missing or invalid API key.")
284def _validation_exception_handler(_: Request, exc: ValidationException) -> Response:
285 """Wrap Litestar's body-parse failures in the Anthropic error envelope."""
286 return _error_response(400, CompletionsErrorCode.INVALID_REQUEST, format_validation(exc))
289def _response_id() -> str:
290 """Anthropic-style ``msg_*`` id."""
291 return f"msg_{uuid.uuid4().hex[:24]}"
294anthropic_router = Router(
295 path="/",
296 route_handlers=[messages_endpoint, count_tokens_endpoint],
297 exception_handlers={ValidationException: _validation_exception_handler},
298)