AGC Anthropic 网关改为桥接 Router 的 Chat Completions
Project CI / AI game creator shell Rust lane 1/2 (push) Failing after 3m28s
Project CI / AI game creator shell Rust lane 2/2 (push) Failing after 1m56s
Project CI / AI game creator shell Rust smoke (push) Successful in 2m23s
Project CI / AI game creator shell Rust crates (push) Successful in 1m34s
Project CI / Frontend tests (push) Successful in 3m11s
Project CI / Repository checks (push) Successful in 3m43s
Project CI / Backend tests (push) Has been cancelled
Project CI / AI game creator shell web tests (push) Has been cancelled
Project CI / Native shell tests (push) Has been cancelled

- llm/anthropic_bridge:新增 Anthropic Messages → OpenAI Chat Completions 请求转换(system、messages、tool_use、tool_result、图片、工具名安全化与还原)
- llm/anthropic_bridge:新增 OpenAI SSE → Anthropic SSE 流式翻译(message_start、content_block_start/delta/stop、message_delta、message_stop,tool_calls 还原成 tool_use)
- llm:/api/llm/anthropic/v1/messages 不再打现网不可用的 Router /v1/messages,改走转换后的 /v1/chat/completions;非流式响应同样翻译回 Anthropic message
- llm:删除只服务旧直通的 Anthropic 流式结算辅助,结算点改为翻译器产出 message_stop 或流结束
- 测试:新增请求转换、工具名还原、流式翻译单测;Anthropic 路由测试改为断言 Chat Completions 上游与 Anthropic 形状回包
This commit is contained in:
kdletters
2026-10-02 15:09:31 +08:00
parent 70aeb79be6
commit 21a8886c4c
2 changed files with 1063 additions and 86 deletions
File diff suppressed because it is too large Load Diff
+190 -86
View File
@@ -30,6 +30,8 @@ use crate::{
pub(crate) const LLM_REQUEST_MAX_BODY_BYTES: usize = 32 * 1024 * 1024;
mod anthropic_bridge;
pub(crate) mod icon_specs;
#[cfg(test)]
@@ -341,7 +343,7 @@ pub async fn proxy_llm_responses(
)
})?
.to_string();
object.insert("model".to_string(), Value::String(selected_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())
@@ -485,7 +487,7 @@ pub async fn proxy_llm_messages(
State(state): State<AppState>,
Extension(request_context): Extension<RequestContext>,
Extension(authenticated): Extension<AuthenticatedAccessToken>,
headers: HeaderMap,
_headers: HeaderMap,
body: Bytes,
) -> Result<Response, Response> {
if body.len() > LLM_REQUEST_MAX_BODY_BYTES {
@@ -556,7 +558,8 @@ pub async fn proxy_llm_messages(
)
})?
.to_string();
object.insert("model".to_string(), Value::String(selected_model));
// 这里不再改写 payload 的 model:出站请求由 Anthropic→Chat 转换器按解析后的模型名重建。
object.remove("model");
let (base_url, api_key, key_id) =
resolve_llm_router_credentials(&state, authenticated.claims().user_id())
@@ -584,26 +587,27 @@ pub async fn proxy_llm_messages(
.with_message(format!("创建 LLM Router 请求客户端失败:{error}")),
)
})?;
let upstream_url = router_protocol_url(&base_url, "messages");
// The Router speaks Anthropic Messages on this path; accept both credential
// headers because the SDK uses `x-api-key` while the gateway also allows
// `Authorization: Bearer`.
let mut request = client
// 账号 Router 的 `/v1/messages` 在现网不可用(最小请求也回 `not implemented`),
// 所以这里把 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| {
llm_error_response(
&request_context,
AppError::from_status(StatusCode::BAD_REQUEST).with_message(message),
)
})?;
let upstream_url = router_protocol_url(&base_url, "chat/completions");
let upstream = client
.post(upstream_url)
.header("x-api-key", api_key.clone())
.bearer_auth(api_key)
.header("content-type", "application/json");
for name in ["anthropic-version", "anthropic-beta", "accept"] {
if let Some(value) = headers.get(name) {
request = request.header(name, value);
}
}
let upstream = request
.body(serde_json::to_vec(&payload).map_err(|error| {
.header("content-type", "application/json")
.body(serde_json::to_vec(&converted).map_err(|error| {
llm_error_response(
&request_context,
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message(format!("序列化 LLM Messages 请求失败:{error}")),
.with_message(format!("序列化 LLM Chat 请求失败:{error}")),
)
})?)
.send()
@@ -647,24 +651,48 @@ pub async fn proxy_llm_messages(
.with_message(format!("读取 LLM Router 响应失败:{error}")),
)
})?;
if status.is_success() {
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 成功但累计额度同步未完成"
);
}
if !status.is_success() {
return build_upstream_response(
status,
&upstream_headers,
Body::from(body),
&request_context,
);
}
return build_upstream_response(
status,
&upstream_headers,
Body::from(body),
&request_context,
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 成功但累计额度同步未完成"
);
}
// 非流式也要交回 Anthropic 形状:Claude 客户端(含 SDK 探针)读的是 message/content。
let upstream_value: Value = serde_json::from_slice(&body).unwrap_or(Value::Null);
let message = anthropic_bridge::chat_completions_body_to_anthropic_message(
&upstream_value,
&selected_model,
&tool_names,
);
let encoded = serde_json::to_vec(&message).map_err(|error| {
llm_error_response(
&request_context,
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message(format!("序列化 Anthropic 响应失败:{error}")),
)
})?;
let mut response = build_upstream_response(
StatusCode::OK,
&upstream_headers,
Body::from(encoded),
&request_context,
)?;
response.headers_mut().insert(
axum::http::header::CONTENT_TYPE,
HeaderValue::from_static("application/json; charset=utf-8"),
);
return Ok(response);
}
if !status.is_success() {
@@ -681,14 +709,16 @@ pub async fn proxy_llm_messages(
);
}
let stream = stream_messages_with_billing(
let stream = stream_chat_completions_as_anthropic(
upstream.bytes_stream(),
selected_model.clone(),
tool_names,
state.clone(),
authenticated.claims().user_id().to_string(),
request_context.request_id().to_string(),
);
build_upstream_response(
status,
StatusCode::OK,
&upstream_headers,
Body::from_stream(stream),
&request_context,
@@ -891,17 +921,6 @@ fn is_responses_done_sse_event(event: &str) -> bool {
.any(|line| matches!(line.trim(), "data: [DONE]" | "[DONE]"))
}
/// Anthropic Messages ends a stream with `message_stop` and has no `[DONE]`
/// sentinel, so the terminal event is the only settle point.
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\"")
})
}
fn finalize_llm_router_terminal_event(
terminal_event: &mut Option<String>,
billing_result: Option<&Result<(), AppError>>,
@@ -1057,50 +1076,46 @@ 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.
fn stream_messages_with_billing(
/// Router 的 Chat Completions SSE → Anthropic Messages SSE(Claude Code 只认后者)。
///
/// 结算点放在 `message_stop` 之后或流结束:Anthropic 的终止事件由翻译器产出,
/// 累计额度同步与 Responses 路径保持同一口径。
fn stream_chat_completions_as_anthropic(
mut upstream: impl futures_util::Stream<Item = Result<Bytes, reqwest::Error>> + Unpin,
model: String,
tool_names: anthropic_bridge::AnthropicToolNames,
state: AppState,
owner_user_id: String,
request_id: String,
) -> impl futures_util::Stream<Item = Result<Bytes, std::io::Error>> {
async_stream::stream! {
let mut pending = String::new();
let mut utf8_pending = Vec::new();
let mut billing_result = None;
let mut translator = anthropic_bridge::ChatToAnthropicStream::new(&model, tool_names);
let mut settled = false;
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 流式响应已完成但累计额度同步未完成"
);
}
let out = translator.push(&bytes);
if !out.is_empty() {
yield Ok(Bytes::from(out));
}
if translator.is_finished() && !settled {
settled = true;
if let Err(error) =
settle_llm_router_usage(&state, owner_user_id.as_str()).await
{
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() {
if let Err(billing_error) = settle_llm_router_usage(
&state,
owner_user_id.as_str(),
)
.await
if !settled {
if let Err(billing_error) =
settle_llm_router_usage(&state, owner_user_id.as_str()).await
{
tracing::error!(
request_id = %request_id,
@@ -1117,8 +1132,21 @@ fn stream_messages_with_billing(
}
}
}
if !pending.is_empty() {
yield Ok(Bytes::from(pending));
if !translator.is_finished() {
let tail = translator.finish();
if !tail.is_empty() {
yield Ok(Bytes::from(tail));
}
}
if !settled {
if let Err(error) = settle_llm_router_usage(&state, owner_user_id.as_str()).await {
tracing::error!(
request_id = %request_id,
user_id = %owner_user_id,
error = %error,
"LLM Router 流式响应结束但累计额度同步未完成"
);
}
}
}
}
@@ -2054,7 +2082,7 @@ mod tests {
let (server_url, captured) = spawn_capturing_mock_server(MockResponse {
status_line: "200 OK",
content_type: "application/json; charset=utf-8",
body: r#"{"id":"msg_api_server_01","type":"message","role":"assistant","content":[{"type":"text","text":"pong"}]}"#.to_string(),
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(),
extra_headers: Vec::new(),
});
let (state, user_id) = seed_authenticated_state(AppConfig {
@@ -2095,8 +2123,14 @@ mod tests {
.expect("body should collect")
.to_bytes();
assert!(
String::from_utf8_lossy(&body).contains("msg_api_server_01"),
"anthropic payload should pass through"
String::from_utf8_lossy(&body).contains("chatcmpl_api_server_01"),
"anthropic 路由必须回 Anthropic 形状的 message(id 透传自上游)"
);
assert!(
String::from_utf8_lossy(&body).contains("\"type\":\"message\"")
&& String::from_utf8_lossy(&body).contains("\"pong\""),
"Anthropic 响应必须包含 role/content 结构:{}",
String::from_utf8_lossy(&body)
);
let captured_request = captured
@@ -2106,19 +2140,89 @@ mod tests {
.expect("upstream request should be captured");
let head = captured_request.to_ascii_lowercase();
assert!(
head.starts_with("post /v1/messages"),
"upstream must receive the Anthropic path: {captured_request}"
head.starts_with("post /v1/chat/completions"),
"网关必须把 Anthropic 请求转成 Router 可用的 Chat Completions:{captured_request}"
);
assert!(
head.contains("x-api-key: router-key"),
head.contains("authorization: bearer router-key"),
"upstream must receive the account Router key: {captured_request}"
);
assert!(
head.contains(r#""role":"user""#) && head.contains(r#""content":"ping""#),
"转换后的 Chat 请求必须带上用户消息:{captured_request}"
);
assert!(
!head.contains(&token.to_ascii_lowercase()),
"platform access token must never reach the Router"
);
}
#[tokio::test]
async fn llm_anthropic_messages_translates_chat_stream_into_anthropic_events() {
let streamed = concat!(
"data: {\"id\":\"chatcmpl_stream_01\",\"choices\":[{\"index\":0,\"delta\":{\"role\":\"assistant\",\"content\":\"po\"},\"finish_reason\":null}]}\n\n",
"data: {\"id\":\"chatcmpl_stream_01\",\"choices\":[{\"index\":0,\"delta\":{\"content\":\"ng\"},\"finish_reason\":null}]}\n\n",
"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 {
status_line: "200 OK",
content_type: "text/event-stream; charset=utf-8",
body: streamed.to_string(),
extra_headers: Vec::new(),
});
let (state, user_id) = seed_authenticated_state(AppConfig {
llm_router_base_url: server_url.clone(),
..AppConfig::default()
})
.await;
install_test_provisioned_router_credential(&user_id, server_url, "router-key");
let token = issue_access_token(&state, &user_id);
let app = build_router(state);
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri("/api/llm/anthropic/v1/messages")
.header("authorization", format!("Bearer {token}"))
.header("content-type", "application/json")
.body(Body::from(
json!({
"max_tokens": 16,
"stream": true,
"messages": [{ "role": "user", "content": "ping" }]
})
.to_string(),
))
.expect("request should build"),
)
.await
.expect("request should succeed");
assert_eq!(response.status(), StatusCode::OK);
let body = response
.into_body()
.collect()
.await
.expect("body should collect")
.to_bytes();
let text = String::from_utf8_lossy(&body).to_string();
for expected in [
"event: message_start",
"event: content_block_delta",
"\"text_delta\"",
"event: message_delta",
"\"stop_reason\":\"end_turn\"",
"event: message_stop",
] {
assert!(
text.contains(expected),
"缺少 Anthropic 事件 {expected}:{text}"
);
}
}
fn spawn_mock_server(responses: Vec<MockResponse>) -> String {
let listener = TcpListener::bind("127.0.0.1:0").expect("listener should bind");
let address = listener.local_addr().expect("listener should have addr");