Protect background workers under DB pool pressure

This commit is contained in:
fawney19
2026-05-22 14:11:47 +08:00
parent 2b32b9a445
commit 8966fd6aac
6 changed files with 236 additions and 3 deletions

View File

@@ -8,7 +8,7 @@ use aether_data_contracts::repository::video_tasks::{
use aether_usage_runtime::{build_upsert_usage_record_from_event, settle_usage_if_needed}; use aether_usage_runtime::{build_upsert_usage_record_from_event, settle_usage_if_needed};
use serde_json::{Map, Value}; use serde_json::{Map, Value};
use tokio::task::JoinHandle; use tokio::task::JoinHandle;
use tracing::{info, warn}; use tracing::{debug, info, warn};
use crate::log_ids::short_request_id; use crate::log_ids::short_request_id;
use crate::usage::{UsageEvent, UsageEventData, UsageEventType}; use crate::usage::{UsageEvent, UsageEventData, UsageEventType};
@@ -148,8 +148,20 @@ pub(crate) fn spawn_video_task_poller(state: AppState) -> Option<JoinHandle<()>>
let mut interval = tokio::time::interval(config.interval); let mut interval = tokio::time::interval(config.interval);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await; interval.tick().await;
let mut deferred_since = None;
loop { loop {
interval.tick().await; interval.tick().await;
if state
.data
.should_defer_maintenance_for_database_pool_pressure(&mut deferred_since)
{
debug!(
event_name = "video_task_poller_deferred",
log_type = "event",
"gateway video task poller deferred because database pool has no idle reserve"
);
continue;
}
if let Err(err) = poll_video_tasks_once(&state, config.batch_size).await { if let Err(err) = poll_video_tasks_once(&state, config.batch_size).await {
warn!( warn!(
event_name = "video_task_poller_tick_failed", event_name = "video_task_poller_tick_failed",

View File

@@ -42,8 +42,12 @@ use aether_data_contracts::repository::usage::{
}; };
use aether_runtime_state::RuntimeQueueStore; use aether_runtime_state::RuntimeQueueStore;
use aether_video_tasks_core::read_data_backed_video_task_response; use aether_video_tasks_core::read_data_backed_video_task_response;
use std::time::{Duration, Instant};
impl GatewayDataState { impl GatewayDataState {
const MAINTENANCE_POOL_IDLE_RESERVE: usize = 1;
const MAINTENANCE_POOL_PRESSURE_MAX_DEFER: Duration = Duration::from_secs(30);
pub(crate) async fn run_database_maintenance( pub(crate) async fn run_database_maintenance(
&self, &self,
table_names: &[&str], table_names: &[&str],
@@ -112,6 +116,47 @@ impl GatewayDataState {
.and_then(|backends| backends.database_pool_summary()) .and_then(|backends| backends.database_pool_summary())
} }
pub(crate) fn database_pool_under_maintenance_pressure(&self) -> bool {
self.database_pool_summary()
.as_ref()
.is_some_and(Self::database_pool_summary_under_maintenance_pressure)
}
pub(crate) fn database_pool_summary_under_maintenance_pressure(
summary: &aether_data::DatabasePoolSummary,
) -> bool {
summary.checked_out > 0 && summary.idle <= Self::MAINTENANCE_POOL_IDLE_RESERVE
}
pub(crate) fn should_defer_maintenance_for_database_pool_pressure(
&self,
deferred_since: &mut Option<Instant>,
) -> bool {
Self::should_defer_maintenance_for_pool_pressure_state(
self.database_pool_under_maintenance_pressure(),
deferred_since,
)
}
pub(crate) fn should_defer_maintenance_for_pool_pressure_state(
pool_under_pressure: bool,
deferred_since: &mut Option<Instant>,
) -> bool {
if !pool_under_pressure {
*deferred_since = None;
return false;
}
let now = Instant::now();
let since = deferred_since.get_or_insert(now);
if now.duration_since(*since) >= Self::MAINTENANCE_POOL_PRESSURE_MAX_DEFER {
*deferred_since = None;
return false;
}
true
}
pub(crate) async fn aggregate_wallet_daily_usage( pub(crate) async fn aggregate_wallet_daily_usage(
&self, &self,
input: &WalletDailyUsageAggregationInput, input: &WalletDailyUsageAggregationInput,

View File

@@ -1,4 +1,5 @@
use std::sync::Arc; use std::sync::Arc;
use std::time::{Duration, Instant};
use aether_crypto::{encrypt_python_fernet_plaintext, DEVELOPMENT_ENCRYPTION_KEY}; use aether_crypto::{encrypt_python_fernet_plaintext, DEVELOPMENT_ENCRYPTION_KEY};
use aether_data::repository::auth::{ use aether_data::repository::auth::{
@@ -51,6 +52,65 @@ fn disabled_gateway_data_state_has_no_backends() {
assert!(!state.has_video_task_reader()); assert!(!state.has_video_task_reader());
} }
#[test]
fn maintenance_pool_pressure_keeps_idle_reserve_for_foreground_work() {
let pressured = aether_data::DatabasePoolSummary {
driver: DatabaseDriver::Postgres,
checked_out: 6,
pool_size: 6,
idle: 0,
max_connections: 20,
usage_rate: 30.0,
};
assert!(GatewayDataState::database_pool_summary_under_maintenance_pressure(&pressured));
let one_idle_left = aether_data::DatabasePoolSummary {
driver: DatabaseDriver::Postgres,
checked_out: 5,
pool_size: 6,
idle: 1,
max_connections: 20,
usage_rate: 25.0,
};
assert!(GatewayDataState::database_pool_summary_under_maintenance_pressure(&one_idle_left));
let idle = aether_data::DatabasePoolSummary {
driver: DatabaseDriver::Postgres,
checked_out: 0,
pool_size: 4,
idle: 4,
max_connections: 20,
usage_rate: 0.0,
};
assert!(!GatewayDataState::database_pool_summary_under_maintenance_pressure(&idle));
}
#[test]
fn maintenance_pool_pressure_deferral_has_timeout() {
let mut deferred_since = None;
assert!(
GatewayDataState::should_defer_maintenance_for_pool_pressure_state(
true,
&mut deferred_since
)
);
assert!(deferred_since.is_some());
assert!(
!GatewayDataState::should_defer_maintenance_for_pool_pressure_state(
false,
&mut deferred_since
)
);
assert!(deferred_since.is_none());
let mut stale_defer = Some(Instant::now() - Duration::from_secs(31));
assert!(
!GatewayDataState::should_defer_maintenance_for_pool_pressure_state(true, &mut stale_defer)
);
assert!(stale_defer.is_none());
}
#[tokio::test] #[tokio::test]
async fn postgres_gateway_data_state_builds_video_task_reader() { async fn postgres_gateway_data_state_builds_video_task_reader() {
let state = GatewayDataState::from_config(GatewayDataConfig::from_postgres_url( let state = GatewayDataState::from_config(GatewayDataConfig::from_postgres_url(

View File

@@ -738,8 +738,21 @@ pub(crate) fn spawn_account_self_check_worker(
let mut interval = tokio::time::interval(config.scan_interval); let mut interval = tokio::time::interval(config.scan_interval);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await; interval.tick().await;
let mut deferred_since = None;
loop { loop {
interval.tick().await; interval.tick().await;
if state
.data
.should_defer_maintenance_for_database_pool_pressure(&mut deferred_since)
{
debug!(
event_name = "maintenance_worker_deferred",
log_type = "ops",
worker = "account_self_check",
"gateway account self-check deferred because database pool has no idle reserve"
);
continue;
}
if let Err(err) = perform_account_self_check_once_with_config(&state, config).await { if let Err(err) = perform_account_self_check_once_with_config(&state, config).await {
warn!( warn!(
error = ?err, error = ?err,

View File

@@ -393,8 +393,21 @@ pub(crate) fn spawn_pool_score_rebuild_worker(
} }
let mut interval = tokio::time::interval(config.interval); let mut interval = tokio::time::interval(config.interval);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
let mut deferred_since = None;
loop { loop {
interval.tick().await; interval.tick().await;
if state
.data
.should_defer_maintenance_for_database_pool_pressure(&mut deferred_since)
{
debug!(
event_name = "maintenance_worker_deferred",
log_type = "ops",
worker = "pool_score_rebuild",
"gateway pool score rebuild deferred because database pool has no idle reserve"
);
continue;
}
match perform_pool_score_rebuild_once_with_config(&state, config).await { match perform_pool_score_rebuild_once_with_config(&state, config).await {
Ok(summary) if summary.scores_upserted > 0 => { Ok(summary) if summary.scores_upserted > 0 => {
info!( info!(

View File

@@ -1,7 +1,8 @@
use std::sync::Arc; use std::sync::Arc;
use std::time::Instant;
use chrono::Utc; use chrono::Utc;
use tracing::warn; use tracing::{debug, warn};
use crate::data::GatewayDataState; use crate::data::GatewayDataState;
use crate::AppState; use crate::AppState;
@@ -46,6 +47,37 @@ fn log_maintenance_worker_failure(
); );
} }
fn should_defer_for_database_pressure(
data: &GatewayDataState,
worker: &'static str,
deferred_since: &mut Option<Instant>,
) -> bool {
let Some(summary) = data.database_pool_summary() else {
*deferred_since = None;
return false;
};
if !GatewayDataState::should_defer_maintenance_for_pool_pressure_state(
GatewayDataState::database_pool_summary_under_maintenance_pressure(&summary),
deferred_since,
) {
return false;
}
debug!(
event_name = "maintenance_worker_deferred",
log_type = "ops",
worker,
driver = %summary.driver,
checked_out = summary.checked_out,
pool_size = summary.pool_size,
idle = summary.idle,
max_connections = summary.max_connections,
usage_rate = summary.usage_rate,
"gateway maintenance worker deferred because database pool has no idle reserve"
);
true
}
pub(crate) fn spawn_audit_cleanup_worker( pub(crate) fn spawn_audit_cleanup_worker(
data: Arc<GatewayDataState>, data: Arc<GatewayDataState>,
) -> Option<tokio::task::JoinHandle<()>> { ) -> Option<tokio::task::JoinHandle<()>> {
@@ -177,8 +209,19 @@ pub(crate) fn spawn_usage_counter_flush_worker(
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await; interval.tick().await;
let mut last_delta_cleanup = tokio::time::Instant::now(); let mut last_delta_cleanup = tokio::time::Instant::now();
let mut usage_counter_flush_deferred_since = None;
let mut usage_counter_delta_cleanup_deferred_since = None;
loop { loop {
if should_defer_for_database_pressure(
&data,
"usage_counter_flush",
&mut usage_counter_flush_deferred_since,
) {
interval.tick().await;
continue;
}
let mut batches = 0_usize; let mut batches = 0_usize;
while batches < USAGE_COUNTER_FLUSH_CATCH_UP_BURST_LIMIT { while batches < USAGE_COUNTER_FLUSH_CATCH_UP_BURST_LIMIT {
match run_usage_counter_flush_once(&data, USAGE_COUNTER_FLUSH_BATCH_SIZE).await { match run_usage_counter_flush_once(&data, USAGE_COUNTER_FLUSH_BATCH_SIZE).await {
@@ -197,7 +240,18 @@ pub(crate) fn spawn_usage_counter_flush_worker(
} }
if last_delta_cleanup.elapsed() >= USAGE_COUNTER_DELTA_CLEANUP_INTERVAL { if last_delta_cleanup.elapsed() >= USAGE_COUNTER_DELTA_CLEANUP_INTERVAL {
if let Err(err) = cleanup_processed_usage_counter_deltas_once( if should_defer_for_database_pressure(
&data,
"usage_counter_delta_cleanup",
&mut usage_counter_delta_cleanup_deferred_since,
) {
debug!(
event_name = "maintenance_worker_deferred",
log_type = "ops",
worker = "usage_counter_delta_cleanup",
"gateway maintenance worker deferred cleanup under database pressure"
);
} else if let Err(err) = cleanup_processed_usage_counter_deltas_once(
&data, &data,
USAGE_COUNTER_DELTA_RETENTION_SECS, USAGE_COUNTER_DELTA_RETENTION_SECS,
USAGE_COUNTER_DELTA_CLEANUP_BATCH_SIZE, USAGE_COUNTER_DELTA_CLEANUP_BATCH_SIZE,
@@ -265,8 +319,16 @@ pub(crate) fn spawn_provider_quota_alert_worker(
let mut interval = tokio::time::interval(PROVIDER_QUOTA_ALERT_INTERVAL); let mut interval = tokio::time::interval(PROVIDER_QUOTA_ALERT_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await; interval.tick().await;
let mut deferred_since = None;
loop { loop {
interval.tick().await; interval.tick().await;
if should_defer_for_database_pressure(
&state.data,
"provider_quota_alert",
&mut deferred_since,
) {
continue;
}
if let Err(err) = perform_provider_quota_alert_once(&state).await { if let Err(err) = perform_provider_quota_alert_once(&state).await {
log_maintenance_worker_failure("provider_quota_alert", "tick", &err); log_maintenance_worker_failure("provider_quota_alert", "tick", &err);
} }
@@ -288,8 +350,16 @@ pub(crate) fn spawn_oauth_token_refresh_worker(
let mut interval = tokio::time::interval(OAUTH_TOKEN_REFRESH_INTERVAL); let mut interval = tokio::time::interval(OAUTH_TOKEN_REFRESH_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await; interval.tick().await;
let mut deferred_since = None;
loop { loop {
interval.tick().await; interval.tick().await;
if should_defer_for_database_pressure(
&state.data,
"oauth_token_refresh",
&mut deferred_since,
) {
continue;
}
if let Err(err) = perform_oauth_token_refresh_once(&state).await { if let Err(err) = perform_oauth_token_refresh_once(&state).await {
log_maintenance_worker_failure("oauth_token_refresh", "tick", &err); log_maintenance_worker_failure("oauth_token_refresh", "tick", &err);
} }
@@ -334,8 +404,12 @@ pub(crate) fn spawn_pending_cleanup_worker(
let mut interval = tokio::time::interval(PENDING_CLEANUP_INTERVAL); let mut interval = tokio::time::interval(PENDING_CLEANUP_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await; interval.tick().await;
let mut deferred_since = None;
loop { loop {
interval.tick().await; interval.tick().await;
if should_defer_for_database_pressure(&data, "pending_cleanup", &mut deferred_since) {
continue;
}
if let Err(err) = run_pending_cleanup_once(&data).await { if let Err(err) = run_pending_cleanup_once(&data).await {
log_maintenance_worker_failure("pending_cleanup", "tick", &err); log_maintenance_worker_failure("pending_cleanup", "tick", &err);
} }
@@ -357,8 +431,16 @@ pub(crate) fn spawn_proxy_node_stale_cleanup_worker(
let mut interval = tokio::time::interval(PROXY_NODE_STALE_SWEEP_INTERVAL); let mut interval = tokio::time::interval(PROXY_NODE_STALE_SWEEP_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await; interval.tick().await;
let mut deferred_since = None;
loop { loop {
interval.tick().await; interval.tick().await;
if should_defer_for_database_pressure(
&data,
"proxy_node_stale_cleanup",
&mut deferred_since,
) {
continue;
}
if let Err(err) = run_proxy_node_stale_cleanup_once(&data).await { if let Err(err) = run_proxy_node_stale_cleanup_once(&data).await {
log_maintenance_worker_failure("proxy_node_stale_cleanup", "tick", &err); log_maintenance_worker_failure("proxy_node_stale_cleanup", "tick", &err);
} }
@@ -407,8 +489,16 @@ pub(crate) fn spawn_proxy_upgrade_rollout_worker(
let mut interval = tokio::time::interval(PROXY_UPGRADE_ROLLOUT_INTERVAL); let mut interval = tokio::time::interval(PROXY_UPGRADE_ROLLOUT_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await; interval.tick().await;
let mut deferred_since = None;
loop { loop {
interval.tick().await; interval.tick().await;
if should_defer_for_database_pressure(
&state.data,
"proxy_upgrade_rollout",
&mut deferred_since,
) {
continue;
}
if let Err(err) = run_proxy_upgrade_rollout_once(&state).await { if let Err(err) = run_proxy_upgrade_rollout_once(&state).await {
log_maintenance_worker_failure("proxy_upgrade_rollout", "tick", &err); log_maintenance_worker_failure("proxy_upgrade_rollout", "tick", &err);
} }