完善错误报告存储与管理员错误契约

使用本地元数据索引支持幂等查找和分页读取

拒绝缺失 events.jsonl 的损坏归档并区分管理员更新错误状态

将 OSS 归档 key 固定为仅含报告 UUID 的路径并跳过 ready 重放上传
This commit is contained in:
2026-09-01 16:28:52 +08:00
parent d9a7260996
commit 4d061b0975
2 changed files with 264 additions and 84 deletions
+259 -84
View File
@@ -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::<StoredErrorReport>(&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<Vec<AdminErrorReportEntry>, String> {
) -> Result<AdminErrorReportListResponse, String> {
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::<StoredErrorReport>(&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::<Vec<_>>();
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<ErrorReportDetail, ErrorReportGetError> {
@@ -656,28 +649,39 @@ impl ErrorReportStore {
&self,
batch_id: &str,
payload: AdminUpdateErrorReportRequest,
) -> Result<StoredErrorReport, String> {
validate_batch_id(batch_id)?;
) -> Result<StoredErrorReport, ErrorReportUpdateError> {
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<Vec<StoredErrorReport>, String> {
let path = self.index_path();
match fs::read(&path) {
Ok(bytes) => {
let reports = serde_json::from_slice::<Vec<StoredErrorReport>>(&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::<StoredErrorReport>(&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<Event>, Option<String>, Vec<String>), 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 =
@@ -46,6 +46,7 @@ pub struct AdminErrorReportListQuery {
pub fingerprint: Option<String>,
pub source: Option<String>,
pub limit: Option<u32>,
pub offset: Option<u32>,
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
@@ -73,6 +74,10 @@ pub struct AdminErrorReportEntry {
#[serde(rename_all = "camelCase")]
pub struct AdminErrorReportListResponse {
pub reports: Vec<AdminErrorReportEntry>,
pub total: u64,
pub offset: u32,
pub limit: u32,
pub has_more: bool,
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]