Files
Genarrative/server-rs/crates/api-server/src/error_reports.rs
T
k88936 422c8931d6
Project CI / Repository checks (push) Successful in 2m45s
Project CI / Frontend tests (push) Successful in 3m22s
Project CI / Backend tests (push) Successful in 5m53s
Project CI / Native shell tests (push) Successful in 16m45s
Feat/AGC错误报告 (#240)
实现:
在rust内存里维护错误事件队列, webview调用tauri command传入, 后台任务agent工具等直接插入
把rust , webview console的日志统一写到AppData文件夹下的(滚动保存的)日志文件.
rust对传入的错误进行筛选,脱敏, 防抖,后通知前端提醒用户.
用户提醒是一个不阻塞的小UI, 展开后可以选择错误上报, 可以附加文字描述
上传时附带最近日志, 错误堆栈等信息

元数据存在数据库, 考虑到字符串信息很难查询, 所以在api server打包成zip存在OSS.
管理页面新增错误报告的查看页面

![shotmd-1788328504.jpg](/attachments/641b43c8-deed-47d8-8c78-36adb9a2548f)
![shotmd-1788328495.jpg](/attachments/33c6ada3-dd4d-44b6-9859-f18f944eb653)
![shotmd-1788328222.jpg](/attachments/59267df9-f624-4e47-bb9b-b69d77f595c4)
![shotmd-1788328230.jpg](/attachments/63ea6a86-9a51-412b-a98c-dce454bcc968)

---------

Co-authored-by: 段舒康 <kdletters@qq.com>
Reviewed-on: http://192.168.35.82/git/GenarrativeAI/Genarrative/pulls/240
Co-authored-by: 王德宇 <kvtodev@outlook.com>
Co-committed-by: 王德宇 <kvtodev@outlook.com>
2026-09-03 10:02:05 +08:00

788 lines
26 KiB
Rust

use crate::{
admin::{AuthenticatedAdmin, require_admin_auth},
api_response::json_success_body,
auth::{AuthenticatedAccessToken, require_bearer_auth},
http_error::AppError,
platform_errors,
request_context::RequestContext,
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},
sync::LazyLock,
};
use time::OffsetDateTime;
use tokio::time::{Duration, interval};
use uuid::Uuid;
use zip::{ZipWriter, write::FileOptions};
const MAX_EVENTS: usize = 100;
const MAX_LOGS: usize = 5;
const MAX_DESCRIPTION_CHARS: usize = 2_000;
const MAX_LOG_CHARS: usize = 2_000_000;
const MAX_BATCH_BYTES: usize = 20 * 1024 * 1024;
const MAX_REQUEST_BODY_BYTES: usize = 24 * 1024 * 1024;
const MAX_EVENT_FIELD_CHARS: usize = 512;
const MAX_EVENTS_JSONL_BYTES: usize = MAX_REQUEST_BODY_BYTES;
const MAX_NOTE_CHARS: usize = 2_000;
static LOCAL_PATH_PATTERN: LazyLock<regex::Regex> = LazyLock::new(|| {
regex::Regex::new(r"(?i)(?:[a-z]:/|/)(?:users|home|private|tmp)/[^\s]+")
.expect("valid report path sanitization pattern")
});
#[derive(Clone, Debug, serde::Serialize)]
#[serde(rename_all = "camelCase")]
struct ErrorReportDetail {
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<String>,
first_source: Option<String>,
review_status: String,
admin_note: Option<String>,
created_at: String,
updated_at: String,
events: Vec<Event>,
log_names: Vec<String>,
user_description: Option<String>,
status: String,
note: Option<String>,
attachment_size_bytes: u64,
}
pub fn router(state: AppState) -> Router<AppState> {
Router::new()
.route(
"/api/error-reports",
post(create_error_report)
.route_layer(middleware::from_fn_with_state(
state.clone(),
require_bearer_auth,
))
.layer(DefaultBodyLimit::max(MAX_REQUEST_BODY_BYTES)),
)
.route(
"/admin/api/error-reports",
get(admin_list_error_reports).route_layer(middleware::from_fn_with_state(
state.clone(),
require_admin_auth,
)),
)
.route(
"/admin/api/error-reports/{batch_id}",
get(admin_get_error_report)
.patch(admin_update_error_report)
.route_layer(middleware::from_fn_with_state(
state.clone(),
require_admin_auth,
)),
)
.route(
"/admin/api/error-reports/{batch_id}/download",
get(admin_download_error_report)
.route_layer(middleware::from_fn_with_state(state, require_admin_auth)),
)
}
pub async fn create_error_report(
State(state): State<AppState>,
Extension(ctx): Extension<RequestContext>,
Extension(auth): Extension<AuthenticatedAccessToken>,
Json(mut payload): Json<CreateErrorReportBatchRequest>,
) -> Result<Json<Value>, AppError> {
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::<Vec<_>>(),
)
.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 digest = format!("{:x}", Sha256::digest(&archive));
let oss = state
.oss_client()
.ok_or_else(|| AppError::from_status(StatusCode::BAD_GATEWAY).with_message("OSS 未配置"))?;
let put_response = 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 object_key = put_response.object_key;
let record = match state
.spacetime_client()
.create_error_report(ErrorReportCreateRecordInput {
batch_id: batch_id.clone(),
submission_id: payload.submission_id.clone(),
idempotency_key: format!("{}:{}", auth.claims().user_id(), payload.submission_id),
user_id: auth.claims().user_id().to_string(),
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()),
})
.await
{
Ok(record) => record,
Err(_) => {
let _ = oss
.delete_object(
state.editor_oss_http_client(),
OssDeleteObjectRequest {
object_key: object_key.clone(),
},
)
.await;
return Err(AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR)
.with_message("报告元数据写入失败"));
}
};
let accepted_event_ids = if record.batch_id != batch_id {
let existing_archive = oss
.get_object(
state.editor_oss_http_client(),
OssGetObjectRequest {
object_key: record.object_key.clone(),
max_bytes: MAX_BATCH_BYTES,
},
)
.await
.map_err(map_error_report_oss_error)?;
let (existing_events, _, _) =
parse_archive(&existing_archive).map_err(|e| internal(format!("报告归档损坏:{e}")))?;
let mut requested_event_ids = payload
.events
.iter()
.map(|event| event.event_id.clone())
.collect::<Vec<_>>();
let existing_event_ids = existing_events
.iter()
.map(|event| event.event_id.clone())
.collect::<Vec<_>>();
let mut existing_event_ids_sorted = existing_event_ids.clone();
requested_event_ids.sort_unstable();
existing_event_ids_sorted.sort_unstable();
if requested_event_ids != existing_event_ids_sorted {
let _ = oss
.delete_object(
state.editor_oss_http_client(),
OssDeleteObjectRequest { object_key },
)
.await;
return Err(AppError::from_status(StatusCode::CONFLICT)
.with_message("相同 submissionId 的错误报告事件集合不一致"));
}
let _ = oss
.delete_object(
state.editor_oss_http_client(),
OssDeleteObjectRequest { object_key },
)
.await;
existing_event_ids
} else {
payload
.events
.into_iter()
.map(|event| event.event_id)
.collect()
};
Ok(json_success_body(
Some(&ctx),
CreateErrorReportBatchResponse {
batch_id: record.batch_id,
submission_id: record.submission_id,
status: "ready".into(),
accepted_event_ids,
created_at: record.created_at,
},
))
}
pub async fn admin_list_error_reports(
State(state): State<AppState>,
Extension(ctx): Extension<RequestContext>,
Extension(admin): Extension<AuthenticatedAdmin>,
Query(q): Query<AdminErrorReportListQuery>,
) -> Result<Json<Value>, AppError> {
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)?;
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 as u64).saturating_add(limit as u64) < total,
},
))
}
pub async fn admin_get_error_report(
State(state): State<AppState>,
Extension(ctx): Extension<RequestContext>,
Extension(admin): Extension<AuthenticatedAdmin>,
AxumPath(batch_id): AxumPath<String>,
) -> Result<Json<Value>, AppError> {
let r = state
.spacetime_client()
.get_error_report(batch_id.clone())
.await
.map_err(map_error_report_spacetime_error)?;
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(map_error_report_oss_error)?;
let (events, desc, logs) =
parse_archive(&bytes).map_err(|e| internal(format!("报告归档损坏:{e}")))?;
record_admin_report_audit(
&state,
&ctx,
admin.session().subject.as_str(),
"view",
Some(&batch_id),
)
.await;
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<AppState>,
Extension(ctx): Extension<RequestContext>,
Extension(admin): Extension<AuthenticatedAdmin>,
AxumPath(batch_id): AxumPath<String>,
Json(p): Json<AdminUpdateErrorReportRequest>,
) -> Result<Json<Value>, AppError> {
if !matches!(p.status.as_str(), "new" | "in-progress" | "resolved") {
return Err(bad_request("状态必须是 new、in-progress 或 resolved"));
}
if p.note
.as_deref()
.is_some_and(|note| note.chars().count() > MAX_NOTE_CHARS)
{
return Err(bad_request("处理备注不能超过 2000 个字符"));
}
let updated = state
.spacetime_client()
.update_error_report(ErrorReportUpdateRecordInput {
batch_id: batch_id.clone(),
status: p.status,
note: p.note,
})
.await
.map_err(internal)?;
record_admin_report_audit(
&state,
&ctx,
admin.session().subject.as_str(),
"update",
Some(&batch_id),
)
.await;
Ok(json_success_body(
Some(&ctx),
AdminErrorReportEntry {
batch_id: updated.batch_id,
event_count: updated.event_count,
log_count: updated.log_count,
fingerprint: updated.first_fingerprint,
source: updated.first_source,
status: updated.review_status.clone(),
user_id: updated.user_id,
created_at: updated.created_at,
updated_at: updated.updated_at,
user_description: None,
attachment_size_bytes: updated.archive_size_bytes,
submission_id: Some(updated.submission_id),
oss_object_key: Some(updated.object_key),
archive_sha256: Some(updated.archive_sha256),
upload_status: Some("ready".into()),
review_status: Some(updated.review_status),
},
))
}
pub async fn admin_download_error_report(
State(state): State<AppState>,
Extension(ctx): Extension<RequestContext>,
Extension(admin): Extension<AuthenticatedAdmin>,
AxumPath(batch_id): AxumPath<String>,
) -> Result<Response, AppError> {
let r = state
.spacetime_client()
.get_error_report(batch_id.clone())
.await
.map_err(map_error_report_spacetime_error)?;
let bytes = state
.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(map_error_report_oss_error)?;
record_admin_report_audit(
&state,
&ctx,
admin.session().subject.as_str(),
"download",
Some(&batch_id),
)
.await;
let mut resp = Response::new(bytes.into());
resp.headers_mut().insert(
header::CONTENT_TYPE,
HeaderValue::from_static("application/zip"),
);
resp.headers_mut().insert(
header::CONTENT_DISPOSITION,
HeaderValue::from_str(&format!("attachment; filename=\"{batch_id}.zip\""))
.map_err(internal)?,
);
Ok(resp)
}
fn bad_request(m: impl Into<String>) -> AppError {
AppError::from_status(StatusCode::BAD_REQUEST).with_message(m)
}
fn internal<E: ToString>(e: E) -> AppError {
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR).with_message(e.to_string())
}
fn map_error_report_oss_error(error: platform_oss::OssError) -> AppError {
match error.kind() {
platform_oss::OssErrorKind::ObjectNotFound => {
AppError::from_status(StatusCode::NOT_FOUND).with_message("报告附件不存在")
}
platform_oss::OssErrorKind::InvalidRequest => {
AppError::from_status(StatusCode::BAD_REQUEST).with_message("报告附件读取请求无效")
}
_ => platform_errors::map_oss_error(error, "error-report"),
}
}
fn map_error_report_spacetime_error(error: spacetime_client::SpacetimeClientError) -> AppError {
if error.to_string() == "错误报告不存在" {
AppError::from_status(StatusCode::NOT_FOUND).with_message("报告不存在")
} else {
internal(error)
}
}
async fn record_admin_report_audit(
state: &AppState,
ctx: &RequestContext,
subject: &str,
action: &'static str,
batch: Option<&str>,
) {
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 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::<HashSet<_>>();
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],
desc: Option<&str>,
logs: &[(String, String)],
) -> Result<Vec<u8>, 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(|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(|e| e.to_string())?;
z.write_all(b"\n").map_err(|e| e.to_string())?;
}
if let Some(d) = desc {
z.start_file("user-description.txt", o)
.map_err(|e| e.to_string())?;
z.write_all(d.as_bytes()).map_err(|e| e.to_string())?;
}
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())?;
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 parse_archive(bytes: &[u8]) -> Result<(Vec<Event>, Option<String>, Vec<String>), String> {
let mut a = zip::ZipArchive::new(Cursor::new(bytes)).map_err(|e| e.to_string())?;
let mut t = String::new();
{
let mut events_file = a
.by_name("events.jsonl")
.map_err(|_| "归档缺少 events.jsonl".to_string())?
.take((MAX_EVENTS_JSONL_BYTES + 1) as u64);
events_file
.read_to_string(&mut t)
.map_err(|e| e.to_string())?;
}
if t.len() > MAX_EVENTS_JSONL_BYTES {
return Err("events.jsonl 解压后超过大小上限".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::<Result<Vec<Event>, _>>()?;
if events.len() > MAX_EVENTS {
return Err("归档事件数量超过上限".to_string());
}
let desc = if let Ok(f) = a.by_name("user-description.txt") {
let mut s = String::new();
f.take((MAX_DESCRIPTION_CHARS + 1) as u64)
.read_to_string(&mut s)
.map_err(|e| e.to_string())?;
if s.chars().count() > MAX_DESCRIPTION_CHARS {
return Err("user-description.txt 解压后超过大小上限".to_string());
}
Some(s)
} else {
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(v: &str) -> String {
sanitize_report_text_with_limit(v, MAX_EVENT_FIELD_CHARS)
}
fn sanitize_report_text_with_limit(v: &str, max: usize) -> String {
let normalized = v
.to_lowercase()
.chars()
.filter(|c| !c.is_whitespace())
.collect::<String>();
if [
"authorization:",
"authorization=",
"bearer",
"token=",
"token:",
"x-api-key:",
"x-api-key=",
"xapikey:",
"xapikey=",
"api-key:",
"api-key=",
"api_key:",
"api_key=",
"apikey:",
"apikey=",
]
.iter()
.any(|m| normalized.contains(m))
{
return "[REDACTED]".into();
}
let normalized = v.replace(['\r', '\n'], " ").replace('\\', "/");
LOCAL_PATH_PATTERN
.replace_all(&normalized, "<absolute-path>/")
.chars()
.take(max)
.collect()
}
#[cfg(test)]
mod tests {
use super::sanitize_report_text_with_limit;
#[test]
fn redacts_credentials_separated_by_unicode_whitespace() {
for value in [
"token\u{00a0}=\u{2003}secret",
"Authorization:\u{2003}Bearer\u{00a0}secret",
"api-key\u{00a0}:\u{2003}secret",
] {
assert_eq!(sanitize_report_text_with_limit(value, 512), "[REDACTED]");
}
}
}
fn sanitize_log_name(v: &str) -> String {
let s = v
.rsplit(['/', '\\'])
.next()
.unwrap_or("")
.chars()
.filter(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '-' | '_'))
.take(128)
.collect::<String>();
if s.is_empty() || s == "." || s == ".." {
"application.log".into()
} else {
s
}
}
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 Some(oss) = state.oss_client() else {
return Ok(());
};
let mut offset = 0_u32;
let mut first_error = None;
loop {
let (rows, _) = state
.spacetime_client()
.list_error_reports(ErrorReportListRecordInput {
status: None,
fingerprint: None,
source: None,
limit: 500,
offset,
})
.await
.map_err(|e| e.to_string())?;
if rows.is_empty() {
break;
}
let page_len = rows.len() as u32;
let mut deleted = 0_u32;
for r in rows {
if r.created_at
.split('.')
.next()
.and_then(|s| s.parse::<i64>().ok())
.and_then(|seconds| OffsetDateTime::from_unix_timestamp(seconds).ok())
.is_some_and(|ts| ts < cutoff)
{
if oss
.delete_object(
state.editor_oss_http_client(),
OssDeleteObjectRequest {
object_key: r.object_key.clone(),
},
)
.await
.is_err()
{
continue;
}
if let Err(error) = state
.spacetime_client()
.delete_error_report(r.batch_id)
.await
{
first_error.get_or_insert_with(|| error.to_string());
} else {
deleted += 1;
}
}
}
if page_len < 500 {
break;
}
if deleted == 0 {
offset = offset.saturating_add(page_len);
}
}
first_error.map_or(Ok(()), Err)
}