diff --git a/docs/project-memory/shared-memory/decision-log.md b/docs/project-memory/shared-memory/decision-log.md index 7ae7628cd..66c524eec 100644 --- a/docs/project-memory/shared-memory/decision-log.md +++ b/docs/project-memory/shared-memory/decision-log.md @@ -7866,4 +7866,4 @@ CI 上 `background_agent_runtime_recovers_stale_running_before_pending_task` 在 - 2026-09-01 追加:`application.log` 不再写结构化错误事件;Rust `app_log!` 和 WebView console 都写入普通文本 raw log,结构化事件仅保留在当前进程内,提交时才生成 ZIP 内的 `events.jsonl`。 - 2026-09-01 review 收口:错误报告修复详情请求竞态、下载 anchor 生命周期、客户端采集脱敏/指纹降级与 4xx 噪声、用户级幂等隔离、`agc` 私有 OSS 前缀越权、日志读取链接检查、ZIP 同名日志和元数据/归档清理一致性;同步在 `review.txt` 标注仍需产品/运维决定的架构项。 - 2026-09-01 追加:api-server 按单实例部署,错误报告 store 保留进程内 Mutex 和同步本地文件 I/O,不引入跨进程锁;ZIP 仅在构建/上传阶段短暂驻留受 20 MiB 上限约束的内存 Vec,随后写入私有本地归档。管理员详情路由不属于 External OpenAPI;不存在返回 404,元数据/ZIP 损坏返回 500。 -- 2026-09-01 追加:错误报告元数据由 `error-reports/index.json` 索引,管理员列表使用 `limit/offset` 并返回 `total/hasMore`;索引缺失时从现有元数据一次性重建。OSS key 固定为 `agc/error-reports/v1/{batchId}.zip`,不含日期;相同 submission 对已 `ready` 报告直接返回。归档缺少 `events.jsonl` 视为损坏,管理员更新不存在目标返回 404。 +- 2026-09-01 追加:错误报告不落本地文件;请求内存构建 ZIP 后直接上传 OSS,成功后写入 SpacetimeDB `error_report` 元数据。OSS key 固定为 `agc/error-reports/v1/{batchId}.zip`,不含日期;同一用户 `userId + submissionId` 幂等。管理员查询 DB,详情/下载按 object key 读取 OSS;每日清理先删 OSS,再删 DB,失败留待下次重试。 diff --git a/docs/technical/【技术方案】AGC错误报告与诊断上传-2026-08-31.md b/docs/technical/【技术方案】AGC错误报告与诊断上传-2026-08-31.md index caae7986a..4c6f8e259 100644 --- a/docs/technical/【技术方案】AGC错误报告与诊断上传-2026-08-31.md +++ b/docs/technical/【技术方案】AGC错误报告与诊断上传-2026-08-31.md @@ -18,14 +18,15 @@ AI Game Creator Shell 采用 IDEA 风格的当前进程错误报告:错误事 - 登录态客户端使用 `POST /api/error-reports`,请求 DTO 位于 `shared-contracts::error_reports`。 - api-server 对请求体设置 24 MiB 上限,并校验 schemaVersion、submissionId、事件/日志数量和 20 MiB 压缩包上限;事件字段、用户说明和日志名/内容均做长度限制与基础脱敏,归档使用 `events.jsonl`(每行一个事件)。结构化事件只保存在当前进程内,用户提交时才生成 `events.jsonl`,不在磁盘单独持久化。submissionId 提供重放幂等。 -- 归档构建会短暂使用一个受 20 MiB 上限约束的内存 `Vec`,随后立即写入私有本地归档;不把 ZIP 长期留在内存。这样既控制峰值,又支持进程重启后后台查看/下载和 OSS 上传失败后的本地取证。 -- 本地 `error-reports/index.json` 保存元数据索引;创建时用索引完成 submission 幂等查找,管理员列表从索引读取并分页。索引缺失时会从现有元数据文件一次性重建。 -- 归档对象使用固定私有 OSS key:`agc/error-reports/v1/{batchId}.zip`;key 只由报告 UUID 决定,不包含时间戳。api-server 先写 `uploading` 元数据,上传成功后记录 `ossObjectKey`、SHA-256、大小和 `ready` 状态;相同 submission 重放对已 `ready` 报告直接返回,不重复上传。完整事件、说明和日志不进入元数据记录。 +- 归档构建只在请求生命周期内使用受 20 MiB 上限约束的内存 `Vec`,随后直接 PUT 到私有 OSS;服务端不写本地报告文件,也不保留本地索引。OSS 上传失败不写入数据库,调用方可稍后重新提交。 +- 归档对象使用固定私有 OSS key:`agc/error-reports/v1/{batchId}.zip`;key 只由报告 UUID 决定,不包含时间戳。上传成功后才写入 SpacetimeDB `error_report` 元数据表;`userId + submissionId` 由唯一幂等键保证重放返回已有记录。完整事件、说明和日志只存在 OSS ZIP。 - `agc` 是服务端专用私有前缀;公共直传票据、通用 object-key 规范化和 legacy 公开路径均拒绝该前缀。归档内同名日志会自动加数字后缀,读取本机诊断日志时拒绝符号链接/非普通文件。 - 后台接口:`GET/PATCH /admin/api/error-reports/{batchId}`、`GET /admin/api/error-reports` 和受保护的 `/download`。列表支持 `limit`/`offset` 分页并返回 `total`、`hasMore`。 - 这些是 api-server 内部登录/管理员路由,不属于 `/api/external/v1`,不纳入 External OpenAPI;管理员详情对不存在返回 404,对归档/元数据损坏返回 500。 - admin viewer 仅接受 error-reports Tab 权限,支持列表筛选、分页、详情、状态 `new/in-progress/resolved`、处理备注和受控下载;不存在的更新目标返回 404,存储损坏返回 500。列表行支持键盘 Enter/Space 打开详情,详情事件预览最多显示 20 条,完整内容通过诊断包下载获取。 -- 当前兼容实现仍在 api-server 配置目录旁保留元数据与本地归档副本,便于无 OSS 配置的开发环境运行;生产配置启用 OSS 后以 OSS 对象为完整内容来源。SpacetimeDB `error_report` 私有表接入及 30 天 OSS/元数据清理 worker 为后续门禁,HTTP DTO 与管理员权限保持不变。 +- 管理员列表、筛选、状态和备注全部读取/更新 SpacetimeDB;详情先读表再从 OSS 下载并解析 ZIP,下载接口直接从 OSS 返回 ZIP。无需新增管理员 DELETE HTTP 接口。每日清理任务删除过期 OSS 对象(成功或对象不存在后再删 DB;失败保留 DB 供下次重试)。旧本地报告不迁移。 + +SpacetimeDB `error_report` 表字段:`batch_id` 主键、`user_id`、`submission_id`、`idempotency_key` 唯一键、`object_key`、`archive_sha256`、`archive_size_bytes`、`event_count`、`log_count`、首个 fingerprint/source、`review_status`、`admin_note`、`created_at`、`updated_at`;索引为 `(user_id, submission_id)`、`created_at`、`review_status`。 ## 验收 diff --git a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md index 672d6c999..6d3292106 100644 --- a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md +++ b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md @@ -633,6 +633,13 @@ Responses 的终态载荷既是工具调用的恢复源,也是正文的恢复 - 说明:外部 OpenAPI 调用使用的账号级 API Key 凭据表,只保存 key prefix、SHA-256 hash、作用域、撤销状态和使用时间;明文 Key 只在 `/api/profile/api-keys` 创建接口返回一次,不进入 SpacetimeDB,且 API Key 管理接口不写入外部 OpenAPI JSON。v1 默认作用域为 `editor:project`、`editor:canvas`、`editor:image-generate`、`editor:asset`;其中 `editor:project` 覆盖项目列表、最近项目、创建、读取、重命名和删除,`editor:canvas` 覆盖默认画布布局保存,`editor:image-generate` 覆盖编辑器现有图片生成、重绘 / 调整、去背景、规范图、宣发素材、图标 spritesheet 生成 / 拆分、UI 设计图素材拆分、角色动画、视频、音效和背景音乐生成,`editor:asset` 覆盖素材直传凭证、素材对象确认、签名读取、账号级素材库和项目画布资源记录操作。 - 索引:`by_external_api_key_owner_user_id` 用于登录态 API Key 列表;`key_hash` 唯一索引用于外部 API 鉴权。 +### `error_report` + +- Rust 结构体:`ErrorReport` +- 源码:`server-rs/crates/spacetime-module/src/error_report.rs` +- 说明:错误报告只在请求内存中构建 ZIP 并上传私有 OSS;SpacetimeDB 仅保存 batch、幂等、对象键、摘要、计数和管理员审核元数据。管理员列表/筛选/更新走 `spacetime-client` facade,详情和下载再从 OSS 读取 ZIP。OSS key 固定为 `agc/error-reports/v1/{batchId}.zip`,不含时间戳;上传失败不写表。 +- 索引:`by_error_report_user_submission` 用于用户提交幂等;`by_error_report_created_at`、`by_error_report_review_status` 用于后台查询与清理。 + ### `admin_account` - Rust 结构体:`AdminAccount` diff --git a/server-rs/crates/api-server/src/error_reports.rs b/server-rs/crates/api-server/src/error_reports.rs index 677ea781a..4e781b9fd 100644 --- a/server-rs/crates/api-server/src/error_reports.rs +++ b/server-rs/crates/api-server/src/error_reports.rs @@ -1,36 +1,3 @@ -use std::{ - collections::HashSet, - fs, - io::{Cursor, Write}, - path::{Path, PathBuf}, - sync::Arc, -}; - -use axum::{ - Json, Router, - extract::{DefaultBodyLimit, Extension, Path as AxumPath, Query, State}, - http::{HeaderValue, StatusCode, header}, - middleware, - response::Response, - routing::{get, post}, -}; -use serde::{Deserialize, Serialize}; -use serde_json::Value; -use sha2::{Digest, Sha256}; -use time::{OffsetDateTime, format_description::well_known::Rfc3339}; -use tokio::sync::Mutex; -use uuid::Uuid; -use zip::{ZipWriter, write::FileOptions}; - -use platform_oss::{LegacyAssetPrefix, OssObjectAccess, OssPutObjectRequest}; -#[cfg(test)] -use shared_contracts::error_reports::ErrorReportLogInput; -use shared_contracts::error_reports::{ - AdminErrorReportEntry, AdminErrorReportListQuery, AdminErrorReportListResponse, - AdminUpdateErrorReportRequest, CreateErrorReportBatchRequest, CreateErrorReportBatchResponse, - Event, -}; - use crate::{ admin::{AuthenticatedAdmin, require_admin_auth}, api_response::json_success_body, @@ -40,6 +7,36 @@ use crate::{ state::AppState, tracking::{TrackingEventDraft, record_tracking_event_after_success}, }; +use axum::{ + Json, Router, + extract::{DefaultBodyLimit, Extension, Path as AxumPath, Query, State}, + http::{HeaderValue, StatusCode, header}, + middleware, + response::Response, + routing::{get, post}, +}; +use platform_oss::{ + LegacyAssetPrefix, OssDeleteObjectRequest, OssGetObjectRequest, OssObjectAccess, + OssPutObjectRequest, +}; +use serde_json::Value; +use sha2::{Digest, Sha256}; +use shared_contracts::error_reports::{ + AdminErrorReportEntry, AdminErrorReportListQuery, AdminErrorReportListResponse, + AdminUpdateErrorReportRequest, CreateErrorReportBatchRequest, CreateErrorReportBatchResponse, + Event, +}; +use spacetime_client::{ + ErrorReportCreateRecordInput, ErrorReportListRecordInput, ErrorReportUpdateRecordInput, +}; +use std::{ + collections::HashSet, + io::{Cursor, Read, Write}, +}; +use time::{OffsetDateTime, format_description::well_known::Rfc3339}; +use tokio::time::{Duration, interval}; +use uuid::Uuid; +use zip::{ZipWriter, write::FileOptions}; const MAX_EVENTS: usize = 100; const MAX_LOGS: usize = 5; @@ -49,71 +46,23 @@ const MAX_BATCH_BYTES: usize = 20 * 1024 * 1024; const MAX_REQUEST_BODY_BYTES: usize = 24 * 1024 * 1024; const MAX_EVENT_FIELD_CHARS: usize = 512; -pub struct ErrorReportStore { - directory: PathBuf, - lock: Mutex<()>, -} - -#[derive(Debug)] -enum ErrorReportGetError { - InvalidId, - NotFound(String), - Internal(String), -} - -#[derive(Debug)] -enum ErrorReportUpdateError { - InvalidId, - NotFound, - InvalidStatus, - Internal(String), -} - -impl std::fmt::Display for ErrorReportUpdateError { - fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - match self { - Self::InvalidId => formatter.write_str("报告标识无效"), - Self::NotFound => formatter.write_str("报告不存在"), - Self::InvalidStatus => formatter.write_str("状态必须是 new、in-progress 或 resolved"), - Self::Internal(message) => formatter.write_str(message), - } - } -} - -impl std::fmt::Display for ErrorReportGetError { - fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - match self { - Self::InvalidId => formatter.write_str("报告标识无效"), - Self::NotFound(message) | Self::Internal(message) => formatter.write_str(message), - } - } -} - -#[derive(Clone, Debug, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub struct StoredErrorReport { - pub batch_id: String, - pub submission_id: String, - pub user_id: String, - pub oss_object_key: Option, - pub archive_sha256: Option, - pub archive_size_bytes: u64, - pub event_count: u32, - pub upload_status: String, - pub review_status: String, - pub admin_note: Option, - pub log_count: u32, - pub first_fingerprint: Option, - pub first_source: Option, - pub created_at: String, - pub updated_at: String, -} - -#[derive(Clone, Debug, Serialize)] +#[derive(Clone, Debug, serde::Serialize)] #[serde(rename_all = "camelCase")] struct ErrorReportDetail { - #[serde(flatten)] - metadata: StoredErrorReport, + batch_id: String, + submission_id: String, + user_id: String, + oss_object_key: String, + archive_sha256: String, + archive_size_bytes: u64, + event_count: u32, + log_count: u32, + first_fingerprint: Option, + first_source: Option, + review_status: String, + admin_note: Option, + created_at: String, + updated_at: String, events: Vec, log_names: Vec, user_description: Option, @@ -158,1093 +107,523 @@ pub fn router(state: AppState) -> Router { pub async fn create_error_report( State(state): State, - Extension(request_context): Extension, - Extension(authenticated): Extension, - Json(payload): Json, + Extension(ctx): Extension, + Extension(auth): Extension, + Json(mut payload): Json, ) -> Result, AppError> { - let user_id = authenticated.claims().user_id().to_string(); - let accepted_event_ids = payload - .events - .iter() - .map(|event| event.event_id.clone()) - .collect::>(); - let stored = state - .error_report_store() - .create(user_id, payload) + let user_id = auth.claims().user_id().to_string(); + validate_payload(&mut payload)?; + let archive = build_error_report_archive( + &payload.events, + payload.user_description.as_deref(), + &payload + .logs + .iter() + .map(|l| { + ( + sanitize_log_name(&l.name), + sanitize_report_text_with_limit(&l.content, MAX_LOG_CHARS), + ) + }) + .collect::>(), + ) + .map_err(bad_request)?; + if archive.len() > MAX_BATCH_BYTES { + return Err(bad_request("报告附件超过 20 MB 上限")); + } + let batch_id = Uuid::new_v4().to_string(); + let object_key = format!("agc/error-reports/v1/{batch_id}.zip"); + let digest = format!("{:x}", Sha256::digest(&archive)); + let oss = state + .oss_client() + .ok_or_else(|| AppError::from_status(StatusCode::BAD_GATEWAY).with_message("OSS 未配置"))?; + oss.put_object( + state.editor_oss_http_client(), + OssPutObjectRequest { + prefix: LegacyAssetPrefix::AgcErrorReports, + path_segments: vec!["error-reports".into(), "v1".into()], + file_name: format!("{batch_id}.zip"), + content_type: Some("application/zip".into()), + access: OssObjectAccess::Private, + metadata: Default::default(), + body: archive.clone(), + }, + ) + .await + .map_err(|_| AppError::from_status(StatusCode::BAD_GATEWAY).with_message("诊断包上传失败"))?; + let record = match state + .spacetime_client() + .create_error_report(ErrorReportCreateRecordInput { + batch_id: batch_id.clone(), + user_id, + submission_id: payload.submission_id.clone(), + idempotency_key: format!("{}:{}", auth.claims().user_id(), payload.submission_id), + object_key: object_key.clone(), + archive_sha256: digest, + archive_size_bytes: archive.len() as u64, + event_count: payload.events.len() as u32, + log_count: payload.logs.len() as u32, + first_fingerprint: payload.events.first().map(|e| e.fingerprint.clone()), + first_source: payload.events.first().map(|e| e.source.clone()), + now_micros: now_micros(), + }) .await - .map_err(|message| AppError::from_status(StatusCode::BAD_REQUEST).with_message(message))?; - let stored = match state.oss_client() { - Some(_oss) if stored.upload_status == "ready" => stored, - Some(oss) => { - let archive = state - .error_report_store() - .read_archive(&stored.batch_id) - .await - .map_err(internal_store_error)?; - let digest = format!("{:x}", Sha256::digest(&archive)); - let response = oss - .put_object( + { + Ok(record) => record, + Err(_) => { + let _ = oss + .delete_object( state.editor_oss_http_client(), - OssPutObjectRequest { - prefix: LegacyAssetPrefix::AgcErrorReports, - path_segments: error_report_oss_path_segments(), - file_name: format!("{}.zip", stored.batch_id), - content_type: Some("application/zip".to_string()), - access: OssObjectAccess::Private, - metadata: Default::default(), - body: archive, + OssDeleteObjectRequest { + object_key: object_key.clone(), }, ) - .await - .map_err(|e| e.to_string()); - let response = match response { - Ok(response) => response, - Err(_error) => { - let _ = state - .error_report_store() - .mark_upload_failed(&stored.batch_id) - .await; - return Err(AppError::from_status(StatusCode::BAD_GATEWAY) - .with_message("诊断包上传失败")); - } - }; - state - .error_report_store() - .mark_uploaded(&stored.batch_id, response.object_key, digest) - .await - .map_err(internal_store_error)? + .await; + return Err(AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR) + .with_message("报告元数据写入失败")); } - None if stored.upload_status == "ready" => stored, - None => state - .error_report_store() - .mark_upload_failed(&stored.batch_id) - .await - .map_err(internal_store_error)?, }; + if record.batch_id != batch_id { + let _ = oss + .delete_object( + state.editor_oss_http_client(), + OssDeleteObjectRequest { object_key }, + ) + .await; + } Ok(json_success_body( - Some(&request_context), + Some(&ctx), CreateErrorReportBatchResponse { - batch_id: stored.batch_id, - submission_id: stored.submission_id, - status: stored.upload_status, - accepted_event_ids, - created_at: stored.created_at, + batch_id: record.batch_id, + submission_id: record.submission_id, + status: "ready".into(), + accepted_event_ids: payload.events.into_iter().map(|e| e.event_id).collect(), + created_at: record.created_at, }, )) } pub async fn admin_list_error_reports( State(state): State, - Extension(request_context): Extension, + Extension(ctx): Extension, Extension(admin): Extension, - Query(query): Query, + Query(q): Query, ) -> Result, AppError> { - let reports = state - .error_report_store() - .list(query) + let limit = q.limit.unwrap_or(100).clamp(1, 500); + let offset = q.offset.unwrap_or(0); + let (rows, total) = state + .spacetime_client() + .list_error_reports(ErrorReportListRecordInput { + status: q.status, + fingerprint: q.fingerprint, + source: q.source, + limit, + offset, + }) .await - .map_err(internal_store_error)?; - record_admin_report_audit( - &state, - &request_context, - admin.session().subject.as_str(), - "list", - None, - ) - .await; - Ok(json_success_body(Some(&request_context), reports)) + .map_err(internal)?; + let reports = rows + .into_iter() + .map(|r| AdminErrorReportEntry { + batch_id: r.batch_id, + event_count: r.event_count, + log_count: r.log_count, + fingerprint: r.first_fingerprint, + source: r.first_source, + status: r.review_status.clone(), + user_id: r.user_id, + created_at: r.created_at, + updated_at: r.updated_at, + user_description: None, + attachment_size_bytes: r.archive_size_bytes, + submission_id: Some(r.submission_id), + oss_object_key: Some(r.object_key), + archive_sha256: Some(r.archive_sha256), + upload_status: Some("ready".into()), + review_status: Some(r.review_status), + }) + .collect(); + record_admin_report_audit(&state, &ctx, admin.session().subject.as_str(), "list", None).await; + Ok(json_success_body( + Some(&ctx), + AdminErrorReportListResponse { + reports, + total, + offset, + limit, + has_more: offset.saturating_add(limit) < total as u32, + }, + )) } pub async fn admin_get_error_report( State(state): State, - Extension(request_context): Extension, + Extension(ctx): Extension, Extension(admin): Extension, AxumPath(batch_id): AxumPath, ) -> Result, AppError> { - let report = state - .error_report_store() - .get(&batch_id) + let r = state + .spacetime_client() + .get_error_report(batch_id.clone()) .await - .map_err(|error| match error { - ErrorReportGetError::InvalidId | ErrorReportGetError::NotFound(_) => { - AppError::from_status(StatusCode::NOT_FOUND).with_message("报告不存在") - } - ErrorReportGetError::Internal(message) => { - AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR).with_message(message) - } - })?; + .map_err(|_| AppError::from_status(StatusCode::NOT_FOUND).with_message("报告不存在"))?; + let bytes = state + .oss_client() + .ok_or_else(|| internal("OSS 未配置"))? + .get_object( + state.editor_oss_http_client(), + OssGetObjectRequest { + object_key: r.object_key.clone(), + max_bytes: MAX_BATCH_BYTES, + }, + ) + .await + .map_err(|_| AppError::from_status(StatusCode::NOT_FOUND).with_message("报告附件不存在"))?; + let (events, desc, logs) = + parse_archive(&bytes).map_err(|e| internal(format!("报告归档损坏:{e}")))?; record_admin_report_audit( &state, - &request_context, + &ctx, admin.session().subject.as_str(), "view", Some(&batch_id), ) .await; - Ok(json_success_body(Some(&request_context), report)) + Ok(json_success_body( + Some(&ctx), + ErrorReportDetail { + batch_id: r.batch_id, + submission_id: r.submission_id, + user_id: r.user_id, + oss_object_key: r.object_key, + archive_sha256: r.archive_sha256, + archive_size_bytes: r.archive_size_bytes, + event_count: r.event_count, + log_count: r.log_count, + first_fingerprint: r.first_fingerprint, + first_source: r.first_source, + review_status: r.review_status.clone(), + admin_note: r.admin_note.clone(), + created_at: r.created_at, + updated_at: r.updated_at, + events, + log_names: logs, + user_description: desc, + status: r.review_status, + note: r.admin_note, + attachment_size_bytes: r.archive_size_bytes, + }, + )) } pub async fn admin_update_error_report( State(state): State, - Extension(request_context): Extension, + Extension(ctx): Extension, Extension(admin): Extension, AxumPath(batch_id): AxumPath, - Json(payload): Json, + Json(p): Json, ) -> Result, AppError> { + if !matches!(p.status.as_str(), "new" | "in-progress" | "resolved") { + return Err(bad_request("状态必须是 new、in-progress 或 resolved")); + } state - .error_report_store() - .update(&batch_id, payload) + .spacetime_client() + .update_error_report(ErrorReportUpdateRecordInput { + batch_id: batch_id.clone(), + status: p.status, + note: p.note, + now_micros: now_micros(), + }) .await - .map_err(|error| match error { - ErrorReportUpdateError::InvalidId | ErrorReportUpdateError::NotFound => { - AppError::from_status(StatusCode::NOT_FOUND).with_message("报告不存在") - } - ErrorReportUpdateError::InvalidStatus => AppError::from_status(StatusCode::BAD_REQUEST) - .with_message("状态必须是 new、in-progress 或 resolved"), - ErrorReportUpdateError::Internal(_) => { - AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR) - .with_message("报告更新失败") - } - })?; - let report = state - .error_report_store() - .get(&batch_id) - .await - .map_err(|error| internal_store_error(error.to_string()))?; + .map_err(internal)?; record_admin_report_audit( &state, - &request_context, + &ctx, admin.session().subject.as_str(), "update", Some(&batch_id), ) .await; - Ok(json_success_body(Some(&request_context), report)) + admin_get_error_report( + State(state), + Extension(ctx), + Extension(admin), + AxumPath(batch_id), + ) + .await } pub async fn admin_download_error_report( State(state): State, - Extension(_request_context): Extension, + Extension(ctx): Extension, Extension(admin): Extension, AxumPath(batch_id): AxumPath, ) -> Result { + let r = state + .spacetime_client() + .get_error_report(batch_id.clone()) + .await + .map_err(|_| AppError::from_status(StatusCode::NOT_FOUND).with_message("报告不存在"))?; let bytes = state - .error_report_store() - .read_archive(&batch_id) + .oss_client() + .ok_or_else(|| internal("OSS 未配置"))? + .get_object( + state.editor_oss_http_client(), + OssGetObjectRequest { + object_key: r.object_key, + max_bytes: MAX_BATCH_BYTES, + }, + ) .await .map_err(|_| AppError::from_status(StatusCode::NOT_FOUND).with_message("报告附件不存在"))?; record_admin_report_audit( &state, - &_request_context, + &ctx, admin.session().subject.as_str(), "download", Some(&batch_id), ) .await; - let mut response = Response::new(bytes.into()); - response.headers_mut().insert( + let mut resp = Response::new(bytes.into()); + resp.headers_mut().insert( header::CONTENT_TYPE, HeaderValue::from_static("application/zip"), ); - response.headers_mut().insert( + resp.headers_mut().insert( header::CONTENT_DISPOSITION, - HeaderValue::from_str(&format!("attachment; filename=\"{batch_id}.zip\"")) - .map_err(|_| AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR))?, + HeaderValue::from_str(&format!("attachment; filename=\"{batch_id}.zip\"")).unwrap(), ); - Ok(response) + Ok(resp) } +fn now_micros() -> i64 { + (OffsetDateTime::now_utc().unix_timestamp_nanos() / 1_000) as i64 +} +fn bad_request(m: impl Into) -> AppError { + AppError::from_status(StatusCode::BAD_REQUEST).with_message(m) +} +fn internal(e: E) -> AppError { + AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR).with_message(e.to_string()) +} async fn record_admin_report_audit( state: &AppState, - request_context: &RequestContext, - admin_subject: &str, + ctx: &RequestContext, + subject: &str, action: &'static str, - batch_id: Option<&str>, + batch: Option<&str>, ) { - let mut draft = - TrackingEventDraft::user("admin_error_report_action", "error-reports", admin_subject); - draft.metadata = serde_json::json!({ "action": action, "batchId": batch_id }); - record_tracking_event_after_success(state, request_context, draft).await; + let mut d = TrackingEventDraft::user("admin_error_report_action", "error-reports", subject); + d.metadata = serde_json::json!({"action":action,"batchId":batch}); + record_tracking_event_after_success(state, ctx, d).await; } -fn internal_store_error(error: String) -> AppError { - AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR).with_message(error) -} - -fn error_report_oss_path_segments() -> Vec { - vec!["error-reports".to_string(), "v1".to_string()] +fn validate_payload(p: &mut CreateErrorReportBatchRequest) -> Result<(), AppError> { + if p.schema_version != 1 { + return Err(bad_request("报告协议版本不支持")); + } + p.submission_id = sanitize_report_text_with_limit(&p.submission_id, 128); + if p.submission_id.is_empty() { + return Err(bad_request("缺少提交幂等标识")); + } + if p.events.is_empty() || p.events.len() > MAX_EVENTS { + return Err(bad_request("报告事件数量必须在 1 到 100 之间")); + } + if p.logs.len() > MAX_LOGS { + return Err(bad_request("日志附件数量超出上限")); + } + for e in &mut p.events { + e.event_id = sanitize_report_text(&e.event_id); + e.fingerprint = sanitize_report_text(&e.fingerprint); + e.source = sanitize_report_text(&e.source); + e.message = sanitize_report_text(&e.message); + e.stack = e + .stack + .take() + .map(|v| sanitize_report_text_with_limit(&v, MAX_LOG_CHARS)); + e.occurred_at = sanitize_report_text(&e.occurred_at); + if e.event_id.is_empty() || e.fingerprint.is_empty() || e.message.is_empty() { + return Err(bad_request("报告事件缺少必要字段")); + } + if e.count == 0 { + e.count = 1; + } + } + let ids = p + .events + .iter() + .map(|e| e.event_id.as_str()) + .collect::>(); + if ids.len() != p.events.len() { + return Err(bad_request("报告事件标识不能重复")); + } + p.user_description = p + .user_description + .take() + .map(|v| sanitize_report_text_with_limit(&v, MAX_DESCRIPTION_CHARS)); + Ok(()) } pub(crate) fn build_error_report_archive( events: &[Event], - user_description: Option<&str>, + desc: Option<&str>, logs: &[(String, String)], ) -> Result, String> { - let cursor = Cursor::new(Vec::new()); - let mut zip = ZipWriter::new(cursor); - let options = FileOptions::<()>::default().compression_method(zip::CompressionMethod::Deflated); - let manifest = serde_json::json!({ - "schemaVersion": 1, - "eventCount": events.len(), - "logCount": logs.len(), - "userDescriptionIncluded": user_description.is_some(), - }); - zip.start_file("manifest.json", options) - .map_err(|error| error.to_string())?; - zip.write_all( - serde_json::to_string_pretty(&manifest) - .map_err(|error| error.to_string())? + let mut z = ZipWriter::new(Cursor::new(Vec::new())); + let o = FileOptions::<()>::default().compression_method(zip::CompressionMethod::Deflated); + z.start_file("manifest.json", o) + .map_err(|e| e.to_string())?; + z.write_all( + serde_json::json!({"schemaVersion":1,"eventCount":events.len(),"logCount":logs.len()}) + .to_string() .as_bytes(), ) - .map_err(|error| error.to_string())?; - zip.start_file("events.jsonl", options) - .map_err(|error| error.to_string())?; - for event in events { - let line = serde_json::to_string(event).map_err(|error| error.to_string())?; - zip.write_all(line.as_bytes()) - .map_err(|error| error.to_string())?; - zip.write_all(b"\n").map_err(|error| error.to_string())?; - } - if let Some(description) = user_description { - zip.start_file("user-description.txt", options) - .map_err(|error| error.to_string())?; - zip.write_all(description.as_bytes()) - .map_err(|error| error.to_string())?; - } - let mut used_log_names = HashSet::new(); - for (name, content) in logs { - let safe_name = name - .chars() - .filter(|character| { - character.is_ascii_alphanumeric() || matches!(character, '.' | '-' | '_') - }) - .collect::(); - if safe_name.is_empty() { - continue; - } - let mut archive_name = safe_name.clone(); - let mut suffix = 2; - while !used_log_names.insert(archive_name.clone()) { - archive_name = format!("{safe_name}.{suffix}"); - suffix += 1; - } - zip.start_file(format!("system-logs/{archive_name}"), options) - .map_err(|error| error.to_string())?; - zip.write_all(content.as_bytes()) - .map_err(|error| error.to_string())?; - } - zip.finish() - .map(|cursor| cursor.into_inner()) - .map_err(|error| error.to_string()) -} - -impl ErrorReportStore { - pub fn new(base_dir: impl Into) -> Arc { - Arc::new(Self { - directory: base_dir.into().join("error-reports"), - lock: Mutex::new(()), - }) - } - - async fn create( - &self, - user_id: String, - mut payload: CreateErrorReportBatchRequest, - ) -> Result { - if payload.schema_version != 1 { - return Err("报告协议版本不支持".to_string()); - } - payload.submission_id = sanitize_report_text_with_limit(&payload.submission_id, 128); - if payload.submission_id.is_empty() { - return Err("缺少提交幂等标识".to_string()); - } - if payload.events.is_empty() || payload.events.len() > MAX_EVENTS { - return Err("报告事件数量必须在 1 到 100 之间".to_string()); - } - if payload.logs.len() > MAX_LOGS { - return Err("日志附件数量超出上限".to_string()); - } - payload.user_description = payload.user_description.map(|value| { - sanitize_report_text_with_limit(&value, MAX_DESCRIPTION_CHARS) - .chars() - .take(MAX_DESCRIPTION_CHARS) - .collect() - }); - for event in &mut payload.events { - event.event_id = - sanitize_report_text_with_limit(&event.event_id, MAX_EVENT_FIELD_CHARS); - event.fingerprint = - sanitize_report_text_with_limit(&event.fingerprint, MAX_EVENT_FIELD_CHARS); - event.source = sanitize_report_text_with_limit(&event.source, MAX_EVENT_FIELD_CHARS); - event.message = sanitize_report_text_with_limit(&event.message, MAX_EVENT_FIELD_CHARS); - event.stack = event - .stack - .take() - .map(|value| sanitize_report_text_with_limit(&value, MAX_LOG_CHARS)); - event.occurred_at = - sanitize_report_text_with_limit(&event.occurred_at, MAX_EVENT_FIELD_CHARS); - if event.event_id.is_empty() || event.fingerprint.is_empty() || event.message.is_empty() - { - return Err("报告事件缺少必要字段".to_string()); - } - if event.count == 0 { - event.count = 1; - } - } - let event_ids = payload - .events - .iter() - .map(|event| event.event_id.as_str()) - .collect::>(); - if event_ids.len() != payload.events.len() { - return Err("报告事件标识不能重复".to_string()); - } - for log in &mut payload.logs { - log.name = sanitize_log_name(&log.name); - log.content = sanitize_report_text_with_limit(&log.content, MAX_LOG_CHARS) - .chars() - .take(MAX_LOG_CHARS) - .collect(); - } - let archive = build_error_report_archive( - &payload.events, - payload.user_description.as_deref(), - &payload - .logs - .iter() - .map(|log| (log.name.clone(), log.content.clone())) - .collect::>(), - )?; - if archive.len() > MAX_BATCH_BYTES { - return Err("报告附件超过 20 MB 上限".to_string()); - } - let now = OffsetDateTime::now_utc() - .format(&Rfc3339) - .map_err(|error| error.to_string())?; - let batch_id = Uuid::new_v4().to_string(); - let stored = StoredErrorReport { - batch_id: batch_id.clone(), - submission_id: payload.submission_id, - user_id, - oss_object_key: None, - archive_sha256: None, - archive_size_bytes: archive.len() as u64, - event_count: payload.events.len() as u32, - upload_status: "uploading".to_string(), - review_status: "new".to_string(), - admin_note: None, - log_count: payload.logs.len() as u32, - first_fingerprint: payload - .events - .first() - .map(|event| event.fingerprint.clone()), - first_source: payload.events.first().map(|event| event.source.clone()), - created_at: now.clone(), - updated_at: now, - }; - let _guard = self.lock.lock().await; - fs::create_dir_all(&self.directory).map_err(|error| error.to_string())?; - set_private_directory_permissions(&self.directory)?; - self.cleanup_expired(); - let mut index = self.load_index()?; - if let Some(existing) = index.iter().find(|existing| { - existing.user_id == stored.user_id && existing.submission_id == stored.submission_id - }) { - return Ok(existing.clone()); - } - let metadata_path = self.metadata_path(&batch_id); - let archive_path = self.archive_path(&batch_id); - write_atomic( - &metadata_path, - &serde_json::to_vec_pretty(&stored).map_err(|error| error.to_string())?, - )?; - write_atomic(&archive_path, &archive)?; - index.push(stored.clone()); - self.write_index(&index)?; - Ok(stored) - } - - async fn list( - &self, - query: AdminErrorReportListQuery, - ) -> Result { - let _guard = self.lock.lock().await; - self.cleanup_expired(); - let mut entries = self.load_index()?; - let limit = query.limit.unwrap_or(100).clamp(1, 500) as usize; - let offset = query.offset.unwrap_or(0).min(u32::MAX) as usize; - entries.retain(|report| { - if query - .status - .as_deref() - .is_some_and(|value| value != report.review_status) - || query - .source - .as_deref() - .is_some_and(|value| report.first_source.as_deref() != Some(value)) - || query - .fingerprint - .as_deref() - .is_some_and(|value| report.first_fingerprint.as_deref() != Some(value)) - { - false - } else { - true - } - }); - entries.sort_by(|left, right| { - OffsetDateTime::parse(&right.created_at, &Rfc3339) - .ok() - .cmp(&OffsetDateTime::parse(&left.created_at, &Rfc3339).ok()) - }); - let total = entries.len() as u64; - let reports = entries - .into_iter() - .skip(offset) - .take(limit) - .map(admin_error_report_entry) - .collect::>(); - Ok(AdminErrorReportListResponse { - reports, - total, - offset: offset as u32, - limit: limit as u32, - has_more: offset.saturating_add(limit) < total as usize, - }) - } - - async fn get(&self, batch_id: &str) -> Result { - validate_batch_id(batch_id).map_err(|_| ErrorReportGetError::InvalidId)?; - let _guard = self.lock.lock().await; - let bytes = fs::read(self.metadata_path(batch_id)).map_err(|error| { - if error.kind() == std::io::ErrorKind::NotFound { - ErrorReportGetError::NotFound(error.to_string()) - } else { - ErrorReportGetError::Internal(error.to_string()) - } - })?; - let metadata: StoredErrorReport = serde_json::from_slice(&bytes) - .map_err(|error| ErrorReportGetError::Internal(format!("报告元数据损坏:{error}")))?; - let archive = fs::read(self.archive_path(batch_id)).map_err(|error| { - if error.kind() == std::io::ErrorKind::NotFound { - ErrorReportGetError::NotFound(error.to_string()) - } else { - ErrorReportGetError::Internal(error.to_string()) - } - })?; - let (events, user_description, log_names) = parse_archive(&archive) - .map_err(|error| ErrorReportGetError::Internal(format!("报告归档损坏:{error}")))?; - let status = metadata.review_status.clone(); - let note = metadata.admin_note.clone(); - let attachment_size_bytes = metadata.archive_size_bytes; - Ok(ErrorReportDetail { - metadata, - events, - log_names, - user_description, - status, - note, - attachment_size_bytes, - }) - } - - async fn update( - &self, - batch_id: &str, - payload: AdminUpdateErrorReportRequest, - ) -> Result { - validate_batch_id(batch_id).map_err(|_| ErrorReportUpdateError::InvalidId)?; - if !matches!(payload.status.as_str(), "new" | "in-progress" | "resolved") { - return Err(ErrorReportUpdateError::InvalidStatus); - } - let _guard = self.lock.lock().await; - let path = self.metadata_path(batch_id); - let bytes = fs::read(&path).map_err(|error| { - if error.kind() == std::io::ErrorKind::NotFound { - ErrorReportUpdateError::NotFound - } else { - ErrorReportUpdateError::Internal(error.to_string()) - } - })?; - let mut report: StoredErrorReport = serde_json::from_slice(&bytes).map_err(|error| { - ErrorReportUpdateError::Internal(format!("报告元数据损坏:{error}")) - })?; - report.review_status = payload.status; - report.admin_note = payload - .note - .map(|value| value.chars().take(MAX_DESCRIPTION_CHARS).collect()); - report.updated_at = OffsetDateTime::now_utc() - .format(&Rfc3339) - .map_err(|error| ErrorReportUpdateError::Internal(error.to_string()))?; - write_atomic( - &path, - &serde_json::to_vec_pretty(&report) - .map_err(|error| ErrorReportUpdateError::Internal(error.to_string()))?, + .map_err(|e| e.to_string())?; + z.start_file("events.jsonl", o).map_err(|e| e.to_string())?; + for e in events { + z.write_all( + serde_json::to_string(e) + .map_err(|e| e.to_string())? + .as_bytes(), ) - .map_err(ErrorReportUpdateError::Internal)?; - refresh_archive_mtime(&self.archive_path(batch_id)); - self.replace_index_entry(report.clone()) - .map_err(ErrorReportUpdateError::Internal)?; - Ok(report) + .map_err(|e| e.to_string())?; + z.write_all(b"\n").map_err(|e| e.to_string())?; } - - async fn read_archive(&self, batch_id: &str) -> Result, String> { - validate_batch_id(batch_id)?; - let _guard = self.lock.lock().await; - fs::read(self.archive_path(batch_id)).map_err(|error| error.to_string()) - } - - async fn mark_uploaded( - &self, - batch_id: &str, - object_key: String, - sha256: String, - ) -> Result { - validate_batch_id(batch_id)?; - let _guard = self.lock.lock().await; - let path = self.metadata_path(batch_id); - let bytes = fs::read(&path).map_err(|error| error.to_string())?; - let mut report: StoredErrorReport = - serde_json::from_slice(&bytes).map_err(|error| error.to_string())?; - report.oss_object_key = Some(object_key); - report.archive_sha256 = Some(sha256); - report.upload_status = "ready".to_string(); - report.updated_at = OffsetDateTime::now_utc() - .format(&Rfc3339) + if let Some(d) = desc { + z.start_file("user-description.txt", o) .map_err(|e| e.to_string())?; - write_atomic( - &path, - &serde_json::to_vec_pretty(&report).map_err(|e| e.to_string())?, - )?; - self.replace_index_entry(report.clone())?; - refresh_archive_mtime(&self.archive_path(batch_id)); - Ok(report) + z.write_all(d.as_bytes()).map_err(|e| e.to_string())?; } - - async fn mark_upload_failed(&self, batch_id: &str) -> Result { - validate_batch_id(batch_id)?; - let _guard = self.lock.lock().await; - let path = self.metadata_path(batch_id); - let bytes = fs::read(&path).map_err(|error| error.to_string())?; - let mut report: StoredErrorReport = - serde_json::from_slice(&bytes).map_err(|error| error.to_string())?; - report.upload_status = "failed".to_string(); - report.updated_at = OffsetDateTime::now_utc() - .format(&Rfc3339) + let mut used = HashSet::new(); + for (n, c) in logs { + let s = sanitize_log_name(n); + let mut a = s.clone(); + let mut i = 2; + while !used.insert(a.clone()) { + a = format!("{s}.{i}"); + i += 1; + } + z.start_file(format!("system-logs/{a}"), o) .map_err(|e| e.to_string())?; - write_atomic( - &path, - &serde_json::to_vec_pretty(&report).map_err(|e| e.to_string())?, - )?; - self.replace_index_entry(report.clone())?; - refresh_archive_mtime(&self.archive_path(batch_id)); - Ok(report) - } - - fn metadata_path(&self, batch_id: &str) -> PathBuf { - self.directory.join(format!("{batch_id}.json")) - } - fn archive_path(&self, batch_id: &str) -> PathBuf { - self.directory.join(format!("{batch_id}.zip")) - } - - fn index_path(&self) -> PathBuf { - self.directory.join("index.json") - } - - fn load_index(&self) -> Result, String> { - let path = self.index_path(); - match fs::read(&path) { - Ok(bytes) => { - let reports = serde_json::from_slice::>(&bytes) - .map_err(|error| format!("错误报告索引损坏:{error}"))?; - Ok(reports - .into_iter() - .filter(|report| self.metadata_path(&report.batch_id).is_file()) - .collect()) - } - Err(error) if error.kind() == std::io::ErrorKind::NotFound => { - let mut reports = Vec::new(); - let Ok(directory) = fs::read_dir(&self.directory) else { - return Ok(reports); - }; - for item in directory.flatten() { - if item.path().extension().and_then(|value| value.to_str()) != Some("json") - || item.path() == path - { - continue; - } - let Ok(bytes) = fs::read(item.path()) else { - continue; - }; - if let Ok(report) = serde_json::from_slice::(&bytes) { - reports.push(report); - } - } - self.write_index(&reports)?; - Ok(reports) - } - Err(error) => Err(error.to_string()), - } - } - - fn write_index(&self, reports: &[StoredErrorReport]) -> Result<(), String> { - write_atomic( - &self.index_path(), - &serde_json::to_vec_pretty(reports).map_err(|error| error.to_string())?, - ) - } - - fn replace_index_entry(&self, report: StoredErrorReport) -> Result<(), String> { - let mut reports = self.load_index()?; - if let Some(existing) = reports - .iter_mut() - .find(|existing| existing.batch_id == report.batch_id) - { - *existing = report; - } else { - reports.push(report); - } - self.write_index(&reports) - } - - fn cleanup_expired(&self) { - let Ok(directory) = fs::read_dir(&self.directory) else { - return; - }; - let cutoff = std::time::SystemTime::now() - .checked_sub(std::time::Duration::from_secs(30 * 24 * 60 * 60)); - for item in directory.flatten() { - let path = item.path(); - if path.extension().and_then(|value| value.to_str()) != Some("json") - || path == self.index_path() - { - continue; - } - let Ok(metadata) = fs::symlink_metadata(&path) else { - continue; - }; - if metadata.file_type().is_symlink() || !metadata.is_file() { - continue; - } - let Some(cutoff) = cutoff else { continue }; - if metadata - .modified() - .ok() - .is_some_and(|modified| modified < cutoff) - { - let _ = fs::remove_file(&path); - let _ = fs::remove_file(path.with_extension("zip")); - } - } - if let Ok(reports) = self.load_index() { - let _ = self.write_index(&reports); - } + z.write_all(c.as_bytes()).map_err(|e| e.to_string())?; } + z.finish() + .map(|c| c.into_inner()) + .map_err(|e| e.to_string()) } - -fn admin_error_report_entry(report: StoredErrorReport) -> AdminErrorReportEntry { - AdminErrorReportEntry { - batch_id: report.batch_id, - event_count: report.event_count, - log_count: report.log_count, - fingerprint: report.first_fingerprint, - source: report.first_source, - status: report.review_status.clone(), - user_id: report.user_id, - created_at: report.created_at, - updated_at: report.updated_at, - user_description: None, - attachment_size_bytes: report.archive_size_bytes, - submission_id: Some(report.submission_id), - oss_object_key: report.oss_object_key, - archive_sha256: report.archive_sha256, - upload_status: Some(report.upload_status), - review_status: Some(report.review_status), - } -} - -fn refresh_archive_mtime(path: &Path) { - let Ok(metadata) = fs::symlink_metadata(path) else { - return; - }; - if metadata.file_type().is_symlink() || !metadata.is_file() { - return; - } - let Ok(file) = fs::OpenOptions::new().append(true).open(path) else { - return; - }; - let _ = file.set_modified(std::time::SystemTime::now()); -} - -fn write_atomic(path: &Path, bytes: &[u8]) -> Result<(), String> { - let temp_path = path.with_extension(format!("tmp-{}", Uuid::new_v4())); - fs::write(&temp_path, bytes).map_err(|error| error.to_string())?; - #[cfg(unix)] - { - use std::os::unix::fs::PermissionsExt; - fs::set_permissions(&temp_path, fs::Permissions::from_mode(0o600)) - .map_err(|error| error.to_string())?; - } - fs::rename(&temp_path, path).map_err(|error| error.to_string()) -} - -fn set_private_directory_permissions(path: &Path) -> Result<(), String> { - #[cfg(unix)] - { - use std::os::unix::fs::PermissionsExt; - fs::set_permissions(path, fs::Permissions::from_mode(0o700)) - .map_err(|error| error.to_string())?; - } - Ok(()) -} - -fn validate_batch_id(batch_id: &str) -> Result<(), String> { - if Uuid::parse_str(batch_id).is_ok() { - Ok(()) +fn parse_archive(bytes: &[u8]) -> Result<(Vec, Option, Vec), String> { + let mut a = zip::ZipArchive::new(Cursor::new(bytes)).map_err(|e| e.to_string())?; + let mut t = String::new(); + a.by_name("events.jsonl") + .map_err(|_| "归档缺少 events.jsonl".to_string())? + .read_to_string(&mut t) + .map_err(|e| e.to_string())?; + let events = t + .lines() + .filter(|l| !l.trim().is_empty()) + .map(|l| serde_json::from_str(l).map_err(|e| e.to_string())) + .collect::, _>>()?; + let desc = if let Ok(mut f) = a.by_name("user-description.txt") { + let mut s = String::new(); + f.read_to_string(&mut s).map_err(|e| e.to_string())?; + Some(s) } else { - Err("报告标识无效".to_string()) + None + }; + let mut names = Vec::new(); + for i in 0..a.len() { + let n = a.by_index(i).map_err(|e| e.to_string())?.name().to_string(); + if let Some(n) = n.strip_prefix("system-logs/") { + names.push(n.to_string()); + } } + Ok((events, desc, names)) } - -fn sanitize_report_text(value: &str) -> String { - sanitize_report_text_with_limit(value, MAX_EVENT_FIELD_CHARS) +fn sanitize_report_text(v: &str) -> String { + sanitize_report_text_with_limit(v, MAX_EVENT_FIELD_CHARS) } - -fn sanitize_report_text_with_limit(value: &str, max_chars: usize) -> String { - let mut sanitized = value.replace(['\r', '\n'], " "); - let lowercase = sanitized.to_ascii_lowercase(); - for marker in [ +fn sanitize_report_text_with_limit(v: &str, max: usize) -> String { + let l = v.to_ascii_lowercase(); + if [ "authorization:", "bearer ", "token=", "token:", "api_key=", "apikey=", - ] { - if lowercase.contains(marker) { - return "[REDACTED]".to_string(); - } + ] + .iter() + .any(|m| l.contains(m)) + { + return "[REDACTED]".into(); } - sanitized = sanitized.replace("/Users/", "/"); - sanitized = sanitized.replace("/home/", "/"); - sanitized = sanitized.replace("C:\\", "/"); - sanitized.chars().take(max_chars).collect() + v.replace(['\r', '\n'], " ") + .replace("/Users/", "/") + .replace("/home/", "/") + .chars() + .take(max) + .collect() } - -fn sanitize_log_name(value: &str) -> String { - let candidate = value +fn sanitize_log_name(v: &str) -> String { + let s = v .rsplit(['/', '\\']) .next() - .unwrap_or("application.log"); - let safe = candidate + .unwrap_or("") .chars() - .filter(|character| { - character.is_ascii_alphanumeric() || matches!(character, '.' | '-' | '_') - }) + .filter(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '-' | '_')) .take(128) .collect::(); - if safe.is_empty() || safe == "." || safe == ".." { - "application.log".to_string() + if s.is_empty() || s == "." || s == ".." { + "application.log".into() } else { - safe + s } } -fn parse_archive(bytes: &[u8]) -> Result<(Vec, Option, Vec), String> { - let mut archive = zip::ZipArchive::new(Cursor::new(bytes)).map_err(|e| e.to_string())?; - let mut events = Vec::new(); - let text = { - let mut file = archive - .by_name("events.jsonl") - .map_err(|_| "归档缺少 events.jsonl".to_string())?; - let mut text = String::new(); - std::io::Read::read_to_string(&mut file, &mut text).map_err(|e| e.to_string())?; - text - }; - for line in text.lines().filter(|line| !line.trim().is_empty()) { - events.push(serde_json::from_str(line).map_err(|e| e.to_string())?); - } - let user_description = if let Ok(mut file) = archive.by_name("user-description.txt") { - let mut text = String::new(); - std::io::Read::read_to_string(&mut file, &mut text).map_err(|e| e.to_string())?; - Some(text) - } else { - None - }; - let mut log_names = Vec::new(); - for index in 0..archive.len() { - let name = archive - .by_index(index) - .map_err(|e| e.to_string())? - .name() - .to_string(); - if let Some(name) = name.strip_prefix("system-logs/") { - log_names.push(name.to_string()); +pub fn spawn_cleanup_worker(state: AppState) { + tokio::spawn(async move { + let mut tick = interval(Duration::from_secs(24 * 60 * 60)); + loop { + tick.tick().await; + let _ = cleanup_expired(&state).await; + } + }); +} +async fn cleanup_expired(state: &AppState) -> Result<(), String> { + let cutoff = OffsetDateTime::now_utc() - time::Duration::days(30); + let (rows, _) = state + .spacetime_client() + .list_error_reports(ErrorReportListRecordInput { + status: None, + fingerprint: None, + source: None, + limit: 500, + offset: 0, + }) + .await + .map_err(|e| e.to_string())?; + for r in rows { + if let Ok(ts) = OffsetDateTime::parse(&r.created_at, &Rfc3339) { + if ts < cutoff { + if let Some(oss) = state.oss_client() { + if oss + .delete_object( + state.editor_oss_http_client(), + platform_oss::OssDeleteObjectRequest { + object_key: r.object_key.clone(), + }, + ) + .await + .is_err() + { + continue; + } + } + let _ = state + .spacetime_client() + .delete_error_report(r.batch_id) + .await; + } } } - Ok((events, user_description, log_names)) -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn archive_contains_manifest_events_description_and_logs() { - let events = vec![Event { - event_id: "event-1".to_string(), - fingerprint: "fingerprint".to_string(), - source: "react-render".to_string(), - message: "失败".to_string(), - stack: Some("Error: 失败".to_string()), - occurred_at: "2026-08-31T00:00:00Z".to_string(), - count: 1, - }]; - let archive = build_error_report_archive( - &events, - Some("用户描述"), - &[("application.log".to_string(), "safe log".to_string())], - ) - .expect("archive should build"); - let mut zip = zip::ZipArchive::new(Cursor::new(archive)).expect("zip should open"); - assert!(zip.by_name("manifest.json").is_ok()); - assert!(zip.by_name("events.jsonl").is_ok()); - assert!(zip.by_name("user-description.txt").is_ok()); - assert!(zip.by_name("system-logs/application.log").is_ok()); - } - - #[test] - fn archive_disambiguates_duplicate_log_names() { - let archive = build_error_report_archive( - &[], - None, - &[ - ("application.log".to_string(), "one".to_string()), - ("application.log".to_string(), "two".to_string()), - ], - ) - .expect("archive should build"); - let mut zip = zip::ZipArchive::new(Cursor::new(archive)).expect("zip should open"); - assert!(zip.by_name("system-logs/application.log").is_ok()); - assert!(zip.by_name("system-logs/application.log.2").is_ok()); - } - - #[test] - fn parse_archive_rejects_missing_events_file() { - let cursor = Cursor::new(Vec::new()); - let mut zip = ZipWriter::new(cursor); - zip.start_file("manifest.json", FileOptions::<()>::default()) - .expect("manifest should start"); - zip.write_all(b"{\"schemaVersion\":1}") - .expect("manifest should write"); - let archive = zip.finish().expect("archive should finish").into_inner(); - let error = parse_archive(&archive).expect_err("missing events should be rejected"); - assert!(error.contains("events.jsonl")); - } - - #[test] - fn error_report_oss_path_uses_only_uuid_filename_segment() { - assert_eq!(error_report_oss_path_segments(), ["error-reports", "v1"]); - } - - #[test] - fn sanitize_report_text_redacts_credentials_and_paths() { - assert_eq!( - sanitize_report_text("Authorization: Bearer secret"), - "[REDACTED]" - ); - assert!(sanitize_report_text("failed at /home/user/project").contains("")); - } - - #[test] - fn sanitize_report_text_with_limit_preserves_large_log_content_budget() { - let input = "x".repeat(MAX_LOG_CHARS + 100); - let output = sanitize_report_text_with_limit(&input, MAX_LOG_CHARS); - assert_eq!(output.len(), MAX_LOG_CHARS); - } - - #[test] - fn sanitize_log_name_strips_paths_and_limits_length() { - assert_eq!( - sanitize_log_name("../../application.log"), - "application.log" - ); - assert_eq!( - sanitize_log_name(r"C:\private\application.log"), - "application.log" - ); - assert_eq!(sanitize_log_name("/tmp/秘密"), "application.log"); - assert_eq!(sanitize_log_name(&"a".repeat(200)).len(), 128); - } - - #[tokio::test] - async fn store_create_sanitizes_event_metadata_and_log_names() { - let directory = - std::env::temp_dir().join(format!("agc-error-reports-test-{}", Uuid::new_v4())); - let store = ErrorReportStore::new(&directory); - let report = store - .create( - "user-1".to_string(), - CreateErrorReportBatchRequest { - schema_version: 1, - submission_id: "submission-1".to_string(), - events: vec![Event { - event_id: "event".to_string(), - fingerprint: "fingerprint".to_string(), - source: "source".to_string(), - message: "message".to_string(), - stack: None, - occurred_at: "now".to_string(), - count: 1, - }], - user_description: None, - logs: vec![ErrorReportLogInput { - name: "../../application.log".to_string(), - content: "x".repeat(3_000), - }], - }, - ) - .await - .expect("report should be stored"); - assert_eq!(report.log_count, 1); - let archive = store - .read_archive(&report.batch_id) - .await - .expect("archive should be readable"); - let mut zip = zip::ZipArchive::new(Cursor::new(archive)).expect("zip should open"); - let mut log = String::new(); - std::io::Read::read_to_string( - &mut zip.by_name("system-logs/application.log").unwrap(), - &mut log, - ) - .expect("log should be readable"); - assert_eq!(log.len(), 3_000); - let _ = fs::remove_dir_all(directory); - } - - #[tokio::test] - async fn store_scopes_submission_idempotency_to_user() { - let directory = - std::env::temp_dir().join(format!("agc-error-reports-idempotency-{}", Uuid::new_v4())); - let store = ErrorReportStore::new(&directory); - let payload = || CreateErrorReportBatchRequest { - schema_version: 1, - submission_id: "same-submission".to_string(), - events: vec![Event { - event_id: "event".to_string(), - fingerprint: "fingerprint".to_string(), - source: "source".to_string(), - message: "message".to_string(), - stack: None, - occurred_at: "now".to_string(), - count: 1, - }], - user_description: None, - logs: Vec::new(), - }; - let first = store - .create("user-1".to_string(), payload()) - .await - .expect("first report should be stored"); - let second = store - .create("user-2".to_string(), payload()) - .await - .expect("second report should be stored"); - assert_ne!(first.batch_id, second.batch_id); - let replay = store - .create("user-1".to_string(), payload()) - .await - .expect("replay should be idempotent"); - assert_eq!(replay.batch_id, first.batch_id); - let _ = fs::remove_dir_all(directory); - } - - #[tokio::test] - async fn store_list_paginates_from_index() { - let directory = - std::env::temp_dir().join(format!("agc-error-reports-list-{}", Uuid::new_v4())); - let store = ErrorReportStore::new(&directory); - for index in 0..3 { - store - .create( - "user-1".to_string(), - CreateErrorReportBatchRequest { - schema_version: 1, - submission_id: format!("submission-{index}"), - events: vec![Event { - event_id: format!("event-{index}"), - fingerprint: "fingerprint".to_string(), - source: "source".to_string(), - message: "message".to_string(), - stack: None, - occurred_at: "now".to_string(), - count: 1, - }], - user_description: None, - logs: Vec::new(), - }, - ) - .await - .expect("report should be stored"); - } - let page = store - .list(AdminErrorReportListQuery { - status: None, - fingerprint: None, - source: None, - limit: Some(2), - offset: Some(1), - }) - .await - .expect("list should succeed"); - assert_eq!(page.total, 3); - assert_eq!(page.offset, 1); - assert_eq!(page.limit, 2); - assert!(!page.has_more); - assert_eq!(page.reports.len(), 2); - assert!(store.index_path().is_file()); - let _ = fs::remove_dir_all(directory); - } - - #[tokio::test] - async fn update_missing_report_returns_not_found_error() { - let directory = - std::env::temp_dir().join(format!("agc-error-reports-update-{}", Uuid::new_v4())); - let store = ErrorReportStore::new(&directory); - let error = store - .update( - &Uuid::new_v4().to_string(), - AdminUpdateErrorReportRequest { - status: "resolved".to_string(), - note: None, - }, - ) - .await - .expect_err("missing report should fail"); - assert!(matches!(error, ErrorReportUpdateError::NotFound)); - let _ = fs::remove_dir_all(directory); - } - - #[tokio::test] - async fn get_distinguishes_missing_and_corrupt_storage() { - let directory = - std::env::temp_dir().join(format!("agc-error-reports-get-{}", Uuid::new_v4())); - let store = ErrorReportStore::new(&directory); - let missing = store - .get(&Uuid::new_v4().to_string()) - .await - .expect_err("missing report should fail"); - assert!(matches!(missing, ErrorReportGetError::NotFound(_))); - - fs::create_dir_all(&store.directory).expect("store directory should exist"); - let batch_id = Uuid::new_v4().to_string(); - fs::write(store.metadata_path(&batch_id), b"not-json").expect("metadata should write"); - let corrupt = store - .get(&batch_id) - .await - .expect_err("corrupt metadata should fail"); - assert!(matches!(corrupt, ErrorReportGetError::Internal(_))); - let _ = fs::remove_dir_all(directory); - } + Ok(()) } diff --git a/server-rs/crates/api-server/src/main.rs b/server-rs/crates/api-server/src/main.rs index 2eb880f48..fbc638bee 100644 --- a/server-rs/crates/api-server/src/main.rs +++ b/server-rs/crates/api-server/src/main.rs @@ -566,6 +566,7 @@ fn spawn_common_app_state_background_workers(state: &AppState) { fn spawn_http_app_state_background_workers(state: &AppState, process_role: ProcessRole) { spawn_common_app_state_background_workers(state); + crate::error_reports::spawn_cleanup_worker(state.clone()); if should_start_profile_recharge_expiration_listener(process_role) { spawn_profile_recharge_expiration_listener(state.clone()); spawn_profile_recharge_refund_reconciliation_worker(state.clone()); diff --git a/server-rs/crates/api-server/src/state.rs b/server-rs/crates/api-server/src/state.rs index fdc715ea0..284b437cb 100644 --- a/server-rs/crates/api-server/src/state.rs +++ b/server-rs/crates/api-server/src/state.rs @@ -44,7 +44,6 @@ use crate::editor_generation_config::{ EditorGenerationPricingConfig, EditorGenerationPricingError, EditorGenerationPricingStore, EditorGenerationPricingUnit, }; -use crate::error_reports::ErrorReportStore; use crate::tracking_outbox::TrackingOutbox; use crate::wallet_refund_outbox::{ProfileWalletRefundOutboxWorker, WalletRefundOutbox}; use crate::wechat::pay::{build_wechat_pay_config, map_wechat_pay_init_error}; @@ -297,7 +296,6 @@ pub struct AppStateInner { #[cfg(any())] puzzle_gallery_cache: PuzzleGalleryCache, tracking_outbox: Option>, - error_report_store: Arc, wallet_refund_outbox: Option>, profile_wallet_refund_outbox_worker: Arc, editor_generation_pricing_store: EditorGenerationPricingStore, @@ -591,12 +589,6 @@ impl AppState { let ai_task_service = AiTaskService::new(InMemoryAiTaskStore::default()); let spacetime_client = SpacetimeClient::new(spacetime_client_config_for_process(&config)); let tracking_outbox = TrackingOutbox::from_config(&config, spacetime_client.clone()); - let error_report_store = ErrorReportStore::new( - config - .tracking_outbox_dir - .parent() - .unwrap_or_else(|| std::path::Path::new("server-rs/.data")), - ); let wallet_refund_outbox = WalletRefundOutbox::from_config(&config, spacetime_client.clone()); let profile_wallet_refund_outbox_worker = @@ -674,7 +666,6 @@ impl AppState { #[cfg(any())] puzzle_gallery_cache: PuzzleGalleryCache::new(), tracking_outbox, - error_report_store, wallet_refund_outbox, profile_wallet_refund_outbox_worker, editor_generation_pricing_store, @@ -1312,10 +1303,6 @@ impl AppState { self.oss_client.as_ref() } - pub fn error_report_store(&self) -> &Arc { - &self.error_report_store - } - pub fn password_entry_service(&self) -> &PasswordEntryService { &self.password_entry_service } diff --git a/server-rs/crates/platform-oss/src/lib.rs b/server-rs/crates/platform-oss/src/lib.rs index 7ff931278..88401b799 100644 --- a/server-rs/crates/platform-oss/src/lib.rs +++ b/server-rs/crates/platform-oss/src/lib.rs @@ -104,6 +104,15 @@ pub struct OssSignedGetObjectUrlRequest { pub struct OssHeadObjectRequest { pub object_key: String, } +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct OssGetObjectRequest { + pub object_key: String, + pub max_bytes: usize, +} +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct OssDeleteObjectRequest { + pub object_key: String, +} #[derive(Clone, Debug, PartialEq, Eq)] pub struct OssPutObjectRequest { @@ -219,6 +228,8 @@ pub struct OssClient { pub enum OssRequestOperation { Put, Head, + Get, + Delete, } #[derive(Clone, Debug, PartialEq, Eq)] @@ -848,6 +859,70 @@ impl OssClient { result } + pub async fn get_object( + &self, + client: &reqwest::Client, + request: OssGetObjectRequest, + ) -> Result, OssError> { + let key = normalize_internal_object_key(&request.object_key)?; + let target = build_object_url(&self.config.bucket, &self.config.endpoint, &key) + .map_err(|e| request_error(OssRequestOperation::Get, &e.to_string()))?; + let response = send_signed_request( + client, + &self.config, + Method::GET, + Some(&key), + target, + OssRequestOperation::Get, + ) + .await?; + if response.status() == reqwest::StatusCode::NOT_FOUND { + return Err(OssError::ObjectNotFound(format!("OSS 对象不存在:{key}"))); + } + if !response.status().is_success() { + return Err(request_status_error( + OssRequestOperation::Get, + response.status().as_u16(), + format!("OSS GET Object 失败,状态码:{}", response.status()), + )); + } + let bytes = response + .bytes() + .await + .map_err(|e| request_error_from_reqwest(OssRequestOperation::Get, e))?; + if bytes.len() > request.max_bytes { + return Err(OssError::InvalidRequest("OSS 对象超过读取上限".to_string())); + } + Ok(bytes.to_vec()) + } + + pub async fn delete_object( + &self, + client: &reqwest::Client, + request: OssDeleteObjectRequest, + ) -> Result<(), OssError> { + let key = normalize_internal_object_key(&request.object_key)?; + let target = build_object_url(&self.config.bucket, &self.config.endpoint, &key) + .map_err(|e| request_error(OssRequestOperation::Delete, &e.to_string()))?; + let response = send_signed_request( + client, + &self.config, + Method::DELETE, + Some(&key), + target, + OssRequestOperation::Delete, + ) + .await?; + if response.status() == reqwest::StatusCode::NOT_FOUND || response.status().is_success() { + return Ok(()); + } + Err(request_status_error( + OssRequestOperation::Delete, + response.status().as_u16(), + format!("OSS DELETE Object 失败,状态码:{}", response.status()), + )) + } + // AI 生成资源默认由服务端上传 OSS,Web 端只拿签名读地址,不直接持有写权限。 pub async fn put_object( &self, @@ -1693,6 +1768,18 @@ fn normalize_editor_agent_messages_object_key(raw: &str) -> Result Result { + let normalized = raw.trim().trim_start_matches('/').trim().to_string(); + validate_object_key_segments(&normalized)?; + if normalized.starts_with("agc/error-reports/v1/") || normalized.starts_with("editor-agent/") { + Ok(normalized) + } else { + Err(OssError::InvalidRequest( + "objectKey 不属于内部对象前缀".to_string(), + )) + } +} + fn validate_object_key_segments(normalized: &str) -> Result<(), OssError> { let segments = normalized.split('/').collect::>(); if segments.len() < 2 {