mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-11 03:39:49 +08:00
feat(gateway): 补齐号池主动探测后台任务 (#342)
- 解析 pool_advanced 主动探测开关与间隔配置 - 对齐 Python 版本的探测间隔默认值和范围限制 - 新增号池配额主动探测 worker - 支持 codex、kiro、antigravity provider - 使用 Redis 时间戳和 provider 级锁避免重复探测 - 接入 gateway 后台任务并补充测试
This commit is contained in:
@@ -1,8 +1,10 @@
|
|||||||
pub(crate) use crate::handlers::admin::{
|
pub(crate) use crate::handlers::admin::{
|
||||||
admin_provider_ops_local_action_response, build_internal_control_error_response,
|
admin_provider_ops_local_action_response, admin_provider_pool_config,
|
||||||
maybe_build_local_admin_pool_response, maybe_build_local_admin_response, AdminAppState,
|
build_internal_control_error_response, maybe_build_local_admin_pool_response,
|
||||||
AdminRequestContext, AdminRouteRequest, AdminRouteResponse, AdminRouteResult,
|
maybe_build_local_admin_response, provider_oauth_runtime_endpoint_for_provider,
|
||||||
AdminStatsTimeRange, AdminStatsUsageFilter,
|
refresh_antigravity_provider_quota_locally, refresh_codex_provider_quota_locally,
|
||||||
|
refresh_kiro_provider_quota_locally, AdminAppState, AdminRequestContext, AdminRouteRequest,
|
||||||
|
AdminRouteResponse, AdminRouteResult, AdminStatsTimeRange, AdminStatsUsageFilter,
|
||||||
};
|
};
|
||||||
|
|
||||||
use crate::handlers::admin::{
|
use crate::handlers::admin::{
|
||||||
|
|||||||
@@ -21,7 +21,12 @@ pub(crate) use self::observability::{
|
|||||||
round_to, AdminStatsTimeRange, AdminStatsUsageFilter,
|
round_to, AdminStatsTimeRange, AdminStatsUsageFilter,
|
||||||
};
|
};
|
||||||
pub(crate) use self::provider::oauth::errors::build_internal_control_error_response;
|
pub(crate) use self::provider::oauth::errors::build_internal_control_error_response;
|
||||||
|
pub(crate) use self::provider::oauth::quota::antigravity::refresh_antigravity_provider_quota_locally;
|
||||||
|
pub(crate) use self::provider::oauth::quota::codex::refresh_codex_provider_quota_locally;
|
||||||
|
pub(crate) use self::provider::oauth::quota::kiro::refresh_kiro_provider_quota_locally;
|
||||||
|
pub(crate) use self::provider::oauth::runtime::provider_oauth_runtime_endpoint_for_provider;
|
||||||
pub(crate) use self::provider::ops::providers::actions::admin_provider_ops_local_action_response;
|
pub(crate) use self::provider::ops::providers::actions::admin_provider_ops_local_action_response;
|
||||||
|
pub(crate) use self::provider::pool::config::admin_provider_pool_config;
|
||||||
pub(crate) use self::provider::pool_admin::maybe_build_local_admin_pool_response;
|
pub(crate) use self::provider::pool_admin::maybe_build_local_admin_pool_response;
|
||||||
pub(crate) use self::provider::{
|
pub(crate) use self::provider::{
|
||||||
maybe_build_local_admin_provider_oauth_response, maybe_build_local_admin_providers_response,
|
maybe_build_local_admin_provider_oauth_response, maybe_build_local_admin_providers_response,
|
||||||
|
|||||||
@@ -246,6 +246,8 @@ pub(crate) fn admin_provider_pool_config_from_config_value(
|
|||||||
rate_limit_cooldown_seconds: 300,
|
rate_limit_cooldown_seconds: 300,
|
||||||
overload_cooldown_seconds: 30,
|
overload_cooldown_seconds: 30,
|
||||||
health_policy_enabled: true,
|
health_policy_enabled: true,
|
||||||
|
probing_enabled: false,
|
||||||
|
probing_interval_minutes: 10,
|
||||||
stream_timeout_threshold: 3,
|
stream_timeout_threshold: 3,
|
||||||
stream_timeout_window_seconds: 1800,
|
stream_timeout_window_seconds: 1800,
|
||||||
stream_timeout_cooldown_seconds: 300,
|
stream_timeout_cooldown_seconds: 300,
|
||||||
@@ -299,6 +301,16 @@ pub(crate) fn admin_provider_pool_config_from_config_value(
|
|||||||
.get("health_policy_enabled")
|
.get("health_policy_enabled")
|
||||||
.and_then(Value::as_bool)
|
.and_then(Value::as_bool)
|
||||||
.unwrap_or(true),
|
.unwrap_or(true),
|
||||||
|
probing_enabled: pool_advanced
|
||||||
|
.get("probing_enabled")
|
||||||
|
.and_then(Value::as_bool)
|
||||||
|
.unwrap_or(false),
|
||||||
|
probing_interval_minutes: pool_advanced
|
||||||
|
.get("probing_interval_minutes")
|
||||||
|
.and_then(json_u64)
|
||||||
|
.filter(|value| *value > 0)
|
||||||
|
.map(|value| value.min(1440))
|
||||||
|
.unwrap_or(10),
|
||||||
stream_timeout_threshold: pool_advanced
|
stream_timeout_threshold: pool_advanced
|
||||||
.get("stream_timeout_threshold")
|
.get("stream_timeout_threshold")
|
||||||
.and_then(json_u64)
|
.and_then(json_u64)
|
||||||
@@ -366,6 +378,8 @@ mod tests {
|
|||||||
"rate_limit_cooldown_seconds": 420,
|
"rate_limit_cooldown_seconds": 420,
|
||||||
"overload_cooldown_seconds": 45,
|
"overload_cooldown_seconds": 45,
|
||||||
"health_policy_enabled": false,
|
"health_policy_enabled": false,
|
||||||
|
"probing_enabled": true,
|
||||||
|
"probing_interval_minutes": 20,
|
||||||
"stream_timeout_threshold": 4,
|
"stream_timeout_threshold": 4,
|
||||||
"stream_timeout_window_seconds": 900,
|
"stream_timeout_window_seconds": 900,
|
||||||
"stream_timeout_cooldown_seconds": 180
|
"stream_timeout_cooldown_seconds": 180
|
||||||
@@ -383,11 +397,34 @@ mod tests {
|
|||||||
assert_eq!(config.rate_limit_cooldown_seconds, 420);
|
assert_eq!(config.rate_limit_cooldown_seconds, 420);
|
||||||
assert_eq!(config.overload_cooldown_seconds, 45);
|
assert_eq!(config.overload_cooldown_seconds, 45);
|
||||||
assert!(!config.health_policy_enabled);
|
assert!(!config.health_policy_enabled);
|
||||||
|
assert!(config.probing_enabled);
|
||||||
|
assert_eq!(config.probing_interval_minutes, 20);
|
||||||
assert_eq!(config.stream_timeout_threshold, 4);
|
assert_eq!(config.stream_timeout_threshold, 4);
|
||||||
assert_eq!(config.stream_timeout_window_seconds, 900);
|
assert_eq!(config.stream_timeout_window_seconds, 900);
|
||||||
assert_eq!(config.stream_timeout_cooldown_seconds, 180);
|
assert_eq!(config.stream_timeout_cooldown_seconds, 180);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn clamps_pool_quota_probe_interval_to_python_range() {
|
||||||
|
let provider = sample_provider(json!({
|
||||||
|
"pool_advanced": {
|
||||||
|
"probing_enabled": true,
|
||||||
|
"probing_interval_minutes": 2000,
|
||||||
|
}
|
||||||
|
}));
|
||||||
|
let config = admin_provider_pool_config(&provider).expect("pool config should exist");
|
||||||
|
assert_eq!(config.probing_interval_minutes, 1440);
|
||||||
|
|
||||||
|
let provider = sample_provider(json!({
|
||||||
|
"pool_advanced": {
|
||||||
|
"probing_enabled": true,
|
||||||
|
"probing_interval_minutes": 0,
|
||||||
|
}
|
||||||
|
}));
|
||||||
|
let config = admin_provider_pool_config(&provider).expect("pool config should exist");
|
||||||
|
assert_eq!(config.probing_interval_minutes, 10);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn parses_zero_sticky_session_ttl_to_disable_sticky_sessions() {
|
fn parses_zero_sticky_session_ttl_to_disable_sticky_sessions() {
|
||||||
let config = admin_provider_pool_config_from_config_value(Some(&json!({
|
let config = admin_provider_pool_config_from_config_value(Some(&json!({
|
||||||
|
|||||||
@@ -612,6 +612,8 @@ mod tests {
|
|||||||
rate_limit_cooldown_seconds: 300,
|
rate_limit_cooldown_seconds: 300,
|
||||||
overload_cooldown_seconds: 30,
|
overload_cooldown_seconds: 30,
|
||||||
health_policy_enabled: true,
|
health_policy_enabled: true,
|
||||||
|
probing_enabled: false,
|
||||||
|
probing_interval_minutes: 10,
|
||||||
stream_timeout_threshold: 3,
|
stream_timeout_threshold: 3,
|
||||||
stream_timeout_window_seconds: 1800,
|
stream_timeout_window_seconds: 1800,
|
||||||
stream_timeout_cooldown_seconds: 300,
|
stream_timeout_cooldown_seconds: 300,
|
||||||
|
|||||||
@@ -37,6 +37,8 @@ pub(crate) struct AdminProviderPoolConfig {
|
|||||||
pub(crate) rate_limit_cooldown_seconds: u64,
|
pub(crate) rate_limit_cooldown_seconds: u64,
|
||||||
pub(crate) overload_cooldown_seconds: u64,
|
pub(crate) overload_cooldown_seconds: u64,
|
||||||
pub(crate) health_policy_enabled: bool,
|
pub(crate) health_policy_enabled: bool,
|
||||||
|
pub(crate) probing_enabled: bool,
|
||||||
|
pub(crate) probing_interval_minutes: u64,
|
||||||
pub(crate) stream_timeout_threshold: u64,
|
pub(crate) stream_timeout_threshold: u64,
|
||||||
pub(crate) stream_timeout_window_seconds: u64,
|
pub(crate) stream_timeout_window_seconds: u64,
|
||||||
pub(crate) stream_timeout_cooldown_seconds: u64,
|
pub(crate) stream_timeout_cooldown_seconds: u64,
|
||||||
|
|||||||
@@ -4,17 +4,18 @@ mod tests;
|
|||||||
|
|
||||||
pub(crate) use runtime::{
|
pub(crate) use runtime::{
|
||||||
cancel_proxy_upgrade_rollout, clear_proxy_upgrade_rollout_conflicts,
|
cancel_proxy_upgrade_rollout, clear_proxy_upgrade_rollout_conflicts,
|
||||||
inspect_proxy_upgrade_rollout, perform_provider_checkin_once,
|
inspect_proxy_upgrade_rollout, perform_pool_quota_probe_once, perform_provider_checkin_once,
|
||||||
record_proxy_upgrade_traffic_success, restore_proxy_upgrade_rollout_skipped_nodes,
|
record_proxy_upgrade_traffic_success, restore_proxy_upgrade_rollout_skipped_nodes,
|
||||||
retry_proxy_upgrade_rollout_node, skip_proxy_upgrade_rollout_node, spawn_audit_cleanup_worker,
|
retry_proxy_upgrade_rollout_node, skip_proxy_upgrade_rollout_node, spawn_audit_cleanup_worker,
|
||||||
spawn_db_maintenance_worker, spawn_gemini_file_mapping_cleanup_worker,
|
spawn_db_maintenance_worker, spawn_gemini_file_mapping_cleanup_worker,
|
||||||
spawn_pending_cleanup_worker, spawn_pool_monitor_worker, spawn_provider_checkin_worker,
|
spawn_pending_cleanup_worker, spawn_pool_monitor_worker, spawn_pool_quota_probe_worker,
|
||||||
spawn_proxy_node_stale_cleanup_worker, spawn_proxy_upgrade_rollout_worker,
|
spawn_provider_checkin_worker, spawn_proxy_node_stale_cleanup_worker,
|
||||||
spawn_request_candidate_cleanup_worker, spawn_stats_aggregation_worker,
|
spawn_proxy_upgrade_rollout_worker, spawn_request_candidate_cleanup_worker,
|
||||||
spawn_stats_hourly_aggregation_worker, spawn_usage_cleanup_worker,
|
spawn_stats_aggregation_worker, spawn_stats_hourly_aggregation_worker,
|
||||||
spawn_wallet_daily_usage_aggregation_worker, start_proxy_upgrade_rollout,
|
spawn_usage_cleanup_worker, spawn_wallet_daily_usage_aggregation_worker,
|
||||||
ProviderCheckinRunSummary, ProxyUpgradeRolloutCancelSummary,
|
start_proxy_upgrade_rollout, PoolQuotaProbeRunSummary, ProviderCheckinRunSummary,
|
||||||
ProxyUpgradeRolloutConflictClearSummary, ProxyUpgradeRolloutNodeActionSummary,
|
ProxyUpgradeRolloutCancelSummary, ProxyUpgradeRolloutConflictClearSummary,
|
||||||
ProxyUpgradeRolloutProbeConfig, ProxyUpgradeRolloutSkippedRestoreSummary,
|
ProxyUpgradeRolloutNodeActionSummary, ProxyUpgradeRolloutProbeConfig,
|
||||||
ProxyUpgradeRolloutStatus, ProxyUpgradeRolloutTrackedNodeState,
|
ProxyUpgradeRolloutSkippedRestoreSummary, ProxyUpgradeRolloutStatus,
|
||||||
|
ProxyUpgradeRolloutTrackedNodeState,
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -27,6 +27,8 @@ mod config;
|
|||||||
mod db_maintenance;
|
mod db_maintenance;
|
||||||
#[path = "runtime/pending_cleanup.rs"]
|
#[path = "runtime/pending_cleanup.rs"]
|
||||||
mod pending_cleanup;
|
mod pending_cleanup;
|
||||||
|
#[path = "runtime/pool_quota_probe.rs"]
|
||||||
|
mod pool_quota_probe;
|
||||||
#[path = "runtime/provider_checkin.rs"]
|
#[path = "runtime/provider_checkin.rs"]
|
||||||
mod provider_checkin;
|
mod provider_checkin;
|
||||||
#[path = "runtime/proxy_node_staleness.rs"]
|
#[path = "runtime/proxy_node_staleness.rs"]
|
||||||
@@ -56,6 +58,11 @@ use audit_cleanup::*;
|
|||||||
use config::*;
|
use config::*;
|
||||||
use db_maintenance::*;
|
use db_maintenance::*;
|
||||||
use pending_cleanup::*;
|
use pending_cleanup::*;
|
||||||
|
pub(crate) use pool_quota_probe::{
|
||||||
|
perform_pool_quota_probe_once, perform_pool_quota_probe_once_with_config,
|
||||||
|
select_pool_quota_probe_key_ids, spawn_pool_quota_probe_worker, PoolQuotaProbeRunSummary,
|
||||||
|
PoolQuotaProbeWorkerConfig,
|
||||||
|
};
|
||||||
pub(crate) use provider_checkin::{perform_provider_checkin_once, ProviderCheckinRunSummary};
|
pub(crate) use provider_checkin::{perform_provider_checkin_once, ProviderCheckinRunSummary};
|
||||||
use proxy_node_staleness::*;
|
use proxy_node_staleness::*;
|
||||||
use proxy_upgrade_rollout::*;
|
use proxy_upgrade_rollout::*;
|
||||||
|
|||||||
@@ -0,0 +1,640 @@
|
|||||||
|
use std::collections::BTreeMap;
|
||||||
|
use std::time::{Duration, SystemTime, UNIX_EPOCH};
|
||||||
|
|
||||||
|
use aether_data::redis::{RedisKvRunner, RedisLockLease, RedisLockRunner};
|
||||||
|
use aether_data_contracts::repository::provider_catalog::{
|
||||||
|
StoredProviderCatalogEndpoint, StoredProviderCatalogKey, StoredProviderCatalogProvider,
|
||||||
|
};
|
||||||
|
use serde_json::Value;
|
||||||
|
use tracing::{debug, info, warn};
|
||||||
|
|
||||||
|
use crate::admin_api::{
|
||||||
|
admin_provider_pool_config, provider_oauth_runtime_endpoint_for_provider,
|
||||||
|
refresh_antigravity_provider_quota_locally, refresh_codex_provider_quota_locally,
|
||||||
|
refresh_kiro_provider_quota_locally, AdminAppState,
|
||||||
|
};
|
||||||
|
use crate::{AppState, GatewayError};
|
||||||
|
|
||||||
|
const POOL_QUOTA_PROBE_REDIS_PREFIX: &str = "ap:quota_probe:last";
|
||||||
|
const POOL_QUOTA_PROBE_DEFAULT_SCAN_INTERVAL_SECONDS: u64 = 60;
|
||||||
|
const POOL_QUOTA_PROBE_MIN_SCAN_INTERVAL_SECONDS: u64 = 15;
|
||||||
|
const POOL_QUOTA_PROBE_DEFAULT_MAX_KEYS_PER_PROVIDER: usize = 50;
|
||||||
|
const POOL_QUOTA_PROBE_PROVIDER_LOCK_TTL_MS: u64 = 30_000;
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub(crate) struct PoolQuotaProbeRunSummary {
|
||||||
|
pub(crate) providers_checked: usize,
|
||||||
|
pub(crate) providers_probed: usize,
|
||||||
|
pub(crate) providers_skipped: usize,
|
||||||
|
pub(crate) selected_keys: usize,
|
||||||
|
pub(crate) succeeded: usize,
|
||||||
|
pub(crate) failed: usize,
|
||||||
|
pub(crate) auto_removed: usize,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PoolQuotaProbeRunSummary {
|
||||||
|
const fn empty() -> Self {
|
||||||
|
Self {
|
||||||
|
providers_checked: 0,
|
||||||
|
providers_probed: 0,
|
||||||
|
providers_skipped: 0,
|
||||||
|
selected_keys: 0,
|
||||||
|
succeeded: 0,
|
||||||
|
failed: 0,
|
||||||
|
auto_removed: 0,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy)]
|
||||||
|
pub(crate) struct PoolQuotaProbeWorkerConfig {
|
||||||
|
pub(crate) scan_interval: Duration,
|
||||||
|
pub(crate) max_keys_per_provider: usize,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PoolQuotaProbeWorkerConfig {
|
||||||
|
fn from_env() -> Self {
|
||||||
|
let scan_interval_seconds = env_u64(
|
||||||
|
"POOL_QUOTA_PROBE_SCAN_INTERVAL_SECONDS",
|
||||||
|
POOL_QUOTA_PROBE_DEFAULT_SCAN_INTERVAL_SECONDS,
|
||||||
|
)
|
||||||
|
.max(POOL_QUOTA_PROBE_MIN_SCAN_INTERVAL_SECONDS);
|
||||||
|
let max_keys_per_provider = env_usize(
|
||||||
|
"POOL_QUOTA_PROBE_MAX_KEYS_PER_PROVIDER",
|
||||||
|
POOL_QUOTA_PROBE_DEFAULT_MAX_KEYS_PER_PROVIDER,
|
||||||
|
);
|
||||||
|
Self {
|
||||||
|
scan_interval: Duration::from_secs(scan_interval_seconds),
|
||||||
|
max_keys_per_provider,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn env_u64(name: &str, default_value: u64) -> u64 {
|
||||||
|
std::env::var(name)
|
||||||
|
.ok()
|
||||||
|
.and_then(|value| value.trim().parse::<u64>().ok())
|
||||||
|
.unwrap_or(default_value)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn env_usize(name: &str, default_value: usize) -> usize {
|
||||||
|
std::env::var(name)
|
||||||
|
.ok()
|
||||||
|
.and_then(|value| value.trim().parse::<usize>().ok())
|
||||||
|
.unwrap_or(default_value)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn now_unix_secs() -> u64 {
|
||||||
|
SystemTime::now()
|
||||||
|
.duration_since(UNIX_EPOCH)
|
||||||
|
.unwrap_or_default()
|
||||||
|
.as_secs()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn provider_supports_quota_probe(provider_type: &str) -> bool {
|
||||||
|
matches!(
|
||||||
|
provider_type.trim().to_ascii_lowercase().as_str(),
|
||||||
|
"codex" | "kiro" | "antigravity"
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn json_number(value: Option<&Value>) -> Option<f64> {
|
||||||
|
let value = value?;
|
||||||
|
if let Some(number) = value.as_f64() {
|
||||||
|
return Some(number);
|
||||||
|
}
|
||||||
|
value
|
||||||
|
.as_str()
|
||||||
|
.map(str::trim)
|
||||||
|
.filter(|value| !value.is_empty())
|
||||||
|
.and_then(|value| value.parse::<f64>().ok())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn extract_quota_updated_at(provider_type: &str, upstream_metadata: Option<&Value>) -> Option<u64> {
|
||||||
|
let metadata = upstream_metadata?.as_object()?;
|
||||||
|
let bucket_name = match provider_type.trim().to_ascii_lowercase().as_str() {
|
||||||
|
"codex" => "codex",
|
||||||
|
"kiro" => "kiro",
|
||||||
|
"antigravity" => "antigravity",
|
||||||
|
_ => return None,
|
||||||
|
};
|
||||||
|
let bucket = metadata.get(bucket_name)?.as_object()?;
|
||||||
|
let mut updated_at = json_number(bucket.get("updated_at"))?;
|
||||||
|
if updated_at <= 0.0 {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
if updated_at > 1_000_000_000_000.0 {
|
||||||
|
updated_at /= 1000.0;
|
||||||
|
}
|
||||||
|
Some(updated_at as u64)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn parse_probe_stamp(raw_value: Option<&str>) -> Option<u64> {
|
||||||
|
let parsed = raw_value
|
||||||
|
.map(str::trim)
|
||||||
|
.filter(|value| !value.is_empty())
|
||||||
|
.and_then(|value| value.parse::<f64>().ok())?;
|
||||||
|
if parsed <= 0.0 {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
Some(parsed as u64)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn select_pool_quota_probe_key_ids(
|
||||||
|
keys: &[StoredProviderCatalogKey],
|
||||||
|
provider_type: &str,
|
||||||
|
now_ts: u64,
|
||||||
|
interval_seconds: u64,
|
||||||
|
last_probe_timestamps: &BTreeMap<String, u64>,
|
||||||
|
limit: usize,
|
||||||
|
) -> Vec<String> {
|
||||||
|
let mut stale = Vec::<(u64, String)>::new();
|
||||||
|
for key in keys {
|
||||||
|
if key.id.trim().is_empty() {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
let quota_updated_ts =
|
||||||
|
extract_quota_updated_at(provider_type, key.upstream_metadata.as_ref());
|
||||||
|
let last_probe_ts = last_probe_timestamps.get(&key.id).copied();
|
||||||
|
let anchor_ts = quota_updated_ts
|
||||||
|
.unwrap_or(0)
|
||||||
|
.max(last_probe_ts.unwrap_or(0));
|
||||||
|
if anchor_ts == 0 || now_ts.saturating_sub(anchor_ts) >= interval_seconds {
|
||||||
|
stale.push((anchor_ts, key.id.clone()));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
stale.sort_by(|left, right| left.0.cmp(&right.0).then_with(|| left.1.cmp(&right.1)));
|
||||||
|
if limit > 0 && stale.len() > limit {
|
||||||
|
stale.truncate(limit);
|
||||||
|
}
|
||||||
|
stale.into_iter().map(|(_, key_id)| key_id).collect()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn probe_stamp_key(provider_id: &str, key_id: &str) -> String {
|
||||||
|
format!("{POOL_QUOTA_PROBE_REDIS_PREFIX}:{provider_id}:{key_id}")
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn load_probe_timestamps(
|
||||||
|
runner: Option<&RedisKvRunner>,
|
||||||
|
provider_id: &str,
|
||||||
|
key_ids: &[String],
|
||||||
|
) -> BTreeMap<String, u64> {
|
||||||
|
let Some(runner) = runner else {
|
||||||
|
return BTreeMap::new();
|
||||||
|
};
|
||||||
|
if key_ids.is_empty() {
|
||||||
|
return BTreeMap::new();
|
||||||
|
}
|
||||||
|
|
||||||
|
let Ok(mut connection) = runner.client().get_multiplexed_async_connection().await else {
|
||||||
|
debug!("gateway pool quota probe: failed to connect redis for stamp read");
|
||||||
|
return BTreeMap::new();
|
||||||
|
};
|
||||||
|
let redis_keys = key_ids
|
||||||
|
.iter()
|
||||||
|
.map(|key_id| runner.keyspace().key(&probe_stamp_key(provider_id, key_id)))
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
let Ok(values) = redis::cmd("MGET")
|
||||||
|
.arg(redis_keys)
|
||||||
|
.query_async::<Vec<Option<String>>>(&mut connection)
|
||||||
|
.await
|
||||||
|
else {
|
||||||
|
debug!("gateway pool quota probe: failed to read redis probe stamps");
|
||||||
|
return BTreeMap::new();
|
||||||
|
};
|
||||||
|
|
||||||
|
key_ids
|
||||||
|
.iter()
|
||||||
|
.zip(values)
|
||||||
|
.filter_map(|(key_id, raw)| {
|
||||||
|
parse_probe_stamp(raw.as_deref()).map(|ts| (key_id.clone(), ts))
|
||||||
|
})
|
||||||
|
.collect()
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn mark_probe_timestamps(
|
||||||
|
runner: Option<&RedisKvRunner>,
|
||||||
|
provider_id: &str,
|
||||||
|
key_ids: &[String],
|
||||||
|
now_ts: u64,
|
||||||
|
interval_seconds: u64,
|
||||||
|
) {
|
||||||
|
let Some(runner) = runner else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
if key_ids.is_empty() {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
let Ok(mut connection) = runner.client().get_multiplexed_async_connection().await else {
|
||||||
|
debug!("gateway pool quota probe: failed to connect redis for stamp write");
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
let ttl_seconds = interval_seconds.saturating_mul(2).max(120);
|
||||||
|
let value = now_ts.to_string();
|
||||||
|
let mut pipeline = redis::pipe();
|
||||||
|
for key_id in key_ids {
|
||||||
|
pipeline
|
||||||
|
.cmd("SETEX")
|
||||||
|
.arg(runner.keyspace().key(&probe_stamp_key(provider_id, key_id)))
|
||||||
|
.arg(ttl_seconds)
|
||||||
|
.arg(&value)
|
||||||
|
.ignore();
|
||||||
|
}
|
||||||
|
if pipeline.query_async::<()>(&mut connection).await.is_err() {
|
||||||
|
debug!("gateway pool quota probe: failed to write redis probe stamps");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn acquire_provider_probe_lock(
|
||||||
|
runner: Option<&RedisLockRunner>,
|
||||||
|
provider_id: &str,
|
||||||
|
) -> Option<RedisLockLease> {
|
||||||
|
let Some(runner) = runner else {
|
||||||
|
return None;
|
||||||
|
};
|
||||||
|
let lock_key = runner
|
||||||
|
.keyspace()
|
||||||
|
.lock_key(&format!("pool_quota_probe:{provider_id}"));
|
||||||
|
let owner = format!("aether-gateway-pool-probe-{}", std::process::id());
|
||||||
|
match runner
|
||||||
|
.try_acquire(
|
||||||
|
&lock_key,
|
||||||
|
&owner,
|
||||||
|
Some(POOL_QUOTA_PROBE_PROVIDER_LOCK_TTL_MS),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
Ok(lease) => lease,
|
||||||
|
Err(err) => {
|
||||||
|
debug!(
|
||||||
|
provider_id,
|
||||||
|
error = %err,
|
||||||
|
"gateway pool quota probe: failed to acquire redis provider lock"
|
||||||
|
);
|
||||||
|
None
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn release_provider_probe_lock(
|
||||||
|
runner: Option<&RedisLockRunner>,
|
||||||
|
lease: Option<RedisLockLease>,
|
||||||
|
) {
|
||||||
|
let (Some(runner), Some(lease)) = (runner, lease) else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
if let Err(err) = runner.release(&lease).await {
|
||||||
|
debug!(
|
||||||
|
error = %err,
|
||||||
|
"gateway pool quota probe: failed to release redis provider lock"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn select_keys_for_provider(
|
||||||
|
state: &AppState,
|
||||||
|
redis_kv: Option<&RedisKvRunner>,
|
||||||
|
redis_lock: Option<&RedisLockRunner>,
|
||||||
|
provider: &StoredProviderCatalogProvider,
|
||||||
|
provider_type: &str,
|
||||||
|
interval_seconds: u64,
|
||||||
|
max_keys_per_provider: usize,
|
||||||
|
now_ts: u64,
|
||||||
|
) -> Result<Vec<StoredProviderCatalogKey>, GatewayError> {
|
||||||
|
let lease = acquire_provider_probe_lock(redis_lock, &provider.id).await;
|
||||||
|
if redis_lock.is_some() && lease.is_none() {
|
||||||
|
return Ok(Vec::new());
|
||||||
|
}
|
||||||
|
|
||||||
|
let result = async {
|
||||||
|
let keys = state
|
||||||
|
.list_provider_catalog_keys_by_provider_ids(std::slice::from_ref(&provider.id))
|
||||||
|
.await?
|
||||||
|
.into_iter()
|
||||||
|
.filter(|key| key.is_active)
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
if keys.is_empty() {
|
||||||
|
return Ok(Vec::new());
|
||||||
|
}
|
||||||
|
|
||||||
|
let key_ids = keys.iter().map(|key| key.id.clone()).collect::<Vec<_>>();
|
||||||
|
let probe_stamps = load_probe_timestamps(redis_kv, &provider.id, &key_ids).await;
|
||||||
|
let selected_ids = select_pool_quota_probe_key_ids(
|
||||||
|
&keys,
|
||||||
|
provider_type,
|
||||||
|
now_ts,
|
||||||
|
interval_seconds,
|
||||||
|
&probe_stamps,
|
||||||
|
max_keys_per_provider,
|
||||||
|
);
|
||||||
|
if selected_ids.is_empty() {
|
||||||
|
return Ok(Vec::new());
|
||||||
|
}
|
||||||
|
|
||||||
|
mark_probe_timestamps(
|
||||||
|
redis_kv,
|
||||||
|
&provider.id,
|
||||||
|
&selected_ids,
|
||||||
|
now_ts,
|
||||||
|
interval_seconds,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
|
||||||
|
let mut keys_by_id = keys
|
||||||
|
.into_iter()
|
||||||
|
.map(|key| (key.id.clone(), key))
|
||||||
|
.collect::<BTreeMap<_, _>>();
|
||||||
|
Ok(selected_ids
|
||||||
|
.into_iter()
|
||||||
|
.filter_map(|key_id| keys_by_id.remove(&key_id))
|
||||||
|
.collect::<Vec<_>>())
|
||||||
|
}
|
||||||
|
.await;
|
||||||
|
|
||||||
|
release_provider_probe_lock(redis_lock, lease).await;
|
||||||
|
result
|
||||||
|
}
|
||||||
|
|
||||||
|
fn endpoint_for_probe(
|
||||||
|
provider_type: &str,
|
||||||
|
endpoints: &[StoredProviderCatalogEndpoint],
|
||||||
|
) -> Option<StoredProviderCatalogEndpoint> {
|
||||||
|
provider_oauth_runtime_endpoint_for_provider(provider_type, endpoints)
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn refresh_provider_probe_keys(
|
||||||
|
admin_state: &AdminAppState<'_>,
|
||||||
|
provider: &StoredProviderCatalogProvider,
|
||||||
|
endpoint: &StoredProviderCatalogEndpoint,
|
||||||
|
provider_type: &str,
|
||||||
|
keys: Vec<StoredProviderCatalogKey>,
|
||||||
|
) -> Result<Option<Value>, GatewayError> {
|
||||||
|
match provider_type {
|
||||||
|
"codex" => {
|
||||||
|
refresh_codex_provider_quota_locally(admin_state, provider, endpoint, keys, None).await
|
||||||
|
}
|
||||||
|
"kiro" => {
|
||||||
|
refresh_kiro_provider_quota_locally(admin_state, provider, endpoint, keys, None).await
|
||||||
|
}
|
||||||
|
"antigravity" => {
|
||||||
|
refresh_antigravity_provider_quota_locally(admin_state, provider, endpoint, keys, None)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
_ => Ok(None),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn update_summary_from_payload(
|
||||||
|
summary: &mut PoolQuotaProbeRunSummary,
|
||||||
|
selected_count: usize,
|
||||||
|
payload: Option<&Value>,
|
||||||
|
) {
|
||||||
|
let Some(payload) = payload else {
|
||||||
|
summary.failed += selected_count;
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
summary.succeeded += payload.get("success").and_then(Value::as_u64).unwrap_or(0) as usize;
|
||||||
|
summary.failed += payload.get("failed").and_then(Value::as_u64).unwrap_or(0) as usize;
|
||||||
|
summary.auto_removed += payload
|
||||||
|
.get("auto_removed")
|
||||||
|
.and_then(Value::as_u64)
|
||||||
|
.unwrap_or(0) as usize;
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn perform_pool_quota_probe_once_with_config(
|
||||||
|
state: &AppState,
|
||||||
|
config: PoolQuotaProbeWorkerConfig,
|
||||||
|
) -> Result<PoolQuotaProbeRunSummary, GatewayError> {
|
||||||
|
if !state.has_provider_catalog_data_reader() || !state.has_provider_catalog_data_writer() {
|
||||||
|
return Ok(PoolQuotaProbeRunSummary::empty());
|
||||||
|
}
|
||||||
|
|
||||||
|
let providers = state
|
||||||
|
.list_provider_catalog_providers(true)
|
||||||
|
.await?
|
||||||
|
.into_iter()
|
||||||
|
.filter_map(|provider| {
|
||||||
|
let provider_type = provider.provider_type.trim().to_ascii_lowercase();
|
||||||
|
if provider_supports_quota_probe(&provider_type) {
|
||||||
|
Some((provider, provider_type))
|
||||||
|
} else {
|
||||||
|
None
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.filter_map(|(provider, provider_type)| {
|
||||||
|
let pool_config = admin_provider_pool_config(&provider)?;
|
||||||
|
if pool_config.probing_enabled {
|
||||||
|
Some((
|
||||||
|
provider,
|
||||||
|
provider_type,
|
||||||
|
pool_config.probing_interval_minutes,
|
||||||
|
))
|
||||||
|
} else {
|
||||||
|
None
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
|
||||||
|
if providers.is_empty() {
|
||||||
|
return Ok(PoolQuotaProbeRunSummary::empty());
|
||||||
|
}
|
||||||
|
|
||||||
|
let provider_ids = providers
|
||||||
|
.iter()
|
||||||
|
.map(|(provider, _, _)| provider.id.clone())
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
let mut endpoints_by_provider = BTreeMap::<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 redis_kv = state.redis_kv_runner();
|
||||||
|
let redis_lock = state.data.oauth_refresh_lock_runner();
|
||||||
|
let admin_state = AdminAppState::new(state);
|
||||||
|
let now_ts = now_unix_secs();
|
||||||
|
let mut summary = PoolQuotaProbeRunSummary {
|
||||||
|
providers_checked: providers.len(),
|
||||||
|
..PoolQuotaProbeRunSummary::empty()
|
||||||
|
};
|
||||||
|
|
||||||
|
for (provider, provider_type, interval_minutes) in providers {
|
||||||
|
let endpoints = endpoints_by_provider
|
||||||
|
.remove(&provider.id)
|
||||||
|
.unwrap_or_default();
|
||||||
|
let Some(endpoint) = endpoint_for_probe(&provider_type, &endpoints) else {
|
||||||
|
summary.providers_skipped += 1;
|
||||||
|
debug!(
|
||||||
|
provider_id = %provider.id,
|
||||||
|
provider_type,
|
||||||
|
"gateway pool quota probe skipped provider without active quota endpoint"
|
||||||
|
);
|
||||||
|
continue;
|
||||||
|
};
|
||||||
|
|
||||||
|
let interval_seconds = interval_minutes.clamp(1, 1440).saturating_mul(60);
|
||||||
|
let keys = select_keys_for_provider(
|
||||||
|
state,
|
||||||
|
redis_kv.as_ref(),
|
||||||
|
redis_lock.as_ref(),
|
||||||
|
&provider,
|
||||||
|
&provider_type,
|
||||||
|
interval_seconds,
|
||||||
|
config.max_keys_per_provider,
|
||||||
|
now_ts,
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
if keys.is_empty() {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
let selected_count = keys.len();
|
||||||
|
summary.providers_probed += 1;
|
||||||
|
summary.selected_keys += selected_count;
|
||||||
|
|
||||||
|
let provider_short_id = provider.id.chars().take(8).collect::<String>();
|
||||||
|
match refresh_provider_probe_keys(&admin_state, &provider, &endpoint, &provider_type, keys)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
Ok(payload) => {
|
||||||
|
update_summary_from_payload(&mut summary, selected_count, payload.as_ref());
|
||||||
|
let probe_success = payload
|
||||||
|
.as_ref()
|
||||||
|
.and_then(|value| value.get("success"))
|
||||||
|
.and_then(Value::as_u64)
|
||||||
|
.unwrap_or(0);
|
||||||
|
let probe_failed = payload
|
||||||
|
.as_ref()
|
||||||
|
.and_then(|value| value.get("failed"))
|
||||||
|
.and_then(Value::as_u64)
|
||||||
|
.unwrap_or(0);
|
||||||
|
info!(
|
||||||
|
provider_id = %provider_short_id,
|
||||||
|
provider_type,
|
||||||
|
selected = selected_count,
|
||||||
|
success = probe_success,
|
||||||
|
failed = probe_failed,
|
||||||
|
"gateway pool quota probe completed"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
Err(err) => {
|
||||||
|
summary.failed += selected_count;
|
||||||
|
warn!(
|
||||||
|
provider_id = %provider_short_id,
|
||||||
|
provider_type,
|
||||||
|
selected = selected_count,
|
||||||
|
error = ?err,
|
||||||
|
"gateway pool quota probe failed"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(summary)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn perform_pool_quota_probe_once(
|
||||||
|
state: &AppState,
|
||||||
|
) -> Result<PoolQuotaProbeRunSummary, GatewayError> {
|
||||||
|
perform_pool_quota_probe_once_with_config(state, PoolQuotaProbeWorkerConfig::from_env()).await
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn spawn_pool_quota_probe_worker(
|
||||||
|
state: AppState,
|
||||||
|
) -> Option<tokio::task::JoinHandle<()>> {
|
||||||
|
if !state.has_provider_catalog_data_reader() || !state.has_provider_catalog_data_writer() {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
|
||||||
|
let config = PoolQuotaProbeWorkerConfig::from_env();
|
||||||
|
Some(tokio::spawn(async move {
|
||||||
|
let mut interval = tokio::time::interval(config.scan_interval);
|
||||||
|
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
|
||||||
|
interval.tick().await;
|
||||||
|
loop {
|
||||||
|
interval.tick().await;
|
||||||
|
if let Err(err) = perform_pool_quota_probe_once_with_config(&state, config).await {
|
||||||
|
warn!(
|
||||||
|
error = ?err,
|
||||||
|
"gateway pool quota probe worker tick failed"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
use serde_json::json;
|
||||||
|
|
||||||
|
fn key(
|
||||||
|
id: &str,
|
||||||
|
provider_id: &str,
|
||||||
|
upstream_metadata: Option<Value>,
|
||||||
|
) -> StoredProviderCatalogKey {
|
||||||
|
let mut key = StoredProviderCatalogKey::new(
|
||||||
|
id.to_string(),
|
||||||
|
provider_id.to_string(),
|
||||||
|
id.to_string(),
|
||||||
|
"oauth".to_string(),
|
||||||
|
None,
|
||||||
|
true,
|
||||||
|
)
|
||||||
|
.expect("key should build");
|
||||||
|
key.upstream_metadata = upstream_metadata;
|
||||||
|
key
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn selects_stale_probe_keys_by_oldest_anchor() {
|
||||||
|
let keys = vec![
|
||||||
|
key(
|
||||||
|
"fresh",
|
||||||
|
"provider-1",
|
||||||
|
Some(json!({ "codex": { "updated_at": 1_990 } })),
|
||||||
|
),
|
||||||
|
key(
|
||||||
|
"old",
|
||||||
|
"provider-1",
|
||||||
|
Some(json!({ "codex": { "updated_at": 1_000 } })),
|
||||||
|
),
|
||||||
|
key("never", "provider-1", None),
|
||||||
|
key(
|
||||||
|
"stamped",
|
||||||
|
"provider-1",
|
||||||
|
Some(json!({ "codex": { "updated_at": 900 } })),
|
||||||
|
),
|
||||||
|
];
|
||||||
|
let stamps = BTreeMap::from([("stamped".to_string(), 1_950)]);
|
||||||
|
|
||||||
|
let selected = select_pool_quota_probe_key_ids(&keys, "codex", 2_000, 600, &stamps, 2);
|
||||||
|
|
||||||
|
assert_eq!(selected, vec!["never".to_string(), "old".to_string()]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn parses_quota_updated_at_seconds_and_milliseconds() {
|
||||||
|
assert_eq!(
|
||||||
|
extract_quota_updated_at(
|
||||||
|
"codex",
|
||||||
|
Some(&json!({ "codex": { "updated_at": 1_700_000_000 } }))
|
||||||
|
),
|
||||||
|
Some(1_700_000_000)
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
extract_quota_updated_at(
|
||||||
|
"kiro",
|
||||||
|
Some(&json!({ "kiro": { "updated_at": 1_700_000_000_000_u64 } }))
|
||||||
|
),
|
||||||
|
Some(1_700_000_000)
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -20,18 +20,19 @@ use super::{
|
|||||||
pending_cleanup_batch_size, pending_cleanup_timeout_minutes, plan_pending_cleanup_batch,
|
pending_cleanup_batch_size, pending_cleanup_timeout_minutes, plan_pending_cleanup_batch,
|
||||||
provider_checkin_schedule, record_proxy_upgrade_traffic_success, run_db_maintenance_with,
|
provider_checkin_schedule, record_proxy_upgrade_traffic_success, run_db_maintenance_with,
|
||||||
run_proxy_upgrade_rollout_once, spawn_audit_cleanup_worker, spawn_db_maintenance_worker,
|
run_proxy_upgrade_rollout_once, spawn_audit_cleanup_worker, spawn_db_maintenance_worker,
|
||||||
spawn_pending_cleanup_worker, spawn_pool_monitor_worker, spawn_provider_checkin_worker,
|
spawn_pending_cleanup_worker, spawn_pool_monitor_worker, spawn_pool_quota_probe_worker,
|
||||||
spawn_proxy_node_stale_cleanup_worker, spawn_proxy_upgrade_rollout_worker,
|
spawn_provider_checkin_worker, spawn_proxy_node_stale_cleanup_worker,
|
||||||
spawn_stats_aggregation_worker, spawn_stats_hourly_aggregation_worker,
|
spawn_proxy_upgrade_rollout_worker, spawn_stats_aggregation_worker,
|
||||||
spawn_usage_cleanup_worker, spawn_wallet_daily_usage_aggregation_worker,
|
spawn_stats_hourly_aggregation_worker, spawn_usage_cleanup_worker,
|
||||||
start_proxy_upgrade_rollout, stats_aggregation_target_day,
|
spawn_wallet_daily_usage_aggregation_worker, start_proxy_upgrade_rollout,
|
||||||
stats_hourly_aggregation_target_hour, summarize_postgres_pool, usage_cleanup_settings,
|
stats_aggregation_target_day, stats_hourly_aggregation_target_hour, summarize_postgres_pool,
|
||||||
usage_cleanup_window, wallet_daily_usage_aggregation_target, AppState, DbMaintenanceRunSummary,
|
usage_cleanup_settings, usage_cleanup_window, wallet_daily_usage_aggregation_target, AppState,
|
||||||
FailedPendingUsageRow, GatewayDataState, ProxyUpgradeRolloutProbeConfig, StalePendingUsageRow,
|
DbMaintenanceRunSummary, FailedPendingUsageRow, GatewayDataState,
|
||||||
UsageCleanupSettings, DELETE_STALE_WALLET_DAILY_USAGE_LEDGERS_SQL,
|
ProxyUpgradeRolloutProbeConfig, StalePendingUsageRow, UsageCleanupSettings,
|
||||||
SELECT_STALE_PENDING_USAGE_BATCH_SQL, UPDATE_FAILED_VOID_STALE_USAGE_SQL,
|
DELETE_STALE_WALLET_DAILY_USAGE_LEDGERS_SQL, SELECT_STALE_PENDING_USAGE_BATCH_SQL,
|
||||||
UPSERT_WALLET_DAILY_USAGE_LEDGER_SQL, USAGE_CLEANUP_HOUR, USAGE_CLEANUP_MINUTE,
|
UPDATE_FAILED_VOID_STALE_USAGE_SQL, UPSERT_WALLET_DAILY_USAGE_LEDGER_SQL, USAGE_CLEANUP_HOUR,
|
||||||
WALLET_DAILY_USAGE_AGGREGATION_HOUR, WALLET_DAILY_USAGE_AGGREGATION_MINUTE,
|
USAGE_CLEANUP_MINUTE, WALLET_DAILY_USAGE_AGGREGATION_HOUR,
|
||||||
|
WALLET_DAILY_USAGE_AGGREGATION_MINUTE,
|
||||||
};
|
};
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
@@ -79,6 +80,15 @@ async fn spawn_pool_monitor_worker_skips_when_postgres_unavailable() {
|
|||||||
assert!(spawn_pool_monitor_worker(Arc::new(GatewayDataState::disabled())).is_none());
|
assert!(spawn_pool_monitor_worker(Arc::new(GatewayDataState::disabled())).is_none());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn spawn_pool_quota_probe_worker_skips_when_provider_catalog_unavailable() {
|
||||||
|
let state = AppState::new()
|
||||||
|
.expect("gateway state should build")
|
||||||
|
.with_data_state_for_tests(GatewayDataState::disabled());
|
||||||
|
|
||||||
|
assert!(spawn_pool_quota_probe_worker(state).is_none());
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn wallet_daily_usage_queries_use_settlement_snapshots_for_wallet_identity() {
|
fn wallet_daily_usage_queries_use_settlement_snapshots_for_wallet_identity() {
|
||||||
assert!(UPSERT_WALLET_DAILY_USAGE_LEDGER_SQL.contains("JOIN usage_settlement_snapshots"));
|
assert!(UPSERT_WALLET_DAILY_USAGE_LEDGER_SQL.contains("JOIN usage_settlement_snapshots"));
|
||||||
|
|||||||
@@ -40,6 +40,7 @@ use crate::maintenance::spawn_db_maintenance_worker;
|
|||||||
use crate::maintenance::spawn_gemini_file_mapping_cleanup_worker;
|
use crate::maintenance::spawn_gemini_file_mapping_cleanup_worker;
|
||||||
use crate::maintenance::spawn_pending_cleanup_worker;
|
use crate::maintenance::spawn_pending_cleanup_worker;
|
||||||
use crate::maintenance::spawn_pool_monitor_worker;
|
use crate::maintenance::spawn_pool_monitor_worker;
|
||||||
|
use crate::maintenance::spawn_pool_quota_probe_worker;
|
||||||
use crate::maintenance::spawn_provider_checkin_worker;
|
use crate::maintenance::spawn_provider_checkin_worker;
|
||||||
use crate::maintenance::spawn_proxy_node_stale_cleanup_worker;
|
use crate::maintenance::spawn_proxy_node_stale_cleanup_worker;
|
||||||
use crate::maintenance::spawn_proxy_upgrade_rollout_worker;
|
use crate::maintenance::spawn_proxy_upgrade_rollout_worker;
|
||||||
@@ -820,6 +821,9 @@ impl AppState {
|
|||||||
if let Some(handle) = spawn_pool_monitor_worker(self.data.clone()) {
|
if let Some(handle) = spawn_pool_monitor_worker(self.data.clone()) {
|
||||||
tasks.push(handle);
|
tasks.push(handle);
|
||||||
}
|
}
|
||||||
|
if let Some(handle) = spawn_pool_quota_probe_worker(self.clone()) {
|
||||||
|
tasks.push(handle);
|
||||||
|
}
|
||||||
if let Some(handle) = spawn_stats_hourly_aggregation_worker(self.data.clone()) {
|
if let Some(handle) = spawn_stats_hourly_aggregation_worker(self.data.clone()) {
|
||||||
tasks.push(handle);
|
tasks.push(handle);
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user