Fix Codex image progress heartbeat merge regressions

This commit is contained in:
fawney19
2026-05-10 02:10:23 +08:00
156 changed files with 11734 additions and 936 deletions
@@ -0,0 +1,15 @@
use crate::handlers::admin::request::{AdminAppState, AdminRequestContext};
use crate::GatewayError;
use axum::body::{Body, Bytes};
use axum::response::Response;
mod routes;
pub(crate) async fn maybe_build_local_admin_background_tasks_response(
state: &AdminAppState<'_>,
request_context: &AdminRequestContext<'_>,
request_body: Option<&Bytes>,
) -> Result<Option<Response<Body>>, GatewayError> {
routes::maybe_build_local_admin_background_tasks_response(state, request_context, request_body)
.await
}
@@ -0,0 +1,373 @@
use crate::handlers::admin::request::{AdminAppState, AdminRequestContext};
use crate::handlers::admin::shared::{
attach_admin_audit_response, query_param_value, unix_secs_to_rfc3339,
};
use crate::task_runtime::{
self, set_cancel_signal, TASK_KEY_PROVIDER_DELETE, TASK_KEY_PROVIDER_OAUTH_BATCH_IMPORT,
};
use crate::GatewayError;
use aether_data_contracts::repository::background_tasks::{
BackgroundTaskKind, BackgroundTaskListQuery, BackgroundTaskStatus,
};
use axum::{
body::{Body, Bytes},
http,
response::{IntoResponse, Response},
Json,
};
use serde_json::json;
const DEFAULT_PAGE_SIZE: usize = 20;
const MAX_PAGE_SIZE: usize = 100;
const DEFAULT_EVENTS_PAGE_SIZE: usize = 50;
pub(super) async fn maybe_build_local_admin_background_tasks_response(
state: &AdminAppState<'_>,
request_context: &AdminRequestContext<'_>,
request_body: Option<&Bytes>,
) -> Result<Option<Response<Body>>, GatewayError> {
if request_context.route_family() != Some("tasks_manage") {
return Ok(None);
}
match request_context.route_kind() {
Some("list_tasks") if request_context.method() == http::Method::GET => {
let query = request_context.query_string();
let page = query_param_value(query, "page")
.and_then(|value| value.parse::<usize>().ok())
.unwrap_or(1)
.max(1);
let page_size = query_param_value(query, "page_size")
.and_then(|value| value.parse::<usize>().ok())
.unwrap_or(DEFAULT_PAGE_SIZE)
.clamp(1, MAX_PAGE_SIZE);
let kind = query_param_value(query, "kind")
.map(|value| BackgroundTaskKind::from_database(value.as_str()))
.transpose()
.map_err(|err| GatewayError::Internal(err.to_string()))?;
let status = query_param_value(query, "status")
.map(|value| BackgroundTaskStatus::from_database(value.as_str()))
.transpose()
.map_err(|err| GatewayError::Internal(err.to_string()))?;
let trigger = query_param_value(query, "trigger");
let task_key_substring = query_param_value(query, "task_key");
let offset = (page - 1).saturating_mul(page_size);
let response = state
.list_background_task_runs(&BackgroundTaskListQuery {
task_key_substring,
kind,
status,
trigger,
offset,
limit: page_size,
})
.await?;
let pages = if response.total == 0 {
0
} else {
(response.total + page_size - 1) / page_size
};
let items = response
.items
.iter()
.map(|run| {
json!({
"id": run.id,
"task_key": run.task_key,
"kind": run.kind.as_database(),
"trigger": run.trigger,
"status": run.status.as_database(),
"attempt": run.attempt,
"max_attempts": run.max_attempts,
"owner_instance": run.owner_instance,
"progress_percent": run.progress_percent,
"progress_message": run.progress_message,
"payload": run.payload_json,
"result": run.result_json,
"error_message": run.error_message,
"cancel_requested": run.cancel_requested,
"created_by": run.created_by,
"created_at": unix_secs_to_rfc3339(run.created_at_unix_secs),
"started_at": run.started_at_unix_secs.and_then(unix_secs_to_rfc3339),
"finished_at": run.finished_at_unix_secs.and_then(unix_secs_to_rfc3339),
"updated_at": unix_secs_to_rfc3339(run.updated_at_unix_secs),
})
})
.collect::<Vec<_>>();
let definitions = task_runtime::task_definitions()
.iter()
.map(|definition| {
json!({
"task_key": definition.key,
"kind": definition.kind.as_str(),
"trigger": definition.trigger,
"max_attempts": definition.retry_policy.max_attempts,
"singleton": definition.singleton,
"persist_history": definition.persist_history,
})
})
.collect::<Vec<_>>();
return Ok(Some(
Json(json!({
"items": items,
"total": response.total,
"page": page,
"page_size": page_size,
"pages": pages,
"definitions": definitions,
}))
.into_response(),
));
}
Some("stats") if request_context.method() == http::Method::GET => {
let stats = state.summarize_background_task_runs().await?;
return Ok(Some(
Json(json!({
"total": stats.total,
"running_count": stats.running_count,
"by_status": stats.by_status,
"by_kind": stats.by_kind,
"registered_tasks": task_runtime::task_definitions().len(),
}))
.into_response(),
));
}
Some("detail") if request_context.method() == http::Method::GET => {
let Some(run_id) = task_id_from_path(request_context.path()) else {
return Ok(Some(
(
http::StatusCode::NOT_FOUND,
Json(json!({"detail":"Task not found"})),
)
.into_response(),
));
};
let Some(run) = state.find_background_task_run(run_id).await? else {
return Ok(Some(
(
http::StatusCode::NOT_FOUND,
Json(json!({"detail":"Task not found"})),
)
.into_response(),
));
};
return Ok(Some(attach_admin_audit_response(
Json(json!({
"id": run.id,
"task_key": run.task_key,
"kind": run.kind.as_database(),
"trigger": run.trigger,
"status": run.status.as_database(),
"attempt": run.attempt,
"max_attempts": run.max_attempts,
"owner_instance": run.owner_instance,
"progress_percent": run.progress_percent,
"progress_message": run.progress_message,
"payload": run.payload_json,
"result": run.result_json,
"error_message": run.error_message,
"cancel_requested": run.cancel_requested,
"created_by": run.created_by,
"created_at": unix_secs_to_rfc3339(run.created_at_unix_secs),
"started_at": run.started_at_unix_secs.and_then(unix_secs_to_rfc3339),
"finished_at": run.finished_at_unix_secs.and_then(unix_secs_to_rfc3339),
"updated_at": unix_secs_to_rfc3339(run.updated_at_unix_secs),
}))
.into_response(),
"admin_task_detail_viewed",
"view_task_detail",
"background_task",
run_id,
)));
}
Some("events") if request_context.method() == http::Method::GET => {
let Some(run_id) = nested_task_id_from_path(request_context.path(), "/events") else {
return Ok(Some(
(
http::StatusCode::NOT_FOUND,
Json(json!({"detail":"Task not found"})),
)
.into_response(),
));
};
let query = request_context.query_string();
let page = query_param_value(query, "page")
.and_then(|value| value.parse::<usize>().ok())
.unwrap_or(1)
.max(1);
let page_size = query_param_value(query, "page_size")
.and_then(|value| value.parse::<usize>().ok())
.unwrap_or(DEFAULT_EVENTS_PAGE_SIZE)
.clamp(1, MAX_PAGE_SIZE);
let offset = (page - 1).saturating_mul(page_size);
let events = state
.list_background_task_events(run_id, offset, page_size)
.await?;
return Ok(Some(
Json(json!({
"items": events.into_iter().map(|event| {
json!({
"id": event.id,
"run_id": event.run_id,
"event_type": event.event_type,
"message": event.message,
"payload": event.payload_json,
"created_at": unix_secs_to_rfc3339(event.created_at_unix_secs),
})
}).collect::<Vec<_>>(),
"page": page,
"page_size": page_size,
}))
.into_response(),
));
}
Some("cancel") if request_context.method() == http::Method::POST => {
let Some(run_id) = nested_task_id_from_path(request_context.path(), "/cancel") else {
return Ok(Some(
(
http::StatusCode::NOT_FOUND,
Json(json!({"detail":"Task not found"})),
)
.into_response(),
));
};
let now = task_runtime::now_unix_secs();
let cancelled = state
.request_cancel_background_task_run(run_id, now)
.await?;
if !cancelled {
return Ok(Some(
(
http::StatusCode::NOT_FOUND,
Json(json!({ "detail": "Task not found" })),
)
.into_response(),
));
}
let _ = set_cancel_signal(state.app(), run_id).await;
task_runtime::append_event_with_logging(
state.app(),
run_id,
"cancel_requested",
"cancel requested by admin",
None,
)
.await;
return Ok(Some(attach_admin_audit_response(
Json(json!({
"id": run_id,
"status": "cancel_requested",
"message": "Task cancellation requested",
}))
.into_response(),
"admin_task_cancel_requested",
"cancel_task",
"background_task",
run_id,
)));
}
Some("trigger") if request_context.method() == http::Method::POST => {
let Some(task_key) = nested_task_id_from_path(request_context.path(), "/trigger")
else {
return Ok(Some(
(
http::StatusCode::NOT_FOUND,
Json(json!({"detail":"Task not found"})),
)
.into_response(),
));
};
let payload = parse_json_payload(request_body)?;
if task_key == TASK_KEY_PROVIDER_DELETE {
let provider_id = payload
.get("provider_id")
.and_then(serde_json::Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| {
GatewayError::Internal(
"admin task trigger provider delete requires provider_id".to_string(),
)
})?;
let Some(run_id) =
task_runtime::submit_provider_delete_task(state, provider_id, Some("admin"))
.await?
else {
return Ok(Some(
(
http::StatusCode::NOT_FOUND,
Json(json!({"detail":"Provider 不存在"})),
)
.into_response(),
));
};
return Ok(Some(attach_admin_audit_response(
Json(json!({
"task_key": task_key,
"run_id": run_id,
"status": "queued",
}))
.into_response(),
"admin_task_triggered",
"trigger_task",
"background_task",
task_key,
)));
}
if task_key == TASK_KEY_PROVIDER_OAUTH_BATCH_IMPORT {
return Ok(Some(
(
http::StatusCode::BAD_REQUEST,
Json(json!({
"detail": "请使用 provider oauth batch import 专用接口触发该任务",
})),
)
.into_response(),
));
}
return Ok(Some(
(
http::StatusCode::BAD_REQUEST,
Json(json!({
"detail": format!("Unsupported task_key: {task_key}"),
})),
)
.into_response(),
));
}
_ => {}
}
Ok(None)
}
fn task_id_from_path(request_path: &str) -> Option<&str> {
let task_id = request_path.strip_prefix("/api/admin/tasks/")?;
if task_id.is_empty() || task_id.contains('/') || task_id == "stats" {
return None;
}
Some(task_id)
}
fn nested_task_id_from_path<'a>(request_path: &'a str, suffix: &str) -> Option<&'a str> {
let task_id = request_path
.strip_prefix("/api/admin/tasks/")?
.strip_suffix(suffix)?;
if task_id.is_empty() || task_id.contains('/') {
return None;
}
Some(task_id)
}
fn parse_json_payload(request_body: Option<&Bytes>) -> Result<serde_json::Value, GatewayError> {
let Some(body) = request_body else {
return Ok(json!({}));
};
if body.is_empty() {
return Ok(json!({}));
}
serde_json::from_slice::<serde_json::Value>(body)
.map_err(|err| GatewayError::Internal(format!("invalid json body: {err}")))
}
@@ -1,7 +1,9 @@
mod background_tasks;
mod gemini_files;
mod routes;
mod video_tasks;
pub(super) use self::background_tasks::maybe_build_local_admin_background_tasks_response;
pub(super) use self::gemini_files::maybe_build_local_admin_gemini_files_response;
pub(super) use self::routes::maybe_build_local_admin_features_response;
pub(crate) use self::video_tasks::maybe_build_local_admin_video_tasks_response;
@@ -1,9 +1,19 @@
use super::{gemini_files, video_tasks};
use super::{background_tasks, gemini_files, video_tasks};
use crate::handlers::admin::request::{AdminRouteRequest, AdminRouteResult};
pub(crate) async fn maybe_build_local_admin_features_response(
request: AdminRouteRequest<'_>,
) -> AdminRouteResult {
if let Some(response) = background_tasks::maybe_build_local_admin_background_tasks_response(
&request.state(),
&request.request_context(),
request.request_body(),
)
.await?
{
return Ok(Some(response));
}
if let Some(response) = video_tasks::maybe_build_local_admin_video_tasks_response(
&request.state(),
&request.request_context(),
@@ -356,6 +356,164 @@ async fn admin_monitoring_trace_request_enriches_proxy_timing_from_usage_audit()
payload["candidates"][0]["extra_data"]["proxy"]["timing"]["response_wait_ms"],
json!(475)
);
assert!(payload["candidates"][0]["extra_data"]
.get("upstream_response")
.is_none());
}
#[tokio::test]
async fn admin_monitoring_trace_request_exposes_request_path_from_usage_audit() {
let mut candidate = sample_candidate(
"cand-used",
"request-1",
0,
RequestCandidateStatus::Failed,
Some(101),
Some(33),
Some(403),
);
candidate.extra_data = Some(json!({
"client_api_format": "gemini:generate_content"
}));
let request_candidates = Arc::new(InMemoryRequestCandidateRepository::seed(vec![candidate]));
let provider_catalog = Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![sample_provider()],
vec![sample_endpoint()],
vec![sample_key()],
));
let mut usage = sample_usage(
"request-1",
"provider-1",
"OpenAI",
40,
0.02,
"failed",
Some(403),
100,
);
usage.candidate_id = Some("cand-used".to_string());
usage.request_metadata = Some(json!({
"request_path": "/v1beta/models/gemini-2.5-pro:generateContent",
"request_query_string": "alt=sse"
}));
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed(vec![usage]));
let data_state =
crate::data::GatewayDataState::with_request_candidate_and_usage_repository_for_tests(
request_candidates,
usage_repository,
)
.with_provider_catalog_reader(provider_catalog);
let state = AppState::new()
.expect("state should build")
.with_data_state_for_tests(data_state);
let context = request_context(http::Method::GET, "/api/admin/monitoring/trace/request-1");
let response = local_monitoring_response(&state, &context)
.await
.expect("handler should not error")
.expect("route should be handled locally");
assert_eq!(response.status(), http::StatusCode::OK);
let body = to_bytes(response.into_body(), usize::MAX)
.await
.expect("body should read");
let payload: serde_json::Value = serde_json::from_slice(&body).expect("json body should parse");
assert_eq!(
payload["request_path"],
json!("/v1beta/models/gemini-2.5-pro:generateContent")
);
assert_eq!(payload["request_query_string"], json!("alt=sse"));
assert_eq!(
payload["request_path_and_query"],
json!("/v1beta/models/gemini-2.5-pro:generateContent?alt=sse")
);
assert_eq!(
payload["candidates"][0]["extra_data"]["request_path_and_query"],
json!("/v1beta/models/gemini-2.5-pro:generateContent?alt=sse")
);
}
#[tokio::test]
async fn admin_monitoring_trace_request_exposes_failed_candidate_upstream_response_boundary() {
let candidate = sample_candidate(
"cand-used",
"request-1",
0,
RequestCandidateStatus::Failed,
Some(101),
Some(33),
Some(302),
);
let request_candidates = Arc::new(InMemoryRequestCandidateRepository::seed(vec![candidate]));
let provider_catalog = Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![sample_provider()],
vec![sample_endpoint()],
vec![sample_key()],
));
let mut usage = sample_usage(
"request-1",
"provider-1",
"OpenAI",
40,
0.02,
"failed",
Some(302),
100,
);
usage.candidate_id = Some("cand-used".to_string());
usage.response_headers = Some(json!({
"location": "/",
"content-type": "text/html"
}));
usage.client_response_headers = Some(json!({
"content-type": "application/json",
"x-aether-upstream-status": "302"
}));
usage.client_response_body = Some(json!({
"error": {
"type": "execution_runtime_non_success_status",
"message": "execution runtime stream returned non-success status 302",
"upstream_status": 302,
"location": "/"
}
}));
usage.request_metadata = Some(json!({
"client_response_status_code": 502
}));
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed(vec![usage]));
let data_state =
crate::data::GatewayDataState::with_request_candidate_and_usage_repository_for_tests(
request_candidates,
usage_repository,
)
.with_provider_catalog_reader(provider_catalog);
let state = AppState::new()
.expect("state should build")
.with_data_state_for_tests(data_state);
let context = request_context(http::Method::GET, "/api/admin/monitoring/trace/request-1");
let response = local_monitoring_response(&state, &context)
.await
.expect("handler should not error")
.expect("route should be handled locally");
assert_eq!(response.status(), http::StatusCode::OK);
let body = to_bytes(response.into_body(), usize::MAX)
.await
.expect("body should read");
let payload: serde_json::Value = serde_json::from_slice(&body).expect("json body should parse");
let extra = &payload["candidates"][0]["extra_data"];
assert_eq!(extra["upstream_response"]["status_code"], json!(302));
assert_eq!(
extra["upstream_response"]["headers"]["location"],
json!("/")
);
assert!(extra["upstream_response"]["body"].is_null());
assert!(extra.get("client_response").is_none());
assert!(extra.get("provider_response").is_none());
}
#[tokio::test]
@@ -1,11 +1,10 @@
use crate::handlers::admin::provider::delete_task::run_admin_provider_delete_task;
use crate::handlers::admin::provider::shared::paths::{
admin_provider_delete_task_parts, admin_provider_id_for_manage_path,
};
use crate::handlers::admin::provider::shared::support::build_admin_provider_delete_task_payload;
use crate::handlers::admin::request::{AdminAppState, AdminRequestContext};
use crate::handlers::admin::shared::attach_admin_audit_response;
use crate::{GatewayError, LocalProviderDeleteTaskState};
use crate::GatewayError;
use axum::{
body::Body,
http,
@@ -13,8 +12,6 @@ use axum::{
Json,
};
use serde_json::json;
use tracing::warn;
use uuid::Uuid;
fn build_admin_provider_not_found_response(detail: impl Into<String>) -> Response<Body> {
(
@@ -60,52 +57,18 @@ pub(crate) async fn maybe_build_local_admin_provider_delete_task_response(
"Provider 不存在",
)));
};
let Some(provider) = state
.read_provider_catalog_providers_by_ids(std::slice::from_ref(&provider_id))
.await?
.into_iter()
.next()
let Some(task_id) =
crate::task_runtime::submit_provider_delete_task(state, &provider_id, Some("admin"))
.await?
else {
return Ok(Some(build_admin_provider_not_found_response(
"提供商不存在",
)));
};
let task_id = Uuid::new_v4().simple().to_string()[..16].to_string();
let pending_task = LocalProviderDeleteTaskState {
task_id: task_id.clone(),
provider_id: provider.id.clone(),
status: "pending".to_string(),
stage: "queued".to_string(),
total_keys: 0,
deleted_keys: 0,
total_endpoints: 0,
deleted_endpoints: 0,
message: "delete task submitted".to_string(),
};
state.put_provider_delete_task(pending_task.clone());
if let Err(err) = state
.run_admin_provider_delete_task(&provider.id, &task_id)
.await
{
warn!(
"gateway admin provider delete task failed for provider {}: {:?}",
provider.id, err
);
state.put_provider_delete_task(LocalProviderDeleteTaskState {
task_id: task_id.clone(),
provider_id: provider.id.clone(),
status: "failed".to_string(),
stage: "failed".to_string(),
total_keys: 0,
deleted_keys: 0,
total_endpoints: 0,
deleted_endpoints: 0,
message: format!("provider delete failed: {err:?}"),
});
}
return Ok(Some(attach_admin_audit_response(
Json(json!({
"task_id": task_id,
"run_id": task_id,
"status": "pending",
"message": "删除任务已提交,提供商已进入后台删除队列",
}))
@@ -113,7 +76,7 @@ pub(crate) async fn maybe_build_local_admin_provider_delete_task_response(
"admin_provider_delete_queued",
"delete_provider",
"provider",
&provider.id,
&provider_id,
)));
}
@@ -15,7 +15,14 @@ use crate::handlers::admin::provider::oauth::state::{
};
use crate::handlers::admin::provider::shared::paths::admin_provider_oauth_batch_import_task_provider_id;
use crate::handlers::admin::request::{AdminAppState, AdminRequestContext};
use crate::task_runtime::{
append_event_with_logging, now_unix_secs, task_definition, update_run_status,
upsert_run_with_logging, TASK_KEY_PROVIDER_OAUTH_BATCH_IMPORT,
};
use crate::GatewayError;
use aether_data_contracts::repository::background_tasks::{
BackgroundTaskKind, BackgroundTaskStatus, UpsertBackgroundTaskRun,
};
use axum::{
body::Bytes,
http,
@@ -133,11 +140,7 @@ pub(in super::super) async fn handle_admin_provider_oauth_start_batch_import_tas
}
let task_id = Uuid::new_v4().to_string();
let created_at = SystemTime::now()
.duration_since(UNIX_EPOCH)
.ok()
.map(|duration| duration.as_secs())
.unwrap_or(0);
let created_at = now_unix_secs();
let submitted_state = build_admin_provider_oauth_batch_task_state(
&task_id,
&provider_id,
@@ -167,6 +170,50 @@ pub(in super::super) async fn handle_admin_provider_oauth_start_batch_import_tas
));
}
if state.has_background_task_data_writer() {
let max_attempts = task_definition(TASK_KEY_PROVIDER_OAUTH_BATCH_IMPORT)
.map(|item| item.retry_policy.max_attempts)
.unwrap_or(1);
let run = UpsertBackgroundTaskRun {
id: task_id.clone(),
task_key: TASK_KEY_PROVIDER_OAUTH_BATCH_IMPORT.to_string(),
kind: BackgroundTaskKind::OnDemand,
trigger: "manual".to_string(),
status: BackgroundTaskStatus::Queued,
attempt: 1,
max_attempts,
owner_instance: Some(state.app().tunnel.local_instance_id().to_string()),
progress_percent: 0,
progress_message: Some("provider oauth batch import queued".to_string()),
payload_json: Some(json!({
"provider_id": provider_id.clone(),
"provider_type": provider_type.clone(),
"total": total,
})),
result_json: None,
error_message: None,
cancel_requested: false,
created_by: Some("admin".to_string()),
created_at_unix_secs: created_at,
started_at_unix_secs: None,
finished_at_unix_secs: None,
updated_at_unix_secs: created_at,
};
let _ = upsert_run_with_logging(state.app(), run).await;
append_event_with_logging(
state.app(),
&task_id,
"queued",
"provider oauth batch import queued",
Some(json!({
"provider_id": provider_id.clone(),
"provider_type": provider_type.clone(),
"total": total,
})),
)
.await;
}
let task_state = state.cloned_app();
let task_id_for_worker = task_id.clone();
let provider_id_for_worker = provider_id.clone();
@@ -201,6 +248,27 @@ pub(in super::super) async fn handle_admin_provider_oauth_start_batch_import_tas
.save_provider_oauth_batch_task_payload(&task_id_for_worker, &processing_state)
.await;
let _ = update_run_status(
&task_state,
&task_id_for_worker,
BackgroundTaskStatus::Running,
Some(1),
Some("provider oauth batch import started".to_string()),
None,
None,
Some(started_at),
None,
)
.await;
append_event_with_logging(
&task_state,
&task_id_for_worker,
"running",
"provider oauth batch import started",
None,
)
.await;
let mut progress_reporter = BatchTaskProgressReporter {
app: task_state.clone(),
task_id: task_id_for_worker.clone(),
@@ -270,6 +338,34 @@ pub(in super::super) async fn handle_admin_provider_oauth_start_batch_import_tas
let _ = AdminAppState::new(&task_state)
.save_provider_oauth_batch_task_payload(&task_id_for_worker, &completed_state)
.await;
let _ = update_run_status(
&task_state,
&task_id_for_worker,
BackgroundTaskStatus::Succeeded,
Some(100),
Some(message),
Some(json!({
"provider_id": provider_id_for_worker,
"provider_type": provider_type_for_worker,
"total": outcome.total,
"success": outcome.success,
"failed": outcome.failed,
"created_count": created_count,
"replaced_count": replaced_count,
})),
None,
None,
Some(finished_at),
)
.await;
append_event_with_logging(
&task_state,
&task_id_for_worker,
"succeeded",
"provider oauth batch import completed",
None,
)
.await;
}
Err(err) => {
let finished_at = SystemTime::now()
@@ -299,6 +395,26 @@ pub(in super::super) async fn handle_admin_provider_oauth_start_batch_import_tas
let _ = AdminAppState::new(&task_state)
.save_provider_oauth_batch_task_payload(&task_id_for_worker, &failed_state)
.await;
let _ = update_run_status(
&task_state,
&task_id_for_worker,
BackgroundTaskStatus::Failed,
Some(100),
Some("provider oauth batch import failed".to_string()),
None,
Some(error_message.clone()),
None,
Some(finished_at),
)
.await;
append_event_with_logging(
&task_state,
&task_id_for_worker,
"failed",
"provider oauth batch import failed",
Some(json!({ "error": error_message.clone() })),
)
.await;
tracing::warn!(
task_id = %task_id_for_worker,
provider_id = %provider_id_for_worker,
@@ -4,6 +4,7 @@ use super::quota::codex::refresh_codex_provider_quota_locally;
use super::quota::kiro::refresh_kiro_provider_quota_locally;
use crate::handlers::admin::request::AdminAppState;
use crate::provider_key_auth::provider_key_is_oauth_managed;
use crate::task_runtime::{spawn_fire_and_forget, TASK_KEY_PROVIDER_OAUTH_ACCOUNT_REFRESH};
use crate::{AppState, GatewayError};
use aether_contracts::ProxySnapshot;
use aether_data_contracts::repository::provider_catalog::{
@@ -170,7 +171,7 @@ pub(crate) fn spawn_provider_oauth_account_state_refresh_after_update(
key_id: String,
proxy_override: Option<ProxySnapshot>,
) {
tokio::spawn(async move {
spawn_fire_and_forget(TASK_KEY_PROVIDER_OAUTH_ACCOUNT_REFRESH, async move {
let _ = refresh_provider_oauth_account_state_after_update(
&AdminAppState::new(&app),
&provider,
@@ -1,5 +1,6 @@
use super::actions::admin_provider_ops_local_action_response;
use crate::handlers::admin::request::AdminAppState;
use crate::task_runtime::{spawn_fire_and_forget, TASK_KEY_PROVIDER_BALANCE_REFRESH};
use serde_json::{json, Value};
use std::collections::HashSet;
use std::time::Duration;
@@ -130,7 +131,7 @@ pub(super) async fn spawn_admin_provider_ops_balance_refresh(
let app = state.cloned_app();
let provider_id = provider_id.to_string();
tokio::spawn(async move {
spawn_fire_and_forget(TASK_KEY_PROVIDER_BALANCE_REFRESH, async move {
let permit = match tokio::time::timeout(
Duration::from_secs(5),
ADMIN_PROVIDER_OPS_BALANCE_REFRESH_SEMAPHORE.acquire(),
@@ -76,6 +76,14 @@ impl<'a> AdminAppState<'a> {
self.app.has_gemini_file_mapping_data_writer()
}
pub(crate) fn has_background_task_data_reader(&self) -> bool {
self.app.has_background_task_data_reader()
}
pub(crate) fn has_background_task_data_writer(&self) -> bool {
self.app.has_background_task_data_writer()
}
pub(crate) fn has_auth_api_key_data_reader(&self) -> bool {
self.app.has_auth_api_key_data_reader()
}
@@ -2,6 +2,59 @@ use super::{AdminAppState, AdminCancelVideoTaskError};
use crate::GatewayError;
impl<'a> AdminAppState<'a> {
pub(crate) async fn list_background_task_runs(
&self,
query: &aether_data_contracts::repository::background_tasks::BackgroundTaskListQuery,
) -> Result<
aether_data_contracts::repository::background_tasks::StoredBackgroundTaskRunPage,
GatewayError,
> {
self.app.list_background_task_runs(query).await
}
pub(crate) async fn find_background_task_run(
&self,
run_id: &str,
) -> Result<
Option<aether_data_contracts::repository::background_tasks::StoredBackgroundTaskRun>,
GatewayError,
> {
self.app.find_background_task_run(run_id).await
}
pub(crate) async fn list_background_task_events(
&self,
run_id: &str,
offset: usize,
limit: usize,
) -> Result<
Vec<aether_data_contracts::repository::background_tasks::StoredBackgroundTaskEvent>,
GatewayError,
> {
self.app
.list_background_task_events(run_id, offset, limit)
.await
}
pub(crate) async fn summarize_background_task_runs(
&self,
) -> Result<
aether_data_contracts::repository::background_tasks::BackgroundTaskSummary,
GatewayError,
> {
self.app.summarize_background_task_runs().await
}
pub(crate) async fn request_cancel_background_task_run(
&self,
run_id: &str,
updated_at_unix_secs: u64,
) -> Result<bool, GatewayError> {
self.app
.request_cancel_background_task_run(run_id, updated_at_unix_secs)
.await
}
pub(crate) async fn list_gemini_file_mappings(
&self,
query: &aether_data::repository::gemini_file_mappings::GeminiFileMappingListQuery,
@@ -18,7 +18,6 @@ use crate::handlers::admin::system::shared::settings::{
};
use crate::handlers::admin::system::shared::smtp::build_admin_smtp_test_payload;
use crate::GatewayError;
use aether_data::repository::system::AdminSystemPurgeTarget;
use axum::{
body::{Body, Bytes},
http,
@@ -26,6 +25,7 @@ use axum::{
Json,
};
use serde_json::json;
use std::time::Instant;
pub(super) async fn maybe_build_local_admin_core_system_response(
state: &AdminAppState<'_>,
@@ -194,15 +194,31 @@ pub(super) async fn maybe_build_local_admin_core_system_response(
)));
}
if let Some((target, action, object_type, object_id)) =
admin_system_purge_target_for_route_kind(decision.route_kind.as_deref())
if decision.route_kind.as_deref() == Some("cleanup_runs")
&& request_method == http::Method::GET
&& request_path == "/api/admin/system/cleanup/runs"
{
let records = crate::maintenance::list_admin_cleanup_run_records(&state.app().data)
.await
.map_err(|err| GatewayError::Internal(err.to_string()))?;
return Ok(Some(Json(json!({ "items": records })).into_response()));
}
if let Some((task_kind, action, object_type, object_id)) =
admin_system_purge_task_for_route_kind(decision.route_kind.as_deref())
{
if request_method != http::Method::POST {
return Ok(None);
}
let task = crate::maintenance::start_admin_system_purge_task(state.cloned_app(), task_kind)
.await?;
return Ok(Some(attach_admin_audit_response(
Json(build_admin_system_purge_payload(state, target).await?).into_response(),
"admin_system_data_purged",
Json(json!({
"message": task.message.clone(),
"task": task,
}))
.into_response(),
"admin_system_purge_task_started",
action,
object_type,
object_id,
@@ -435,100 +451,60 @@ pub(super) async fn maybe_build_local_admin_core_system_response(
Ok(None)
}
fn admin_system_purge_target_for_route_kind(
fn admin_system_purge_task_for_route_kind(
route_kind: Option<&str>,
) -> Option<(
AdminSystemPurgeTarget,
crate::maintenance::AdminCleanupTaskKind,
&'static str,
&'static str,
&'static str,
)> {
match route_kind {
Some("purge_config") => Some((
AdminSystemPurgeTarget::Config,
"purge_system_config",
crate::maintenance::AdminCleanupTaskKind::Config,
"purge_system_config_async",
"system_config",
"global",
)),
Some("purge_users") => Some((
AdminSystemPurgeTarget::Users,
"purge_non_admin_users",
crate::maintenance::AdminCleanupTaskKind::Users,
"purge_non_admin_users_async",
"users",
"non_admin",
)),
Some("purge_usage") => Some((
AdminSystemPurgeTarget::Usage,
"purge_usage_records",
crate::maintenance::AdminCleanupTaskKind::Usage,
"purge_usage_records_async",
"usage",
"all",
)),
Some("purge_audit_logs") => Some((
AdminSystemPurgeTarget::AuditLogs,
"purge_audit_logs",
crate::maintenance::AdminCleanupTaskKind::AuditLogs,
"purge_audit_logs_async",
"audit_logs",
"all",
)),
Some("purge_request_bodies") => Some((
AdminSystemPurgeTarget::RequestBodies,
"purge_request_bodies",
Some("purge_request_bodies") | Some("purge_request_bodies_task") => Some((
crate::maintenance::AdminCleanupTaskKind::RequestBodies,
"purge_request_bodies_async",
"request_bodies",
"all",
)),
Some("purge_stats") => Some((AdminSystemPurgeTarget::Stats, "purge_stats", "stats", "all")),
Some("purge_stats") => Some((
crate::maintenance::AdminCleanupTaskKind::Stats,
"purge_stats_async",
"stats",
"all",
)),
_ => None,
}
}
async fn build_admin_system_purge_payload(
state: &AdminAppState<'_>,
target: AdminSystemPurgeTarget,
) -> Result<serde_json::Value, GatewayError> {
let summary = state.purge_admin_system_data(target).await?;
let total = summary.total();
let affected = summary.affected.clone();
if target == AdminSystemPurgeTarget::Stats {
let rebuild = state.rebuild_admin_stats_once().await?;
let message = if rebuild.capped {
format!(
"统计聚合已清空,已重建 {} 个小时桶和 {} 个日桶,仍有历史统计待后台任务继续重建",
rebuild.hourly_buckets, rebuild.daily_buckets
)
} else {
format!(
"统计聚合已清空并重建,删除 {} 行,重建 {} 个小时桶和 {} 个日桶",
total, rebuild.hourly_buckets, rebuild.daily_buckets
)
};
return Ok(json!({
"message": message,
"deleted": affected,
"rebuilt": {
"hourly_buckets": rebuild.hourly_buckets,
"daily_buckets": rebuild.daily_buckets,
"capped": rebuild.capped,
},
}));
}
let (message, count_key) = match target {
AdminSystemPurgeTarget::Config => ("系统配置已清空", "deleted"),
AdminSystemPurgeTarget::Users => ("非管理员用户已清空", "deleted"),
AdminSystemPurgeTarget::Usage => ("使用记录已清空", "deleted"),
AdminSystemPurgeTarget::AuditLogs => ("审计日志已清空", "deleted"),
AdminSystemPurgeTarget::RequestBodies => ("请求/响应体已清空", "cleaned"),
AdminSystemPurgeTarget::Stats => unreachable!("stats handled above"),
};
Ok(json!({
"message": format!("{message},影响 {} 行", total),
count_key: affected,
}))
}
async fn build_admin_system_cleanup_payload(
state: &AdminAppState<'_>,
) -> Result<serde_json::Value, GatewayError> {
let started_at_unix_secs = chrono::Utc::now().timestamp().max(0) as u64;
let started_at = Instant::now();
let summary = state.run_admin_system_cleanup_once().await?;
let cleaned = json!({
"audit_logs": summary.audit_logs_deleted,
@@ -558,6 +534,17 @@ async fn build_admin_system_cleanup_payload(
.saturating_add(summary.usage.keys_cleaned)
.saturating_add(summary.usage.records_deleted);
crate::maintenance::record_completed_cleanup_run(
&state.app().data,
"system_cleanup",
"manual",
started_at_unix_secs,
started_at,
cleaned.clone(),
format!("系统清理已执行,影响 {total} 项"),
)
.await;
Ok(json!({
"message": format!("系统清理已执行,影响 {} 项", total),
"cleaned": cleaned,