合并 origin/master 到 AGC 渲染层下沉分支并完成三处对齐

- 活动回合事实源统一到 Direct 线程管理器:删除 direct_runtime 的第二份快照,接单、进度内容变化与收口各广播一次活动回合变更事件
- 平台维护态判定移入 Rust 并在渲染层只订阅单一事件:新增 platform_maintenance 模块与各平台 facade 错误分支的分类入口
- 封面生成请求补 generationInputs.source,保持队列回填后仍能拿到平台素材 ID
- 渲染层按 master 5398a53e6 退役诊断详情入口:删除 agentRuntimeErrorDetail 与「查看详情」交互及其专属用例
- 冲突收口:nginx SPA 白名单、.gitignore、mobile 检查脚本、capabilities 描述取上游,两个已退役计划随上游删除,文档保留双方条目
- 新增并回写本里程碑取证、decision-log 与 pitfalls 的 2026-09-28 记录
This commit is contained in:
kdletters
2026-09-28 15:07:30 +08:00
403 changed files with 22103 additions and 16372 deletions
+50 -1
View File
@@ -2215,7 +2215,16 @@ fn admin_permission_requirement(_method: &Method, path: &str) -> AdminPermission
"/admin/api/editor-assets" => AnyTab(&["editor-assets"]),
"/admin/api/assets/read-url" => AnyTab(&["editor-assets", "editor-showcase"]),
path if path.starts_with("/admin/api/editor-showcase/") => AnyTab(&["editor-showcase"]),
path if path.starts_with("/admin/api/game-distribution/") => AnyTab(&["editor-showcase"]),
path if path.starts_with("/admin/api/game-distribution/reviews") => {
AnyTab(&["editor-showcase"])
}
path if path.starts_with("/admin/api/game-distribution/versions/") => {
AnyTab(&["editor-showcase"])
}
// 游戏管理页与审核页共享 games/* 面(列表、恢复、安全下架)。
path if path.starts_with("/admin/api/game-distribution/games") => {
AnyTab(&["editor-showcase", "game-management"])
}
"/admin/api/profile/redeem-codes" | "/admin/api/profile/redeem-codes/disable" => {
AnyTab(&["redeem"])
}
@@ -7435,6 +7444,11 @@ mod tests {
Method::GET,
"/admin/api/editor-showcase/assets",
),
(
"game-management",
Method::GET,
"/admin/api/game-distribution/games",
),
("editor-assets", Method::GET, "/admin/api/editor-assets"),
];
@@ -7453,6 +7467,41 @@ mod tests {
}
}
#[test]
fn game_management_tab_is_separate_from_game_review_queue() {
// 游戏管理页可以读全量游戏与恢复;待审队列仍只属于游戏审核页。
assert!(
enforce_admin_request_permission(
"member",
&["game-management".to_string()],
&[],
&Method::GET,
"/admin/api/game-distribution/games",
)
.is_ok()
);
assert!(
enforce_admin_request_permission(
"member",
&["game-management".to_string()],
&[],
&Method::POST,
"/admin/api/game-distribution/games/game_1/restore",
)
.is_ok()
);
assert!(
enforce_admin_request_permission(
"member",
&["game-management".to_string()],
&[],
&Method::GET,
"/admin/api/game-distribution/reviews",
)
.is_err()
);
}
#[test]
fn wallet_consumption_reconcile_requires_its_standalone_action_permission() {
assert!(
+363 -14
View File
@@ -7,26 +7,208 @@ use axum::{
extract::{Extension, State},
http::StatusCode,
};
use module_runtime::AgcModelCatalog;
use module_runtime::{
AGC_MODEL_CATALOG_CONFLICT, AGC_MODEL_CATALOG_NOT_INITIALIZED, AgcModelCatalog,
};
use shared_contracts::admin::{AdminAgcModel, AdminAgcModelCatalog};
use spacetime_client::SpacetimeClientError;
use std::time::Duration;
use tracing::warn;
/// 目录未初始化时对外统一的失败文案:目录只能来自上游同步或后台保存。
pub(crate) const AGC_MODEL_CATALOG_NOT_INITIALIZED_MESSAGE: &str =
"模型目录未初始化,服务端正在尝试从上游同步,请稍后重试";
/// 上游模型列表请求超时与响应大小上限;越界按同步失败处理。
///
/// 同步发生在启动期、且在开始对外服务之前,超时必须足够短:上游挂起时不能让
/// 每个 API/All 实例都延迟三十秒才可用。单次失败只记录 error,下次启动会重试。
const AGC_MODEL_LIST_REQUEST_TIMEOUT: Duration = Duration::from_secs(10);
const AGC_MODEL_LIST_MAX_BYTES: usize = 1024 * 1024;
/// 只读 `revision`:存量目录内容不合法时,覆盖写入仍需对齐乐观锁版本。
#[derive(serde::Deserialize)]
struct StoredCatalogRevision {
revision: u64,
}
pub(crate) async fn load_catalog(state: &AppState) -> Result<AgcModelCatalog, AppError> {
let json = state
.spacetime_client()
.read_agc_model_catalog()
.await
.map_err(|_| {
AppError::from_status(StatusCode::SERVICE_UNAVAILABLE).with_message("模型目录暂不可用")
})?;
let catalog: AgcModelCatalog = serde_json::from_str(&json).map_err(|_| {
AppError::from_status(StatusCode::SERVICE_UNAVAILABLE).with_message("模型目录格式无效")
})?;
catalog.validate().map_err(|message| {
AppError::from_status(StatusCode::SERVICE_UNAVAILABLE).with_message(message)
let stored = read_stored_catalog(state).await.map_err(|_| {
AppError::from_status(StatusCode::SERVICE_UNAVAILABLE).with_message("模型目录暂不可用")
})?;
let Some(json) = stored else {
return Err(uninitialized_error());
};
parse_catalog(&json).map_err(|_| uninitialized_error())
}
fn uninitialized_error() -> AppError {
AppError::from_status(StatusCode::SERVICE_UNAVAILABLE)
.with_message(AGC_MODEL_CATALOG_NOT_INITIALIZED_MESSAGE)
}
/// 解析并校验目录内容;解析或校验失败都按“未初始化”处理,由启动期重新同步。
fn parse_catalog(json: &str) -> Result<AgcModelCatalog, String> {
let catalog: AgcModelCatalog =
serde_json::from_str(json).map_err(|_| "模型目录格式无效".to_string())?;
catalog.validate()?;
Ok(catalog)
}
async fn read_stored_catalog(state: &AppState) -> Result<Option<String>, SpacetimeClientError> {
match state.spacetime_client().read_agc_model_catalog().await {
Ok(json) => Ok(Some(json)),
Err(SpacetimeClientError::Procedure(message))
if message == AGC_MODEL_CATALOG_NOT_INITIALIZED =>
{
Ok(None)
}
Err(error) => Err(error),
}
}
/// 启动期确保目录已初始化:未初始化时从上游模型列表重建,失败只返回错误由调用方记录。
///
/// 目录已可用时不做任何写入;只有缺失、结构与当前定义不符或校验不通过才重建,因此上游模型
/// 变化不会自动覆盖后台维护过的目录。
pub(crate) async fn ensure_agc_model_catalog_initialized(state: &AppState) -> Result<(), String> {
let stored = read_stored_catalog(state)
.await
.map_err(|error| format!("读取 AGC 模型目录失败:{error}"))?;
let revision = match stored.as_deref() {
Some(json) => match parse_catalog(json) {
Ok(_) => return Ok(()),
Err(message) => {
warn!(
error = %message,
"AGC 模型目录内容与当前定义不符,按未初始化处理并从上游重建"
);
stored_catalog_revision(json)
}
},
None => Some(0),
};
let revision = revision.ok_or_else(|| {
"存量 AGC 模型目录缺少可解析的 revision,需要先清理该行再重启".to_string()
})?;
let models = fetch_upstream_model_names(state).await?;
let catalog = AgcModelCatalog::from_upstream_models(models, revision)?;
let payload =
serde_json::to_string(&catalog).map_err(|_| "AGC 模型目录序列化失败".to_string())?;
match state
.spacetime_client()
.save_agc_model_catalog(payload)
.await
{
Ok(saved) => {
let saved: AgcModelCatalog = serde_json::from_str(&saved)
.map_err(|_| "AGC 模型目录写回结果格式无效".to_string())?;
tracing::info!(
revision = saved.revision,
model_count = saved.models.len(),
"已按上游模型列表初始化 AGC 模型目录"
);
Ok(())
}
// 多实例同时启动时只有一个写入成功:接受既有目录,但仍要确认它可用,
// 否则会静默地把「每次启动都冲突、目录一直不可用」变成没有任何线索的黑洞。
Err(SpacetimeClientError::Procedure(message)) if message == AGC_MODEL_CATALOG_CONFLICT => {
let stored = read_stored_catalog(state)
.await
.map_err(|error| format!("写入冲突后重读 AGC 模型目录失败:{error}"))?;
if stored
.as_deref()
.map(parse_catalog)
.is_some_and(|result| result.is_ok())
{
warn!("AGC 模型目录写入冲突:已接受其它实例写入的目录");
Ok(())
} else {
Err("AGC 模型目录写入冲突后仍不可用,需要人工检查该行内容与 revision".to_string())
}
}
Err(error) => Err(format!("写入 AGC 模型目录失败:{error}")),
}
}
fn stored_catalog_revision(json: &str) -> Option<u64> {
serde_json::from_str::<StoredCatalogRevision>(json)
.ok()
.map(|stored| stored.revision)
}
/// 上游在售模型列表:`GET {Router 控制面}/api/pricing?group=taonier` 的 `data[].model_name`。
///
/// 取“该分组可见的在售模型”,而不是管理面模型注册表:注册表里会残留已下线、没有路由绑定的
/// 条目(例如已从上游移除的 `gpt-6-astra`/`gpt-6-luna`),而定价列表就是 AGC 账号实际能调用的集合。
/// 该端点是公开只读接口,不需要管理凭据。
async fn fetch_upstream_model_names(state: &AppState) -> Result<Vec<String>, String> {
crate::external_api_keys::ensure_llm_router_url_allowed(state)?;
let origin =
crate::external_api_keys::router_control_origin(&state.config.llm_router_base_url)?;
let url = format!(
"{origin}/api/pricing?group={}",
crate::external_api_keys::LLM_ROUTER_TOKEN_GROUP
);
let client = reqwest::Client::builder()
.timeout(AGC_MODEL_LIST_REQUEST_TIMEOUT)
.redirect(reqwest::redirect::Policy::none())
.build()
.map_err(|error| format!("构建 LLM Router 客户端失败:{error}"))?;
let response = client
.get(url)
.send()
.await
.map_err(|error| format!("请求上游模型列表失败:{error}"))?;
let status = response.status();
if !status.is_success() {
return Err(format!("上游模型列表返回 HTTP {status}"));
}
let bytes = read_bounded_json_body(response).await?;
let payload: serde_json::Value =
serde_json::from_slice(&bytes).map_err(|_| "上游模型列表格式无效".to_string())?;
parse_upstream_model_names(&payload)
}
async fn read_bounded_json_body(mut response: reqwest::Response) -> Result<Vec<u8>, String> {
// 先按 Content-Length 快速拒绝,再流式累加做兜底:不信任上游声明的长度,
// 逐块累计超阈值立即中断,避免 `bytes()` 一次性分配任意大小响应撑爆内存。
if response
.content_length()
.is_some_and(|length| length > AGC_MODEL_LIST_MAX_BYTES as u64)
{
return Err("上游模型列表响应超过大小上限".to_string());
}
let mut bytes = Vec::new();
while let Some(chunk) = response
.chunk()
.await
.map_err(|error| format!("读取上游模型列表失败:{error}"))?
{
if bytes.len().saturating_add(chunk.len()) > AGC_MODEL_LIST_MAX_BYTES {
return Err("上游模型列表响应超过大小上限".to_string());
}
bytes.extend_from_slice(chunk.as_ref());
}
Ok(bytes)
}
fn parse_upstream_model_names(payload: &serde_json::Value) -> Result<Vec<String>, String> {
let data = payload
.get("data")
.and_then(serde_json::Value::as_array)
.ok_or_else(|| "上游模型列表缺少 data 数组".to_string())?;
let models = data
.iter()
.filter_map(|entry| entry.get("model_name").and_then(serde_json::Value::as_str))
.map(str::to_string)
.collect::<Vec<_>>();
if models.iter().all(|model| model.trim().is_empty()) {
return Err("上游模型列表为空".to_string());
}
Ok(models)
}
pub async fn admin_get_agc_models(
State(state): State<AppState>,
Extension(context): Extension<RequestContext>,
@@ -44,6 +226,8 @@ pub async fn admin_save_agc_models(
Extension(_admin): Extension<AuthenticatedAdmin>,
Json(payload): Json<AdminAgcModelCatalog>,
) -> Result<Json<serde_json::Value>, AppError> {
// 目录只来自上游同步:未初始化时后台写入同样失败关闭,避免出现第二条绕过同步的写入口。
load_catalog(&state).await?;
let catalog = AgcModelCatalog {
revision: payload.revision,
default_model_id: payload.default_model_id,
@@ -68,7 +252,7 @@ pub async fn admin_save_agc_models(
.save_agc_model_catalog(payload)
.await
.map_err(|error| {
if matches!(error, spacetime_client::SpacetimeClientError::Procedure(ref message) if message == module_runtime::AGC_MODEL_CATALOG_CONFLICT) {
if matches!(&error, SpacetimeClientError::Procedure(message) if message == AGC_MODEL_CATALOG_CONFLICT) {
AppError::from_status(StatusCode::CONFLICT).with_message("模型目录已被更新,请重新读取")
} else {
AppError::from_status(StatusCode::SERVICE_UNAVAILABLE).with_message("保存模型目录失败,请稍后重试")
@@ -95,3 +279,168 @@ fn catalog_dto(catalog: AgcModelCatalog) -> AdminAgcModelCatalog {
.collect(),
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn upstream_model_names_come_from_pricing_data_array() {
let payload = json!({
"auto_groups": ["default"],
"data": [
{"model_name": "glm-5.3", "model_ratio": 1.0},
{"model_name": "deepseek-flash", "model_ratio": 0.075},
{"model_ratio": 1.0}
]
});
assert_eq!(
parse_upstream_model_names(&payload).unwrap(),
vec!["glm-5.3".to_string(), "deepseek-flash".to_string()]
);
assert_eq!(
parse_upstream_model_names(&json!({"data": []})).unwrap_err(),
"上游模型列表为空"
);
assert_eq!(
parse_upstream_model_names(&json!({"data": [{"model_name": " "}]})).unwrap_err(),
"上游模型列表为空"
);
assert!(parse_upstream_model_names(&json!({"object": "list"})).is_err());
}
#[test]
fn stored_catalog_revision_reads_row_revision() {
assert_eq!(
stored_catalog_revision(
r#"{"revision":4,"defaultModelId":"quality","models":[{"id":"quality","alias":"高质量","modelId":"gpt-6-astra","enabled":true}]}"#
),
Some(4)
);
assert_eq!(stored_catalog_revision("not json"), None);
assert_eq!(stored_catalog_revision(r#"{"models":[]}"#), None);
}
#[test]
fn catalog_parsing_marks_unusable_content_as_uninitialized() {
// 后台保存过的目录结构必须能直接解析。
let catalog = AgcModelCatalog::from_upstream_models(
vec!["deepseek-v4-pro".to_string(), "glm-5.3".to_string()],
4,
)
.unwrap();
assert_eq!(
parse_catalog(&serde_json::to_string(&catalog).unwrap()).unwrap(),
catalog
);
// 结构或内容不合法(例如被外部工具改过)都按未初始化处理,由启动期重新同步。
assert!(
parse_catalog(r#"{"revision":1,"defaultModel":"deepseek-v4-pro","models":[]}"#)
.is_err()
);
assert!(parse_catalog(
r#"{"revision":1,"defaultModelId":"quality","models":[{"id":"quality","alias":"高质量","modelId":"gpt-6-astra","enabled":false}]}"#
)
.is_err());
}
struct MockModelListServer {
base_url: String,
captured: std::sync::Arc<std::sync::Mutex<Option<String>>>,
_handle: std::thread::JoinHandle<()>,
}
fn spawn_mock_model_list_server(status_line: &str, body: &str) -> MockModelListServer {
use std::io::{Read, Write};
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("mock listener binds");
let address = listener.local_addr().expect("mock address");
let captured = std::sync::Arc::new(std::sync::Mutex::new(None));
let captured_for_thread = std::sync::Arc::clone(&captured);
let response = format!(
"HTTP/1.1 {status_line}\r\ncontent-type: application/json; charset=utf-8\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}",
body.len()
);
let handle = std::thread::spawn(move || {
let (mut stream, _) = listener.accept().expect("mock accept");
let mut buffer = [0u8; 8192];
let read = stream.read(&mut buffer).unwrap_or_default();
*captured_for_thread.lock().expect("captured lock") =
Some(String::from_utf8_lossy(&buffer[..read]).to_string());
let _ = stream.write_all(response.as_bytes());
let _ = stream.flush();
});
MockModelListServer {
base_url: format!("http://{address}/v1"),
captured,
_handle: handle,
}
}
fn model_list_state(base_url: &str) -> AppState {
AppState::new(crate::config::AppConfig {
llm_router_base_url: base_url.to_string(),
..crate::config::AppConfig::default()
})
.expect("state should build")
}
#[tokio::test]
async fn fetch_upstream_model_names_reads_group_pricing_without_credentials() {
let server = spawn_mock_model_list_server(
"200 OK",
&json!({"data": [{"model_name": "glm-5.3"}, {"model_name": "deepseek-flash"}]})
.to_string(),
);
let state = model_list_state(&server.base_url);
assert_eq!(
fetch_upstream_model_names(&state).await.unwrap(),
vec!["glm-5.3".to_string(), "deepseek-flash".to_string()]
);
let request = server
.captured
.lock()
.expect("captured lock")
.clone()
.expect("mock server should capture request");
// 控制面路径由 base_url 推导(去掉 /v1),并显式带 AGC 账号所在分组。
assert!(
request.starts_with("GET /api/pricing?group=taonier HTTP/1.1"),
"{request}"
);
// 定价列表是公开只读接口:不得把任何凭据发过去。
assert!(
!request.to_ascii_lowercase().contains("authorization:"),
"{request}"
);
}
#[tokio::test]
async fn fetch_upstream_model_names_fails_closed_when_upstream_unavailable_or_empty() {
let unauthorized = spawn_mock_model_list_server("401 Unauthorized", "{}");
let state = model_list_state(&unauthorized.base_url);
assert_eq!(
fetch_upstream_model_names(&state).await.unwrap_err(),
"上游模型列表返回 HTTP 401 Unauthorized"
);
let empty = spawn_mock_model_list_server("200 OK", &json!({"data": []}).to_string());
let state = model_list_state(&empty.base_url);
assert_eq!(
fetch_upstream_model_names(&state).await.unwrap_err(),
"上游模型列表为空"
);
let failing = spawn_mock_model_list_server("500 Internal Server Error", "{}");
let state = model_list_state(&failing.base_url);
assert_eq!(
fetch_upstream_model_names(&state).await.unwrap_err(),
"上游模型列表返回 HTTP 500 Internal Server Error"
);
}
}
+38
View File
@@ -2749,6 +2749,44 @@ mod tests {
}
}
#[tokio::test]
async fn editor_scene_generation_rejects_invalid_idempotency_key_before_queueing() {
let state = AppState::new(AppConfig {
external_generation_mode: ExternalGenerationMode::Queue,
..AppConfig::default()
})
.expect("state should build");
let seed_user = seed_phone_user_with_password(&state, "13800138232", TEST_PASSWORD).await;
let token = sign_test_user_token(&state, &seed_user, "sess_editor_scene_idempotency");
state.fail_test_editor_generation_enqueue();
let app = build_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri("/api/editor/scenes/generations")
.header("authorization", format!("Bearer {token}"))
.header("content-type", "application/json")
.header("idempotency-key", "contains space")
.body(Body::from(
serde_json::json!({
"sceneContent": "雨夜小镇",
"stylePreset": "anime",
"generationInputs": { "source": "ai-game-creator-client" },
})
.to_string(),
))
.expect("request should build"),
)
.await
.expect("request should complete");
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
assert_eq!(state.test_editor_generation_enqueue_attempts(), 0);
let body = response.into_body().collect().await.unwrap().to_bytes();
assert!(String::from_utf8_lossy(&body).contains("Idempotency-Key"));
}
#[tokio::test]
async fn editor_scene_generation_rejects_inline_data_url_before_queueing() {
let state = AppState::new(AppConfig {
@@ -826,6 +826,19 @@ mod tests {
editor_generation_idempotency_namespace(&game_creator),
GAME_CREATOR_CLIENT_GENERATION_DEDUPE_PREFIX
);
let scene = crate::editor_project::build_editor_scene_image_generation_payload(
serde_json::from_value(json!({
"sceneContent": "雨夜小镇",
"stylePreset": "anime",
"generationInputs": { "source": GAME_CREATOR_CLIENT_GENERATION_SOURCE },
}))
.expect("scene request should deserialize"),
)
.expect("scene payload should build");
assert_eq!(
editor_generation_idempotency_namespace(&scene),
GAME_CREATOR_CLIENT_GENERATION_DEDUPE_PREFIX,
);
assert_eq!(
editor_generation_idempotency_namespace(&ordinary),
EXTERNAL_API_GENERATION_DEDUPE_PREFIX
+124 -13
View File
@@ -2658,12 +2658,23 @@ fn build_editor_scene_generation_inputs(
{
fields.push(json!({ "id": "customStyle", "title": "自定义画风", "value": custom_style }));
}
json!({
let mut inputs = json!({
"version": 2,
"action": "scene.generate",
"fields": fields,
"references": references,
})
});
// AGC 账号任务依赖来源标记选择幂等命名空间和可下载的队列结果;其余字段仍由服务端重建。
if payload
.generation_inputs
.as_ref()
.and_then(|inputs| inputs.get("source"))
.and_then(Value::as_str)
== Some(GAME_CREATOR_CLIENT_GENERATION_SOURCE)
{
inputs["source"] = json!(GAME_CREATOR_CLIENT_GENERATION_SOURCE);
}
inputs
}
fn normalize_editor_scene_optional_text<'a>(value: Option<&'a str>, default: &'a str) -> &'a str {
@@ -2692,13 +2703,9 @@ fn normalize_editor_scene_asset_label(asset_label: Option<String>) -> String {
resolve_editor_generated_asset_label(asset_label, "游戏场景")
}
pub async fn generate_editor_scene(
State(state): State<AppState>,
Extension(request_context): Extension<RequestContext>,
Extension(authenticated): Extension<AuthenticatedAccessToken>,
payload: Result<Json<EditorSceneGenerateRequest>, JsonRejection>,
) -> Result<Json<Value>, AppError> {
let Json(payload) = parse_editor_generation_json_payload(payload)?;
pub(crate) fn build_editor_scene_image_generation_payload(
payload: EditorSceneGenerateRequest,
) -> Result<EditorImageGenerationRequest, AppError> {
let generation_options = normalize_editor_scene_generation_options(
payload.model.as_deref(),
payload.aspect_ratio.as_deref(),
@@ -2717,8 +2724,7 @@ pub async fn generate_editor_scene(
}))
})?;
let generation_inputs = build_editor_scene_generation_inputs(&payload, &generation_options);
let caller = EditorGenerationCaller::from_authenticated(&authenticated);
let image_payload = EditorImageGenerationRequest {
Ok(EditorImageGenerationRequest {
prompt,
size: None,
kind: Some("scene".to_string()),
@@ -2737,14 +2743,27 @@ pub async fn generate_editor_scene(
asset_label: Some(normalize_editor_scene_asset_label(payload.asset_label)),
source_resource_id: None,
canvas_completion: payload.canvas_completion,
};
})
}
pub async fn generate_editor_scene(
State(state): State<AppState>,
Extension(request_context): Extension<RequestContext>,
Extension(authenticated): Extension<AuthenticatedAccessToken>,
headers: HeaderMap,
payload: Result<Json<EditorSceneGenerateRequest>, JsonRejection>,
) -> Result<Json<Value>, AppError> {
let Json(payload) = parse_editor_generation_json_payload(payload)?;
let idempotency_key = optional_editor_idempotency_key(&headers)?;
let image_payload = build_editor_scene_image_generation_payload(payload)?;
let caller = EditorGenerationCaller::from_authenticated(&authenticated);
if !state.config.external_generation_mode.is_inline() {
let queue_job = enqueue_editor_image_generation_for_owner(
&state,
&request_context,
&caller,
image_payload,
None,
idempotency_key,
)
.await?;
return Ok(json_success_body(
@@ -13269,6 +13288,98 @@ mod tests {
thread,
};
#[test]
fn scene_image_generation_payload_uses_shared_scene_contract() {
let payload = build_editor_scene_image_generation_payload(EditorSceneGenerateRequest {
scene_content: "雨夜小镇的街道".to_string(),
style_preset: "custom".to_string(),
custom_style: Some("像素水彩混合".to_string()),
model: None,
aspect_ratio: None,
image_size: None,
reference_image_srcs: vec!["resource-ref-1".to_string()],
project_id: Some("project-1".to_string()),
generation_inputs: None,
asset_folder_id: Some("folder-1".to_string()),
asset_label: Some(" ".to_string()),
canvas_completion: None,
})
.expect("合法场景意图必须组装成功");
assert_eq!(payload.kind.as_deref(), Some("scene"));
assert_eq!(payload.asset_kind.as_deref(), Some("scene"));
assert_eq!(payload.aspect_ratio.as_deref(), Some("16:9"));
assert_eq!(payload.image_size.as_deref(), Some("1K"));
assert_eq!(payload.asset_label.as_deref(), Some("游戏场景"));
assert_eq!(
payload.reference_image_srcs.as_deref(),
Some(["resource-ref-1".to_string()].as_slice())
);
assert!(payload.prompt.contains("雨夜小镇的街道"));
assert!(payload.prompt.contains("像素水彩混合"));
let generation_inputs = payload
.generation_inputs
.as_ref()
.expect("场景 generationInputs 必须存在");
assert_eq!(generation_inputs["version"], json!(2));
assert_eq!(generation_inputs["action"], json!("scene.generate"));
assert_eq!(
generation_inputs["references"],
json!([{ "id": "reference" }])
);
}
#[test]
fn scene_generation_preserves_agc_downloadable_queue_result() {
let request = serde_json::from_value::<EditorSceneGenerateRequest>(json!({
"sceneContent": "雨夜小镇",
"stylePreset": "anime",
"referenceImageSrcs": ["art-spec-resource"],
"generationInputs": {
"source": GAME_CREATOR_CLIENT_GENERATION_SOURCE,
"action": "client-supplied-action",
"fields": [{ "id": "prompt", "value": "不可信配方" }],
"references": [{ "refId": "untrusted-resource" }],
},
}))
.expect("scene request should deserialize");
let mut image_payload = build_editor_scene_image_generation_payload(request)
.expect("scene payload should build");
image_payload.generation_inputs =
sanitize_editor_queued_generation_inputs(image_payload.generation_inputs);
let inputs = image_payload.generation_inputs.as_ref().unwrap();
assert_eq!(inputs["source"], GAME_CREATOR_CLIENT_GENERATION_SOURCE);
assert_eq!(inputs["action"], "scene.generate");
assert_eq!(inputs["fields"][0]["value"], "雨夜小镇");
assert_eq!(inputs["references"], json!([{ "id": "reference" }]));
let mut job = atomic_editor_generation_job_fixture();
job.request_payload_json = serde_json::to_string(&image_payload).unwrap();
let context = EditorGenerationQueueResultContext::from_job(&job);
assert_eq!(
context.consumer,
EditorGenerationQueueConsumer::GameCreatorResourceEditor
);
let result: Value = serde_json::from_str(
&serialize_atomic_editor_generation_job_result(
&context,
&json!({
"ok": true,
"objectKey": "generated/scene.png",
"resource": {
"resourceId": "scene-resource",
"objectKey": "generated/scene.png",
"assetObjectId": "scene-object",
},
}),
)
.expect("AGC scene result should serialize"),
)
.unwrap();
assert_eq!(result["result"]["objectKey"], "generated/scene.png");
assert_eq!(result["result"]["resource"]["resourceId"], "scene-resource");
}
#[test]
fn background_removal_options_preserve_queue_parameters_and_legacy_identity() {
for (fields, mode, color) in [
@@ -49,7 +49,7 @@ const EXTERNAL_API_KEY_SCOPES: [&str; 4] = [
const LLM_ROUTER_TOKEN_IDENTIFIER: &str = "agc_auto_generate";
/// Router 用户(账号)与它名下固定 Token / API Key 都归属同一分组 `taonier`。
const LLM_ROUTER_USER_GROUP: &str = "taonier";
const LLM_ROUTER_TOKEN_GROUP: &str = "taonier";
pub(crate) const LLM_ROUTER_TOKEN_GROUP: &str = "taonier";
const LLM_ROUTER_API_KEY_SCOPES: [&str; 1] = ["llm:responses"];
const LLM_ROUTER_SUBSCRIPTION_PLAN_ID: i64 = 1;
const LLM_ROUTER_SUBSCRIPTION_RENEWAL_THRESHOLD_SECONDS: i64 = 24 * 60 * 60;
@@ -1624,7 +1624,7 @@ async fn ensure_router_token_contract(
Ok(())
}
fn router_control_origin(base_url: &str) -> Result<String, String> {
pub(crate) fn router_control_origin(base_url: &str) -> Result<String, String> {
let mut url = reqwest::Url::parse(base_url.trim_end_matches('/'))
.map_err(|error| format!("LLM Router 地址无效:{error}"))?;
let is_loopback = url.host_str().is_some_and(|host| {
@@ -1643,7 +1643,11 @@ fn router_control_origin(base_url: &str) -> Result<String, String> {
Ok(url.to_string().trim_end_matches('/').to_string())
}
fn ensure_llm_router_target_allowed(state: &AppState) -> Result<(), String> {
/// 只校验 LLM Router 目标地址是否允许(官方路由 / loopback、scheme),不校验固定模型。
///
/// 与具体模型无关的调用(例如按分组定价列表同步 AGC 模型目录)用这个入口,
/// 避免被“必须使用官方固定模型”的哨兵常量挡住。
pub(crate) fn ensure_llm_router_url_allowed(state: &AppState) -> Result<(), String> {
let base_url = state.config.llm_router_base_url.trim_end_matches('/');
let url =
reqwest::Url::parse(base_url).map_err(|error| format!("LLM Router 地址无效:{error}"))?;
@@ -1664,9 +1668,6 @@ fn ensure_llm_router_target_allowed(state: &AppState) -> Result<(), String> {
if base_url != OFFICIAL_LLM_ROUTER_BASE_URL {
return Err("生产环境 LLM Router 必须使用官方固定路由".to_string());
}
if state.config.llm_router_model.trim() != OFFICIAL_LLM_ROUTER_MODEL {
return Err("生产环境 LLM Router 必须使用官方固定模型".to_string());
}
if url.scheme() != "https" {
return Err("生产环境 LLM Router 只允许 HTTPS 地址".to_string());
}
@@ -1674,9 +1675,6 @@ fn ensure_llm_router_target_allowed(state: &AppState) -> Result<(), String> {
}
if base_url == OFFICIAL_LLM_ROUTER_BASE_URL {
if state.config.llm_router_model.trim() != OFFICIAL_LLM_ROUTER_MODEL {
return Err("LLM Router 必须使用官方固定模型".to_string());
}
if url.scheme() != "https" {
return Err("官方 LLM Router 只允许 HTTPS 地址".to_string());
}
@@ -1698,6 +1696,19 @@ fn ensure_llm_router_target_allowed(state: &AppState) -> Result<(), String> {
Ok(())
}
pub(crate) fn ensure_llm_router_target_allowed(state: &AppState) -> Result<(), String> {
ensure_llm_router_url_allowed(state)?;
if state.config.llm_router_model.trim() != OFFICIAL_LLM_ROUTER_MODEL {
if state.config.is_production() {
return Err("生产环境 LLM Router 必须使用官方固定模型".to_string());
}
if state.config.llm_router_base_url.trim_end_matches('/') == OFFICIAL_LLM_ROUTER_BASE_URL {
return Err("LLM Router 必须使用官方固定模型".to_string());
}
}
Ok(())
}
fn router_username_for_owner(owner_user_id: &str) -> String {
// New API 的 User.Username 校验上限是 20 个字符。保留可读前缀后只
// 能放 11 个字符;使用完整 owner id 做 SHA-256,再编码成 8 字节的
@@ -7,7 +7,9 @@ use axum::{
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use shared_contracts::assets::EditorCanvasGenerationCompletionPayload;
use shared_contracts::assets::{
EditorCanvasGenerationCompletionPayload, EditorSceneGenerateRequest,
};
use shared_contracts::external_generation::{
ExternalEditorGenerationJobResponse, ExternalEditorGenerationSubmissionResponse,
ExternalGenerationJobStatus,
@@ -38,12 +40,13 @@ use crate::{
EditorCanvasViewportPayload, EditorGenerationCaller, EditorImageEditRequest,
EditorImageGenerationRequest, EditorProjectListQuery, EditorProjectListView,
EditorProjectPayload, EditorProjectResourcePayload, EditorProjectSummaryListResponse,
EditorUiDesignAssetExtractionRequest, current_utc_micros,
editor_asset_folder_payload_from_record, editor_asset_library_payload_from_record,
editor_asset_payload_from_record, editor_idempotent_create_id,
editor_project_payload_from_record, editor_project_resource_payload_from_record,
editor_project_summary_from_record, enqueue_editor_background_removal_for_owner,
enqueue_editor_image_edit_for_owner, enqueue_editor_image_generation_for_owner,
EditorUiDesignAssetExtractionRequest, build_editor_scene_image_generation_payload,
current_utc_micros, editor_asset_folder_payload_from_record,
editor_asset_library_payload_from_record, editor_asset_payload_from_record,
editor_idempotent_create_id, editor_project_payload_from_record,
editor_project_resource_payload_from_record, editor_project_summary_from_record,
enqueue_editor_background_removal_for_owner, enqueue_editor_image_edit_for_owner,
enqueue_editor_image_generation_for_owner,
enqueue_editor_ui_design_asset_extraction_for_owner,
ensure_generic_editor_image_generation_contract, map_editor_project_error,
normalize_editor_persisted_media_src, normalize_optional_string,
@@ -750,6 +753,29 @@ pub async fn generate_external_editor_image(
Ok(external_generation_accepted_response(&request_context, job))
}
pub async fn generate_external_editor_scene(
State(state): State<AppState>,
Extension(request_context): Extension<RequestContext>,
Extension(principal): Extension<ExternalApiPrincipal>,
headers: HeaderMap,
payload: Result<Json<EditorSceneGenerateRequest>, JsonRejection>,
) -> Result<Response, AppError> {
let Json(payload) = parse_editor_generation_json_payload(payload)?;
require_scope(&principal, SCOPE_EDITOR_IMAGE_GENERATE)?;
let idempotency_key = require_idempotency_key(&headers)?;
let project_id = payload.project_id.clone();
let image_payload = build_editor_scene_image_generation_payload(payload)?;
let job = enqueue_editor_image_generation_for_owner(
&state,
&request_context,
&editor_generation_caller(&principal, project_id),
image_payload,
Some(idempotency_key),
)
.await?;
Ok(external_generation_accepted_response(&request_context, job))
}
pub async fn edit_external_editor_image(
State(state): State<AppState>,
Extension(request_context): Extension<RequestContext>,
@@ -1792,6 +1818,61 @@ mod tests {
.await;
}
#[tokio::test]
async fn external_scene_generation_rejects_invalid_intent_before_queueing() {
let state = AppState::new(crate::config::AppConfig::default())
.expect("external scene test state should build");
state.fail_test_editor_generation_enqueue();
let app = Router::new()
.route(
"/api/external/v1/editor/scenes/generations",
post(generate_external_editor_scene),
)
.layer(Extension(request_context(false)))
.layer(Extension(ExternalApiPrincipal::for_test(
"user-external-scene",
&[SCOPE_EDITOR_IMAGE_GENERATE],
)))
.with_state(state.clone());
for (case_name, request_body) in [
(
"empty sceneContent",
json!({"sceneContent": " ", "stylePreset": "anime"}),
),
(
"unknown stylePreset",
json!({"sceneContent": "雨夜小镇", "stylePreset": "oil-painting"}),
),
(
"custom without customStyle",
json!({"sceneContent": "雨夜小镇", "stylePreset": "custom"}),
),
("missing sceneContent", json!({"stylePreset": "anime"})),
] {
let response = app
.clone()
.oneshot(
axum::http::Request::builder()
.method("POST")
.uri("/api/external/v1/editor/scenes/generations")
.header("content-type", "application/json")
.header(IDEMPOTENCY_KEY_HEADER, "scene-contract-test")
.body(Body::from(request_body.to_string()))
.expect("external scene request should build"),
)
.await
.expect("external scene response should return");
assert_eq!(response.status(), StatusCode::BAD_REQUEST, "{case_name}");
assert_eq!(
state.test_editor_generation_enqueue_attempts(),
0,
"{case_name} must fail before queueing",
);
}
}
#[tokio::test]
async fn external_background_removal_rejects_invalid_parameters_before_queueing() {
let state = AppState::new(crate::config::AppConfig::default())
@@ -2221,6 +2302,21 @@ mod tests {
}
}
#[test]
fn external_openapi_documents_dedicated_scene_generation_route() {
let parsed: Value = serde_json::from_str(OPENAPI_JSON).expect("openapi json should parse");
let operation = &parsed["paths"]["/api/external/v1/editor/scenes/generations"]["post"];
assert_eq!(
operation["requestBody"]["content"]["application/json"]["schema"]["$ref"],
json!("#/components/schemas/EditorSceneGenerationRequest")
);
assert_eq!(
parsed["components"]["schemas"]["EditorSceneGenerationRequest"]["required"],
json!(["sceneContent", "stylePreset"])
);
}
#[test]
fn exported_openapi_json_contains_external_editor_routes_and_security() {
let parsed: Value = serde_json::from_str(OPENAPI_JSON).expect("openapi json should parse");
+356 -38
View File
@@ -25,6 +25,8 @@ use tower::ServiceExt;
use crate::{modules, request_context::RequestContext, state::AppState};
mod semantic;
const OPENAPI_JSON: &str =
include_str!("../../../../docs/openapi/genarrative-external-v1.openapi.json");
const SKILL_MD: &str =
@@ -54,7 +56,15 @@ const SKILL_REQUESTS_AND_OUTPUTS_URI: &str =
"genarrative://external-editor/skill/references/requests-and-outputs.md";
const MAX_MCP_REST_RESPONSE_BYTES: usize = 4 * 1024 * 1024;
const MCP_INSTRUCTIONS: &str = r#"陶泥儿外部编辑器工具。先创建或复用画布项目,并创建与画布同名的素材文件夹;生成结果应同时写入画布和素材库。参考本地文件时先走上传票据和对象确认,不要把 Data URL、Blob URL 或临时签名 URL写入生成参数。所有生成工具都是异步提交:必须提供 idempotencyKey,提交后按 pollAfterMs 调用 get_external_editor_generation_job,只有 status=completed 时消费 result;查询超时不能重新提交。图集生成必须显式声明 sliceMode,没有默认值:需求要求等分网格、固定槽位或指定行列数时用 grid 并提供来自需求的 gridX/gridY,自由排布或数量不定时用 connected-components(可用 sliceCount 约束张数),connected-components 不接受 gridX/gridY;缺失、越界或自相矛盾在计费前返回 400。warning 表示主结果可用但存在降级,sliceWarning 表示完整透明图集可用但切片未完成。详细说明、OpenAPI、Skill 主入口和分主题 references 见 resources/list;需要本地文件编排或不支持 MCP 时再下载 skill.zip。"#;
const MCP_INSTRUCTIONS: &str = r#"陶泥儿提供画布项目管理、素材库管理,以及图片、角色动画、视频和音频生成能力。
按用户任务需要创建或复用项目、素材文件夹,不默认创建。生成结果需要进入画布或素材库时,使用对应生成工具支持的目标字段;已有落库结果不要重复登记。
生成操作会产生费用,采用异步提交。每次独立生成使用稳定的 idempotencyKey;取得任务 ID 后,使用 check_generation 按 pollAfterMs 查询,直到 completed 或 failed。查询超时不代表生成失败,不要因此重新提交或更换幂等键。
本地参考文件通过 prepare_asset_upload 获取上传票据,由调用方实际上传后确认对象。后续引用遵循各工具要求,使用稳定的对象键或资源、素材 ID;不要把临时下载 URL 当作持久引用。
以实际返回的结果和告警判断完成情况,部分产物成功不代表所有处理步骤成功。具体参数以工具 schema 和说明为准;需要详细流程、示例或 API 契约时,通过 resources/list 查找相关文档。"#;
#[derive(Clone, Debug)]
struct McpOperation {
@@ -125,7 +135,11 @@ impl ServerHandler for GenarrativeExternalMcp {
_context: McpRequestContext<RoleServer>,
) -> Result<ListToolsResult, ErrorData> {
Ok(ListToolsResult::with_all_items(
MCP_OPERATIONS.iter().map(mcp_operation_tool).collect(),
MCP_OPERATIONS
.iter()
.map(mcp_operation_tool)
.chain(semantic::TOOLS.iter().map(|entry| entry.tool.clone()))
.collect(),
))
}
@@ -134,6 +148,7 @@ impl ServerHandler for GenarrativeExternalMcp {
.iter()
.find(|operation| operation.tool_name == name)
.map(mcp_operation_tool)
.or_else(|| semantic::find(name).map(|entry| entry.tool.clone()))
}
async fn call_tool(
@@ -141,12 +156,30 @@ impl ServerHandler for GenarrativeExternalMcp {
request: CallToolRequestParams,
context: McpRequestContext<RoleServer>,
) -> Result<CallToolResult, ErrorData> {
if let Some(tool) = semantic::find(request.name.as_ref()) {
let result = match tool.prepare(request.arguments.unwrap_or_default()) {
Ok(call) => {
dispatch_operation(
call.operation,
call.arguments,
&context,
call.optional_idempotency_key,
)
.await
}
Err(error) => Err(error),
};
return Ok(match result {
Ok(value) => CallToolResult::structured(value),
Err(value) => CallToolResult::structured_error(value),
});
}
let operation = MCP_OPERATIONS
.iter()
.find(|operation| operation.tool_name == request.name.as_ref())
.ok_or_else(|| ErrorData::invalid_params("未知的陶泥儿外部 API 工具", None))?;
let arguments = request.arguments.unwrap_or_default();
match dispatch_operation(operation, arguments, &context).await {
match dispatch_operation(operation, arguments, &context, None).await {
Ok(value) => Ok(CallToolResult::structured(value)),
Err(value) => Ok(CallToolResult::structured_error(value)),
}
@@ -475,6 +508,7 @@ async fn dispatch_operation(
operation: &McpOperation,
arguments: Map<String, Value>,
context: &McpRequestContext<RoleServer>,
optional_idempotency_key: Option<axum::http::HeaderValue>,
) -> Result<Value, Value> {
validate_required_body(operation, &arguments)?;
@@ -498,6 +532,51 @@ async fn dispatch_operation(
.cloned()
.ok_or_else(|| json!({"error": "Authorization 请求头缺失"}))?;
let request = build_operation_request(
operation,
&arguments,
authorization,
request_context,
optional_idempotency_key,
)?;
let response = modules::external_api::router(state.clone())
.with_state(state)
.oneshot(request)
.await
.unwrap_or_else(|never| match never {});
let status = response.status();
let bytes = response
.into_body()
.collect()
.await
.map_err(|_| json!({"error": "读取外部 API 响应失败"}))?
.to_bytes();
if bytes.len() > MAX_MCP_REST_RESPONSE_BYTES {
return Err(json!({"error": "外部 API 响应超过 MCP 返回上限"}));
}
let payload = serde_json::from_slice::<Value>(&bytes).unwrap_or_else(|_| {
json!({
"status": status.as_u16(),
"message": "外部 API 返回了非 JSON 响应"
})
});
if status.is_success() {
Ok(unwrap_external_api_success_payload(payload))
} else {
Err(json!({
"status": status.as_u16(),
"response": payload,
}))
}
}
fn build_operation_request(
operation: &McpOperation,
arguments: &Map<String, Value>,
authorization: axum::http::HeaderValue,
request_context: RequestContext,
optional_idempotency_key: Option<axum::http::HeaderValue>,
) -> Result<Request<Body>, Value> {
let mut path = operation.path_template.clone();
if let Some(path_parameters) = arguments.get("pathParameters").and_then(Value::as_object) {
for (name, value) in path_parameters {
@@ -532,37 +611,12 @@ async fn dispatch_operation(
"application/json".parse().expect("valid content type"),
);
}
apply_operation_headers(operation, &arguments, request.headers_mut())?;
apply_operation_headers(operation, arguments, request.headers_mut())?;
if let Some(key) = optional_idempotency_key {
request.headers_mut().insert("idempotency-key", key);
}
let response = modules::external_api::router(state.clone())
.with_state(state)
.oneshot(request)
.await
.unwrap_or_else(|never| match never {});
let status = response.status();
let bytes = response
.into_body()
.collect()
.await
.map_err(|_| json!({"error": "读取外部 API 响应失败"}))?
.to_bytes();
if bytes.len() > MAX_MCP_REST_RESPONSE_BYTES {
return Err(json!({"error": "外部 API 响应超过 MCP 返回上限"}));
}
let payload = serde_json::from_slice::<Value>(&bytes).unwrap_or_else(|_| {
json!({
"status": status.as_u16(),
"message": "外部 API 返回了非 JSON 响应"
})
});
if status.is_success() {
Ok(unwrap_external_api_success_payload(payload))
} else {
Err(json!({
"status": status.as_u16(),
"response": payload,
}))
}
Ok(request)
}
fn apply_operation_headers(
@@ -977,15 +1031,16 @@ mod tests {
let tools = MCP_OPERATIONS
.iter()
.map(mcp_operation_tool)
.chain(semantic::TOOLS.iter().map(|entry| entry.tool.clone()))
.collect::<Vec<_>>();
let serialized = serde_json::to_vec(&tools).expect("tool catalog should serialize");
assert!(serialized.len() < 512 * 1024);
for operation in MCP_OPERATIONS.iter() {
let serialized = serde_json::to_string(&operation.input_schema)
for tool in tools {
let serialized = serde_json::to_string(&tool.input_schema)
.expect("tool input schema should serialize");
assert!(serialized.len() < 64 * 1024, "{}", operation.tool_name);
assert!(!serialized.contains("\"$ref\""), "{}", operation.tool_name);
assert_eq!(operation.input_schema.get("type"), Some(&json!("object")));
assert!(serialized.len() < 64 * 1024, "{}", tool.name);
assert!(!serialized.contains("\"$ref\""), "{}", tool.name);
assert_eq!(tool.input_schema.get("type"), Some(&json!("object")));
}
}
@@ -1024,6 +1079,269 @@ mod tests {
})),
json!({"operationId": "task-1", "status": "queued"})
);
let result = json!({
"operationId": "task-1", "status": "completed",
"result": {"objectKey": "media/sheet.png", "warning": {"code": "source_preserved"},
"sliceWarning": {"code": "slice_failed"}, "project": {"revision": 7}}
});
assert_eq!(
unwrap_external_api_success_payload(json!({"ok": true, "data": result})),
result
);
}
fn rpc_request(method: &str, params: Value) -> Request<Body> {
Request::builder()
.method(Method::POST)
.uri("/api/external/v1/mcp")
.header(HOST, "localhost")
.header(CONTENT_TYPE, "application/json")
.header(ACCEPT, "application/json, text/event-stream")
.header("mcp-protocol-version", "2025-11-25")
.body(Body::from(
json!({"jsonrpc": "2.0", "id": 1, "method": method, "params": params}).to_string(),
))
.unwrap()
}
async fn rpc_payload(response: axum::response::Response) -> Value {
assert_eq!(response.status(), StatusCode::OK);
serde_json::from_slice(&response.into_body().collect().await.unwrap().to_bytes()).unwrap()
}
#[tokio::test]
async fn semantic_catalog_appends_tools_without_changing_legacy_definitions() {
let payload = rpc_payload(
service()
.oneshot(rpc_request("tools/list", json!({})))
.await
.unwrap()
.map(Body::new),
)
.await;
let tools = payload["result"]["tools"].as_array().unwrap();
assert_eq!(tools.len(), MCP_OPERATIONS.len() + 15);
for op in MCP_OPERATIONS.iter() {
let expected = serde_json::to_value(mcp_operation_tool(op)).unwrap();
assert_eq!(
tools.iter().find(|tool| tool["name"] == op.tool_name),
Some(&expected)
);
}
for entry in semantic::TOOLS.iter() {
assert_eq!(
GenarrativeExternalMcp.get_tool(&entry.tool.name),
Some(entry.tool.clone())
);
}
}
#[tokio::test]
async fn semantic_invalid_arguments_fail_before_http_context_or_side_effects() {
for (name, arguments) in [
(
"modify_image",
json!({"action": "edit", "input": {"prompt": "修改", "sourceImageSrc": "wrong-reference"}, "idempotencyKey": "test"}),
),
(
"delete_resources",
json!({"action": "delete_project", "input": {}}),
),
(
"manage_canvas_projects",
json!({"action": "rename", "input": {"projectId": "project", "title": "新名"}, "idempotencyKey": "not-supported"}),
),
("generate_image", json!({"prompt": "test"})),
] {
let payload = rpc_payload(
service()
.oneshot(rpc_request(
"tools/call",
json!({"name": name, "arguments": arguments}),
))
.await
.unwrap()
.map(Body::new),
)
.await;
assert_eq!(payload["result"]["isError"], true, "{name}: {payload}");
assert!(payload["result"]["structuredContent"]["error"].is_string());
assert!(!payload.to_string().contains("上下文缺失"));
}
}
#[tokio::test]
async fn semantic_adapter_builds_real_rest_paths_bodies_and_optional_headers() {
for (name, args, method, path, body, key) in [
(
"manage_canvas_projects",
json!({"action":"create","input":{},"idempotencyKey":"create-project"}),
Method::POST,
"/api/external/v1/editor/projects",
json!({}),
Some("create-project"),
),
(
"manage_canvas_projects",
json!({"action":"rename","input":{"projectId":"project/a","title":"新名"}}),
Method::PATCH,
"/api/external/v1/editor/projects/project%2Fa/metadata",
json!({"title":"新名"}),
None,
),
(
"find_assets",
json!({"action":"get_download_url","input":{"objectKey":"images/a b.png","expireSeconds":60}}),
Method::GET,
"/api/external/v1/assets/read-url?expireSeconds=60&objectKey=images%2Fa+b.png",
Value::Null,
None,
),
(
"find_canvas_projects",
json!({"action":"list","input":{}}),
Method::GET,
"/api/external/v1/editor/projects?view=summary",
Value::Null,
None,
),
(
"modify_image",
json!({"action":"variation","input":{"prompt":"变体","referenceImageSrcs":["ref"]},"idempotencyKey":"same-generation"}),
Method::POST,
"/api/external/v1/editor/images/generations",
json!({"prompt":"变体","referenceImageSrcs":["ref"],"kind":"quick-edit"}),
Some("same-generation"),
),
] {
let call = semantic::find(name)
.unwrap()
.prepare(args.as_object().unwrap().clone())
.unwrap();
let context = RequestContext::new(
"test-request".into(),
"POST /api/external/v1/mcp".into(),
std::time::Duration::ZERO,
false,
);
let request = build_operation_request(
call.operation,
&call.arguments,
"Bearer fixture".parse().unwrap(),
context,
call.optional_idempotency_key,
)
.unwrap();
assert_eq!(request.method(), method);
assert_eq!(request.uri().to_string(), path);
assert_eq!(request.headers()[AUTHORIZATION], "Bearer fixture");
assert_eq!(
request
.headers()
.get("idempotency-key")
.map(|v| v.to_str().unwrap()),
key
);
assert_eq!(
request
.extensions()
.get::<RequestContext>()
.unwrap()
.request_id(),
"test-request"
);
let bytes = request.into_body().collect().await.unwrap().to_bytes();
if body.is_null() {
assert!(bytes.is_empty());
} else {
assert_eq!(serde_json::from_slice::<Value>(&bytes).unwrap(), body);
}
}
let operation = MCP_OPERATIONS
.iter()
.find(|op| op.operation_id == "createEditorProject")
.unwrap();
let mut headers = HeaderMap::new();
apply_operation_headers(
operation,
&json!({"idempotencyKey": "legacy-ignored"})
.as_object()
.unwrap()
.clone(),
&mut headers,
)
.unwrap();
assert!(
headers.get("idempotency-key").is_none(),
"old optional-header behavior must remain unchanged"
);
}
#[tokio::test]
async fn semantic_calls_reuse_rest_scope_checks_and_structured_errors() {
use crate::state::external_api_auth::ExternalApiKeyAuthenticator;
use futures_util::future::BoxFuture;
use spacetime_client::{
ExternalApiKeyAuthenticateRecordInput, ExternalApiKeyRecord, SpacetimeClientError,
};
use std::sync::atomic::{AtomicUsize, Ordering};
struct NoScopes(AtomicUsize);
impl ExternalApiKeyAuthenticator for NoScopes {
fn authenticate_external_api_key(
&self,
_: ExternalApiKeyAuthenticateRecordInput,
) -> BoxFuture<'_, Result<ExternalApiKeyRecord, SpacetimeClientError>> {
self.0.fetch_add(1, Ordering::Relaxed);
Box::pin(async {
Ok(ExternalApiKeyRecord {
key_id: "fixture-key".into(),
owner_user_id: "owner-from-store".into(),
name: "测试".into(),
key_prefix: "tnr_sk_fixture".into(),
scopes: vec![],
created_at: "0.000000Z".into(),
last_used_at: None,
revoked_at: None,
updated_at: "0.000000Z".into(),
})
})
}
}
let auth = Arc::new(NoScopes(AtomicUsize::new(0)));
let state = AppState::new(AppConfig::default())
.unwrap()
.with_external_api_auth_state(crate::state::ExternalApiAuthState::new(auth.clone()));
let router = modules::external_api::router(state.clone())
.with_state(state)
.layer(middleware::from_fn(attach_request_context));
for (name, arguments) in [
(
"delete_resources",
json!({"action":"delete_project","input":{"projectId":"fixture-project"}}),
),
(
"delete_editor_project",
json!({"pathParameters":{"projectId":"fixture-project"}}),
),
] {
let mut request = rpc_request("tools/call", json!({"name":name,"arguments":arguments}));
request
.headers_mut()
.insert(AUTHORIZATION, "Bearer tnr_sk_fixture".parse().unwrap());
let payload = rpc_payload(router.clone().oneshot(request).await.unwrap()).await;
assert_eq!(payload["result"]["isError"], true, "{payload}");
assert_eq!(
payload["result"]["structuredContent"]["status"], 403,
"{payload}"
);
assert!(!payload.to_string().contains("owner-from-store"));
}
assert_eq!(
auth.0.load(Ordering::Relaxed),
4,
"outer MCP and inner REST both authenticate"
);
}
#[tokio::test]
@@ -0,0 +1,460 @@
//! 语义入口只负责操作选择和参数位置转换,业务校验与副作用仍由 External router 承担。
use super::*;
use axum::http::HeaderValue;
pub(super) static TOOLS: LazyLock<Vec<SemanticTool>> = LazyLock::new(build_tools);
pub(super) struct SemanticTool {
pub(super) tool: Tool,
actions: Vec<Action>,
}
struct Action {
name: Option<&'static str>,
operation: &'static McpOperation,
input_schema: Value,
key_schema: Option<Value>,
fixed_body: Map<String, Value>,
destructive: bool,
}
pub(super) struct PreparedCall {
pub(super) operation: &'static McpOperation,
pub(super) arguments: Map<String, Value>,
pub(super) optional_idempotency_key: Option<HeaderValue>,
}
pub(super) fn find(name: &str) -> Option<&'static SemanticTool> {
TOOLS.iter().find(|entry| entry.tool.name == name)
}
fn build_tools() -> Vec<SemanticTool> {
let openapi: Value = serde_json::from_str(OPENAPI_JSON).expect("embedded OpenAPI must parse");
let descriptions: Value = serde_json::from_str(include_str!(
"../../prompts/external_mcp/semantic_tools.json"
))
.expect("semantic tool descriptions must parse");
let definitions: &[(&str, &[(&str, &str)])] = &[
(
"find_canvas_projects",
&[
("list", "listEditorProjects"),
("recent", "loadRecentEditorProject"),
("get", "getEditorProject"),
],
),
(
"manage_canvas_projects",
&[
("create", "createEditorProject"),
("rename", "renameEditorProject"),
],
),
(
"find_assets",
&[
("list_library", "getEditorAssetLibrary"),
("get_project_resources", "getEditorProject"),
("get_download_url", "getExternalAssetReadUrl"),
],
),
(
"prepare_asset_upload",
&[
("create_upload_ticket", "createExternalDirectUploadTicket"),
("confirm_upload", "confirmExternalAssetObject"),
],
),
("generate_image", &[("", "generateExternalEditorImage")]),
(
"modify_image",
&[
("edit", "editExternalEditorImage"),
("variation", "generateExternalEditorImage"),
("remove_background", "removeExternalEditorImageBackground"),
],
),
(
"generate_icon_spritesheet",
&[("", "generateExternalEditorIconSpritesheet")],
),
(
"extract_ui_assets",
&[("", "extractExternalEditorUiDesignAssets")],
),
(
"generate_character_animation",
&[("", "generateExternalEditorCharacterAnimation")],
),
("generate_video", &[("", "generateExternalEditorVideo")]),
(
"generate_audio",
&[
("sound_effect", "generateExternalEditorSoundEffect"),
("background_music", "generateExternalEditorBackgroundMusic"),
],
),
(
"edit_canvas",
&[
("get", "getEditorProject"),
("save_layout", "saveEditorProjectCanvas"),
("register_resource", "createEditorProjectResource"),
],
),
(
"organize_asset_library",
&[
("create_folder", "createEditorAssetFolder"),
("update_folder", "updateEditorAssetFolder"),
("create_asset", "createEditorAsset"),
("update_asset", "updateEditorAsset"),
],
),
(
"check_generation",
&[("", "getExternalEditorGenerationJob")],
),
(
"delete_resources",
&[
("delete_project", "deleteEditorProject"),
("delete_folder", "deleteEditorAssetFolder"),
("delete_asset", "deleteEditorAsset"),
],
),
];
definitions
.iter()
.map(|(name, operations)| {
let actions = operations
.iter()
.map(|(action, operation)| {
Action::new((!action.is_empty()).then_some(*action), operation, &openapi)
})
.collect::<Vec<_>>();
let read_only = actions.iter().all(|a| a.operation.method == Method::GET);
let generation = actions.iter().any(|a| a.operation.requires_idempotency_key);
let destructive = actions.iter().any(|a| a.destructive);
let schema = tool_schema(&actions);
let mut tool = Tool::new(
name.to_string(),
descriptions[name]["description"]
.as_str()
.expect("tool description")
.to_string(),
Arc::new(schema.as_object().expect("object schema").clone()),
);
tool.title = Some(
descriptions[name]["title"]
.as_str()
.expect("tool title")
.to_string(),
);
tool.annotations = Some(
ToolAnnotations::new()
.read_only(read_only)
.destructive(destructive)
.idempotent(
read_only || actions.iter().all(|a| a.operation.requires_idempotency_key),
)
.open_world(
generation || *name == "prepare_asset_upload" || *name == "find_assets",
),
);
SemanticTool { tool, actions }
})
.collect()
}
impl Action {
fn new(name: Option<&'static str>, operation_id: &str, openapi: &Value) -> Self {
// 按实际副作用声明;POST 也可能覆盖已有记录,参数名不能代表风险。
let destructive = match operation_id {
"listEditorProjects"
| "loadRecentEditorProject"
| "getEditorProject"
| "getEditorAssetLibrary"
| "getExternalAssetReadUrl"
| "getExternalEditorGenerationJob"
| "createEditorProject"
| "createEditorProjectResource"
| "createEditorAssetFolder"
| "createEditorAsset"
| "createExternalDirectUploadTicket" => false,
// 对象确认允许更新同一 owner 的已有对象元数据。
"confirmExternalAssetObject"
| "renameEditorProject"
| "saveEditorProjectCanvas"
| "updateEditorAssetFolder"
| "updateEditorAsset"
| "deleteEditorProject"
| "deleteEditorAssetFolder"
| "deleteEditorAsset" => true,
// 生成完成可修改已有画布状态,编辑与抠图还支持原位替换。
"generateExternalEditorImage"
| "editExternalEditorImage"
| "removeExternalEditorImageBackground"
| "generateExternalEditorIconSpritesheet"
| "extractExternalEditorUiDesignAssets"
| "generateExternalEditorCharacterAnimation"
| "generateExternalEditorVideo"
| "generateExternalEditorSoundEffect"
| "generateExternalEditorBackgroundMusic" => true,
_ => panic!("semantic operation must declare destructive risk: {operation_id}"),
};
let operation = MCP_OPERATIONS
.iter()
.find(|op| op.operation_id == operation_id)
.expect("semantic tools must map to existing operations");
let wrapped = &operation.input_schema["properties"];
// 保留 body 的 if/then/allOf 等约束;仅合并位置包装,不重建字段定义。
let mut input_schema = wrapped
.get("body")
.cloned()
.unwrap_or_else(|| json!({"type": "object", "properties": {}}));
let mut required = input_schema
.get("required")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default();
for location in ["pathParameters", "queryParameters"] {
if let Some(schema) = wrapped.get(location) {
for (name, field) in schema["properties"]
.as_object()
.expect("parameter properties")
{
assert!(
input_schema["properties"].get(name).is_none(),
"ambiguous field {name}"
);
input_schema["properties"][name] = inline_openapi_schema(openapi, field, 0);
}
required.extend(
schema
.get("required")
.and_then(Value::as_array)
.into_iter()
.flatten()
.cloned(),
);
}
}
let mut fixed_body = Map::new();
if name == Some("variation") {
input_schema["properties"]
.as_object_mut()
.unwrap()
.remove("kind");
required.retain(|field| field != "kind");
fixed_body.insert("kind".into(), json!("quick-edit"));
}
if operation_id == "confirmExternalAssetObject" {
input_schema["properties"]
.as_object_mut()
.unwrap()
.remove("ownerUserId");
}
input_schema["required"] = Value::Array(required);
input_schema["additionalProperties"] = json!(false);
let path = operation.path_template.split('?').next().unwrap();
let path_item = &openapi["paths"][path];
let rest = &path_item[operation.method.as_str().to_ascii_lowercase()];
let key_schema = path_item
.get("parameters")
.and_then(Value::as_array)
.into_iter()
.flatten()
.chain(
rest.get("parameters")
.and_then(Value::as_array)
.into_iter()
.flatten(),
)
.filter_map(|p| resolve_openapi_reference(openapi, p))
.find(|p| p["in"] == "header" && p["name"] == "Idempotency-Key")
.map(|p| {
let mut schema = inline_openapi_schema(openapi, &p["schema"], 0);
if let Some(description) = p.get("description") {
schema["description"] = description.clone();
}
schema
});
Self {
name,
operation,
input_schema,
key_schema,
fixed_body,
destructive,
}
}
fn call_schema(&self) -> Value {
let mut schema = match self.name {
Some(name) => json!({
"type": "object",
"properties": {"action": {"type": "string", "const": name}, "input": self.input_schema},
"required": ["action", "input"],
"additionalProperties": false
}),
None => self.input_schema.clone(),
};
if let Some(key) = &self.key_schema {
schema["properties"]["idempotencyKey"] = key.clone();
if self.operation.requires_idempotency_key {
schema["required"]
.as_array_mut()
.unwrap()
.push(json!("idempotencyKey"));
}
}
schema
}
}
fn tool_schema(actions: &[Action]) -> Value {
if actions[0].name.is_none() {
return actions[0].call_schema();
}
let mut schema = json!({
"type": "object",
"properties": {
"action": {"type": "string", "enum": actions.iter().map(|a| a.name.unwrap()).collect::<Vec<_>>()},
"input": {"type": "object"}
},
"required": ["action", "input"],
"additionalProperties": false,
"oneOf": actions.iter().map(Action::call_schema).collect::<Vec<_>>()
});
if let Some(key) = actions.iter().find_map(|a| a.key_schema.as_ref()) {
schema["properties"]["idempotencyKey"] = key.clone();
}
schema
}
impl SemanticTool {
pub(super) fn prepare(&self, mut arguments: Map<String, Value>) -> Result<PreparedCall, Value> {
let action = if self.actions[0].name.is_none() {
&self.actions[0]
} else {
let name = arguments
.get("action")
.and_then(Value::as_str)
.ok_or_else(|| json!({"error": "必须提供字符串 action"}))?;
self.actions
.iter()
.find(|a| a.name == Some(name))
.ok_or_else(|| json!({"error": "未知 action"}))?
};
validate_fields(&action.call_schema(), &arguments)?;
let key = arguments.remove("idempotencyKey");
let mut optional_idempotency_key = None;
if let Some(key) = &key {
let key = key
.as_str()
.ok_or_else(|| json!({"error": "idempotencyKey 必须是字符串"}))?;
if key.is_empty() || key.len() > 128 || !key.bytes().all(|c| (b'!'..=b'~').contains(&c))
{
return Err(json!({"error": "idempotencyKey 必须为 1–128 个非空白 ASCII 字符"}));
}
if !action.operation.requires_idempotency_key {
optional_idempotency_key = Some(
HeaderValue::from_str(key)
.map_err(|_| json!({"error": "idempotencyKey 不是合法 HTTP 头值"}))?,
);
}
}
let input = if action.name.is_some() {
arguments
.remove("input")
.and_then(|value| value.as_object().cloned())
.ok_or_else(|| json!({"error": "input 必须是 JSON 对象"}))?
} else {
arguments
};
validate_fields(&action.input_schema, &input)?;
let wrapped = &action.operation.input_schema["properties"];
let mut mapped = Map::new();
for location in ["pathParameters", "queryParameters", "body"] {
if let Some(schema) = wrapped.get(location) {
let mut fields = input
.iter()
.filter(|(name, _)| schema["properties"].get(*name).is_some())
.map(|(name, value)| (name.clone(), value.clone()))
.collect::<Map<_, _>>();
if location == "body" {
fields.extend(action.fixed_body.clone());
}
// 有请求体的操作始终发送对象,包括无字段的项目创建。
if location == "body" || !fields.is_empty() {
mapped.insert(location.into(), Value::Object(fields));
}
}
}
if action.operation.requires_idempotency_key {
if let Some(key) = key {
mapped.insert("idempotencyKey".into(), key);
}
}
Ok(PreparedCall {
operation: action.operation,
arguments: mapped,
optional_idempotency_key,
})
}
}
// 只校验适配层结构和直接字段,不实现第二套业务 schema 验证器。
// 嵌套字段与跨字段条件在现有 REST DTO/业务入口中校验,完整 schema 仍向客户端提供。
fn validate_fields(schema: &Value, input: &Map<String, Value>) -> Result<(), Value> {
let properties = schema["properties"]
.as_object()
.expect("input schema properties");
for name in schema["required"]
.as_array()
.into_iter()
.flatten()
.filter_map(Value::as_str)
{
if !input.contains_key(name) {
return Err(json!({"error": "缺少必填字段", "field": name}));
}
}
for (name, value) in input {
let field = properties
.get(name)
.ok_or_else(|| json!({"error": "当前操作不接受此字段", "field": name}))?;
let matches_type = |kind: &str| match kind {
"string" => value.is_string(),
"object" => value.is_object(),
"array" => value.is_array(),
"boolean" => value.is_boolean(),
"number" => value.is_number(),
"integer" => {
value.is_i64() || value.is_u64() || value.as_f64().is_some_and(|v| v.fract() == 0.0)
}
"null" => value.is_null(),
_ => true,
};
let valid_type = match &field["type"] {
Value::String(kind) => matches_type(kind),
Value::Array(kinds) => kinds.iter().filter_map(Value::as_str).any(matches_type),
_ => true,
};
if !valid_type
|| field
.get("enum")
.and_then(Value::as_array)
.is_some_and(|values| !values.contains(value))
|| field.get("const").is_some_and(|expected| expected != value)
{
return Err(json!({"error": "字段类型或取值不符合当前操作", "field": name}));
}
}
Ok(())
}
#[cfg(test)]
mod tests;
File diff suppressed because it is too large Load Diff
+110 -23
View File
@@ -37,18 +37,24 @@ mod model_catalog_tests {
use super::*;
#[test]
fn public_catalog_only_exposes_alias_and_stable_id() {
let mut catalog = module_runtime::AgcModelCatalog::default();
catalog.revision = 7;
fn public_catalog_exposes_stable_id_and_upstream_alias() {
let mut catalog = module_runtime::AgcModelCatalog::from_upstream_models(
vec!["gpt-5.6-sol".to_string(), "gpt-5.6-terra".to_string()],
7,
)
.expect("catalog should build");
catalog.models[1].enabled = false;
let payload = serde_json::to_value(public_model_catalog(catalog)).unwrap();
// 客户端拿到稳定标识 + 别名(别名就是上游原始模型名),实际模型名不下发。
assert_eq!(
payload["models"],
json!([{"id": "quality", "displayName": "高质量"}])
json!([{"id": "gpt-5-6-sol", "displayName": "gpt-5.6-sol"}])
);
assert_eq!(payload["defaultModelId"], "quality");
assert_eq!(payload["defaultModelId"], "gpt-5-6-sol");
assert_eq!(payload["revision"], json!(7));
assert!(!payload.to_string().contains("gpt-"));
assert!(payload.get("defaultModel").is_none());
assert!(payload["models"][0].get("enabled").is_none());
assert!(payload["models"][0].get("modelId").is_none());
}
}
@@ -194,6 +200,7 @@ fn public_model_catalog(catalog: module_runtime::AgcModelCatalog) -> LlmModelsRe
.filter(|model| model.enabled)
.map(|model| LlmModelSummary {
id: model.id,
// 初始目录里别名就是上游原始模型名(不再填“高质量/快速”这类人工别名)。
display_name: model.alias,
})
.collect(),
@@ -201,6 +208,29 @@ fn public_model_catalog(catalog: module_runtime::AgcModelCatalog) -> LlmModelsRe
}
}
/// 测试用目录:两项。上游模型名带 `.`,标识是它的 slug —— 既验证「客户端只回传目录标识」,
/// 也验证标识 → 实际模型名的映射;默认项是排序后的第一项,`TEST_AGC_MODEL_ID` 不是默认项。
#[cfg(test)]
pub(crate) const TEST_AGC_MODEL_ID: &str = "test-router-model";
#[cfg(test)]
pub(crate) const TEST_AGC_MODEL_MODEL_ID: &str = "test-router.model";
#[cfg(test)]
pub(crate) const TEST_AGC_MODEL_DEFAULT_ID: &str = "test-router-default";
#[cfg(test)]
pub(crate) const TEST_AGC_MODEL_DEFAULT_MODEL_ID: &str = "test-router.default";
#[cfg(test)]
pub(crate) fn test_agc_model_catalog() -> module_runtime::AgcModelCatalog {
module_runtime::AgcModelCatalog::from_upstream_models(
vec![
TEST_AGC_MODEL_DEFAULT_MODEL_ID.to_string(),
TEST_AGC_MODEL_MODEL_ID.to_string(),
],
0,
)
.expect("test catalog should build")
}
async fn load_llm_catalog(
state: &AppState,
owner: &str,
@@ -211,7 +241,7 @@ async fn load_llm_catalog(
.expect("fixture lock")
.contains_key(owner)
{
return Ok(module_runtime::AgcModelCatalog::default());
return Ok(test_agc_model_catalog());
}
let _ = owner;
crate::agc_models::load_catalog(state).await
@@ -283,9 +313,8 @@ pub async fn proxy_llm_responses(
] {
object.remove(field);
}
// The AGC client may select a model from the server-provided Router
// directory. Older callers without the reserved marker remain pinned to
// the official default model.
// AGC 客户端可以在服务端目录内选择模型;`model` 就是上游原始模型名。
// 老客户端存的历史稳定标识与目录外模型一律拒绝,不回退其它模型。
let agc_client = headers
.get("x-genarrative-client")
.and_then(|value| value.to_str().ok())
@@ -293,16 +322,13 @@ pub async fn proxy_llm_responses(
let catalog = load_llm_catalog(&state, authenticated.claims().user_id())
.await
.map_err(|error| llm_error_response(&request_context, error))?;
let selected_id = if agc_client {
requested_model
.as_deref()
.filter(|id| *id != "platform-default")
let requested_model = if agc_client {
requested_model.as_deref()
} else {
None
}
.unwrap_or(&catalog.default_model_id);
};
let selected_model = catalog
.resolve(selected_id)
.resolve_requested(requested_model)
.map_err(|message| {
llm_error_response(
&request_context,
@@ -847,7 +873,7 @@ async fn resolve_llm_router_client(
let catalog = load_llm_catalog(state, owner_user_id)
.await
.map_err(|_| "模型目录暂不可用".to_string())?;
let model = catalog.resolve(&catalog.default_model_id)?;
let model = catalog.resolve_requested(None)?;
let config = platform_llm::LlmConfig::new(
platform_llm::LlmProvider::OpenAiCompatible,
base_url.to_string(),
@@ -1304,11 +1330,14 @@ mod tests {
}
#[tokio::test]
async fn llm_responses_proxy_forces_official_model_and_keeps_router_key_server_side() {
async fn llm_responses_without_agc_marker_uses_catalog_default_and_keeps_router_key_server_side()
{
let (server_url, captured_request) = spawn_capturing_mock_server(MockResponse {
status_line: "200 OK",
content_type: "application/json; charset=utf-8",
body: r#"{"id":"resp_proxy_01","model":"gpt-6-astra","output":[]}"#.to_string(),
body: format!(
r#"{{"id":"resp_proxy_01","model":"{TEST_AGC_MODEL_DEFAULT_MODEL_ID}","output":[]}}"#
),
extra_headers: Vec::new(),
});
let (state, user_id) = seed_authenticated_state(AppConfig {
@@ -1373,12 +1402,64 @@ mod tests {
.expect("upstream request body");
let upstream_payload: Value =
serde_json::from_str(upstream_body).expect("upstream body should be json");
assert_eq!(upstream_payload["model"], "gpt-6-astra");
assert_eq!(upstream_payload["model"], TEST_AGC_MODEL_DEFAULT_MODEL_ID);
assert_ne!(upstream_payload["model"], "client-must-not-control");
}
#[tokio::test]
async fn llm_responses_rejects_upstream_names_and_unknown_catalog_ids() {
async fn llm_responses_forwards_catalog_model_selected_by_agc_client() {
let (server_url, captured_request) = spawn_capturing_mock_server(MockResponse {
status_line: "200 OK",
content_type: "application/json; charset=utf-8",
body: format!(
r#"{{"id":"resp_proxy_02","model":"{TEST_AGC_MODEL_MODEL_ID}","output":[]}}"#
),
extra_headers: Vec::new(),
});
let (state, user_id) = seed_authenticated_state(AppConfig {
llm_router_base_url: server_url.clone(),
llm_router_api_key_encryption_secret: Some("fixture-encryption-secret".to_string()),
..AppConfig::default()
})
.await;
install_test_provisioned_router_credential(&user_id, server_url, "fixture-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/responses")
.header("authorization", format!("Bearer {token}"))
.header("x-genarrative-client", "agc")
.header("content-type", "application/json")
.body(Body::from(
json!({"model": TEST_AGC_MODEL_ID, "input": "hello"}).to_string(),
))
.expect("request should build"),
)
.await
.expect("request should succeed");
assert_eq!(response.status(), StatusCode::OK);
let upstream_request = captured_request
.lock()
.expect("captured request lock")
.clone()
.expect("mock server should capture upstream request");
let (_, upstream_body) = upstream_request
.split_once("\r\n\r\n")
.expect("upstream request body");
let upstream_payload: Value =
serde_json::from_str(upstream_body).expect("upstream body should be json");
// 客户端只能回传目录标识,服务端映射成上游实际模型名;默认项不参与。
assert_eq!(upstream_payload["model"], TEST_AGC_MODEL_MODEL_ID);
assert_ne!(upstream_payload["model"], TEST_AGC_MODEL_DEFAULT_MODEL_ID);
}
#[tokio::test]
async fn llm_responses_rejects_models_outside_catalog() {
let (state, user_id) = seed_authenticated_state(AppConfig::default()).await;
install_test_provisioned_router_credential(
&user_id,
@@ -1387,7 +1468,13 @@ mod tests {
);
let token = issue_access_token(&state, &user_id);
let app = build_router(state);
for model in ["gpt-6-astra", "unlisted"] {
// 历史稳定标识、目录外名称、以及「直接拿上游实际模型名当标识」都必须拒绝。
for model in [
"quality",
"gpt-6-astra",
"unlisted",
TEST_AGC_MODEL_MODEL_ID,
] {
let response = app
.clone()
.oneshot(
+14
View File
@@ -500,6 +500,10 @@ fn should_initialize_editor_generation_pricing_for_startup(process_role: Process
process_role.runs_http()
}
fn should_initialize_agc_model_catalog_for_startup(process_role: ProcessRole) -> bool {
process_role.runs_http()
}
async fn run_http_role(config: AppConfig) -> Result<(), io::Error> {
let bind_address = config.bind_socket_addr();
let listen_backlog = config.listen_backlog;
@@ -764,6 +768,16 @@ async fn try_restore_app_state_for_startup(
))
})?;
}
// AGC 模型目录只来自上游同步或后台保存;这里同步失败不阻塞启动,由下一次启动重试,
// 未初始化期间 AGC 相关接口失败关闭。
if should_initialize_agc_model_catalog_for_startup(process_role) {
if let Err(error) = crate::agc_models::ensure_agc_model_catalog_initialized(&state).await {
error!(
error = %error,
"AGC 模型目录未初始化:本次启动未从上游同步到模型列表,AGC 目录与对话接口将失败关闭,下次启动会重试"
);
}
}
Ok(state)
}
@@ -19,12 +19,12 @@ use crate::{
delete_external_editor_project, edit_external_editor_image,
extract_external_editor_ui_design_assets, generate_external_editor_background_music,
generate_external_editor_character_animation, generate_external_editor_icon_spritesheet,
generate_external_editor_image, generate_external_editor_sound_effect,
generate_external_editor_video, get_external_editor_asset_library,
get_external_editor_generation_job, get_external_editor_project,
list_external_editor_projects, load_recent_external_editor_project, openapi_json,
remove_external_editor_image_background, rename_external_editor_project,
save_external_editor_canvas, update_external_editor_asset,
generate_external_editor_image, generate_external_editor_scene,
generate_external_editor_sound_effect, generate_external_editor_video,
get_external_editor_asset_library, get_external_editor_generation_job,
get_external_editor_project, list_external_editor_projects,
load_recent_external_editor_project, openapi_json, remove_external_editor_image_background,
rename_external_editor_project, save_external_editor_canvas, update_external_editor_asset,
update_external_editor_asset_folder,
},
external_mcp,
@@ -111,6 +111,10 @@ pub fn router(state: AppState) -> Router<AppState> {
"/api/external/v1/editor/images/generations",
post(generate_external_editor_image),
),
(
"/api/external/v1/editor/scenes/generations",
post(generate_external_editor_scene),
),
(
"/api/external/v1/editor/images/edits",
post(edit_external_editor_image),
@@ -257,6 +261,7 @@ mod route_contract_tests {
),
("/api/external/v1/generations/{operation_id}", &["GET"]),
("/api/external/v1/editor/images/generations", &["POST"]),
("/api/external/v1/editor/scenes/generations", &["POST"]),
("/api/external/v1/editor/images/edits", &["POST"]),
(
"/api/external/v1/editor/images/background-removals",
@@ -31,10 +31,12 @@ use shared_contracts::game_distribution::{
GameDistributionPublishMetadataSuggestion, GameDistributionPublishMetadataSuggestionRequest,
};
use spacetime_client::{
GameDistributionApproveRecordInput, GameDistributionCancelVersionRecordInput,
GameDistributionGameRecord, GameDistributionGetGameRecordInput,
GameDistributionPublicGameListRecordInput, GameDistributionPublicGameRecord,
GameDistributionRejectRecordInput, GameDistributionSubmitReviewRecordInput,
GameDistributionAdminGameListRecordInput, GameDistributionAdminGameRecord,
GameDistributionAdminVersionRecord, GameDistributionApproveRecordInput,
GameDistributionCancelVersionRecordInput, GameDistributionGameRecord,
GameDistributionGetGameRecordInput, GameDistributionPublicGameListRecordInput,
GameDistributionPublicGameRecord, GameDistributionRejectRecordInput,
GameDistributionRestoreRecordInput, GameDistributionSubmitReviewRecordInput,
GameDistributionSuspendRecordInput, GameDistributionUnpublishRecordInput,
GameDistributionVersionRecord, SpacetimeClientError,
};
@@ -60,6 +62,8 @@ pub(crate) const MAX_PACKAGE_CHUNK_REQUEST_BODY_BYTES: usize = PACKAGE_UPLOAD_CH
/// 分片偏移由客户端显式声明,服务端以对象当前长度为唯一权威。
const PACKAGE_UPLOAD_OFFSET_HEADER: &str = "x-genarrative-upload-offset";
const MAX_LIST_LIMIT: u32 = 48;
/// 后台游戏管理页全量列表上限,与 spacetime-module 的 admin game list limit 保持同口径。
const MAX_ADMIN_GAME_LIST_LIMIT: u32 = 200;
const MAX_IDEMPOTENCY_KEY_CHARS: usize = 128;
const MAX_PACKAGE_MANIFEST_JSON_BYTES: usize = 2 * 1024 * 1024;
/// 首版截图上限,与主规范冻结口径一致。
@@ -156,8 +160,6 @@ struct AdminReviewRequest {
expected_publication_revision: u64,
#[serde(default)]
review_reason: Option<String>,
#[serde(default)]
entry_url: Option<String>,
}
#[derive(Debug, Deserialize)]
@@ -176,6 +178,17 @@ struct AdminSuspendRequest {
reason: Option<String>,
}
#[derive(Debug, Deserialize)]
struct AdminGameListQuery {
limit: Option<u32>,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct AdminRestoreGameRequest {
expected_publication_revision: u64,
}
pub fn router(state: AppState) -> Router<AppState> {
let protected = Router::new()
.route(
@@ -243,10 +256,15 @@ pub fn router(state: AppState) -> Router<AppState> {
"/admin/api/game-distribution/versions/{version_id}",
get(admin_get_version),
)
.route("/admin/api/game-distribution/games", get(admin_list_games))
.route(
"/admin/api/game-distribution/games/{game_id}/suspend",
post(admin_suspend_game),
)
.route(
"/admin/api/game-distribution/games/{game_id}/restore",
post(admin_restore_game),
)
.route_layer(middleware::from_fn_with_state(
state.clone(),
require_admin_auth,
@@ -1395,6 +1413,37 @@ async fn admin_list_reviews(
))
}
async fn admin_list_games(
State(state): State<AppState>,
Extension(ctx): Extension<RequestContext>,
Extension(_admin): Extension<AuthenticatedAdmin>,
Query(query): Query<AdminGameListQuery>,
) -> Result<Json<Value>, AppError> {
let limit = query
.limit
.unwrap_or(MAX_ADMIN_GAME_LIST_LIMIT)
.min(MAX_ADMIN_GAME_LIST_LIMIT);
let games = state
.spacetime_client()
.list_admin_game_distribution_games(GameDistributionAdminGameListRecordInput { limit })
.await
.map_err(map_spacetime_error)?;
info!(
request_id = ctx.request_id(),
operation = "admin_games_listed",
games = games.len(),
limit,
elapsed_ms = ctx.elapsed(),
"后台读取全量发行游戏"
);
Ok(json_success_body(
Some(&ctx),
json!({
"games": games.iter().map(admin_game_payload).collect::<Vec<_>>(),
}),
))
}
async fn admin_review_version(
State(state): State<AppState>,
Extension(ctx): Extension<RequestContext>,
@@ -1416,25 +1465,21 @@ async fn admin_review_version(
decision.as_str(),
payload.expected_publication_revision,
payload.review_reason.as_deref(),
payload.entry_url.as_deref(),
))
.map_err(|error| internal(error.to_string()))?,
);
let (version, replayed) = if decision == "approve" {
// 回滚窗口里“关闭新版本激活”,但拒绝审核与安全下架必须始终可用。
ensure_publish_enabled(&state, None).await?;
let entry_url = payload
.entry_url
.as_deref()
.ok_or_else(|| bad_request("审核通过必须提供发行网关 HTTPS 入口"))?;
validate_release_entry_url(entry_url, !state.config.is_production())?;
// 发行入口由部署模板和 gameId 派生,管理员不填地址,也不做二次确认。
let entry_url = derive_release_entry_url(&state, &version_id).await?;
state
.spacetime_client()
.approve_game_distribution_version(GameDistributionApproveRecordInput {
version_id,
admin_user_id,
expected_publication_revision: payload.expected_publication_revision,
entry_url: entry_url.to_string(),
entry_url,
idempotency_key,
request_digest,
now_micros: now_micros(),
@@ -1506,34 +1551,30 @@ async fn admin_get_version(
))
}
/// 校验管理员提交的发行入口。
/// 审核通过时派生的发行入口:平台同源路径 `/games/{gameId}/`。
///
/// 生产环境只接受绝对 HTTPS 地址;非生产环境额外允许 http 回环地址,口径与前端
/// `normalizeGameEntryUrl` 一致,便于本地把发行网关跑在 127.0.0.1 上验证内嵌游玩。
/// 任何环境都拒绝凭据、query 和 fragment,也不允许服务端自行拼默认地址。
fn validate_release_entry_url(value: &str, allow_loopback_http: bool) -> Result<(), AppError> {
let parsed =
url::Url::parse(value.trim()).map_err(|_| bad_request("发行入口必须是有效 URL"))?;
let host = parsed.host_str();
let scheme_allowed = parsed.scheme() == "https"
|| (allow_loopback_http
&& parsed.scheme() == "http"
&& matches!(
host,
Some("127.0.0.1") | Some("localhost") | Some("[::1]") | Some("::1")
));
if !scheme_allowed
|| host.is_none()
|| parsed.username() != ""
|| parsed.password().is_some()
|| parsed.query().is_some()
|| parsed.fragment().is_some()
/// 存相对路径而不是绝对 URL,部署侧就不需要提供发行域名;dev / release / 预览环境
/// 口径一致,由客户端按当前 origin 解析成绝对地址后再交给 iframe。
async fn derive_release_entry_url(state: &AppState, version_id: &str) -> Result<String, AppError> {
let version = state
.spacetime_client()
.get_game_distribution_version(version_id.to_string())
.await
.map_err(map_spacetime_error)?
.ok_or_else(|| AppError::from_status(StatusCode::NOT_FOUND))?;
build_release_entry_url(&version.game_id)
}
/// 发行入口固定走平台同源路径,游戏标识必须能安全落在路径段里。
fn build_release_entry_url(game_id: &str) -> Result<String, AppError> {
if game_id.is_empty()
|| !game_id.chars().all(|character| {
character.is_ascii_alphanumeric() || character == '-' || character == '_'
})
{
return Err(bad_request(
"发行入口必须是无凭据、无查询参数的 HTTPS URL;仅非生产环境允许回环 http",
));
return Err(internal("游戏标识不适用于发行路径"));
}
Ok(())
Ok(format!("/games/{game_id}/"))
}
async fn admin_suspend_game(
@@ -1597,6 +1638,51 @@ async fn admin_suspend_game(
))
}
async fn admin_restore_game(
State(state): State<AppState>,
Extension(ctx): Extension<RequestContext>,
Extension(admin): Extension<AuthenticatedAdmin>,
headers: HeaderMap,
Path(game_id): Path<String>,
Json(payload): Json<AdminRestoreGameRequest>,
) -> Result<Json<Value>, AppError> {
let idempotency_key = idempotency_key(&headers)?;
let admin_user_id = admin.session().subject.clone();
let request_digest = compute_request_digest(
&serde_json::to_vec(&(game_id.as_str(), payload.expected_publication_revision))
.map_err(|error| internal(error.to_string()))?,
);
let log_game_id = game_id.clone();
let log_admin_user_id = admin_user_id.clone();
let game = state
.spacetime_client()
.restore_game_distribution_game(GameDistributionRestoreRecordInput {
game_id,
admin_user_id,
expected_publication_revision: payload.expected_publication_revision,
idempotency_key,
request_digest,
now_micros: now_micros(),
})
.await
.map_err(map_spacetime_error)?;
info!(
request_id = ctx.request_id(),
operation = "game_restored",
game_id = %log_game_id,
admin_user_id = %log_admin_user_id,
publication_revision = game.0.publication_revision,
visibility = %game.0.visibility,
replayed = game.1,
elapsed_ms = ctx.elapsed(),
"管理员恢复已下架游戏"
);
Ok(json_success_body(
Some(&ctx),
json!({ "game": game_payload(&game.0), "replayed": game.1 }),
))
}
async fn record_upload_failure(
state: &AppState,
owner_user_id: &str,
@@ -1766,6 +1852,51 @@ fn public_game_payload(game: GameDistributionPublicGameRecord) -> Value {
payload
}
/// 后台游戏管理页的游戏行:作者名/头像由 spacetime 事务内读时联账号表得到。
fn admin_game_payload(game: &GameDistributionAdminGameRecord) -> Value {
json!({
"gameId": game.game_id,
"title": game.title,
"author": {
"id": game.owner_user_id,
"name": game.author_name.as_deref().unwrap_or("未知作者"),
"avatarUrl": game.author_avatar_url,
},
"status": game.visibility,
"versionCount": game.version_count,
"playCount": game.play_count,
"activeVersionId": game.active_version_id,
"publicationRevision": game.publication_revision,
"createdAt": game.created_at,
"updatedAt": game.updated_at,
"versions": game
.versions
.iter()
.map(|version| admin_game_version_payload(&game.game_id, version))
.collect::<Vec<_>>(),
})
}
fn admin_game_version_payload(
game_id: &str,
version: &GameDistributionAdminVersionRecord,
) -> Value {
json!({
"versionId": version.version_id,
"gameId": game_id,
"versionNumber": version.version_number,
"status": version.status,
"reviewReason": version.review_reason,
"packageBytes": version.package_bytes,
"packageSha256": version.package_sha256,
"entryUrl": version.entry_url,
"createdAt": version.created_at,
"updatedAt": version.updated_at,
"reviewedAt": version.reviewed_at,
"publishedAt": version.published_at,
})
}
fn game_payload(game: &GameDistributionGameRecord) -> Value {
let tags = serde_json::from_str::<Vec<String>>(&game.tags_json).unwrap_or_default();
let screenshots = game
@@ -2713,6 +2844,51 @@ mod tests {
assert_eq!(recovery_action_for_status("unknown_status"), "none");
}
#[tokio::test]
async fn admin_game_management_routes_are_mounted() {
use axum::{body::Body, http::Request};
use tower::ServiceExt;
let app = crate::app::build_router(
crate::state::AppState::new(crate::config::AppConfig::default())
.expect("测试状态应可构建"),
);
// 全量列表与恢复都必须先过管理员鉴权,未带 token 时在进入业务前被拒。
let unauthenticated_list = app
.clone()
.oneshot(
Request::builder()
.uri("/admin/api/game-distribution/games")
.body(Body::empty())
.expect("请求"),
)
.await
.expect("路由响应");
// 测试态没有启用后台运行时,鉴权中间件会在 503 处失败关闭;关键是不能 404。
assert!(matches!(
unauthenticated_list.status(),
StatusCode::UNAUTHORIZED | StatusCode::SERVICE_UNAVAILABLE
));
let unauthenticated_restore = app
.oneshot(
Request::builder()
.method("POST")
.uri("/admin/api/game-distribution/games/game_1/restore")
.header("content-type", "application/json")
.header("Idempotency-Key", "restore-1")
.body(Body::from(r#"{"expectedPublicationRevision":1}"#))
.expect("请求"),
)
.await
.expect("路由响应");
assert!(matches!(
unauthenticated_restore.status(),
StatusCode::UNAUTHORIZED | StatusCode::SERVICE_UNAVAILABLE
));
}
#[tokio::test]
async fn version_readback_and_cancel_routes_are_mounted() {
use axum::{body::Body, http::Request};
@@ -2803,51 +2979,19 @@ mod tests {
}
#[test]
fn approve_requires_credential_free_https_entry_url() {
validate_release_entry_url(
"https://games.example.test/releases/game_1/index.html",
false,
)
.expect("发行入口");
for invalid in [
"/releases/game_1/index.html",
"http://games.example.test/releases/game_1/index.html",
"http://127.0.0.1:10001/releases/game_1/index.html",
"https://user:pass@games.example.test/index.html",
"https://games.example.test/index.html?token=1",
"https://games.example.test/index.html#x",
] {
assert_eq!(
validate_release_entry_url(invalid, false)
.expect_err("生产环境非法发行入口应被拒绝")
.status_code(),
StatusCode::BAD_REQUEST,
"未拒绝的发行入口:{invalid}"
);
}
fn release_entry_url_is_same_origin_path_with_game_id() {
assert_eq!(
build_release_entry_url("game_1").expect("派生发行入口"),
"/games/game_1/"
);
}
#[test]
fn non_production_release_entry_allows_loopback_http_only() {
for allowed in [
"http://127.0.0.1:10001/api/game-distribution/releases/game_1/index.html",
"http://localhost:10001/api/game-distribution/releases/game_1/index.html",
"https://games.example.test/releases/game_1/index.html",
] {
validate_release_entry_url(allowed, true).expect("非生产环境应接受回环 http");
}
for invalid in [
"http://games.example.test/releases/game_1/index.html",
"http://192.168.1.10:10001/index.html",
"http://127.0.0.1:10001/index.html?token=1",
"http://user:pass@127.0.0.1:10001/index.html",
] {
assert_eq!(
validate_release_entry_url(invalid, true)
.expect_err("非生产环境也不能放宽回环之外的地址")
.status_code(),
StatusCode::BAD_REQUEST,
"未拒绝的发行入口:{invalid}"
fn release_entry_rejects_game_id_that_is_not_path_safe() {
for invalid in ["", "../escape", "game/1", "game 1"] {
assert!(
build_release_entry_url(invalid).is_err(),
"未拒绝的游戏标识:{invalid}"
);
}
}