diff --git a/server-rs/crates/api-server/src/llm/mod.rs b/server-rs/crates/api-server/src/llm/mod.rs index 8112afec3..94f3be0bd 100644 --- a/server-rs/crates/api-server/src/llm/mod.rs +++ b/server-rs/crates/api-server/src/llm/mod.rs @@ -30,6 +30,14 @@ use crate::{ pub(crate) const LLM_REQUEST_MAX_BODY_BYTES: usize = 32 * 1024 * 1024; +/// 原生 Anthropic 直通的首包等待上限。 +/// +/// 有些部署(按账号分组)只有 Anthropic 原生渠道,这时直通 `/v1/messages` 是正确路径; +/// 另一些部署只有 OpenAI 形状的渠道,`/v1/messages` 会快速失败或长时间不响应。这里给直通 +/// 一个有限等待窗口:拿到 2xx 就用原生,否则回退到协议桥接。 +const LLM_ANTHROPIC_NATIVE_FIRST_BYTE_TIMEOUT: std::time::Duration = + std::time::Duration::from_secs(20); + mod anthropic_bridge; pub(crate) mod icon_specs; @@ -415,11 +423,11 @@ pub async fn proxy_llm_responses( ); } } - let upstream_headers = upstream.headers().clone(); let is_stream = payload .get("stream") .and_then(Value::as_bool) .unwrap_or(false); + let upstream_headers = upstream.headers().clone(); if !is_stream { let body = upstream.bytes().await.map_err(|error| { llm_error_response( @@ -487,7 +495,7 @@ pub async fn proxy_llm_messages( State(state): State, Extension(request_context): Extension, Extension(authenticated): Extension, - _headers: HeaderMap, + headers: HeaderMap, body: Bytes, ) -> Result { if body.len() > LLM_REQUEST_MAX_BODY_BYTES { @@ -558,8 +566,8 @@ pub async fn proxy_llm_messages( ) })? .to_string(); - // 这里不再改写 payload 的 model:出站请求由 Anthropic→Chat 转换器按解析后的模型名重建。 - object.remove("model"); + // 原生直通要用目录解析后的上游模型名;桥接路径由转换器显式接收同一个名字。 + object.insert("model".to_string(), Value::String(selected_model.clone())); let (base_url, api_key, key_id) = resolve_llm_router_credentials(&state, authenticated.claims().user_id()) @@ -587,9 +595,122 @@ pub async fn proxy_llm_messages( .with_message(format!("创建 LLM Router 请求客户端失败:{error}")), ) })?; - // 账号 Router 的 `/v1/messages` 在现网不可用(最小请求也回 `not implemented`), - // 所以这里把 Anthropic Messages 转成 Router 可用的 Chat Completions,再把回程的 - // SSE 翻译回 Anthropic 事件流。Claude Code 只认 Anthropic 形状,转换发生在网关内。 + // 先试原生 Anthropic 直通:账号所在分组有 Anthropic 渠道时这是最短、最保真的路径。 + // 直通失败(非 2xx / 首包超时)再回退到协议桥接——现网有些分组只有 OpenAI 形状的 + // 渠道,`/v1/messages` 在那里会直接报 `not implemented`。 + let is_stream = payload + .get("stream") + .and_then(Value::as_bool) + .unwrap_or(false); + let native_body = payload.to_string(); + let mut native_request = client + .post(router_protocol_url(&base_url, "messages")) + .header("x-api-key", api_key.clone()) + .bearer_auth(api_key.clone()) + .header("content-type", "application/json"); + for name in ["anthropic-version", "anthropic-beta", "accept"] { + if let Some(value) = headers.get(name) { + native_request = native_request.header(name, value); + } + } + let native = tokio::time::timeout( + LLM_ANTHROPIC_NATIVE_FIRST_BYTE_TIMEOUT, + native_request.body(native_body).send(), + ) + .await; + let native_response = match native { + Ok(Ok(response)) if response.status().is_success() => Some(response), + Ok(Ok(response)) => { + tracing::warn!( + request_id = request_context.request_id(), + user_id = %authenticated.claims().user_id(), + model = %selected_model, + status = response.status().as_u16(), + "LLM Router 原生 Anthropic 直通失败,回退协议桥接" + ); + None + } + Ok(Err(error)) => { + tracing::warn!( + request_id = request_context.request_id(), + user_id = %authenticated.claims().user_id(), + model = %selected_model, + error = %error, + "LLM Router 原生 Anthropic 请求失败,回退协议桥接" + ); + None + } + Err(_) => { + tracing::warn!( + request_id = request_context.request_id(), + user_id = %authenticated.claims().user_id(), + model = %selected_model, + "LLM Router 原生 Anthropic 首包超时,回退协议桥接" + ); + None + } + }; + if let Some(upstream) = native_response { + let status = upstream.status(); + let upstream_headers = upstream.headers().clone(); + if matches!(status, StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN) { + if let Err(error) = crate::external_api_keys::revoke_llm_router_account( + &state, + authenticated.claims().user_id(), + &key_id, + ) + .await + { + tracing::warn!( + request_id = request_context.request_id(), + user_id = %authenticated.claims().user_id(), + key_id = %key_id, + error = %error, + "LLM Router 返回确定鉴权失败,但本地账号 Key 失效标记未完成" + ); + } + } + if !is_stream { + let body = upstream.bytes().await.map_err(|error| { + llm_error_response( + &request_context, + AppError::from_status(StatusCode::BAD_GATEWAY) + .with_message(format!("读取 LLM Router 响应失败:{error}")), + ) + })?; + if let Err(error) = + settle_llm_router_usage(&state, authenticated.claims().user_id()).await + { + tracing::error!( + request_id = request_context.request_id(), + user_id = %authenticated.claims().user_id(), + error = %error, + "LLM Router Messages 成功但累计额度同步未完成" + ); + } + return build_upstream_response( + status, + &upstream_headers, + Body::from(body), + &request_context, + ); + } + let stream = stream_messages_with_billing( + upstream.bytes_stream(), + state.clone(), + authenticated.claims().user_id().to_string(), + request_context.request_id().to_string(), + ); + return build_upstream_response( + status, + &upstream_headers, + Body::from_stream(stream), + &request_context, + ); + } + + // 回退路径:把 Anthropic Messages 转成 Router 可用的 Chat Completions,再把回程的 + // SSE 翻译回 Anthropic 事件流(Claude Code 只认 Anthropic 形状)。 let (converted, tool_names) = anthropic_bridge::anthropic_messages_to_chat_completions(&payload, &selected_model) .map_err(|message| { @@ -1081,6 +1202,77 @@ fn stream_responses_with_billing( /// /// The bytes are forwarded verbatim so the Claude Agent SDK keeps parsing the /// same event shapes it received before the platform gateway was introduced. +/// 原生 Anthropic 直通的流式结算:`message_stop` 是唯一结算点。 +fn stream_messages_with_billing( + mut upstream: impl futures_util::Stream> + Unpin, + state: AppState, + owner_user_id: String, + request_id: String, +) -> impl futures_util::Stream> { + async_stream::stream! { + let mut pending = String::new(); + let mut utf8_pending = Vec::new(); + let mut billing_result: Option> = None; + while let Some(chunk) = upstream.next().await { + match chunk { + Ok(bytes) => { + append_utf8_chunk(&mut pending, &mut utf8_pending, bytes.as_ref()); + pending = pending.replace("\r\n", "\n"); + while let Some(separator) = pending.find("\n\n") { + let raw_event = pending[..separator + 2].to_string(); + let event = pending[..separator].to_string(); + pending.drain(..separator + 2); + if is_messages_terminal_sse_event(&event) && billing_result.is_none() { + billing_result = Some( + settle_llm_router_usage(&state, owner_user_id.as_str()).await, + ); + if let Some(Err(error)) = billing_result.as_ref() { + tracing::error!( + request_id = %request_id, + user_id = %owner_user_id, + error = %error, + "LLM Router 流式响应已完成但累计额度同步未完成" + ); + } + } + yield Ok(Bytes::from(raw_event)); + } + } + Err(error) => { + if billing_result.is_none() + && let Err(billing_error) = + settle_llm_router_usage(&state, owner_user_id.as_str()).await + { + tracing::error!( + request_id = %request_id, + user_id = %owner_user_id, + error = %billing_error, + "LLM Router 流式响应中断且累计额度同步未完成" + ); + } + yield Err(std::io::Error::other(format!( + "LLM Router 响应流读取失败:{error}" + ))); + return; + } + } + } + if !pending.is_empty() { + yield Ok(Bytes::from(pending)); + } + } +} + +/// Anthropic Messages 以 `message_stop` 结束,没有 `[DONE]` 哨兵。 +fn is_messages_terminal_sse_event(event: &str) -> bool { + event.lines().any(|line| { + let line = line.trim(); + line == "event: message_stop" + || line.contains("\"type\":\"message_stop\"") + || line.contains("\"type\": \"message_stop\"") + }) +} + /// Router 的 Chat Completions SSE → Anthropic Messages SSE(Claude Code 只认后者)。 /// /// 结算点放在 `message_stop` 之后或流结束:Anthropic 的终止事件由翻译器产出, @@ -2084,7 +2276,7 @@ mod tests { #[tokio::test] async fn llm_anthropic_messages_uses_account_router_credential() { - let (server_url, captured) = spawn_capturing_mock_server(MockResponse { + let (server_url, captured) = spawn_native_rejecting_capturing_mock_server(MockResponse { status_line: "200 OK", content_type: "application/json; charset=utf-8", body: r#"{"id":"chatcmpl_api_server_01","object":"chat.completion","choices":[{"index":0,"message":{"role":"assistant","content":"pong"},"finish_reason":"stop"}],"usage":{"prompt_tokens":5,"completion_tokens":2}}"#.to_string(), @@ -2170,7 +2362,7 @@ mod tests { "data: {\"id\":\"chatcmpl_stream_01\",\"choices\":[{\"index\":0,\"delta\":{},\"finish_reason\":\"stop\"}],\"usage\":{\"prompt_tokens\":5,\"completion_tokens\":2}}\n\n", "data: [DONE]\n\n", ); - let (server_url, _captured) = spawn_capturing_mock_server(MockResponse { + let (server_url, _captured) = spawn_native_rejecting_capturing_mock_server(MockResponse { status_line: "200 OK", content_type: "text/event-stream; charset=utf-8", body: streamed.to_string(), @@ -2259,6 +2451,36 @@ mod tests { (format!("http://{address}"), captured) } + /// 原生 Anthropic 直通先被拒(500),再由桥接路径命中;捕获的是第二条(桥接)请求。 + fn spawn_native_rejecting_capturing_mock_server( + response: MockResponse, + ) -> (String, Arc>>) { + let listener = TcpListener::bind("127.0.0.1:0").expect("listener should bind"); + let address = listener.local_addr().expect("listener should have addr"); + let captured = Arc::new(Mutex::new(None)); + let captured_for_thread = Arc::clone(&captured); + + thread::spawn(move || { + let (mut native_stream, _) = listener.accept().expect("native request should connect"); + let _ = read_request(&mut native_stream); + write_response( + &mut native_stream, + MockResponse { + status_line: "500 Internal Server Error", + content_type: "application/json; charset=utf-8", + body: r#"{"error":"native anthropic route unavailable"}"#.to_string(), + extra_headers: Vec::new(), + }, + ); + let (mut stream, _) = listener.accept().expect("request should connect"); + let request = read_request(&mut stream); + *captured_for_thread.lock().expect("captured request lock") = Some(request); + write_response(&mut stream, response); + }); + + (format!("http://{address}"), captured) + } + fn read_request(stream: &mut std::net::TcpStream) -> String { stream .set_read_timeout(Some(StdDuration::from_secs(1)))