Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 7 additions & 1 deletion .config/nextest.toml
Original file line number Diff line number Diff line change
Expand Up @@ -43,12 +43,18 @@ filter = 'package(api) & binary(e2e_all) & test(/^admin_pricing_changes::test_(r
test-group = "e2e-db"
threads-required = "num-test-threads"

[[profile.default.overrides]]
filter = 'package(api) & binary(e2e_all) & test(/^admin_analytics::test_admin_platform_billing_summary$/)'
test-group = "e2e-db"
# Exact before/after platform totals require exclusive access to the shared DB.
threads-required = "num-test-threads"

[[profile.default.overrides]]
filter = "package(api) & binary(e2e_all)"
test-group = "e2e-db"

[[profile.default.overrides]]
filter = "package(database) & binary(spend_counter_backfill)"
filter = "package(database) & (binary(spend_counter_backfill) | binary(counter_readers))"
test-group = "e2e-db"

[[profile.default.scripts]]
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -199,7 +199,7 @@ jobs:
cache-key: e2e

- name: Run e2e tests
run: cargo nextest run --test e2e_all --test spend_counter_backfill
run: cargo nextest run --test e2e_all --test spend_counter_backfill --test counter_readers
Comment thread
henrypark133 marked this conversation as resolved.
env:
POSTGRES_PRIMARY_APP_ID: ${{ secrets.POSTGRES_PRIMARY_APP_ID }}
DATABASE_HOST: localhost
Expand Down
3 changes: 3 additions & 0 deletions crates/api/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,9 @@ async fn main() {
.await
.expect("Usage reporting index prerequisites are not satisfied");
}
database::ensure_spend_counters_ready(database.pool())
Comment thread
henrypark133 marked this conversation as resolved.
.await
.expect("Spend counters are incomplete; run the spend-counter backfill before starting Cloud API");
let auth_components = init_auth_services(database.clone(), &config);

// Initialize OpenTelemetry pipeline
Expand Down
131 changes: 69 additions & 62 deletions crates/api/tests/e2e_all/api_keys.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,9 @@
use crate::common::*;
use database::models::RecordUsageRequest;
use database::repositories::{
OrganizationServiceUsageRepository, OrganizationUsageRepository, RecordServiceUsageRequest,
};
use services::usage::InferenceType;

