mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-04 08:27:46 +08:00
fix(codex): restore recovered quota status and stabilize lifecycle tests
This commit is contained in:
@@ -11,6 +11,7 @@ aether-contracts.workspace = true
|
||||
aether-data-contracts.workspace = true
|
||||
aether-pool-core.workspace = true
|
||||
aether-provider-transport.workspace = true
|
||||
chrono.workspace = true
|
||||
serde_json.workspace = true
|
||||
url.workspace = true
|
||||
uuid.workspace = true
|
||||
|
||||
@@ -38,11 +38,12 @@ pub use providers::{
|
||||
XAI_BILLING_PATH, XAI_USER_PATH,
|
||||
};
|
||||
pub use quota::{
|
||||
provider_pool_key_account_quota_exhausted, provider_pool_key_minimum_quota_reached,
|
||||
provider_pool_key_model_quota_exhausted, provider_pool_key_model_quota_hard_blocked,
|
||||
provider_pool_key_quota_hard_blocked, provider_pool_key_scheduling_label,
|
||||
provider_pool_member_quota_snapshot, provider_pool_quota_metadata_provider_type,
|
||||
provider_pool_quota_metadata_updated_at, provider_pool_quota_snapshot_updated_at,
|
||||
provider_pool_codex_metadata_has_account_quota, provider_pool_key_account_quota_exhausted,
|
||||
provider_pool_key_minimum_quota_reached, provider_pool_key_model_quota_exhausted,
|
||||
provider_pool_key_model_quota_hard_blocked, provider_pool_key_quota_hard_blocked,
|
||||
provider_pool_key_scheduling_label, provider_pool_member_quota_snapshot,
|
||||
provider_pool_quota_metadata_provider_type, provider_pool_quota_metadata_updated_at,
|
||||
provider_pool_quota_snapshot_updated_at,
|
||||
};
|
||||
pub use quota_refresh::ProviderPoolQuotaRequestSpec;
|
||||
pub use service::ProviderPoolService;
|
||||
@@ -1600,6 +1601,245 @@ mod tests {
|
||||
assert!(!signals.quota_exhausted);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_newer_flat_quota_metadata_clears_stale_exhausted_snapshot() {
|
||||
let now = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.expect("system time should be after unix epoch")
|
||||
.as_secs();
|
||||
let service = ProviderPoolService::with_builtin_adapters();
|
||||
// Headers and quota refreshes persist the flat metadata shape while the
|
||||
// derived snapshot may still describe the previous quota observation.
|
||||
for used_percent in [83.0, 99.0, 100.0] {
|
||||
let mut key = sample_key(Some(json!({
|
||||
"codex": {
|
||||
"updated_at": now,
|
||||
"primary_used_percent": used_percent,
|
||||
"primary_reset_at": now + 3600
|
||||
}
|
||||
})));
|
||||
key.status_snapshot = Some(json!({
|
||||
"quota": {
|
||||
"provider_type": "codex",
|
||||
"observed_at": now - 60,
|
||||
"updated_at": now - 60,
|
||||
"code": "exhausted",
|
||||
"exhausted": true,
|
||||
"allowed": false,
|
||||
"limit_reached": true,
|
||||
"reset_at": now + 3600,
|
||||
"windows": [{
|
||||
"code": "weekly",
|
||||
"scope": "account",
|
||||
"used_ratio": 1.0,
|
||||
"remaining_ratio": 0.0,
|
||||
"is_exhausted": true,
|
||||
"reset_at": now + 3600
|
||||
}]
|
||||
}
|
||||
}));
|
||||
for model in [None, Some("gpt-5.4"), Some("o3")] {
|
||||
let signals = service.member_signals("codex", &key, None, model);
|
||||
assert_eq!(signals.quota_exhausted, used_percent >= 100.0);
|
||||
assert!(!signals.quota_hard_blocked);
|
||||
assert_eq!(
|
||||
provider_pool_key_minimum_quota_reached(&key, "codex", model),
|
||||
used_percent >= 99.0
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_account_and_model_quota_agree_on_explicit_allow_and_deny_flags() {
|
||||
let service = ProviderPoolService::with_builtin_adapters();
|
||||
for (flags, exhausted) in [
|
||||
(json!({ "allowed": true }), false),
|
||||
(json!({ "limit_reached": false }), false),
|
||||
(json!({ "allowed": true, "limit_reached": true }), true),
|
||||
(json!({ "allowed": false, "limit_reached": false }), true),
|
||||
] {
|
||||
for use_windows in [false, true] {
|
||||
let mut metadata = flags.clone();
|
||||
metadata["updated_at"] = json!(300);
|
||||
if use_windows {
|
||||
metadata["windows"] = json!([{ "code": "weekly", "used_ratio": 1.0 }]);
|
||||
} else {
|
||||
metadata["primary_used_percent"] = json!(100.0);
|
||||
}
|
||||
let key = sample_key(Some(json!({ "codex": metadata })));
|
||||
for model in [None, Some("gpt-5.4"), Some("o3")] {
|
||||
assert_eq!(
|
||||
service
|
||||
.member_signals("codex", &key, None, model)
|
||||
.quota_exhausted,
|
||||
exhausted,
|
||||
"flags={flags}, use_windows={use_windows}, model={model:?}"
|
||||
);
|
||||
// The opt-in reserve still protects a numerically full
|
||||
// window even when the upstream reports it as allowed.
|
||||
assert!(provider_pool_key_minimum_quota_reached(
|
||||
&key, "codex", model
|
||||
));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_model_quota_honors_latest_account_refusal_without_blocking_spark() {
|
||||
let service = ProviderPoolService::with_builtin_adapters();
|
||||
for include_window in [false, true] {
|
||||
let mut metadata = json!({
|
||||
"updated_at": 300,
|
||||
"allowed": false,
|
||||
"limit_reached": true,
|
||||
"spark_primary_used_percent": 17.0
|
||||
});
|
||||
if include_window {
|
||||
metadata["primary_used_percent"] = json!(83.0);
|
||||
}
|
||||
let mut key = sample_key(Some(json!({ "codex": metadata })));
|
||||
for include_snapshot in [false, true] {
|
||||
if include_snapshot {
|
||||
key.status_snapshot = Some(json!({
|
||||
"quota": {
|
||||
"provider_type": "codex",
|
||||
"observed_at": 200,
|
||||
"exhausted": false,
|
||||
"windows": [{ "code": "weekly", "used_ratio": 0.83 }]
|
||||
}
|
||||
}));
|
||||
}
|
||||
assert!(
|
||||
service
|
||||
.member_signals("codex", &key, None, Some("gpt-5.4"))
|
||||
.quota_exhausted
|
||||
);
|
||||
assert!(
|
||||
!service
|
||||
.member_signals("codex", &key, None, Some("gpt-5.3-codex-spark"))
|
||||
.quota_exhausted
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_quota_source_freshness_supports_all_timestamp_formats() {
|
||||
let service = ProviderPoolService::with_builtin_adapters();
|
||||
for observed_at in [
|
||||
json!(1_700_000_200_u64),
|
||||
json!(1_700_000_200_000_u64),
|
||||
json!("1700000200000"),
|
||||
json!("2023-11-14T22:16:40Z"),
|
||||
] {
|
||||
for use_windows in [false, true] {
|
||||
let metadata = if use_windows {
|
||||
json!({
|
||||
"updated_at": observed_at,
|
||||
"windows": [{ "code": "weekly", "used_ratio": 0.83 }]
|
||||
})
|
||||
} else {
|
||||
json!({ "updated_at": observed_at, "primary_used_percent": 83.0 })
|
||||
};
|
||||
let mut key = sample_key(Some(json!({ "codex": metadata })));
|
||||
key.status_snapshot = Some(json!({
|
||||
"quota": {
|
||||
"provider_type": "codex",
|
||||
"observed_at": "2023-11-14T22:15:00Z",
|
||||
"exhausted": true,
|
||||
"allowed": false,
|
||||
"windows": [{ "code": "weekly", "used_ratio": 1.0 }]
|
||||
}
|
||||
}));
|
||||
for model in [None, Some("gpt-5.4")] {
|
||||
let signals = service.member_signals("codex", &key, None, model);
|
||||
assert!(!signals.quota_exhausted, "timestamp={observed_at}");
|
||||
assert!(!signals.quota_hard_blocked, "timestamp={observed_at}");
|
||||
assert!(!provider_pool_key_minimum_quota_reached(
|
||||
&key, "codex", model
|
||||
));
|
||||
}
|
||||
|
||||
// Swap freshness while retaining the same observations. An
|
||||
// older usable bucket cannot erase a later exhausted snapshot.
|
||||
key.status_snapshot.as_mut().unwrap()["quota"]["observed_at"] =
|
||||
json!("2023-11-14T22:18:20Z");
|
||||
assert!(provider_pool_key_account_quota_exhausted(&key, "codex"));
|
||||
assert!(provider_pool_key_quota_hard_blocked(&key, "codex"));
|
||||
assert_eq!(
|
||||
provider_pool_key_model_quota_exhausted(&key, "codex", "gpt-5.4"),
|
||||
Some(true)
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_account_quota_recovery_requires_new_account_observation() {
|
||||
for metadata in [
|
||||
json!({ "updated_at": 100, "primary_used_percent": 83.0 }),
|
||||
json!({ "primary_used_percent": 83.0 }),
|
||||
json!({ "updated_at": 300, "plan_type": "plus" }),
|
||||
json!({ "updated_at": 300, "credits_unlimited": false }),
|
||||
json!({ "updated_at": 300, "spark_primary_used_percent": 83.0 }),
|
||||
json!({ "updated_at": 300, "windows": [{ "code": "weekly" }] }),
|
||||
] {
|
||||
let mut key = sample_key(Some(json!({ "codex": metadata })));
|
||||
key.status_snapshot = Some(json!({
|
||||
"quota": {
|
||||
"provider_type": "codex",
|
||||
"observed_at": 200,
|
||||
"exhausted": true,
|
||||
"allowed": false,
|
||||
"limit_reached": true,
|
||||
"windows": [{ "code": "weekly", "used_ratio": 1.0 }]
|
||||
}
|
||||
}));
|
||||
assert!(provider_pool_key_account_quota_exhausted(&key, "codex"));
|
||||
assert!(provider_pool_key_quota_hard_blocked(&key, "codex"));
|
||||
assert_eq!(
|
||||
provider_pool_key_model_quota_exhausted(&key, "codex", "gpt-5.4"),
|
||||
Some(true)
|
||||
);
|
||||
}
|
||||
|
||||
let mut key = sample_key(Some(json!({
|
||||
"codex": {
|
||||
"updated_at": 300,
|
||||
"primary_used_percent": 83.0,
|
||||
"allowed": false,
|
||||
"limit_reached": true
|
||||
}
|
||||
})));
|
||||
key.status_snapshot = Some(json!({
|
||||
"quota": { "provider_type": "codex", "updated_at": 200, "exhausted": false }
|
||||
}));
|
||||
assert!(provider_pool_key_account_quota_exhausted(&key, "codex"));
|
||||
assert!(provider_pool_key_quota_hard_blocked(&key, "codex"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_account_updates_do_not_override_independent_model_quota() {
|
||||
for (spark_ratio, account_percent) in [(0.83, 100.0), (1.0, 83.0)] {
|
||||
let mut key = sample_key(Some(json!({
|
||||
"codex": { "updated_at": 300, "primary_used_percent": account_percent }
|
||||
})));
|
||||
key.status_snapshot = Some(json!({
|
||||
"quota": {
|
||||
"provider_type": "codex",
|
||||
"updated_at": 200,
|
||||
"windows": [{ "code": "spark_5h", "used_ratio": spark_ratio }]
|
||||
}
|
||||
}));
|
||||
assert_eq!(
|
||||
provider_pool_key_model_quota_exhausted(&key, "codex", "gpt-5.3-codex-spark"),
|
||||
Some(spark_ratio >= 1.0)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_explicit_quota_block_is_hard_until_reset() {
|
||||
let now = std::time::SystemTime::now()
|
||||
|
||||
@@ -10,10 +10,11 @@ use crate::provider::{
|
||||
ProviderPoolMemberInput,
|
||||
};
|
||||
use crate::quota::{
|
||||
provider_pool_current_unix_secs, provider_pool_json_bool, provider_pool_json_f64,
|
||||
provider_pool_member_quota_snapshot, provider_pool_metadata_bucket,
|
||||
provider_pool_model_quota_exhausted, provider_pool_quota_snapshot_exhausted_decision,
|
||||
provider_pool_reset_deadline_elapsed, provider_pool_timestamp_unix_secs,
|
||||
provider_pool_codex_metadata_has_account_quota, provider_pool_current_unix_secs,
|
||||
provider_pool_json_bool, provider_pool_json_f64, provider_pool_member_quota_snapshot,
|
||||
provider_pool_metadata_bucket, provider_pool_model_quota_exhausted,
|
||||
provider_pool_quota_snapshot_exhausted_decision, provider_pool_reset_deadline_elapsed,
|
||||
provider_pool_source_account_quota_exhausted, provider_pool_timestamp_unix_secs,
|
||||
};
|
||||
use crate::quota_refresh::ProviderPoolQuotaRequestSpec;
|
||||
|
||||
@@ -54,6 +55,9 @@ impl ProviderPoolAdapter for CodexProviderPoolAdapter {
|
||||
}) {
|
||||
return exhausted;
|
||||
}
|
||||
if let Some(bucket) = codex_newer_account_quota_metadata(input.key, input.provider_type) {
|
||||
return quota_exhausted_from_bucket(bucket);
|
||||
}
|
||||
if let Some(quota_snapshot) =
|
||||
provider_pool_member_quota_snapshot(input.key, input.provider_type)
|
||||
{
|
||||
@@ -108,15 +112,12 @@ fn codex_explicit_quota_block_active(
|
||||
key: &aether_data_contracts::repository::provider_catalog::StoredProviderCatalogKey,
|
||||
provider_type: &str,
|
||||
) -> bool {
|
||||
if let Some(bucket) = codex_newer_account_quota_metadata(key, provider_type) {
|
||||
return codex_explicit_quota_block_from_bucket(bucket);
|
||||
}
|
||||
let Some(quota_snapshot) = provider_pool_member_quota_snapshot(key, provider_type) else {
|
||||
return provider_pool_metadata_bucket(key.upstream_metadata.as_ref(), provider_type)
|
||||
.is_some_and(|bucket| {
|
||||
(provider_pool_json_bool(bucket.get("allowed")) == Some(false)
|
||||
|| provider_pool_json_bool(bucket.get("limit_reached")) == Some(true))
|
||||
&& !["primary", "secondary"]
|
||||
.into_iter()
|
||||
.any(|prefix| codex_window_reset_elapsed(bucket, prefix))
|
||||
});
|
||||
.is_some_and(codex_explicit_quota_block_from_bucket);
|
||||
};
|
||||
let explicitly_blocked = provider_pool_json_bool(quota_snapshot.get("allowed")) == Some(false)
|
||||
|| provider_pool_json_bool(quota_snapshot.get("limit_reached")) == Some(true);
|
||||
@@ -133,6 +134,35 @@ fn codex_explicit_quota_block_active(
|
||||
})
|
||||
}
|
||||
|
||||
fn codex_explicit_quota_block_from_bucket(bucket: &Map<String, Value>) -> bool {
|
||||
(provider_pool_json_bool(bucket.get("allowed")) == Some(false)
|
||||
|| provider_pool_json_bool(bucket.get("limit_reached")) == Some(true))
|
||||
&& !["primary", "secondary"]
|
||||
.into_iter()
|
||||
.any(|prefix| codex_window_reset_elapsed(bucket, prefix))
|
||||
}
|
||||
|
||||
/// A successful refresh can update raw metadata before the status snapshot.
|
||||
/// Do not keep an old account-level block once a newer quota observation exists.
|
||||
/// Identity-only or model-only updates cannot clear an account quota decision.
|
||||
fn codex_newer_account_quota_metadata<'a>(
|
||||
key: &'a aether_data_contracts::repository::provider_catalog::StoredProviderCatalogKey,
|
||||
provider_type: &str,
|
||||
) -> Option<&'a Map<String, Value>> {
|
||||
let snapshot = provider_pool_member_quota_snapshot(key, provider_type)?;
|
||||
let metadata = provider_pool_metadata_bucket(key.upstream_metadata.as_ref(), provider_type)?;
|
||||
if !provider_pool_codex_metadata_has_account_quota(metadata) {
|
||||
return None;
|
||||
}
|
||||
let metadata_observed_at = provider_pool_timestamp_unix_secs(metadata.get("observed_at"))
|
||||
.or_else(|| provider_pool_timestamp_unix_secs(metadata.get("updated_at")))?;
|
||||
let snapshot_observed_at = provider_pool_timestamp_unix_secs(snapshot.get("observed_at"))
|
||||
.or_else(|| provider_pool_timestamp_unix_secs(snapshot.get("updated_at")));
|
||||
snapshot_observed_at
|
||||
.is_none_or(|observed_at| metadata_observed_at >= observed_at)
|
||||
.then_some(metadata)
|
||||
}
|
||||
|
||||
fn build_codex_wham_headers(
|
||||
resolved_oauth_auth: Option<(String, String)>,
|
||||
decrypted_api_key: Option<&str>,
|
||||
@@ -322,11 +352,14 @@ pub(crate) fn quota_exhausted_from_bucket(bucket: &Map<String, Value>) -> bool {
|
||||
if provider_pool_json_bool(bucket.get("credits_unlimited")) == Some(true) {
|
||||
return false;
|
||||
}
|
||||
let account_windows_exhausted = provider_pool_source_account_quota_exhausted(bucket);
|
||||
let has_window_data = provider_pool_json_f64(bucket.get("primary_used_percent")).is_some()
|
||||
|| provider_pool_json_f64(bucket.get("secondary_used_percent")).is_some();
|
||||
|| provider_pool_json_f64(bucket.get("secondary_used_percent")).is_some()
|
||||
|| account_windows_exhausted.is_some();
|
||||
if !has_window_data && provider_pool_json_bool(bucket.get("has_credits")) == Some(false) {
|
||||
return true;
|
||||
}
|
||||
codex_window_used_percent_exhausted(bucket, "primary")
|
||||
account_windows_exhausted == Some(true)
|
||||
|| codex_window_used_percent_exhausted(bucket, "primary")
|
||||
|| codex_window_used_percent_exhausted(bucket, "secondary")
|
||||
}
|
||||
|
||||
@@ -87,6 +87,7 @@ fn provider_pool_quota_reaches_reserve(
|
||||
provider_model_name: Option<&str>,
|
||||
reserve_ratio: f64,
|
||||
) -> Option<bool> {
|
||||
let is_codex = provider_type.trim().eq_ignore_ascii_case("codex");
|
||||
let requested = provider_model_name.map(provider_pool_identifier_tokens);
|
||||
if reserve_ratio <= 0.0 && requested.as_ref().is_some_and(|tokens| tokens.is_empty()) {
|
||||
return None;
|
||||
@@ -103,12 +104,19 @@ fn provider_pool_quota_reaches_reserve(
|
||||
let mut resolved = None::<(Option<u64>, bool)>;
|
||||
let mut resolved_specific_bucket = false;
|
||||
for source in sources.into_iter().flatten() {
|
||||
let observed_at = provider_pool_timestamp_unix_secs(source.get("observed_at"))
|
||||
.or_else(|| provider_pool_timestamp_unix_secs(source.get("updated_at")));
|
||||
let account_signal = (is_codex && reserve_ratio <= 0.0)
|
||||
.then(|| provider_pool_codex_account_quota_signal(source, observed_at))
|
||||
.flatten();
|
||||
let mut windows = provider_pool_collect_quota_windows(source);
|
||||
if reserve_ratio > 0.0 && provider_type.trim().eq_ignore_ascii_case("codex") {
|
||||
if is_codex {
|
||||
// Raw refresh metadata can be newer than the materialized windows.
|
||||
// Include all four legacy slots before selecting the model bucket.
|
||||
for prefix in ["primary", "secondary", "spark_primary", "spark_secondary"] {
|
||||
let Some(used_percent) = source.get(&format!("{prefix}_used_percent")) else {
|
||||
let Some(used_percent) =
|
||||
provider_pool_json_f64(source.get(&format!("{prefix}_used_percent")))
|
||||
else {
|
||||
continue;
|
||||
};
|
||||
if provider_pool_json_f64(source.get(&format!("{prefix}_window_minutes")))
|
||||
@@ -118,7 +126,7 @@ fn provider_pool_quota_reaches_reserve(
|
||||
}
|
||||
let mut window = Map::from_iter([
|
||||
("code".to_string(), json!(prefix)),
|
||||
("used_percent".to_string(), used_percent.clone()),
|
||||
("used_percent".to_string(), json!(used_percent)),
|
||||
]);
|
||||
for field in [
|
||||
"reset_at",
|
||||
@@ -132,12 +140,21 @@ fn provider_pool_quota_reaches_reserve(
|
||||
}
|
||||
windows.push(window);
|
||||
}
|
||||
windows.retain(provider_pool_quota_window_has_observation);
|
||||
if let Some(exhausted) = account_signal {
|
||||
if !windows.iter().any(provider_pool_window_is_generic) {
|
||||
// A flags-only quota refresh is still a newer account
|
||||
// observation, but must not override a model-specific bucket.
|
||||
windows.push(Map::from_iter([
|
||||
("code".to_string(), json!("account")),
|
||||
("is_exhausted".to_string(), json!(exhausted)),
|
||||
]));
|
||||
}
|
||||
}
|
||||
}
|
||||
if windows.is_empty() {
|
||||
continue;
|
||||
}
|
||||
let observed_at = provider_pool_timestamp_unix_secs(source.get("observed_at"))
|
||||
.or_else(|| provider_pool_timestamp_unix_secs(source.get("updated_at")));
|
||||
let model_matches = windows
|
||||
.iter()
|
||||
.filter(|window| {
|
||||
@@ -155,7 +172,7 @@ fn provider_pool_quota_reaches_reserve(
|
||||
provider_pool_explicit_model_windows_exhausted(model_matches, observed_at)
|
||||
};
|
||||
if resolved.is_none()
|
||||
|| (reserve_ratio > 0.0 && !resolved_specific_bucket)
|
||||
|| (is_codex && !resolved_specific_bucket)
|
||||
|| provider_pool_should_replace_model_quota_resolution(
|
||||
resolved.as_ref().and_then(|(observed_at, _)| *observed_at),
|
||||
observed_at,
|
||||
@@ -183,7 +200,7 @@ fn provider_pool_quota_reaches_reserve(
|
||||
let exhausted =
|
||||
provider_pool_any_window_exhausted(family_matches, observed_at, reserve_ratio);
|
||||
if resolved.is_none()
|
||||
|| (reserve_ratio > 0.0 && !resolved_specific_bucket)
|
||||
|| (is_codex && !resolved_specific_bucket)
|
||||
|| provider_pool_should_replace_model_quota_resolution(
|
||||
resolved.as_ref().and_then(|(observed_at, _)| *observed_at),
|
||||
observed_at,
|
||||
@@ -197,7 +214,7 @@ fn provider_pool_quota_reaches_reserve(
|
||||
|
||||
// A newer account observation does not update an independent model
|
||||
// bucket. Only compare freshness between applicable model sources.
|
||||
if reserve_ratio > 0.0 && resolved_specific_bucket {
|
||||
if is_codex && resolved_specific_bucket {
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -211,8 +228,9 @@ fn provider_pool_quota_reaches_reserve(
|
||||
.filter(|window| provider_pool_window_is_generic(window))
|
||||
.collect::<Vec<_>>();
|
||||
if !generic_matches.is_empty() {
|
||||
let exhausted =
|
||||
provider_pool_any_window_exhausted(generic_matches, observed_at, reserve_ratio);
|
||||
let exhausted = account_signal.unwrap_or_else(|| {
|
||||
provider_pool_any_window_exhausted(generic_matches, observed_at, reserve_ratio)
|
||||
});
|
||||
if resolved.is_none()
|
||||
|| provider_pool_should_replace_model_quota_resolution(
|
||||
resolved.as_ref().and_then(|(observed_at, _)| *observed_at),
|
||||
@@ -227,6 +245,37 @@ fn provider_pool_quota_reaches_reserve(
|
||||
resolved.map(|(_, exhausted)| exhausted)
|
||||
}
|
||||
|
||||
fn provider_pool_codex_account_quota_signal(
|
||||
source: &Map<String, Value>,
|
||||
observed_at: Option<u64>,
|
||||
) -> Option<bool> {
|
||||
let allowed = provider_pool_json_bool(source.get("allowed"));
|
||||
let limit_reached = provider_pool_json_bool(source.get("limit_reached"));
|
||||
if allowed == Some(false) || limit_reached == Some(true) {
|
||||
let reset_elapsed = provider_pool_current_unix_secs().is_some_and(|now| {
|
||||
provider_pool_reset_deadline_elapsed(source, observed_at, now)
|
||||
|| ["primary", "secondary"].into_iter().any(|prefix| {
|
||||
let window = [
|
||||
"reset_at",
|
||||
"next_reset_at",
|
||||
"reset_seconds",
|
||||
"reset_after_seconds",
|
||||
]
|
||||
.into_iter()
|
||||
.filter_map(|field| {
|
||||
source
|
||||
.get(&format!("{prefix}_{field}"))
|
||||
.map(|value| (field.to_string(), value.clone()))
|
||||
})
|
||||
.collect::<Map<_, _>>();
|
||||
provider_pool_reset_deadline_elapsed(&window, observed_at, now)
|
||||
})
|
||||
});
|
||||
return (!reset_elapsed).then_some(true);
|
||||
}
|
||||
(allowed == Some(true) || limit_reached == Some(false)).then_some(false)
|
||||
}
|
||||
|
||||
fn provider_pool_should_replace_model_quota_resolution(
|
||||
previous_observed_at: Option<u64>,
|
||||
next_observed_at: Option<u64>,
|
||||
@@ -282,6 +331,69 @@ fn provider_pool_any_window_exhausted(
|
||||
})
|
||||
}
|
||||
|
||||
fn provider_pool_quota_window_has_observation(window: &Map<String, Value>) -> bool {
|
||||
["is_exhausted", "exhausted"]
|
||||
.into_iter()
|
||||
.any(|field| provider_pool_json_bool(window.get(field)).is_some())
|
||||
|| [
|
||||
"used_ratio",
|
||||
"usage_ratio",
|
||||
"used_percent",
|
||||
"remaining_ratio",
|
||||
"remaining_fraction",
|
||||
"remaining_percent",
|
||||
]
|
||||
.into_iter()
|
||||
.any(|field| provider_pool_json_f64(window.get(field)).is_some())
|
||||
|| (provider_pool_json_f64(
|
||||
window
|
||||
.get("remaining")
|
||||
.or_else(|| window.get("remaining_value")),
|
||||
)
|
||||
.is_some()
|
||||
&& provider_pool_json_f64(
|
||||
window
|
||||
.get("limit")
|
||||
.or_else(|| window.get("limit_value"))
|
||||
.or_else(|| window.get("total")),
|
||||
)
|
||||
.is_some_and(|limit| limit > 0.0))
|
||||
}
|
||||
|
||||
pub(crate) fn provider_pool_source_account_quota_exhausted(
|
||||
source: &Map<String, Value>,
|
||||
) -> Option<bool> {
|
||||
let windows = provider_pool_collect_quota_windows(source);
|
||||
let account_windows = windows
|
||||
.iter()
|
||||
.filter(|window| provider_pool_window_is_generic(window))
|
||||
.filter(|window| provider_pool_quota_window_has_observation(window))
|
||||
.collect::<Vec<_>>();
|
||||
if account_windows.is_empty() {
|
||||
return None;
|
||||
}
|
||||
let observed_at = provider_pool_timestamp_unix_secs(source.get("observed_at"))
|
||||
.or_else(|| provider_pool_timestamp_unix_secs(source.get("updated_at")));
|
||||
Some(provider_pool_any_window_exhausted(
|
||||
account_windows,
|
||||
observed_at,
|
||||
0.0,
|
||||
))
|
||||
}
|
||||
|
||||
/// Whether Codex metadata contains an account quota observation rather than an
|
||||
/// identity-only update or an independent model quota bucket.
|
||||
pub fn provider_pool_codex_metadata_has_account_quota(source: &Map<String, Value>) -> bool {
|
||||
["primary_used_percent", "secondary_used_percent"]
|
||||
.into_iter()
|
||||
.any(|field| provider_pool_json_f64(source.get(field)).is_some())
|
||||
|| provider_pool_source_account_quota_exhausted(source).is_some()
|
||||
|| ["allowed", "limit_reached", "has_credits"]
|
||||
.into_iter()
|
||||
.any(|field| provider_pool_json_bool(source.get(field)).is_some())
|
||||
|| provider_pool_json_bool(source.get("credits_unlimited")) == Some(true)
|
||||
}
|
||||
|
||||
/// Whether an applicable Codex window has at most 1% remaining. Missing quota
|
||||
/// data and windows whose reset has elapsed do not trigger this opt-in guard.
|
||||
pub fn provider_pool_key_minimum_quota_reached(
|
||||
@@ -700,7 +812,16 @@ pub(crate) fn provider_pool_json_f64(value: Option<&Value>) -> Option<f64> {
|
||||
}
|
||||
|
||||
pub(crate) fn provider_pool_timestamp_unix_secs(value: Option<&Value>) -> Option<u64> {
|
||||
let mut timestamp = provider_pool_json_f64(value)?;
|
||||
let mut timestamp = match provider_pool_json_f64(value) {
|
||||
Some(timestamp) => timestamp,
|
||||
None => {
|
||||
return chrono::DateTime::parse_from_rfc3339(value?.as_str()?.trim())
|
||||
.ok()?
|
||||
.timestamp()
|
||||
.try_into()
|
||||
.ok();
|
||||
}
|
||||
};
|
||||
if timestamp <= 0.0 {
|
||||
return None;
|
||||
}
|
||||
|
||||
@@ -10055,6 +10055,24 @@ mod tests {
|
||||
2,
|
||||
"the later direct caller should receive its own bounded write attempt"
|
||||
);
|
||||
// Direct persistence can finish before the submission worker joins the
|
||||
// barrier handoff and accounts for its completed slot.
|
||||
timeout(Duration::from_secs(1), async {
|
||||
loop {
|
||||
let snapshot = runtime.metrics_snapshot();
|
||||
let submission = &runtime.lifecycle_submission.state;
|
||||
if snapshot.terminal_submission_pending == 0
|
||||
&& snapshot.ordered_lifecycle_pending == 0
|
||||
&& snapshot.lifecycle_submission_pending == 0
|
||||
&& submission.admission.available_permits() == submission.capacity
|
||||
{
|
||||
break;
|
||||
}
|
||||
sleep(Duration::from_millis(1)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("failed terminal submission accounting and admission should drain");
|
||||
let snapshot = runtime.metrics_snapshot();
|
||||
assert_eq!(snapshot.terminal_submission_pending, 0);
|
||||
assert_eq!(snapshot.ordered_lifecycle_pending, 0);
|
||||
@@ -10162,24 +10180,43 @@ mod tests {
|
||||
.await;
|
||||
|
||||
assert_eq!(remaining_policy_panics.load(Ordering::Acquire), 0);
|
||||
let records = records.lock().expect("records lock");
|
||||
assert_eq!(
|
||||
records
|
||||
.iter()
|
||||
.filter(|record| record.request_id == healthy_request_id)
|
||||
.count(),
|
||||
2,
|
||||
"the same terminal shard should continue processing healthy requests"
|
||||
);
|
||||
assert!(
|
||||
records
|
||||
.iter()
|
||||
.filter(|record| record.request_id == failed_request_id)
|
||||
.count()
|
||||
== 1,
|
||||
"only the later healthy attempt should persist for the panicked request"
|
||||
);
|
||||
drop(records);
|
||||
{
|
||||
let records = records.lock().expect("records lock");
|
||||
assert_eq!(
|
||||
records
|
||||
.iter()
|
||||
.filter(|record| record.request_id == healthy_request_id)
|
||||
.count(),
|
||||
2,
|
||||
"the same terminal shard should continue processing healthy requests"
|
||||
);
|
||||
assert!(
|
||||
records
|
||||
.iter()
|
||||
.filter(|record| record.request_id == failed_request_id)
|
||||
.count()
|
||||
== 1,
|
||||
"only the later healthy attempt should persist for the panicked request"
|
||||
);
|
||||
}
|
||||
// The final direct attempt also submits a barrier whose worker may
|
||||
// account for completion after the persistence call has returned.
|
||||
timeout(Duration::from_secs(1), async {
|
||||
loop {
|
||||
let snapshot = runtime.metrics_snapshot();
|
||||
let submission = &runtime.lifecycle_submission.state;
|
||||
if snapshot.terminal_submission_pending == 0
|
||||
&& snapshot.ordered_lifecycle_pending == 0
|
||||
&& snapshot.lifecycle_submission_pending == 0
|
||||
&& submission.admission.available_permits() == submission.capacity
|
||||
{
|
||||
break;
|
||||
}
|
||||
sleep(Duration::from_millis(1)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("panicked terminal submission accounting and admission should drain");
|
||||
let snapshot = runtime.metrics_snapshot();
|
||||
assert_eq!(snapshot.terminal_submission_pending, 0);
|
||||
assert_eq!(snapshot.ordered_lifecycle_pending, 0);
|
||||
|
||||
Reference in New Issue
Block a user