use std::{
collections::{BTreeMap, HashMap, VecDeque},
io::Write,
sync::{LazyLock, Mutex},
time::{Instant, SystemTime, UNIX_EPOCH},
};
use axum::{
Json, Router,
body::{Body, Bytes},
extract::{
DefaultBodyLimit, Extension, Path, Query, Request, State,
rejection::{JsonRejection, QueryRejection},
},
http::{HeaderMap, HeaderValue, StatusCode, header},
middleware::{self, Next},
response::Response,
routing::{get, post, put},
};
use flate2::{Compression as GzipCompression, write::GzEncoder};
use module_game_distribution::{
GAME_DISTRIBUTION_COLLECTION_PAGE_LIMIT_DEFAULT, GAME_DISTRIBUTION_COLLECTION_PAGE_LIMIT_MAX,
MAX_PACKAGE_BYTES, MAX_PROJECT_BUNDLE_BYTES, ProjectBundleError, ProjectBundleManifest,
ReleaseAssetError, ReleasePackageError, ReleasePackageManifest, compute_request_digest,
extract_release_asset, game_distribution_collection_page_limit, normalize_review_comment,
normalize_review_moderation_reason, release_asset_content_type, validate_project_bundle_zip,
validate_release_zip, validate_review_list_status,
};
use platform_auth::read_refresh_session_token;
use platform_llm::{EDITOR_AGENT_GPT5_MODEL, LlmMessage, LlmRunRequest};
use platform_oss::{
OssAppendInternalObjectRequest, OssDeleteObjectRequest, OssGetObjectRequest,
OssInternalPutObjectRequest, OssObjectAccess,
};
use serde::Deserialize;
use serde_json::{Value, json};
use sha2::{Digest, Sha256};
use shared_contracts::admin::{
AdminGameReview, AdminGameReviewDetailResponse, AdminGameReviewGame, AdminGameReviewGameInfo,
AdminGameReviewGamesQuery, AdminGameReviewGamesResponse, AdminGameReviewModerationRequest,
AdminGameReviewModerationResponse, AdminGameReviewOperation, AdminGameReviewsQuery,
AdminGameReviewsResponse,
};
use shared_contracts::game_distribution::{
GAME_DISTRIBUTION_CATEGORIES, GAME_DISTRIBUTION_VERSION_NUMBER_CONFLICT,
GameDistributionAuthor, GameDistributionCollectionState, GameDistributionCreateGameRequest,
GameDistributionCreateVersionRequest, GameDistributionDerivedResponse,
GameDistributionForkAuthorization, GameDistributionForkSource, GameDistributionForkSourceKind,
GameDistributionForkSourceResponse, GameDistributionInputMode, GameDistributionLineageNode,
GameDistributionLineageResponse, GameDistributionMyReviewResponse, GameDistributionOrientation,
GameDistributionPublishMetadataSuggestion, GameDistributionPublishMetadataSuggestionRequest,
GameDistributionRatingSummary, GameDistributionReview, GameDistributionReviewsResponse,
GameDistributionSaveReviewRequest, GameDistributionSaveReviewResponse,
GameDistributionSetForkAuthorizationRequest, GameDistributionUpdateGameMetadataRequest,
GameDistributionVisibility,
};
use spacetime_client::{
GameDistributionAdminGameListRecordInput, GameDistributionAdminGameRecord,
GameDistributionAdminUserReviewListRecordInput, GameDistributionAdminUserReviewRecord,
GameDistributionAdminVersionRecord, GameDistributionApproveRecordInput,
GameDistributionCancelVersionRecordInput, GameDistributionCollectGameRecordInput,
GameDistributionDeleteGameRecordInput, GameDistributionDerivedGamesRecord,
GameDistributionForkSourceRecord, GameDistributionGameRecord,
GameDistributionGetGameRecordInput, GameDistributionLineageNodeRecord,
GameDistributionLineageTreeRecord, GameDistributionOwnerGameRecord,
GameDistributionPublicGameListRecordInput, GameDistributionPublicGameRecord,
GameDistributionRatingSummaryRecord, GameDistributionRejectRecordInput,
GameDistributionRestoreRecordInput, GameDistributionReviewGameListRecordInput,
GameDistributionReviewModerationOperationRecord, GameDistributionReviewModerationRecordInput,
GameDistributionSetForkAuthorizationRecordInput, GameDistributionSubmitReviewRecordInput,
GameDistributionSuspendRecordInput, GameDistributionUncollectGameRecordInput,
GameDistributionUnpublishRecordInput, GameDistributionUpdateMetadataRecordInput,
GameDistributionUserReviewRecord, GameDistributionVersionRecord, SpacetimeClientError,
};
use tracing::{debug, info, warn};
use uuid::Uuid;
use crate::{
admin::{AuthenticatedAdmin, require_admin_auth},
api_response::json_success_body,
auth::{AuthenticatedAccessToken, optional_access_token_from_headers, require_bearer_auth},
game_play_counter::{GamePlayOutcome, GamePlayReport},
http_error::AppError,
platform_errors::{map_llm_error, map_oss_error},
request_context::{RequestContext, client_ip_from_headers},
state::AppState,
tracking::{TrackingEventDraft, record_tracking_event_after_success},
};
pub(crate) const MAX_PACKAGE_REQUEST_BODY_BYTES: usize = MAX_PACKAGE_BYTES as usize + 1024;
/// 分片续传的固定分片大小:200 MiB 上限下最多 25 片,单片远低于反代放行量。
/// 客户端只能使用服务端下发的值,不得自行改变分片边界,否则权威偏移会立刻对不上。
pub(crate) const PACKAGE_UPLOAD_CHUNK_BYTES: usize = 8 * 1024 * 1024;
/// 分片路由的请求体放行量:分片大小 + 1 KiB 头部余量。
pub(crate) const MAX_PACKAGE_CHUNK_REQUEST_BODY_BYTES: usize = PACKAGE_UPLOAD_CHUNK_BYTES + 1024;
/// 工程源包整包 PUT 的请求体放行量:上限与发行包同值(两者共用同一条上传链路,反代与
/// Pingora 的放行量就是按 200 MiB 校准的),只多留 1 KiB 头部余量。
pub(crate) const MAX_PROJECT_BUNDLE_REQUEST_BODY_BYTES: usize =
MAX_PROJECT_BUNDLE_BYTES as usize + 1024;
/// 分片偏移由客户端显式声明,服务端以对象当前长度为唯一权威。
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;
/// 后台作品状态过滤的白名单;`deleted` 表示已软删除的作品。
const ADMIN_GAME_LIST_STATUSES: [&str; 4] = ["published", "unpublished", "suspended", "deleted"];
const MAX_IDEMPOTENCY_KEY_CHARS: usize = 128;
const MAX_PACKAGE_MANIFEST_JSON_BYTES: usize = 2 * 1024 * 1024;
/// 首版截图上限,与主规范冻结口径一致。
const MAX_GAME_SCREENSHOTS: usize = 6;
const GAME_DISTRIBUTION_OBJECT_PREFIX: &str = "agc/project-snapshots/v1/game-distribution/";
const GAME_DISTRIBUTION_PUBLISHED_STATUS: &str = "published";
/// 发行包 PUT 的尝试次数与退避,口径与 `platform-oss` 的可重试分类一致。
const GAME_DISTRIBUTION_OSS_PUT_MAX_ATTEMPTS: usize = 3;
const GAME_DISTRIBUTION_OSS_PUT_RETRY_DELAYS_MS: [u64; 2] = [250, 500];
const RELEASE_PACKAGE_CACHE_MAX_ENTRIES: usize = 4;
/// 缓存字节预算必须比单个发行包上限大出一档,否则 200 MiB 档的包只能刚好自占整份预算,
/// 任何并发的小包都会被立刻挤掉。
const RELEASE_PACKAGE_CACHE_MAX_BYTES: usize = 256 * 1024 * 1024;
/// 发行包运行在 `sandbox="allow-scripts"` 的 opaque origin 中,浏览器原生 storage
/// 会抛 `SecurityError`。不授予 `allow-same-origin`(否则同源脚本可能移除 sandbox),
/// 而是在游戏脚本前安装本次运行期的同步兼容存储;它不接触平台 Cookie、DOM 或账号数据。
const RELEASE_STORAGE_BOOTSTRAP: &str = concat!(
"",
);
/// 发行静态资源的进程内缓存。
///
/// 单个资源取自整个 ZIP,若每个请求都重新下载整包会拖垮发行网关;缓存只保存已通过
/// 校验的私有包字节,键是对象键,超出条目或字节预算时按插入顺序淘汰。
/// 值用 `Bytes` 而不是 `Vec`:整包下发(M2a 取件通道)要直接把缓存里的字节交给响应体,
/// `Bytes::clone` 只加引用计数,不会为了一个 200 MiB 的包再复制一份。
static RELEASE_PACKAGE_CACHE: LazyLock> =
LazyLock::new(|| Mutex::new(ReleasePackageCache::default()));
#[derive(Default)]
struct ReleasePackageCache {
packages: HashMap,
order: VecDeque,
total_bytes: usize,
}
impl ReleasePackageCache {
fn get(&self, object_key: &str) -> Option {
self.packages.get(object_key).cloned()
}
fn insert(&mut self, object_key: String, bytes: Bytes) {
self.insert_with_limits(
object_key,
bytes,
RELEASE_PACKAGE_CACHE_MAX_ENTRIES,
RELEASE_PACKAGE_CACHE_MAX_BYTES,
);
}
fn insert_with_limits(
&mut self,
object_key: String,
bytes: Bytes,
max_entries: usize,
max_bytes: usize,
) {
if self.packages.contains_key(&object_key) {
return;
}
// 单个包超过缓存预算时直接不缓存,避免一次插入把整个进程内存顶满。
if bytes.len() > max_bytes {
return;
}
while self.order.len() >= max_entries
|| self.total_bytes.saturating_add(bytes.len()) > max_bytes
{
let Some(evicted) = self.order.pop_front() else {
break;
};
if let Some(previous) = self.packages.remove(&evicted) {
self.total_bytes = self.total_bytes.saturating_sub(previous.len());
}
}
self.total_bytes = self.total_bytes.saturating_add(bytes.len());
self.order.push_back(object_key.clone());
self.packages.insert(object_key, bytes);
}
}
#[derive(Debug, Deserialize)]
struct GameListQuery {
#[serde(alias = "keyword")]
search: Option,
category: Option,
#[serde(rename = "authorId")]
author_id: Option,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct UserReviewListQuery {
page: Option,
page_size: Option,
}
impl UserReviewListQuery {
fn pagination(self) -> Result<(u32, u32), AppError> {
let page = self.page.unwrap_or(1);
let page_size = self.page_size.unwrap_or(20);
if page == 0 || !(1..=50).contains(&page_size) {
return Err(AppError::from_status(StatusCode::BAD_REQUEST)
.with_message("页码须从 1 开始,每页条数须为 1–50"));
}
Ok((page, page_size))
}
}
#[derive(Debug, Deserialize)]
struct AdminReviewListQuery {
limit: Option,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct PublicationRevisionRequest {
expected_publication_revision: u64,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct AdminReviewRequest {
decision: String,
expected_publication_revision: u64,
#[serde(default)]
review_reason: Option,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct CancelVersionRequest {
expected_publication_revision: u64,
#[serde(default)]
reason: Option,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct AdminSuspendRequest {
expected_publication_revision: u64,
#[serde(default)]
reason: Option,
}
#[derive(Debug, Deserialize)]
struct AdminGameListQuery {
limit: Option,
/// 关键词:匹配标题、gameId 或作者 user ID。
keyword: Option,
#[serde(alias = "ownerUserId")]
owner: Option,
/// `published` / `unpublished` / `suspended` / `deleted`;为空时排除已软删除游戏。
status: Option,
cursor: Option,
}
/// 「我的收藏(收录)」列表的查询串:`limit` + `cursor`。
///
/// 分页口径与后台列表**同构但数值不同**:默认 20(网格一屏)、上限 50(约束单响应体积)。
/// 两处口径不必相等——后台是运营表格视图,需要一次扫读更多行;用户态网格按屏取数。
/// 数值本身只在 `module_game_distribution` 里定义一次,这里引用常量而不是再抄一遍。
#[derive(Debug, Deserialize)]
struct MyCollectionsQuery {
limit: Option,
cursor: Option,
}
/// 作者软删除游戏:CAS 修订号走查询串,删除本身没有请求体。
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct DeleteOwnerGameQuery {
#[serde(alias = "expected_publication_revision")]
expected_publication_revision: u64,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct AdminRestoreGameRequest {
expected_publication_revision: u64,
}
pub fn router(state: AppState) -> Router {
let admin_user_reviews = Router::new()
.route(
"/admin/api/game-distribution/user-review-games",
get(admin_review_games),
)
.route(
"/admin/api/game-distribution/user-reviews",
get(admin_user_review_list),
)
.route(
"/admin/api/game-distribution/user-reviews/{review_id}",
get(admin_user_review_detail),
)
.route(
"/admin/api/game-distribution/user-reviews/{review_id}/moderation",
post(admin_moderate_user_review),
)
.route_layer(middleware::from_fn_with_state(
state.clone(),
require_admin_auth,
))
.route_layer(middleware::from_fn(add_no_store_response_headers));
let user_reviews = Router::new()
.route(
"/api/game-distribution/games/{game_id}/my-review",
get(get_my_review).put(save_my_review),
)
.route_layer(middleware::from_fn_with_state(
state.clone(),
require_bearer_auth,
))
.route(
"/api/game-distribution/games/{game_id}/reviews",
get(list_user_reviews),
)
.route_layer(middleware::from_fn(add_no_store_response_headers));
// Fork 取件通道(M2a/M2b):要求 Bearer 登录,但**不叠加发布灰度**——灰度只针对「发布」,
// 任何登录用户都应该能改编已授权的作品。三个 handler 共用同一条校验,规则只写一遍;
// `/project` 失败关闭:只有选定资产确为工程源包时才服务,绝不悄悄回落成品包。
let fork_sources = Router::new()
.route(
"/api/game-distribution/games/{game_id}/fork-source",
get(get_fork_source),
)
.route(
"/api/game-distribution/games/{game_id}/fork-source/package",
get(get_fork_source_package),
)
.route(
"/api/game-distribution/games/{game_id}/fork-source/project",
get(get_fork_source_project),
)
.route_layer(middleware::from_fn_with_state(
state.clone(),
require_bearer_auth,
))
.route_layer(middleware::from_fn(add_no_store_response_headers));
// 收藏(收录):登录用户的真实用户态投影,绝不虚构收藏状态。
// - PUT 需要 `Idempotency-Key`(重放与「重复收藏」要靠收据区分);
// - DELETE 按确定性主键 `{userId}:{gameId}` 删除,天然幂等,所以**不要求**幂等键;
// - 读路径是 per-user 数据,整组 `no-store`,不给任何共享缓存留缝。
let collections = Router::new()
.route(
"/api/game-distribution/games/{game_id}/collection",
put(collect_game).delete(uncollect_game),
)
.route(
"/api/game-distribution/my-collections",
get(list_my_collections),
)
.route_layer(middleware::from_fn_with_state(
state.clone(),
require_bearer_auth,
))
.route_layer(middleware::from_fn(add_no_store_response_headers));
let protected = Router::new()
.route(
"/api/game-distribution/publish-metadata/suggestions",
post(suggest_publish_metadata),
)
.route("/api/game-distribution/games", post(create_game))
.route(
"/api/game-distribution/games/{game_id}/versions",
post(create_version),
)
.route(
"/api/game-distribution/versions/{version_id}/package",
put(upload_package).layer(DefaultBodyLimit::max(MAX_PACKAGE_REQUEST_BODY_BYTES)),
)
.route(
"/api/game-distribution/versions/{version_id}/package/upload-state",
get(package_upload_state),
)
.route(
"/api/game-distribution/versions/{version_id}/package/chunk",
put(upload_package_chunk)
.layer(DefaultBodyLimit::max(MAX_PACKAGE_CHUNK_REQUEST_BODY_BYTES)),
)
.route(
"/api/game-distribution/versions/{version_id}/package/complete",
post(complete_package_upload),
)
.route(
"/api/game-distribution/versions/{version_id}/package/reset",
post(reset_package_upload),
)
// 工程源包上行族(M2b):与发行包族逐条对齐,只是资产换成作者的工程源包。
// 载体类型一律 `application/octet-stream`(见技术方案 §3.4),分片边界与偏移头
// 直接复用发行包那一套——两者上限同值,客户端只能有一套偏移语义。
.route(
"/api/game-distribution/versions/{version_id}/project-bundle",
put(upload_project_bundle)
.layer(DefaultBodyLimit::max(MAX_PROJECT_BUNDLE_REQUEST_BODY_BYTES)),
)
.route(
"/api/game-distribution/versions/{version_id}/project-bundle/upload-state",
get(project_bundle_upload_state),
)
.route(
"/api/game-distribution/versions/{version_id}/project-bundle/chunk",
put(upload_project_bundle_chunk)
.layer(DefaultBodyLimit::max(MAX_PACKAGE_CHUNK_REQUEST_BODY_BYTES)),
)
.route(
"/api/game-distribution/versions/{version_id}/project-bundle/complete",
post(complete_project_bundle_upload),
)
.route(
"/api/game-distribution/versions/{version_id}/project-bundle/reset",
post(reset_project_bundle_upload),
)
.route(
"/api/game-distribution/versions/{version_id}/submit",
post(submit_version),
)
.route(
"/api/game-distribution/versions/{version_id}",
get(get_owner_version),
)
.route(
"/api/game-distribution/versions/{version_id}/cancel",
post(cancel_version),
)
.route("/api/game-distribution/my-games", get(list_my_games))
.route(
"/api/game-distribution/my-games/{game_id}",
get(get_owner_game)
.patch(update_owner_game_metadata)
.delete(delete_owner_game),
)
.route(
"/api/game-distribution/games/{game_id}/unpublish",
post(unpublish_game),
)
.route(
"/api/game-distribution/games/{game_id}/fork-authorization",
put(set_fork_authorization),
)
.route_layer(middleware::from_fn_with_state(
state.clone(),
require_bearer_auth,
));
let admin = Router::new()
.route(
"/admin/api/game-distribution/reviews",
get(admin_list_reviews),
)
.route(
"/admin/api/game-distribution/versions/{version_id}/review",
post(admin_review_version),
)
.route(
"/admin/api/game-distribution/versions/{version_id}",
get(admin_get_version),
)
.route(
"/admin/api/game-distribution/versions/{version_id}/preview-session",
post(admin_create_version_preview_session),
)
.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,
));
let public_games = Router::new()
.route("/api/game-distribution/games", get(list_games))
.route("/api/game-distribution/games/{game_id}", get(get_game))
.route(
"/api/game-distribution/games/{game_id}/lineage",
get(get_game_lineage),
)
.route(
"/api/game-distribution/games/{game_id}/derived",
get(get_game_derived_games),
)
.route(
"/api/game-distribution/games/{game_id}/plays",
post(record_game_play),
)
.route_layer(middleware::from_fn(add_no_store_response_headers));
Router::new()
.route(
"/api/game-distribution/releases/{game_id}/{*asset_path}",
get(serve_release_asset),
)
// 根路径等价于入口页:生产由发行来源(每游戏 origin)把 `/` 映射到 index.html,
// 本地直连网关或入口直接填网关地址时也必须能打开游戏。
.route(
"/api/game-distribution/releases/{game_id}",
get(serve_release_entry),
)
.route(
"/api/game-distribution/releases/{game_id}/",
get(serve_release_entry),
)
.route(
"/api/game-distribution/admin-previews/{preview_token}/{*asset_path}",
get(serve_admin_version_preview_asset),
)
.route(
"/api/game-distribution/admin-previews/{preview_token}",
get(serve_admin_version_preview_entry),
)
.route(
"/api/game-distribution/admin-previews/{preview_token}/",
get(serve_admin_version_preview_entry),
)
.merge(public_games)
.merge(protected)
.merge(user_reviews)
.merge(fork_sources)
.merge(collections)
.merge(admin_user_reviews)
.merge(admin)
}
async fn add_no_store_response_headers(request: Request, next: Next) -> Response {
let mut response = next.run(request).await;
response
.headers_mut()
.insert(header::CACHE_CONTROL, HeaderValue::from_static("no-store"));
response
}
async fn list_user_reviews(
State(state): State,
Extension(ctx): Extension,
Path(game_id): Path,
query: Result, QueryRejection>,
) -> Result, AppError> {
let Query(query) = query.map_err(|_| {
AppError::from_status(StatusCode::BAD_REQUEST).with_message("分页参数须为整数")
})?;
let (page, page_size) = query.pagination()?;
let result = state
.spacetime_client()
.list_game_distribution_user_reviews(game_id, page, page_size)
.await
.map_err(map_spacetime_error)?;
Ok(json_success_body(
Some(&ctx),
GameDistributionReviewsResponse {
reviews: result
.reviews
.into_iter()
.map(user_review_payload)
.collect::, _>>()?,
page: result.page,
page_size: result.page_size,
total: result.total,
total_pages: result.total_pages,
rating_summary: rating_summary_payload(result.rating_summary),
},
))
}
async fn get_my_review(
State(state): State,
Extension(ctx): Extension,
Extension(authenticated): Extension,
Path(game_id): Path,
) -> Result, AppError> {
let review = state
.spacetime_client()
.get_game_distribution_my_review(game_id, authenticated.claims().user_id().to_string())
.await
.map_err(map_spacetime_error)?;
Ok(json_success_body(
Some(&ctx),
GameDistributionMyReviewResponse {
review: review.map(user_review_payload).transpose()?,
},
))
}
async fn save_my_review(
State(state): State,
Extension(ctx): Extension,
Extension(authenticated): Extension,
Path(game_id): Path,
payload: Result, JsonRejection>,
) -> Result, AppError> {
let Json(payload) = payload.map_err(|_| {
AppError::from_status(StatusCode::BAD_REQUEST)
.with_message("评价请求须包含整数评分与可选文字评论")
})?;
let comment = normalize_review_comment(payload.score, &payload.comment).map_err(|error| {
AppError::from_status(StatusCode::UNPROCESSABLE_ENTITY).with_message(error.to_string())
})?;
let result = state
.spacetime_client()
.save_game_distribution_my_review(
game_id,
authenticated.claims().user_id().to_string(),
payload.score,
comment,
)
.await
.map_err(map_spacetime_error)?;
Ok(json_success_body(
Some(&ctx),
GameDistributionSaveReviewResponse {
review: user_review_payload(result.review)?,
rating_summary: rating_summary_payload(result.rating_summary),
},
))
}
fn user_review_payload(
review: GameDistributionUserReviewRecord,
) -> Result {
Ok(GameDistributionReview {
id: review.review_id,
game_id: review.game_id,
author: GameDistributionAuthor {
id: review.author_id,
name: review.author_name,
avatar_url: review.author_avatar_url,
},
score: review.score,
comment: review.comment,
is_hidden: review.is_hidden,
created_at: user_review_timestamp(review.created_at_micros)?,
updated_at: user_review_timestamp(review.updated_at_micros)?,
})
}
fn user_review_timestamp(micros: i64) -> Result {
time::OffsetDateTime::from_unix_timestamp_nanos(i128::from(micros) * 1_000)
.map_err(|_| AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR))
.and_then(|timestamp| {
shared_kernel::format_rfc3339(timestamp)
.map_err(|_| AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR))
})
}
fn rating_summary_payload(
summary: GameDistributionRatingSummaryRecord,
) -> GameDistributionRatingSummary {
GameDistributionRatingSummary {
average_score: summary.average_score,
rating_count: summary.rating_count,
}
}
async fn admin_review_games(
State(state): State,
Extension(ctx): Extension,
Extension(_admin): Extension,
query: Result, QueryRejection>,
) -> Result, AppError> {
let Query(query) = query.map_err(|_| bad_request("游戏查询参数不合法"))?;
let (page, page_size) = UserReviewListQuery {
page: query.page,
page_size: query.page_size,
}
.pagination()?;
let result = state
.spacetime_client()
.list_game_distribution_review_games(GameDistributionReviewGameListRecordInput {
query: normalize_optional(query.query),
page,
page_size,
})
.await
.map_err(map_user_review_admin_error)?;
Ok(json_success_body(
Some(&ctx),
AdminGameReviewGamesResponse {
games: result
.games
.into_iter()
.map(|game| AdminGameReviewGame {
game_id: game.game_id,
title: game.title,
status: game.status,
})
.collect(),
page: result.page,
page_size: result.page_size,
total: result.total,
total_pages: result.total_pages,
},
))
}
async fn admin_user_review_list(
State(state): State,
Extension(ctx): Extension,
Extension(_admin): Extension,
query: Result, QueryRejection>,
) -> Result, AppError> {
let Query(query) = query.map_err(|_| bad_request("评价查询参数不合法"))?;
let input = admin_user_review_list_input(query)?;
let result = state
.spacetime_client()
.list_admin_game_distribution_user_reviews(input)
.await
.map_err(map_user_review_admin_error)?;
Ok(json_success_body(
Some(&ctx),
AdminGameReviewsResponse {
reviews: result
.reviews
.into_iter()
.map(admin_user_review_payload)
.collect::, _>>()?,
page: result.page,
page_size: result.page_size,
total: result.total,
total_pages: result.total_pages,
},
))
}
fn admin_user_review_list_input(
query: AdminGameReviewsQuery,
) -> Result {
let (page, page_size) = UserReviewListQuery {
page: query.page,
page_size: query.page_size,
}
.pagination()?;
let status = query.status.unwrap_or_else(|| "all".to_string());
validate_review_list_status(&status).map_err(|error| bad_request(error.to_string()))?;
Ok(GameDistributionAdminUserReviewListRecordInput {
game_id: normalize_optional(query.game_id),
user_id: normalize_optional(query.user_id),
keyword: normalize_optional(query.keyword),
status,
page,
page_size,
})
}
async fn admin_user_review_detail(
State(state): State,
Extension(ctx): Extension,
Extension(_admin): Extension,
Path(review_id): Path,
) -> Result, AppError> {
let result = state
.spacetime_client()
.get_admin_game_distribution_user_review(review_id)
.await
.map_err(map_user_review_admin_error)?;
Ok(json_success_body(
Some(&ctx),
AdminGameReviewDetailResponse {
review: admin_user_review_payload(result.review)?,
operations: result
.operations
.into_iter()
.map(review_moderation_operation_payload)
.collect::, _>>()?,
},
))
}
async fn admin_moderate_user_review(
State(state): State,
Extension(ctx): Extension,
Extension(admin): Extension,
headers: HeaderMap,
Path(review_id): Path,
payload: Result, JsonRejection>,
) -> Result, AppError> {
let Json(payload) = payload.map_err(|_| bad_request("评价管理请求字段不合法"))?;
let input = review_moderation_input(
review_id,
admin.session().subject.clone(),
&headers,
payload,
)?;
let result = state
.spacetime_client()
.moderate_game_distribution_user_review(input)
.await
.map_err(map_user_review_admin_error)?;
Ok(json_success_body(
Some(&ctx),
AdminGameReviewModerationResponse {
review: result.review.map(admin_user_review_payload).transpose()?,
operation: review_moderation_operation_payload(result.operation)?,
replayed: result.replayed,
},
))
}
fn review_moderation_input(
review_id: String,
admin_user_id: String,
headers: &HeaderMap,
payload: AdminGameReviewModerationRequest,
) -> Result {
let idempotency_key = idempotency_key(headers)?;
if !matches!(payload.action.as_str(), "hide" | "restore" | "delete") {
return Err(bad_request("管理动作必须为 hide、restore 或 delete"));
}
let expected_created_at_micros = shared_kernel::parse_rfc3339(&payload.expected_created_at)
.map(shared_kernel::offset_datetime_to_unix_micros)
.map_err(|_| bad_request("目标评价创建时间格式不合法"))?;
let reason = normalize_review_moderation_reason(&payload.action, payload.reason.as_deref())
.map_err(|error| {
AppError::from_status(StatusCode::UNPROCESSABLE_ENTITY).with_message(error.to_string())
})?;
Ok(GameDistributionReviewModerationRecordInput {
review_id,
admin_user_id,
action: payload.action,
idempotency_key,
expected_created_at_micros,
reason,
})
}
fn admin_user_review_payload(
record: GameDistributionAdminUserReviewRecord,
) -> Result {
let review = user_review_payload(record.review)?;
Ok(AdminGameReview {
id: review.id,
game_id: review.game_id,
game: AdminGameReviewGameInfo {
title: record.game_title,
status: record.game_status,
},
author: review.author,
score: review.score,
comment: review.comment,
is_hidden: review.is_hidden,
created_at: review.created_at,
updated_at: review.updated_at,
})
}
fn review_moderation_operation_payload(
record: GameDistributionReviewModerationOperationRecord,
) -> Result {
Ok(AdminGameReviewOperation {
id: record.operation_id,
review_id: record.review_id,
game_id: record.game_id,
user_id: record.user_id,
review_created_at: user_review_timestamp(record.review_created_at_micros)?,
action: record.action,
admin_user_id: record.admin_user_id,
reason: record.reason,
created_at: user_review_timestamp(record.created_at_micros)?,
})
}
fn map_user_review_admin_error(error: SpacetimeClientError) -> AppError {
if let SpacetimeClientError::Procedure(message) = &error {
let status = if message.starts_with("REVIEW_NOT_FOUND") {
Some(StatusCode::NOT_FOUND)
} else if message.starts_with("REVIEW_CONFLICT")
|| message.starts_with("REVIEW_IDEMPOTENCY_CONFLICT")
{
Some(StatusCode::CONFLICT)
} else if message.starts_with("REVIEW_VALIDATION") {
Some(StatusCode::UNPROCESSABLE_ENTITY)
} else if message.starts_with("REVIEW_BAD_REQUEST") {
Some(StatusCode::BAD_REQUEST)
} else {
None
};
if let Some(status) = status {
return AppError::from_status(status).with_message(message.clone());
}
}
map_spacetime_error(error)
}
/// 发行网关根路径:等价于请求该游戏的 `index.html`。
async fn serve_release_entry(
state: State,
headers: HeaderMap,
Path(game_id): Path,
) -> Result {
serve_release_asset(state, headers, Path((game_id, "index.html".to_string()))).await
}
/// 公开发行网关。
///
/// 只服务当前已公开版本的游戏文件,路径必须在白名单内容类型内;私有 ZIP 对象和
/// 未公开版本不会因为知道 ID 而可读。
async fn serve_release_asset(
State(state): State,
headers: HeaderMap,
Path((game_id, asset_path)): Path<(String, String)>,
) -> Result {
// 发行文件必须由独立来源提供。带上平台 Cookie 的请求说明它正落在主站来源上,
// 此时同源脚本可以读到平台会话,必须直接关闭而不是降级服务。
if headers.contains_key(header::COOKIE) {
debug!(
operation = "release_rejected",
game_id = %game_id,
reason = "cookie_present",
"发行资源请求带平台 Cookie,已拒绝"
);
return Err(AppError::from_status(StatusCode::FORBIDDEN)
.with_message("发行资源必须在独立来源上请求"));
}
let asset_path = asset_path.trim_start_matches('/').to_string();
let content_type = release_asset_content_type(&asset_path).ok_or_else(|| {
debug!(
operation = "release_rejected",
game_id = %game_id,
asset_path = %asset_path,
reason = "unsupported_extension",
"发行资源扩展名不在白名单内"
);
AppError::from_status(StatusCode::NOT_FOUND)
})?;
let public_game = state
.spacetime_client()
.get_public_game_distribution_game(game_id.clone())
.await
.map_err(map_spacetime_error)?
.ok_or_else(|| {
debug!(
operation = "release_rejected",
game_id = %game_id,
asset_path = %asset_path,
reason = "not_public",
"游戏没有公开可玩版本"
);
AppError::from_status(StatusCode::NOT_FOUND)
})?;
let version = public_game.current_version.ok_or_else(|| {
debug!(
operation = "release_rejected",
game_id = %game_id,
asset_path = %asset_path,
reason = "no_active_version",
"游戏缺少当前公开版本"
);
AppError::from_status(StatusCode::NOT_FOUND)
})?;
if version.status != GAME_DISTRIBUTION_PUBLISHED_STATUS {
debug!(
operation = "release_rejected",
game_id = %game_id,
version_id = %version.version_id,
asset_path = %asset_path,
reason = "version_not_published",
status = %version.status,
"请求的版本不是公开状态"
);
return Err(AppError::from_status(StatusCode::NOT_FOUND));
}
let package = release_package_bytes(&state, &game_id, &version.version_id).await?;
let etag = release_asset_etag(&version.version_id, &asset_path);
release_package_asset_response(
&package,
&asset_path,
ReleaseAssetResponseInput {
content_type,
cache_control: "public, max-age=60, must-revalidate",
etag: Some(&etag),
if_none_match: headers.get(header::IF_NONE_MATCH),
accept_encoding: headers.get(header::ACCEPT_ENCODING),
},
)
}
/// 单次发行资源响应的输入:内容类型、缓存策略、条件请求与压缩协商。
struct ReleaseAssetResponseInput<'a> {
content_type: &'static str,
cache_control: &'static str,
/// 为空表示该响应不可缓存(后台试玩会话是 `no-store`),条件请求与 ETag 都不参与。
etag: Option<&'a str>,
if_none_match: Option<&'a HeaderValue>,
accept_encoding: Option<&'a HeaderValue>,
}
/// 低于这个字节数的响应不做 gzip:压缩头与字典开销会把收益吃掉。
const RELEASE_COMPRESSION_MIN_BYTES: usize = 1024;
/// 发行资源的强 ETag:同一版本同一路径的字节在包被冻结后不会改变。
///
/// 用摘要而不是拼字符串:资源路径来自 URL,可能带上 ETag 的保留字符(引号、反斜杠)。
fn release_asset_etag(version_id: &str, asset_path: &str) -> String {
let mut hasher = Sha256::new();
hasher.update(version_id.as_bytes());
hasher.update([0]);
hasher.update(asset_path.as_bytes());
let digest = hasher.finalize();
format!("\"{}\"", hex::encode(&digest[..16]))
}
/// `If-None-Match` 命中判定:支持 `*`、弱校验前缀 `W/` 与逗号分隔列表。
fn if_none_match_matches(header: Option<&HeaderValue>, etag: &str) -> bool {
let Some(value) = header.and_then(|value| value.to_str().ok()) else {
return false;
};
value.split(',').any(|candidate| {
let candidate = candidate.trim();
candidate == "*" || candidate.trim_start_matches("W/") == etag
})
}
/// 客户端是否接受 gzip;`gzip;q=0` 与 `*;q=0` 表示明确拒绝。
fn accepts_gzip_encoding(header: Option<&HeaderValue>) -> bool {
let Some(value) = header.and_then(|value| value.to_str().ok()) else {
return false;
};
value.split(',').any(|entry| {
let mut parts = entry.split(';');
let coding = parts.next().unwrap_or_default().trim();
if !coding.eq_ignore_ascii_case("gzip") && coding != "*" {
return false;
}
!parts.any(|parameter| {
// 权重参数名大小写不敏感(RFC 9110 §12.5.3):`Q=0` 与 `q=0` 一样是明确拒绝。
let parameter = parameter.trim();
let Some((name, value)) = parameter.split_once('=') else {
return false;
};
name.eq_ignore_ascii_case("q")
&& value
.trim()
.parse::()
.is_ok_and(|quality| quality <= 0.0)
})
})
}
/// 只压文本类发行资源;图片、音频、视频本身已是压缩格式,再压一遍只是浪费 CPU。
fn is_compressible_release_content_type(content_type: &str) -> bool {
content_type.starts_with("text/")
|| content_type.contains("javascript")
|| content_type.contains("json")
|| content_type.contains("svg")
|| content_type.contains("application/wasm")
}
/// 超过这个体积改用更快的压缩级别。
///
/// 单文件上限是 64 MiB,而发行网关是公开无鉴权端点:release 构建下 level 6 实测
/// 1.3 MiB→10 ms、8 MiB→59 ms、64 MiB→522 ms 纯 CPU,每请求重算会占满 worker 线程。
/// 大文件改用 level 1(zlib 端实测约为 level 6 的 1/3 耗时、压缩比只差约 3%),
/// 小文件仍用 level 6 拿更好的比例。
const RELEASE_COMPRESSION_FAST_ABOVE_BYTES: usize = 2 * 1024 * 1024;
fn gzip_release_asset(content: &[u8]) -> Option> {
let level = if content.len() >= RELEASE_COMPRESSION_FAST_ABOVE_BYTES {
GzipCompression::fast()
} else {
GzipCompression::new(6)
};
let mut encoder = GzEncoder::new(Vec::new(), level);
encoder.write_all(content).ok()?;
encoder.finish().ok()
}
fn release_package_asset_response(
package: &[u8],
asset_path: &str,
input: ReleaseAssetResponseInput<'_>,
) -> Result {
let ReleaseAssetResponseInput {
content_type,
cache_control,
etag,
if_none_match,
accept_encoding,
} = input;
let content = match extract_release_asset(package, asset_path) {
Ok(content) => content,
Err(ReleaseAssetError::FileTooLarge) => {
return Err(AppError::from_status(StatusCode::PAYLOAD_TOO_LARGE)
.with_message("发行资源超过响应上限"));
}
Err(_) => return Err(AppError::from_status(StatusCode::NOT_FOUND)),
};
// 条件请求必须在确认资源存在于包内之后判定:`If-None-Match: *` 只表示「任一份表示
// 存在就复用」,对包内不存在的路径仍要 404。
if let Some(etag) = etag
&& if_none_match_matches(if_none_match, etag)
{
return Ok(release_asset_not_modified_response(etag, cache_control));
}
let content = normalize_release_asset_references(content, content_type);
let content = inject_release_storage_bootstrap(content, content_type);
// 发行包里的字节是 ZIP 解压后的原文:不压缩时 Phaser 4 的 1.31 MiB 引擎包会原样
// 走完用户网络,而它在包内本来就是 deflate 压缩的。
let compressed = (accepts_gzip_encoding(accept_encoding)
&& is_compressible_release_content_type(content_type)
&& content.len() >= RELEASE_COMPRESSION_MIN_BYTES)
.then(|| gzip_release_asset(&content))
.flatten();
let (body, content_encoding) = match compressed {
Some(compressed) => (compressed, Some("gzip")),
None => (content, None),
};
Ok(release_asset_response_with_cache(
body,
content_type,
cache_control,
etag,
content_encoding,
))
}
fn normalize_release_asset_references(content: Vec, content_type: &str) -> Vec {
if !(content_type.starts_with("text/html")
|| content_type.starts_with("text/css")
|| content_type.contains("javascript"))
{
return content;
}
let mut source = match String::from_utf8(content) {
Ok(source) => source,
Err(error) => return error.into_bytes(),
};
for root in ["assets", "game", "ui"] {
source = source.replace(&format!("\"/{root}/"), &format!("\"{root}/"));
source = source.replace(&format!("'/{root}/"), &format!("'{root}/"));
source = source.replace(&format!("`/{root}/"), &format!("`{root}/"));
source = source.replace(&format!("url(/{root}/"), &format!("url({root}/"));
}
source.into_bytes()
}
fn inject_release_storage_bootstrap(content: Vec, content_type: &str) -> Vec {
if !content_type.starts_with("text/html") {
return content;
}
let mut result = Vec::with_capacity(RELEASE_STORAGE_BOOTSTRAP.len() + content.len());
result.extend_from_slice(RELEASE_STORAGE_BOOTSTRAP.as_bytes());
result.extend_from_slice(&content);
result
}
/// 读取(并按**对象键**缓存)已确认的私有资产。
///
/// 缓存键就是对象键,因此资产种类天然分开:同一 (作品, 版本) 的成品包与工程源包落在不同
/// 对象键上,各自取回自己那份字节,不会串味。上限由调用方给(发行包 200 MiB、工程源包
/// 同为 200 MiB),`max_bytes` 是 OSS 读路径的硬闸门。
async fn release_asset_bytes(
state: &AppState,
object_key: &str,
max_bytes: u64,
) -> Result {
let cache = &*RELEASE_PACKAGE_CACHE;
if let Some(cached) = cache.lock().ok().and_then(|guard| guard.get(object_key)) {
return Ok(cached);
}
let oss = state.project_snapshot_oss_client().ok_or_else(|| {
AppError::from_status(StatusCode::SERVICE_UNAVAILABLE).with_message("游戏发行包 OSS 未配置")
})?;
let bytes = oss
.get_object(
state.editor_oss_http_client(),
OssGetObjectRequest {
object_key: object_key.to_string(),
max_bytes: max_bytes as usize,
},
)
.await
.map_err(|error| {
if matches!(error, platform_oss::OssError::ObjectNotFound(_)) {
AppError::from_status(StatusCode::NOT_FOUND)
} else {
map_oss_error(error, "aliyun-oss")
}
})?;
// `Bytes::from(Vec)` 直接接管所有权,不再复制整包;之后每次命中缓存只加引用计数。
let bytes = Bytes::from(bytes);
if let Ok(mut guard) = cache.lock() {
guard.insert(object_key.to_string(), bytes.clone());
}
Ok(bytes)
}
/// 读取(并按对象键缓存)已确认的私有发行包;键与上限与既有行为逐字节一致。
async fn release_package_bytes(
state: &AppState,
game_id: &str,
version_id: &str,
) -> Result {
let object_key = game_distribution_package_object_key(game_id, version_id);
release_asset_bytes(state, &object_key, MAX_PACKAGE_BYTES).await
}
fn release_asset_response_with_cache(
content: Vec,
content_type: &'static str,
cache_control: &'static str,
etag: Option<&str>,
content_encoding: Option<&'static str>,
) -> Response {
let mut response = Response::new(Body::from(content));
let headers = response.headers_mut();
headers.insert(header::CONTENT_TYPE, HeaderValue::from_static(content_type));
headers.insert(
header::X_CONTENT_TYPE_OPTIONS,
HeaderValue::from_static("nosniff"),
);
headers.insert(
header::REFERRER_POLICY,
HeaderValue::from_static("no-referrer"),
);
headers.insert(
header::CACHE_CONTROL,
HeaderValue::from_static(cache_control),
);
// 发行资源可能带 gzip 表示;共享缓存必须按 Accept-Encoding 分桶,否则会把压缩体
// 发给不支持它的客户端。
headers.insert(header::VARY, HeaderValue::from_static("accept-encoding"));
if let Some(etag) = etag {
if let Ok(value) = HeaderValue::from_str(etag) {
headers.insert(header::ETAG, value);
}
}
if let Some(content_encoding) = content_encoding {
headers.insert(
header::CONTENT_ENCODING,
HeaderValue::from_static(content_encoding),
);
}
// 发行文档运行在 allow-scripts 的 opaque origin 沙箱里,其同包资源请求不再与
// 网关同源;CORP 必须允许跨来源,ES modules 还需要不带 credentials 的 CORS,
// 否则游戏自己的脚本会被浏览器拦下(实测 net::ERR_BLOCKED_BY_RESPONSE)。
// 这些是公开静态文件,放宽 CORP 不涉及凭据。
headers.insert(
header::HeaderName::from_static("cross-origin-resource-policy"),
HeaderValue::from_static("cross-origin"),
);
headers.insert(
header::ACCESS_CONTROL_ALLOW_ORIGIN,
HeaderValue::from_static("*"),
);
if content_type.starts_with("text/html") {
// 发行 HTML 走与主站不同的来源并强制最小权限策略;包内 meta 不能放宽。
headers.insert(
header::CONTENT_SECURITY_POLICY,
HeaderValue::from_static(
"default-src 'none'; script-src 'self' 'unsafe-inline'; style-src 'self' 'unsafe-inline'; img-src 'self' data: blob:; media-src 'self' data: blob:; font-src 'self' data:; connect-src 'self'; worker-src 'none'; object-src 'none'; frame-src 'none'; form-action 'none'; base-uri 'none'",
),
);
}
response
}
/// 条件请求命中:只回报校验器与缓存策略,不带正文。
///
/// 发行包在版本冻结后不可变,浏览器在 `max-age=60, must-revalidate` 之后只要拿到同一个
/// ETag 就能停在 304,不必把整个引擎包再下一次。
fn release_asset_not_modified_response(etag: &str, cache_control: &'static str) -> Response {
let mut response = Response::new(Body::empty());
*response.status_mut() = StatusCode::NOT_MODIFIED;
let headers = response.headers_mut();
headers.insert(
header::CACHE_CONTROL,
HeaderValue::from_static(cache_control),
);
headers.insert(header::VARY, HeaderValue::from_static("accept-encoding"));
if let Ok(value) = HeaderValue::from_str(etag) {
headers.insert(header::ETAG, value);
}
response
}
async fn list_games(
State(state): State,
Extension(ctx): Extension,
Query(query): Query,
) -> Result, AppError> {
let author_id = query
.author_id
.as_deref()
.map(module_auth::creator::normalize_user_id)
.transpose()
.map_err(|error| {
AppError::from_status(StatusCode::BAD_REQUEST).with_message(error.to_string())
})?;
let games = state
.spacetime_client()
.list_game_distribution_games(GameDistributionPublicGameListRecordInput {
search: normalize_optional(query.search),
category: normalize_optional(query.category),
limit: MAX_LIST_LIMIT,
author_id,
})
.await
.map_err(map_spacetime_error)?;
let games = games
.into_iter()
.map(public_game_payload)
.collect::>();
Ok(json_success_body(
Some(&ctx),
json!({ "games": games, "nextCursor": Value::Null }),
))
}
async fn get_game(
State(state): State,
Extension(ctx): Extension,
headers: HeaderMap,
Path(game_id): Path,
) -> Result, AppError> {
let game = state
.spacetime_client()
.get_public_game_distribution_game(game_id.clone())
.await
.map_err(map_spacetime_error)?
.ok_or_else(|| AppError::from_status(StatusCode::NOT_FOUND))?;
// 可选登录态:登录时才追加 `collected`(真实投影,不是前端本地状态)。
//
// 无效 / 过期 token 按**匿名**处理(与 `record_game_play` 同一取舍):公开详情是匿名可读的,
// 不能因为客户端带着一个过期 token 就把整页变成 401;客户端刷新 token 后会重新拉取。
// 这也意味着「带了 token 但被按匿名处理」时返回的响应与纯匿名完全相同——即**不带**该键。
let collected = match optional_access_token_from_headers(
&state,
format!("/api/game-distribution/games/{game_id}"),
headers,
ctx.request_id().to_string(),
)
.await
{
Ok(Some(authenticated)) => {
let user_id = authenticated.claims().user_id().to_string();
Some(
state
.spacetime_client()
.is_game_distribution_collected(game_id.clone(), user_id)
.await
.map_err(map_spacetime_error)?,
)
}
Ok(None) => None,
Err(error) => {
debug!(error = %error, "公开详情忽略无效 bearer,按匿名返回");
None
}
};
Ok(json_success_body(
Some(&ctx),
public_game_detail_payload(game, collected),
))
}
/// 公开详情负载:在公开目录那条 `public_game_payload` 之上按可选登录态追加 `collected`。
///
/// 抽成函数而不是写在 handler 里,一是让「登录才加键、匿名不加键」能被单测钉住,二是它同时是
/// DTO parity 脚本登记的响应构建器(证明这条路径确实会发出 `collected`)。
///
/// `collected == None`(匿名 / 无效 token 按匿名)时**不加键**,而不是发 `false`:`false` 会把
/// 「未登录」说成「没收藏」,客户端无法区分,会渲染出错误的收藏按钮态。
fn public_game_detail_payload(
game: GameDistributionPublicGameRecord,
collected: Option,
) -> Value {
let mut payload = public_game_payload(game);
if let Some(collected) = collected {
if let Value::Object(object) = &mut payload {
object.insert("collected".to_string(), json!(collected));
}
}
payload
}
/// 收藏(收录)的请求摘要:只绑定 `(user_id, game_id)`。
///
/// 服务端的收据键是 `(user_id, action, idempotency_key)`,而 `Idempotency-Key` 由客户端生成、
/// 可能在不同作品之间复用。摘要里带上这对组合,才让「同一个 key 撞到**不同作品**」被判成同键
/// 不同请求(409),而不是把另一个作品的收藏结果重放给当前请求者。
///
/// **同键换用户不会是 409、也不会互相命中**:收据键本身含 `user_id`,两个用户各带相同
/// `Idempotency-Key` 时读到的是各自(不存在)的收据,各自按新请求处理并 200 成功——端到端
/// 实测如此(`spacetime-module` 侧同步写明「不同用户 / 不同作品即使共用同一个
/// `Idempotency-Key` 也不会互相命中」);本条注释曾把「换用户」一并错写成 409。
fn collection_request_digest(user_id: &str, game_id: &str) -> Result {
Ok(compute_request_digest(
&serde_json::to_vec(&(user_id, game_id)).map_err(|error| internal(error.to_string()))?,
))
}
/// 收藏(收录)某作品:`PUT /games/{gameId}/collection`。
///
/// 幂等语义:必须带 `Idempotency-Key`。同键重放返回同一结果并带 `replayed = true`;
/// 同键**换作品**是 409(摘要绑定 `(user_id, game_id)`);同键**换用户**是 200 各自成功
/// (收据键 `(user_id, action, idempotency_key)` 按用户隔离,两个用户不会命中彼此的收据);
/// 重复收藏(不同键)不会产生第二行,仍然成功。
async fn collect_game(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(game_id): Path,
) -> Result, AppError> {
let user_id = auth.claims().user_id().to_string();
let idempotency_key = idempotency_key(&headers)?;
let request_digest = collection_request_digest(user_id.as_str(), game_id.as_str())?;
let collection = state
.spacetime_client()
.set_game_distribution_collection(GameDistributionCollectGameRecordInput {
game_id: game_id.clone(),
user_id: user_id.clone(),
idempotency_key,
request_digest,
now_micros: now_micros(),
})
.await
.map_err(map_spacetime_error)?;
info!(
request_id = ctx.request_id(),
operation = "game_collection_set",
game_id = %game_id,
replayed = collection.replayed,
elapsed_ms = ctx.elapsed(),
"收藏游戏作品"
);
Ok(json_success_body(
Some(&ctx),
GameDistributionCollectionState {
collected: collection.collected,
replayed: Some(collection.replayed),
},
))
}
/// 取消收藏(收录):`DELETE /games/{gameId}/collection`。
///
/// **不要求 `Idempotency-Key`**:删除按确定性主键 `{userId}:{gameId}` 执行,重复调用结果完全
/// 相同(不存在也算成功),没有「重放 vs 新意图」需要区分——幂等键只在请求本身无法表达意图时
/// 才有意义。也**不要求作品仍公开**:下架后拒绝取消只会给用户留下清理不掉的脏行。
async fn uncollect_game(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
Path(game_id): Path,
) -> Result, AppError> {
let user_id = auth.claims().user_id().to_string();
let collection = state
.spacetime_client()
.unset_game_distribution_collection(GameDistributionUncollectGameRecordInput {
game_id: game_id.clone(),
user_id,
})
.await
.map_err(map_spacetime_error)?;
info!(
request_id = ctx.request_id(),
operation = "game_collection_unset",
game_id = %game_id,
elapsed_ms = ctx.elapsed(),
"取消收藏游戏作品"
);
Ok(json_success_body(
Some(&ctx),
GameDistributionCollectionState {
collected: collection.collected,
// 取消没有幂等键,因此不下发 `replayed`(`skip_serializing_if` 会略过该键)。
replayed: None,
},
))
}
/// 「我的收藏(收录)」:`GET /my-collections?limit=&cursor=`。
///
/// 逐条用公开目录同一份 `public_game_payload` 组装,形状与公开目录一致(`games` + `nextCursor`)。
/// 只返回当前公开可读的作品;已下架 / 软删除的收藏**只是不在响应里**,行不删除——作品重新
/// 公开后会自动回来。
///
/// 分页:默认 20、上限 50,超界**截断**(与后台列表一致:客户端拿到一个完整页,而不是重试错误);
/// 游标格式非法由模块侧报错并在这里透传成 400(不吞掉、也不自己造一种 200 的空页)。
async fn list_my_collections(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
Query(query): Query,
) -> Result, AppError> {
let user_id = auth.claims().user_id().to_string();
// 与模块侧同一套归一化(同一函数),因此日志里的 `limit` 就是真正生效的页大小。
let limit = my_collections_page_limit(query.limit);
let cursor = normalize_optional(query.cursor);
let (games, next_cursor) = state
.spacetime_client()
.list_game_distribution_collections(user_id, limit, cursor.clone())
.await
.map_err(map_spacetime_error)?;
info!(
request_id = ctx.request_id(),
operation = "game_collections_listed",
games = games.len(),
limit,
max_limit = GAME_DISTRIBUTION_COLLECTION_PAGE_LIMIT_MAX,
has_cursor = cursor.is_some(),
has_more = next_cursor.is_some(),
elapsed_ms = ctx.elapsed(),
"读取我的收藏列表"
);
Ok(json_success_body(
Some(&ctx),
my_collections_payload(games, next_cursor),
))
}
/// 「我的收藏」响应负载:分页字段与公开目录逐字一致(`games` + `nextCursor`)。
///
/// `nextCursor` 是**真实**游标:还有下一页时给出,最后一页为 `null`,客户端据此决定是否继续拉。
///
/// 同步纯函数:既是 handler 的组装点,也是 DTO parity 脚本登记的响应构建器。
fn my_collections_payload(
games: Vec,
next_cursor: Option,
) -> Value {
let games = games
.into_iter()
.map(public_game_payload)
.collect::>();
json!({ "games": games, "nextCursor": next_cursor })
}
/// 查询串的 `limit` → 生效页大小:缺省取默认 20,超界截断到上限 50(`0` 也取默认)。
///
/// 归一化本身委托给 `module_game_distribution::game_distribution_collection_page_limit`,
/// 与事务里真正切页用的是**同一个函数**,避免「日志写 50、实际发了 20」这类漂移。
fn my_collections_page_limit(limit: Option) -> u32 {
game_distribution_collection_page_limit(
limit.unwrap_or(GAME_DISTRIBUTION_COLLECTION_PAGE_LIMIT_DEFAULT),
) as u32
}
/// 读接口的「不可读即 404」:族谱与衍生列表的锚点必须公开可读,否则按「不存在」处理。
///
/// 抽成函数是为了让这条映射可被单测钉住(不可读 → 404),而不是散落在两个 handler 里。
fn lineage_read_or_not_found(value: Option) -> Result {
value.ok_or_else(|| AppError::from_status(StatusCode::NOT_FOUND))
}
/// 公开创作族谱:以该作品的母版为顶返回整棵树。公开只读、匿名可读。
///
/// 锚点必须公开可读:不存在、未公开或已软删除一律 404,与公开详情同一口径。这里**不**
/// 用「空标题节点」代替 404——那等于对外确认该 gameId 存在、它是第几代、它在血缘里的位置。
async fn get_game_lineage(
State(state): State,
Extension(ctx): Extension,
Path(game_id): Path,
) -> Result, AppError> {
let tree = lineage_read_or_not_found(
state
.spacetime_client()
.get_game_distribution_lineage(game_id)
.await
.map_err(map_spacetime_error)?,
)?;
Ok(json_success_body(Some(&ctx), lineage_tree_payload(tree)))
}
/// 公开「被改编」列表:该作品的直接衍生作品(只含未软删除且已公开的直接子代)。
///
/// 与公开详情 `forkCount` 同口径,因此条数与「被改编 N」一致;锚点不可公开读取时 404。
async fn get_game_derived_games(
State(state): State,
Extension(ctx): Extension,
Path(game_id): Path,
) -> Result, AppError> {
let derived = lineage_read_or_not_found(
state
.spacetime_client()
.list_game_distribution_derived_games(game_id)
.await
.map_err(map_spacetime_error)?,
)?;
Ok(json_success_body(
Some(&ctx),
derived_games_payload(derived),
))
}
/// 线上可见性字符串 → DTO 枚举;未知取值按「未公开」保守解释,不会因此多暴露任何信息。
fn game_distribution_visibility(value: &str) -> GameDistributionVisibility {
match value {
"published" => GameDistributionVisibility::Published,
"suspended" => GameDistributionVisibility::Suspended,
_ => GameDistributionVisibility::Unpublished,
}
}
fn lineage_node_payload(node: GameDistributionLineageNodeRecord) -> GameDistributionLineageNode {
GameDistributionLineageNode {
game_id: node.game_id,
title: node.title,
author_name: node.author_name,
generation: node.generation,
parent_game_id: node.parent_game_id,
play_count: node.play_count,
status: game_distribution_visibility(node.status.as_str()),
}
}
/// 族谱 / 衍生列表走结构化 DTO 而不是手拼 JSON:节点字段本来就少且固定,
/// 用 DTO 才能让「不发对象键、不发素材键」由类型保证,而不是靠人肉回忆。
fn lineage_tree_payload(
tree: GameDistributionLineageTreeRecord,
) -> GameDistributionLineageResponse {
GameDistributionLineageResponse {
root_game_id: tree.root_game_id,
root: tree.root.map(lineage_node_payload),
nodes: tree.nodes.into_iter().map(lineage_node_payload).collect(),
truncated: tree.truncated,
}
}
fn derived_games_payload(
derived: GameDistributionDerivedGamesRecord,
) -> GameDistributionDerivedResponse {
GameDistributionDerivedResponse {
game_id: derived.game_id,
nodes: derived
.nodes
.into_iter()
.map(lineage_node_payload)
.collect(),
truncated: derived.truncated,
}
}
/// 一次游玩上报的请求体;只有匿名身份需要 `clientId`,登录身份由 bearer 决定。
#[derive(Debug, Default, Deserialize)]
#[serde(rename_all = "camelCase")]
struct RecordGamePlayRequest {
#[serde(default)]
client_id: Option,
}
/// 记录一次「开始游戏」。
///
/// 公开端点:登录用户按 `userId` 去重,匿名按 `clientId`(缺失时回退 `IP + UA`)去重;
/// 命中 30 分钟去重窗口或超过 `IP + game` 限流时不增加计数。计数只进内存缓冲,
/// 立即返回 `recorded`,任何失败都不影响游玩本身。
async fn record_game_play(
State(state): State,
Extension(ctx): Extension,
Path(game_id): Path,
headers: HeaderMap,
body: Bytes,
) -> Result, AppError> {
let game_id = game_id.trim().to_string();
if game_id.is_empty() {
return Err(AppError::from_status(StatusCode::NOT_FOUND));
}
let client_ip = client_ip_from_headers(&headers);
// 先在内存里挡掉明显超限的请求,避免它们也去打一次 SpacetimeDB;真正计数时 record 会再判一次。
if state
.game_play_counter()
.is_rate_limited(&game_id, &client_ip, Instant::now())
{
return Err(AppError::from_status(StatusCode::TOO_MANY_REQUESTS));
}
// 非公开 / 已下架 / 已暂停的游戏不计数,按不存在返回。
let is_public = state
.spacetime_client()
.get_public_game_distribution_game(game_id.clone())
.await
.map_err(map_spacetime_error)?
.is_some();
if !is_public {
return Err(AppError::from_status(StatusCode::NOT_FOUND));
}
let user_agent = user_agent_tag(&headers);
let authenticated = optional_access_token_from_headers(
&state,
format!("/api/game-distribution/games/{game_id}/plays"),
headers,
ctx.request_id().to_string(),
)
.await
.unwrap_or_else(|error| {
// 可选 bearer:无效 token 按匿名处理,绝不能因为它挡掉一次真实游玩。
debug!(error = %error, "游戏游玩计数忽略无效 bearer,按匿名计数");
None
});
let identity = authenticated
.as_ref()
.map(|token| format!("user:{}", token.claims().user_id()))
.or_else(|| request_client_id(&body).map(|client_id| format!("client:{client_id}")))
.unwrap_or_else(|| format!("ip:{client_ip}|ua:{user_agent}"));
let outcome = state.game_play_counter().record(
GamePlayReport {
game_id: &game_id,
identity: &identity,
client_ip: &client_ip,
},
Instant::now(),
);
if outcome == GamePlayOutcome::RateLimited {
return Err(AppError::from_status(StatusCode::TOO_MANY_REQUESTS));
}
Ok(json_success_body(
Some(&ctx),
json!({ "recorded": outcome == GamePlayOutcome::Counted }),
))
}
fn request_client_id(body: &Bytes) -> Option {
if body.is_empty() {
return None;
}
let request = serde_json::from_slice::(body).ok()?;
request
.client_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(|value| value.chars().take(128).collect())
}
fn user_agent_tag(headers: &HeaderMap) -> String {
headers
.get(header::USER_AGENT)
.and_then(|value| value.to_str().ok())
.map(str::trim)
.filter(|value| !value.is_empty())
.unwrap_or("unknown")
.chars()
.take(64)
.collect()
}
/// 作者自有游戏列表:只返回当前认证主体名下的游戏与最近版本状态。
async fn list_my_games(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
) -> Result, AppError> {
let games = state
.spacetime_client()
.list_owner_game_distribution_games(
spacetime_client::GameDistributionOwnerGameListRecordInput {
owner_user_id: auth.claims().user_id().to_string(),
limit: MAX_LIST_LIMIT,
},
)
.await
.map_err(map_spacetime_error)?;
let payload = games
.into_iter()
.map(owner_game_entry_payload)
.collect::>();
Ok(json_success_body(Some(&ctx), json!({ "games": payload })))
}
/// 作者读取自己名下单个游戏的详情,与 `list_my_games` 的条目同形。
///
/// 作者管理页要能打开「审核中 / 被驳回 / 已下架 / 已撤回」的作品,公开详情只服务已公开
/// 投影,所以作者视角必须走这条 owner 作用域路由,否则作者点自己的作品只会拿到 404。
/// 游戏不存在或不属于当前主体都返回 404,避免用错误码区分"别人的游戏"和"不存在的游戏"。
async fn get_owner_game(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
Path(game_id): Path,
) -> Result, AppError> {
let owner_user_id = auth.claims().user_id().to_string();
let game = state
.spacetime_client()
.get_game_distribution_game(GameDistributionGetGameRecordInput {
game_id: game_id.clone(),
owner_user_id: Some(owner_user_id.clone()),
})
.await
.map_err(map_spacetime_error)?
.ok_or_else(|| AppError::from_status(StatusCode::NOT_FOUND))?;
let (versions, fork_count) = load_owner_game_versions(&state, owner_user_id, &game_id).await?;
Ok(json_success_body(
Some(&ctx),
json!({
"game": owner_game_entry_payload(GameDistributionOwnerGameRecord {
game,
versions,
fork_count,
}),
}),
))
}
/// 读取当前主体名下某个游戏的版本列表。
///
/// 目前复用作者自有列表 procedure(版本随游戏聚合返回),因此与 `/my-games` 共享同一个
/// 条数上限:单作者作品数超过上限时该游戏没有版本记录。放量前需要补一条按 `game_id`
/// 精确列版本的 procedure。
async fn load_owner_game_versions(
state: &AppState,
owner_user_id: String,
game_id: &str,
) -> Result<(Vec, u64), AppError> {
let games = state
.spacetime_client()
.list_owner_game_distribution_games(
spacetime_client::GameDistributionOwnerGameListRecordInput {
owner_user_id,
limit: MAX_LIST_LIMIT,
},
)
.await
.map_err(map_spacetime_error)?;
Ok(games
.into_iter()
.find(|entry| entry.game.game_id == game_id)
// 「被改编 N」与版本列表来自同一份 owner 聚合投影,详情页与列表页因此不会各算一套。
.map(|entry| (entry.versions, entry.fork_count))
.unwrap_or_default())
}
/// 把资料编辑请求映射成创建请求,以复用同一套资料校验与素材归属解析。
///
/// `localProjectId` 只属于创建语义,编辑资料不参与游戏身份复用,因此固定为空;
/// 改编来源同理——资料编辑不是建立血缘的入口(血缘在创建作品时一次性写入且不可变),
/// 所以 `fork` 也固定为 `None`。
fn game_metadata_update_as_create_request(
payload: &GameDistributionUpdateGameMetadataRequest,
) -> GameDistributionCreateGameRequest {
GameDistributionCreateGameRequest {
local_project_id: None,
title: payload.title.clone(),
summary: payload.summary.clone(),
description: payload.description.clone(),
category: payload.category.clone(),
tags: payload.tags.clone(),
cover_asset_id: payload.cover_asset_id.clone(),
screenshots: payload.screenshots.clone(),
device_support: payload.device_support.clone(),
input_modes: payload.input_modes.clone(),
orientation: payload.orientation,
// 这两个字段只服务创建语义:资料编辑既不建立血缘也不改授权档位(授权只能靠
// `set_fork_authorization` 单向提升),所以复用创建校验时固定取默认值。
fork_authorization: GameDistributionForkAuthorization::Forbidden,
fork: None,
}
}
/// 编辑游戏资料的审计草稿:谁在什么时候把哪个作品的哪些展示字段改成了什么。
fn build_game_metadata_update_audit(
owner_user_id: &str,
game_id: &str,
title: &str,
category: &str,
expected_publication_revision: u64,
) -> TrackingEventDraft {
let mut draft = TrackingEventDraft::user(
"game_distribution_game_metadata_updated",
"game-distribution",
owner_user_id,
);
draft.metadata = json!({
"gameId": game_id,
"title": title,
"category": category,
"expectedPublicationRevision": expected_publication_revision,
});
draft
}
/// 软删除游戏的审计草稿:删除动作不可逆,必须能回答"谁在什么时候删了哪个作品"。
fn build_game_delete_audit(
owner_user_id: &str,
game_id: &str,
title: &str,
expected_publication_revision: u64,
) -> TrackingEventDraft {
let mut draft = TrackingEventDraft::user(
"game_distribution_game_deleted",
"game-distribution",
owner_user_id,
);
draft.metadata = json!({
"gameId": game_id,
"title": title,
"expectedPublicationRevision": expected_publication_revision,
});
draft
}
/// 作者编辑自己名下游戏的展示资料。
///
/// 资料在 game 级立即生效:页面、公开目录与详情下一次读取即换新;随版本冻结的包摘要与
/// 资料快照不受影响,下一次审核通过仍会用新版本的冻结资料覆盖游戏行。写入要求
/// `Idempotency-Key` 与 `expectedPublicationRevision` CAS,并落 `tracking_event` 审计。
async fn update_owner_game_metadata(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(game_id): Path,
Json(payload): Json,
) -> Result, AppError> {
ensure_publish_enabled(&state, Some(auth.claims().user_id())).await?;
let idempotency_key = idempotency_key(&headers)?;
let owner_user_id = auth.claims().user_id().to_string();
// 资料校验与素材归属解析复用创建游戏同一套:编辑不能绕过"必须有封面/截图必须是本人图片"。
let create_metadata = game_metadata_update_as_create_request(&payload);
validate_game_metadata(&create_metadata)?;
let (cover_asset_id, cover_object_key, screenshots) =
resolve_owned_game_media(&state, owner_user_id.as_str(), &create_metadata).await?;
let audit = build_game_metadata_update_audit(
owner_user_id.as_str(),
game_id.as_str(),
payload.title.as_str(),
payload.category.as_str(),
payload.expected_publication_revision,
);
let request_digest = compute_request_digest(
&serde_json::to_vec(&(game_id.as_str(), &payload))
.map_err(|error| internal(error.to_string()))?,
);
let log_owner_user_id = owner_user_id.clone();
let log_title = payload.title.clone();
let game = state
.spacetime_client()
.update_game_distribution_game_metadata(GameDistributionUpdateMetadataRecordInput {
game_id,
owner_user_id,
expected_publication_revision: payload.expected_publication_revision,
title: payload.title,
summary: payload.summary,
description: payload.description,
category: payload.category,
tags_json: serde_json::to_string(&payload.tags)
.map_err(|error| internal(error.to_string()))?,
cover_asset_id: Some(cover_asset_id),
cover_object_key: Some(cover_object_key),
screenshots_json: Some(
serde_json::to_string(&screenshots).map_err(|error| internal(error.to_string()))?,
),
device_support_desktop: payload.device_support.desktop,
device_support_mobile: payload.device_support.mobile,
device_support_touch: payload.device_support.touch,
input_modes_json: serde_json::to_string(&payload.input_modes)
.map_err(|error| internal(error.to_string()))?,
orientation: orientation_wire_value(payload.orientation)?,
idempotency_key,
request_digest,
now_micros: now_micros(),
})
.await
.map_err(map_spacetime_error)?;
record_tracking_event_after_success(&state, &ctx, audit).await;
info!(
request_id = ctx.request_id(),
operation = "game_metadata_updated",
game_id = %game.0.game_id,
owner_user_id = %log_owner_user_id,
title = %log_title,
publication_revision = game.0.publication_revision,
replayed = game.1,
elapsed_ms = ctx.elapsed(),
"作者更新游戏展示资料"
);
Ok(json_success_body(
Some(&ctx),
json!({ "game": game_payload(&game.0), "replayed": game.1 }),
))
}
/// 作者软删除自己名下的游戏。
///
/// 删除只标记 `deleted_at` 并把公开投影下线,保留版本行、发行包与冻结资料;作者视图、
/// 公开目录/详情、发行网关与后台默认列表都不再返回该作品。与下架一致,该动作不受发布
/// 灰度开关约束(收紧投稿时仍必须允许作者撤下自己的内容)。
async fn delete_owner_game(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(game_id): Path,
Query(query): Query,
) -> Result, AppError> {
let idempotency_key = idempotency_key(&headers)?;
let owner_user_id = auth.claims().user_id().to_string();
let request_digest = compute_request_digest(
&serde_json::to_vec(&(game_id.as_str(), query.expected_publication_revision))
.map_err(|error| internal(error.to_string()))?,
);
let log_owner_user_id = owner_user_id.clone();
let log_expected_revision = query.expected_publication_revision;
let game = state
.spacetime_client()
.delete_game_distribution_game(GameDistributionDeleteGameRecordInput {
game_id: game_id.clone(),
owner_user_id,
expected_publication_revision: query.expected_publication_revision,
idempotency_key,
request_digest,
now_micros: now_micros(),
})
.await
.map_err(map_spacetime_error)?;
let audit = build_game_delete_audit(
log_owner_user_id.as_str(),
game_id.as_str(),
game.0.title.as_str(),
log_expected_revision,
);
record_tracking_event_after_success(&state, &ctx, audit).await;
warn!(
request_id = ctx.request_id(),
operation = "game_deleted",
game_id = %game_id,
owner_user_id = %log_owner_user_id,
expected_publication_revision = log_expected_revision,
publication_revision = game.0.publication_revision,
replayed = game.1,
elapsed_ms = ctx.elapsed(),
"作者软删除游戏"
);
Ok(json_success_body(
Some(&ctx),
json!({ "game": game_payload(&game.0), "replayed": game.1 }),
))
}
async fn create_game(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
payload: Result, JsonRejection>,
) -> Result, AppError> {
// 未知共创档位(例如 `"allowed"`)在进入业务前就被 serde 拦下:必须映射成平台信封的 400,
// 而不是让 axum 默认的 `JsonRejection` 回 422 纯文本(前端错误处理依赖信封)。框架文本只用来
// 选文案,绝不回传:至少不会把「请求体里哪一段长什么样」透给客户端。
let Json(payload) = payload.map_err(|rejection| {
if rejection.body_text().contains("forkAuthorization") {
bad_request("共创授权档位不合法,只接受 forbidden / nonCommercial / full")
} else {
bad_request("创建作品请求字段不合法")
}
})?;
ensure_publish_enabled(&state, Some(auth.claims().user_id())).await?;
let idempotency_key = idempotency_key(&headers)?;
validate_game_metadata(&payload)?;
// 创建游戏时就把封面/截图的归属与类型校验掉:否则游戏行会先落一个不属于当前作者
// 或根本不存在的素材 ID,直到创建版本才失败,留下无法解释的半成品资料。
resolve_owned_game_media(&state, auth.claims().user_id(), &payload).await?;
let now = now_micros();
let game_id = format!("game_{}", Uuid::new_v4().simple());
let request_digest = compute_request_digest(
&serde_json::to_vec(&payload).map_err(|error| internal(error.to_string()))?,
);
let game = state
.spacetime_client()
.create_game_distribution_game(spacetime_client::GameDistributionCreateGameRecordInput {
game_id,
owner_user_id: auth.claims().user_id().to_string(),
title: payload.title,
summary: payload.summary,
description: payload.description,
category: payload.category,
tags_json: serde_json::to_string(&payload.tags)
.map_err(|error| internal(error.to_string()))?,
cover_asset_id: payload.cover_asset_id,
author_name: None,
author_avatar_url: None,
device_support_desktop: payload.device_support.desktop,
device_support_mobile: payload.device_support.mobile,
device_support_touch: payload.device_support.touch,
input_modes_json: serde_json::to_string(&payload.input_modes)
.map_err(|error| internal(error.to_string()))?,
orientation: orientation_wire_value(payload.orientation)?,
fork_authorization: fork_authorization_value(payload.fork_authorization),
idempotency_key,
request_digest,
now_micros: now,
local_project_id: normalize_local_project_id(payload.local_project_id.as_deref())?,
forked_from_game_id: payload
.fork
.as_ref()
.map(|fork| fork.parent_game_id.clone()),
forked_from_version_id: payload
.fork
.as_ref()
.map(|fork| fork.parent_version_id.clone()),
})
.await
.map_err(map_spacetime_error)?;
Ok(json_success_body(Some(&ctx), game_payload(&game.0)))
}
async fn create_version(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(game_id): Path,
Json(payload): Json,
) -> Result, AppError> {
ensure_publish_enabled(&state, Some(auth.claims().user_id())).await?;
let idempotency_key = idempotency_key(&headers)?;
validate_version_declaration(&payload)?;
let now = now_micros();
let version_id = format!("gamever_{}", Uuid::new_v4().simple());
let metadata_json =
resolve_version_metadata_json(&state, auth.claims().user_id(), &payload.game_metadata)
.await?;
let request_digest = compute_request_digest(
&serde_json::to_vec(&(game_id.as_str(), &payload, metadata_json.as_str()))
.map_err(|error| internal(error.to_string()))?,
);
let version = state
.spacetime_client()
.create_game_distribution_version(
spacetime_client::GameDistributionCreateVersionRecordInput {
game_id,
owner_user_id: auth.claims().user_id().to_string(),
version_id,
version_number: payload.version_number,
metadata_json,
package_sha256: payload.package_sha256,
package_bytes: payload.package_bytes,
package_file_count: payload.package_file_count,
package_entry_path: payload.package_entry_path,
local_project_id: normalize_local_project_id(payload.local_project_id.as_deref())?,
idempotency_key,
request_digest,
now_micros: now,
},
)
.await
.map_err(map_spacetime_error)?;
Ok(json_success_body(
Some(&ctx),
private_version_payload(&version.0),
))
}
async fn upload_package(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(version_id): Path,
body: Bytes,
) -> Result, AppError> {
require_zip_content_type(&headers)?;
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let idempotency_key = idempotency_key(&headers)?;
let expected = state
.spacetime_client()
.get_owner_game_distribution_version(owner_user_id.clone(), version_id.clone())
.await
.map_err(map_spacetime_error)?
.ok_or_else(|| AppError::from_status(StatusCode::NOT_FOUND))?;
let manifest = match validate_release_zip(&body) {
Ok(manifest) => manifest,
Err(error) => {
let reason = format!("{error:?}");
warn!(
request_id = ctx.request_id(),
operation = "package_rejected",
game_id = %expected.game_id,
version_id = %version_id,
code = "PACKAGE_VALIDATION_FAILED",
reason = %reason,
uploaded_bytes = body.len(),
elapsed_ms = ctx.elapsed(),
"发行包校验失败"
);
let mapped = map_package_error(error);
record_upload_failure(
&state,
&owner_user_id,
&version_id,
&idempotency_key,
"PACKAGE_VALIDATION_FAILED",
reason,
)
.await;
return Err(mapped);
}
};
let package_object_key = format!(
"{GAME_DISTRIBUTION_OBJECT_PREFIX}{}/{version_id}.zip",
expected.game_id
);
let oss = state.project_snapshot_oss_client().ok_or_else(|| {
AppError::from_status(StatusCode::SERVICE_UNAVAILABLE).with_message("游戏发行包 OSS 未配置")
})?;
let existing = oss
.head_internal_object(state.editor_oss_http_client(), &package_object_key)
.await
.map_err(|error| map_oss_error(error, "aliyun-oss"))?;
let skipped = match existing {
Some(existing) if existing.content_length == manifest.package_bytes => true,
Some(_) => {
let error = AppError::from_status(StatusCode::CONFLICT)
.with_code("PACKAGE_OBJECT_MISMATCH")
.with_message("发行包对象已存在但体积不一致");
record_upload_failure(
&state,
&owner_user_id,
&version_id,
&idempotency_key,
"PACKAGE_OBJECT_MISMATCH",
"发行包对象已存在但体积不一致".to_string(),
)
.await;
return Err(error);
}
None => false,
};
if !skipped {
// 单次 100 MiB 档 PUT 在本机实测 12 秒上下(上限提升到 200 MiB 后单次耗时与失败
// 暴露面同步放大),偶发传输失败会让作者白传一次;
// 这里按 platform-oss 既有的可重试分类做受控重试(只重试传输/超时/408/429/5xx)。
oss.put_internal_object_with_retry(
state.editor_oss_http_client(),
OssInternalPutObjectRequest {
object_key: package_object_key.clone(),
content_type: Some("application/zip".to_string()),
access: OssObjectAccess::Private,
metadata: BTreeMap::new(),
body: body.to_vec(),
},
GAME_DISTRIBUTION_OSS_PUT_MAX_ATTEMPTS,
&GAME_DISTRIBUTION_OSS_PUT_RETRY_DELAYS_MS,
)
.await
.map_err(|error| map_oss_error(error, "aliyun-oss"))?;
}
confirm_validated_package(
&state,
&ctx,
&owner_user_id,
&version_id,
&expected,
&manifest,
package_object_key,
&idempotency_key,
skipped,
)
.await
}
/// 校验通过后的共同收口:声明比对 → 确认 → 结构化事件。
///
/// 整包 `PUT` 与分片续传的完成动作共用这条路径,两种入口的校验、幂等与事件口径必须一致;
/// 任何入口都不得绕过它直接写版本状态。
#[allow(clippy::too_many_arguments)]
async fn confirm_validated_package(
state: &AppState,
ctx: &RequestContext,
owner_user_id: &str,
version_id: &str,
expected: &GameDistributionVersionRecord,
manifest: &ReleasePackageManifest,
package_object_key: String,
idempotency_key: &str,
oss_put_skipped: bool,
) -> Result, AppError> {
if manifest.package_sha256 != expected.package_sha256
|| manifest.package_bytes != expected.package_bytes
|| u32::try_from(manifest.files.len()).unwrap_or(u32::MAX) != expected.package_file_count
|| expected.package_entry_path != "index.html"
{
warn!(
request_id = ctx.request_id(),
operation = "package_rejected",
game_id = %expected.game_id,
version_id = %version_id,
code = "PACKAGE_MISMATCH",
declared_bytes = expected.package_bytes,
actual_bytes = manifest.package_bytes,
declared_file_count = expected.package_file_count,
actual_file_count = u32::try_from(manifest.files.len()).unwrap_or(u32::MAX),
elapsed_ms = ctx.elapsed(),
"发行包与版本声明不一致"
);
let error = AppError::from_status(StatusCode::CONFLICT)
.with_code("PACKAGE_MISMATCH")
.with_details(json!({
"provider": "game-distribution",
"message": "发行包摘要、体积、文件数或入口与版本声明不一致",
}));
record_upload_failure(
state,
owner_user_id,
version_id,
idempotency_key,
"PACKAGE_MISMATCH",
"发行包摘要、体积、文件数或入口与版本声明不一致".to_string(),
)
.await;
return Err(error);
}
let package_manifest_json = package_manifest_json(manifest)?;
let request_digest = compute_request_digest(
&serde_json::to_vec(&(version_id, manifest.package_sha256.as_str()))
.map_err(|error| internal(error.to_string()))?,
);
let log_game_id = expected.game_id.clone();
let log_package_bytes = manifest.package_bytes;
let log_file_count = u32::try_from(manifest.files.len()).unwrap_or(u32::MAX);
let log_sha_prefix = manifest.package_sha256.chars().take(12).collect::();
let confirmed = state
.spacetime_client()
.confirm_game_distribution_package(
spacetime_client::GameDistributionConfirmPackageRecordInput {
version_id: version_id.to_string(),
owner_user_id: owner_user_id.to_string(),
package_sha256: manifest.package_sha256.clone(),
package_bytes: manifest.package_bytes,
package_file_count: u32::try_from(manifest.files.len()).unwrap_or(u32::MAX),
package_entry_path: "index.html".to_string(),
package_object_key,
package_manifest_json,
idempotency_key: idempotency_key.to_string(),
request_digest,
updated_at_micros: now_micros(),
},
)
.await
.map_err(map_spacetime_error)?;
info!(
request_id = ctx.request_id(),
operation = "package_confirmed",
game_id = %log_game_id,
version_id = %confirmed.0.version_id,
package_bytes = log_package_bytes,
file_count = log_file_count,
sha256_prefix = %log_sha_prefix,
oss_put_skipped,
elapsed_ms = ctx.elapsed(),
"发行包已确认"
);
Ok(json_success_body(
Some(ctx),
json!({ "versionId": confirmed.0.version_id, "status": confirmed.0.status }),
))
}
/// 分片续传的状态查询:客户端拿到的「已收字节」来自 OSS 对象事实,不依赖本地记录,
/// 因此进程重启、换机器或换网络后都能从权威偏移继续。
async fn package_upload_state(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
Path(version_id): Path,
) -> Result, AppError> {
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let version = load_owner_version_or_404(&state, owner_user_id, version_id.clone()).await?;
let oss = game_distribution_oss_client(&state)?;
let object_key = game_distribution_package_object_key(&version.game_id, &version_id);
let received_bytes = staged_package_bytes(&state, oss, &object_key).await?;
Ok(json_success_body(
Some(&ctx),
json!({
"versionId": version_id,
"status": version.status,
"chunkBytes": PACKAGE_UPLOAD_CHUNK_BYTES,
"declaredPackageBytes": version.package_bytes,
"receivedBytes": received_bytes,
}),
))
}
/// 分片写入。
///
/// 客户端声明的偏移必须等于服务端已收字节;不一致时返回 409 与权威偏移,
/// 由客户端按权威偏移续传 —— 这样重放与乱序都不会造成重复写入。
async fn upload_package_chunk(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(version_id): Path,
body: Bytes,
) -> Result, AppError> {
require_octet_stream_content_type(&headers, "发行包分片必须使用 application/octet-stream")?;
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
// 分片级重放由偏移语义保证,这里仍要求幂等键,保持与其它写入口一致的调用约定。
let _idempotency_key = idempotency_key(&headers)?;
let offset = package_upload_offset(&headers, "发行包")?;
if body.is_empty() {
return Err(bad_request("发行包分片内容不能为空"));
}
if body.len() > PACKAGE_UPLOAD_CHUNK_BYTES {
return Err(AppError::from_status(StatusCode::PAYLOAD_TOO_LARGE)
.with_code("PACKAGE_CHUNK_TOO_LARGE")
.with_message("发行包分片超过服务端下发的大小"));
}
let version = load_owner_version_or_404(&state, owner_user_id, version_id.clone()).await?;
let chunk_bytes = u64::try_from(body.len()).unwrap_or(u64::MAX);
let end = offset
.checked_add(chunk_bytes)
.ok_or_else(|| bad_request("发行包分片偏移溢出"))?;
if end > version.package_bytes {
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PACKAGE_UPLOAD_EXCEEDS_DECLARED")
.with_details(json!({
"provider": "game-distribution",
"declaredPackageBytes": version.package_bytes,
"receivedBytes": offset,
"message": "分片写入会超过版本声明的发行包大小",
})));
}
let oss = game_distribution_oss_client(&state)?;
let object_key = game_distribution_package_object_key(&version.game_id, &version_id);
let received_bytes = staged_package_bytes(&state, oss, &object_key).await?;
if offset != received_bytes {
warn!(
request_id = ctx.request_id(),
operation = "package_chunk_offset_mismatch",
game_id = %version.game_id,
version_id = %version_id,
declared_offset = offset,
received_bytes,
"发行包分片偏移与服务端已收字节不一致"
);
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PACKAGE_UPLOAD_OFFSET_MISMATCH")
.with_message("分片偏移与服务端已收字节不一致,请按权威偏移续传")
.with_details(json!({
"provider": "game-distribution",
"receivedBytes": received_bytes,
})));
}
let append_result = oss
.append_internal_object_with_retry(
state.editor_oss_http_client(),
OssAppendInternalObjectRequest {
object_key: object_key.clone(),
content_type: Some("application/zip".to_string()),
access: OssObjectAccess::Private,
position: offset,
body: body.to_vec(),
},
GAME_DISTRIBUTION_OSS_PUT_MAX_ATTEMPTS,
&GAME_DISTRIBUTION_OSS_PUT_RETRY_DELAYS_MS,
)
.await;
let appended = match append_result {
Ok(appended) => appended,
Err(error) => {
// 追加失败也可能是「同偏移的并发写入先赢了一片」:先读权威已收字节,
// 只要长度已经前进就按偏移冲突返回,让客户端按权威偏移续传,
// 而不是把一个可恢复的并发结果报成上游故障。
if let Ok(authoritative) = staged_package_bytes(&state, oss, &object_key).await
&& authoritative > offset
{
warn!(
request_id = ctx.request_id(),
operation = "package_chunk_offset_lost_race",
game_id = %version.game_id,
version_id = %version_id,
declared_offset = offset,
received_bytes = authoritative,
"并发写入已推进已收字节,按偏移冲突返回权威位置"
);
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PACKAGE_UPLOAD_OFFSET_MISMATCH")
.with_message("分片偏移与服务端已收字节不一致,请按权威偏移续传")
.with_details(json!({
"provider": "game-distribution",
"receivedBytes": authoritative,
})));
}
return Err(map_oss_error(error, "aliyun-oss"));
}
};
info!(
request_id = ctx.request_id(),
operation = "package_chunk_stored",
game_id = %version.game_id,
version_id = %version_id,
offset,
chunk_bytes = appended.appended_bytes,
received_bytes = appended.next_position,
elapsed_ms = ctx.elapsed(),
"发行包分片已写入"
);
Ok(json_success_body(
Some(&ctx),
json!({
"versionId": version_id,
"chunkBytes": PACKAGE_UPLOAD_CHUNK_BYTES,
"receivedBytes": appended.next_position,
}),
))
}
/// 分片续传的完成动作:全部字节到齐后才回读整包、校验并确认。
///
/// 校验失败时删除半包对象并把版本落到 `upload_failed`,避免半包留在对象键上拖住后续重传。
async fn complete_package_upload(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(version_id): Path,
) -> Result, AppError> {
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let idempotency_key = idempotency_key(&headers)?;
let version =
load_owner_version_or_404(&state, owner_user_id.clone(), version_id.clone()).await?;
let oss = game_distribution_oss_client(&state)?;
let object_key = game_distribution_package_object_key(&version.game_id, &version_id);
let received_bytes = staged_package_bytes(&state, oss, &object_key).await?;
if received_bytes == 0 {
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PACKAGE_UPLOAD_NOT_STARTED")
.with_details(json!({
"provider": "game-distribution",
"declaredPackageBytes": version.package_bytes,
"receivedBytes": 0,
"message": "该版本还没有任何已收分片",
})));
}
if received_bytes != version.package_bytes {
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PACKAGE_UPLOAD_INCOMPLETE")
.with_message("发行包分片尚未收齐")
.with_details(json!({
"provider": "game-distribution",
"declaredPackageBytes": version.package_bytes,
"receivedBytes": received_bytes,
})));
}
let body = oss
.get_object(
state.editor_oss_http_client(),
OssGetObjectRequest {
object_key: object_key.clone(),
max_bytes: MAX_PACKAGE_BYTES as usize,
},
)
.await
.map_err(|error| map_oss_error(error, "aliyun-oss"))?;
let manifest = match validate_release_zip(&body) {
Ok(manifest) => manifest,
Err(error) => {
let reason = format!("{error:?}");
warn!(
request_id = ctx.request_id(),
operation = "package_rejected",
game_id = %version.game_id,
version_id = %version_id,
code = "PACKAGE_VALIDATION_FAILED",
reason = %reason,
uploaded_bytes = body.len(),
elapsed_ms = ctx.elapsed(),
"发行包校验失败"
);
let mapped = map_package_error(error);
if let Err(delete_error) = oss
.delete_object(
state.editor_oss_http_client(),
OssDeleteObjectRequest {
object_key: object_key.clone(),
},
)
.await
{
warn!(
request_id = ctx.request_id(),
operation = "package_staging_delete_failed",
version_id = %version_id,
error = %delete_error,
"校验失败的半包对象删除失败,需要人工确认对象键状态"
);
}
record_upload_failure(
&state,
&owner_user_id,
&version_id,
&idempotency_key,
"PACKAGE_VALIDATION_FAILED",
reason,
)
.await;
return Err(mapped);
}
};
confirm_validated_package(
&state,
&ctx,
&owner_user_id,
&version_id,
&version,
&manifest,
object_key,
&idempotency_key,
true,
)
.await
}
/// 显式重置分片会话:删除半包对象并把已收字节归零。
///
/// 只有尚未确认过发行包的版本能重置;已确认的版本必须新建版本,不能在半包之上续写不同字节。
async fn reset_package_upload(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(version_id): Path,
) -> Result, AppError> {
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let _idempotency_key = idempotency_key(&headers)?;
let version = load_owner_version_or_404(&state, owner_user_id, version_id.clone()).await?;
if !matches!(version.status.as_str(), "awaiting_upload" | "upload_failed") {
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PACKAGE_UPLOAD_RESET_NOT_ALLOWED")
.with_message("该版本已经确认过发行包,重新上传请新建版本"));
}
let oss = game_distribution_oss_client(&state)?;
let object_key = game_distribution_package_object_key(&version.game_id, &version_id);
oss.delete_object(
state.editor_oss_http_client(),
OssDeleteObjectRequest {
object_key: object_key.clone(),
},
)
.await
.map_err(|error| map_oss_error(error, "aliyun-oss"))?;
info!(
request_id = ctx.request_id(),
operation = "package_upload_reset",
game_id = %version.game_id,
version_id = %version_id,
elapsed_ms = ctx.elapsed(),
"发行包分片会话已重置"
);
Ok(json_success_body(
Some(&ctx),
json!({ "versionId": version_id, "receivedBytes": 0 }),
))
}
/// 工程源包上传的阶段门:5 条路由共用这一条规则,禁止在 handler 里各写一遍。
///
/// 两道判定缺一不可:
/// - **先判「已确认过工程包」**(`project_bundle_bytes > 0`)→ 409,文案必须含「已存在」,
/// 与 `map_spacetime_error` 的「已存在」子串映射同口径。顺序不能颠倒:确认工程包**不驱动**
/// 版本状态机,已确认的版本仍可能停在 `awaiting_upload`,只判状态会把它当成可写,放任第二次
/// 上传覆盖已确认的摘要与字节。
/// - **再判阶段**:只有 `awaiting_upload` / `upload_failed` 能写(与发行包确认同一道门);
/// 已提交、验证中、待审核、已拒绝、已公开、已撤回、已取消一律 409——版本不可变。
fn ensure_project_bundle_uploadable(
version: &GameDistributionVersionRecord,
) -> Result<(), AppError> {
if version.project_bundle_bytes > 0 {
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PROJECT_BUNDLE_ALREADY_EXISTS")
.with_message("同一版本已存在工程源包,换内容必须新建版本"));
}
if !matches!(version.status.as_str(), "awaiting_upload" | "upload_failed") {
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PROJECT_BUNDLE_UPLOAD_NOT_ALLOWED")
.with_message(format!(
"版本状态 {} 不允许上传工程源包,未公开前才能写一次",
version.status
)));
}
Ok(())
}
/// 工程源包校验失败的映射:与发行包同形(422 + 稳定错误码 + 具体原因),只是错误码换成工程包。
fn map_project_bundle_error(error: ProjectBundleError) -> AppError {
AppError::from_status(StatusCode::UNPROCESSABLE_ENTITY)
.with_code("PROJECT_BUNDLE_VALIDATION_FAILED")
.with_details(json!({ "provider": "game-distribution", "reason": format!("{error:?}") }))
}
/// 工程源包整包上传(一次 PUT):语义逐条镜像发行包整包上传,只是资产与对象键换成工程源包。
///
/// 载体类型是 `application/octet-stream`(技术方案 §3.4):作者侧打包器已经产出 zip 字节,
/// 不需要客户端再声明 `application/zip`;**服务端仍独立跑工程包门禁**,不信任客户端。
async fn upload_project_bundle(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(version_id): Path,
body: Bytes,
) -> Result, AppError> {
require_octet_stream_content_type(&headers, "工程源包必须使用 application/octet-stream")?;
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let idempotency_key = idempotency_key(&headers)?;
let expected =
load_owner_version_or_404(&state, owner_user_id.clone(), version_id.clone()).await?;
ensure_project_bundle_uploadable(&expected)?;
let manifest = match validate_project_bundle_zip(&body) {
Ok(manifest) => manifest,
Err(error) => {
let reason = format!("{error:?}");
warn!(
request_id = ctx.request_id(),
operation = "project_bundle_rejected",
game_id = %expected.game_id,
version_id = %version_id,
code = "PROJECT_BUNDLE_VALIDATION_FAILED",
reason = %reason,
uploaded_bytes = body.len(),
elapsed_ms = ctx.elapsed(),
"工程源包校验失败"
);
let mapped = map_project_bundle_error(error);
record_upload_failure(
&state,
&owner_user_id,
&version_id,
&idempotency_key,
"PROJECT_BUNDLE_VALIDATION_FAILED",
reason,
)
.await;
return Err(mapped);
}
};
let bundle_object_key =
game_distribution_project_bundle_object_key(&expected.game_id, &version_id);
let oss = game_distribution_oss_client(&state)?;
let existing = oss
.head_internal_object(state.editor_oss_http_client(), &bundle_object_key)
.await
.map_err(|error| map_oss_error(error, "aliyun-oss"))?;
let skipped = match existing {
Some(existing) if existing.content_length == manifest.bundle_bytes => true,
Some(_) => {
let error = AppError::from_status(StatusCode::CONFLICT)
.with_code("PROJECT_BUNDLE_OBJECT_MISMATCH")
.with_message("工程源包对象已存在但体积不一致");
record_upload_failure(
&state,
&owner_user_id,
&version_id,
&idempotency_key,
"PROJECT_BUNDLE_OBJECT_MISMATCH",
"工程源包对象已存在但体积不一致".to_string(),
)
.await;
return Err(error);
}
None => false,
};
if !skipped {
// 重试口径与发行包一致:只重试 platform-oss 认定的可重试分类(传输/超时/408/429/5xx)。
oss.put_internal_object_with_retry(
state.editor_oss_http_client(),
OssInternalPutObjectRequest {
object_key: bundle_object_key.clone(),
content_type: Some("application/zip".to_string()),
access: OssObjectAccess::Private,
metadata: BTreeMap::new(),
body: body.to_vec(),
},
GAME_DISTRIBUTION_OSS_PUT_MAX_ATTEMPTS,
&GAME_DISTRIBUTION_OSS_PUT_RETRY_DELAYS_MS,
)
.await
.map_err(|error| map_oss_error(error, "aliyun-oss"))?;
}
confirm_validated_project_bundle(
&state,
&ctx,
&owner_user_id,
&version_id,
&expected.game_id,
&manifest,
bundle_object_key,
&idempotency_key,
skipped,
)
.await
}
/// 工程源包校验通过后的共同收口:整包 PUT 与分片 complete 共用这条路径。
///
/// 幂等摘要的组织方式与发行包 complete 同形(`(version_id, sha256)` 序列化后取摘要),
/// 只是字段换成工程源包;同 key 重放由模块事务按摘要识别并返回 `replayed = true`。
#[allow(clippy::too_many_arguments)]
async fn confirm_validated_project_bundle(
state: &AppState,
ctx: &RequestContext,
owner_user_id: &str,
version_id: &str,
game_id: &str,
manifest: &ProjectBundleManifest,
bundle_object_key: String,
idempotency_key: &str,
oss_put_skipped: bool,
) -> Result, AppError> {
let request_digest = compute_request_digest(
&serde_json::to_vec(&(version_id, manifest.bundle_sha256.as_str()))
.map_err(|error| internal(error.to_string()))?,
);
let log_bundle_bytes = manifest.bundle_bytes;
let log_file_count = u32::try_from(manifest.files.len()).unwrap_or(u32::MAX);
let log_sha_prefix = manifest.bundle_sha256.chars().take(12).collect::();
let confirmed = state
.spacetime_client()
.confirm_game_distribution_project_bundle(
spacetime_client::GameDistributionConfirmProjectBundleRecordInput {
version_id: version_id.to_string(),
owner_user_id: owner_user_id.to_string(),
project_bundle_object_key: bundle_object_key,
project_bundle_bytes: manifest.bundle_bytes,
project_bundle_sha256: manifest.bundle_sha256.clone(),
idempotency_key: idempotency_key.to_string(),
request_digest,
updated_at_micros: now_micros(),
},
)
.await
.map_err(map_spacetime_error)?;
info!(
request_id = ctx.request_id(),
operation = "project_bundle_confirmed",
game_id = %game_id,
version_id = %confirmed.0.version_id,
project_bundle_bytes = log_bundle_bytes,
file_count = log_file_count,
sha256_prefix = %log_sha_prefix,
replayed = confirmed.1,
oss_put_skipped,
elapsed_ms = ctx.elapsed(),
"工程源包已确认"
);
Ok(json_success_body(
Some(ctx),
json!({ "versionId": confirmed.0.version_id, "status": confirmed.0.status }),
))
}
/// 工程源包分片续传的状态查询:已收字节同样取自 OSS 对象事实,因此进程重启、换机器或换网络
/// 后都能从权威偏移继续。工程源包没有「创建版本时预登记的大小」,因此响应里没有 declared* 键。
async fn project_bundle_upload_state(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
Path(version_id): Path,
) -> Result, AppError> {
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let version = load_owner_version_or_404(&state, owner_user_id, version_id.clone()).await?;
ensure_project_bundle_uploadable(&version)?;
let oss = game_distribution_oss_client(&state)?;
let object_key = game_distribution_project_bundle_object_key(&version.game_id, &version_id);
let received_bytes = staged_package_bytes(&state, oss, &object_key).await?;
Ok(json_success_body(
Some(&ctx),
json!({
"versionId": version_id,
"status": version.status,
"chunkBytes": PACKAGE_UPLOAD_CHUNK_BYTES,
"receivedBytes": received_bytes,
}),
))
}
/// 工程源包分片写入:偏移语义、分片边界与并发处理逐条对齐发行包分片。
///
/// 工程源包没有预登记大小,因此「不得超过」的闸门是合同上限 `MAX_PROJECT_BUNDLE_BYTES`
/// (与发行包上限同值);真实体积由 complete 时的整包校验确定。
async fn upload_project_bundle_chunk(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(version_id): Path,
body: Bytes,
) -> Result, AppError> {
require_octet_stream_content_type(&headers, "工程源包分片必须使用 application/octet-stream")?;
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
// 分片级重放由偏移语义保证,这里仍要求幂等键,保持与其它写入口一致的调用约定。
let _idempotency_key = idempotency_key(&headers)?;
let offset = package_upload_offset(&headers, "工程源包")?;
if body.is_empty() {
return Err(bad_request("工程源包分片内容不能为空"));
}
if body.len() > PACKAGE_UPLOAD_CHUNK_BYTES {
return Err(AppError::from_status(StatusCode::PAYLOAD_TOO_LARGE)
.with_code("PROJECT_BUNDLE_CHUNK_TOO_LARGE")
.with_message("工程源包分片超过服务端下发的大小"));
}
let version = load_owner_version_or_404(&state, owner_user_id, version_id.clone()).await?;
ensure_project_bundle_uploadable(&version)?;
let chunk_bytes = u64::try_from(body.len()).unwrap_or(u64::MAX);
let end = offset
.checked_add(chunk_bytes)
.ok_or_else(|| bad_request("工程源包分片偏移溢出"))?;
if end > MAX_PROJECT_BUNDLE_BYTES {
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PROJECT_BUNDLE_UPLOAD_EXCEEDS_LIMIT")
.with_details(json!({
"provider": "game-distribution",
"maxProjectBundleBytes": MAX_PROJECT_BUNDLE_BYTES,
"receivedBytes": offset,
"message": "分片写入会超过工程源包体积上限",
})));
}
let oss = game_distribution_oss_client(&state)?;
let object_key = game_distribution_project_bundle_object_key(&version.game_id, &version_id);
let received_bytes = staged_package_bytes(&state, oss, &object_key).await?;
if offset != received_bytes {
warn!(
request_id = ctx.request_id(),
operation = "project_bundle_chunk_offset_mismatch",
game_id = %version.game_id,
version_id = %version_id,
declared_offset = offset,
received_bytes,
"工程源包分片偏移与服务端已收字节不一致"
);
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PROJECT_BUNDLE_UPLOAD_OFFSET_MISMATCH")
.with_message("分片偏移与服务端已收字节不一致,请按权威偏移续传")
.with_details(json!({
"provider": "game-distribution",
"receivedBytes": received_bytes,
})));
}
let append_result = oss
.append_internal_object_with_retry(
state.editor_oss_http_client(),
OssAppendInternalObjectRequest {
object_key: object_key.clone(),
content_type: Some("application/zip".to_string()),
access: OssObjectAccess::Private,
position: offset,
body: body.to_vec(),
},
GAME_DISTRIBUTION_OSS_PUT_MAX_ATTEMPTS,
&GAME_DISTRIBUTION_OSS_PUT_RETRY_DELAYS_MS,
)
.await;
let appended = match append_result {
Ok(appended) => appended,
Err(error) => {
// 同偏移的并发写入可能先赢了一片:只要权威长度已经前进,就按偏移冲突返回,
// 让客户端按权威偏移续传,而不是把一个可恢复的并发结果报成上游故障。
if let Ok(authoritative) = staged_package_bytes(&state, oss, &object_key).await
&& authoritative > offset
{
warn!(
request_id = ctx.request_id(),
operation = "project_bundle_chunk_offset_lost_race",
game_id = %version.game_id,
version_id = %version_id,
declared_offset = offset,
received_bytes = authoritative,
"并发写入已推进已收字节,按偏移冲突返回权威位置"
);
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PROJECT_BUNDLE_UPLOAD_OFFSET_MISMATCH")
.with_message("分片偏移与服务端已收字节不一致,请按权威偏移续传")
.with_details(json!({
"provider": "game-distribution",
"receivedBytes": authoritative,
})));
}
return Err(map_oss_error(error, "aliyun-oss"));
}
};
info!(
request_id = ctx.request_id(),
operation = "project_bundle_chunk_stored",
game_id = %version.game_id,
version_id = %version_id,
offset,
chunk_bytes = appended.appended_bytes,
received_bytes = appended.next_position,
elapsed_ms = ctx.elapsed(),
"工程源包分片已写入"
);
Ok(json_success_body(
Some(&ctx),
json!({
"versionId": version_id,
"chunkBytes": PACKAGE_UPLOAD_CHUNK_BYTES,
"receivedBytes": appended.next_position,
}),
))
}
/// 工程源包分片续传的完成动作:全部字节到齐后才回读整包、独立校验并确认。
///
/// 校验失败时照抄发行包 complete 的处理:删除半包对象、记一次上传失败、返回既有错误形状
/// (422 + `PROJECT_BUNDLE_VALIDATION_FAILED`),避免半包留在对象键上拖住后续重传。
async fn complete_project_bundle_upload(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(version_id): Path,
) -> Result, AppError> {
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let idempotency_key = idempotency_key(&headers)?;
let version =
load_owner_version_or_404(&state, owner_user_id.clone(), version_id.clone()).await?;
ensure_project_bundle_uploadable(&version)?;
let oss = game_distribution_oss_client(&state)?;
let object_key = game_distribution_project_bundle_object_key(&version.game_id, &version_id);
let received_bytes = staged_package_bytes(&state, oss, &object_key).await?;
if received_bytes == 0 {
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PROJECT_BUNDLE_UPLOAD_NOT_STARTED")
.with_details(json!({
"provider": "game-distribution",
"receivedBytes": 0,
"message": "该版本还没有任何已收工程源包分片",
})));
}
let body = oss
.get_object(
state.editor_oss_http_client(),
OssGetObjectRequest {
object_key: object_key.clone(),
max_bytes: MAX_PROJECT_BUNDLE_BYTES as usize,
},
)
.await
.map_err(|error| map_oss_error(error, "aliyun-oss"))?;
// 工程源包没有预登记大小,「收齐」只能由「HEAD 的权威长度 == 读回的字节数」证明;
// 两者不一致说明读回期间对象被并发改写,按未收齐拒绝,让客户端重新对齐偏移。
if u64::try_from(body.len()).unwrap_or(u64::MAX) != received_bytes {
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_code("PROJECT_BUNDLE_UPLOAD_INCOMPLETE")
.with_message("工程源包分片尚未收齐")
.with_details(json!({
"provider": "game-distribution",
"receivedBytes": received_bytes,
"readBytes": body.len(),
})));
}
let manifest = match validate_project_bundle_zip(&body) {
Ok(manifest) => manifest,
Err(error) => {
let reason = format!("{error:?}");
warn!(
request_id = ctx.request_id(),
operation = "project_bundle_rejected",
game_id = %version.game_id,
version_id = %version_id,
code = "PROJECT_BUNDLE_VALIDATION_FAILED",
reason = %reason,
uploaded_bytes = body.len(),
elapsed_ms = ctx.elapsed(),
"工程源包校验失败"
);
let mapped = map_project_bundle_error(error);
if let Err(delete_error) = oss
.delete_object(
state.editor_oss_http_client(),
OssDeleteObjectRequest {
object_key: object_key.clone(),
},
)
.await
{
warn!(
request_id = ctx.request_id(),
operation = "project_bundle_staging_delete_failed",
version_id = %version_id,
error = %delete_error,
"校验失败的半包对象删除失败,需要人工确认对象键状态"
);
}
record_upload_failure(
&state,
&owner_user_id,
&version_id,
&idempotency_key,
"PROJECT_BUNDLE_VALIDATION_FAILED",
reason,
)
.await;
return Err(mapped);
}
};
confirm_validated_project_bundle(
&state,
&ctx,
&owner_user_id,
&version_id,
&version.game_id,
&manifest,
object_key,
&idempotency_key,
true,
)
.await
}
/// 显式重置工程源包分片会话:删除暂存对象并把已收字节归零。
///
/// 与发行包 reset 同一道门(共用 `ensure_project_bundle_uploadable`):只有尚未确认过工程源包
/// 且仍处于上传档位的版本能重置;已确认的版本换内容必须新建版本。
async fn reset_project_bundle_upload(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(version_id): Path,
) -> Result, AppError> {
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let _idempotency_key = idempotency_key(&headers)?;
let version = load_owner_version_or_404(&state, owner_user_id, version_id.clone()).await?;
ensure_project_bundle_uploadable(&version)?;
let oss = game_distribution_oss_client(&state)?;
let object_key = game_distribution_project_bundle_object_key(&version.game_id, &version_id);
oss.delete_object(
state.editor_oss_http_client(),
OssDeleteObjectRequest {
object_key: object_key.clone(),
},
)
.await
.map_err(|error| map_oss_error(error, "aliyun-oss"))?;
info!(
request_id = ctx.request_id(),
operation = "project_bundle_upload_reset",
game_id = %version.game_id,
version_id = %version_id,
elapsed_ms = ctx.elapsed(),
"工程源包分片会话已重置"
);
Ok(json_success_body(
Some(&ctx),
json!({ "versionId": version_id, "receivedBytes": 0 }),
))
}
fn game_distribution_oss_client(state: &AppState) -> Result<&platform_oss::OssClient, AppError> {
state.project_snapshot_oss_client().ok_or_else(|| {
AppError::from_status(StatusCode::SERVICE_UNAVAILABLE).with_message("游戏发行包 OSS 未配置")
})
}
fn game_distribution_package_object_key(game_id: &str, version_id: &str) -> String {
format!("{GAME_DISTRIBUTION_OBJECT_PREFIX}{game_id}/{version_id}.zip")
}
/// 工程源包对象键。
///
/// **必须**与发行包键(`…/{version_id}.zip`)不同:同一 (作品, 版本) 的两份资产(成品包与
/// 工程源包)会先后上传,共用键会让后传的那份覆盖前一份,已确认的摘要与字节随即变成谎话,
/// 发行网关与取件通道也会读到另一份资产。这里靠 `.project.zip` 后缀区分,两个键都在同一
/// 前缀族下,便于生命周期策略统一。
fn game_distribution_project_bundle_object_key(game_id: &str, version_id: &str) -> String {
format!("{GAME_DISTRIBUTION_OBJECT_PREFIX}{game_id}/{version_id}.project.zip")
}
/// 已收字节的权威来源:对象存在时的长度;确定不存在时是 0,其它失败按上游错误上报。
async fn staged_package_bytes(
state: &AppState,
oss: &platform_oss::OssClient,
object_key: &str,
) -> Result {
let head = oss
.head_internal_object(state.editor_oss_http_client(), object_key)
.await
.map_err(|error| map_oss_error(error, "aliyun-oss"))?;
Ok(head.map(|object| object.content_length).unwrap_or(0))
}
/// `application/octet-stream` 是发行包分片与工程源包(整包与分片)共同的载体类型;
/// 错误文案由调用方给,避免把「发行包分片」这句话安在工程源包上。
fn require_octet_stream_content_type(
headers: &HeaderMap,
message: &'static str,
) -> Result<(), AppError> {
let content_type = headers
.get(header::CONTENT_TYPE)
.and_then(|value| value.to_str().ok())
.map(|value| {
value
.split(';')
.next()
.unwrap_or_default()
.trim()
.to_ascii_lowercase()
});
if content_type.as_deref() != Some("application/octet-stream") {
return Err(bad_request(message));
}
Ok(())
}
/// 分片偏移头:发行包分片与工程源包分片共用同一个头名与解析口径(两者上限同值,客户端
/// 只能有一套偏移语义);`asset` 只用于错误文案,避免把「发行包」安在工程源包上。
fn package_upload_offset(headers: &HeaderMap, asset: &'static str) -> Result {
let raw = headers
.get(PACKAGE_UPLOAD_OFFSET_HEADER)
.and_then(|value| value.to_str().ok())
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| bad_request(format!("缺少{asset}分片偏移")))?;
raw.parse::()
.map_err(|_| bad_request(format!("{asset}分片偏移必须是非负整数")))
}
async fn submit_version(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(version_id): Path,
Json(payload): Json,
) -> Result<(StatusCode, Json), AppError> {
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let idempotency_key = idempotency_key(&headers)?;
// 与其它作者入口同口径:版本不存在或不属于当前主体都按 404 处理,
// 不能用 403 区分“别人的版本”,否则送审入口会泄露版本是否存在。
let version =
load_owner_version_or_404(&state, owner_user_id.clone(), version_id.clone()).await?;
let game = state
.spacetime_client()
.get_game_distribution_game(GameDistributionGetGameRecordInput {
game_id: version.game_id.clone(),
owner_user_id: Some(owner_user_id.clone()),
})
.await
.map_err(map_spacetime_error)?
.ok_or_else(|| AppError::from_status(StatusCode::NOT_FOUND))?;
let request_digest = compute_request_digest(
&serde_json::to_vec(&(version_id.as_str(), payload.expected_publication_revision))
.map_err(|error| internal(error.to_string()))?,
);
let log_game_id = version.game_id.clone();
let log_version_number = version.version_number;
let log_revision = payload.expected_publication_revision;
let submitted = state
.spacetime_client()
.submit_game_distribution_version_for_review(GameDistributionSubmitReviewRecordInput {
version_id,
owner_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 = "version_submitted",
game_id = %log_game_id,
version_id = %submitted.0.version_id,
version_number = log_version_number,
publication_revision = log_revision,
replayed = submitted.1,
elapsed_ms = ctx.elapsed(),
"版本已送审"
);
Ok((
StatusCode::ACCEPTED,
json_success_body(
Some(&ctx),
json!({
"game": game_payload(&game),
"version": private_version_payload(&submitted.0),
"replayed": submitted.1,
}),
),
))
}
/// 作者回读单个版本的私有状态与恢复动作。
///
/// 版本不存在或不属于当前主体都返回 404,避免用错误码区分“别人的版本”和“不存在的版本”。
async fn get_owner_version(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
Path(version_id): Path,
) -> Result, AppError> {
let owner_user_id = auth.claims().user_id().to_string();
let version = load_owner_version_or_404(&state, owner_user_id.clone(), version_id).await?;
let game = state
.spacetime_client()
.get_game_distribution_game(GameDistributionGetGameRecordInput {
game_id: version.game_id.clone(),
owner_user_id: Some(owner_user_id),
})
.await
.map_err(map_spacetime_error)?
.ok_or_else(|| AppError::from_status(StatusCode::NOT_FOUND))?;
Ok(json_success_body(
Some(&ctx),
version_detail_payload(&version, &game),
))
}
/// 作者撤回尚未公开的版本。
///
/// 只能撤回自己名下、且未参与当前公开投影的版本;`expectedPublicationRevision` 以
/// 游戏公开修订号做 CAS,过期请求返回 409,已公开版本改用下架。
async fn cancel_version(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(version_id): Path,
Json(payload): Json,
) -> Result, AppError> {
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let idempotency_key = idempotency_key(&headers)?;
let version =
load_owner_version_or_404(&state, owner_user_id.clone(), version_id.clone()).await?;
if version.publication_revision != payload.expected_publication_revision {
return Err(
AppError::from_status(StatusCode::CONFLICT).with_details(json!({
"provider": "game-distribution",
"code": "PUBLICATION_CONFLICT",
"message": "游戏的公开修订号已变化,请刷新后重试",
})),
);
}
let game = state
.spacetime_client()
.get_game_distribution_game(GameDistributionGetGameRecordInput {
game_id: version.game_id.clone(),
owner_user_id: Some(owner_user_id.clone()),
})
.await
.map_err(map_spacetime_error)?
.ok_or_else(|| AppError::from_status(StatusCode::NOT_FOUND))?;
let reason = payload
.reason
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty());
let request_digest = compute_request_digest(
&serde_json::to_vec(&(
version_id.as_str(),
payload.expected_publication_revision,
reason,
))
.map_err(|error| internal(error.to_string()))?,
);
let (version, replayed) = state
.spacetime_client()
.cancel_game_distribution_version(GameDistributionCancelVersionRecordInput {
version_id,
owner_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 = "version_cancelled",
game_id = %version.game_id,
version_id = %version.version_id,
version_number = version.version_number,
replayed,
elapsed_ms = ctx.elapsed(),
"版本已撤回"
);
Ok(json_success_body(
Some(&ctx),
json!({
"game": game_payload(&game),
"version": private_version_payload(&version),
"replayed": replayed,
}),
))
}
async fn unpublish_game(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(game_id): Path,
Json(payload): Json,
) -> Result, AppError> {
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let log_expected_revision = payload.expected_publication_revision;
let log_game_id = game_id.clone();
let idempotency_key = idempotency_key(&headers)?;
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 game = state
.spacetime_client()
.unpublish_game_distribution_game(GameDistributionUnpublishRecordInput {
game_id,
owner_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_unpublished",
game_id = %log_game_id,
expected_publication_revision = log_expected_revision,
visibility = %game.0.visibility,
publication_revision = game.0.publication_revision,
active_version_id = game.0.active_version_id.as_deref().unwrap_or(""),
replayed = game.1,
elapsed_ms = ctx.elapsed(),
"作者下架游戏,公开入口已关闭"
);
Ok(json_success_body(
Some(&ctx),
json!({ "game": game_payload(&game.0), "replayed": game.1 }),
))
}
/// 作者提升作品的共创授权档位:只升不降,降级与未知档位由领域层拒绝。
async fn set_fork_authorization(
State(state): State,
Extension(ctx): Extension,
Extension(auth): Extension,
headers: HeaderMap,
Path(game_id): Path,
payload: Result, JsonRejection>,
) -> Result, AppError> {
// 未知档位(例如 `"allowed"`)在进入业务前就被 serde 拦下:这里必须把它映射成平台信封的
// 400,而不是让 axum 的默认 `JsonRejection` 直接回 422 纯文本——合同要求未知档位是 400,
// 且前端错误处理依赖信封(同文件的评价保存、评价管理两处同写法)。
// 领域层的 `FORK_AUTHORIZATION_UNKNOWN`(map_spacetime_error 里映射 400)仍然可达:
// procedure 路径读到库里存的未知档位字符串时依旧由它兜底。
let Json(payload) = payload.map_err(|_| {
AppError::from_status(StatusCode::BAD_REQUEST)
.with_message("共创授权档位不合法,只接受 forbidden / nonCommercial / full")
})?;
let owner_user_id = auth.claims().user_id().to_string();
ensure_publish_enabled(&state, Some(owner_user_id.as_str())).await?;
let idempotency_key = idempotency_key(&headers)?;
let request_digest = compute_request_digest(
&serde_json::to_vec(&(
game_id.as_str(),
payload.expected_fork_authorization,
payload.fork_authorization,
))
.map_err(|error| internal(error.to_string()))?,
);
let game = state
.spacetime_client()
.set_game_distribution_fork_authorization(GameDistributionSetForkAuthorizationRecordInput {
game_id,
owner_user_id,
fork_authorization: fork_authorization_value(payload.fork_authorization),
expected_fork_authorization: fork_authorization_value(
payload.expected_fork_authorization,
),
idempotency_key,
request_digest,
now_micros: now_micros(),
})
.await
.map_err(map_spacetime_error)?;
Ok(json_success_body(
Some(&ctx),
json!({ "game": game_payload(&game.0), "replayed": game.1 }),
))
}
/// 取件校验通过后的目标:内容下发只需要**选定资产**的版本身份与摘要。
#[derive(Debug, PartialEq, Eq)]
struct ForkSourceTarget {
version_id: String,
/// 选定资产:有工程源包时 `Project`(优先),否则回落 `Package`。
source: GameDistributionForkSourceKind,
sha256: String,
bytes: u64,
}
/// 取件校验的纯映射:把「读到了什么」折成 HTTP 语义。
///
/// 顺序即合同顺序,也与 procedure 侧创建血缘时的判定顺序一致:行不存在 → 404;行存在但不可作
/// 来源(已软删除 / 未公开 / 没有当前公开版本)→ 409;授权为禁止或**未知档位** → 403
/// (未知按「禁止共创」解释,与 `resolve_game_distribution_fork_declaration_tx` 同口径)。
/// 抽成纯函数是为了让这条映射可被单测钉住,而不是散落在两个 handler 里各写一遍。
///
/// 选定资产(M2b):`project_bundle_bytes > 0 && project_bundle_sha256.is_some()` 才算「有工程
/// 源包」,此时优先 `Project`;否则回落 `Package`。**失败关闭**:工程包的字节数与摘要必须成对
/// ——「字节数 > 0 但摘要为空」这种半写行按「没有工程包」处理并回落成品包,而不是把取件指向一个
/// 摘不出来、客户端也无法校验的资产;没有任何可用资产时 409,不发半截信息。
fn fork_source_target(
record: GameDistributionForkSourceRecord,
) -> Result {
if !record.found {
return Err(AppError::from_status(StatusCode::NOT_FOUND).with_code("FORK_SOURCE_NOT_FOUND"));
}
if !record.available {
return Err(
AppError::from_status(StatusCode::CONFLICT).with_code("FORK_SOURCE_NOT_AVAILABLE")
);
}
let authorization =
module_game_distribution::ForkAuthorization::parse(record.fork_authorization.as_str())
.unwrap_or(module_game_distribution::ForkAuthorization::Forbidden);
if !authorization.allows_fork() {
return Err(AppError::from_status(StatusCode::FORBIDDEN).with_code("FORK_NOT_AUTHORIZED"));
}
let Some(version_id) = record.version_id else {
return Err(
AppError::from_status(StatusCode::CONFLICT).with_code("FORK_SOURCE_NOT_AVAILABLE")
);
};
// 有工程源包(字节数与摘要成对)→ 优先取源码包;半写行与缺失都按「没有工程包」处理。
if let (Some(sha256), bytes) = (record.project_bundle_sha256, record.project_bundle_bytes)
&& bytes > 0
{
return Ok(ForkSourceTarget {
version_id,
source: GameDistributionForkSourceKind::Project,
sha256,
bytes,
});
}
// 回落成品包:同样要求字节数与摘要成对,否则失败关闭。
let (Some(sha256), Some(bytes)) = (record.package_sha256, record.package_bytes) else {
return Err(
AppError::from_status(StatusCode::CONFLICT).with_code("FORK_SOURCE_NOT_AVAILABLE")
);
};
Ok(ForkSourceTarget {
version_id,
source: GameDistributionForkSourceKind::Package,
sha256,
bytes,
})
}
/// 取件通道的公共校验:两个 handler 都走它,规则只写一遍。
async fn resolve_fork_source(
state: &AppState,
game_id: &str,
) -> Result {
let record = state
.spacetime_client()
.get_game_distribution_fork_source(game_id.to_string())
.await
.map_err(map_spacetime_error)?;
fork_source_target(record)
}
/// 游戏标识必须能安全落在 URL 路径段里;发行入口与取件下载路径共用同一条判据。
fn is_path_safe_game_id(game_id: &str) -> bool {
!game_id.is_empty()
&& game_id.chars().all(|character| {
character.is_ascii_alphanumeric() || character == '-' || character == '_'
})
}
/// 取件下载路径:同源相对路径,**绝不下发 OSS 对象键**;路径按选定资产指向对应资产。
fn build_fork_source_download_path(
game_id: &str,
source: GameDistributionForkSourceKind,
) -> Result