mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-07 01:47:47 +08:00
feat: unify user analytics and optimize overview aggregation
Merge user accounts and usage reporting into one page with a combined ranking and account table, shared precise time ranges, and simpler range labels. Parse overview metadata once through a schema-only view migration and disable JIT locally for bucket rebuilds. Preserve automatic backfills. Add redacted OAuth refresh diagnostics, bucket failure context, and regression coverage. Resolve strict Clippy warnings.
This commit is contained in:
+78
@@ -0,0 +1,78 @@
|
||||
-- Replace only the read model; existing facts, queues and projections are untouched.
|
||||
-- Preserve billing/attribution semantics while avoiding repeated JSON parsing.
|
||||
CREATE OR REPLACE VIEW public.usage_analytics_facts_v1 AS
|
||||
SELECT u.request_id, COALESCE(u.id, u.request_id) AS id, u.created_at,
|
||||
CASE WHEN identity.owner_id IS NOT NULL AND identity.is_standalone=false THEN identity.owner_id END AS actor_user_id,
|
||||
identity.owner_id AS credential_owner_id,
|
||||
CASE WHEN identity.owner_id IS NULL THEN 'unknown' WHEN identity.is_standalone THEN 'standalone'
|
||||
WHEN NOT identity.is_standalone THEN 'employee' ELSE 'unknown' END AS attribution_kind,
|
||||
CASE WHEN identity.owner_id IS NULL THEN 'unknown' WHEN identity.is_standalone THEN 'standalone_key'
|
||||
WHEN NOT identity.is_standalone THEN 'user_account' ELSE 'unknown' END AS attribution_source,
|
||||
COALESCE(a.record_kind, 'request') AS record_kind, a.parent_request_id,
|
||||
u.api_key_id, u.model, u.target_model, u.provider_id, u.provider_name,
|
||||
u.api_format, u.endpoint_kind, u.request_type, u.is_stream, u.has_format_conversion,
|
||||
u.status, u.status_code, u.error_category, u.failure_origin, u.failure_stage, u.failure_reason,
|
||||
u.failure_schema_version, u.response_time_ms, u.first_byte_time_ms,
|
||||
COALESCE(s.billing_status, u.billing_status) AS settlement_status,
|
||||
COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb AS usage_available,
|
||||
COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
AND (s.billing_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled') AS pricing_available,
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
THEN b.input_tokens END AS input_tokens,
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
THEN b.output_tokens END AS output_tokens,
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
THEN b.total_tokens END AS total_tokens,
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
THEN b.cache_read_input_tokens END AS cache_read_input_tokens,
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
THEN b.cache_creation_input_tokens END AS cache_creation_input_tokens,
|
||||
CASE WHEN COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
AND (s.billing_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled')
|
||||
THEN round(COALESCE(s.billing_total_cost_usd::numeric, u.total_cost_usd::numeric), 8) END AS rated_amount,
|
||||
CASE WHEN COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
AND (s.billing_actual_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled')
|
||||
THEN round(COALESCE(s.billing_actual_total_cost_usd::numeric, u.actual_total_cost_usd::numeric), 8) END AS billable_amount,
|
||||
s.quota_covered_amount_usd AS quota_covered_amount,
|
||||
s.wallet_consumed_amount_usd AS wallet_consumed_amount,
|
||||
s.wallet_debit_amount_usd AS wallet_debit_amount,
|
||||
s.wallet_recharge_debit_usd AS wallet_recharge_debit_amount,
|
||||
s.wallet_gift_debit_usd AS wallet_gift_debit_amount,
|
||||
s.wallet_overdraft_usd AS wallet_overdraft_amount,
|
||||
s.allocation_status, s.finalized_at AS settled_at,
|
||||
CASE WHEN s.billing_total_cost_usd IS NOT NULL THEN 'settlement_snapshot' ELSE 'legacy_float' END AS amount_source,
|
||||
b.upstream_is_stream,
|
||||
CASE WHEN metadata.value #>> '{analytics_measurement,source}' IN ('reported','estimated','mixed')
|
||||
THEN metadata.value #>> '{analytics_measurement,source}' ELSE 'unknown' END AS token_source,
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_read_cost_usd IS NOT NULL
|
||||
THEN round(s.input_price_per_1m::numeric * b.cache_read_input_tokens::numeric / 1000000,8) END AS cache_estimated_full_cost_amount,
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_read_cost_usd IS NOT NULL
|
||||
THEN round(s.billing_cache_read_cost_usd::numeric,8) END AS cache_read_cost_amount,
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_creation_cost_usd IS NOT NULL
|
||||
THEN round(s.billing_cache_creation_cost_usd::numeric,8) END AS cache_creation_cost_amount
|
||||
FROM public.usage u
|
||||
-- OFFSET 0 keeps this projection from being flattened: large metadata is
|
||||
-- detoasted and parsed once per request, rather than once per metric expression.
|
||||
CROSS JOIN LATERAL (SELECT u.request_metadata::jsonb AS value OFFSET 0) metadata
|
||||
LEFT JOIN public.usage_settlement_snapshots s USING (request_id)
|
||||
LEFT JOIN public.usage_attribution_snapshots a USING (request_id)
|
||||
JOIN public.usage_billing_facts b USING (request_id)
|
||||
LEFT JOIN public.api_keys k ON k.id=u.api_key_id
|
||||
CROSS JOIN LATERAL (
|
||||
SELECT CASE WHEN a.request_id IS NOT NULL THEN a.credential_owner_id
|
||||
WHEN EXISTS (SELECT 1 FROM public.users WHERE id=u.user_id AND NOT is_deleted) THEN u.user_id END AS owner_id,
|
||||
COALESCE(k.is_standalone,
|
||||
CASE WHEN jsonb_typeof(metadata.value #> '{analytics_attribution,is_standalone}')='boolean'
|
||||
THEN (metadata.value #>> '{analytics_attribution,is_standalone}')::boolean END,
|
||||
CASE WHEN jsonb_typeof(metadata.value->'api_key_is_standalone')='boolean'
|
||||
THEN (metadata.value->>'api_key_is_standalone')::boolean END,
|
||||
CASE WHEN a.attribution_source='user_account' THEN false
|
||||
WHEN a.attribution_source='standalone_key' THEN true END,
|
||||
CASE WHEN u.api_key_id IS NULL THEN false END) AS is_standalone
|
||||
) identity;
|
||||
@@ -88,6 +88,7 @@ impl SqlxUsageReadRepository {
|
||||
for row in rows {
|
||||
let granularity: String = row.try_get("granularity").map_postgres_err()?;
|
||||
let bucket: DateTime<Utc> = row.try_get("bucket_start").map_postgres_err()?;
|
||||
let started = std::time::Instant::now();
|
||||
match self
|
||||
.rebuild_merged_overview_bucket(&granularity, bucket)
|
||||
.await
|
||||
@@ -95,8 +96,20 @@ impl SqlxUsageReadRepository {
|
||||
Ok(true) => published += 1,
|
||||
Ok(false) => {}
|
||||
Err(error) => {
|
||||
let error_detail = error.to_string().chars().take(500).collect::<String>();
|
||||
tracing::warn!(
|
||||
event_name = "overview_bucket_rebuild_failed",
|
||||
log_type = "ops",
|
||||
projection_version = "overview-v2",
|
||||
granularity = %granularity,
|
||||
bucket_start = %bucket,
|
||||
elapsed_ms = started.elapsed().as_millis() as u64,
|
||||
retry_after_secs = 600,
|
||||
error = %error_detail,
|
||||
"overview bucket rebuild failed; retry deferred"
|
||||
);
|
||||
sqlx::query("UPDATE stats_bucket_state SET last_error=$3,last_failed_at=NOW() WHERE projection_version='overview-v2' AND granularity=$1 AND bucket_start=$2")
|
||||
.bind(&granularity).bind(bucket).bind(error.to_string().chars().take(500).collect::<String>())
|
||||
.bind(&granularity).bind(bucket).bind(error_detail)
|
||||
.execute(&self.pool).await.map_postgres_err()?;
|
||||
}
|
||||
}
|
||||
@@ -144,6 +157,13 @@ impl SqlxUsageReadRepository {
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_postgres_err()?;
|
||||
// These bounded aggregates have many expressions but run only once per
|
||||
// bucket. JIT compilation consumes a significant part of their timeout.
|
||||
// Keep the setting transaction-local so other pool users retain theirs.
|
||||
sqlx::query("SET LOCAL jit = off")
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_postgres_err()?;
|
||||
let acquired: bool =
|
||||
sqlx::query_scalar("SELECT pg_try_advisory_xact_lock(hashtextextended($1, 19))")
|
||||
.bind(format!("overview-v2:{granularity}:{}", bucket.timestamp()))
|
||||
|
||||
@@ -33,12 +33,17 @@ The schema workspace has three normal source areas:
|
||||
| `drivers/postgres/` | Current maintenance fragments for executable SQL. | Edit only for deployment compatibility, ordering, or generator gaps. |
|
||||
| `bootstrap/postgres/` | Source fragments for the Postgres empty-database bootstrap snapshot. | Edit here when the bootstrap snapshot changes, then rebuild `aether-data` so `build.rs` regenerates the embedded snapshot. |
|
||||
|
||||
Everything else is output:
|
||||
Generated schema and the composed baseline are outputs:
|
||||
|
||||
| Path | Role | Edit policy |
|
||||
|---|---|---|
|
||||
| `generated/postgres/` | Machine-written SQL emitted from `logical/*.toml` for audit and drift detection. | Do not edit; regenerate with `compose_schema.sh generate`. |
|
||||
| `../../adapters/postgres/migrations/` | Runtime SQL artifacts embedded by each database adapter. | Regenerate through `compose_schema.sh compose`; do not edit independently. |
|
||||
| `../../adapters/postgres/migrations/20260403000000_baseline.sql` | Composed PostgreSQL baseline embedded by the adapter. | Regenerate through `compose_schema.sh compose`; do not edit independently. |
|
||||
|
||||
Later incremental migrations are maintained directly under
|
||||
`../../adapters/postgres/migrations/`; they have no compose target. Add a new
|
||||
version for an upgrade and preserve the checksums of already-applied scripts.
|
||||
Keep any corresponding maintained bootstrap definitions in sync.
|
||||
|
||||
`generated/**` is deliberately checked in so reviews and CI can see exactly
|
||||
what the logical schema compiler emits for each driver. It is not a fourth SQL
|
||||
@@ -131,6 +136,7 @@ cannot run them inside a transaction.
|
||||
| `20260921020000` | Add the attribution-owner lookup index concurrently on existing databases. |
|
||||
| `20260921020100` | Create the usage metadata actor index concurrently. |
|
||||
| `20261001000000` | Remove deleted-user attribution from dashboard activity on future user deletion; schema-only upgrade without rewriting historical rows. |
|
||||
| `20261004000000` | Parse request metadata once per overview fact; replace only the view definition without rewriting facts or statistics. |
|
||||
|
||||
Do not remove an applied migration after folding its changes into an earlier
|
||||
schema definition. Existing databases retain its version in `_sqlx_migrations`
|
||||
|
||||
@@ -215,23 +215,23 @@ SELECT u.request_id, COALESCE(u.id, u.request_id) AS id, u.created_at,
|
||||
u.status, u.status_code, u.error_category, u.failure_origin, u.failure_stage, u.failure_reason,
|
||||
u.failure_schema_version, u.response_time_ms, u.first_byte_time_ms,
|
||||
COALESCE(s.billing_status, u.billing_status) AS settlement_status,
|
||||
COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb AS usage_available,
|
||||
COALESCE(u.request_metadata::jsonb->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb AS usage_available,
|
||||
COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
AND (s.billing_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled') AS pricing_available,
|
||||
CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
THEN b.input_tokens END AS input_tokens,
|
||||
CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
THEN b.output_tokens END AS output_tokens,
|
||||
CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
THEN b.total_tokens END AS total_tokens,
|
||||
CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
THEN b.cache_read_input_tokens END AS cache_read_input_tokens,
|
||||
CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
THEN b.cache_creation_input_tokens END AS cache_creation_input_tokens,
|
||||
CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
CASE WHEN COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
AND (s.billing_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled')
|
||||
THEN round(COALESCE(s.billing_total_cost_usd::numeric, u.total_cost_usd::numeric), 8) END AS rated_amount,
|
||||
CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
CASE WHEN COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb
|
||||
AND (s.billing_actual_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled')
|
||||
THEN round(COALESCE(s.billing_actual_total_cost_usd::numeric, u.actual_total_cost_usd::numeric), 8) END AS billable_amount,
|
||||
s.quota_covered_amount_usd AS quota_covered_amount,
|
||||
@@ -243,21 +243,24 @@ SELECT u.request_id, COALESCE(u.id, u.request_id) AS id, u.created_at,
|
||||
s.allocation_status, s.finalized_at AS settled_at,
|
||||
CASE WHEN s.billing_total_cost_usd IS NOT NULL THEN 'settlement_snapshot' ELSE 'legacy_float' END AS amount_source,
|
||||
b.upstream_is_stream,
|
||||
CASE WHEN u.request_metadata #>> '{analytics_measurement,source}' IN ('reported','estimated','mixed')
|
||||
THEN u.request_metadata #>> '{analytics_measurement,source}' ELSE 'unknown' END AS token_source,
|
||||
CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND COALESCE(u.request_metadata::jsonb->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb
|
||||
CASE WHEN metadata.value #>> '{analytics_measurement,source}' IN ('reported','estimated','mixed')
|
||||
THEN metadata.value #>> '{analytics_measurement,source}' ELSE 'unknown' END AS token_source,
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_read_cost_usd IS NOT NULL
|
||||
THEN round(s.input_price_per_1m::numeric * b.cache_read_input_tokens::numeric / 1000000,8) END AS cache_estimated_full_cost_amount,
|
||||
CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND COALESCE(u.request_metadata::jsonb->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_read_cost_usd IS NOT NULL
|
||||
THEN round(s.billing_cache_read_cost_usd::numeric,8) END AS cache_read_cost_amount,
|
||||
CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND COALESCE(u.request_metadata::jsonb->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb
|
||||
CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb
|
||||
AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_creation_cost_usd IS NOT NULL
|
||||
THEN round(s.billing_cache_creation_cost_usd::numeric,8) END AS cache_creation_cost_amount
|
||||
FROM public.usage u
|
||||
-- OFFSET 0 keeps this projection from being flattened: large metadata is
|
||||
-- detoasted and parsed once per request, rather than once per metric expression.
|
||||
CROSS JOIN LATERAL (SELECT u.request_metadata::jsonb AS value OFFSET 0) metadata
|
||||
LEFT JOIN public.usage_settlement_snapshots s USING (request_id)
|
||||
LEFT JOIN public.usage_attribution_snapshots a USING (request_id)
|
||||
JOIN public.usage_billing_facts b USING (request_id)
|
||||
@@ -266,10 +269,10 @@ CROSS JOIN LATERAL (
|
||||
SELECT CASE WHEN a.request_id IS NOT NULL THEN a.credential_owner_id
|
||||
WHEN EXISTS (SELECT 1 FROM public.users WHERE id=u.user_id AND NOT is_deleted) THEN u.user_id END AS owner_id,
|
||||
COALESCE(k.is_standalone,
|
||||
CASE WHEN jsonb_typeof(u.request_metadata::jsonb #> '{analytics_attribution,is_standalone}')='boolean'
|
||||
THEN (u.request_metadata #>> '{analytics_attribution,is_standalone}')::boolean END,
|
||||
CASE WHEN jsonb_typeof(u.request_metadata::jsonb->'api_key_is_standalone')='boolean'
|
||||
THEN (u.request_metadata->>'api_key_is_standalone')::boolean END,
|
||||
CASE WHEN jsonb_typeof(metadata.value #> '{analytics_attribution,is_standalone}')='boolean'
|
||||
THEN (metadata.value #>> '{analytics_attribution,is_standalone}')::boolean END,
|
||||
CASE WHEN jsonb_typeof(metadata.value->'api_key_is_standalone')='boolean'
|
||||
THEN (metadata.value->>'api_key_is_standalone')::boolean END,
|
||||
CASE WHEN a.attribution_source='user_account' THEN false
|
||||
WHEN a.attribution_source='standalone_key' THEN true END,
|
||||
CASE WHEN u.api_key_id IS NULL THEN false END) AS is_standalone
|
||||
|
||||
@@ -31,6 +31,7 @@ mod dashboard_user_anonymization;
|
||||
mod legacy_overview_upgrade;
|
||||
mod migration_deadlines;
|
||||
mod overview_dirty_events;
|
||||
mod overview_fact_metadata;
|
||||
mod overview_migration_safety;
|
||||
mod policy_nulls;
|
||||
mod provider_expenses;
|
||||
@@ -1594,6 +1595,7 @@ fn pending_migrations_from_applied_skips_versions_already_applied() {
|
||||
20260921020100,
|
||||
20260923000000,
|
||||
20261001000000,
|
||||
20261004000000,
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,180 @@
|
||||
use super::*;
|
||||
use serde_json::Value;
|
||||
|
||||
const OPTIMIZATION_VERSION: i64 = 20261004000000;
|
||||
|
||||
async fn read_facts(pool: &PgPool) -> Value {
|
||||
query_scalar(
|
||||
"SELECT jsonb_agg(to_jsonb(f) ORDER BY request_id) FROM usage_analytics_facts_v1 f",
|
||||
)
|
||||
.fetch_one(pool)
|
||||
.await
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
async fn stored_rows(pool: &PgPool) -> Vec<Value> {
|
||||
let mut rows = Vec::new();
|
||||
for table in [
|
||||
"usage",
|
||||
"usage_settlement_snapshots",
|
||||
"usage_attribution_snapshots",
|
||||
"stats_overview_dirty_events",
|
||||
"stats_bucket_state",
|
||||
"stats_overview_hourly",
|
||||
"stats_overview_daily",
|
||||
"dashboard_request_contributions",
|
||||
"dashboard_stats_total",
|
||||
"dashboard_activity_minute",
|
||||
] {
|
||||
rows.push(
|
||||
query_scalar(&format!(
|
||||
"SELECT COALESCE(jsonb_agg(to_jsonb(t) ORDER BY to_jsonb(t)), '[]') FROM {table} t"
|
||||
))
|
||||
.fetch_one(pool)
|
||||
.await
|
||||
.unwrap(),
|
||||
);
|
||||
}
|
||||
rows
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn overview_fact_metadata_optimization_preserves_facts_without_backfill() {
|
||||
let Some(server) = ManagedPostgresServer::try_start().await.unwrap() else {
|
||||
return;
|
||||
};
|
||||
let mut connection = PgConnection::connect(server.database_url()).await.unwrap();
|
||||
connection.ensure_migrations_table().await.unwrap();
|
||||
for migration in POSTGRES_MIGRATOR
|
||||
.iter()
|
||||
.filter(|migration| migration.version < OPTIMIZATION_VERSION)
|
||||
{
|
||||
connection.apply(migration).await.unwrap();
|
||||
}
|
||||
let pool = PgPool::connect(server.database_url()).await.unwrap();
|
||||
sqlx::raw_sql(
|
||||
r#"
|
||||
INSERT INTO users(id,username,email_verified) VALUES('owner','owner',false);
|
||||
INSERT INTO api_keys(id,user_id,key_hash,is_standalone)
|
||||
VALUES('employee-key','owner',repeat('e',64),false),
|
||||
('standalone-key','owner',repeat('f',64),true);
|
||||
"#,
|
||||
)
|
||||
.execute(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// Preserve distinctions between absent/null flags, booleans and strings,
|
||||
// malformed nested shapes, key overrides and pre-attribution legacy rows.
|
||||
let metadata_cases = [
|
||||
None,
|
||||
Some("null"),
|
||||
Some("[]"),
|
||||
Some("false"),
|
||||
Some("\"metadata\""),
|
||||
Some("{}"),
|
||||
Some(r#"{"usage_available":false,"usage_pricing_available":false}"#),
|
||||
Some(r#"{"usage_available":null,"usage_pricing_available":null}"#),
|
||||
Some(r#"{"usage_available":"false","usage_pricing_available":"false"}"#),
|
||||
Some(
|
||||
r#"{"analytics_attribution":{"is_standalone":true},"analytics_measurement":{"source":"reported"}}"#,
|
||||
),
|
||||
Some(
|
||||
r#"{"analytics_attribution":{"is_standalone":false},"analytics_measurement":{"source":"estimated"}}"#,
|
||||
),
|
||||
Some(
|
||||
r#"{"analytics_attribution":{"is_standalone":"true"},"api_key_is_standalone":false,"analytics_measurement":{"source":"mixed"}}"#,
|
||||
),
|
||||
Some(
|
||||
r#"{"analytics_attribution":[],"api_key_is_standalone":true,"analytics_measurement":{"source":false}}"#,
|
||||
),
|
||||
Some(
|
||||
r#"{"analytics_attribution":null,"api_key_is_standalone":"true","analytics_measurement":[]}"#,
|
||||
),
|
||||
Some(
|
||||
r#"{"usage_available":true,"usage_available":false,"analytics_attribution":{"is_standalone":false},"analytics_attribution":{"is_standalone":true}}"#,
|
||||
),
|
||||
];
|
||||
for (index, metadata) in metadata_cases.into_iter().enumerate() {
|
||||
for key in [None, Some("employee-key"), Some("standalone-key")] {
|
||||
let request = format!("metadata-{index}-{}", key.unwrap_or("legacy"));
|
||||
query(
|
||||
r#"INSERT INTO usage(id,request_id,user_id,api_key_id,model,provider_name,
|
||||
status,billing_status,input_tokens,output_tokens,total_tokens,
|
||||
cache_read_input_tokens,cache_creation_input_tokens,total_cost_usd,
|
||||
actual_total_cost_usd,response_time_ms,first_byte_time_ms,is_stream,
|
||||
created_at,request_metadata)
|
||||
VALUES($1,$1,'owner',$2,'test','test','completed','settled',100,10,110,
|
||||
20,5,0.12345678,0.11111111,1000,100,true,'2026-01-01 00:00:00+00',$3::json)"#,
|
||||
)
|
||||
.bind(&request)
|
||||
.bind(key)
|
||||
.bind(metadata)
|
||||
.execute(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
if index % 2 == 0 {
|
||||
query(
|
||||
r#"INSERT INTO usage_settlement_snapshots(request_id,billing_status,
|
||||
billing_input_tokens,billing_effective_input_tokens,billing_output_tokens,
|
||||
billing_cache_read_tokens,billing_cache_creation_tokens,
|
||||
billing_total_cost_usd,billing_actual_total_cost_usd,input_price_per_1m,
|
||||
billing_cache_read_cost_usd,billing_cache_creation_cost_usd,
|
||||
quota_covered_amount_usd,wallet_consumed_amount_usd,wallet_debit_amount_usd,
|
||||
wallet_recharge_debit_usd,wallet_gift_debit_usd,wallet_overdraft_usd,
|
||||
allocation_status)
|
||||
VALUES($1,'settled',200,175,20,15,10,0.3,0.25,2,0.00001,0.00002,
|
||||
0.05,0.2,0.2,0.1,0.1,0,'complete')"#,
|
||||
)
|
||||
.bind(&request)
|
||||
.execute(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
if key.is_none() {
|
||||
query("DELETE FROM usage_attribution_snapshots WHERE request_id=$1")
|
||||
.bind(&request)
|
||||
.execute(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
}
|
||||
}
|
||||
// Keep a large nested payload so parity also covers toasted JSON metadata.
|
||||
query("UPDATE usage SET request_metadata=json_build_object('usage_available',true,'payload',repeat('metadata payload ',4096)) WHERE request_id='metadata-5-legacy'")
|
||||
.execute(&pool).await.unwrap();
|
||||
let before_facts = read_facts(&pool).await;
|
||||
let before_rows = stored_rows(&pool).await;
|
||||
|
||||
// The definition-only upgrade must not touch historical projections/queues.
|
||||
let mut blocked_history = pool.begin().await.unwrap();
|
||||
query("LOCK TABLE stats_bucket_state,stats_overview_dirty_events,stats_overview_hourly,stats_overview_daily IN ACCESS EXCLUSIVE MODE")
|
||||
.execute(&mut *blocked_history).await.unwrap();
|
||||
query("SET lock_timeout='500ms'")
|
||||
.execute(&mut connection)
|
||||
.await
|
||||
.unwrap();
|
||||
let migration = POSTGRES_MIGRATOR
|
||||
.iter()
|
||||
.find(|migration| migration.version == OPTIMIZATION_VERSION)
|
||||
.unwrap();
|
||||
connection.apply(migration).await.unwrap();
|
||||
blocked_history.rollback().await.unwrap();
|
||||
|
||||
assert_eq!(read_facts(&pool).await, before_facts);
|
||||
assert_eq!(stored_rows(&pool).await, before_rows);
|
||||
|
||||
// The maintained bootstrap fragment must install the same read model.
|
||||
let bootstrap =
|
||||
include_str!("../../../../schema/bootstrap/postgres/190_overview_analytics.sql");
|
||||
let view_start = bootstrap
|
||||
.find("CREATE OR REPLACE VIEW public.usage_analytics_facts_v1 AS")
|
||||
.unwrap();
|
||||
sqlx::raw_sql(&bootstrap[view_start..])
|
||||
.execute(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(read_facts(&pool).await, before_facts);
|
||||
assert_eq!(stored_rows(&pool).await, before_rows);
|
||||
pool.close().await;
|
||||
}
|
||||
Reference in New Issue
Block a user