From 21a8886c4c5d9c35090de65b5dee5357afe4096c Mon Sep 17 00:00:00 2001 From: kdletters <61648117+kdletters@users.noreply.github.com> Date: Fri, 2 Oct 2026 15:09:31 +0800 Subject: [PATCH 1/4] =?UTF-8?q?AGC=20Anthropic=20=E7=BD=91=E5=85=B3?= =?UTF-8?q?=E6=94=B9=E4=B8=BA=E6=A1=A5=E6=8E=A5=20Router=20=E7=9A=84=20Cha?= =?UTF-8?q?t=20Completions?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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 形状回包 --- .../api-server/src/llm/anthropic_bridge.rs | 873 ++++++++++++++++++ server-rs/crates/api-server/src/llm/mod.rs | 276 ++++-- 2 files changed, 1063 insertions(+), 86 deletions(-) create mode 100644 server-rs/crates/api-server/src/llm/anthropic_bridge.rs diff --git a/server-rs/crates/api-server/src/llm/anthropic_bridge.rs b/server-rs/crates/api-server/src/llm/anthropic_bridge.rs new file mode 100644 index 000000000..741f4a6eb --- /dev/null +++ b/server-rs/crates/api-server/src/llm/anthropic_bridge.rs @@ -0,0 +1,873 @@ +//! Anthropic Messages ⇄ OpenAI Chat Completions 的网关侧桥接。 +//! +//! 背景(2026-10-02 实测):AGC 的 Claude Code 执行器只会说 Anthropic Messages,而账号 +//! Router 的 `/v1/messages` 在现网不可用——最小的 Messages 请求也返回 +//! `500 not implemented`,Claude Code 的真实请求则是 60 秒收不到任何字节后被断开。 +//! 同一个账号、同一个 Router 的 `/v1/chat/completions` 是通的,所以网关在这里做一次 +//! 协议转换:入站 Anthropic Messages → 出站 OpenAI Chat Completions;回程再把 OpenAI 的 +//! SSE 翻译回 Anthropic 事件流(Claude Code 只认后者)。 +//! +//! 这一层只做协议形状转换,不做模型选择、鉴权和计费:模型名、Router 凭据与额度结算仍由 +//! 调用方沿用既有路径。 + +use serde_json::{Map, Value, json}; +use std::collections::HashMap; + +/// 工具名映射:OpenAI function name 只允许 `[A-Za-z0-9_-]` 且 ≤64 字符,而 Anthropic 侧的 +/// MCP 工具名带点号(例如 `mcp__agc__client.session.info`)。发出去前替换成安全名,回程用 +/// 同一份映射还原,保证模型看到的工具名与宿主注册的一致。 +#[derive(Debug, Default, Clone)] +pub(crate) struct AnthropicToolNames { + to_original: HashMap, +} + +impl AnthropicToolNames { + fn from_tools(tools: &[Value]) -> Self { + let mut to_original = HashMap::new(); + for tool in tools { + let Some(name) = tool.get("name").and_then(Value::as_str) else { + continue; + }; + to_original.insert(chat_tool_name(name), name.to_string()); + } + Self { to_original } + } + + pub(crate) fn original(&self, name: &str) -> String { + self.to_original + .get(name) + .cloned() + .unwrap_or_else(|| name.to_string()) + } +} + +/// Anthropic 工具名 → OpenAI function name。 +pub(crate) fn chat_tool_name(name: &str) -> String { + let mut sanitized: String = name + .chars() + .map(|value| { + if value.is_ascii_alphanumeric() || value == '_' || value == '-' { + value + } else { + '_' + } + }) + .collect(); + if sanitized.is_empty() { + sanitized.push_str("tool"); + } + if sanitized.len() > 64 { + // 截断会撞名,所以用稳定指纹补在尾部。 + let digest = blake_short(name); + sanitized.truncate(64 - digest.len() - 1); + sanitized.push('_'); + sanitized.push_str(&digest); + } + sanitized +} + +fn blake_short(value: &str) -> String { + use sha2::{Digest, Sha256}; + let digest = Sha256::digest(value.as_bytes()); + digest[..4] + .iter() + .map(|byte| format!("{byte:02x}")) + .collect() +} + +fn text_of(value: &Value) -> Option { + match value { + Value::String(text) => Some(text.clone()), + Value::Array(parts) => { + let mut text = String::new(); + for part in parts { + match part.get("type").and_then(Value::as_str) { + Some("text") => { + if let Some(value) = part.get("text").and_then(Value::as_str) { + text.push_str(value); + } + } + Some("image") | Some("document") => text.push_str("[附件]"), + _ => {} + } + } + Some(text) + } + _ => None, + } +} + +fn system_text(system: &Value) -> Option { + match system { + Value::Null => None, + Value::String(text) => (!text.trim().is_empty()).then(|| text.clone()), + Value::Array(blocks) => { + let mut text = String::new(); + for block in blocks { + let is_text = matches!( + block.get("type").and_then(Value::as_str), + Some("text") | None + ); + if !is_text { + continue; + } + if let Some(value) = block.get("text").and_then(Value::as_str) { + if !text.is_empty() { + text.push_str("\n\n"); + } + text.push_str(value); + } + } + (!text.trim().is_empty()).then_some(text) + } + _ => None, + } +} + +fn image_part(block: &Value) -> Option { + let source = block.get("source")?; + let url = match source.get("type").and_then(Value::as_str) { + Some("base64") => { + let media_type = source + .get("media_type") + .and_then(Value::as_str) + .unwrap_or("image/png"); + let data = source.get("data").and_then(Value::as_str)?; + format!("data:{media_type};base64,{data}") + } + Some("url") => source.get("url").and_then(Value::as_str)?.to_string(), + _ => return None, + }; + Some(json!({ "type": "image_url", "image_url": { "url": url } })) +} + +fn tool_parameters(input_schema: Option<&Value>) -> Value { + match input_schema { + Some(schema @ Value::Object(_)) => schema.clone(), + _ => json!({ "type": "object", "properties": {} }), + } +} + +/// Anthropic Messages 请求 → OpenAI Chat Completions 请求体。 +pub(crate) fn anthropic_messages_to_chat_completions( + request: &Value, + model: &str, +) -> Result<(Value, AnthropicToolNames), String> { + let object = request + .as_object() + .ok_or_else(|| "Anthropic 请求体必须是 JSON 对象".to_string())?; + let messages = object + .get("messages") + .and_then(Value::as_array) + .ok_or_else(|| "Anthropic 请求缺少 messages".to_string())?; + let tools = object + .get("tools") + .and_then(Value::as_array) + .cloned() + .unwrap_or_default(); + let tool_names = AnthropicToolNames::from_tools(&tools); + + let mut out_messages: Vec = Vec::new(); + if let Some(system) = object.get("system").and_then(system_text) { + // 只有真的带了 system 才补:Claude Code 每条请求都带,空白 system 会被上游当成噪声。 + out_messages.push(json!({ "role": "system", "content": system })); + } + + let mut pending_parts: Vec = Vec::new(); + let mut pending_tool_calls: Vec = Vec::new(); + let mut pending_role = "user"; + + fn flush( + out: &mut Vec, + role: &mut &str, + parts: &mut Vec, + tool_calls: &mut Vec, + ) { + if *role == "assistant" { + let text = parts + .iter() + .filter_map(|part| part.get("text").and_then(Value::as_str)) + .collect::>() + .join(""); + if text.is_empty() && tool_calls.is_empty() { + return; + } + let mut message = Map::new(); + message.insert("role".into(), json!("assistant")); + message.insert( + "content".into(), + if text.is_empty() { + Value::Null + } else { + json!(text) + }, + ); + if !tool_calls.is_empty() { + message.insert("tool_calls".into(), json!(tool_calls)); + tool_calls.clear(); + } + out.push(Value::Object(message)); + parts.clear(); + return; + } + if parts.is_empty() { + return; + } + let only_text = parts.iter().all(|part| { + part.get("type").and_then(Value::as_str) == Some("text") && parts.len() == 1 + }); + let content = if only_text { + json!(parts[0].get("text").and_then(Value::as_str).unwrap_or("")) + } else { + json!(parts) + }; + out.push(json!({ "role": "user", "content": content })); + parts.clear(); + } + + for message in messages { + let role = message + .get("role") + .and_then(Value::as_str) + .unwrap_or("user"); + let is_assistant = role == "assistant"; + if is_assistant != (pending_role == "assistant") && !pending_parts.is_empty() { + flush( + &mut out_messages, + &mut pending_role, + &mut pending_parts, + &mut pending_tool_calls, + ); + } + pending_role = if is_assistant { "assistant" } else { "user" }; + + let content = message.get("content").cloned().unwrap_or(Value::Null); + let blocks: Vec = match content { + Value::String(text) => { + pending_parts.push(json!({ "type": "text", "text": text })); + continue; + } + Value::Array(blocks) => blocks, + _ => continue, + }; + + for block in blocks { + match block.get("type").and_then(Value::as_str) { + Some("text") | None => { + if let Some(text) = block.get("text").and_then(Value::as_str) { + if !text.is_empty() { + pending_parts.push(json!({ "type": "text", "text": text })); + } + } + } + Some("image") => { + if let Some(part) = image_part(&block) { + pending_parts.push(part); + } + } + Some("tool_use") => { + let id = block + .get("id") + .and_then(Value::as_str) + .unwrap_or("call_unknown"); + let name = block.get("name").and_then(Value::as_str).unwrap_or("tool"); + let arguments = block.get("input").cloned().unwrap_or_else(|| json!({})); + pending_tool_calls.push(json!({ + "id": id, + "type": "function", + "function": { + "name": chat_tool_name(name), + "arguments": arguments.to_string(), + } + })); + } + Some("tool_result") => { + // 工具结果在 OpenAI 里是独立的 `tool` 消息:先把已积累的文本/图片收口, + // 再按顺序发出去,保持与 Anthropic 侧一致的时序。 + flush( + &mut out_messages, + &mut pending_role, + &mut pending_parts, + &mut pending_tool_calls, + ); + let tool_use_id = block + .get("tool_use_id") + .and_then(Value::as_str) + .unwrap_or("call_unknown"); + let mut text = String::new(); + let mut images: Vec = Vec::new(); + match block.get("content") { + Some(Value::Array(parts)) => { + for part in parts { + match part.get("type").and_then(Value::as_str) { + Some("text") => { + if let Some(value) = + part.get("text").and_then(Value::as_str) + { + text.push_str(value); + } + } + Some("image") => { + if let Some(image) = image_part(part) { + images.push(image); + } + } + _ => {} + } + } + } + Some(other) => { + if let Some(value) = text_of(other) { + text.push_str(&value); + } + } + None => {} + } + if text.is_empty() && images.is_empty() { + text.push_str("(无输出)"); + } + out_messages.push(json!({ + "role": "tool", + "tool_call_id": tool_use_id, + "content": text, + })); + if !images.is_empty() { + // OpenAI 的工具消息不能带图片:图片另起一条 user 消息,模型仍能看到。 + let mut parts = vec![json!({ + "type": "text", + "text": format!("工具 {tool_use_id} 的截图结果:") + })]; + parts.extend(images); + out_messages.push(json!({ "role": "user", "content": parts })); + } + } + _ => {} + } + } + if !is_assistant { + flush( + &mut out_messages, + &mut pending_role, + &mut pending_parts, + &mut pending_tool_calls, + ); + } + } + flush( + &mut out_messages, + &mut pending_role, + &mut pending_parts, + &mut pending_tool_calls, + ); + + let mut body = Map::new(); + body.insert("model".into(), json!(model)); + body.insert("messages".into(), json!(out_messages)); + body.insert( + "stream".into(), + json!( + object + .get("stream") + .and_then(Value::as_bool) + .unwrap_or(false) + ), + ); + if object.get("stream").and_then(Value::as_bool) == Some(true) { + body.insert("stream_options".into(), json!({ "include_usage": true })); + } + for key in ["max_tokens", "temperature", "top_p"] { + if let Some(value) = object.get(key) { + if !value.is_null() { + body.insert(key.to_string(), value.clone()); + } + } + } + if let Some(stops) = object.get("stop_sequences").and_then(Value::as_array) { + if !stops.is_empty() { + body.insert("stop".into(), json!(stops)); + } + } + if !tools.is_empty() { + let converted: Vec = tools + .iter() + .filter_map(|tool| { + let name = tool.get("name").and_then(Value::as_str)?; + let mut function = Map::new(); + function.insert("name".into(), json!(chat_tool_name(name))); + if let Some(description) = tool.get("description").and_then(Value::as_str) { + function.insert("description".into(), json!(description)); + } + function.insert( + "parameters".into(), + tool_parameters(tool.get("input_schema")), + ); + Some(json!({ "type": "function", "function": function })) + }) + .collect(); + body.insert("tools".into(), json!(converted)); + } + if let Some(choice) = object.get("tool_choice") { + let kind = choice.get("type").and_then(Value::as_str).unwrap_or("auto"); + let converted = match kind { + "any" => json!("required"), + "none" => json!("none"), + "tool" => { + let name = choice + .get("name") + .and_then(Value::as_str) + .unwrap_or_default(); + json!({ "type": "function", "function": { "name": chat_tool_name(name) } }) + } + _ => json!("auto"), + }; + body.insert("tool_choice".into(), converted); + } + Ok((Value::Object(body), tool_names)) +} + +fn stop_reason_of(finish_reason: Option<&str>, saw_tool_call: bool) -> &'static str { + match finish_reason { + Some("length") => "max_tokens", + Some("tool_calls") => "tool_use", + Some("stop") => { + if saw_tool_call { + "tool_use" + } else { + "end_turn" + } + } + _ => { + if saw_tool_call { + "tool_use" + } else { + "end_turn" + } + } + } +} + +fn content_blocks_of_message(message: &Value, tool_names: &AnthropicToolNames) -> Vec { + let mut blocks = Vec::new(); + if let Some(text) = message.get("content").and_then(|value| match value { + Value::String(text) => Some(text.clone()), + Value::Null => None, + other => text_of(other), + }) { + if !text.is_empty() { + blocks.push(json!({ "type": "text", "text": text })); + } + } + if let Some(calls) = message.get("tool_calls").and_then(Value::as_array) { + for call in calls { + let name = call + .pointer("/function/name") + .and_then(Value::as_str) + .unwrap_or("tool"); + let arguments = call + .pointer("/function/arguments") + .and_then(Value::as_str) + .and_then(|raw| serde_json::from_str::(raw).ok()) + .unwrap_or_else(|| json!({})); + blocks.push(json!({ + "type": "tool_use", + "id": call.get("id").and_then(Value::as_str).unwrap_or("call_unknown"), + "name": tool_names.original(name), + "input": arguments, + })); + } + } + blocks +} + +/// OpenAI Chat Completions 响应体 → Anthropic Messages 响应体(非流式)。 +pub(crate) fn chat_completions_body_to_anthropic_message( + value: &Value, + model: &str, + tool_names: &AnthropicToolNames, +) -> Value { + let choice = value + .get("choices") + .and_then(Value::as_array) + .and_then(|choices| choices.first()); + let message = choice + .and_then(|choice| choice.get("message")) + .cloned() + .unwrap_or(Value::Null); + let blocks = content_blocks_of_message(&message, tool_names); + let saw_tool_call = !message + .get("tool_calls") + .and_then(Value::as_array) + .unwrap_or(&Vec::new()) + .is_empty(); + let finish_reason = choice + .and_then(|choice| choice.get("finish_reason")) + .and_then(Value::as_str); + json!({ + "id": value.get("id").and_then(Value::as_str).unwrap_or("msg_agc_bridge"), + "type": "message", + "role": "assistant", + "model": model, + "content": blocks, + "stop_reason": stop_reason_of(finish_reason, saw_tool_call), + "stop_sequence": Value::Null, + "usage": { + "input_tokens": value.pointer("/usage/prompt_tokens").and_then(Value::as_u64).unwrap_or(0), + "output_tokens": value.pointer("/usage/completion_tokens").and_then(Value::as_u64).unwrap_or(0), + } + }) +} + +fn sse_event(name: &str, payload: &Value) -> Vec { + format!("event: {name}\ndata: {payload}\n\n").into_bytes() +} + +/// OpenAI Chat Completions SSE → Anthropic Messages SSE 的状态机。 +/// +/// 拆成同步状态机而不是直接的 async 流,是为了能在单元测试里逐块喂字节、逐块断言事件。 +#[derive(Debug)] +pub(crate) struct ChatToAnthropicStream { + tool_names: AnthropicToolNames, + model: String, + message_id: String, + buffer: String, + started: bool, + finished: bool, + next_index: u32, + open_block: Option, + open_tool_index: Option, + tool_blocks: HashMap, + usage_input: u64, + usage_output: u64, + saw_tool_call: bool, + finish_reason: Option, +} + +impl ChatToAnthropicStream { + pub(crate) fn new(model: &str, tool_names: AnthropicToolNames) -> Self { + Self { + tool_names, + model: model.to_string(), + message_id: format!("msg_agc_{}", uuid::Uuid::new_v4().simple()), + buffer: String::new(), + started: false, + finished: false, + next_index: 0, + open_block: None, + open_tool_index: None, + tool_blocks: HashMap::new(), + usage_input: 0, + usage_output: 0, + saw_tool_call: false, + finish_reason: None, + } + } + + /// 喂入一段上游字节,返回可以立刻下发给 Claude Code 的 Anthropic SSE 字节。 + pub(crate) fn push(&mut self, chunk: &[u8]) -> Vec { + self.buffer.push_str(&String::from_utf8_lossy(chunk)); + let mut out = Vec::new(); + loop { + let Some(index) = self.buffer.find('\n') else { + break; + }; + let line = self.buffer[..index].trim_end_matches('\r').to_string(); + self.buffer.drain(..=index); + let Some(payload) = line.strip_prefix("data:") else { + continue; + }; + let payload = payload.trim(); + if payload.is_empty() { + continue; + } + if payload == "[DONE]" { + out.extend(self.finish()); + continue; + } + let Ok(value) = serde_json::from_str::(payload) else { + continue; + }; + out.extend(self.handle_chunk(&value)); + } + out + } + + /// 收尾:关闭未闭合的块并写出 `message_delta` / `message_stop`。 + pub(crate) fn finish(&mut self) -> Vec { + if self.finished { + return Vec::new(); + } + self.finished = true; + let mut out = Vec::new(); + if !self.started { + out.extend(self.start_message()); + } + if let Some(index) = self.open_block.take() { + out.extend(sse_event( + "content_block_stop", + &json!({ "type": "content_block_stop", "index": index }), + )); + } + let stop_reason = stop_reason_of(self.finish_reason.as_deref(), self.saw_tool_call); + out.extend(sse_event( + "message_delta", + &json!({ + "type": "message_delta", + "delta": { "stop_reason": stop_reason, "stop_sequence": Value::Null }, + "usage": { "output_tokens": self.usage_output }, + }), + )); + out.extend(sse_event( + "message_stop", + &json!({ "type": "message_stop" }), + )); + out + } + + pub(crate) fn is_finished(&self) -> bool { + self.finished + } + + fn start_message(&mut self) -> Vec { + self.started = true; + sse_event( + "message_start", + &json!({ + "type": "message_start", + "message": { + "id": self.message_id, + "type": "message", + "role": "assistant", + "model": self.model, + "content": [], + "stop_reason": Value::Null, + "stop_sequence": Value::Null, + "usage": { "input_tokens": self.usage_input, "output_tokens": 0 }, + } + }), + ) + } + + fn handle_chunk(&mut self, value: &Value) -> Vec { + let mut out = Vec::new(); + if !self.started { + if let Some(id) = value.get("id").and_then(Value::as_str) { + self.message_id = id.to_string(); + } + out.extend(self.start_message()); + } + if let Some(prompt) = value + .pointer("/usage/prompt_tokens") + .and_then(Value::as_u64) + { + self.usage_input = prompt; + } + if let Some(completion) = value + .pointer("/usage/completion_tokens") + .and_then(Value::as_u64) + { + self.usage_output = completion; + } + + let Some(choice) = value + .get("choices") + .and_then(Value::as_array) + .and_then(|choices| choices.first()) + else { + return out; + }; + if let Some(reason) = choice.get("finish_reason").and_then(Value::as_str) { + self.finish_reason = Some(reason.to_string()); + } + let Some(delta) = choice.get("delta") else { + return out; + }; + + if let Some(text) = delta.get("content").and_then(Value::as_str) { + if !text.is_empty() { + if self.open_block.is_none() { + let index = self.next_index; + self.next_index += 1; + self.open_block = Some(index); + out.extend(sse_event( + "content_block_start", + &json!({ + "type": "content_block_start", + "index": index, + "content_block": { "type": "text", "text": "" } + }), + )); + } + if let Some(index) = self.open_block { + out.extend(sse_event( + "content_block_delta", + &json!({ + "type": "content_block_delta", + "index": index, + "delta": { "type": "text_delta", "text": text } + }), + )); + } + } + } + + if let Some(calls) = delta.get("tool_calls").and_then(Value::as_array) { + for call in calls { + let slot = call.get("index").and_then(Value::as_u64).unwrap_or(0); + self.saw_tool_call = true; + let name = call + .pointer("/function/name") + .and_then(Value::as_str) + .map(|name| self.tool_names.original(name)); + if let Some(name) = name { + // 新的工具调用块:先关掉上一个块(文本或上一个工具)。 + if let Some(index) = self.open_block.take() { + out.extend(sse_event( + "content_block_stop", + &json!({ "type": "content_block_stop", "index": index }), + )); + } + let index = self.next_index; + self.next_index += 1; + self.tool_blocks.insert(slot, index); + self.open_block = Some(index); + self.open_tool_index = Some(slot); + out.extend(sse_event( + "content_block_start", + &json!({ + "type": "content_block_start", + "index": index, + "content_block": { + "type": "tool_use", + "id": call.get("id").and_then(Value::as_str).unwrap_or("call_unknown"), + "name": name, + "input": {} + } + }), + )); + } + if let Some(arguments) = call.pointer("/function/arguments").and_then(Value::as_str) + { + if !arguments.is_empty() { + let index = self.tool_blocks.get(&slot).copied().or(self.open_block); + if let Some(index) = index { + out.extend(sse_event( + "content_block_delta", + &json!({ + "type": "content_block_delta", + "index": index, + "delta": { + "type": "input_json_delta", + "partial_json": arguments + } + }), + )); + } + } + } + } + } + out + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn mcp_tool_names_survive_the_openai_name_constraint() { + // MCP 工具名带点号,OpenAI function name 不允许;转换后仍要能还原。 + let name = "mcp__agc__client.session.info"; + let sanitized = chat_tool_name(name); + assert_eq!(sanitized, "mcp__agc__client_session_info"); + let tools = vec![json!({ "name": name, "input_schema": { "type": "object" } })]; + let names = AnthropicToolNames::from_tools(&tools); + assert_eq!(names.original(&sanitized), name); + } + + #[test] + fn anthropic_request_converts_system_tools_and_tool_results() { + let request = json!({ + "model": "ignored", + "system": [{ "type": "text", "text": "系统提示", "cache_control": { "type": "ephemeral" } }], + "max_tokens": 128, + "stream": true, + "tools": [{ "name": "agc_read_file", "description": "读文件", "input_schema": { "type": "object", "properties": { "path": { "type": "string" } } } }], + "messages": [ + { "role": "user", "content": [{ "type": "text", "text": "改一下" }] }, + { "role": "assistant", "content": [ + { "type": "text", "text": "先读文件" }, + { "type": "tool_use", "id": "call_1", "name": "agc_read_file", "input": { "path": "a.ts" } } + ] }, + { "role": "user", "content": [ + { "type": "tool_result", "tool_use_id": "call_1", "content": [{ "type": "text", "text": "文件内容" }] }, + { "type": "text", "text": "继续" } + ] } + ] + }); + let (converted, names) = + anthropic_messages_to_chat_completions(&request, "claude-opus-5-5").expect("convert"); + assert_eq!(converted["model"], "claude-opus-5-5"); + assert_eq!(converted["stream"], true); + assert_eq!(converted["stream_options"]["include_usage"], true); + assert_eq!(converted["max_tokens"], 128); + let messages = converted["messages"].as_array().expect("messages"); + assert_eq!(messages[0]["role"], "system"); + assert_eq!(messages[0]["content"], "系统提示"); + // assistant 的 tool_use 必须落到 tool_calls,且函数名经过安全化。 + let assistant = messages + .iter() + .find(|message| message["role"] == "assistant") + .expect("assistant message"); + assert_eq!( + assistant["tool_calls"][0]["function"]["name"], + "agc_read_file" + ); + assert_eq!(names.original("agc_read_file"), "agc_read_file"); + // tool_result 必须变成独立 tool 消息,且原样带上 tool_call_id。 + let tool_message = messages + .iter() + .find(|message| message["role"] == "tool") + .expect("tool message"); + assert_eq!(tool_message["tool_call_id"], "call_1"); + assert_eq!(tool_message["content"], "文件内容"); + assert_eq!(converted["tools"][0]["function"]["name"], "agc_read_file"); + } + + #[test] + fn chat_stream_translates_text_and_tool_calls() { + let tools = vec![json!({ "name": "mcp__agc__client.session.info" })]; + let names = AnthropicToolNames::from_tools(&tools); + let mut stream = ChatToAnthropicStream::new("claude-opus-5-5", names); + + let first = stream.push( + b"data: {\"id\":\"chatcmpl_1\",\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hello\"},\"finish_reason\":null}]}\n\n", + ); + let first = String::from_utf8_lossy(&first).to_string(); + assert!(first.contains("event: message_start"), "{first}"); + assert!( + first.contains("\"text_delta\"") && first.contains("\"text\":\"hello\""), + "{first}" + ); + + let second = stream.push( + b"data: {\"id\":\"chatcmpl_1\",\"choices\":[{\"index\":0,\"delta\":{\"tool_calls\":[{\"index\":0,\"id\":\"call_9\",\"function\":{\"name\":\"mcp__agc__client_session_info\",\"arguments\":\"{\\\"path\\\":\"}}]},\"finish_reason\":null}]}\n\n", + ); + let second = String::from_utf8_lossy(&second).to_string(); + // 工具名必须还原成宿主注册的原名(点号/下划线都要还原回 MCP 命名)。 + assert!( + second.contains("\"name\":\"mcp__agc__client.session.info\""), + "{second}" + ); + assert!(second.contains("input_json_delta"), "{second}"); + + let mut tail = stream.push( + b"data: {\"id\":\"chatcmpl_1\",\"choices\":[{\"index\":0,\"delta\":{},\"finish_reason\":\"tool_calls\"}],\"usage\":{\"prompt_tokens\":7,\"completion_tokens\":3}}\n\n", + ); + tail.extend(stream.push(b"data: [DONE]\n\n")); + let tail = String::from_utf8_lossy(&tail).to_string(); + assert!(tail.contains("\"stop_reason\":\"tool_use\""), "{tail}"); + assert!(tail.contains("event: message_stop"), "{tail}"); + assert!(stream.is_finished()); + } +} diff --git a/server-rs/crates/api-server/src/llm/mod.rs b/server-rs/crates/api-server/src/llm/mod.rs index e43cf2e8f..f231079e8 100644 --- a/server-rs/crates/api-server/src/llm/mod.rs +++ b/server-rs/crates/api-server/src/llm/mod.rs @@ -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, Extension(request_context): Extension, Extension(authenticated): Extension, - headers: HeaderMap, + _headers: HeaderMap, body: Bytes, ) -> Result { 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, 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> + Unpin, + model: String, + tool_names: anthropic_bridge::AnthropicToolNames, 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 = 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) -> String { let listener = TcpListener::bind("127.0.0.1:0").expect("listener should bind"); let address = listener.local_addr().expect("listener should have addr"); From 580c74759a04f76a143803fdb5f0f18e41176ba8 Mon Sep 17 00:00:00 2001 From: kdletters <61648117+kdletters@users.noreply.github.com> Date: Fri, 2 Oct 2026 15:29:42 +0800 Subject: [PATCH 2/4] =?UTF-8?q?Anthropic=20=E6=A1=A5=E6=8E=A5=E5=AF=B9?= =?UTF-8?q?=E4=B8=8A=E6=B8=B8=E4=B8=80=E5=BE=8B=E6=B5=81=E5=BC=8F=EF=BC=8C?= =?UTF-8?q?=E5=B9=B6=E8=A1=A5=E9=9D=9E=E6=B5=81=E5=BC=8F=E8=81=9A=E5=90=88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - anthropic_bridge:向 Router 一律请求 stream=true,兼容只接受流式 chat 的上游 - anthropic_bridge:新增 SSE 聚合,客户端要非流式时在网关侧拼回 OpenAI 形状再转 Anthropic message - llm:非流式分支按上游返回形状选择直读 JSON 或聚合 SSE - 测试:新增聚合单测(文本 + tool_calls + usage → Anthropic tool_use 还原) --- .../api-server/src/llm/anthropic_bridge.rs | 172 ++++++++++++++++-- server-rs/crates/api-server/src/llm/mod.rs | 9 +- 2 files changed, 167 insertions(+), 14 deletions(-) diff --git a/server-rs/crates/api-server/src/llm/anthropic_bridge.rs b/server-rs/crates/api-server/src/llm/anthropic_bridge.rs index 741f4a6eb..9e39c2125 100644 --- a/server-rs/crates/api-server/src/llm/anthropic_bridge.rs +++ b/server-rs/crates/api-server/src/llm/anthropic_bridge.rs @@ -363,18 +363,10 @@ pub(crate) fn anthropic_messages_to_chat_completions( let mut body = Map::new(); body.insert("model".into(), json!(model)); body.insert("messages".into(), json!(out_messages)); - body.insert( - "stream".into(), - json!( - object - .get("stream") - .and_then(Value::as_bool) - .unwrap_or(false) - ), - ); - if object.get("stream").and_then(Value::as_bool) == Some(true) { - body.insert("stream_options".into(), json!({ "include_usage": true })); - } + // 上游要求流式(Router 侧非流式 chat 会被打回 `400 Stream must be set to true`), + // 所以对 Router 一律按流式请求;客户端要不要流式由网关这边决定(非流式就地聚合)。 + body.insert("stream".into(), json!(true)); + body.insert("stream_options".into(), json!({ "include_usage": true })); for key in ["max_tokens", "temperature", "top_p"] { if let Some(value) = object.get(key) { if !value.is_null() { @@ -521,6 +513,129 @@ fn sse_event(name: &str, payload: &Value) -> Vec { format!("event: {name}\ndata: {payload}\n\n").into_bytes() } +/// 把上游的 Chat Completions SSE 聚合成一份非流式响应体。 +/// +/// 客户端要非流式(例如平台自己的一次性调用)时,网关仍按流式向上游取数——这家上游不接受 +/// 非流式 chat——然后在这里拼回 OpenAI 形状,再交给同一条 Anthropic 转换。 +pub(crate) fn assemble_chat_completions_stream(bytes: &[u8]) -> Value { + let text = String::from_utf8_lossy(bytes); + let mut id = "chatcmpl_agc_bridge".to_string(); + let mut content = String::new(); + let mut tool_calls: Vec = Vec::new(); + let mut arguments: HashMap = HashMap::new(); + let mut finish_reason: Option = None; + let mut prompt_tokens = 0_u64; + let mut completion_tokens = 0_u64; + + for line in text.lines() { + let Some(payload) = line.trim().strip_prefix("data:") else { + continue; + }; + let payload = payload.trim(); + if payload.is_empty() || payload == "[DONE]" { + continue; + } + let Ok(value) = serde_json::from_str::(payload) else { + continue; + }; + if let Some(value_id) = value.get("id").and_then(Value::as_str) { + id = value_id.to_string(); + } + if let Some(prompt) = value + .pointer("/usage/prompt_tokens") + .and_then(Value::as_u64) + { + prompt_tokens = prompt; + } + if let Some(completion) = value + .pointer("/usage/completion_tokens") + .and_then(Value::as_u64) + { + completion_tokens = completion; + } + let Some(choice) = value + .get("choices") + .and_then(Value::as_array) + .and_then(|choices| choices.first()) + else { + continue; + }; + if let Some(reason) = choice.get("finish_reason").and_then(Value::as_str) { + finish_reason = Some(reason.to_string()); + } + let Some(delta) = choice.get("delta") else { + continue; + }; + if let Some(chunk) = delta.get("content").and_then(Value::as_str) { + content.push_str(chunk); + } + if let Some(calls) = delta.get("tool_calls").and_then(Value::as_array) { + for call in calls { + let slot = call.get("index").and_then(Value::as_u64).unwrap_or(0); + if let Some(fragment) = call.pointer("/function/arguments").and_then(Value::as_str) + { + arguments.entry(slot).or_default().push_str(fragment); + } + let name = call + .pointer("/function/name") + .and_then(Value::as_str) + .map(str::to_string); + let existing = tool_calls + .iter_mut() + .find(|entry| entry.get("index").and_then(Value::as_u64) == Some(slot)); + match existing { + Some(entry) => { + if let Some(name) = name { + entry["function"]["name"] = json!(name); + } + if let Some(call_id) = call.get("id").and_then(Value::as_str) { + entry["id"] = json!(call_id); + } + } + None => tool_calls.push(json!({ + "index": slot, + "id": call.get("id").and_then(Value::as_str).unwrap_or("call_unknown"), + "type": "function", + "function": { "name": name.unwrap_or_else(|| "tool".to_string()), "arguments": "" }, + })), + } + } + } + } + + for call in tool_calls.iter_mut() { + let slot = call.get("index").and_then(Value::as_u64).unwrap_or(0); + call["function"]["arguments"] = json!(arguments.get(&slot).cloned().unwrap_or_default()); + if let Some(object) = call.as_object_mut() { + object.remove("index"); + } + } + + let mut message = Map::new(); + message.insert("role".into(), json!("assistant")); + message.insert( + "content".into(), + if content.is_empty() { + Value::Null + } else { + json!(content) + }, + ); + if !tool_calls.is_empty() { + message.insert("tool_calls".into(), json!(tool_calls)); + } + json!({ + "id": id, + "object": "chat.completion", + "choices": [{ + "index": 0, + "message": Value::Object(message), + "finish_reason": finish_reason.unwrap_or_else(|| "stop".to_string()), + }], + "usage": { "prompt_tokens": prompt_tokens, "completion_tokens": completion_tokens } + }) +} + /// OpenAI Chat Completions SSE → Anthropic Messages SSE 的状态机。 /// /// 拆成同步状态机而不是直接的 async 流,是为了能在单元测试里逐块喂字节、逐块断言事件。 @@ -870,4 +985,37 @@ mod tests { assert!(tail.contains("event: message_stop"), "{tail}"); assert!(stream.is_finished()); } + + #[test] + fn chat_stream_assembles_into_a_non_streaming_response() { + let streamed = concat!( + "data: {\"id\":\"chatcmpl_asm\",\"choices\":[{\"index\":0,\"delta\":{\"content\":\"po\"},\"finish_reason\":null}]}\n\n", + "data: {\"id\":\"chatcmpl_asm\",\"choices\":[{\"index\":0,\"delta\":{\"content\":\"ng\"},\"finish_reason\":null}]}\n\n", + "data: {\"id\":\"chatcmpl_asm\",\"choices\":[{\"index\":0,\"delta\":{\"tool_calls\":[{\"index\":0,\"id\":\"call_1\",\"function\":{\"name\":\"mcp__agc__client_session_info\",\"arguments\":\"{\\\"a\\\":\"}}]},\"finish_reason\":null}]}\n\n", + "data: {\"id\":\"chatcmpl_asm\",\"choices\":[{\"index\":0,\"delta\":{\"tool_calls\":[{\"index\":0,\"function\":{\"arguments\":\"1}\"}}]},\"finish_reason\":\"tool_calls\"}]}\n\n", + "data: {\"id\":\"chatcmpl_asm\",\"choices\":[],\"usage\":{\"prompt_tokens\":9,\"completion_tokens\":4}}\n\n", + "data: [DONE]\n\n", + ); + let assembled = assemble_chat_completions_stream(streamed.as_bytes()); + assert_eq!(assembled["id"], "chatcmpl_asm"); + assert_eq!(assembled["choices"][0]["message"]["content"], "pong"); + assert_eq!(assembled["choices"][0]["finish_reason"], "tool_calls"); + assert_eq!( + assembled["choices"][0]["message"]["tool_calls"][0]["function"]["arguments"], + "{\"a\":1}" + ); + assert_eq!(assembled["usage"]["prompt_tokens"], 9); + + // 聚合结果必须能被同一条 Anthropic 转换吃下(工具名还原 + stop_reason 映射)。 + let tools = vec![json!({ "name": "mcp__agc__client.session.info" })]; + let names = AnthropicToolNames::from_tools(&tools); + let message = chat_completions_body_to_anthropic_message(&assembled, "m", &names); + assert_eq!(message["stop_reason"], "tool_use"); + assert_eq!(message["content"][1]["type"], "tool_use"); + assert_eq!( + message["content"][1]["name"], + "mcp__agc__client.session.info" + ); + assert_eq!(message["content"][1]["input"]["a"], 1); + } } diff --git a/server-rs/crates/api-server/src/llm/mod.rs b/server-rs/crates/api-server/src/llm/mod.rs index f231079e8..8112afec3 100644 --- a/server-rs/crates/api-server/src/llm/mod.rs +++ b/server-rs/crates/api-server/src/llm/mod.rs @@ -668,8 +668,13 @@ pub async fn proxy_llm_messages( "LLM Router Messages 成功但累计额度同步未完成" ); } - // 非流式也要交回 Anthropic 形状:Claude 客户端(含 SDK 探针)读的是 message/content。 - let upstream_value: Value = serde_json::from_slice(&body).unwrap_or(Value::Null); + // 上游一律流式(这家 Router 拒绝非流式 chat),所以这里先把 SSE 聚合成一份 + // OpenAI 形状,再交回 Anthropic message:Claude 客户端(含 SDK 探针)读 message/content。 + let upstream_value = if body.starts_with(b"{\"") { + serde_json::from_slice::(&body).unwrap_or(Value::Null) + } else { + anthropic_bridge::assemble_chat_completions_stream(&body) + }; let message = anthropic_bridge::chat_completions_body_to_anthropic_message( &upstream_value, &selected_model, From 3e9bc76c2c097229c1261e5253f573816948eb8b Mon Sep 17 00:00:00 2001 From: kdletters <61648117+kdletters@users.noreply.github.com> Date: Fri, 2 Oct 2026 15:49:01 +0800 Subject: [PATCH 3/4] =?UTF-8?q?Anthropic=20=E8=B7=AF=E7=94=B1=E6=94=B9?= =?UTF-8?q?=E5=9B=9E=E5=8E=9F=E7=94=9F=E4=BC=98=E5=85=88=E3=80=81=E6=A1=A5?= =?UTF-8?q?=E6=8E=A5=E5=85=9C=E5=BA=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - llm:/api/llm/anthropic/* 先按 Anthropic 原样直通 Router 的 /v1/messages(20 秒首包上限),拿到 2xx 就直接透传 - llm:直通返回非 2xx、连接失败或首包超时才回退到 Anthropic→Chat 协议桥接,并记录 model/status 便于定位 - 测试:两条 Anthropic 路由用例改为「原生先被拒 + 桥接命中」的两段式 mock --- server-rs/crates/api-server/src/llm/mod.rs | 240 ++++++++++++++++++++- 1 file changed, 231 insertions(+), 9 deletions(-) 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))) From 3a0d92b622a3c67b8d7355382f7b023579750a76 Mon Sep 17 00:00:00 2001 From: kdletters <61648117+kdletters@users.noreply.github.com> Date: Fri, 2 Oct 2026 16:14:56 +0800 Subject: [PATCH 4/4] =?UTF-8?q?Anthropic=20=E7=9B=B4=E9=80=9A=E5=8F=AA?= =?UTF-8?q?=E7=94=A8=20x-api-key=20=E5=87=BA=E7=A4=BA=20Router=20=E5=87=AD?= =?UTF-8?q?=E6=8D=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - llm:原生 Anthropic 直通不再同时带 Authorization: Bearer,只发 Anthropic 规范的 x-api-key - 实测:Router 对同一次 /v1/messages 请求,只带 x-api-key 时 200(4.7s 返回真实回复),同时带 Bearer 时 79 秒后被断开 --- server-rs/crates/api-server/src/llm/mod.rs | 1 - 1 file changed, 1 deletion(-) diff --git a/server-rs/crates/api-server/src/llm/mod.rs b/server-rs/crates/api-server/src/llm/mod.rs index 94f3be0bd..351ffbb52 100644 --- a/server-rs/crates/api-server/src/llm/mod.rs +++ b/server-rs/crates/api-server/src/llm/mod.rs @@ -606,7 +606,6 @@ pub async fn proxy_llm_messages( 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) {