mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-01 17:00:21 +08:00
fix(data): 同步模型调用次数至 global_models 读模型
- 模型列表改为读取 global_models.usage_count,避免每次扫描 usage 明细表 - 通过历史 backfill 与 usage upsert delta 维护该读模型 - 详情页保留实时 facts 兜底,统一排除 pending/streaming 状态
This commit is contained in:
@@ -0,0 +1,25 @@
|
|||||||
|
WITH aggregated AS (
|
||||||
|
SELECT
|
||||||
|
usage.model,
|
||||||
|
COUNT(*)::BIGINT AS usage_count
|
||||||
|
FROM usage_billing_facts AS usage
|
||||||
|
WHERE usage.model IS NOT NULL
|
||||||
|
AND BTRIM(usage.model) <> ''
|
||||||
|
AND usage.status NOT IN ('pending', 'streaming')
|
||||||
|
GROUP BY usage.model
|
||||||
|
),
|
||||||
|
refreshed AS (
|
||||||
|
SELECT
|
||||||
|
gm.id,
|
||||||
|
COALESCE(aggregated.usage_count, 0)::BIGINT AS usage_count
|
||||||
|
FROM global_models AS gm
|
||||||
|
LEFT JOIN aggregated
|
||||||
|
ON aggregated.model = gm.name
|
||||||
|
)
|
||||||
|
UPDATE global_models AS gm
|
||||||
|
SET
|
||||||
|
usage_count = refreshed.usage_count,
|
||||||
|
updated_at = NOW()
|
||||||
|
FROM refreshed
|
||||||
|
WHERE gm.id = refreshed.id
|
||||||
|
AND COALESCE(gm.usage_count, 0) <> refreshed.usage_count;
|
||||||
@@ -277,7 +277,15 @@ mod tests {
|
|||||||
.into_iter()
|
.into_iter()
|
||||||
.map(|item| item.version)
|
.map(|item| item.version)
|
||||||
.collect::<Vec<_>>();
|
.collect::<Vec<_>>();
|
||||||
assert_eq!(versions, vec![20260422110000, 20260422120000]);
|
assert_eq!(
|
||||||
|
versions,
|
||||||
|
vec![
|
||||||
|
20260422110000,
|
||||||
|
20260422120000,
|
||||||
|
20260504120000,
|
||||||
|
20260505120000
|
||||||
|
]
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -289,7 +297,10 @@ mod tests {
|
|||||||
.into_iter()
|
.into_iter()
|
||||||
.map(|item| item.version)
|
.map(|item| item.version)
|
||||||
.collect::<Vec<_>>();
|
.collect::<Vec<_>>();
|
||||||
assert_eq!(versions, vec![20260422120000]);
|
assert_eq!(
|
||||||
|
versions,
|
||||||
|
vec![20260422120000, 20260504120000, 20260505120000]
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
@@ -531,6 +542,25 @@ mod tests {
|
|||||||
.await
|
.await
|
||||||
.expect("api key fixture should insert");
|
.expect("api key fixture should insert");
|
||||||
|
|
||||||
|
query(
|
||||||
|
r#"
|
||||||
|
INSERT INTO public.global_models (
|
||||||
|
id,
|
||||||
|
name,
|
||||||
|
display_name,
|
||||||
|
usage_count
|
||||||
|
) VALUES (
|
||||||
|
'global-model-backfill-1',
|
||||||
|
'gpt-4o-mini',
|
||||||
|
'GPT-4o Mini',
|
||||||
|
77
|
||||||
|
)
|
||||||
|
"#,
|
||||||
|
)
|
||||||
|
.execute(&pool)
|
||||||
|
.await
|
||||||
|
.expect("global model fixture should insert");
|
||||||
|
|
||||||
query(
|
query(
|
||||||
r#"
|
r#"
|
||||||
INSERT INTO public.usage (
|
INSERT INTO public.usage (
|
||||||
@@ -597,10 +627,11 @@ mod tests {
|
|||||||
let pending_before = pending_backfills(&pool)
|
let pending_before = pending_backfills(&pool)
|
||||||
.await
|
.await
|
||||||
.expect("pending backfills should load");
|
.expect("pending backfills should load");
|
||||||
assert_eq!(pending_before.len(), 3);
|
assert_eq!(pending_before.len(), 4);
|
||||||
assert_eq!(pending_before[0].version, 20260422110000);
|
assert_eq!(pending_before[0].version, 20260422110000);
|
||||||
assert_eq!(pending_before[1].version, 20260422120000);
|
assert_eq!(pending_before[1].version, 20260422120000);
|
||||||
assert_eq!(pending_before[2].version, 20260504120000);
|
assert_eq!(pending_before[2].version, 20260504120000);
|
||||||
|
assert_eq!(pending_before[3].version, 20260505120000);
|
||||||
|
|
||||||
run_backfills(&pool)
|
run_backfills(&pool)
|
||||||
.await
|
.await
|
||||||
@@ -618,7 +649,12 @@ mod tests {
|
|||||||
.expect("applied backfill versions should load");
|
.expect("applied backfill versions should load");
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
applied_versions,
|
applied_versions,
|
||||||
vec![20260422110000, 20260422120000, 20260504120000]
|
vec![
|
||||||
|
20260422110000,
|
||||||
|
20260422120000,
|
||||||
|
20260504120000,
|
||||||
|
20260505120000
|
||||||
|
]
|
||||||
);
|
);
|
||||||
|
|
||||||
let api_key_total_requests: i64 = query_scalar(
|
let api_key_total_requests: i64 = query_scalar(
|
||||||
@@ -629,6 +665,14 @@ mod tests {
|
|||||||
.expect("api key total requests should load");
|
.expect("api key total requests should load");
|
||||||
assert_eq!(api_key_total_requests, 1);
|
assert_eq!(api_key_total_requests, 1);
|
||||||
|
|
||||||
|
let global_model_usage_count: i64 = query_scalar(
|
||||||
|
"SELECT COALESCE(usage_count, 0)::BIGINT FROM public.global_models WHERE name = 'gpt-4o-mini'",
|
||||||
|
)
|
||||||
|
.fetch_one(&pool)
|
||||||
|
.await
|
||||||
|
.expect("global model usage count should load");
|
||||||
|
assert_eq!(global_model_usage_count, 1);
|
||||||
|
|
||||||
let api_key_total_tokens: i64 = query_scalar(
|
let api_key_total_tokens: i64 = query_scalar(
|
||||||
"SELECT COALESCE(total_tokens, 0)::BIGINT FROM public.api_keys WHERE id = 'api-key-backfill-1'",
|
"SELECT COALESCE(total_tokens, 0)::BIGINT FROM public.api_keys WHERE id = 'api-key-backfill-1'",
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -405,7 +405,7 @@ SELECT
|
|||||||
gm.config,
|
gm.config,
|
||||||
COALESCE(gm_stats.provider_count, 0) AS provider_count,
|
COALESCE(gm_stats.provider_count, 0) AS provider_count,
|
||||||
COALESCE(gm_stats.active_provider_count, 0) AS active_provider_count,
|
COALESCE(gm_stats.active_provider_count, 0) AS active_provider_count,
|
||||||
COALESCE(gm.usage_count, 0)::bigint AS usage_count,
|
COALESCE(usage_stats.usage_count, gm.usage_count, 0)::bigint AS usage_count,
|
||||||
EXTRACT(EPOCH FROM gm.created_at)::bigint AS created_at_unix_ms,
|
EXTRACT(EPOCH FROM gm.created_at)::bigint AS created_at_unix_ms,
|
||||||
EXTRACT(EPOCH FROM gm.updated_at)::bigint AS updated_at_unix_secs
|
EXTRACT(EPOCH FROM gm.updated_at)::bigint AS updated_at_unix_secs
|
||||||
FROM global_models gm
|
FROM global_models gm
|
||||||
@@ -423,6 +423,12 @@ LEFT JOIN (
|
|||||||
JOIN providers p ON p.id = m.provider_id
|
JOIN providers p ON p.id = m.provider_id
|
||||||
GROUP BY m.global_model_id
|
GROUP BY m.global_model_id
|
||||||
) gm_stats ON gm_stats.global_model_id = gm.id
|
) gm_stats ON gm_stats.global_model_id = gm.id
|
||||||
|
LEFT JOIN LATERAL (
|
||||||
|
SELECT NULLIF(COUNT(*), 0)::bigint AS usage_count
|
||||||
|
FROM usage_billing_facts AS usage
|
||||||
|
WHERE usage.model = gm.name
|
||||||
|
AND usage.status NOT IN ('pending', 'streaming')
|
||||||
|
) usage_stats ON TRUE
|
||||||
WHERE gm.id = $1
|
WHERE gm.id = $1
|
||||||
LIMIT 1
|
LIMIT 1
|
||||||
"#,
|
"#,
|
||||||
@@ -452,7 +458,7 @@ SELECT
|
|||||||
gm.config,
|
gm.config,
|
||||||
COALESCE(gm_stats.provider_count, 0) AS provider_count,
|
COALESCE(gm_stats.provider_count, 0) AS provider_count,
|
||||||
COALESCE(gm_stats.active_provider_count, 0) AS active_provider_count,
|
COALESCE(gm_stats.active_provider_count, 0) AS active_provider_count,
|
||||||
COALESCE(gm.usage_count, 0)::bigint AS usage_count,
|
COALESCE(usage_stats.usage_count, gm.usage_count, 0)::bigint AS usage_count,
|
||||||
EXTRACT(EPOCH FROM gm.created_at)::bigint AS created_at_unix_ms,
|
EXTRACT(EPOCH FROM gm.created_at)::bigint AS created_at_unix_ms,
|
||||||
EXTRACT(EPOCH FROM gm.updated_at)::bigint AS updated_at_unix_secs
|
EXTRACT(EPOCH FROM gm.updated_at)::bigint AS updated_at_unix_secs
|
||||||
FROM global_models gm
|
FROM global_models gm
|
||||||
@@ -470,6 +476,12 @@ LEFT JOIN (
|
|||||||
JOIN providers p ON p.id = m.provider_id
|
JOIN providers p ON p.id = m.provider_id
|
||||||
GROUP BY m.global_model_id
|
GROUP BY m.global_model_id
|
||||||
) gm_stats ON gm_stats.global_model_id = gm.id
|
) gm_stats ON gm_stats.global_model_id = gm.id
|
||||||
|
LEFT JOIN LATERAL (
|
||||||
|
SELECT NULLIF(COUNT(*), 0)::bigint AS usage_count
|
||||||
|
FROM usage_billing_facts AS usage
|
||||||
|
WHERE usage.model = gm.name
|
||||||
|
AND usage.status NOT IN ('pending', 'streaming')
|
||||||
|
) usage_stats ON TRUE
|
||||||
WHERE gm.name = $1
|
WHERE gm.name = $1
|
||||||
LIMIT 1
|
LIMIT 1
|
||||||
"#,
|
"#,
|
||||||
@@ -1205,7 +1217,10 @@ fn map_provider_active_global_model_row(
|
|||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{SqlxGlobalModelReadRepository, LIST_ADMIN_PROVIDER_MODELS_PREFIX};
|
use super::{
|
||||||
|
SqlxGlobalModelReadRepository, LIST_ADMIN_GLOBAL_MODELS_PREFIX,
|
||||||
|
LIST_ADMIN_PROVIDER_MODELS_PREFIX,
|
||||||
|
};
|
||||||
use crate::postgres::{PostgresPoolConfig, PostgresPoolFactory};
|
use crate::postgres::{PostgresPoolConfig, PostgresPoolFactory};
|
||||||
|
|
||||||
const ADMIN_PROVIDER_MODEL_REQUIRED_COLUMNS: &[&str] = &[
|
const ADMIN_PROVIDER_MODEL_REQUIRED_COLUMNS: &[&str] = &[
|
||||||
@@ -1241,6 +1256,64 @@ mod tests {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn admin_global_model_sql_projects_usage_count_from_billing_facts() {
|
||||||
|
let source = include_str!("sql.rs");
|
||||||
|
|
||||||
|
assert!(
|
||||||
|
!LIST_ADMIN_GLOBAL_MODELS_PREFIX.contains("usage_billing_facts"),
|
||||||
|
"admin global model list should read the maintained usage_count field"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
source.contains("WHERE usage.model = gm.name"),
|
||||||
|
"admin global model detail lookups should count usage facts for the selected model"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
source.contains("AND usage.status NOT IN ('pending', 'streaming')"),
|
||||||
|
"admin global model detail usage_count should use the maintained read-model status scope"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
source.contains("COALESCE(usage_stats.usage_count, gm.usage_count, 0)::bigint AS usage_count"),
|
||||||
|
"admin global model usage_count should prefer actual usage facts with stored-count fallback"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn admin_global_model_list_sql_counts_usage_without_full_fact_aggregate() {
|
||||||
|
assert!(
|
||||||
|
LIST_ADMIN_GLOBAL_MODELS_PREFIX.contains("COALESCE(gm.usage_count, 0)::bigint AS usage_count"),
|
||||||
|
"admin global model list should read maintained usage_count without per-request fact scans"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
!LIST_ADMIN_GLOBAL_MODELS_PREFIX.contains("GROUP BY usage.model"),
|
||||||
|
"admin global model list must not aggregate the full usage fact table before pagination"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn global_model_usage_count_read_model_has_backfill_and_delta_maintenance() {
|
||||||
|
let backfill_sql =
|
||||||
|
include_str!("../../../backfills/20260505120000_rebuild_global_model_usage_count.sql");
|
||||||
|
let delta_sql = include_str!("../usage/sql/queries/apply_global_model_usage_delta_sql.sql");
|
||||||
|
|
||||||
|
assert!(
|
||||||
|
backfill_sql.contains("UPDATE global_models AS gm"),
|
||||||
|
"global model usage_count backfill should refresh the read model"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
backfill_sql.contains("FROM usage_billing_facts AS usage"),
|
||||||
|
"global model usage_count backfill should rebuild from canonical usage facts"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
delta_sql.contains("UPDATE global_models"),
|
||||||
|
"usage writes should maintain the global model usage_count read model"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
delta_sql.contains("WHERE name = $1"),
|
||||||
|
"global model usage_count delta should target models by canonical model name"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn repository_constructs_from_lazy_pool() {
|
async fn repository_constructs_from_lazy_pool() {
|
||||||
let factory = PostgresPoolFactory::new(PostgresPoolConfig {
|
let factory = PostgresPoolFactory::new(PostgresPoolConfig {
|
||||||
|
|||||||
@@ -89,6 +89,41 @@ impl ApiKeyUsageDelta {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq, Default)]
|
||||||
|
pub(crate) struct ModelUsageContribution {
|
||||||
|
pub model: String,
|
||||||
|
pub request_count: i64,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq, Default)]
|
||||||
|
pub(crate) struct ModelUsageDelta {
|
||||||
|
pub request_count: i64,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ModelUsageDelta {
|
||||||
|
pub(crate) fn between(before: &ModelUsageContribution, after: &ModelUsageContribution) -> Self {
|
||||||
|
Self {
|
||||||
|
request_count: after.request_count - before.request_count,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn addition(after: &ModelUsageContribution) -> Self {
|
||||||
|
Self {
|
||||||
|
request_count: after.request_count,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn removal(before: &ModelUsageContribution) -> Self {
|
||||||
|
Self {
|
||||||
|
request_count: -before.request_count,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn is_noop(&self) -> bool {
|
||||||
|
self.request_count == 0
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, PartialEq, Default)]
|
#[derive(Debug, Clone, PartialEq, Default)]
|
||||||
pub(crate) struct ProviderApiKeyUsageContribution {
|
pub(crate) struct ProviderApiKeyUsageContribution {
|
||||||
pub key_id: String,
|
pub key_id: String,
|
||||||
@@ -265,6 +300,23 @@ pub(crate) fn provider_api_key_usage_contribution(
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(crate) fn model_usage_contribution(
|
||||||
|
usage: &StoredRequestUsageAudit,
|
||||||
|
) -> Option<ModelUsageContribution> {
|
||||||
|
if matches!(usage.status.as_str(), "pending" | "streaming") {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
let model = usage.model.trim();
|
||||||
|
if model.is_empty() {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
|
||||||
|
Some(ModelUsageContribution {
|
||||||
|
model: model.to_string(),
|
||||||
|
request_count: 1,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) fn api_key_usage_contribution(
|
pub(crate) fn api_key_usage_contribution(
|
||||||
usage: &StoredRequestUsageAudit,
|
usage: &StoredRequestUsageAudit,
|
||||||
) -> Option<ApiKeyUsageContribution> {
|
) -> Option<ApiKeyUsageContribution> {
|
||||||
@@ -292,9 +344,10 @@ pub(crate) fn api_key_usage_contribution(
|
|||||||
mod tests {
|
mod tests {
|
||||||
use super::{
|
use super::{
|
||||||
api_key_usage_contribution, incoming_usage_can_recover_terminal_failure,
|
api_key_usage_contribution, incoming_usage_can_recover_terminal_failure,
|
||||||
provider_api_key_usage_contribution, provider_api_key_usage_is_error,
|
model_usage_contribution, provider_api_key_usage_contribution,
|
||||||
provider_api_key_usage_is_success, strip_deprecated_usage_display_fields,
|
provider_api_key_usage_is_error, provider_api_key_usage_is_success,
|
||||||
usage_can_recover_terminal_failure, StoredRequestUsageAudit, UpsertUsageRecord,
|
strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure, ModelUsageDelta,
|
||||||
|
StoredRequestUsageAudit, UpsertUsageRecord,
|
||||||
};
|
};
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -570,4 +623,75 @@ mod tests {
|
|||||||
assert_eq!(contribution.total_cost_usd, 0.25);
|
assert_eq!(contribution.total_cost_usd, 0.25);
|
||||||
assert_eq!(contribution.last_used_at_unix_secs, Some(123));
|
assert_eq!(contribution.last_used_at_unix_secs, Some(123));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn model_usage_contribution_tracks_terminal_requests_only() {
|
||||||
|
let completed = StoredRequestUsageAudit::new(
|
||||||
|
"usage-1".to_string(),
|
||||||
|
"request-1".to_string(),
|
||||||
|
Some("user-1".to_string()),
|
||||||
|
Some("api-key-1".to_string()),
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
"OpenAI".to_string(),
|
||||||
|
" gpt-5.5 ".to_string(),
|
||||||
|
None,
|
||||||
|
Some("provider-1".to_string()),
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
false,
|
||||||
|
false,
|
||||||
|
12,
|
||||||
|
8,
|
||||||
|
20,
|
||||||
|
0.25,
|
||||||
|
0.25,
|
||||||
|
Some(200),
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
Some(120),
|
||||||
|
None,
|
||||||
|
"completed".to_string(),
|
||||||
|
"settled".to_string(),
|
||||||
|
123,
|
||||||
|
124,
|
||||||
|
Some(125),
|
||||||
|
)
|
||||||
|
.expect("usage should build");
|
||||||
|
let contribution =
|
||||||
|
model_usage_contribution(&completed).expect("completed usage should count");
|
||||||
|
assert_eq!(contribution.model, "gpt-5.5");
|
||||||
|
assert_eq!(contribution.request_count, 1);
|
||||||
|
|
||||||
|
let mut streaming = completed.clone();
|
||||||
|
streaming.status = "streaming".to_string();
|
||||||
|
assert!(model_usage_contribution(&streaming).is_none());
|
||||||
|
|
||||||
|
let mut pending = completed;
|
||||||
|
pending.status = "pending".to_string();
|
||||||
|
assert!(model_usage_contribution(&pending).is_none());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn model_usage_delta_handles_model_changes() {
|
||||||
|
let before = super::ModelUsageContribution {
|
||||||
|
model: "gpt-5.4".to_string(),
|
||||||
|
request_count: 1,
|
||||||
|
};
|
||||||
|
let after = super::ModelUsageContribution {
|
||||||
|
model: "gpt-5.5".to_string(),
|
||||||
|
request_count: 1,
|
||||||
|
};
|
||||||
|
|
||||||
|
assert_eq!(ModelUsageDelta::removal(&before).request_count, -1);
|
||||||
|
assert_eq!(ModelUsageDelta::addition(&after).request_count, 1);
|
||||||
|
assert!(ModelUsageDelta::between(&before, &before).is_noop());
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -36,7 +36,8 @@ use uuid::Uuid;
|
|||||||
|
|
||||||
use super::{
|
use super::{
|
||||||
api_key_usage_contribution, incoming_usage_can_recover_terminal_failure,
|
api_key_usage_contribution, incoming_usage_can_recover_terminal_failure,
|
||||||
provider_api_key_usage_contribution, strip_deprecated_usage_display_fields, ApiKeyUsageDelta,
|
model_usage_contribution, provider_api_key_usage_contribution,
|
||||||
|
strip_deprecated_usage_display_fields, ApiKeyUsageDelta, ModelUsageDelta,
|
||||||
ProviderApiKeyUsageDelta, StoredProviderApiKeyUsageSummary, StoredProviderUsageSummary,
|
ProviderApiKeyUsageDelta, StoredProviderApiKeyUsageSummary, StoredProviderUsageSummary,
|
||||||
StoredRequestUsageAudit, StoredUsageDailySummary, UpsertUsageRecord, UsageAuditListQuery,
|
StoredRequestUsageAudit, StoredUsageDailySummary, UpsertUsageRecord, UsageAuditListQuery,
|
||||||
UsageDailyHeatmapQuery, UsageReadRepository, UsageWriteRepository,
|
UsageDailyHeatmapQuery, UsageReadRepository, UsageWriteRepository,
|
||||||
@@ -1219,6 +1220,9 @@ const REBUILD_API_KEY_USAGE_STATS_SQL: &str =
|
|||||||
const APPLY_PROVIDER_API_KEY_USAGE_DELTA_SQL: &str =
|
const APPLY_PROVIDER_API_KEY_USAGE_DELTA_SQL: &str =
|
||||||
include_str!("queries/apply_provider_api_key_usage_delta_sql.sql");
|
include_str!("queries/apply_provider_api_key_usage_delta_sql.sql");
|
||||||
|
|
||||||
|
const APPLY_GLOBAL_MODEL_USAGE_DELTA_SQL: &str =
|
||||||
|
include_str!("queries/apply_global_model_usage_delta_sql.sql");
|
||||||
|
|
||||||
const RESET_PROVIDER_API_KEY_USAGE_STATS_SQL: &str =
|
const RESET_PROVIDER_API_KEY_USAGE_STATS_SQL: &str =
|
||||||
include_str!("queries/reset_provider_api_key_usage_stats_sql.sql");
|
include_str!("queries/reset_provider_api_key_usage_stats_sql.sql");
|
||||||
|
|
||||||
@@ -7268,6 +7272,40 @@ ORDER BY "usage".user_id ASC
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let before_model_contribution =
|
||||||
|
previous_usage.as_ref().and_then(model_usage_contribution);
|
||||||
|
let after_model_contribution = model_usage_contribution(&stored);
|
||||||
|
match (
|
||||||
|
before_model_contribution.as_ref(),
|
||||||
|
after_model_contribution.as_ref(),
|
||||||
|
) {
|
||||||
|
(Some(before), Some(after)) if before.model == after.model => {
|
||||||
|
let delta = ModelUsageDelta::between(before, after);
|
||||||
|
apply_global_model_usage_delta_in_tx(tx, before.model.as_str(), &delta)
|
||||||
|
.await?;
|
||||||
|
}
|
||||||
|
_ => {
|
||||||
|
if let Some(before) = before_model_contribution.as_ref() {
|
||||||
|
let delta = ModelUsageDelta::removal(before);
|
||||||
|
apply_global_model_usage_delta_in_tx(
|
||||||
|
tx,
|
||||||
|
before.model.as_str(),
|
||||||
|
&delta,
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
}
|
||||||
|
if let Some(after) = after_model_contribution.as_ref() {
|
||||||
|
let delta = ModelUsageDelta::addition(after);
|
||||||
|
apply_global_model_usage_delta_in_tx(
|
||||||
|
tx,
|
||||||
|
after.model.as_str(),
|
||||||
|
&delta,
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Ok(stored)
|
Ok(stored)
|
||||||
}) as BoxFuture<'_, Result<StoredRequestUsageAudit, DataLayerError>>
|
}) as BoxFuture<'_, Result<StoredRequestUsageAudit, DataLayerError>>
|
||||||
})
|
})
|
||||||
@@ -7686,6 +7724,33 @@ async fn apply_provider_api_key_usage_delta_in_tx(
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn apply_global_model_usage_delta_in_tx(
|
||||||
|
tx: &mut sqlx::Transaction<'_, Postgres>,
|
||||||
|
model: &str,
|
||||||
|
delta: &ModelUsageDelta,
|
||||||
|
) -> Result<(), DataLayerError> {
|
||||||
|
let model = model.trim();
|
||||||
|
if model.is_empty() {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
if delta.is_noop() {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
|
||||||
|
sqlx::query(APPLY_GLOBAL_MODEL_USAGE_DELTA_SQL)
|
||||||
|
.bind(model)
|
||||||
|
.bind(i32::try_from(delta.request_count).map_err(|_| {
|
||||||
|
DataLayerError::UnexpectedValue(format!(
|
||||||
|
"global_models.usage_count delta exceeds i32: {}",
|
||||||
|
delta.request_count
|
||||||
|
))
|
||||||
|
})?)
|
||||||
|
.execute(&mut **tx)
|
||||||
|
.await
|
||||||
|
.map_postgres_err()?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
// Build the usage read model from the split storage layout.
|
// Build the usage read model from the split storage layout.
|
||||||
//
|
//
|
||||||
// Query projections already prefer the newer audit/snapshot owners and only fall back to
|
// Query projections already prefer the newer audit/snapshot owners and only fall back to
|
||||||
|
|||||||
@@ -0,0 +1,5 @@
|
|||||||
|
UPDATE global_models
|
||||||
|
SET
|
||||||
|
usage_count = GREATEST(COALESCE(usage_count, 0) + $2, 0),
|
||||||
|
updated_at = NOW()
|
||||||
|
WHERE name = $1
|
||||||
Reference in New Issue
Block a user