Fix affinity key parsing and counters

This commit is contained in:
RWDai
2026-05-07 15:08:07 +08:00
parent 96f9be26da
commit d5d93eda11
4 changed files with 427 additions and 120 deletions

View File

@@ -2,61 +2,84 @@ use super::cache_types::AdminMonitoringCacheAffinityRecord;
use crate::cache::SchedulerAffinityTarget; use crate::cache::SchedulerAffinityTarget;
use crate::handlers::admin::request::AdminAppState; use crate::handlers::admin::request::AdminAppState;
use crate::GatewayError; use crate::GatewayError;
use aether_ai_formats::FormatId;
use std::time::Duration; use std::time::Duration;
fn parse_admin_monitoring_cache_affinity_key(raw_key: &str) -> Option<(String, String, String)> { #[derive(Debug, Clone, Copy)]
let parts = raw_key.split(':').collect::<Vec<_>>(); enum AdminMonitoringAffinityKeyKind {
let start = parts Cache,
.iter() Scheduler,
.position(|segment| *segment == "cache_affinity")?;
let affinity_key = parts.get(start + 1)?.trim();
if affinity_key.is_empty() {
return None;
}
let api_format = parts
.get(start + 2)
.map(|value| value.trim())
.filter(|value| !value.is_empty())
.unwrap_or("unknown")
.to_string();
let model_name = parts
.get(start + 3..)
.filter(|segments| !segments.is_empty())
.map(|segments| segments.join(":"))
.unwrap_or_else(|| "unknown".to_string());
Some((affinity_key.to_string(), api_format, model_name))
} }
fn parse_admin_monitoring_scheduler_affinity_key( #[derive(Debug, Clone, PartialEq, Eq)]
raw_key: &str, struct ParsedAdminMonitoringAffinityKey {
) -> Option<(String, String, String)> { affinity_key: String,
let parts = raw_key.split(':').collect::<Vec<_>>(); api_format: String,
let start = parts model_name: String,
.iter() client_family: Option<String>,
.position(|segment| *segment == "scheduler_affinity")?; session_hash: Option<String>,
let affinity_key = parts.get(start + 1)?.trim(); }
if affinity_key.is_empty() {
fn split_admin_monitoring_api_format_and_model(
segments: &[&str],
key_kind: AdminMonitoringAffinityKeyKind,
) -> Option<(String, String)> {
if segments.len() < 2 {
return None; return None;
} }
let remaining = parts.get(start + 2..)?; let max_api_format_segments = segments.len().saturating_sub(1).min(3);
if remaining.len() < 2 { for segment_count in (2..=max_api_format_segments).rev() {
return None; let candidate = segments[..segment_count]
.iter()
.map(|segment| segment.trim())
.collect::<Vec<_>>()
.join(":");
if FormatId::parse(&candidate).is_some() {
let model_name = segments[segment_count..]
.iter()
.map(|segment| segment.trim())
.filter(|segment| !segment.is_empty())
.collect::<Vec<_>>()
.join(":");
if model_name.is_empty() {
return None;
}
return Some((candidate, model_name));
}
} }
let (api_format, model_name_parts) = if remaining.len() == 2 { if matches!(key_kind, AdminMonitoringAffinityKeyKind::Scheduler)
(remaining[0].trim().to_string(), &remaining[1..]) && segments.len() >= 3
} else { && is_known_admin_monitoring_api_format_family(segments[0])
( {
format!("{}:{}", remaining[0].trim(), remaining[1].trim()), let api_format_kind = segments.get(1)?.trim();
&remaining[2..], if api_format_kind.is_empty() {
) return None;
}; }
if api_format.trim().is_empty() { let model_name = segments[2..]
return None; .iter()
.map(|segment| segment.trim())
.filter(|segment| !segment.is_empty())
.collect::<Vec<_>>()
.join(":");
if model_name.is_empty() {
return None;
}
return Some((
format!(
"{}:{api_format_kind}",
segments[0].trim().to_ascii_lowercase()
),
model_name,
));
} }
let model_name = model_name_parts let api_format = segments.first()?.trim();
if api_format.is_empty() {
return None;
}
let model_name = segments[1..]
.iter() .iter()
.map(|segment| segment.trim()) .map(|segment| segment.trim())
.filter(|segment| !segment.is_empty()) .filter(|segment| !segment.is_empty())
@@ -66,7 +89,105 @@ fn parse_admin_monitoring_scheduler_affinity_key(
return None; return None;
} }
Some((affinity_key.to_string(), api_format, model_name)) Some((api_format.to_string(), model_name))
}
fn is_known_admin_monitoring_api_format_family(value: &str) -> bool {
matches!(
value.trim().to_ascii_lowercase().as_str(),
"openai" | "claude" | "gemini" | "jina" | "doubao"
)
}
fn parse_admin_monitoring_cache_affinity_key(
raw_key: &str,
) -> Option<ParsedAdminMonitoringAffinityKey> {
let parts = raw_key.split(':').collect::<Vec<_>>();
let start = parts
.iter()
.position(|segment| *segment == "cache_affinity")?;
let affinity_key = parts.get(start + 1)?.trim();
if affinity_key.is_empty() {
return None;
}
let remaining = parts.get(start + 2..)?;
let (api_format, model_name) = split_admin_monitoring_api_format_and_model(
remaining,
AdminMonitoringAffinityKeyKind::Cache,
)
.unwrap_or_else(|| ("unknown".to_string(), "unknown".to_string()));
Some(ParsedAdminMonitoringAffinityKey {
affinity_key: affinity_key.to_string(),
api_format,
model_name,
client_family: None,
session_hash: None,
})
}
fn parse_admin_monitoring_scheduler_affinity_key(
raw_key: &str,
) -> Option<ParsedAdminMonitoringAffinityKey> {
let parts = raw_key.split(':').collect::<Vec<_>>();
let start = parts
.iter()
.position(|segment| *segment == "scheduler_affinity")?;
if parts.get(start + 1).is_some_and(|segment| *segment == "v2") {
return parse_admin_monitoring_scheduler_affinity_v2_key(&parts, start);
}
let affinity_key = parts.get(start + 1)?.trim();
if affinity_key.is_empty() {
return None;
}
let remaining = parts.get(start + 2..)?;
let (api_format, model_name) = split_admin_monitoring_api_format_and_model(
remaining,
AdminMonitoringAffinityKeyKind::Scheduler,
)?;
Some(ParsedAdminMonitoringAffinityKey {
affinity_key: affinity_key.to_string(),
api_format,
model_name,
client_family: None,
session_hash: None,
})
}
fn parse_admin_monitoring_scheduler_affinity_v2_key(
parts: &[&str],
start: usize,
) -> Option<ParsedAdminMonitoringAffinityKey> {
let affinity_key = parts.get(start + 2)?.trim();
if affinity_key.is_empty() {
return None;
}
let remaining = parts.get(start + 3..)?;
if remaining.len() < 4 {
return None;
}
let session_hash = remaining.last()?.trim();
let client_family = remaining.get(remaining.len().saturating_sub(2))?.trim();
if session_hash.is_empty() || client_family.is_empty() {
return None;
}
let api_model_segments = &remaining[..remaining.len() - 2];
let (api_format, model_name) = split_admin_monitoring_api_format_and_model(
api_model_segments,
AdminMonitoringAffinityKeyKind::Scheduler,
)?;
Some(ParsedAdminMonitoringAffinityKey {
affinity_key: affinity_key.to_string(),
api_format,
model_name,
client_family: Some(client_family.to_ascii_lowercase()),
session_hash: Some(session_hash.to_string()),
})
} }
pub(super) fn admin_monitoring_scheduler_affinity_cache_key( pub(super) fn admin_monitoring_scheduler_affinity_cache_key(
@@ -78,9 +199,47 @@ pub(super) fn admin_monitoring_scheduler_affinity_cache_key(
if affinity_key.is_empty() || api_format.is_empty() || model_name.is_empty() { if affinity_key.is_empty() || api_format.is_empty() || model_name.is_empty() {
return None; return None;
} }
Some(format!( match (
"scheduler_affinity:{affinity_key}:{api_format}:{model_name}" record
)) .client_family
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty()),
record
.session_hash
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty()),
) {
(Some(client_family), Some(session_hash)) => Some(format!(
"scheduler_affinity:v2:{affinity_key}:{api_format}:{model_name}:{client_family}:{session_hash}"
)),
_ => Some(format!(
"scheduler_affinity:{affinity_key}:{api_format}:{model_name}"
)),
}
}
fn admin_monitoring_json_string_field(
object: &serde_json::Map<String, serde_json::Value>,
field: &str,
) -> Option<String> {
object
.get(field)
.and_then(serde_json::Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
fn admin_monitoring_json_request_count(
object: &serde_json::Map<String, serde_json::Value>,
) -> Option<u64> {
object.get("request_count").and_then(|value| {
value
.as_u64()
.or_else(|| value.as_i64().and_then(|number| u64::try_from(number).ok()))
})
} }
pub(super) fn admin_monitoring_cache_affinity_record( pub(super) fn admin_monitoring_cache_affinity_record(
@@ -89,35 +248,22 @@ pub(super) fn admin_monitoring_cache_affinity_record(
) -> Option<AdminMonitoringCacheAffinityRecord> { ) -> Option<AdminMonitoringCacheAffinityRecord> {
let payload = serde_json::from_str::<serde_json::Value>(raw_value).ok()?; let payload = serde_json::from_str::<serde_json::Value>(raw_value).ok()?;
let object = payload.as_object()?; let object = payload.as_object()?;
let (affinity_key, parsed_api_format, parsed_model_name) = let parsed = parse_admin_monitoring_cache_affinity_key(raw_key)?;
parse_admin_monitoring_cache_affinity_key(raw_key)?; let api_format = admin_monitoring_json_string_field(object, "api_format")
let api_format = object .unwrap_or_else(|| parsed.api_format.clone());
.get("api_format") let model_name = admin_monitoring_json_string_field(object, "model_name")
.and_then(serde_json::Value::as_str) .unwrap_or_else(|| parsed.model_name.clone());
.map(str::trim) let request_count = admin_monitoring_json_request_count(object);
.filter(|value| !value.is_empty())
.unwrap_or(parsed_api_format.as_str())
.to_string();
let model_name = object
.get("model_name")
.and_then(serde_json::Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.unwrap_or(parsed_model_name.as_str())
.to_string();
let request_count = object
.get("request_count")
.and_then(|value| {
value
.as_u64()
.or_else(|| value.as_i64().and_then(|number| u64::try_from(number).ok()))
})
.unwrap_or(0);
Some(AdminMonitoringCacheAffinityRecord { Some(AdminMonitoringCacheAffinityRecord {
raw_key: raw_key.to_string(), raw_key: raw_key.to_string(),
affinity_key, affinity_key: parsed.affinity_key,
api_format, api_format,
model_name, model_name,
client_family: admin_monitoring_json_string_field(object, "client_family")
.map(|value| value.to_ascii_lowercase())
.or(parsed.client_family),
session_hash: admin_monitoring_json_string_field(object, "session_hash")
.or(parsed.session_hash),
provider_id: object provider_id: object
.get("provider_id") .get("provider_id")
.and_then(serde_json::Value::as_str) .and_then(serde_json::Value::as_str)
@@ -132,7 +278,8 @@ pub(super) fn admin_monitoring_cache_affinity_record(
.map(ToOwned::to_owned), .map(ToOwned::to_owned),
created_at: object.get("created_at").cloned(), created_at: object.get("created_at").cloned(),
expire_at: object.get("expire_at").cloned(), expire_at: object.get("expire_at").cloned(),
request_count, request_count: request_count.unwrap_or(0),
request_count_known: request_count.is_some(),
}) })
} }
@@ -143,23 +290,25 @@ pub(super) fn admin_monitoring_scheduler_affinity_record(
ttl: Duration, ttl: Duration,
now_unix_secs: u64, now_unix_secs: u64,
) -> Option<AdminMonitoringCacheAffinityRecord> { ) -> Option<AdminMonitoringCacheAffinityRecord> {
let (affinity_key, api_format, model_name) = let parsed = parse_admin_monitoring_scheduler_affinity_key(cache_key)?;
parse_admin_monitoring_scheduler_affinity_key(cache_key)?;
let age_secs = age.as_secs(); let age_secs = age.as_secs();
let created_at = now_unix_secs.saturating_sub(age_secs); let created_at = now_unix_secs.saturating_sub(age_secs);
let expire_at = created_at.saturating_add(ttl.as_secs()); let expire_at = created_at.saturating_add(ttl.as_secs());
Some(AdminMonitoringCacheAffinityRecord { Some(AdminMonitoringCacheAffinityRecord {
raw_key: cache_key.to_string(), raw_key: cache_key.to_string(),
affinity_key, affinity_key: parsed.affinity_key,
api_format, api_format: parsed.api_format,
model_name, model_name: parsed.model_name,
client_family: parsed.client_family,
session_hash: parsed.session_hash,
provider_id: Some(target.provider_id.clone()), provider_id: Some(target.provider_id.clone()),
endpoint_id: Some(target.endpoint_id.clone()), endpoint_id: Some(target.endpoint_id.clone()),
key_id: Some(target.key_id.clone()), key_id: Some(target.key_id.clone()),
created_at: Some(serde_json::json!(created_at)), created_at: Some(serde_json::json!(created_at)),
expire_at: Some(serde_json::json!(expire_at)), expire_at: Some(serde_json::json!(expire_at)),
request_count: 0, request_count: 0,
request_count_known: false,
}) })
} }
@@ -169,36 +318,23 @@ pub(super) fn admin_monitoring_scheduler_affinity_record_from_raw(
) -> Option<AdminMonitoringCacheAffinityRecord> { ) -> Option<AdminMonitoringCacheAffinityRecord> {
let payload = serde_json::from_str::<serde_json::Value>(raw_value).ok()?; let payload = serde_json::from_str::<serde_json::Value>(raw_value).ok()?;
let object = payload.as_object()?; let object = payload.as_object()?;
let (affinity_key, parsed_api_format, parsed_model_name) = let parsed = parse_admin_monitoring_scheduler_affinity_key(raw_key)?;
parse_admin_monitoring_scheduler_affinity_key(raw_key)?; let api_format = admin_monitoring_json_string_field(object, "api_format")
let api_format = object .unwrap_or_else(|| parsed.api_format.clone());
.get("api_format") let model_name = admin_monitoring_json_string_field(object, "model_name")
.and_then(serde_json::Value::as_str) .unwrap_or_else(|| parsed.model_name.clone());
.map(str::trim) let request_count = admin_monitoring_json_request_count(object);
.filter(|value| !value.is_empty())
.unwrap_or(parsed_api_format.as_str())
.to_string();
let model_name = object
.get("model_name")
.and_then(serde_json::Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.unwrap_or(parsed_model_name.as_str())
.to_string();
let request_count = object
.get("request_count")
.and_then(|value| {
value
.as_u64()
.or_else(|| value.as_i64().and_then(|number| u64::try_from(number).ok()))
})
.unwrap_or(0);
Some(AdminMonitoringCacheAffinityRecord { Some(AdminMonitoringCacheAffinityRecord {
raw_key: raw_key.to_string(), raw_key: raw_key.to_string(),
affinity_key, affinity_key: parsed.affinity_key,
api_format, api_format,
model_name, model_name,
client_family: admin_monitoring_json_string_field(object, "client_family")
.map(|value| value.to_ascii_lowercase())
.or(parsed.client_family),
session_hash: admin_monitoring_json_string_field(object, "session_hash")
.or(parsed.session_hash),
provider_id: object provider_id: object
.get("provider_id") .get("provider_id")
.and_then(serde_json::Value::as_str) .and_then(serde_json::Value::as_str)
@@ -213,7 +349,8 @@ pub(super) fn admin_monitoring_scheduler_affinity_record_from_raw(
.map(ToOwned::to_owned), .map(ToOwned::to_owned),
created_at: object.get("created_at").cloned(), created_at: object.get("created_at").cloned(),
expire_at: object.get("expire_at").cloned(), expire_at: object.get("expire_at").cloned(),
request_count, request_count: request_count.unwrap_or(0),
request_count_known: request_count.is_some(),
}) })
} }
@@ -227,10 +364,15 @@ pub(super) fn clear_admin_monitoring_scheduler_affinity_entries(
state: &AdminAppState<'_>, state: &AdminAppState<'_>,
records: &[AdminMonitoringCacheAffinityRecord], records: &[AdminMonitoringCacheAffinityRecord],
) { ) {
let scheduler_keys = records let mut scheduler_keys = std::collections::BTreeSet::new();
.iter() for record in records {
.filter_map(admin_monitoring_scheduler_affinity_cache_key) if record.raw_key.contains("scheduler_affinity:") {
.collect::<std::collections::BTreeSet<_>>(); scheduler_keys.insert(record.raw_key.clone());
}
if let Some(cache_key) = admin_monitoring_scheduler_affinity_cache_key(record) {
scheduler_keys.insert(cache_key);
}
}
for scheduler_key in scheduler_keys { for scheduler_key in scheduler_keys {
let _ = state let _ = state
.as_ref() .as_ref()
@@ -238,6 +380,135 @@ pub(super) fn clear_admin_monitoring_scheduler_affinity_entries(
} }
} }
#[cfg(test)]
mod tests {
use super::{
admin_monitoring_scheduler_affinity_record_from_raw,
parse_admin_monitoring_cache_affinity_key, parse_admin_monitoring_scheduler_affinity_key,
};
use serde_json::json;
#[test]
fn parses_legacy_scheduler_affinity_key_with_multi_segment_api_format() {
let parsed = parse_admin_monitoring_scheduler_affinity_key(
"scheduler_affinity:user-key-1:openai:chat:gpt-4.1",
)
.expect("legacy scheduler key should parse");
assert_eq!(parsed.affinity_key, "user-key-1");
assert_eq!(parsed.api_format, "openai:chat");
assert_eq!(parsed.model_name, "gpt-4.1");
assert_eq!(parsed.client_family, None);
assert_eq!(parsed.session_hash, None);
}
#[test]
fn parses_v2_scheduler_affinity_key_with_client_family() {
let parsed = parse_admin_monitoring_scheduler_affinity_key(
"scheduler_affinity:v2:user-key-1:openai:responses:gpt-5.5:codex:d35efdccd572e9c5430e17d663dfde4cce83ea7e6665ac332f565780c98b1dff",
)
.expect("v2 scheduler key should parse");
assert_eq!(parsed.affinity_key, "user-key-1");
assert_eq!(parsed.api_format, "openai:responses");
assert_eq!(parsed.model_name, "gpt-5.5");
assert_eq!(parsed.client_family.as_deref(), Some("codex"));
assert_eq!(
parsed.session_hash.as_deref(),
Some("d35efdccd572e9c5430e17d663dfde4cce83ea7e6665ac332f565780c98b1dff")
);
}
#[test]
fn parses_v2_scheduler_affinity_key_with_three_segment_api_format() {
let parsed = parse_admin_monitoring_scheduler_affinity_key(
"scheduler_affinity:v2:user-key-1:openai:responses:compact:gpt-5.5:opencode:sessionhash",
)
.expect("compact v2 scheduler key should parse");
assert_eq!(parsed.affinity_key, "user-key-1");
assert_eq!(parsed.api_format, "openai:responses:compact");
assert_eq!(parsed.model_name, "gpt-5.5");
assert_eq!(parsed.client_family.as_deref(), Some("opencode"));
assert_eq!(parsed.session_hash.as_deref(), Some("sessionhash"));
}
#[test]
fn parses_scheduler_affinity_key_with_unregistered_two_segment_api_format() {
let parsed = parse_admin_monitoring_scheduler_affinity_key(
"scheduler_affinity:v2:user-key-1:openai:image:gpt-image-1:opencode:sessionhash",
)
.expect("image v2 scheduler key should parse");
assert_eq!(parsed.affinity_key, "user-key-1");
assert_eq!(parsed.api_format, "openai:image");
assert_eq!(parsed.model_name, "gpt-image-1");
assert_eq!(parsed.client_family.as_deref(), Some("opencode"));
assert_eq!(parsed.session_hash.as_deref(), Some("sessionhash"));
}
#[test]
fn keeps_legacy_cache_affinity_single_segment_api_format() {
let parsed = parse_admin_monitoring_cache_affinity_key(
"cache_affinity:user-key-1:openai:model-alpha",
)
.expect("cache affinity key should parse");
assert_eq!(parsed.affinity_key, "user-key-1");
assert_eq!(parsed.api_format, "openai");
assert_eq!(parsed.model_name, "model-alpha");
}
#[test]
fn scheduler_raw_record_exposes_client_family_and_known_count() {
let raw_value = json!({
"provider_id": "provider-1",
"endpoint_id": "endpoint-1",
"key_id": "provider-key-1",
"created_at": 1710000000u64,
"expire_at": 1710000300u64,
"request_count": 3u64,
})
.to_string();
let record = admin_monitoring_scheduler_affinity_record_from_raw(
"scheduler_affinity:v2:user-key-1:openai:responses:gpt-5.5:codex:sessionhash",
&raw_value,
)
.expect("raw scheduler affinity should parse");
assert_eq!(record.affinity_key, "user-key-1");
assert_eq!(record.api_format, "openai:responses");
assert_eq!(record.model_name, "gpt-5.5");
assert_eq!(record.client_family.as_deref(), Some("codex"));
assert_eq!(record.session_hash.as_deref(), Some("sessionhash"));
assert_eq!(record.request_count, 3);
assert!(record.request_count_known);
}
#[test]
fn scheduler_raw_record_treats_zero_request_count_as_known() {
let raw_value = json!({
"provider_id": "provider-1",
"endpoint_id": "endpoint-1",
"key_id": "provider-key-1",
"created_at": 1710000000u64,
"expire_at": 1710000300u64,
"request_count": 0u64,
})
.to_string();
let record = admin_monitoring_scheduler_affinity_record_from_raw(
"scheduler_affinity:v2:user-key-1:openai:responses:gpt-5.5:codex:sessionhash",
&raw_value,
)
.expect("raw scheduler affinity should parse");
assert_eq!(record.request_count, 0);
assert!(record.request_count_known);
}
}
#[cfg(test)] #[cfg(test)]
pub(super) fn delete_admin_monitoring_cache_affinity_entries_for_tests( pub(super) fn delete_admin_monitoring_cache_affinity_entries_for_tests(
state: &AdminAppState<'_>, state: &AdminAppState<'_>,

View File

@@ -216,6 +216,9 @@ async fn list_admin_monitoring_cache_affinity_records_matching(
runner runner
.keyspace() .keyspace()
.key(&format!("scheduler_affinity:{affinity_key}:*")), .key(&format!("scheduler_affinity:{affinity_key}:*")),
runner
.keyspace()
.key(&format!("scheduler_affinity:v2:{affinity_key}:*")),
] ]
}) })
.collect::<Vec<_>>() .collect::<Vec<_>>()

View File

@@ -4,12 +4,15 @@ pub(super) struct AdminMonitoringCacheAffinityRecord {
pub(super) affinity_key: String, pub(super) affinity_key: String,
pub(super) api_format: String, pub(super) api_format: String,
pub(super) model_name: String, pub(super) model_name: String,
pub(super) client_family: Option<String>,
pub(super) session_hash: Option<String>,
pub(super) provider_id: Option<String>, pub(super) provider_id: Option<String>,
pub(super) endpoint_id: Option<String>, pub(super) endpoint_id: Option<String>,
pub(super) key_id: Option<String>, pub(super) key_id: Option<String>,
pub(super) created_at: Option<serde_json::Value>, pub(super) created_at: Option<serde_json::Value>,
pub(super) expire_at: Option<serde_json::Value>, pub(super) expire_at: Option<serde_json::Value>,
pub(super) request_count: u64, pub(super) request_count: u64,
pub(super) request_count_known: bool,
} }
pub(super) struct AdminMonitoringCacheSnapshot { pub(super) struct AdminMonitoringCacheSnapshot {

View File

@@ -67,6 +67,7 @@ impl AppState {
}; };
let cache_key = cache_key.to_string(); let cache_key = cache_key.to_string();
let namespaced_cache_key = runner.keyspace().key(&cache_key);
let provider_id = target.provider_id.clone(); let provider_id = target.provider_id.clone();
let endpoint_id = target.endpoint_id.clone(); let endpoint_id = target.endpoint_id.clone();
let key_id = target.key_id.clone(); let key_id = target.key_id.clone();
@@ -75,19 +76,48 @@ impl AppState {
let expire_at = now_unix_secs.saturating_add(ttl_seconds); let expire_at = now_unix_secs.saturating_add(ttl_seconds);
handle.spawn(async move { handle.spawn(async move {
let payload = serde_json::json!({ let Ok(mut connection) = runner.client().get_multiplexed_async_connection().await
"provider_id": provider_id, else {
"endpoint_id": endpoint_id,
"key_id": key_id,
"created_at": now_unix_secs,
"expire_at": expire_at,
"request_count": 0,
});
let Ok(serialized) = serde_json::to_string(&payload) else {
return; return;
}; };
let _ = runner let script = r#"
.setex(&cache_key, &serialized, Some(ttl_seconds)) local existing = redis.call('GET', KEYS[1])
local request_count = 0
local created_at = tonumber(ARGV[4])
if existing then
local ok, payload = pcall(cjson.decode, existing)
if ok and type(payload) == 'table' then
if type(payload['request_count']) == 'number' then
request_count = payload['request_count']
end
if type(payload['created_at']) == 'number' then
created_at = payload['created_at']
end
end
end
request_count = request_count + 1
local payload = {
provider_id = ARGV[1],
endpoint_id = ARGV[2],
key_id = ARGV[3],
created_at = created_at,
expire_at = tonumber(ARGV[5]),
request_count = request_count
}
redis.call('SETEX', KEYS[1], tonumber(ARGV[6]), cjson.encode(payload))
return request_count
"#;
let _ = redis::cmd("EVAL")
.arg(script)
.arg(1)
.arg(&namespaced_cache_key)
.arg(&provider_id)
.arg(&endpoint_id)
.arg(&key_id)
.arg(now_unix_secs)
.arg(expire_at)
.arg(ttl_seconds)
.query_async::<i64>(&mut connection)
.await; .await;
}); });
} }