mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 09:20:22 +08:00
Add session-scoped affinity deletion
This commit is contained in:
@@ -197,10 +197,13 @@ pub(super) async fn build_admin_monitoring_cache_affinities_response(
|
|||||||
"global_model_id": affinity.model_name,
|
"global_model_id": affinity.model_name,
|
||||||
"model_name": affinity.model_name,
|
"model_name": affinity.model_name,
|
||||||
"model_display_name": serde_json::Value::Null,
|
"model_display_name": serde_json::Value::Null,
|
||||||
|
"client_family": affinity.client_family,
|
||||||
|
"session_hash": affinity.session_hash,
|
||||||
"api_format": affinity.api_format,
|
"api_format": affinity.api_format,
|
||||||
"created_at": affinity.created_at,
|
"created_at": affinity.created_at,
|
||||||
"expire_at": affinity.expire_at,
|
"expire_at": affinity.expire_at,
|
||||||
"request_count": affinity.request_count,
|
"request_count": affinity.request_count,
|
||||||
|
"request_count_known": affinity.request_count_known,
|
||||||
}));
|
}));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -318,9 +321,12 @@ pub(super) async fn build_admin_monitoring_cache_affinity_response(
|
|||||||
"key_id": item.key_id,
|
"key_id": item.key_id,
|
||||||
"api_format": item.api_format,
|
"api_format": item.api_format,
|
||||||
"model_name": item.model_name,
|
"model_name": item.model_name,
|
||||||
|
"client_family": item.client_family,
|
||||||
|
"session_hash": item.session_hash,
|
||||||
"created_at": item.created_at,
|
"created_at": item.created_at,
|
||||||
"expire_at": item.expire_at,
|
"expire_at": item.expire_at,
|
||||||
"request_count": item.request_count,
|
"request_count": item.request_count,
|
||||||
|
"request_count_known": item.request_count_known,
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
.collect::<Vec<_>>();
|
.collect::<Vec<_>>();
|
||||||
|
|||||||
@@ -20,6 +20,50 @@ use aether_admin::observability::monitoring::{
|
|||||||
};
|
};
|
||||||
use axum::{body::Body, response::Response};
|
use axum::{body::Body, response::Response};
|
||||||
|
|
||||||
|
#[derive(Debug, Default)]
|
||||||
|
struct AdminMonitoringCacheAffinityDeleteFilter {
|
||||||
|
client_family: Option<String>,
|
||||||
|
session_hash: Option<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
fn admin_monitoring_cache_affinity_delete_filter_from_query(
|
||||||
|
query: Option<&str>,
|
||||||
|
) -> AdminMonitoringCacheAffinityDeleteFilter {
|
||||||
|
let Some(query) = query else {
|
||||||
|
return AdminMonitoringCacheAffinityDeleteFilter::default();
|
||||||
|
};
|
||||||
|
let mut filter = AdminMonitoringCacheAffinityDeleteFilter::default();
|
||||||
|
for (key, value) in url::form_urlencoded::parse(query.as_bytes()) {
|
||||||
|
let value = value.trim();
|
||||||
|
if value.is_empty() {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
match key.as_ref() {
|
||||||
|
"client_family" => filter.client_family = Some(value.to_ascii_lowercase()),
|
||||||
|
"session_hash" => filter.session_hash = Some(value.to_string()),
|
||||||
|
_ => {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
filter
|
||||||
|
}
|
||||||
|
|
||||||
|
fn admin_monitoring_cache_affinity_matches_delete_filter(
|
||||||
|
item: &super::super::cache_types::AdminMonitoringCacheAffinityRecord,
|
||||||
|
filter: &AdminMonitoringCacheAffinityDeleteFilter,
|
||||||
|
) -> bool {
|
||||||
|
if let Some(client_family) = filter.client_family.as_deref() {
|
||||||
|
if item.client_family.as_deref() != Some(client_family) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if let Some(session_hash) = filter.session_hash.as_deref() {
|
||||||
|
if item.session_hash.as_deref() != Some(session_hash) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
true
|
||||||
|
}
|
||||||
|
|
||||||
pub(in super::super) async fn build_admin_monitoring_cache_affinity_delete_response(
|
pub(in super::super) async fn build_admin_monitoring_cache_affinity_delete_response(
|
||||||
state: &AdminAppState<'_>,
|
state: &AdminAppState<'_>,
|
||||||
request_context: &AdminRequestContext<'_>,
|
request_context: &AdminRequestContext<'_>,
|
||||||
@@ -41,31 +85,33 @@ pub(in super::super) async fn build_admin_monitoring_cache_affinity_delete_respo
|
|||||||
|
|
||||||
let target_affinity_keys =
|
let target_affinity_keys =
|
||||||
std::iter::once(affinity_key.clone()).collect::<std::collections::BTreeSet<_>>();
|
std::iter::once(affinity_key.clone()).collect::<std::collections::BTreeSet<_>>();
|
||||||
let target_affinity =
|
let delete_filter = admin_monitoring_cache_affinity_delete_filter_from_query(
|
||||||
|
request_context.request_query_string.as_deref(),
|
||||||
|
);
|
||||||
|
let target_affinities =
|
||||||
list_admin_monitoring_cache_affinity_records_by_affinity_keys(state, &target_affinity_keys)
|
list_admin_monitoring_cache_affinity_records_by_affinity_keys(state, &target_affinity_keys)
|
||||||
.await?
|
.await?
|
||||||
.into_iter()
|
.into_iter()
|
||||||
.find(|item| {
|
.filter(|item| {
|
||||||
item.affinity_key == affinity_key
|
item.affinity_key == affinity_key
|
||||||
&& item.endpoint_id.as_deref() == Some(endpoint_id.as_str())
|
&& item.endpoint_id.as_deref() == Some(endpoint_id.as_str())
|
||||||
&& item.model_name == model_id
|
&& item.model_name == model_id
|
||||||
&& item.api_format.eq_ignore_ascii_case(&api_format)
|
&& item.api_format.eq_ignore_ascii_case(&api_format)
|
||||||
});
|
&& admin_monitoring_cache_affinity_matches_delete_filter(item, &delete_filter)
|
||||||
let Some(target_affinity) = target_affinity else {
|
})
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
if target_affinities.is_empty() {
|
||||||
return Ok(admin_monitoring_not_found_response(
|
return Ok(admin_monitoring_not_found_response(
|
||||||
"未找到指定的缓存亲和性记录",
|
"未找到指定的缓存亲和性记录",
|
||||||
));
|
));
|
||||||
};
|
}
|
||||||
|
|
||||||
let _ = delete_admin_monitoring_cache_affinity_raw_keys(
|
let raw_keys = target_affinities
|
||||||
state,
|
.iter()
|
||||||
std::slice::from_ref(&target_affinity.raw_key),
|
.map(|item| item.raw_key.clone())
|
||||||
)
|
.collect::<Vec<_>>();
|
||||||
.await?;
|
let _ = delete_admin_monitoring_cache_affinity_raw_keys(state, &raw_keys).await?;
|
||||||
clear_admin_monitoring_scheduler_affinity_entries(
|
clear_admin_monitoring_scheduler_affinity_entries(state, &target_affinities);
|
||||||
state,
|
|
||||||
std::slice::from_ref(&target_affinity),
|
|
||||||
);
|
|
||||||
|
|
||||||
let mut api_key_by_id = admin_monitoring_list_export_api_key_records_by_ids(
|
let mut api_key_by_id = admin_monitoring_list_export_api_key_records_by_ids(
|
||||||
state,
|
state,
|
||||||
|
|||||||
@@ -339,6 +339,146 @@ async fn admin_monitoring_cache_affinities_and_delete_use_runtime_scheduler_affi
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn admin_monitoring_cache_affinities_parse_session_scoped_scheduler_affinity_cache() {
|
||||||
|
let user_repository = Arc::new(
|
||||||
|
InMemoryUserReadRepository::seed_auth_users(vec![sample_monitoring_auth_user("user-1")])
|
||||||
|
.with_export_users(vec![sample_monitoring_export_user("user-1")]),
|
||||||
|
);
|
||||||
|
let auth_repository = Arc::new(
|
||||||
|
InMemoryAuthApiKeySnapshotRepository::default().with_export_records(vec![
|
||||||
|
sample_monitoring_export_api_key("user-1", "user-key-1"),
|
||||||
|
]),
|
||||||
|
);
|
||||||
|
let provider_catalog = Arc::new(InMemoryProviderCatalogReadRepository::seed(
|
||||||
|
vec![sample_provider()],
|
||||||
|
vec![sample_monitoring_catalog_endpoint()],
|
||||||
|
vec![sample_monitoring_catalog_key()],
|
||||||
|
));
|
||||||
|
let state = AppState::new()
|
||||||
|
.expect("state should build")
|
||||||
|
.with_data_state_for_tests(
|
||||||
|
crate::data::GatewayDataState::with_provider_catalog_reader_for_tests(provider_catalog)
|
||||||
|
.with_user_reader(user_repository)
|
||||||
|
.with_auth_api_key_reader(auth_repository),
|
||||||
|
);
|
||||||
|
let client_session = aether_scheduler_core::ClientSessionAffinity::new(
|
||||||
|
Some("Codex".to_string()),
|
||||||
|
Some("account=acct-1;session=session-1".to_string()),
|
||||||
|
);
|
||||||
|
let other_client_session = aether_scheduler_core::ClientSessionAffinity::new(
|
||||||
|
Some("Codex".to_string()),
|
||||||
|
Some("account=acct-1;session=session-2".to_string()),
|
||||||
|
);
|
||||||
|
let affinity_cache_key =
|
||||||
|
aether_scheduler_core::build_scheduler_affinity_cache_key_for_api_key_id_with_client_session(
|
||||||
|
"user-key-1",
|
||||||
|
"openai:responses",
|
||||||
|
"gpt-5.5",
|
||||||
|
Some(&client_session),
|
||||||
|
)
|
||||||
|
.expect("session scheduler affinity cache key should build");
|
||||||
|
let other_affinity_cache_key =
|
||||||
|
aether_scheduler_core::build_scheduler_affinity_cache_key_for_api_key_id_with_client_session(
|
||||||
|
"user-key-1",
|
||||||
|
"openai:responses",
|
||||||
|
"gpt-5.5",
|
||||||
|
Some(&other_client_session),
|
||||||
|
)
|
||||||
|
.expect("other session scheduler affinity cache key should build");
|
||||||
|
assert!(affinity_cache_key
|
||||||
|
.starts_with("scheduler_affinity:v2:user-key-1:openai:responses:gpt-5.5:codex:"));
|
||||||
|
let session_hash = affinity_cache_key
|
||||||
|
.rsplit(':')
|
||||||
|
.next()
|
||||||
|
.expect("session hash should exist")
|
||||||
|
.to_string();
|
||||||
|
state.scheduler_affinity_cache.insert(
|
||||||
|
affinity_cache_key.clone(),
|
||||||
|
crate::cache::SchedulerAffinityTarget {
|
||||||
|
provider_id: "provider-1".to_string(),
|
||||||
|
endpoint_id: "endpoint-1".to_string(),
|
||||||
|
key_id: "provider-key-1".to_string(),
|
||||||
|
},
|
||||||
|
crate::scheduler::affinity::SCHEDULER_AFFINITY_TTL,
|
||||||
|
128,
|
||||||
|
);
|
||||||
|
state.scheduler_affinity_cache.insert(
|
||||||
|
other_affinity_cache_key.clone(),
|
||||||
|
crate::cache::SchedulerAffinityTarget {
|
||||||
|
provider_id: "provider-1".to_string(),
|
||||||
|
endpoint_id: "endpoint-1".to_string(),
|
||||||
|
key_id: "provider-key-1".to_string(),
|
||||||
|
},
|
||||||
|
crate::scheduler::affinity::SCHEDULER_AFFINITY_TTL,
|
||||||
|
128,
|
||||||
|
);
|
||||||
|
|
||||||
|
let list_response = local_monitoring_response(
|
||||||
|
&state,
|
||||||
|
&request_context(
|
||||||
|
http::Method::GET,
|
||||||
|
"/api/admin/monitoring/cache/affinities?keyword=alice&limit=20&offset=0",
|
||||||
|
),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("handler should not error")
|
||||||
|
.expect("route should be handled locally");
|
||||||
|
assert_eq!(list_response.status(), http::StatusCode::OK);
|
||||||
|
let list_body = to_bytes(list_response.into_body(), usize::MAX)
|
||||||
|
.await
|
||||||
|
.expect("body should read");
|
||||||
|
let list_payload: serde_json::Value =
|
||||||
|
serde_json::from_slice(&list_body).expect("json body should parse");
|
||||||
|
let items = list_payload["data"]["items"]
|
||||||
|
.as_array()
|
||||||
|
.expect("items should be an array");
|
||||||
|
let item = items
|
||||||
|
.iter()
|
||||||
|
.find(|item| item["session_hash"] == json!(session_hash))
|
||||||
|
.expect("session-scoped item should be listed");
|
||||||
|
assert_eq!(list_payload["data"]["meta"]["total"], json!(2));
|
||||||
|
assert_eq!(item["affinity_key"], json!("user-key-1"));
|
||||||
|
assert_eq!(item["username"], json!("alice"));
|
||||||
|
assert_eq!(item["api_format"], json!("openai:responses"));
|
||||||
|
assert_eq!(item["model_name"], json!("gpt-5.5"));
|
||||||
|
assert_eq!(item["client_family"], json!("codex"));
|
||||||
|
assert_eq!(item["provider_name"], json!("OpenAI"));
|
||||||
|
assert_eq!(item["key_name"], json!("prod-key"));
|
||||||
|
assert_eq!(item["request_count"], json!(0));
|
||||||
|
assert_eq!(item["request_count_known"], json!(false));
|
||||||
|
|
||||||
|
let delete_response = local_monitoring_response(
|
||||||
|
&state,
|
||||||
|
&request_context(
|
||||||
|
http::Method::DELETE,
|
||||||
|
&format!(
|
||||||
|
"/api/admin/monitoring/cache/affinity/user-key-1/endpoint-1/gpt-5.5/openai:responses?client_family=codex&session_hash={session_hash}"
|
||||||
|
),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("handler should not error")
|
||||||
|
.expect("route should be handled locally");
|
||||||
|
assert_eq!(delete_response.status(), http::StatusCode::OK);
|
||||||
|
assert_eq!(
|
||||||
|
state.read_scheduler_affinity_target(
|
||||||
|
&affinity_cache_key,
|
||||||
|
crate::scheduler::affinity::SCHEDULER_AFFINITY_TTL,
|
||||||
|
),
|
||||||
|
None
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
state
|
||||||
|
.read_scheduler_affinity_target(
|
||||||
|
&other_affinity_cache_key,
|
||||||
|
crate::scheduler::affinity::SCHEDULER_AFFINITY_TTL,
|
||||||
|
)
|
||||||
|
.is_some(),
|
||||||
|
"deleting one session-scoped row should keep sibling sessions"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn admin_monitoring_cache_users_delete_returns_local_payload_from_test_store() {
|
async fn admin_monitoring_cache_users_delete_returns_local_payload_from_test_store() {
|
||||||
let user_repository = Arc::new(
|
let user_repository = Arc::new(
|
||||||
|
|||||||
Reference in New Issue
Block a user