chore(proxy): guard metrics retention cleanup

This commit is contained in:
fawney19
2026-05-08 22:57:00 +08:00
parent a703acd1fe
commit e21fb58479
16 changed files with 562 additions and 75 deletions

View File

@@ -1073,11 +1073,16 @@ impl GatewayDataState {
&self,
retain_1m_from_unix_secs: u64,
retain_1h_from_unix_secs: u64,
delete_limit: usize,
) -> Result<super::ProxyNodeMetricsCleanupSummary, DataLayerError> {
match &self.proxy_node_writer {
Some(repository) => {
repository
.cleanup_proxy_node_metrics(retain_1m_from_unix_secs, retain_1h_from_unix_secs)
.cleanup_proxy_node_metrics(
retain_1m_from_unix_secs,
retain_1h_from_unix_secs,
delete_limit,
)
.await
}
None => Ok(super::ProxyNodeMetricsCleanupSummary::default()),

View File

@@ -533,6 +533,8 @@ async fn build_admin_system_cleanup_payload(
let cleaned = json!({
"audit_logs": summary.audit_logs_deleted,
"request_candidates": summary.request_candidates_deleted,
"proxy_node_metrics_1m": summary.proxy_node_metrics.deleted_1m_rows,
"proxy_node_metrics_1h": summary.proxy_node_metrics.deleted_1h_rows,
"pending_failed": summary.pending_failed,
"pending_recovered": summary.pending_recovered,
"usage_body_externalized": summary.usage.body_externalized,
@@ -545,6 +547,8 @@ async fn build_admin_system_cleanup_payload(
let total = summary
.audit_logs_deleted
.saturating_add(summary.request_candidates_deleted)
.saturating_add(summary.proxy_node_metrics.deleted_1m_rows)
.saturating_add(summary.proxy_node_metrics.deleted_1h_rows)
.saturating_add(summary.pending_failed)
.saturating_add(summary.pending_recovered)
.saturating_add(summary.usage.body_externalized)

View File

@@ -132,6 +132,8 @@ struct UsageCleanupSettings {
pub(crate) struct AdminSystemCleanupSummary {
pub(crate) audit_logs_deleted: usize,
pub(crate) request_candidates_deleted: usize,
pub(crate) proxy_node_metrics:
aether_data::repository::proxy_nodes::ProxyNodeMetricsCleanupSummary,
pub(crate) pending_failed: usize,
pub(crate) pending_recovered: usize,
pub(crate) usage: UsageCleanupSummary,
@@ -149,12 +151,14 @@ pub(crate) async fn run_admin_system_cleanup_once(
) -> Result<AdminSystemCleanupSummary, aether_data::DataLayerError> {
let audit_logs_deleted = cleanup_audit_logs_once(data).await?;
let request_candidates_deleted = cleanup_request_candidates_once(data).await?;
let proxy_node_metrics = cleanup_proxy_node_metrics_once(data).await?;
let pending = cleanup_stale_pending_requests_once(data).await?;
let usage = perform_usage_cleanup_once(data).await?;
Ok(AdminSystemCleanupSummary {
audit_logs_deleted,
request_candidates_deleted,
proxy_node_metrics,
pending_failed: pending.failed,
pending_recovered: pending.recovered,
usage,

View File

@@ -3,18 +3,107 @@ use aether_data_contracts::DataLayerError;
use crate::data::GatewayDataState;
use super::now_unix_secs;
use super::{now_unix_secs, system_config_bool, system_config_u64, system_config_usize};
const PROXY_NODE_METRICS_1M_RETENTION_SECS: u64 = 30 * 24 * 60 * 60;
const PROXY_NODE_METRICS_1H_RETENTION_SECS: u64 = 180 * 24 * 60 * 60;
const SECS_PER_DAY: u64 = 24 * 60 * 60;
const PROXY_NODE_METRICS_1M_RETENTION_DAYS_DEFAULT: u64 = 30;
const PROXY_NODE_METRICS_1H_RETENTION_DAYS_DEFAULT: u64 = 180;
const PROXY_NODE_METRICS_RETENTION_DAYS_MIN: u64 = 1;
const PROXY_NODE_METRICS_1M_RETENTION_DAYS_MAX: u64 = 365;
const PROXY_NODE_METRICS_1H_RETENTION_DAYS_MAX: u64 = 1_095;
const PROXY_NODE_METRICS_CLEANUP_BATCH_SIZE_DEFAULT: usize = 5_000;
const PROXY_NODE_METRICS_CLEANUP_BATCH_SIZE_MAX: usize = 50_000;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) struct ProxyNodeMetricsCleanupSettings {
pub retain_1m_days: u64,
pub retain_1h_days: u64,
pub batch_size: usize,
}
pub(super) async fn proxy_node_metrics_cleanup_settings(
data: &GatewayDataState,
) -> Result<ProxyNodeMetricsCleanupSettings, DataLayerError> {
let retain_1m_days = system_config_u64(
data,
"proxy_node_metrics_1m_retention_days",
PROXY_NODE_METRICS_1M_RETENTION_DAYS_DEFAULT,
)
.await?
.clamp(
PROXY_NODE_METRICS_RETENTION_DAYS_MIN,
PROXY_NODE_METRICS_1M_RETENTION_DAYS_MAX,
);
let retain_1h_days = system_config_u64(
data,
"proxy_node_metrics_1h_retention_days",
PROXY_NODE_METRICS_1H_RETENTION_DAYS_DEFAULT,
)
.await?
.clamp(retain_1m_days, PROXY_NODE_METRICS_1H_RETENTION_DAYS_MAX);
let cleanup_batch_size =
system_config_usize(data, "proxy_node_metrics_cleanup_batch_size", 0).await?;
let fallback_batch_size = system_config_usize(
data,
"cleanup_batch_size",
PROXY_NODE_METRICS_CLEANUP_BATCH_SIZE_DEFAULT,
)
.await?;
let batch_size = (if cleanup_batch_size > 0 {
cleanup_batch_size
} else {
fallback_batch_size
})
.clamp(1, PROXY_NODE_METRICS_CLEANUP_BATCH_SIZE_MAX);
Ok(ProxyNodeMetricsCleanupSettings {
retain_1m_days,
retain_1h_days,
batch_size,
})
}
pub(super) async fn cleanup_proxy_node_metrics_once(
data: &GatewayDataState,
) -> Result<ProxyNodeMetricsCleanupSummary, DataLayerError> {
let now = now_unix_secs();
data.cleanup_proxy_node_metrics(
now.saturating_sub(PROXY_NODE_METRICS_1M_RETENTION_SECS),
now.saturating_sub(PROXY_NODE_METRICS_1H_RETENTION_SECS),
)
.await
cleanup_proxy_node_metrics_at(data, now_unix_secs()).await
}
pub(super) async fn cleanup_proxy_node_metrics_at(
data: &GatewayDataState,
now_unix_secs: u64,
) -> Result<ProxyNodeMetricsCleanupSummary, DataLayerError> {
if !system_config_bool(data, "enable_auto_cleanup", true).await? {
return Ok(ProxyNodeMetricsCleanupSummary::default());
}
let settings = proxy_node_metrics_cleanup_settings(data).await?;
let retain_1m_from_unix_secs =
now_unix_secs.saturating_sub(settings.retain_1m_days.saturating_mul(SECS_PER_DAY));
let retain_1h_from_unix_secs =
now_unix_secs.saturating_sub(settings.retain_1h_days.saturating_mul(SECS_PER_DAY));
let mut summary = ProxyNodeMetricsCleanupSummary::default();
loop {
let deleted = data
.cleanup_proxy_node_metrics(
retain_1m_from_unix_secs,
retain_1h_from_unix_secs,
settings.batch_size,
)
.await?;
summary.deleted_1m_rows = summary
.deleted_1m_rows
.saturating_add(deleted.deleted_1m_rows);
summary.deleted_1h_rows = summary
.deleted_1h_rows
.saturating_add(deleted.deleted_1h_rows);
if deleted.deleted_1m_rows < settings.batch_size
&& deleted.deleted_1h_rows < settings.batch_size
{
break;
}
}
Ok(summary)
}

View File

@@ -3,8 +3,8 @@ use std::sync::{Arc, Mutex};
use std::time::Duration;
use aether_data::repository::proxy_nodes::{
InMemoryProxyNodeRepository, ProxyNodeHeartbeatMutation, ProxyNodeReadRepository,
ProxyNodeWriteRepository, StoredProxyNode,
bucket_start_unix_secs, InMemoryProxyNodeRepository, ProxyNodeHeartbeatMutation,
ProxyNodeMetricsStep, ProxyNodeReadRepository, ProxyNodeWriteRepository, StoredProxyNode,
};
use aether_runtime::bounded_queue;
use axum::extract::ws::Message;
@@ -14,23 +14,25 @@ use serde_json::json;
use tokio::sync::watch;
use super::{
advance_proxy_upgrade_rollout_once, cleanup_audit_logs_with, cleanup_stale_proxy_nodes_once,
inspect_proxy_upgrade_rollout, next_daily_run_after, next_db_maintenance_run_after,
next_stats_aggregation_run_after, next_stats_hourly_aggregation_run_after,
pending_cleanup_batch_size, pending_cleanup_timeout_minutes, plan_pending_cleanup_batch,
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,
spawn_oauth_token_refresh_worker, spawn_pending_cleanup_worker, spawn_pool_monitor_worker,
spawn_pool_quota_probe_worker, spawn_provider_checkin_worker,
advance_proxy_upgrade_rollout_once, cleanup_audit_logs_with, cleanup_proxy_node_metrics_at,
cleanup_proxy_node_metrics_once, cleanup_stale_proxy_nodes_once, inspect_proxy_upgrade_rollout,
next_daily_run_after, next_db_maintenance_run_after, next_stats_aggregation_run_after,
next_stats_hourly_aggregation_run_after, pending_cleanup_batch_size,
pending_cleanup_timeout_minutes, plan_pending_cleanup_batch, provider_checkin_schedule,
proxy_node_metrics_cleanup_settings, record_proxy_upgrade_traffic_success,
run_db_maintenance_with, run_proxy_upgrade_rollout_once, spawn_audit_cleanup_worker,
spawn_db_maintenance_worker, spawn_oauth_token_refresh_worker, spawn_pending_cleanup_worker,
spawn_pool_monitor_worker, spawn_pool_quota_probe_worker, spawn_provider_checkin_worker,
spawn_proxy_node_stale_cleanup_worker, spawn_proxy_upgrade_rollout_worker,
spawn_stats_aggregation_worker, spawn_stats_hourly_aggregation_worker,
spawn_usage_cleanup_worker, spawn_wallet_daily_usage_aggregation_worker,
start_proxy_upgrade_rollout, stats_aggregation_target_day,
stats_hourly_aggregation_target_hour, summarize_database_pool, usage_cleanup_settings,
usage_cleanup_window, wallet_daily_usage_aggregation_target, AppState, DbMaintenanceRunSummary,
FailedPendingUsageRow, GatewayDataState, ProxyUpgradeRolloutProbeConfig, StalePendingUsageRow,
UsageCleanupSettings, USAGE_CLEANUP_HOUR, USAGE_CLEANUP_MINUTE,
WALLET_DAILY_USAGE_AGGREGATION_HOUR, WALLET_DAILY_USAGE_AGGREGATION_MINUTE,
FailedPendingUsageRow, GatewayDataState, ProxyNodeMetricsCleanupSettings,
ProxyUpgradeRolloutProbeConfig, StalePendingUsageRow, UsageCleanupSettings, USAGE_CLEANUP_HOUR,
USAGE_CLEANUP_MINUTE, WALLET_DAILY_USAGE_AGGREGATION_HOUR,
WALLET_DAILY_USAGE_AGGREGATION_MINUTE,
};
#[tokio::test]
@@ -734,6 +736,176 @@ async fn usage_cleanup_settings_resolve_batch_and_delete_toggle() {
);
}
#[tokio::test]
async fn proxy_node_metrics_cleanup_settings_use_dedicated_retention_and_batch_limits() {
let data = GatewayDataState::disabled().with_system_config_values_for_tests([
("cleanup_batch_size".to_string(), json!(250)),
("proxy_node_metrics_1m_retention_days".to_string(), json!(0)),
("proxy_node_metrics_1h_retention_days".to_string(), json!(7)),
(
"proxy_node_metrics_cleanup_batch_size".to_string(),
json!(100_000),
),
]);
let settings = proxy_node_metrics_cleanup_settings(&data)
.await
.expect("proxy metrics cleanup settings should resolve");
assert_eq!(
settings,
ProxyNodeMetricsCleanupSettings {
retain_1m_days: 1,
retain_1h_days: 7,
batch_size: 50_000,
}
);
}
#[tokio::test]
async fn proxy_node_metrics_cleanup_settings_fallback_to_global_batch_size() {
let data = GatewayDataState::disabled().with_system_config_values_for_tests([
("cleanup_batch_size".to_string(), json!(250)),
(
"proxy_node_metrics_cleanup_batch_size".to_string(),
json!(0),
),
]);
let settings = proxy_node_metrics_cleanup_settings(&data)
.await
.expect("proxy metrics cleanup settings should resolve");
assert_eq!(
settings,
ProxyNodeMetricsCleanupSettings {
retain_1m_days: 30,
retain_1h_days: 180,
batch_size: 250,
}
);
}
#[tokio::test]
async fn proxy_node_metrics_cleanup_deletes_expired_buckets_in_batches() {
let repository = Arc::new(InMemoryProxyNodeRepository::seed(vec![
sample_connected_proxy_node("node-metrics-1", 30, 1),
sample_connected_proxy_node("node-metrics-2", 30, 1),
sample_connected_proxy_node("node-metrics-3", 30, 1),
]));
let data = GatewayDataState::with_proxy_node_repository_for_tests(Arc::clone(&repository))
.with_system_config_values_for_tests([
("enable_auto_cleanup".to_string(), json!(true)),
("proxy_node_metrics_1m_retention_days".to_string(), json!(1)),
("proxy_node_metrics_1h_retention_days".to_string(), json!(1)),
(
"proxy_node_metrics_cleanup_batch_size".to_string(),
json!(1),
),
]);
for (idx, node_id) in ["node-metrics-1", "node-metrics-2", "node-metrics-3"]
.into_iter()
.enumerate()
{
repository
.apply_heartbeat(&ProxyNodeHeartbeatMutation {
node_id: node_id.to_string(),
heartbeat_interval: Some(30),
active_connections: Some(i32::try_from(idx + 1).unwrap()),
total_requests_delta: None,
avg_latency_ms: None,
failed_requests_delta: None,
dns_failures_delta: None,
stream_errors_delta: None,
proxy_metadata: Some(json!({
"tunnel_metrics": {
"connect_errors": idx + 1,
"disconnects": 0,
"error_events_total": 0,
"ws_in_bytes": idx + 1,
"ws_out_bytes": idx + 1,
"ws_in_frames": idx + 1,
"ws_out_frames": idx + 1,
"heartbeat_rtt_last_ms": 10
}
})),
proxy_version: Some("1.0.0".to_string()),
})
.await
.expect("heartbeat should write metrics");
}
let now = chrono::Utc::now().timestamp().max(0) as u64;
let old_bucket = bucket_start_unix_secs(now, ProxyNodeMetricsStep::OneMinute);
let cleanup = data
.cleanup_proxy_node_metrics(
old_bucket.saturating_add(60),
old_bucket.saturating_add(3_600),
1,
)
.await
.expect("direct cleanup should delete a limited batch");
assert_eq!(cleanup.deleted_1m_rows, 1);
assert_eq!(cleanup.deleted_1h_rows, 1);
let cleanup = cleanup_proxy_node_metrics_at(&data, now.saturating_add(2 * 86_400))
.await
.expect("runtime cleanup should loop over batches");
assert_eq!(cleanup.deleted_1m_rows, 2);
assert_eq!(cleanup.deleted_1h_rows, 2);
}
#[tokio::test]
async fn proxy_node_metrics_cleanup_respects_auto_cleanup_toggle() {
let repository = Arc::new(InMemoryProxyNodeRepository::seed(vec![
sample_connected_proxy_node("node-metrics-disabled", 30, 1),
]));
let data = GatewayDataState::with_proxy_node_repository_for_tests(Arc::clone(&repository))
.with_system_config_values_for_tests([
("enable_auto_cleanup".to_string(), json!(false)),
("proxy_node_metrics_1m_retention_days".to_string(), json!(1)),
("proxy_node_metrics_1h_retention_days".to_string(), json!(1)),
(
"proxy_node_metrics_cleanup_batch_size".to_string(),
json!(1),
),
]);
repository
.apply_heartbeat(&ProxyNodeHeartbeatMutation {
node_id: "node-metrics-disabled".to_string(),
heartbeat_interval: Some(30),
active_connections: Some(1),
total_requests_delta: None,
avg_latency_ms: None,
failed_requests_delta: None,
dns_failures_delta: None,
stream_errors_delta: None,
proxy_metadata: Some(json!({
"tunnel_metrics": {
"connect_errors": 1,
"disconnects": 0,
"error_events_total": 0,
"ws_in_bytes": 1,
"ws_out_bytes": 1,
"ws_in_frames": 1,
"ws_out_frames": 1,
"heartbeat_rtt_last_ms": 10
}
})),
proxy_version: Some("1.0.0".to_string()),
})
.await
.expect("heartbeat should write metrics");
let cleanup = cleanup_proxy_node_metrics_once(&data)
.await
.expect("runtime cleanup should short-circuit");
assert_eq!(cleanup.deleted_1m_rows, 0);
assert_eq!(cleanup.deleted_1h_rows, 0);
}
#[test]
fn usage_cleanup_window_uses_non_overlapping_ranges() {
let now_utc = "2026-03-18T03:00:00Z"

View File

@@ -675,10 +675,15 @@ impl AppState {
&self,
retain_1m_from_unix_secs: u64,
retain_1h_from_unix_secs: u64,
delete_limit: usize,
) -> Result<aether_data::repository::proxy_nodes::ProxyNodeMetricsCleanupSummary, GatewayError>
{
self.data
.cleanup_proxy_node_metrics(retain_1m_from_unix_secs, retain_1h_from_unix_secs)
.cleanup_proxy_node_metrics(
retain_1m_from_unix_secs,
retain_1h_from_unix_secs,
delete_limit,
)
.await
.map_err(|err| GatewayError::Internal(err.to_string()))
}