use crate::app::AppState; use crate::context::RequestContext; use crate::error::WebError; use crate::routes::query_support::resolve_effective_workspace_id; use crate::routes::snapshot_support::{load_projection_snapshot, ProjectionSnapshotSpec}; use axum::body::{Body, Bytes}; use axum::extract::Query; use axum::extract::{Extension, State}; use axum::http::{header, HeaderMap, HeaderName, HeaderValue, Request, StatusCode}; use axum::response::{IntoResponse, Response}; use axum::Json; use core_protocol::KernelProjectionKind; use futures_util::{StreamExt, TryStreamExt}; use serde::Deserialize; use serde_json::{json, Value}; use std::env; use std::fs; use std::path::PathBuf; use std::process::Command; use std::time::Duration; #[derive(Debug, Deserialize)] #[serde(rename_all = "camelCase")] pub struct CompatSidebarQuery { pub workspace_id: Option, } pub async fn next_ai_agent_run( State(state): State, Extension(context): Extension, request: Request, ) -> Result { let (parts, body) = request.into_parts(); let body = axum::body::to_bytes(body, 10 * 1024 * 1024) .await .map_err(|error| WebError::internal(format!("读取 AI 请求体失败: {error}")))?; let payload: Value = serde_json::from_slice(&body).map_err(|error| { WebError::bad_request_code( "ai_agent_bad_request", format!("AI 请求体不是合法 JSON: {error}"), ) .with_context(&context) .with_header("x-mnote-web-owner", "mnote-web") })?; if state.config().enable_legacy_next_compat { if let Some(base_url) = state.config().legacy_next_base_url.as_deref() { if let Ok(response) = proxy_ai_agent_run_to_next(base_url, &context, &parts.headers, body.clone()).await { return Ok(response); } } } 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); } let Some(backend_url) = resolve_ai_orchestrator_backend_url() else { return Err(WebError::bad_gateway_code( "ai_orchestrator_unavailable", "未配置 BACKEND_URL,无法连接 document AI orchestrator。", ) .with_context(&context) .with_header("x-mnote-web-owner", "mnote-web")); }; let upstream_url = reqwest::Url::parse(&format!( "{}/api/v1/ai-agent/document/run", backend_url.trim_end_matches('/') )) .map_err(|error| WebError::internal(format!("AI orchestrator URL 非法: {error}")))?; let forward_payload = build_document_ai_orchestrator_payload(&state, &context, &payload); let client = reqwest::Client::builder() .connect_timeout(Duration::from_secs(5)) .redirect(reqwest::redirect::Policy::none()) .build() .map_err(|error| { WebError::internal(format!("AI orchestrator HTTP 客户端创建失败: {error}")) })?; let mut upstream_request = client.post(upstream_url).json(&forward_payload); upstream_request = apply_ai_forward_headers(upstream_request, &parts.headers, &context); if let Some(api_key) = read_env_or_dotenv("MNOTE_AI_ORCHESTRATOR_API_KEY") { upstream_request = upstream_request.header("x-mnote-ai-key", api_key); } let upstream_response = upstream_request.send().await.map_err(|error| { WebError::bad_gateway_code( "ai_orchestrator_proxy_error", format!("document AI orchestrator 请求失败: {error}"), ) .with_context(&context) .with_header("x-mnote-web-owner", "mnote-web") .with_header("x-upstream-service", "document-ai-orchestrator") })?; if !upstream_response.status().is_success() { let status = upstream_response.status(); let body = upstream_response .text() .await .unwrap_or_else(|error| format!("读取 upstream 错误响应失败: {error}")); return Err(WebError::bad_gateway_code( "ai_orchestrator_upstream_error", format!( "document AI orchestrator 返回 HTTP {}: {}", status.as_u16(), body ), ) .with_context(&context) .with_header("x-mnote-web-owner", "mnote-web") .with_header("x-upstream-service", "document-ai-orchestrator")); } 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 { 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 = 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 { text.split("\n\n") .filter_map(|frame| { let data = frame .lines() .filter_map(|line| line.strip_prefix("data:")) .map(str::trim_start) .collect::>() .join("\n"); if data.trim().is_empty() { return None; } serde_json::from_str::(&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, ) -> Result { 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, headers: &HeaderMap, body_bytes: Bytes, ) -> Result { let upstream_url = reqwest::Url::parse(&format!( "{}/api/ai-agent/run", base_url.trim_end_matches('/') )) .map_err(|error| WebError::internal(format!("Next AI route URL 非法: {error}")))?; let client = reqwest::Client::builder() .connect_timeout(Duration::from_secs(5)) .redirect(reqwest::redirect::Policy::none()) .build() .map_err(|error| WebError::internal(format!("Next AI HTTP 客户端创建失败: {error}")))?; let mut upstream_request = client.post(upstream_url).body(body_bytes.clone()); for (name, value) in headers.iter() { if is_hop_by_hop_header(name.as_str()) || name == header::HOST || name == header::CONTENT_LENGTH || name == header::COOKIE || name == header::AUTHORIZATION { continue; } upstream_request = upstream_request.header(name.as_str(), value.as_bytes()); } if let Some(cookie) = context.auth.cookie_header.as_deref() { upstream_request = upstream_request.header(header::COOKIE, cookie); } if let Some(authorization) = context.auth.authorization.as_deref() { upstream_request = upstream_request.header(header::AUTHORIZATION, authorization); } upstream_request = upstream_request.header("x-request-id", context.trace.request_id.as_str()); upstream_request = upstream_request.header("x-trace-id", context.trace.trace_id.as_str()); upstream_request = upstream_request.header("x-mnote-source-channel", "rust_web_route"); upstream_request = upstream_request.header("x-mnote-source-client", "mnote-web"); upstream_request = upstream_request.header("x-mnote-actor-id", context.auth.actor_id.as_str()); upstream_request = upstream_request.header("x-mnote-actor-type", context.auth.actor_type.as_str()); if let Some(workspace_id) = context.workspace.workspace_id.as_deref() { upstream_request = upstream_request.header("x-mnote-workspace-id", workspace_id); } if let Some(session_id) = context.auth.session_id.as_deref() { upstream_request = upstream_request.header("x-mnote-session-id", session_id); } let upstream_response = upstream_request.send().await.map_err(|error| { WebError::bad_gateway_code( "next_ai_proxy_error", format!("Next /api/ai-agent/run 请求失败: {error}"), ) .with_context(context) .with_header("x-mnote-web-owner", "mnote-web") .with_header("x-upstream-service", "next-ai-route") })?; if !upstream_response.status().is_success() { let status = upstream_response.status(); let body = upstream_response .text() .await .unwrap_or_else(|error| format!("读取 Next AI 错误响应失败: {error}")); return Err(WebError::bad_gateway_code( "next_ai_upstream_error", format!( "Next /api/ai-agent/run 返回 HTTP {}: {}", status.as_u16(), body ), ) .with_context(context) .with_header("x-mnote-web-owner", "mnote-web") .with_header("x-upstream-service", "next-ai-route")); } build_ai_sse_proxy_response(upstream_response, context) } fn resolve_repo_root() -> PathBuf { PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../../..") } fn build_mnote_cli_args(context: &RequestContext, payload: &Value) -> Vec { let ai = payload .get("options") .and_then(|value| value.get("ai")) .cloned() .unwrap_or(Value::Null); let runtime_context = payload.get("context").cloned().unwrap_or(Value::Null); let document_id = runtime_context .get("documentId") .and_then(Value::as_str) .unwrap_or("current"); let workspace_id = runtime_context .get("workspaceId") .and_then(Value::as_str) .or(context.workspace.workspace_id.as_deref()); let session_id = ai .get("sessionId") .and_then(Value::as_str) .filter(|value| !value.trim().is_empty()) .map(str::to_string) .unwrap_or_else(|| format!("ai-{}", context.trace.request_id)); let args_json = json!({ "pageId": document_id, "documentId": document_id, "workspaceId": workspace_id, "provider": ai.get("provider").cloned().unwrap_or(Value::Null), "modelKey": ai.get("modelKey").cloned().unwrap_or(Value::Null), "profileId": ai.get("profileId").cloned().unwrap_or(Value::Null), "selectedUids": runtime_context.get("selectedUids").cloned().unwrap_or(Value::Null), "pageOptions": runtime_context.get("pageOptions").cloned().unwrap_or(Value::Null), }) .to_string(); vec![ "run".into(), "--quiet".into(), "--manifest-path".into(), resolve_repo_root() .join("rust") .join("Cargo.toml") .to_string_lossy() .to_string(), "-p".into(), "mnote-cli".into(), "--".into(), "--json".into(), "--validate-only".into(), "--dry-run".into(), "--actor-id".into(), context.auth.actor_id.clone(), "--actor-type".into(), context.auth.actor_type.clone(), "--session-id".into(), session_id, "--reason".into(), "ai-agent-run:mnote-web-rust-host".into(), "tool".into(), "run".into(), "--tool-name".into(), "doc_get".into(), "--kind".into(), "query".into(), "--mode".into(), "explain-plan".into(), "--args-json".into(), args_json, ] } async fn run_local_mnote_cli_ai_host( state: &AppState, context: &RequestContext, payload: &Value, ) -> Result { let stream = payload .get("stream") .and_then(Value::as_bool) .unwrap_or(true); let args = build_mnote_cli_args(context, payload); let repo_root = resolve_repo_root(); let actor_id = if context.auth.actor_id.trim().is_empty() || context.auth.actor_id == "anonymous" { state.config().dev_user_id.clone() } else { context.auth.actor_id.clone() }; let actor_type = if context.auth.actor_type.trim().is_empty() || context.auth.actor_type == "anonymous" { "user".to_string() } else { context.auth.actor_type.clone() }; let dev_email = state.config().dev_user_email.clone(); let dev_name = state.config().dev_user_name.clone(); let output = tokio::task::spawn_blocking(move || { Command::new("cargo") .args(args) .current_dir(repo_root) .env("CARGO_TERM_COLOR", "never") .env( "RUSTUP_TOOLCHAIN", env::var("RUSTUP_TOOLCHAIN") .ok() .filter(|value| !value.trim().is_empty()) .unwrap_or_else(|| "1.89.0".into()), ) .env("DEV_USER_ID", actor_id) .env("DEV_USER_EMAIL", dev_email) .env("DEV_USER_NAME", dev_name) .env("MNOTE_CLI_ALLOW_CREATE_PAGE", "1") .env("MNOTE_CLI_ALLOW_EDIT", "1") .env("MNOTE_ACTOR_TYPE", actor_type) .output() }) .await .map_err(|error| { WebError::bad_gateway_code( "mnote_cli_host_join_error", format!("mnote-cli host join 失败: {error}"), ) .with_context(context) .with_header("x-mnote-web-owner", "mnote-web") })? .map_err(|error| { WebError::bad_gateway_code( "mnote_cli_host_spawn_error", format!("mnote-cli host 启动失败: {error}"), ) .with_context(context) .with_header("x-mnote-web-owner", "mnote-web") })?; let stdout = String::from_utf8_lossy(&output.stdout).trim().to_string(); let stderr = String::from_utf8_lossy(&output.stderr).trim().to_string(); if !stream { if !output.status.success() { return Err(WebError::bad_gateway_code( "mnote_cli_host_failed", if stderr.is_empty() { "mnote-cli 执行失败".into() } else { stderr }, ) .with_context(context) .with_header("x-mnote-web-owner", "mnote-web")); } let mut response = Json(json!({ "ok": true, "bridgeOwner": "mnote-cli", "text": if stdout.is_empty() { "mnote-cli 无输出" } else { &stdout }, })) .into_response(); stamp_owner_header(response.headers_mut()); response.headers_mut().insert( HeaderName::from_static("x-mnote-ai-execution-owner"), HeaderValue::from_static("mnote-cli"), ); return Ok(response); } let body = if output.status.success() { format!( "event: ready\ndata: {}\n\nevent: assistant_message\ndata: {}\n\nevent: completion\ndata: {}\n\n", json!({"ok": true, "bridgeOwner": "mnote-cli"}).to_string(), json!({"text": if stdout.is_empty() { "mnote-cli 无输出" } else { &stdout }}).to_string(), json!({"ok": true, "text": if stdout.is_empty() { "mnote-cli 无输出" } else { &stdout }, "steps": 1}).to_string(), ) } else { format!( "event: ready\ndata: {}\n\nevent: error\ndata: {}\n\n", json!({"ok": true, "bridgeOwner": "mnote-cli"}).to_string(), json!({"ok": false, "message": if stderr.is_empty() { "mnote-cli 执行失败" } else { &stderr }}).to_string(), ) }; 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", "mnote-cli") .body(Body::from(body)) .map_err(|error| WebError::internal(format!("mnote-cli SSE 响应构造失败: {error}")))?; stamp_owner_header(response.headers_mut()); Ok(response) } fn build_document_ai_orchestrator_payload( state: &AppState, context: &RequestContext, payload: &Value, ) -> Value { let ai_options = payload.get("options").and_then(|value| value.get("ai")); let user_id = if context.auth.actor_id.trim().is_empty() || context.auth.actor_id == "anonymous" { state.config().dev_user_id.clone() } else { context.auth.actor_id.clone() }; json!({ "userId": user_id, "sessionId": ai_options .and_then(|value| value.get("sessionId")) .and_then(Value::as_str) .map(str::trim) .filter(|value| !value.is_empty()), "model": ai_options .and_then(|value| value.get("model")) .and_then(Value::as_str) .map(str::trim) .filter(|value| !value.is_empty()), "modelKey": ai_options .and_then(|value| value.get("modelKey")) .and_then(Value::as_str) .map(str::trim) .filter(|value| !value.is_empty()), "profileId": ai_options .and_then(|value| value.get("profileId")) .and_then(Value::as_str) .map(str::trim) .filter(|value| !value.is_empty()), "maxSteps": payload.get("maxSteps").filter(|value| !value.is_null()).cloned().unwrap_or_else(|| json!(10)), "messages": payload.get("messages").cloned().unwrap_or_else(|| json!([])), "context": build_document_ai_context(payload.get("context")), }) } fn build_document_ai_context(context: Option<&Value>) -> Value { let get = |key: &str| { context .and_then(|value| value.get(key)) .cloned() .unwrap_or(Value::Null) }; json!({ "source": get("source"), "action": get("action"), "documentId": get("documentId"), "workspaceId": get("workspaceId"), "selectedBlockId": get("selectedBlockId"), "selectedBlockIndex": get("selectedBlockIndex"), "selectedUids": get("selectedUids"), "selectedText": get("selectedText"), "selection": get("selection"), "tiptapDocument": get("tiptapDocument"), "documentBlocks": get("documentBlocks"), "node": get("node"), "subtree": get("subtree"), "outline": get("outline"), "evidence": get("evidence"), "pageOptions": get("pageOptions"), }) } fn apply_ai_forward_headers( mut request: reqwest::RequestBuilder, headers: &HeaderMap, context: &RequestContext, ) -> reqwest::RequestBuilder { for (name, value) in headers.iter() { if is_hop_by_hop_header(name.as_str()) || name == header::HOST || name == header::CONTENT_LENGTH || name == header::CONTENT_TYPE { continue; } request = request.header(name.as_str(), value.as_bytes()); } request = request.header(header::CONTENT_TYPE.as_str(), "application/json"); request = request.header("x-request-id", context.trace.request_id.as_str()); request = request.header("x-trace-id", context.trace.trace_id.as_str()); request = request.header("x-mnote-source-channel", "rust_web_route"); request = request.header("x-mnote-source-client", "mnote-web"); request = request.header("x-mnote-actor-id", context.auth.actor_id.as_str()); request.header("x-mnote-actor-type", context.auth.actor_type.as_str()) } fn build_ai_sse_proxy_response( upstream_response: reqwest::Response, context: &RequestContext, ) -> Result { let status = StatusCode::from_u16(upstream_response.status().as_u16()).unwrap_or(StatusCode::OK); let ready = Bytes::from(format!( "event: ready\ndata: {}\n\n", json!({"ok": true, "requestId": context.trace.request_id}).to_string() )); let upstream_stream = upstream_response .bytes_stream() .map_err(std::io::Error::other); let body_stream = futures_util::stream::once(async move { Ok::(ready) }) .chain(upstream_stream); let mut response = Response::builder() .status(status) .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") .body(Body::from_stream(body_stream)) .map_err(|error| WebError::internal(format!("AI SSE 响应构造失败: {error}")))?; stamp_owner_header(response.headers_mut()); Ok(response) } fn stamp_owner_header(headers: &mut HeaderMap) { if let Ok(name) = HeaderName::from_lowercase(b"x-mnote-web-owner") { headers.insert(name, HeaderValue::from_static("mnote-web")); } } fn resolve_ai_orchestrator_backend_url() -> Option { read_env_or_dotenv("BACKEND_INTERNAL_URL") .or_else(|| read_env_or_dotenv("BACKEND_URL")) .map(|value| value.trim().trim_end_matches('/').to_string()) .filter(|value| !value.is_empty()) } fn read_env_or_dotenv(key: &str) -> Option { if let Ok(value) = env::var(key) { let trimmed = value.trim().trim_matches('"').to_string(); if !trimmed.is_empty() { return Some(trimmed); } } let root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")) .join("../../..") .join(".env.all"); let content = fs::read_to_string(root).ok()?; for line in content.lines() { let line = line.trim_end_matches('\r'); if line.starts_with('#') || line.trim().is_empty() { continue; } let Some((k, v)) = line.split_once('=') else { continue; }; if k.trim() != key { continue; } let trimmed = v.trim().trim_matches('"').to_string(); if !trimmed.is_empty() { return Some(trimmed); } } None } fn is_hop_by_hop_header(name: &str) -> bool { matches!( name.to_ascii_lowercase().as_str(), "connection" | "keep-alive" | "proxy-authenticate" | "proxy-authorization" | "te" | "trailers" | "transfer-encoding" | "upgrade" ) } pub async fn next_sidebar( State(state): State, Extension(context): Extension, Query(query): Query, ) -> Result<(StatusCode, Json), WebError> { let effective_workspace_id = resolve_effective_workspace_id(&context, query.workspace_id.as_deref(), true)? .expect("workspace_required 已确保存在"); let snapshot = load_projection_snapshot( state.config(), &context, &ProjectionSnapshotSpec { workspace_id: &effective_workspace_id, root_node_id: None, depth: None, query: None, max_results: None, projection: KernelProjectionKind::SidebarTree, }, ) .await?; let mut dataset_object = snapshot.dataset.as_object().cloned().ok_or_else(|| { WebError::bad_gateway_code("convex_bad_response", "sidebar.dataset.list 返回值不是对象") .with_context(&context) .with_header("x-error-phase", "compat_sidebar_shape") .with_header("x-upstream-service", "convex") })?; dataset_object.insert("kernel_sidebar_projection".into(), snapshot.projection); Ok(( StatusCode::OK, Json(json!({ "ok": true, "boundary": "next_sidebar_compat", "requestId": context.trace.request_id, "traceId": context.trace.trace_id, "workspaceId": effective_workspace_id, "result": Value::Object(dataset_object), })), )) } #[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::{HeaderMap, Method, Request, StatusCode, Uri}; use axum::response::IntoResponse; use serde_json::json; use tokio::net::TcpListener; use tower::util::ServiceExt; fn app() -> axum::Router { 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:3100".into()), enable_legacy_next_compat: true, 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: Some(r#"{"sidebar:datasetList":{"active_workspace_id":"ws_demo","workspaces":[],"documents":[{"id":"page_root","workspace_id":"ws_demo","title":"工作区首页","parent_id":null,"sort_order":0,"is_starred":true,"is_template":false,"created_at":"2026-04-16T00:00:00Z","updated_at":"2026-04-16T00:00:00Z"},{"id":"page_child","workspace_id":"ws_demo","title":"子页面","parent_id":"page_root","sort_order":1,"is_starred":false,"is_template":false,"created_at":"2026-04-16T00:00:00Z","updated_at":"2026-04-16T00:00:00Z"}],"trashed_documents":[],"media_assets":[],"trashed_media_assets":[],"mindmap_assets":[],"trashed_mindmap_assets":[],"table_assets":[],"trashed_table_assets":[],"mindmap_docs":[],"mindmap_asset_children":{}},"bridgeLogs:listWorkspaceOverview":{"workspace_id":"ws_demo","command_logs":[],"domain_events":[],"next_cursor":null,"has_more":false,"filters":{"command_status":null,"event_status":null,"target_page_id":null,"target_block_id":null,"aggregate_type":null,"aggregate_id":null},"generated_at":"2026-04-16T00:00:00Z"}}"#.into()), mutation_fixtures_json: None, dev_user_id: "dev-user".into(), dev_user_name: "开发用户".into(), dev_user_email: "dev@mnote.local".into(), })) } #[tokio::test] async fn compat_sidebar_route_returns_dataset_and_projection() { let response = app() .oneshot( Request::builder() .uri("/api/compat/next/sidebar?workspaceId=ws_demo") .body(Body::empty()) .expect("request"), ) .await .expect("response"); assert_eq!(response.status(), StatusCode::OK); } #[tokio::test] async fn direct_ai_agent_run_is_owned_by_rust_web_when_next_compat_disabled() { 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:3100".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":"online"}}}"#)) .expect("request"), ) .await .expect("response"); assert_ne!(response.status(), StatusCode::SERVICE_UNAVAILABLE); assert_eq!( response .headers() .get("x-mnote-web-owner") .and_then(|value| value.to_str().ok()), Some("mnote-web") ); 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("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::().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( "/api/ai-agent/run", axum::routing::post(|| async move { ( [( axum::http::header::CONTENT_TYPE, "text/event-stream; charset=utf-8", )], "event: assistant_message\ndata: {\"text\":\"hello from next ai\"}\n\n", ) .into_response() }), ); let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind next"); let addr = listener.local_addr().expect("local addr"); let server = tokio::spawn(async move { axum::serve(listener, next_app).await.expect("serve next"); }); 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(format!("http://{}", addr)), enable_legacy_next_compat: true, 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":"hermes"}}}"#)) .expect("request"), ) .await .expect("response"); assert_eq!(response.status(), StatusCode::OK); assert_eq!( response .headers() .get("x-mnote-web-owner") .and_then(|value| value.to_str().ok()), Some("mnote-web") ); 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("assistant_message")); assert!(text.contains("hello from next ai")); server.abort(); } }