Anthropic 桥接对上游一律流式,并补非流式聚合
Project CI / AI game creator shell Rust lane 2/2 (push) Failing after 2m19s
Project CI / AI game creator shell Rust lane 1/2 (push) Failing after 3m4s
Project CI / AI game creator shell Rust smoke (push) Successful in 2m30s
Project CI / AI game creator shell Rust crates (push) Successful in 1m17s
Project CI / Frontend tests (push) Successful in 3m14s
Project CI / Repository checks (push) Successful in 3m49s
Project CI / Backend tests (push) Successful in 5m45s
Project CI / AI game creator shell web tests (push) Successful in 1m51s
Project CI / Native shell tests (push) Successful in 6m47s
Project CI / AI game creator shell Rust lane 2/2 (push) Failing after 2m19s
Project CI / AI game creator shell Rust lane 1/2 (push) Failing after 3m4s
Project CI / AI game creator shell Rust smoke (push) Successful in 2m30s
Project CI / AI game creator shell Rust crates (push) Successful in 1m17s
Project CI / Frontend tests (push) Successful in 3m14s
Project CI / Repository checks (push) Successful in 3m49s
Project CI / Backend tests (push) Successful in 5m45s
Project CI / AI game creator shell web tests (push) Successful in 1m51s
Project CI / Native shell tests (push) Successful in 6m47s
- anthropic_bridge:向 Router 一律请求 stream=true,兼容只接受流式 chat 的上游 - anthropic_bridge:新增 SSE 聚合,客户端要非流式时在网关侧拼回 OpenAI 形状再转 Anthropic message - llm:非流式分支按上游返回形状选择直读 JSON 或聚合 SSE - 测试:新增聚合单测(文本 + tool_calls + usage → Anthropic tool_use 还原)
This commit is contained in:
@@ -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<u8> {
|
||||
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<Value> = Vec::new();
|
||||
let mut arguments: HashMap<u64, String> = HashMap::new();
|
||||
let mut finish_reason: Option<String> = 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::<Value>(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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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::<Value>(&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,
|
||||
|
||||
Reference in New Issue
Block a user