checkpoint before gfm ast parser design
This commit is contained in:
@@ -53,6 +53,21 @@ pub async fn next_ai_agent_run(
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(provider) = explicit_agent_provider(&payload) {
|
||||
if provider == "hermes" {
|
||||
return run_direct_hermes_agent(&context, &payload).await;
|
||||
}
|
||||
return Err(WebError::bad_gateway_code(
|
||||
"ai_provider_bridge_unavailable",
|
||||
format!(
|
||||
"{provider} provider 需要可用的 Next AI bridge,不能静默降级到本地页面工具 host。"
|
||||
),
|
||||
)
|
||||
.with_context(&context)
|
||||
.with_header("x-mnote-web-owner", "mnote-web")
|
||||
.with_header("x-mnote-ai-execution-owner", "provider-bridge-unavailable"));
|
||||
}
|
||||
|
||||
if let Ok(response) = run_local_mnote_cli_ai_host(&state, &context, &payload).await {
|
||||
return Ok(response);
|
||||
}
|
||||
@@ -118,6 +133,304 @@ pub async fn next_ai_agent_run(
|
||||
build_ai_sse_proxy_response(upstream_response, &context)
|
||||
}
|
||||
|
||||
fn explicit_agent_provider(payload: &Value) -> Option<&'static str> {
|
||||
let provider = payload
|
||||
.get("options")
|
||||
.and_then(|value| value.get("ai"))
|
||||
.and_then(|value| value.get("provider"))
|
||||
.and_then(Value::as_str)?
|
||||
.trim()
|
||||
.to_ascii_lowercase();
|
||||
match provider.as_str() {
|
||||
"hermes" => Some("hermes"),
|
||||
"codex" => Some("codex"),
|
||||
"claudecode" => Some("claudecode"),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_direct_hermes_agent(
|
||||
context: &RequestContext,
|
||||
payload: &Value,
|
||||
) -> Result<Response, WebError> {
|
||||
let base_url = resolve_hermes_api_base_url();
|
||||
let api_key = read_env_or_dotenv("MNOTE_HERMES_API_KEY").ok_or_else(|| {
|
||||
WebError::bad_gateway_code(
|
||||
"hermes_bridge_unconfigured",
|
||||
"未配置 MNOTE_HERMES_API_KEY,无法直连 Hermes API。",
|
||||
)
|
||||
.with_context(context)
|
||||
.with_header("x-mnote-web-owner", "mnote-web")
|
||||
.with_header("x-mnote-ai-execution-owner", "hermes")
|
||||
})?;
|
||||
let client = reqwest::Client::builder()
|
||||
.connect_timeout(Duration::from_secs(5))
|
||||
.timeout(Duration::from_secs(180))
|
||||
.redirect(reqwest::redirect::Policy::none())
|
||||
.build()
|
||||
.map_err(|error| WebError::internal(format!("Hermes HTTP 客户端创建失败: {error}")))?;
|
||||
|
||||
let start_url = reqwest::Url::parse(&format!("{}/v1/runs", base_url.trim_end_matches('/')))
|
||||
.map_err(|error| WebError::internal(format!("Hermes runs URL 非法: {error}")))?;
|
||||
let start_response = client
|
||||
.post(start_url)
|
||||
.bearer_auth(&api_key)
|
||||
.json(&build_hermes_run_payload(payload))
|
||||
.send()
|
||||
.await
|
||||
.map_err(|error| {
|
||||
WebError::bad_gateway_code(
|
||||
"hermes_bridge_start_error",
|
||||
format!("Hermes run 启动请求失败: {error}"),
|
||||
)
|
||||
.with_context(context)
|
||||
.with_header("x-mnote-web-owner", "mnote-web")
|
||||
.with_header("x-upstream-service", "hermes-api")
|
||||
.with_header("x-mnote-ai-execution-owner", "hermes")
|
||||
})?;
|
||||
|
||||
if !start_response.status().is_success() {
|
||||
let status = start_response.status();
|
||||
let body = start_response
|
||||
.text()
|
||||
.await
|
||||
.unwrap_or_else(|error| format!("读取 Hermes 错误响应失败: {error}"));
|
||||
return Err(WebError::bad_gateway_code(
|
||||
"hermes_bridge_start_rejected",
|
||||
format!("Hermes run 启动返回 HTTP {}: {}", status.as_u16(), body),
|
||||
)
|
||||
.with_context(context)
|
||||
.with_header("x-mnote-web-owner", "mnote-web")
|
||||
.with_header("x-upstream-service", "hermes-api")
|
||||
.with_header("x-mnote-ai-execution-owner", "hermes"));
|
||||
}
|
||||
|
||||
let started: Value = start_response.json().await.map_err(|error| {
|
||||
WebError::bad_gateway_code(
|
||||
"hermes_bridge_bad_start_response",
|
||||
format!("Hermes run 启动响应不是合法 JSON: {error}"),
|
||||
)
|
||||
.with_context(context)
|
||||
.with_header("x-mnote-web-owner", "mnote-web")
|
||||
.with_header("x-upstream-service", "hermes-api")
|
||||
.with_header("x-mnote-ai-execution-owner", "hermes")
|
||||
})?;
|
||||
let run_id = started
|
||||
.get("run_id")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.ok_or_else(|| {
|
||||
WebError::bad_gateway_code(
|
||||
"hermes_bridge_missing_run_id",
|
||||
"Hermes run 启动响应缺少 run_id。",
|
||||
)
|
||||
.with_context(context)
|
||||
.with_header("x-mnote-web-owner", "mnote-web")
|
||||
.with_header("x-upstream-service", "hermes-api")
|
||||
.with_header("x-mnote-ai-execution-owner", "hermes")
|
||||
})?;
|
||||
|
||||
let mut events_url =
|
||||
reqwest::Url::parse(&format!("{}/v1/runs/", base_url.trim_end_matches('/')))
|
||||
.map_err(|error| WebError::internal(format!("Hermes events URL 非法: {error}")))?;
|
||||
events_url
|
||||
.path_segments_mut()
|
||||
.map_err(|_| WebError::internal("Hermes events URL 不支持路径拼接"))?
|
||||
.pop_if_empty()
|
||||
.push(run_id)
|
||||
.push("events");
|
||||
|
||||
let events_response = client
|
||||
.get(events_url)
|
||||
.bearer_auth(api_key)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|error| {
|
||||
WebError::bad_gateway_code(
|
||||
"hermes_bridge_events_error",
|
||||
format!("Hermes 事件流请求失败: {error}"),
|
||||
)
|
||||
.with_context(context)
|
||||
.with_header("x-mnote-web-owner", "mnote-web")
|
||||
.with_header("x-upstream-service", "hermes-api")
|
||||
.with_header("x-mnote-ai-execution-owner", "hermes")
|
||||
})?;
|
||||
|
||||
if !events_response.status().is_success() {
|
||||
let status = events_response.status();
|
||||
let body = events_response
|
||||
.text()
|
||||
.await
|
||||
.unwrap_or_else(|error| format!("读取 Hermes 事件流错误响应失败: {error}"));
|
||||
return Err(WebError::bad_gateway_code(
|
||||
"hermes_bridge_events_rejected",
|
||||
format!("Hermes 事件流返回 HTTP {}: {}", status.as_u16(), body),
|
||||
)
|
||||
.with_context(context)
|
||||
.with_header("x-mnote-web-owner", "mnote-web")
|
||||
.with_header("x-upstream-service", "hermes-api")
|
||||
.with_header("x-mnote-ai-execution-owner", "hermes"));
|
||||
}
|
||||
|
||||
let events_text = events_response.text().await.map_err(|error| {
|
||||
WebError::bad_gateway_code(
|
||||
"hermes_bridge_events_read_error",
|
||||
format!("读取 Hermes 事件流失败: {error}"),
|
||||
)
|
||||
.with_context(context)
|
||||
.with_header("x-mnote-web-owner", "mnote-web")
|
||||
.with_header("x-upstream-service", "hermes-api")
|
||||
.with_header("x-mnote-ai-execution-owner", "hermes")
|
||||
})?;
|
||||
|
||||
build_hermes_sse_response(context, parse_sse_data_json_values(&events_text))
|
||||
}
|
||||
|
||||
fn resolve_hermes_api_base_url() -> String {
|
||||
read_env_or_dotenv("MNOTE_HERMES_API_BASE_URL")
|
||||
.unwrap_or_else(|| "http://127.0.0.1:8642".into())
|
||||
.trim()
|
||||
.trim_end_matches('/')
|
||||
.to_string()
|
||||
}
|
||||
|
||||
fn build_hermes_run_payload(payload: &Value) -> Value {
|
||||
let messages = payload
|
||||
.get("messages")
|
||||
.and_then(Value::as_array)
|
||||
.cloned()
|
||||
.unwrap_or_default();
|
||||
let input: Vec<Value> = messages
|
||||
.iter()
|
||||
.map(|message| {
|
||||
json!({
|
||||
"role": message.get("role").and_then(Value::as_str).unwrap_or("user"),
|
||||
"content": message.get("content").and_then(Value::as_str).unwrap_or(""),
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
let conversation_history = if input.len() > 1 {
|
||||
input[..input.len() - 1].to_vec()
|
||||
} else {
|
||||
Vec::new()
|
||||
};
|
||||
let session_id = payload
|
||||
.get("options")
|
||||
.and_then(|value| value.get("ai"))
|
||||
.and_then(|value| value.get("sessionId"))
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty());
|
||||
|
||||
json!({
|
||||
"input": input,
|
||||
"conversation_history": conversation_history,
|
||||
"session_id": session_id,
|
||||
})
|
||||
}
|
||||
|
||||
fn parse_sse_data_json_values(text: &str) -> Vec<Value> {
|
||||
text.split("\n\n")
|
||||
.filter_map(|frame| {
|
||||
let data = frame
|
||||
.lines()
|
||||
.filter_map(|line| line.strip_prefix("data:"))
|
||||
.map(str::trim_start)
|
||||
.collect::<Vec<_>>()
|
||||
.join("\n");
|
||||
if data.trim().is_empty() {
|
||||
return None;
|
||||
}
|
||||
serde_json::from_str::<Value>(&data).ok()
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn sse_frame(event: &str, data: Value) -> String {
|
||||
format!("event: {event}\ndata: {}\n\n", data.to_string())
|
||||
}
|
||||
|
||||
fn build_hermes_sse_response(
|
||||
context: &RequestContext,
|
||||
events: Vec<Value>,
|
||||
) -> Result<Response, WebError> {
|
||||
let mut body = String::new();
|
||||
let mut assistant_text = String::new();
|
||||
let mut completed = false;
|
||||
body.push_str(&sse_frame(
|
||||
"ready",
|
||||
json!({"ok": true, "bridgeOwner": "hermes", "requestId": context.trace.request_id}),
|
||||
));
|
||||
for event in events {
|
||||
let event_name = event.get("event").and_then(Value::as_str).unwrap_or("");
|
||||
match event_name {
|
||||
"message.delta" => {
|
||||
if let Some(delta) = event.get("delta").and_then(Value::as_str) {
|
||||
assistant_text.push_str(delta);
|
||||
body.push_str(&sse_frame("assistant_delta", json!({"text": delta})));
|
||||
}
|
||||
}
|
||||
"run.completed" => {
|
||||
completed = true;
|
||||
let output = event
|
||||
.get("output")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.unwrap_or_else(|| assistant_text.trim());
|
||||
let text = if output.is_empty() {
|
||||
"(无输出)"
|
||||
} else {
|
||||
output
|
||||
};
|
||||
body.push_str(&sse_frame("assistant_message", json!({"text": text})));
|
||||
body.push_str(&sse_frame(
|
||||
"completion",
|
||||
json!({"ok": true, "text": text, "steps": 1}),
|
||||
));
|
||||
}
|
||||
"run.failed" => {
|
||||
completed = true;
|
||||
let message = event
|
||||
.get("error")
|
||||
.and_then(Value::as_str)
|
||||
.unwrap_or("Hermes 执行失败");
|
||||
body.push_str(&sse_frame(
|
||||
"error",
|
||||
json!({"ok": false, "message": message}),
|
||||
));
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
if !completed {
|
||||
let text = assistant_text.trim();
|
||||
let text = if text.is_empty() {
|
||||
"(无输出)"
|
||||
} else {
|
||||
text
|
||||
};
|
||||
body.push_str(&sse_frame("assistant_message", json!({"text": text})));
|
||||
body.push_str(&sse_frame(
|
||||
"completion",
|
||||
json!({"ok": true, "text": text, "steps": 1}),
|
||||
));
|
||||
}
|
||||
|
||||
let mut response = Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header(header::CONTENT_TYPE, "text/event-stream; charset=utf-8")
|
||||
.header(header::CACHE_CONTROL, "no-cache, no-transform")
|
||||
.header(header::CONNECTION, "keep-alive")
|
||||
.header("x-accel-buffering", "no")
|
||||
.header("x-mnote-ai-execution-owner", "hermes")
|
||||
.body(Body::from(body))
|
||||
.map_err(|error| WebError::internal(format!("Hermes SSE 响应构造失败: {error}")))?;
|
||||
stamp_owner_header(response.headers_mut());
|
||||
Ok(response)
|
||||
}
|
||||
|
||||
async fn proxy_ai_agent_run_to_next(
|
||||
base_url: &str,
|
||||
context: &RequestContext,
|
||||
@@ -619,9 +932,14 @@ pub async fn next_sidebar(
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use crate::app::{build_app, AppConfig, AppState};
|
||||
use crate::context::RequestContext;
|
||||
use crate::routes::compat::{
|
||||
build_hermes_run_payload, build_hermes_sse_response, parse_sse_data_json_values,
|
||||
};
|
||||
use axum::body::Body;
|
||||
use axum::http::{Request, StatusCode};
|
||||
use axum::http::{HeaderMap, Method, Request, StatusCode, Uri};
|
||||
use axum::response::IntoResponse;
|
||||
use serde_json::json;
|
||||
use tokio::net::TcpListener;
|
||||
use tower::util::ServiceExt;
|
||||
|
||||
@@ -709,6 +1027,108 @@ mod tests {
|
||||
assert!(!text.contains("legacy_next_compat_disabled"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn hermes_payload_uses_page_messages_and_session_id() {
|
||||
let payload = json!({
|
||||
"messages": [
|
||||
{"role": "system", "content": "系统约束"},
|
||||
{"role": "user", "content": "总结当前页面"}
|
||||
],
|
||||
"options": {
|
||||
"ai": {
|
||||
"provider": "hermes",
|
||||
"sessionId": "hermes-session-1"
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
let hermes_payload = build_hermes_run_payload(&payload);
|
||||
|
||||
assert_eq!(hermes_payload["session_id"], "hermes-session-1");
|
||||
assert_eq!(hermes_payload["input"][1]["content"], "总结当前页面");
|
||||
assert_eq!(hermes_payload["conversation_history"][0]["role"], "system");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn hermes_sse_is_translated_to_page_ai_events() {
|
||||
let events = parse_sse_data_json_values(
|
||||
"data: {\"event\":\"message.delta\",\"delta\":\"Hel\"}\n\n\
|
||||
data: {\"event\":\"message.delta\",\"delta\":\"lo\"}\n\n\
|
||||
data: {\"event\":\"run.completed\",\"output\":\"Hello\"}\n\n",
|
||||
);
|
||||
let context = RequestContext::from_http_parts(
|
||||
&Method::POST,
|
||||
&"/api/ai-agent/run".parse::<Uri>().expect("uri"),
|
||||
&HeaderMap::new(),
|
||||
);
|
||||
|
||||
let response = build_hermes_sse_response(&context, events).expect("response");
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
assert_eq!(
|
||||
response
|
||||
.headers()
|
||||
.get("x-mnote-ai-execution-owner")
|
||||
.and_then(|value| value.to_str().ok()),
|
||||
Some("hermes")
|
||||
);
|
||||
let body = axum::body::to_bytes(response.into_body(), usize::MAX)
|
||||
.await
|
||||
.expect("body");
|
||||
let text = String::from_utf8(body.to_vec()).expect("utf8");
|
||||
assert!(text.contains("event: assistant_message"));
|
||||
assert!(text.contains("event: completion"));
|
||||
assert!(text.contains("Hello"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn explicit_agent_provider_does_not_fallback_to_mnote_cli_when_next_unavailable() {
|
||||
let response = build_app(AppState::new(AppConfig {
|
||||
service_name: "mnote-web".into(),
|
||||
service_version: "0.1.0".into(),
|
||||
bind_addr: "127.0.0.1:0".into(),
|
||||
public_bind_addr: "127.0.0.1:3000".into(),
|
||||
legacy_next_base_url: Some("http://127.0.0.1:9".into()),
|
||||
enable_legacy_next_compat: false,
|
||||
enable_debug_shell_routes: false,
|
||||
hermes_base_path: "/api/hermes".into(),
|
||||
compat_next_base_path: "/api/compat/next".into(),
|
||||
convex_url: None,
|
||||
convex_admin_key: None,
|
||||
allow_dev_fixtures: true,
|
||||
query_fixtures_json: None,
|
||||
mutation_fixtures_json: None,
|
||||
dev_user_id: "dev-user".into(),
|
||||
dev_user_name: "开发用户".into(),
|
||||
dev_user_email: "dev@mnote.local".into(),
|
||||
}))
|
||||
.oneshot(
|
||||
Request::builder()
|
||||
.method("POST")
|
||||
.uri("/api/ai-agent/run")
|
||||
.header("content-type", "application/json")
|
||||
.body(Body::from(r#"{"stream":true,"scope":"document","messages":[{"role":"user","content":"ping"}],"context":{"documentId":"doc-1"},"options":{"ai":{"provider":"codex"}}}"#))
|
||||
.expect("request"),
|
||||
)
|
||||
.await
|
||||
.expect("response");
|
||||
|
||||
assert_eq!(response.status(), StatusCode::BAD_GATEWAY);
|
||||
assert_eq!(
|
||||
response
|
||||
.headers()
|
||||
.get("x-mnote-ai-execution-owner")
|
||||
.and_then(|value| value.to_str().ok()),
|
||||
Some("provider-bridge-unavailable")
|
||||
);
|
||||
let body = axum::body::to_bytes(response.into_body(), usize::MAX)
|
||||
.await
|
||||
.expect("body");
|
||||
let text = String::from_utf8(body.to_vec()).expect("utf8");
|
||||
assert!(text.contains("codex"));
|
||||
assert!(!text.contains("mnote-cli"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn ai_agent_run_proxies_to_next_ai_route_when_legacy_next_compat_enabled() {
|
||||
let next_app = axum::Router::new().route(
|
||||
|
||||
Reference in New Issue
Block a user