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

1"""HTTP routes for the Anthropic-compatible ``/v1/messages`` surface.""" 

2 

3from __future__ import annotations 

4 

5import asyncio 

6import logging 

7import uuid 

8from collections.abc import AsyncGenerator 

9 

10from litestar import Request, Router, post 

11from litestar.background_tasks import BackgroundTask 

12from litestar.exceptions import NotAuthorizedException, ValidationException 

13from litestar.response import Response, Stream 

14 

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 

60 

61log = logging.getLogger(__name__) 

62 

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" 

66 

67 

68async def _auth_before_request(request: Request) -> Response | None: 

69 """Reject an unauthenticated caller before Litestar parses the body. 

70 

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) 

75 

76 

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. 

81 

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)) 

94 

95 preflight = await _preflight_resolved_model(req) 

96 if isinstance(preflight, Response): 

97 return preflight 

98 resolved_model = preflight 

99 

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 ) 

109 

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 ) 

124 

125 

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. 

130 

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)) 

145 

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 ) 

157 

158 

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 

164 

165 

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) 

172 

173 

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) 

180 

181 

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") 

203 

204 

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. 

214 

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() 

245 

246 

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) 

251 

252 

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 ) 

266 

267 

268def _auth_failure(request: Request) -> Response | None: 

269 """Return a 401 envelope if the bearer token is missing/wrong, else None. 

270 

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.") 

282 

283 

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)) 

287 

288 

289def _response_id() -> str: 

290 """Anthropic-style ``msg_*`` id.""" 

291 return f"msg_{uuid.uuid4().hex[:24]}" 

292 

293 

294anthropic_router = Router( 

295 path="/", 

296 route_handlers=[messages_endpoint, count_tokens_endpoint], 

297 exception_handlers={ValidationException: _validation_exception_handler}, 

298)