mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 01:10:23 +08:00
fix(gateway): 修复 balance 刷新去重键冲突、usage 状态回退与测试竞态
- balance_cache: 引入实例级 refresh key 防止多实例共享进程级 HashSet 冲突 - InMemoryUsageRepo: 阻止 pending/streaming 状态覆盖已终结(completed/failed/cancelled)记录 - usage 同步测试: 等待条件从 is_some() 改为检查 status=="completed" 避免竞态 - wallet 测试: 增加轮询等待 wallet 扣款完成 - 整理 import 语句与 tests 模块位置
This commit is contained in:
@@ -4,4 +4,5 @@ mod filters;
|
|||||||
|
|
||||||
pub(super) use aggregations::admin_usage_aggregation_by_user_json;
|
pub(super) use aggregations::admin_usage_aggregation_by_user_json;
|
||||||
pub(super) use cache_affinity::list_recent_completed_usage_for_cache_affinity;
|
pub(super) use cache_affinity::list_recent_completed_usage_for_cache_affinity;
|
||||||
pub(super) use filters::{admin_usage_api_key_names, admin_usage_provider_key_names};
|
pub(super) use filters::admin_usage_api_key_names;
|
||||||
|
pub(super) use filters::admin_usage_provider_key_names;
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
use super::analytics::{admin_usage_api_key_names, admin_usage_provider_key_names};
|
use super::analytics::admin_usage_api_key_names;
|
||||||
|
use super::analytics::admin_usage_provider_key_names;
|
||||||
use super::replay::{
|
use super::replay::{
|
||||||
admin_usage_curl_headers, admin_usage_curl_url, admin_usage_headers_from_value,
|
admin_usage_curl_headers, admin_usage_curl_url, admin_usage_headers_from_value,
|
||||||
admin_usage_id_from_action_path, admin_usage_id_from_detail_path,
|
admin_usage_id_from_action_path, admin_usage_id_from_detail_path,
|
||||||
|
|||||||
@@ -107,40 +107,6 @@ pub(super) fn build_admin_usage_curl_response(
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod tests {
|
|
||||||
use super::admin_usage_body_value_from_sources;
|
|
||||||
use serde_json::json;
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn resolved_reference_body_wins_over_inline_fallback() {
|
|
||||||
let inline_body = json!({
|
|
||||||
"truncated": true,
|
|
||||||
"reason": "usage_capture_limits_exceeded"
|
|
||||||
});
|
|
||||||
let ref_body = json!({
|
|
||||||
"messages": [{"role": "user", "content": "real request body"}]
|
|
||||||
});
|
|
||||||
|
|
||||||
assert_eq!(
|
|
||||||
admin_usage_body_value_from_sources(Some(ref_body.clone()), Some(&inline_body)),
|
|
||||||
Some(ref_body)
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn inline_body_is_used_when_reference_body_is_unavailable() {
|
|
||||||
let inline_body = json!({
|
|
||||||
"messages": [{"role": "user", "content": "fallback inline body"}]
|
|
||||||
});
|
|
||||||
|
|
||||||
assert_eq!(
|
|
||||||
admin_usage_body_value_from_sources(None, Some(&inline_body)),
|
|
||||||
Some(inline_body)
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(super) fn build_admin_usage_detail_payload(
|
pub(super) fn build_admin_usage_detail_payload(
|
||||||
item: &StoredRequestUsageAudit,
|
item: &StoredRequestUsageAudit,
|
||||||
users_by_id: &BTreeMap<String, aether_data::repository::users::StoredUserSummary>,
|
users_by_id: &BTreeMap<String, aether_data::repository::users::StoredUserSummary>,
|
||||||
@@ -415,3 +381,37 @@ pub(super) fn admin_usage_build_curl_command(
|
|||||||
) -> String {
|
) -> String {
|
||||||
aether_admin::observability::usage::admin_usage_build_curl_command(url, headers, body)
|
aether_admin::observability::usage::admin_usage_build_curl_command(url, headers, body)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::admin_usage_body_value_from_sources;
|
||||||
|
use serde_json::json;
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn resolved_reference_body_wins_over_inline_fallback() {
|
||||||
|
let inline_body = json!({
|
||||||
|
"truncated": true,
|
||||||
|
"reason": "usage_capture_limits_exceeded"
|
||||||
|
});
|
||||||
|
let ref_body = json!({
|
||||||
|
"messages": [{"role": "user", "content": "real request body"}]
|
||||||
|
});
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
admin_usage_body_value_from_sources(Some(ref_body.clone()), Some(&inline_body)),
|
||||||
|
Some(ref_body)
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn inline_body_is_used_when_reference_body_is_unavailable() {
|
||||||
|
let inline_body = json!({
|
||||||
|
"messages": [{"role": "user", "content": "fallback inline body"}]
|
||||||
|
});
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
admin_usage_body_value_from_sources(None, Some(&inline_body)),
|
||||||
|
Some(inline_body)
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
use super::super::stats::{AdminStatsTimeRange, AdminStatsUsageFilter};
|
use super::super::stats::{AdminStatsTimeRange, AdminStatsUsageFilter};
|
||||||
use super::analytics::{admin_usage_api_key_names, admin_usage_provider_key_names};
|
use super::analytics::admin_usage_api_key_names;
|
||||||
|
use super::analytics::admin_usage_provider_key_names;
|
||||||
use crate::handlers::admin::request::{AdminAppState, AdminRequestContext};
|
use crate::handlers::admin::request::{AdminAppState, AdminRequestContext};
|
||||||
use crate::handlers::admin::shared::query_param_value;
|
use crate::handlers::admin::shared::query_param_value;
|
||||||
use crate::GatewayError;
|
use crate::GatewayError;
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ use tokio::sync::{Mutex, Semaphore};
|
|||||||
use tracing::{debug, warn};
|
use tracing::{debug, warn};
|
||||||
|
|
||||||
const ADMIN_PROVIDER_OPS_BALANCE_CACHE_PREFIX: &str = "provider_ops:balance:";
|
const ADMIN_PROVIDER_OPS_BALANCE_CACHE_PREFIX: &str = "provider_ops:balance:";
|
||||||
|
const ADMIN_PROVIDER_OPS_BALANCE_REFRESH_PREFIX: &str = "provider_ops:balance_refresh:";
|
||||||
const ADMIN_PROVIDER_OPS_BALANCE_CACHE_TTL_SECS: u64 = 86_400;
|
const ADMIN_PROVIDER_OPS_BALANCE_CACHE_TTL_SECS: u64 = 86_400;
|
||||||
const ADMIN_PROVIDER_OPS_BALANCE_AUTH_FAILED_CACHE_TTL_SECS: u64 = 60;
|
const ADMIN_PROVIDER_OPS_BALANCE_AUTH_FAILED_CACHE_TTL_SECS: u64 = 60;
|
||||||
const ADMIN_PROVIDER_OPS_BALANCE_REFRESH_CONCURRENCY: usize = 3;
|
const ADMIN_PROVIDER_OPS_BALANCE_REFRESH_CONCURRENCY: usize = 3;
|
||||||
@@ -139,8 +140,9 @@ pub(super) async fn spawn_admin_provider_ops_balance_refresh(
|
|||||||
state: &AdminAppState<'_>,
|
state: &AdminAppState<'_>,
|
||||||
provider_id: &str,
|
provider_id: &str,
|
||||||
) {
|
) {
|
||||||
|
let refresh_key = admin_provider_ops_balance_refresh_key(state, provider_id);
|
||||||
let mut guard = ADMIN_PROVIDER_OPS_REFRESHING_PROVIDERS.lock().await;
|
let mut guard = ADMIN_PROVIDER_OPS_REFRESHING_PROVIDERS.lock().await;
|
||||||
if !guard.insert(provider_id.to_string()) {
|
if !guard.insert(refresh_key.clone()) {
|
||||||
debug!(provider_id, "provider ops balance refresh already running");
|
debug!(provider_id, "provider ops balance refresh already running");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
@@ -162,12 +164,12 @@ pub(super) async fn spawn_admin_provider_ops_balance_refresh(
|
|||||||
error = %err,
|
error = %err,
|
||||||
"provider ops balance refresh semaphore closed"
|
"provider ops balance refresh semaphore closed"
|
||||||
);
|
);
|
||||||
finish_refresh_provider(&provider_id).await;
|
finish_refresh_provider(&refresh_key).await;
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
Err(_) => {
|
Err(_) => {
|
||||||
debug!(provider_id = %provider_id, "provider ops balance refresh skipped by concurrency limit");
|
debug!(provider_id = %provider_id, "provider ops balance refresh skipped by concurrency limit");
|
||||||
finish_refresh_provider(&provider_id).await;
|
finish_refresh_provider(&refresh_key).await;
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@@ -186,7 +188,7 @@ pub(super) async fn spawn_admin_provider_ops_balance_refresh(
|
|||||||
"failed to load provider for balance refresh"
|
"failed to load provider for balance refresh"
|
||||||
);
|
);
|
||||||
drop(permit);
|
drop(permit);
|
||||||
finish_refresh_provider(&provider_id).await;
|
finish_refresh_provider(&refresh_key).await;
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@@ -204,7 +206,7 @@ pub(super) async fn spawn_admin_provider_ops_balance_refresh(
|
|||||||
"failed to load endpoints for balance refresh"
|
"failed to load endpoints for balance refresh"
|
||||||
);
|
);
|
||||||
drop(permit);
|
drop(permit);
|
||||||
finish_refresh_provider(&provider_id).await;
|
finish_refresh_provider(&refresh_key).await;
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -223,7 +225,7 @@ pub(super) async fn spawn_admin_provider_ops_balance_refresh(
|
|||||||
.await;
|
.await;
|
||||||
store_admin_provider_ops_balance_cache(&admin_state, &provider_id, &payload).await;
|
store_admin_provider_ops_balance_cache(&admin_state, &provider_id, &payload).await;
|
||||||
drop(permit);
|
drop(permit);
|
||||||
finish_refresh_provider(&provider_id).await;
|
finish_refresh_provider(&refresh_key).await;
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -239,11 +241,24 @@ fn balance_cache_ttl_seconds(payload: &Value) -> Option<u64> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn finish_refresh_provider(provider_id: &str) {
|
fn admin_provider_ops_balance_refresh_key(state: &AdminAppState<'_>, provider_id: &str) -> String {
|
||||||
|
let raw_key = format!("{ADMIN_PROVIDER_OPS_BALANCE_REFRESH_PREFIX}{provider_id}");
|
||||||
|
if let Some(runner) = state.redis_kv_runner() {
|
||||||
|
format!(
|
||||||
|
"{:p}:{}",
|
||||||
|
state.app(),
|
||||||
|
runner.keyspace().key(raw_key.as_str())
|
||||||
|
)
|
||||||
|
} else {
|
||||||
|
format!("{:p}:{raw_key}", state.app())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn finish_refresh_provider(refresh_key: &str) {
|
||||||
ADMIN_PROVIDER_OPS_REFRESHING_PROVIDERS
|
ADMIN_PROVIDER_OPS_REFRESHING_PROVIDERS
|
||||||
.lock()
|
.lock()
|
||||||
.await
|
.await
|
||||||
.remove(provider_id);
|
.remove(refresh_key);
|
||||||
}
|
}
|
||||||
|
|
||||||
fn admin_provider_ops_action_response(
|
fn admin_provider_ops_action_response(
|
||||||
|
|||||||
@@ -2,9 +2,8 @@ use crate::handlers::admin::provider::shared::payloads::AdminProviderCreateReque
|
|||||||
use crate::handlers::admin::provider::shared::support::{
|
use crate::handlers::admin::provider::shared::support::{
|
||||||
normalize_provider_billing_type, parse_optional_rfc3339_unix_secs,
|
normalize_provider_billing_type, parse_optional_rfc3339_unix_secs,
|
||||||
};
|
};
|
||||||
use crate::handlers::admin::provider::write::normalize::{
|
use crate::handlers::admin::provider::write::normalize::normalize_pool_advanced_config;
|
||||||
normalize_pool_advanced_config, normalize_provider_type_input,
|
use crate::handlers::admin::provider::write::normalize::normalize_provider_type_input;
|
||||||
};
|
|
||||||
use crate::handlers::admin::request::AdminAppState;
|
use crate::handlers::admin::request::AdminAppState;
|
||||||
use crate::handlers::admin::shared::normalize_json_object;
|
use crate::handlers::admin::shared::normalize_json_object;
|
||||||
use aether_data_contracts::repository::provider_catalog::StoredProviderCatalogProvider;
|
use aether_data_contracts::repository::provider_catalog::StoredProviderCatalogProvider;
|
||||||
|
|||||||
@@ -2,9 +2,8 @@ use crate::handlers::admin::provider::shared::payloads::AdminProviderUpdatePatch
|
|||||||
use crate::handlers::admin::provider::shared::support::{
|
use crate::handlers::admin::provider::shared::support::{
|
||||||
normalize_provider_billing_type, parse_optional_rfc3339_unix_secs,
|
normalize_provider_billing_type, parse_optional_rfc3339_unix_secs,
|
||||||
};
|
};
|
||||||
use crate::handlers::admin::provider::write::normalize::{
|
use crate::handlers::admin::provider::write::normalize::normalize_pool_advanced_config;
|
||||||
normalize_pool_advanced_config, normalize_provider_type_input,
|
use crate::handlers::admin::provider::write::normalize::normalize_provider_type_input;
|
||||||
};
|
|
||||||
use crate::handlers::admin::request::AdminAppState;
|
use crate::handlers::admin::request::AdminAppState;
|
||||||
use crate::handlers::admin::shared::normalize_json_object;
|
use crate::handlers::admin::shared::normalize_json_object;
|
||||||
use aether_data_contracts::repository::provider_catalog::StoredProviderCatalogProvider;
|
use aether_data_contracts::repository::provider_catalog::StoredProviderCatalogProvider;
|
||||||
|
|||||||
@@ -29,8 +29,8 @@ use super::{
|
|||||||
usage_cleanup_window, wallet_daily_usage_aggregation_target, AppState, DbMaintenanceRunSummary,
|
usage_cleanup_window, wallet_daily_usage_aggregation_target, AppState, DbMaintenanceRunSummary,
|
||||||
FailedPendingUsageRow, GatewayDataState, ProxyUpgradeRolloutProbeConfig, StalePendingUsageRow,
|
FailedPendingUsageRow, GatewayDataState, ProxyUpgradeRolloutProbeConfig, StalePendingUsageRow,
|
||||||
UsageCleanupSettings, DELETE_STALE_WALLET_DAILY_USAGE_LEDGERS_SQL,
|
UsageCleanupSettings, DELETE_STALE_WALLET_DAILY_USAGE_LEDGERS_SQL,
|
||||||
SELECT_STALE_PENDING_USAGE_BATCH_SQL, UPSERT_WALLET_DAILY_USAGE_LEDGER_SQL,
|
SELECT_STALE_PENDING_USAGE_BATCH_SQL, UPDATE_FAILED_VOID_STALE_USAGE_SQL,
|
||||||
UPDATE_FAILED_VOID_STALE_USAGE_SQL, USAGE_CLEANUP_HOUR, USAGE_CLEANUP_MINUTE,
|
UPSERT_WALLET_DAILY_USAGE_LEDGER_SQL, USAGE_CLEANUP_HOUR, USAGE_CLEANUP_MINUTE,
|
||||||
WALLET_DAILY_USAGE_AGGREGATION_HOUR, WALLET_DAILY_USAGE_AGGREGATION_MINUTE,
|
WALLET_DAILY_USAGE_AGGREGATION_HOUR, WALLET_DAILY_USAGE_AGGREGATION_MINUTE,
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -81,11 +81,8 @@ async fn spawn_pool_monitor_worker_skips_when_postgres_unavailable() {
|
|||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn wallet_daily_usage_queries_use_settlement_snapshots_for_wallet_identity() {
|
fn wallet_daily_usage_queries_use_settlement_snapshots_for_wallet_identity() {
|
||||||
assert!(
|
assert!(UPSERT_WALLET_DAILY_USAGE_LEDGER_SQL.contains("JOIN usage_settlement_snapshots"));
|
||||||
UPSERT_WALLET_DAILY_USAGE_LEDGER_SQL.contains("JOIN usage_settlement_snapshots")
|
assert!(UPSERT_WALLET_DAILY_USAGE_LEDGER_SQL.contains("usage_settlement_snapshots.wallet_id"));
|
||||||
);
|
|
||||||
assert!(UPSERT_WALLET_DAILY_USAGE_LEDGER_SQL
|
|
||||||
.contains("usage_settlement_snapshots.wallet_id"));
|
|
||||||
assert!(DELETE_STALE_WALLET_DAILY_USAGE_LEDGERS_SQL.contains("JOIN usage_settlement_snapshots"));
|
assert!(DELETE_STALE_WALLET_DAILY_USAGE_LEDGERS_SQL.contains("JOIN usage_settlement_snapshots"));
|
||||||
assert!(DELETE_STALE_WALLET_DAILY_USAGE_LEDGERS_SQL
|
assert!(DELETE_STALE_WALLET_DAILY_USAGE_LEDGERS_SQL
|
||||||
.contains("usage_settlement_snapshots.wallet_id = ledgers.wallet_id"));
|
.contains("usage_settlement_snapshots.wallet_id = ledgers.wallet_id"));
|
||||||
|
|||||||
@@ -23,10 +23,9 @@ use super::{
|
|||||||
NULLIFY_REQUEST_CANDIDATE_API_KEY_BATCH_SQL, NULLIFY_USAGE_API_KEY_BATCH_SQL,
|
NULLIFY_REQUEST_CANDIDATE_API_KEY_BATCH_SQL, NULLIFY_USAGE_API_KEY_BATCH_SQL,
|
||||||
SELECT_EXPIRED_ACTIVE_API_KEYS_SQL, SELECT_USAGE_BODY_COMPRESSION_BATCH_SQL,
|
SELECT_EXPIRED_ACTIVE_API_KEYS_SQL, SELECT_USAGE_BODY_COMPRESSION_BATCH_SQL,
|
||||||
SELECT_USAGE_BODY_COMPRESSION_ROW_SQL, SELECT_USAGE_HEADER_BATCH_SQL,
|
SELECT_USAGE_BODY_COMPRESSION_ROW_SQL, SELECT_USAGE_HEADER_BATCH_SQL,
|
||||||
SELECT_USAGE_LEGACY_BODY_REF_METADATA_BATCH_SQL,
|
SELECT_USAGE_LEGACY_BODY_REF_METADATA_BATCH_SQL, SELECT_USAGE_STALE_BODY_BATCH_SQL,
|
||||||
SELECT_USAGE_STALE_BODY_BATCH_SQL, UPDATE_USAGE_BODY_COMPRESSION_SQL,
|
UPDATE_USAGE_BODY_COMPRESSION_SQL, UPDATE_USAGE_REQUEST_METADATA_SQL,
|
||||||
UPDATE_USAGE_REQUEST_METADATA_SQL, UPSERT_USAGE_BODY_BLOB_SQL,
|
UPSERT_USAGE_BODY_BLOB_SQL, UPSERT_USAGE_HTTP_AUDIT_BODY_REFS_SQL,
|
||||||
UPSERT_USAGE_HTTP_AUDIT_BODY_REFS_SQL,
|
|
||||||
};
|
};
|
||||||
|
|
||||||
pub(super) async fn perform_usage_cleanup_once(
|
pub(super) async fn perform_usage_cleanup_once(
|
||||||
|
|||||||
@@ -111,7 +111,10 @@ async fn gateway_records_usage_for_execution_runtime_sync_when_runtime_enabled()
|
|||||||
.find_by_request_id("req-usage-sync-123")
|
.find_by_request_id("req-usage-sync-123")
|
||||||
.await
|
.await
|
||||||
.expect("usage lookup should succeed");
|
.expect("usage lookup should succeed");
|
||||||
if stored.is_some() {
|
if stored
|
||||||
|
.as_ref()
|
||||||
|
.is_some_and(|usage| usage.status == "completed")
|
||||||
|
{
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
||||||
|
|||||||
@@ -309,7 +309,10 @@ async fn gateway_truncates_deep_request_echo_for_local_openai_chat_sync_usage()
|
|||||||
.find_by_request_id("trace-openai-chat-local-report-sync-deep-123")
|
.find_by_request_id("trace-openai-chat-local-report-sync-deep-123")
|
||||||
.await
|
.await
|
||||||
.expect("usage lookup should succeed");
|
.expect("usage lookup should succeed");
|
||||||
if stored_usage.is_some() {
|
if stored_usage
|
||||||
|
.as_ref()
|
||||||
|
.is_some_and(|usage| usage.status == "completed")
|
||||||
|
{
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
||||||
|
|||||||
@@ -145,7 +145,10 @@ async fn gateway_settles_wallet_for_completed_execution_runtime_sync_usage() {
|
|||||||
.find_by_request_id("req-usage-wallet-sync-123")
|
.find_by_request_id("req-usage-wallet-sync-123")
|
||||||
.await
|
.await
|
||||||
.expect("usage lookup should succeed");
|
.expect("usage lookup should succeed");
|
||||||
if stored.is_some() {
|
if stored
|
||||||
|
.as_ref()
|
||||||
|
.is_some_and(|usage| usage.status == "completed")
|
||||||
|
{
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
||||||
@@ -154,11 +157,20 @@ async fn gateway_settles_wallet_for_completed_execution_runtime_sync_usage() {
|
|||||||
assert_eq!(stored.status, "completed");
|
assert_eq!(stored.status, "completed");
|
||||||
assert_eq!(stored.total_tokens, 1500);
|
assert_eq!(stored.total_tokens, 1500);
|
||||||
|
|
||||||
let wallet = wallet_repository
|
let mut wallet = None;
|
||||||
.find(WalletLookupKey::UserId("user-usage-sync-123"))
|
for _ in 0..50 {
|
||||||
.await
|
wallet = wallet_repository
|
||||||
.expect("wallet lookup should succeed")
|
.find(WalletLookupKey::UserId("user-usage-sync-123"))
|
||||||
.expect("wallet should exist");
|
.await
|
||||||
|
.expect("wallet lookup should succeed");
|
||||||
|
if wallet.as_ref().is_some_and(|wallet| {
|
||||||
|
wallet.balance < 10.0 || wallet.gift_balance < 2.0 || wallet.total_consumed > 0.0
|
||||||
|
}) {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
||||||
|
}
|
||||||
|
let wallet = wallet.expect("wallet should exist");
|
||||||
assert!(wallet.balance < 10.0 || wallet.gift_balance < 2.0);
|
assert!(wallet.balance < 10.0 || wallet.gift_balance < 2.0);
|
||||||
assert!(wallet.total_consumed > 0.0);
|
assert!(wallet.total_consumed > 0.0);
|
||||||
|
|
||||||
|
|||||||
@@ -100,6 +100,14 @@ impl InMemoryUsageReadRepository {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn usage_status_is_finalized(status: &str) -> bool {
|
||||||
|
matches!(status, "completed" | "failed" | "cancelled")
|
||||||
|
}
|
||||||
|
|
||||||
|
fn usage_status_is_lifecycle(status: &str) -> bool {
|
||||||
|
matches!(status, "pending" | "streaming")
|
||||||
|
}
|
||||||
|
|
||||||
#[async_trait]
|
#[async_trait]
|
||||||
impl UsageReadRepository for InMemoryUsageReadRepository {
|
impl UsageReadRepository for InMemoryUsageReadRepository {
|
||||||
async fn find_by_id(
|
async fn find_by_id(
|
||||||
@@ -439,6 +447,12 @@ impl UsageWriteRepository for InMemoryUsageReadRepository {
|
|||||||
})
|
})
|
||||||
.unwrap_or_default();
|
.unwrap_or_default();
|
||||||
let existing = by_request_id.get(&usage.request_id);
|
let existing = by_request_id.get(&usage.request_id);
|
||||||
|
if existing.is_some_and(|existing| {
|
||||||
|
usage_status_is_finalized(existing.status.as_str())
|
||||||
|
&& usage_status_is_lifecycle(usage.status.as_str())
|
||||||
|
}) {
|
||||||
|
return Ok(existing.expect("existing usage should be present").clone());
|
||||||
|
}
|
||||||
|
|
||||||
let request_metadata = usage
|
let request_metadata = usage
|
||||||
.request_metadata
|
.request_metadata
|
||||||
@@ -690,6 +704,158 @@ mod tests {
|
|||||||
assert_eq!(usage.total_tokens, 150);
|
assert_eq!(usage.total_tokens, 150);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn stale_pending_update_does_not_regress_finalized_usage() {
|
||||||
|
let repository = InMemoryUsageReadRepository::default();
|
||||||
|
repository
|
||||||
|
.upsert(UpsertUsageRecord {
|
||||||
|
request_id: "req-finalized-1".to_string(),
|
||||||
|
user_id: Some("user-1".to_string()),
|
||||||
|
api_key_id: Some("api-key-1".to_string()),
|
||||||
|
username: None,
|
||||||
|
api_key_name: None,
|
||||||
|
provider_name: "OpenAI".to_string(),
|
||||||
|
model: "gpt-5".to_string(),
|
||||||
|
target_model: None,
|
||||||
|
provider_id: Some("provider-1".to_string()),
|
||||||
|
provider_endpoint_id: Some("endpoint-1".to_string()),
|
||||||
|
provider_api_key_id: Some("provider-key-1".to_string()),
|
||||||
|
request_type: Some("chat".to_string()),
|
||||||
|
api_format: Some("openai:chat".to_string()),
|
||||||
|
api_family: Some("openai".to_string()),
|
||||||
|
endpoint_kind: Some("chat".to_string()),
|
||||||
|
endpoint_api_format: Some("openai:chat".to_string()),
|
||||||
|
provider_api_family: Some("openai".to_string()),
|
||||||
|
provider_endpoint_kind: Some("chat".to_string()),
|
||||||
|
has_format_conversion: Some(false),
|
||||||
|
is_stream: Some(false),
|
||||||
|
input_tokens: Some(3),
|
||||||
|
output_tokens: Some(5),
|
||||||
|
total_tokens: Some(8),
|
||||||
|
cache_creation_input_tokens: None,
|
||||||
|
cache_creation_ephemeral_5m_input_tokens: None,
|
||||||
|
cache_creation_ephemeral_1h_input_tokens: None,
|
||||||
|
cache_read_input_tokens: None,
|
||||||
|
cache_creation_cost_usd: None,
|
||||||
|
cache_read_cost_usd: None,
|
||||||
|
output_price_per_1m: None,
|
||||||
|
total_cost_usd: None,
|
||||||
|
actual_total_cost_usd: None,
|
||||||
|
status_code: Some(200),
|
||||||
|
error_message: None,
|
||||||
|
error_category: None,
|
||||||
|
response_time_ms: Some(45),
|
||||||
|
first_byte_time_ms: None,
|
||||||
|
status: "completed".to_string(),
|
||||||
|
billing_status: "pending".to_string(),
|
||||||
|
request_headers: None,
|
||||||
|
request_body: None,
|
||||||
|
request_body_ref: None,
|
||||||
|
provider_request_headers: None,
|
||||||
|
provider_request_body: None,
|
||||||
|
provider_request_body_ref: None,
|
||||||
|
response_headers: None,
|
||||||
|
response_body: None,
|
||||||
|
response_body_ref: None,
|
||||||
|
client_response_headers: None,
|
||||||
|
client_response_body: None,
|
||||||
|
client_response_body_ref: None,
|
||||||
|
candidate_id: None,
|
||||||
|
candidate_index: None,
|
||||||
|
key_name: None,
|
||||||
|
planner_kind: None,
|
||||||
|
route_family: None,
|
||||||
|
route_kind: None,
|
||||||
|
execution_path: None,
|
||||||
|
local_execution_runtime_miss_reason: None,
|
||||||
|
request_metadata: None,
|
||||||
|
finalized_at_unix_secs: Some(101),
|
||||||
|
created_at_unix_ms: Some(100),
|
||||||
|
updated_at_unix_secs: 101,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.expect("completed usage should upsert");
|
||||||
|
|
||||||
|
repository
|
||||||
|
.upsert(UpsertUsageRecord {
|
||||||
|
request_id: "req-finalized-1".to_string(),
|
||||||
|
user_id: Some("user-1".to_string()),
|
||||||
|
api_key_id: Some("api-key-1".to_string()),
|
||||||
|
username: None,
|
||||||
|
api_key_name: None,
|
||||||
|
provider_name: "OpenAI".to_string(),
|
||||||
|
model: "gpt-5".to_string(),
|
||||||
|
target_model: None,
|
||||||
|
provider_id: Some("provider-1".to_string()),
|
||||||
|
provider_endpoint_id: Some("endpoint-1".to_string()),
|
||||||
|
provider_api_key_id: Some("provider-key-1".to_string()),
|
||||||
|
request_type: Some("chat".to_string()),
|
||||||
|
api_format: Some("openai:chat".to_string()),
|
||||||
|
api_family: Some("openai".to_string()),
|
||||||
|
endpoint_kind: Some("chat".to_string()),
|
||||||
|
endpoint_api_format: Some("openai:chat".to_string()),
|
||||||
|
provider_api_family: Some("openai".to_string()),
|
||||||
|
provider_endpoint_kind: Some("chat".to_string()),
|
||||||
|
has_format_conversion: Some(false),
|
||||||
|
is_stream: Some(false),
|
||||||
|
input_tokens: None,
|
||||||
|
output_tokens: None,
|
||||||
|
total_tokens: None,
|
||||||
|
cache_creation_input_tokens: None,
|
||||||
|
cache_creation_ephemeral_5m_input_tokens: None,
|
||||||
|
cache_creation_ephemeral_1h_input_tokens: None,
|
||||||
|
cache_read_input_tokens: None,
|
||||||
|
cache_creation_cost_usd: None,
|
||||||
|
cache_read_cost_usd: None,
|
||||||
|
output_price_per_1m: None,
|
||||||
|
total_cost_usd: None,
|
||||||
|
actual_total_cost_usd: None,
|
||||||
|
status_code: None,
|
||||||
|
error_message: None,
|
||||||
|
error_category: None,
|
||||||
|
response_time_ms: None,
|
||||||
|
first_byte_time_ms: None,
|
||||||
|
status: "pending".to_string(),
|
||||||
|
billing_status: "pending".to_string(),
|
||||||
|
request_headers: None,
|
||||||
|
request_body: None,
|
||||||
|
request_body_ref: None,
|
||||||
|
provider_request_headers: None,
|
||||||
|
provider_request_body: None,
|
||||||
|
provider_request_body_ref: None,
|
||||||
|
response_headers: None,
|
||||||
|
response_body: None,
|
||||||
|
response_body_ref: None,
|
||||||
|
client_response_headers: None,
|
||||||
|
client_response_body: None,
|
||||||
|
client_response_body_ref: None,
|
||||||
|
candidate_id: None,
|
||||||
|
candidate_index: None,
|
||||||
|
key_name: None,
|
||||||
|
planner_kind: None,
|
||||||
|
route_family: None,
|
||||||
|
route_kind: None,
|
||||||
|
execution_path: None,
|
||||||
|
local_execution_runtime_miss_reason: None,
|
||||||
|
request_metadata: None,
|
||||||
|
finalized_at_unix_secs: None,
|
||||||
|
created_at_unix_ms: Some(100),
|
||||||
|
updated_at_unix_secs: 102,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.expect("stale pending usage should upsert");
|
||||||
|
|
||||||
|
let stored = repository
|
||||||
|
.find_by_request_id("req-finalized-1")
|
||||||
|
.await
|
||||||
|
.expect("usage lookup should succeed")
|
||||||
|
.expect("usage should exist");
|
||||||
|
assert_eq!(stored.status, "completed");
|
||||||
|
assert_eq!(stored.status_code, Some(200));
|
||||||
|
assert_eq!(stored.total_tokens, 8);
|
||||||
|
assert_eq!(stored.finalized_at_unix_secs, Some(101));
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn seed_hydrates_legacy_body_ref_metadata_into_typed_fields() {
|
async fn seed_hydrates_legacy_body_ref_metadata_into_typed_fields() {
|
||||||
let repository = InMemoryUsageReadRepository::seed(vec![StoredRequestUsageAudit {
|
let repository = InMemoryUsageReadRepository::seed(vec![StoredRequestUsageAudit {
|
||||||
|
|||||||
Reference in New Issue
Block a user