mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 01:10:23 +08:00
fix(gateway): preserve affinity during adaptive retries
This commit is contained in:
@@ -360,7 +360,16 @@ async fn record_attempt_failure_effect(
|
|||||||
}
|
}
|
||||||
|
|
||||||
if let Some(cache_key) = local_scheduler_affinity_cache_key(context.report_context) {
|
if let Some(cache_key) = local_scheduler_affinity_cache_key(context.report_context) {
|
||||||
let _ = state.remove_scheduler_affinity_cache_entry(&cache_key);
|
let Some(failed_target) = local_scheduler_affinity_target(context.plan) else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
if state
|
||||||
|
.read_scheduler_affinity_target(&cache_key, SCHEDULER_AFFINITY_TTL)
|
||||||
|
.as_ref()
|
||||||
|
== Some(&failed_target)
|
||||||
|
{
|
||||||
|
let _ = state.remove_scheduler_affinity_cache_entry(&cache_key);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -448,7 +457,10 @@ async fn record_adaptive_rate_limit_effect(
|
|||||||
updated_key.status_snapshot = Some(projection.status_snapshot);
|
updated_key.status_snapshot = Some(projection.status_snapshot);
|
||||||
updated_key.updated_at_unix_secs = Some(observed_at_unix_secs);
|
updated_key.updated_at_unix_secs = Some(observed_at_unix_secs);
|
||||||
|
|
||||||
if let Err(err) = state.update_provider_catalog_key(&updated_key).await {
|
if let Err(err) = state
|
||||||
|
.update_provider_catalog_key_runtime_state(&updated_key)
|
||||||
|
.await
|
||||||
|
{
|
||||||
warn!(
|
warn!(
|
||||||
"gateway orchestration effects: failed to persist adaptive rate-limit projection for provider {} endpoint {} key {}: {:?}",
|
"gateway orchestration effects: failed to persist adaptive rate-limit projection for provider {} endpoint {} key {}: {:?}",
|
||||||
context.plan.provider_id, context.plan.endpoint_id, context.plan.key_id, err
|
context.plan.provider_id, context.plan.endpoint_id, context.plan.key_id, err
|
||||||
@@ -496,7 +508,10 @@ async fn record_adaptive_success_effect(
|
|||||||
updated_key.status_snapshot = Some(projection.status_snapshot);
|
updated_key.status_snapshot = Some(projection.status_snapshot);
|
||||||
updated_key.updated_at_unix_secs = Some(observed_at_unix_secs);
|
updated_key.updated_at_unix_secs = Some(observed_at_unix_secs);
|
||||||
|
|
||||||
if let Err(err) = state.update_provider_catalog_key(&updated_key).await {
|
if let Err(err) = state
|
||||||
|
.update_provider_catalog_key_runtime_state(&updated_key)
|
||||||
|
.await
|
||||||
|
{
|
||||||
warn!(
|
warn!(
|
||||||
"gateway orchestration effects: failed to persist adaptive success projection for provider {} endpoint {} key {}: {:?}",
|
"gateway orchestration effects: failed to persist adaptive success projection for provider {} endpoint {} key {}: {:?}",
|
||||||
context.plan.provider_id, context.plan.endpoint_id, context.plan.key_id, err
|
context.plan.provider_id, context.plan.endpoint_id, context.plan.key_id, err
|
||||||
@@ -1429,6 +1444,50 @@ mod tests {
|
|||||||
.is_some());
|
.is_some());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn attempt_failure_keeps_scheduler_affinity_for_non_affinity_candidate() {
|
||||||
|
let state = AppState::new().expect("gateway state should build");
|
||||||
|
let plan = sample_plan();
|
||||||
|
let report_context = json!({
|
||||||
|
"api_key_id": "api-key-1",
|
||||||
|
"client_api_format": "openai:chat",
|
||||||
|
"model": "gpt-5",
|
||||||
|
});
|
||||||
|
let cache_key =
|
||||||
|
build_scheduler_affinity_cache_key_for_api_key_id("api-key-1", "openai:chat", "gpt-5")
|
||||||
|
.expect("scheduler affinity cache key should build");
|
||||||
|
let affinity_target = SchedulerAffinityTarget {
|
||||||
|
provider_id: "prov-2".to_string(),
|
||||||
|
endpoint_id: "ep-2".to_string(),
|
||||||
|
key_id: "key-2".to_string(),
|
||||||
|
};
|
||||||
|
|
||||||
|
state.remember_scheduler_affinity_target(
|
||||||
|
&cache_key,
|
||||||
|
affinity_target.clone(),
|
||||||
|
SCHEDULER_AFFINITY_TTL,
|
||||||
|
16,
|
||||||
|
);
|
||||||
|
|
||||||
|
apply_local_execution_effect(
|
||||||
|
&state,
|
||||||
|
LocalExecutionEffectContext {
|
||||||
|
plan: &plan,
|
||||||
|
report_context: Some(&report_context),
|
||||||
|
},
|
||||||
|
LocalExecutionEffect::AttemptFailure(LocalAttemptFailureEffect {
|
||||||
|
status_code: 524,
|
||||||
|
classification: LocalFailoverClassification::RetryUpstreamFailure,
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
state.read_scheduler_affinity_target(cache_key.as_str(), SCHEDULER_AFFINITY_TTL),
|
||||||
|
Some(affinity_target)
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn attempt_failure_keeps_scheduler_affinity_for_non_failure_status() {
|
async fn attempt_failure_keeps_scheduler_affinity_for_non_failure_status() {
|
||||||
let state = AppState::new().expect("gateway state should build");
|
let state = AppState::new().expect("gateway state should build");
|
||||||
@@ -2018,6 +2077,21 @@ mod tests {
|
|||||||
async fn adaptive_rate_limit_effect_updates_adaptive_key_observation() {
|
async fn adaptive_rate_limit_effect_updates_adaptive_key_observation() {
|
||||||
let state = adaptive_state();
|
let state = adaptive_state();
|
||||||
let plan = sample_plan();
|
let plan = sample_plan();
|
||||||
|
let cache_key =
|
||||||
|
build_scheduler_affinity_cache_key_for_api_key_id("api-key-1", "openai:chat", "gpt-5")
|
||||||
|
.expect("scheduler affinity cache key should build");
|
||||||
|
let target = SchedulerAffinityTarget {
|
||||||
|
provider_id: plan.provider_id.clone(),
|
||||||
|
endpoint_id: plan.endpoint_id.clone(),
|
||||||
|
key_id: plan.key_id.clone(),
|
||||||
|
};
|
||||||
|
state.remember_scheduler_affinity_target(
|
||||||
|
&cache_key,
|
||||||
|
target.clone(),
|
||||||
|
SCHEDULER_AFFINITY_TTL,
|
||||||
|
16,
|
||||||
|
);
|
||||||
|
let initial_epoch = state.scheduler_affinity_epoch();
|
||||||
|
|
||||||
apply_local_execution_effect(
|
apply_local_execution_effect(
|
||||||
&state,
|
&state,
|
||||||
@@ -2081,6 +2155,11 @@ mod tests {
|
|||||||
.and_then(|value| value.get("enforcement_active")),
|
.and_then(|value| value.get("enforcement_active")),
|
||||||
Some(&json!(false))
|
Some(&json!(false))
|
||||||
);
|
);
|
||||||
|
assert_eq!(state.scheduler_affinity_epoch(), initial_epoch);
|
||||||
|
assert_eq!(
|
||||||
|
state.read_scheduler_affinity_target(cache_key.as_str(), SCHEDULER_AFFINITY_TTL),
|
||||||
|
Some(target)
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
@@ -2211,6 +2290,21 @@ mod tests {
|
|||||||
.expect("request candidate should build")],
|
.expect("request candidate should build")],
|
||||||
);
|
);
|
||||||
let plan = sample_plan();
|
let plan = sample_plan();
|
||||||
|
let cache_key =
|
||||||
|
build_scheduler_affinity_cache_key_for_api_key_id("api-key-1", "openai:chat", "gpt-5")
|
||||||
|
.expect("scheduler affinity cache key should build");
|
||||||
|
let target = SchedulerAffinityTarget {
|
||||||
|
provider_id: plan.provider_id.clone(),
|
||||||
|
endpoint_id: plan.endpoint_id.clone(),
|
||||||
|
key_id: plan.key_id.clone(),
|
||||||
|
};
|
||||||
|
state.remember_scheduler_affinity_target(
|
||||||
|
&cache_key,
|
||||||
|
target.clone(),
|
||||||
|
SCHEDULER_AFFINITY_TTL,
|
||||||
|
16,
|
||||||
|
);
|
||||||
|
let initial_epoch = state.scheduler_affinity_epoch();
|
||||||
|
|
||||||
apply_local_execution_effect(
|
apply_local_execution_effect(
|
||||||
&state,
|
&state,
|
||||||
@@ -2242,5 +2336,10 @@ mod tests {
|
|||||||
.and_then(Value::as_str),
|
.and_then(Value::as_str),
|
||||||
Some("high_utilization")
|
Some("high_utilization")
|
||||||
);
|
);
|
||||||
|
assert_eq!(state.scheduler_affinity_epoch(), initial_epoch);
|
||||||
|
assert_eq!(
|
||||||
|
state.read_scheduler_affinity_target(cache_key.as_str(), SCHEDULER_AFFINITY_TTL),
|
||||||
|
Some(target)
|
||||||
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user