perf(gateway): scale request hot paths for 20k streams

Shard and singleflight hot-path caches, batch and prioritize candidate and usage lifecycle persistence, and extend database and pressure-test instrumentation for 20k concurrent streams.
This commit is contained in:
elky
2026-07-22 02:11:08 +08:00
parent 7756c0913f
commit fc92c4f431
124 changed files with 36325 additions and 3217 deletions
+430 -118
View File
@@ -6,7 +6,7 @@ use aether_data_contracts::repository::usage::{StoredRequestUsageAudit, UpsertUs
use aether_data_contracts::DataLayerError;
use aether_runtime_state::{RuntimeQueueEntry, RuntimeQueueStore};
use async_trait::async_trait;
use tokio::sync::mpsc;
use tokio::sync::{mpsc, Notify};
use tracing::warn;
use crate::executor::spawn_on_usage_background_runtime;
@@ -40,10 +40,59 @@ pub trait ManualProxyNodeCounter: Send + Sync {
#[async_trait]
pub trait UsageRecordWriter: Send + Sync {
/// Native batch support is opt-in; the default preserves one-row writes for other backends.
fn supports_first_byte_usage_batch(&self) -> bool {
false
}
/// Stable identity for the underlying first-byte writer. Implementations that opt into
/// batching must return the same value for clones backed by the same repository.
fn first_byte_usage_writer_identity(&self) -> Option<usize> {
None
}
/// Native pending batching is opt-in because it must retain the complete usage audit write
/// contract, not just the base lifecycle row.
fn supports_pending_usage_batch(&self) -> bool {
false
}
/// Stable identity for clones backed by the same pending usage repository.
fn pending_usage_writer_identity(&self) -> Option<usize> {
None
}
async fn upsert_usage_record(
&self,
record: UpsertUsageRecord,
) -> Result<Option<StoredRequestUsageAudit>, DataLayerError>;
async fn upsert_first_byte_usage_record(
&self,
record: UpsertUsageRecord,
) -> Result<(), DataLayerError> {
self.upsert_usage_record(record).await.map(|_| ())
}
async fn upsert_first_byte_usage_records(
&self,
records: Vec<UpsertUsageRecord>,
) -> Result<(), DataLayerError> {
for record in records {
self.upsert_first_byte_usage_record(record).await?;
}
Ok(())
}
async fn upsert_pending_usage_records(
&self,
records: Vec<UpsertUsageRecord>,
) -> Result<(), DataLayerError> {
for record in records {
self.upsert_usage_record(record).await?;
}
Ok(())
}
}
pub struct UsageDataEventRecorder<T> {
@@ -126,16 +175,24 @@ pub struct UsageQueueWorker {
#[derive(Clone, Default)]
pub(crate) struct UsageWorkerControl {
shutdown: Arc<AtomicBool>,
shutdown_notify: Arc<Notify>,
}
impl UsageWorkerControl {
pub(crate) fn request_shutdown(&self) {
self.shutdown.store(true, Ordering::Release);
self.shutdown_notify.notify_one();
}
fn should_shutdown(&self) -> bool {
self.shutdown.load(Ordering::Acquire)
}
async fn wait_for_shutdown(&self) {
while !self.should_shutdown() {
self.shutdown_notify.notified().await;
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
@@ -307,77 +364,106 @@ impl UsageQueueWorker {
reclaim_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
reclaim_interval.tick().await;
let mut reclaim_due = false;
loop {
if self.should_shutdown() {
break;
}
tokio::select! {
_ = reclaim_interval.tick() => {
match self.queue.claim_stale(&self.consumer, "0-0").await {
Ok(entries) => {
self.report_reclaimed(entries.len());
if let Err(err) = self.process_entries(entries).await {
self.report_process_failed();
warn!(
event_name = "usage_worker_reclaim_process_failed",
log_type = "ops",
worker_consumer = %self.consumer,
worker_group = %self.config.consumer_group,
error = %err,
"usage worker failed while reclaiming stale entries"
);
}
}
Err(err) => {
self.report_reclaim_failed();
warn!(
event_name = "usage_worker_reclaim_failed",
log_type = "ops",
worker_consumer = %self.consumer,
worker_group = %self.config.consumer_group,
error = %err,
"usage worker failed to reclaim stale entries"
);
let result = {
let mut read_future = Box::pin(self.queue.read_group(&self.consumer));
loop {
tokio::select! {
biased;
// A command already delivered by Redis remains in the PEL and is recovered
// by a subsequent worker reclaim after the idle period.
_ = self.wait_for_shutdown() => return,
_ = reclaim_interval.tick(), if !reclaim_due => {
// Do not reclaim while XREADGROUP is in flight. Redis can add an entry
// to this consumer's PEL before delivering the response; claiming it in
// that window would make both paths process the same stream entry.
reclaim_due = true;
}
result = &mut read_future => break result,
}
}
result = self.queue.read_group(&self.consumer) => {
match result {
Ok(entries) => {
self.report_read(entries.len());
if entries.is_empty() && self.should_shutdown() {
break;
}
if let Err(err) = self.process_entries(entries).await {
self.report_process_failed();
warn!(
event_name = "usage_worker_process_failed",
log_type = "ops",
worker_consumer = %self.consumer,
worker_group = %self.config.consumer_group,
error = %err,
"usage worker failed to process queue entries"
);
tokio::time::sleep(Duration::from_millis(250)).await;
}
if self.should_shutdown() {
break;
}
}
Err(err) => {
self.report_read_failed();
warn!(
event_name = "usage_worker_read_failed",
log_type = "ops",
worker_consumer = %self.consumer,
worker_group = %self.config.consumer_group,
error = %err,
"usage worker failed to read queue"
);
tokio::time::sleep(Duration::from_millis(500)).await;
}
};
match result {
Ok(entries) => {
self.report_read(entries.len());
if let Err(err) = self.process_entries(entries).await {
self.report_process_failed();
warn!(
event_name = "usage_worker_process_failed",
log_type = "ops",
worker_consumer = %self.consumer,
worker_group = %self.config.consumer_group,
error = %err,
"usage worker failed to process queue entries"
);
tokio::time::sleep(Duration::from_millis(250)).await;
}
}
Err(err) => {
self.report_read_failed();
warn!(
event_name = "usage_worker_read_failed",
log_type = "ops",
worker_consumer = %self.consumer,
worker_group = %self.config.consumer_group,
error = %err,
"usage worker failed to read queue"
);
tokio::time::sleep(Duration::from_millis(500)).await;
}
}
if self.should_shutdown() {
break;
}
if reclaim_due {
reclaim_due = false;
self.reclaim_stale_entries().await;
}
}
}
async fn wait_for_shutdown(&self) {
match self.control.as_ref() {
Some(control) => control.wait_for_shutdown().await,
None => std::future::pending().await,
}
}
async fn reclaim_stale_entries(&self) {
match self.queue.claim_stale(&self.consumer, "0-0").await {
Ok(entries) => {
self.report_reclaimed(entries.len());
if let Err(err) = self.process_entries(entries).await {
self.report_process_failed();
warn!(
event_name = "usage_worker_reclaim_process_failed",
log_type = "ops",
worker_consumer = %self.consumer,
worker_group = %self.config.consumer_group,
error = %err,
"usage worker failed while reclaiming stale entries"
);
}
}
Err(err) => {
self.report_reclaim_failed();
warn!(
event_name = "usage_worker_reclaim_failed",
log_type = "ops",
worker_consumer = %self.consumer,
worker_group = %self.config.consumer_group,
error = %err,
"usage worker failed to reclaim stale entries"
);
}
}
}
@@ -588,13 +674,14 @@ where
pub async fn write_event_record<T>(data: &T, event: &UsageEvent) -> Result<(), DataLayerError>
where
T: UsageRecordWriter + UsageSettlementWriter + ManualProxyNodeCounter + Send + Sync,
T: UsageRecordWriter + UsageSettlementWriter + Send + Sync,
{
let record = build_upsert_usage_record_from_event(event)?;
if let Some(stored) = data.upsert_usage_record(record).await? {
settle_usage_if_needed(data, &stored).await?;
}
increment_manual_proxy_node_from_event(data, event).await;
// Manual proxy traffic is counted at the actual transport-attempt boundary. Usage events are
// replayable, so emitting that side effect here would count normal requests and reclaims twice.
Ok(())
}
@@ -621,56 +708,6 @@ where
}
}
async fn increment_manual_proxy_node_from_event<T>(data: &T, event: &UsageEvent)
where
T: ManualProxyNodeCounter + Send + Sync,
{
let is_terminal = matches!(
event.event_type,
crate::UsageEventType::Completed | crate::UsageEventType::Failed
);
if !is_terminal {
return;
}
let Some(node_id) = extract_manual_proxy_node_id(event) else {
return;
};
let failed = matches!(event.event_type, crate::UsageEventType::Failed);
let failed_delta = if failed { 1i64 } else { 0i64 };
let latency_ms = event.data.response_time_ms.map(|v| v as i64);
if let Err(err) = data
.increment_manual_proxy_node_requests(&node_id, 1, failed_delta, latency_ms)
.await
{
warn!(
event_name = "manual_proxy_node_increment_failed",
log_type = "ops",
node_id = %node_id,
error = ?err,
"failed to increment manual proxy node request count"
);
}
}
fn extract_manual_proxy_node_id(event: &UsageEvent) -> Option<String> {
let metadata = event.data.request_metadata.as_ref()?;
let proxy = metadata.get("proxy")?.as_object()?;
let mode = proxy
.get("mode")
.and_then(serde_json::Value::as_str)
.unwrap_or("")
.trim();
if mode == "tunnel" {
return None;
}
proxy
.get("node_id")
.and_then(serde_json::Value::as_str)
.map(str::trim)
.filter(|v| !v.is_empty())
.map(String::from)
}
fn consumer_name(worker_index: Option<usize>) -> String {
let host = std::env::var("HOSTNAME")
.ok()
@@ -685,7 +722,8 @@ fn consumer_name(worker_index: Option<usize>) -> String {
#[cfg(test)]
mod tests {
use std::sync::atomic::{AtomicBool, Ordering};
use std::collections::BTreeMap;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
@@ -694,13 +732,17 @@ mod tests {
};
use aether_data_contracts::repository::usage::{StoredRequestUsageAudit, UpsertUsageRecord};
use aether_data_contracts::DataLayerError;
use aether_runtime_state::{MemoryRuntimeStateConfig, RuntimeQueueStore, RuntimeState};
use aether_runtime_state::{
MemoryRuntimeStateConfig, RuntimeQueueEntry, RuntimeQueueReclaimConfig, RuntimeQueueStats,
RuntimeQueueStore, RuntimeState,
};
use async_trait::async_trait;
use tokio::sync::Notify;
use super::{
build_usage_queue_worker_with_record_gate, usage_event_record_error_is_permanent,
write_event_record, ManualProxyNodeCounter, UsageEventRecorder, UsageQueueWorker,
UsageRecordWriter,
UsageRecordWriter, UsageWorkerControl,
};
use crate::runtime::UsageWorkerRecordConcurrencyGate;
use crate::UsageBillingEventEnricher;
@@ -714,6 +756,7 @@ mod tests {
records: Mutex<Vec<UpsertUsageRecord>>,
settlements: Mutex<Vec<UsageSettlementInput>>,
enrich_calls: Mutex<Vec<String>>,
manual_proxy_counter_calls: AtomicUsize,
}
#[derive(Default)]
@@ -729,6 +772,129 @@ mod tests {
records: Mutex<Vec<String>>,
}
struct ReadReclaimRaceProbeQueue {
entry: RuntimeQueueEntry,
read_calls: AtomicUsize,
first_read_cancelled: AtomicBool,
read_completed: AtomicUsize,
release_read: Notify,
reclaim_calls: AtomicUsize,
acked: AtomicBool,
}
impl ReadReclaimRaceProbeQueue {
fn new(entry: RuntimeQueueEntry) -> Self {
Self {
entry,
read_calls: AtomicUsize::new(0),
first_read_cancelled: AtomicBool::new(false),
read_completed: AtomicUsize::new(0),
release_read: Notify::new(),
reclaim_calls: AtomicUsize::new(0),
acked: AtomicBool::new(false),
}
}
}
struct FirstReadDropGuard<'a> {
cancelled: &'a AtomicBool,
completed: bool,
}
impl Drop for FirstReadDropGuard<'_> {
fn drop(&mut self) {
if !self.completed {
self.cancelled.store(true, Ordering::Release);
}
}
}
#[async_trait]
impl RuntimeQueueStore for ReadReclaimRaceProbeQueue {
async fn ensure_consumer_group(
&self,
_stream: &str,
_group: &str,
_start_id: &str,
) -> Result<(), DataLayerError> {
Ok(())
}
async fn append_fields_with_maxlen(
&self,
_stream: &str,
_fields: &BTreeMap<String, String>,
_maxlen: Option<usize>,
) -> Result<String, DataLayerError> {
Ok("0-0".to_string())
}
async fn read_group(
&self,
_stream: &str,
_group: &str,
_consumer: &str,
_count: usize,
_block_ms: Option<u64>,
) -> Result<Vec<RuntimeQueueEntry>, DataLayerError> {
let call_index = self.read_calls.fetch_add(1, Ordering::AcqRel);
let mut first_read_guard = (call_index == 0).then(|| FirstReadDropGuard {
cancelled: &self.first_read_cancelled,
completed: false,
});
self.release_read.notified().await;
if let Some(guard) = first_read_guard.as_mut() {
guard.completed = true;
}
self.read_completed.fetch_add(1, Ordering::AcqRel);
Ok((call_index == 0)
.then(|| self.entry.clone())
.into_iter()
.collect())
}
async fn claim_stale(
&self,
_stream: &str,
_group: &str,
_consumer: &str,
_start_id: &str,
_config: RuntimeQueueReclaimConfig,
) -> Result<Vec<RuntimeQueueEntry>, DataLayerError> {
self.reclaim_calls.fetch_add(1, Ordering::AcqRel);
Ok((!self.acked.load(Ordering::Acquire))
.then(|| self.entry.clone())
.into_iter()
.collect())
}
async fn ack(
&self,
_stream: &str,
_group: &str,
ids: &[String],
) -> Result<usize, DataLayerError> {
if ids.iter().any(|id| id == &self.entry.id) {
self.acked.store(true, Ordering::Release);
Ok(1)
} else {
Ok(0)
}
}
async fn delete(&self, _stream: &str, ids: &[String]) -> Result<usize, DataLayerError> {
Ok(ids.len())
}
async fn stats(
&self,
_stream: &str,
_group: Option<&str>,
) -> Result<RuntimeQueueStats, DataLayerError> {
Ok(RuntimeQueueStats::default())
}
}
#[async_trait]
impl UsageRecordWriter for TestUsageStore {
async fn upsert_usage_record(
@@ -813,6 +979,8 @@ mod tests {
_failed_delta: i64,
_latency_ms: Option<i64>,
) -> Result<(), aether_data_contracts::DataLayerError> {
self.manual_proxy_counter_calls
.fetch_add(1, Ordering::AcqRel);
Ok(())
}
}
@@ -979,6 +1147,28 @@ mod tests {
assert_eq!(settlements[0].request_id, "req-worker-123");
}
#[tokio::test]
async fn replayable_usage_write_does_not_duplicate_transport_owned_proxy_counter() {
let store = TestUsageStore::default();
let mut event = sample_event();
event.data.request_metadata = Some(serde_json::json!({
"proxy": {"mode": "manual", "node_id": "manual-node-1"}
}));
write_event_record(&store, &event)
.await
.expect("first usage write should succeed");
write_event_record(&store, &event)
.await
.expect("replayed usage write should succeed");
assert_eq!(
store.manual_proxy_counter_calls.load(Ordering::Acquire),
0,
"proxy traffic belongs to the transport attempt, not the replayable usage worker"
);
}
#[tokio::test]
async fn data_event_recorder_enriches_terminal_event_before_write() {
let store = Arc::new(TestUsageStore::default());
@@ -1137,6 +1327,128 @@ mod tests {
assert_eq!(store.records.lock().expect("records lock").len(), 4);
}
#[tokio::test]
async fn usage_worker_defers_reclaim_until_inflight_read_is_processed() {
let event = sample_event();
let queue = Arc::new(ReadReclaimRaceProbeQueue::new(RuntimeQueueEntry {
id: "1-0".to_string(),
fields: event
.to_stream_fields()
.expect("usage event should serialize"),
}));
let queue_runner: Arc<dyn RuntimeQueueStore> = queue.clone();
let config = UsageRuntimeConfig {
enabled: true,
stream_key: "usage:test:worker:read-reclaim-race".to_string(),
consumer_group: "usage:test:worker:read-reclaim-race-group".to_string(),
dlq_stream_key: "usage:test:worker:read-reclaim-race-dlq".to_string(),
consumer_batch_size: 1,
consumer_block_ms: 1_000,
reclaim_interval_ms: 10,
reclaim_idle_ms: 1,
reclaim_count: 1,
..UsageRuntimeConfig::default()
};
let recorder = Arc::new(SelectiveFailingRecorder::default());
let worker_recorder: Arc<dyn UsageEventRecorder> = recorder.clone();
let control = UsageWorkerControl::default();
let (telemetry_tx, _telemetry_rx) = tokio::sync::mpsc::channel(8);
let worker = UsageQueueWorker::new(queue_runner, worker_recorder, config, None)
.expect("worker should build")
.with_supervisor(control.clone(), telemetry_tx);
let handle = tokio::spawn(worker.run());
tokio::time::timeout(Duration::from_secs(1), async {
while queue.read_calls.load(Ordering::Acquire) == 0 {
tokio::task::yield_now().await;
}
})
.await
.expect("worker should start the blocking read");
tokio::time::sleep(Duration::from_millis(50)).await;
assert_eq!(
queue.reclaim_calls.load(Ordering::Acquire),
0,
"reclaim must not run while XREADGROUP can still return the same PEL entry"
);
assert!(recorder.calls.lock().expect("calls lock").is_empty());
queue.release_read.notify_one();
tokio::time::timeout(Duration::from_secs(1), async {
while !queue.acked.load(Ordering::Acquire)
|| queue.reclaim_calls.load(Ordering::Acquire) == 0
{
tokio::task::yield_now().await;
}
})
.await
.expect("read entry should be processed before deferred reclaim runs");
assert_eq!(
recorder.calls.lock().expect("calls lock").as_slice(),
[event.request_id.as_str()],
"the stream entry must be recorded exactly once"
);
assert_eq!(queue.read_completed.load(Ordering::Acquire), 1);
assert!(!queue.first_read_cancelled.load(Ordering::Acquire));
control.request_shutdown();
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("worker should stop promptly")
.expect("worker task should not panic");
}
#[tokio::test]
async fn usage_worker_shutdown_cancels_blocking_read_promptly() {
let event = sample_event();
let queue = Arc::new(ReadReclaimRaceProbeQueue::new(RuntimeQueueEntry {
id: "2-0".to_string(),
fields: event
.to_stream_fields()
.expect("usage event should serialize"),
}));
let queue_runner: Arc<dyn RuntimeQueueStore> = queue.clone();
let config = UsageRuntimeConfig {
enabled: true,
stream_key: "usage:test:worker:shutdown-read".to_string(),
consumer_group: "usage:test:worker:shutdown-read-group".to_string(),
dlq_stream_key: "usage:test:worker:shutdown-read-dlq".to_string(),
consumer_batch_size: 1,
consumer_block_ms: 60_000,
reclaim_interval_ms: 10,
reclaim_idle_ms: 1,
reclaim_count: 1,
..UsageRuntimeConfig::default()
};
let recorder: Arc<dyn UsageEventRecorder> = Arc::new(SelectiveFailingRecorder::default());
let control = UsageWorkerControl::default();
let (telemetry_tx, _telemetry_rx) = tokio::sync::mpsc::channel(8);
let worker = UsageQueueWorker::new(queue_runner, recorder, config, None)
.expect("worker should build")
.with_supervisor(control.clone(), telemetry_tx);
let handle = tokio::spawn(worker.run());
tokio::time::timeout(Duration::from_secs(1), async {
while queue.read_calls.load(Ordering::Acquire) == 0 {
tokio::task::yield_now().await;
}
})
.await
.expect("worker should start the blocking read");
control.request_shutdown();
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("shutdown should interrupt the blocking read")
.expect("worker task should not panic");
assert!(queue.first_read_cancelled.load(Ordering::Acquire));
assert_eq!(queue.read_completed.load(Ordering::Acquire), 0);
assert_eq!(queue.reclaim_calls.load(Ordering::Acquire), 0);
}
#[test]
fn usage_event_record_error_classifies_permanent_failures() {
assert!(usage_event_record_error_is_permanent(