// ============================================
// API Key Creation and Management Tests
Expand Down Expand Up @@ -128,76 +133,78 @@ async fn test_list_workspace_api_keys_orders_by_usage() {
.unwrap()
.get(0);

for (api_key_id, total_cost) in [
(service_spend_key_id, 100_000_000_i64),
(inference_spend_key_id, 300_000_000_i64),
drop(client);
let inference_repository = OrganizationUsageRepository::new(database.pool().clone());
let inference = RecordUsageRequest {
organization_id,
workspace_id,
api_key_id: service_spend_key_id,
model_id,
model_name,
input_tokens: 10,
output_tokens: 10,
input_cost: 1,
output_cost: 1,
total_cost: 100_000_000,
inference_type: InferenceType::ChatCompletion.as_str().to_string(),
ttft_ms: None,
avg_itl_ms: None,
inference_id: Some(uuid::Uuid::new_v4()),
provider_request_id: None,
stop_reason: None,
response_id: None,
image_count: None,
cache_read_tokens: 0,
cache_write_tokens: 0,
billing_details: None,
service_tier: None,
context_band: None,
served_provider_tier: None,
served_provider_type: None,
served_via_fallback: false,
};
inference_repository
.record_usage(inference.clone())
.await
.unwrap();
inference_repository
.record_usage(RecordUsageRequest {
api_key_id: inference_spend_key_id,
total_cost: 300_000_000,
inference_id: Some(uuid::Uuid::new_v4()),
..inference
})
.await
.unwrap();
OrganizationServiceUsageRepository::new(database.pool().clone())
.record_usage(&RecordServiceUsageRequest {
organization_id,
workspace_id,
api_key_id: service_spend_key_id,
service_id,
quantity: 1,
total_cost: 400_000_000,
inference_id: None,
})
.await
.unwrap();

let client = database.pool().get().await.unwrap();

for (api_key_id, hours) in [
(service_spend_key_id, 3_i64),
(inference_spend_key_id, 2_i64),
(unused_key_id, 1_i64),
] {
client
.execute(
r#"
INSERT INTO organization_usage_log (
id, organization_id, workspace_id, api_key_id,
model_id, model_name, input_tokens, output_tokens,
total_tokens, input_cost, output_cost, total_cost,
inference_type, created_at
) VALUES ($1, $2, $3, $4, $5, $6, 10, 10, 20, 1, 1, $7,
'chat_completion', NOW())
"#,
&[
&uuid::Uuid::new_v4(),
&organization_id,
&workspace_id,
&api_key_id,
&model_id,
&model_name,
&total_cost,
],
"UPDATE api_keys SET created_at = NOW() + ($2::BIGINT * INTERVAL '1 hour') WHERE id = $1",
&[&api_key_id, &hours],
)
.await
.unwrap();
}

client
.execute(
r#"
INSERT INTO organization_service_usage_log (
id, organization_id, workspace_id, api_key_id,
service_id, quantity, total_cost, inference_id, created_at
) VALUES ($1, $2, $3, $4, $5, 1, 400000000, NULL, NOW())
"#,
&[
&uuid::Uuid::new_v4(),
&organization_id,
&workspace_id,
&service_spend_key_id,
&service_id,
],
)
.await
.unwrap();

client
.execute(
"UPDATE api_keys SET created_at = NOW() + INTERVAL '3 hours' WHERE id = $1",
&[&service_spend_key_id],
)
.await
.unwrap();
client
.execute(
"UPDATE api_keys SET created_at = NOW() + INTERVAL '2 hours' WHERE id = $1",
&[&inference_spend_key_id],
)
.await
.unwrap();
client
.execute(
"UPDATE api_keys SET created_at = NOW() + INTERVAL '1 hour' WHERE id = $1",
&[&unused_key_id],
)
.await
.unwrap();

let default_response = server
.get(format!("/v1/workspaces/{}/api-keys?limit=3", workspace.id).as_str())
.add_header("Authorization", format!("Bearer {}", get_session_id()))
Expand Down
12 changes: 6 additions & 6 deletions crates/database/src/repositories/analytics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -545,16 +545,16 @@ impl AnalyticsRepository for PgAnalyticsRepository {
let paying_org_count: i64 = limits_row.get(2);
let granted_org_count: i64 = limits_row.get(3);

// All-time consumed cost. `total` (from the cached balance) is ALL usage
// (inference + services); the inference/service splits come from their logs
// and reconcile to the total.
// All-time consumed cost. Keep the legacy total independently; historical
// adjustments can differ from the balance splits.
let consumed_row = client
.query_one(
r#"
SELECT
(SELECT COALESCE(SUM(total_spent), 0) FROM organization_balance)::bigint as total_nano,
(SELECT COALESCE(SUM(total_cost), 0) FROM organization_usage_log)::bigint as inference_nano,
(SELECT COALESCE(SUM(total_cost), 0) FROM organization_service_usage_log)::bigint as service_nano
COALESCE(SUM(total_spent), 0)::bigint as total_nano,
COALESCE(SUM(inference_spent), 0)::bigint as inference_nano,
COALESCE(SUM(service_spent), 0)::bigint as service_nano
FROM organization_balance
"#,
&[],
)
Expand Down
21 changes: 5 additions & 16 deletions crates/database/src/repositories/api_key.rs
Original file line number Diff line number Diff line change
Expand Up @@ -253,8 +253,8 @@ impl ApiKeyRepository {
Ok(row.get::<_, i64>("count"))
}

/// List API keys for a workspace with usage data
/// This is the primary method to list API keys, using an efficient JOIN query
/// List API keys for a workspace with usage data.
/// This is the primary method to list API keys, using the spend counter JOIN.
pub async fn list_by_workspace_paginated(
&self,
workspace_id: Uuid,
Expand Down Expand Up @@ -305,22 +305,11 @@ impl ApiKeyRepository {
ak.deleted_at,
ak.spend_limit,
(
COALESCE(inference_usage.total_cost, 0)
+ COALESCE(service_usage.total_cost, 0)
COALESCE(spend.inference_spent, 0)
+ COALESCE(spend.service_spent, 0)
)::BIGINT as usage
FROM api_keys ak
LEFT JOIN (
SELECT api_key_id, COALESCE(SUM(total_cost), 0)::BIGINT AS total_cost
FROM organization_usage_log
WHERE workspace_id = $1
GROUP BY api_key_id
) inference_usage ON ak.id = inference_usage.api_key_id
LEFT JOIN (
SELECT api_key_id, COALESCE(SUM(total_cost), 0)::BIGINT AS total_cost
FROM organization_service_usage_log
WHERE workspace_id = $1
GROUP BY api_key_id
) service_usage ON ak.id = service_usage.api_key_id
LEFT JOIN api_key_spend spend ON ak.id = spend.api_key_id
WHERE ak.workspace_id = $1 AND ak.deleted_at IS NULL
ORDER BY {order_by_column} {order_dir}{tie_breaker}
LIMIT $2 OFFSET $3
Expand Down
13 changes: 7 additions & 6 deletions crates/database/src/repositories/organization_usage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ impl OrganizationUsageRepository {
}
}

/// Get total spend for a specific API key
/// Get inference-only spend for a specific API key for admission-limit checks.
pub async fn get_api_key_spend(&self, api_key_id: Uuid) -> Result<i64> {
let row = retry_db!("get_api_key_spend", {
let client = self
Expand All @@ -68,18 +68,19 @@ impl OrganizationUsageRepository {
client
.query_one(
r#"
SELECT COALESCE(SUM(total_cost), 0)::BIGINT as total_spend
FROM organization_usage_log
WHERE api_key_id = $1
SELECT COALESCE(
(SELECT inference_spent FROM api_key_spend WHERE api_key_id = $1),
0
)::BIGINT as inference_spend
"#,
&[&api_key_id],
)
.await
.map_err(map_db_error)
})?;

let total_spend: i64 = row.get("total_spend");
Ok(total_spend)
let inference_spend: i64 = row.get("inference_spend");
Ok(inference_spend)
}

/// Record usage and update balance atomically.
Expand Down
Loading
Loading