diff --git a/server-rs/crates/api-server/src/error_reports.rs b/server-rs/crates/api-server/src/error_reports.rs index 0419e7df1..dc6cb6805 100644 --- a/server-rs/crates/api-server/src/error_reports.rs +++ b/server-rs/crates/api-server/src/error_reports.rs @@ -51,8 +51,6 @@ const MAX_EVENT_FIELD_CHARS: usize = 512; pub struct ErrorReportStore { directory: PathBuf, - // 当前 api-server 按单实例部署;进程内锁足以串行化本地元数据/归档操作, - // 不引入跨进程锁或内存索引,避免把单实例低频诊断路径复杂化。 lock: Mutex<()>, } @@ -63,6 +61,25 @@ enum ErrorReportGetError { 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 { @@ -157,6 +174,7 @@ pub async fn create_error_report( .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() @@ -164,15 +182,7 @@ pub async fn create_error_report( .await .map_err(internal_store_error)?; let digest = format!("{:x}", Sha256::digest(&archive)); - let now = OffsetDateTime::now_utc(); - let date = now.date(); - let key_segments = vec![ - "error-reports".to_string(), - "v1".to_string(), - format!("{:04}", date.year()), - format!("{:02}", u8::from(date.month())), - format!("{:02}", date.day()), - ]; + let key_segments = vec!["error-reports".to_string(), "v1".to_string()]; let response = oss .put_object( state.editor_oss_http_client(), @@ -205,6 +215,7 @@ pub async fn create_error_report( .await .map_err(internal_store_error)? } + None if stored.upload_status == "ready" => stored, None => state .error_report_store() .mark_upload_failed(&stored.batch_id) @@ -242,10 +253,7 @@ pub async fn admin_list_error_reports( None, ) .await; - Ok(json_success_body( - Some(&request_context), - AdminErrorReportListResponse { reports }, - )) + Ok(json_success_body(Some(&request_context), reports)) } pub async fn admin_get_error_report( @@ -288,7 +296,17 @@ pub async fn admin_update_error_report( .error_report_store() .update(&batch_id, payload) .await - .map_err(|message| AppError::from_status(StatusCode::BAD_REQUEST).with_message(message))?; + .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) @@ -525,21 +543,11 @@ impl ErrorReportStore { fs::create_dir_all(&self.directory).map_err(|error| error.to_string())?; set_private_directory_permissions(&self.directory)?; self.cleanup_expired(); - if let Ok(directory) = fs::read_dir(&self.directory) { - for item in directory.flatten() { - if item.path().extension().and_then(|value| value.to_str()) != Some("json") { - continue; - } - if let Ok(bytes) = fs::read(item.path()) { - if let Ok(existing) = serde_json::from_slice::(&bytes) { - if existing.user_id == stored.user_id - && existing.submission_id == stored.submission_id - { - return Ok(existing); - } - } - } - } + 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); @@ -548,32 +556,21 @@ impl ErrorReportStore { &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, String> { + ) -> Result { let _guard = self.lock.lock().await; self.cleanup_expired(); - let mut entries = Vec::new(); + let mut entries = self.load_index()?; let limit = query.limit.unwrap_or(100).clamp(1, 500) as usize; - let directory = match fs::read_dir(&self.directory) { - Ok(directory) => directory, - Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(entries), - Err(error) => return Err(error.to_string()), - }; - for item in directory.flatten() { - if item.path().extension().and_then(|value| value.to_str()) != Some("json") { - continue; - } - let Ok(bytes) = fs::read(item.path()) else { - continue; - }; - let Ok(report) = serde_json::from_slice::(&bytes) else { - continue; - }; + let offset = query.offset.unwrap_or(0).min(u32::MAX) as usize; + entries.retain(|report| { if query .status .as_deref() @@ -587,34 +584,30 @@ impl ErrorReportStore { .as_deref() .is_some_and(|value| report.first_fingerprint.as_deref() != Some(value)) { - continue; + false + } else { + true } - entries.push(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), - }); - } + }); entries.sort_by(|left, right| { OffsetDateTime::parse(&right.created_at, &Rfc3339) .ok() .cmp(&OffsetDateTime::parse(&left.created_at, &Rfc3339).ok()) }); - entries.truncate(limit); - Ok(entries) + 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 { @@ -656,28 +649,39 @@ impl ErrorReportStore { &self, batch_id: &str, payload: AdminUpdateErrorReportRequest, - ) -> Result { - validate_batch_id(batch_id)?; + ) -> Result { + validate_batch_id(batch_id).map_err(|_| ErrorReportUpdateError::InvalidId)?; if !matches!(payload.status.as_str(), "new" | "in-progress" | "resolved") { - return Err("状态必须是 new、in-progress 或 resolved".to_string()); + 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| error.to_string())?; - let mut report: StoredErrorReport = - serde_json::from_slice(&bytes).map_err(|error| error.to_string())?; + 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| error.to_string())?; + .map_err(|error| ErrorReportUpdateError::Internal(error.to_string()))?; write_atomic( &path, - &serde_json::to_vec_pretty(&report).map_err(|error| error.to_string())?, - )?; + &serde_json::to_vec_pretty(&report) + .map_err(|error| ErrorReportUpdateError::Internal(error.to_string()))?, + ) + .map_err(ErrorReportUpdateError::Internal)?; refresh_archive_mtime(&self.archive_path(batch_id)); + self.replace_index_entry(report.clone()) + .map_err(ErrorReportUpdateError::Internal)?; Ok(report) } @@ -709,6 +713,7 @@ impl ErrorReportStore { &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) } @@ -728,6 +733,7 @@ impl ErrorReportStore { &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) } @@ -739,6 +745,66 @@ impl ErrorReportStore { 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; @@ -747,7 +813,9 @@ impl ErrorReportStore { .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") { + if path.extension().and_then(|value| value.to_str()) != Some("json") + || path == self.index_path() + { continue; } let Ok(metadata) = fs::symlink_metadata(&path) else { @@ -766,6 +834,30 @@ impl ErrorReportStore { let _ = fs::remove_file(path.with_extension("zip")); } } + if let Ok(reports) = self.load_index() { + let _ = self.write_index(&reports); + } + } +} + +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), } } @@ -859,12 +951,16 @@ fn sanitize_log_name(value: &str) -> String { 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(); - if let Ok(mut file) = archive.by_name("events.jsonl") { + 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())?; - for line in text.lines().filter(|line| !line.trim().is_empty()) { - events.push(serde_json::from_str(line).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(); @@ -931,6 +1027,19 @@ mod tests { 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 sanitize_report_text_redacts_credentials_and_paths() { assert_eq!( @@ -1043,6 +1152,72 @@ mod tests { 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 = diff --git a/server-rs/crates/shared-contracts/src/error_reports.rs b/server-rs/crates/shared-contracts/src/error_reports.rs index 033b4f993..e68cdcb88 100644 --- a/server-rs/crates/shared-contracts/src/error_reports.rs +++ b/server-rs/crates/shared-contracts/src/error_reports.rs @@ -46,6 +46,7 @@ pub struct AdminErrorReportListQuery { pub fingerprint: Option, pub source: Option, pub limit: Option, + pub offset: Option, } #[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] @@ -73,6 +74,10 @@ pub struct AdminErrorReportEntry { #[serde(rename_all = "camelCase")] pub struct AdminErrorReportListResponse { pub reports: Vec, + pub total: u64, + pub offset: u32, + pub limit: u32, + pub has_more: bool, } #[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]