fix(gateway): persist quota refresh from strong catalog reads

This commit is contained in:
ZheFox
2026-09-03 13:16:40 +08:00
parent 76fb8905c9
commit 40a5e1470d
4 changed files with 163 additions and 11 deletions
@@ -213,6 +213,21 @@ impl GatewayDataState {
self
}
#[cfg(test)]
pub(crate) fn with_cached_provider_catalog_reader_for_tests<T>(
mut self,
repository: Arc<T>,
) -> Self
where
T: ProviderCatalogReadRepository + 'static,
{
let inner: Arc<dyn ProviderCatalogReadRepository> = repository;
self.provider_catalog_reader = Some(Arc::new(
super::provider_catalog_cache::CachedProviderCatalogReadRepository::new(inner),
));
self
}
#[cfg(test)]
pub(crate) fn with_request_candidate_reader(
mut self,
@@ -240,7 +240,16 @@ fn merge_upstream_metadata(
.unwrap_or_default();
if let Some(update_object) = updates.as_object() {
for (key, value) in update_object {
merged.insert(key.clone(), value.clone());
let mut next = value.clone();
if let (Some(current_namespace), Some(next_namespace)) = (
merged.get(key).and_then(serde_json::Value::as_object),
next.as_object_mut(),
) {
let mut combined = current_namespace.clone();
combined.extend(next_namespace.clone());
next = serde_json::Value::Object(combined);
}
merged.insert(key.clone(), next);
}
}
serde_json::Value::Object(merged)
@@ -1306,7 +1315,8 @@ where
F: std::future::Future<Output = ()>,
{
let Some(mut latest_key) = state
.read_provider_catalog_keys_by_ids(&[key_id.to_string()])
.app()
.list_provider_catalog_keys_by_ids_strong(&[key_id.to_string()])
.await?
.into_iter()
.next()
@@ -1350,9 +1360,18 @@ where
let metadata_updates = metadata_update
.and_then(serde_json::Value::as_object)
.map(|updates| {
let merged = latest_key
.upstream_metadata
.as_ref()
.and_then(serde_json::Value::as_object);
updates
.iter()
.map(|(namespace, value)| (namespace.clone(), value.clone()))
.keys()
.filter_map(|namespace| {
merged
.and_then(|metadata| metadata.get(namespace))
.cloned()
.map(|value| (namespace.clone(), value))
})
.collect::<Vec<_>>()
})
.unwrap_or_default();
@@ -1384,7 +1403,7 @@ where
} else {
serde_json::json!({})
};
let mut expected = observed_upstream_metadata
let expected = observed_upstream_metadata
.as_ref()
.and_then(serde_json::Value::as_object)
.and_then(|metadata| metadata.get(namespace))
@@ -2807,4 +2826,121 @@ mod tests {
json!({"remaining":4})
);
}
#[tokio::test]
async fn quota_refresh_strong_read_bypasses_stale_provider_catalog_cache() {
let mut key = StoredProviderCatalogKey::new(
"key-antigravity-stale-cache".to_string(),
"provider-antigravity-stale-cache".to_string(),
"Antigravity stale cache".to_string(),
"oauth".to_string(),
None,
true,
)
.expect("key should build");
key.upstream_metadata = Some(json!({
"antigravity": {
"project_id": "project-1",
"quota_by_model": {
"gemini-3.7-flash-tiered": {"remaining_fraction": 0.9}
}
}
}));
let repository = Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![],
vec![],
vec![key],
));
let data =
GatewayDataState::with_provider_catalog_repository_for_tests(Arc::clone(&repository))
.with_cached_provider_catalog_reader_for_tests(Arc::clone(&repository));
let app = AppState::new()
.expect("app should build")
.with_data_state_for_tests(data);
let admin_state = AdminAppState::new(&app);
let key_ids = ["key-antigravity-stale-cache".to_string()];
let cached = app
.read_provider_catalog_keys_by_ids(&key_ids)
.await
.expect("initial cached read should succeed");
assert_eq!(
cached[0].upstream_metadata.as_ref().unwrap()["antigravity"]["quota_by_model"]
["gemini-3.7-flash-tiered"]["remaining_fraction"],
json!(0.9)
);
let current_namespace = json!({
"project_id": "project-1",
"model_fetch_revision": 2,
"quota_by_model": {
"gemini-3.7-flash-tiered": {"remaining_fraction": 0.7}
}
});
assert!(repository
.upsert_key_upstream_metadata_namespace(
"key-antigravity-stale-cache",
"antigravity",
&current_namespace,
None,
)
.await
.expect("out-of-band metadata update should succeed"));
let still_cached = app
.read_provider_catalog_keys_by_ids(&key_ids)
.await
.expect("stale cached read should succeed");
assert_eq!(
still_cached[0].upstream_metadata.as_ref().unwrap()["antigravity"]["quota_by_model"]
["gemini-3.7-flash-tiered"]["remaining_fraction"],
json!(0.9),
"regression setup must keep the ordinary read stale"
);
let metadata_update = json!({
"antigravity": {
"project_id": "project-1",
"quota_by_model": {
"gemini-3.7-flash-tiered": {"remaining_fraction": 0.6}
},
"quota_groups": [{
"display_name": "Gemini models",
"buckets": [{"bucket_id": "gemini-weekly", "window": "weekly"}]
}]
}
});
assert!(persist_provider_quota_refresh_state(
&admin_state,
"key-antigravity-stale-cache",
Some(&metadata_update),
None,
None,
None,
)
.await
.expect("quota refresh persistence should not error"));
let stored = repository
.list_keys_by_ids(&key_ids)
.await
.expect("key should reload")
.pop()
.expect("key should exist");
assert_eq!(
stored.upstream_metadata.as_ref().unwrap()["antigravity"]["quota_groups"][0]["buckets"]
[0]["bucket_id"],
json!("gemini-weekly")
);
assert_eq!(
stored.upstream_metadata.as_ref().unwrap()["antigravity"]["model_fetch_revision"],
json!(2),
"quota refresh must preserve fields written by another Antigravity metadata producer"
);
assert_eq!(
stored.upstream_metadata.as_ref().unwrap()["antigravity"]["quota_by_model"]
["gemini-3.7-flash-tiered"]["remaining_fraction"],
json!(0.6)
);
}
}
@@ -2536,7 +2536,7 @@ async fn gateway_refreshes_admin_provider_quota_locally_for_antigravity_with_tru
.upstream_metadata
.as_ref()
.and_then(|value| value.get("antigravity"))
.and_then(|value| value.get("models"))
.and_then(|value| value.get("quota_by_model"))
.and_then(|value| value.get("claude-sonnet-4"))
.and_then(|value| value.get("remaining_fraction")),
Some(&json!(0.25))
@@ -2546,7 +2546,7 @@ async fn gateway_refreshes_admin_provider_quota_locally_for_antigravity_with_tru
.upstream_metadata
.as_ref()
.and_then(|value| value.get("antigravity"))
.and_then(|value| value.get("models"))
.and_then(|value| value.get("quota_by_model"))
.and_then(|value| value.get("claude-sonnet-4"))
.and_then(|value| value.get("used_percent")),
Some(&json!(75.0))
+5 -4
View File
@@ -282,7 +282,7 @@ pub fn parse_antigravity_usage_response(
"is_forbidden": false,
"forbidden_reason": serde_json::Value::Null,
"forbidden_at": serde_json::Value::Null,
"models": quota_by_model,
"quota_by_model": quota_by_model,
}))
}
@@ -7026,16 +7026,17 @@ mod tests {
.expect("antigravity quota should parse");
assert_eq!(
parsed["models"]["RateLimitResetCredit_05cbb6eeeb9c81918e011d8300f9ebfb"]
parsed["quota_by_model"]["RateLimitResetCredit_05cbb6eeeb9c81918e011d8300f9ebfb"]
["display_name"],
json!("Key-1")
);
assert_eq!(
parsed["models"]["RateLimitResetCredit_05cbb6eeeb9c81918e011d8300f9ebfb"]["reset_time"],
parsed["quota_by_model"]["RateLimitResetCredit_05cbb6eeeb9c81918e011d8300f9ebfb"]
["reset_time"],
json!("2030-01-01T00:00:00Z")
);
assert_eq!(
parsed["models"]["gemini-3-pro-preview"]["display_name"],
parsed["quota_by_model"]["gemini-3-pro-preview"]["display_name"],
json!("Gemini 3 Pro Preview")
);
}