Merge branch 'pr-503'

# Conflicts:
#	apps/aether-gateway/src/handlers/admin/provider/summary/value.rs
#	apps/aether-gateway/src/lib.rs
#	apps/aether-gateway/src/maintenance/mod.rs
#	apps/aether-gateway/src/maintenance/runtime/workers.rs
#	frontend/src/api/endpoints/types/provider.ts
This commit is contained in:
fawney19
2026-05-22 00:19:23 +08:00
36 changed files with 2595 additions and 604 deletions
@@ -0,0 +1,463 @@
use std::collections::HashMap;
use aether_data_contracts::repository::provider_catalog::{
StoredProviderCatalogEndpoint, StoredProviderCatalogProvider,
};
use futures_util::stream::{self, StreamExt};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use tracing::warn;
use crate::admin_api::{
admin_provider_ops_local_action_response, store_admin_provider_ops_balance_cache, AdminAppState,
};
use crate::important_notification::{
important_notification_dispatch_ready, send_important_notification, ImportantNotification,
};
use crate::{AppState, GatewayError};
use super::PROVIDER_QUOTA_ALERT_CONCURRENCY;
const PROVIDER_QUOTA_ALERT_STATE_PREFIX: &str = "provider_ops:quota_alert:";
const PROVIDER_QUOTA_ALERT_DEFAULT_FETCH_INTERVAL_SECS: u64 = 30;
const PROVIDER_QUOTA_ALERT_MIN_FETCH_INTERVAL_SECS: u64 = 30;
const PROVIDER_QUOTA_ALERT_REPEAT_COOLDOWN_SECS: u64 = 24 * 60 * 60;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct ProviderQuotaAlertRunSummary {
pub(crate) checked: usize,
pub(crate) alerted: usize,
pub(crate) skipped: usize,
pub(crate) failed: usize,
}
#[derive(Debug, Clone)]
struct ProviderQuotaAlertTarget {
provider: StoredProviderCatalogProvider,
config: ProviderQuotaAlertConfig,
}
#[derive(Debug, Clone, Copy)]
struct ProviderQuotaAlertConfig {
enabled: bool,
threshold_amount: f64,
fetch_interval_seconds: u64,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
struct ProviderQuotaAlertRuntimeState {
#[serde(default)]
last_checked_at: u64,
#[serde(default)]
last_available: Option<f64>,
#[serde(default)]
below_threshold: bool,
#[serde(default)]
last_notified_at: Option<u64>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ProviderQuotaAlertStatus {
Checked,
Alerted,
Skipped,
Failed,
}
pub(crate) async fn perform_provider_quota_alert_once(
state: &AppState,
) -> Result<ProviderQuotaAlertRunSummary, GatewayError> {
if !state.has_provider_catalog_data_reader() {
return Ok(ProviderQuotaAlertRunSummary {
checked: 0,
alerted: 0,
skipped: 0,
failed: 0,
});
}
if !important_notification_dispatch_ready(state).await? {
return Ok(ProviderQuotaAlertRunSummary {
checked: 0,
alerted: 0,
skipped: 0,
failed: 0,
});
}
let now_unix_secs = now_unix_secs();
let targets = select_provider_quota_alert_targets(state, now_unix_secs).await?;
if targets.is_empty() {
return Ok(ProviderQuotaAlertRunSummary {
checked: 0,
alerted: 0,
skipped: 0,
failed: 0,
});
}
let provider_ids = targets
.iter()
.map(|target| target.provider.id.clone())
.collect::<Vec<_>>();
let mut endpoints_by_provider = HashMap::<String, Vec<StoredProviderCatalogEndpoint>>::new();
for endpoint in state
.list_provider_catalog_endpoints_by_provider_ids(&provider_ids)
.await?
{
endpoints_by_provider
.entry(endpoint.provider_id.clone())
.or_default()
.push(endpoint);
}
let mut results = stream::iter(targets.into_iter().map(|target| {
let state = state.clone();
let provider_id = target.provider.id.clone();
let endpoints = endpoints_by_provider
.get(&provider_id)
.cloned()
.unwrap_or_default();
async move { run_provider_quota_alert_for_provider(&state, target, endpoints).await }
}))
.buffer_unordered(PROVIDER_QUOTA_ALERT_CONCURRENCY);
let mut summary = ProviderQuotaAlertRunSummary {
checked: 0,
alerted: 0,
skipped: 0,
failed: 0,
};
while let Some(status) = results.next().await {
match status {
ProviderQuotaAlertStatus::Checked => summary.checked += 1,
ProviderQuotaAlertStatus::Alerted => {
summary.checked += 1;
summary.alerted += 1;
}
ProviderQuotaAlertStatus::Skipped => summary.skipped += 1,
ProviderQuotaAlertStatus::Failed => summary.failed += 1,
}
}
Ok(summary)
}
async fn select_provider_quota_alert_targets(
state: &AppState,
now_unix_secs: u64,
) -> Result<Vec<ProviderQuotaAlertTarget>, GatewayError> {
let providers = state
.list_provider_catalog_providers(true)
.await?
.into_iter()
.filter_map(|provider| {
let config = provider_quota_alert_config(&provider)?;
(config.enabled).then_some(ProviderQuotaAlertTarget { provider, config })
})
.collect::<Vec<_>>();
let mut due = Vec::new();
for target in providers {
let runtime = read_quota_alert_runtime_state(state, &target.provider.id).await;
let last_checked_at = runtime
.as_ref()
.map(|state| state.last_checked_at)
.unwrap_or(0);
if now_unix_secs.saturating_sub(last_checked_at)
>= target
.config
.fetch_interval_seconds
.max(PROVIDER_QUOTA_ALERT_MIN_FETCH_INTERVAL_SECS)
{
due.push(target);
}
}
Ok(due)
}
async fn run_provider_quota_alert_for_provider(
state: &AppState,
target: ProviderQuotaAlertTarget,
endpoints: Vec<StoredProviderCatalogEndpoint>,
) -> ProviderQuotaAlertStatus {
let provider_id = target.provider.id.clone();
let admin_state = AdminAppState::new(state);
let payload = admin_provider_ops_local_action_response(
&admin_state,
&provider_id,
Some(&target.provider),
&endpoints,
"query_balance",
None,
)
.await;
store_admin_provider_ops_balance_cache(&admin_state, &provider_id, &payload).await;
let now_unix_secs = now_unix_secs();
let payload_status_success = payload.get("status").and_then(Value::as_str) == Some("success");
let Some(total_available) = extract_total_available(&payload) else {
warn!(
provider_id = %provider_id,
payload = %payload,
"provider quota alert skipped because balance payload has no total_available"
);
write_checked_runtime_state_without_balance(state, &provider_id, now_unix_secs).await;
return if payload_status_success {
ProviderQuotaAlertStatus::Skipped
} else {
ProviderQuotaAlertStatus::Failed
};
};
let previous = read_quota_alert_runtime_state(state, &provider_id).await;
let should_notify = provider_quota_alert_should_notify(
now_unix_secs,
total_available,
target.config.threshold_amount,
previous.as_ref(),
);
let mut next = ProviderQuotaAlertRuntimeState {
last_checked_at: now_unix_secs,
last_available: Some(total_available),
below_threshold: total_available <= target.config.threshold_amount,
last_notified_at: previous.and_then(|state| state.last_notified_at),
};
if should_notify {
let report = send_important_notification(
state,
build_provider_quota_alert_notification(
&target.provider,
total_available,
target.config.threshold_amount,
),
)
.await;
let delivered = match &report {
Ok(report) if report.success => true,
Ok(report) => {
warn!(
provider_id = %provider_id,
report = ?report,
"provider quota alert notification did not reach any channel"
);
false
}
Err(err) => {
warn!(
provider_id = %provider_id,
error = ?err,
"provider quota alert notification failed"
);
false
}
};
if delivered {
next.last_notified_at = Some(now_unix_secs);
write_quota_alert_runtime_state(state, &provider_id, &next).await;
return ProviderQuotaAlertStatus::Alerted;
}
write_quota_alert_runtime_state(state, &provider_id, &next).await;
return ProviderQuotaAlertStatus::Failed;
}
write_quota_alert_runtime_state(state, &provider_id, &next).await;
ProviderQuotaAlertStatus::Checked
}
fn provider_quota_alert_config(
provider: &StoredProviderCatalogProvider,
) -> Option<ProviderQuotaAlertConfig> {
let quota_alert = provider
.config
.as_ref()
.and_then(Value::as_object)
.and_then(|config| config.get("provider_ops"))
.and_then(Value::as_object)
.and_then(|provider_ops| provider_ops.get("quota_alert"))
.and_then(Value::as_object)?;
let enabled = quota_alert
.get("enabled")
.and_then(Value::as_bool)
.unwrap_or(false);
let threshold_amount = quota_alert
.get("threshold_amount")
.and_then(value_as_f64)
.filter(|value| value.is_finite() && *value >= 0.0)
.unwrap_or(0.0);
let fetch_interval_seconds = quota_alert
.get("fetch_interval_seconds")
.and_then(Value::as_u64)
.unwrap_or(PROVIDER_QUOTA_ALERT_DEFAULT_FETCH_INTERVAL_SECS)
.max(PROVIDER_QUOTA_ALERT_MIN_FETCH_INTERVAL_SECS);
Some(ProviderQuotaAlertConfig {
enabled,
threshold_amount,
fetch_interval_seconds,
})
}
fn provider_quota_alert_should_notify(
now_unix_secs: u64,
total_available: f64,
threshold_amount: f64,
previous: Option<&ProviderQuotaAlertRuntimeState>,
) -> bool {
if total_available > threshold_amount {
return false;
}
let Some(previous) = previous else {
return true;
};
if !previous.below_threshold {
return true;
}
previous
.last_notified_at
.map(|last| now_unix_secs.saturating_sub(last) >= PROVIDER_QUOTA_ALERT_REPEAT_COOLDOWN_SECS)
.unwrap_or(true)
}
fn extract_total_available(payload: &Value) -> Option<f64> {
if payload.get("status").and_then(Value::as_str) != Some("success") {
return None;
}
payload
.get("data")
.and_then(|data| data.get("total_available"))
.and_then(value_as_f64)
.filter(|value| value.is_finite())
}
fn value_as_f64(value: &Value) -> Option<f64> {
value.as_f64().or_else(|| {
value
.as_str()
.and_then(|raw| raw.trim().parse::<f64>().ok())
})
}
fn build_provider_quota_alert_notification(
provider: &StoredProviderCatalogProvider,
total_available: f64,
threshold_amount: f64,
) -> ImportantNotification {
let title = format!("提供商额度提醒:{}", provider.name);
let body = format!(
"提供商 `{}` 当前剩余额度为 `{:.4}`,已低于或等于提醒阈值 `{:.4}`。\n\nProvider ID: `{}`",
provider.name, total_available, threshold_amount, provider.id
);
let text_body = format!(
"提供商 {} 当前剩余额度为 {:.4},已低于或等于提醒阈值 {:.4}。\n\nProvider ID: {}",
provider.name, total_available, threshold_amount, provider.id
);
ImportantNotification {
title,
markdown_body: body,
text_body,
}
}
async fn read_quota_alert_runtime_state(
state: &AppState,
provider_id: &str,
) -> Option<ProviderQuotaAlertRuntimeState> {
let key = provider_quota_alert_state_key(provider_id);
state
.runtime_kv_get(&key)
.await
.ok()
.flatten()
.and_then(|raw| serde_json::from_str::<ProviderQuotaAlertRuntimeState>(&raw).ok())
}
async fn write_checked_runtime_state_without_balance(
state: &AppState,
provider_id: &str,
now_unix_secs: u64,
) {
let mut next = read_quota_alert_runtime_state(state, provider_id)
.await
.unwrap_or_default();
next.last_checked_at = now_unix_secs;
write_quota_alert_runtime_state(state, provider_id, &next).await;
}
async fn write_quota_alert_runtime_state(
state: &AppState,
provider_id: &str,
runtime_state: &ProviderQuotaAlertRuntimeState,
) {
let Ok(serialized) = serde_json::to_string(runtime_state) else {
return;
};
if let Err(err) = state
.runtime_state()
.kv_set(
&provider_quota_alert_state_key(provider_id),
serialized,
None,
)
.await
{
warn!(
error = %err,
provider_id,
"failed to write provider quota alert runtime state"
);
}
}
fn provider_quota_alert_state_key(provider_id: &str) -> String {
format!("{PROVIDER_QUOTA_ALERT_STATE_PREFIX}{provider_id}")
}
fn now_unix_secs() -> u64 {
chrono::Utc::now().timestamp().max(0) as u64
}
#[cfg(test)]
mod tests {
use super::{
provider_quota_alert_should_notify, ProviderQuotaAlertRuntimeState,
PROVIDER_QUOTA_ALERT_REPEAT_COOLDOWN_SECS,
};
#[test]
fn quota_alert_notifies_on_first_drop_and_after_cooldown() {
assert!(provider_quota_alert_should_notify(100, 3.0, 5.0, None));
assert!(provider_quota_alert_should_notify(
100,
3.0,
5.0,
Some(&ProviderQuotaAlertRuntimeState {
below_threshold: false,
..ProviderQuotaAlertRuntimeState::default()
})
));
assert!(!provider_quota_alert_should_notify(
100,
3.0,
5.0,
Some(&ProviderQuotaAlertRuntimeState {
below_threshold: true,
last_notified_at: Some(90),
..ProviderQuotaAlertRuntimeState::default()
})
));
assert!(provider_quota_alert_should_notify(
100 + PROVIDER_QUOTA_ALERT_REPEAT_COOLDOWN_SECS,
3.0,
5.0,
Some(&ProviderQuotaAlertRuntimeState {
below_threshold: true,
last_notified_at: Some(100),
..ProviderQuotaAlertRuntimeState::default()
})
));
}
#[test]
fn quota_alert_does_not_notify_above_threshold() {
assert!(!provider_quota_alert_should_notify(100, 6.0, 5.0, None));
}
}
@@ -10,22 +10,22 @@ use super::{
cleanup_processed_usage_counter_deltas_once, duration_until_next_daily_run,
duration_until_next_db_maintenance_run, duration_until_next_stats_aggregation_run,
duration_until_next_stats_hourly_aggregation_run, maintenance_timezone, parse_hhmm_time,
perform_oauth_token_refresh_once, provider_checkin_schedule, run_audit_cleanup_once,
run_db_maintenance_once, run_gemini_file_mapping_cleanup_once, run_pending_cleanup_once,
run_pool_monitor_once, run_provider_checkin_once, run_proxy_node_metrics_cleanup_once,
run_proxy_node_stale_cleanup_once, run_proxy_upgrade_rollout_once,
run_request_candidate_cleanup_once, run_stats_aggregation_once,
perform_oauth_token_refresh_once, perform_provider_quota_alert_once, provider_checkin_schedule,
run_audit_cleanup_once, run_db_maintenance_once, run_gemini_file_mapping_cleanup_once,
run_pending_cleanup_once, run_pool_monitor_once, run_provider_checkin_once,
run_proxy_node_metrics_cleanup_once, run_proxy_node_stale_cleanup_once,
run_proxy_upgrade_rollout_once, run_request_candidate_cleanup_once, run_stats_aggregation_once,
run_stats_hourly_aggregation_once, run_usage_cleanup_once, run_usage_counter_flush_once,
run_wallet_daily_usage_aggregation_once, AUDIT_LOG_CLEANUP_INTERVAL,
GEMINI_FILE_MAPPING_CLEANUP_INTERVAL, OAUTH_TOKEN_REFRESH_INTERVAL, PENDING_CLEANUP_INTERVAL,
POOL_MONITOR_INTERVAL, PROVIDER_CHECKIN_DEFAULT_TIME, PROXY_NODE_METRICS_CLEANUP_HOUR,
PROXY_NODE_METRICS_CLEANUP_MINUTE, PROXY_NODE_STALE_SWEEP_INTERVAL,
PROXY_UPGRADE_ROLLOUT_INTERVAL, REQUEST_CANDIDATE_CLEANUP_INTERVAL, USAGE_CLEANUP_HOUR,
USAGE_CLEANUP_MINUTE, USAGE_COUNTER_DELTA_CLEANUP_BATCH_SIZE,
USAGE_COUNTER_DELTA_CLEANUP_INTERVAL, USAGE_COUNTER_DELTA_RETENTION_SECS,
USAGE_COUNTER_FLUSH_BATCH_SIZE, USAGE_COUNTER_FLUSH_CATCH_UP_BURST_LIMIT,
USAGE_COUNTER_FLUSH_INTERVAL, WALLET_DAILY_USAGE_AGGREGATION_HOUR,
WALLET_DAILY_USAGE_AGGREGATION_MINUTE,
POOL_MONITOR_INTERVAL, PROVIDER_CHECKIN_DEFAULT_TIME, PROVIDER_QUOTA_ALERT_INTERVAL,
PROXY_NODE_METRICS_CLEANUP_HOUR, PROXY_NODE_METRICS_CLEANUP_MINUTE,
PROXY_NODE_STALE_SWEEP_INTERVAL, PROXY_UPGRADE_ROLLOUT_INTERVAL,
REQUEST_CANDIDATE_CLEANUP_INTERVAL, USAGE_CLEANUP_HOUR, USAGE_CLEANUP_MINUTE,
USAGE_COUNTER_DELTA_CLEANUP_BATCH_SIZE, USAGE_COUNTER_DELTA_CLEANUP_INTERVAL,
USAGE_COUNTER_DELTA_RETENTION_SECS, USAGE_COUNTER_FLUSH_BATCH_SIZE,
USAGE_COUNTER_FLUSH_CATCH_UP_BURST_LIMIT, USAGE_COUNTER_FLUSH_INTERVAL,
WALLET_DAILY_USAGE_AGGREGATION_HOUR, WALLET_DAILY_USAGE_AGGREGATION_MINUTE,
};
const STATS_DAILY_CATCH_UP_BURST_LIMIT: usize = 14;
@@ -254,6 +254,26 @@ pub(crate) fn spawn_provider_checkin_worker(
}))
}
pub(crate) fn spawn_provider_quota_alert_worker(
state: AppState,
) -> Option<tokio::task::JoinHandle<()>> {
if !state.has_provider_catalog_data_reader() {
return None;
}
Some(tokio::spawn(async move {
let mut interval = tokio::time::interval(PROVIDER_QUOTA_ALERT_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await;
loop {
interval.tick().await;
if let Err(err) = perform_provider_quota_alert_once(&state).await {
log_maintenance_worker_failure("provider_quota_alert", "tick", &err);
}
}
}))
}
pub(crate) fn spawn_oauth_token_refresh_worker(
state: AppState,
) -> Option<tokio::task::JoinHandle<()>> {