[Feat] Litellm x CloudZero Integration - Cost Tracking (#14296)

* fix: just pull LiteLLM_DailyUserSpend

* get the team_id from user daily spend table

* cloudzero_dry_run_export

* fix CZ endpoints

* trace entity_id

* fix: get_usage_data

* fix get_usage_data

* fix _create_cbf_record

* fix get_usage_data

* ensure start and end time is used for exporting data

* fix init_cloudzero_background_job

* fix CloudZeroExportRequest

* fix initialize_cloudzero_export_job

* fix initialize_cloudzero_export_job

* allow init with env + config.yaml for cloudzero

* fix: init CZ through config.yaml

* fix DRY run on CZ

* TestCloudZeroDryRunEndpoint

* fix: CLOUDZERO_EXPORT_INTERVAL_MINUTES

* fix init_cloudzero_background_job

* fix exporting data

* fix transform

* stash cloudzero docs

* docs: CloudZero

* ruff fix

* fix rendering key alias

* fix polars
This commit is contained in:
Ishaan Jaff 2025-09-05 21:29:41 -07:00 • committed by GitHub
parent 07ba3ff036
commit 5310bba35b
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
14 changed files with 719 additions and 196 deletions

View file

@ -1292,6 +1292,7 @@ jobs:
pip install "tokenizers==0.20.0"
pip install "uvloop==0.21.0"
pip install "fastuuid==0.12.0"
pip install "polars==1.31.0"
pip install jsonschema
- setup_litellm_enterprise_pip
- run:

View file

@ -0,0 +1,209 @@
import Tabs from '@theme/Tabs';
import TabItem from '@theme/TabItem';
# CloudZero Integration
LiteLLM provides an integration with CloudZero's AnyCost API, allowing you to export your LLM usage data to CloudZero for cost tracking analysis.
## Overview
| Property | Details |
|----------|---------|
| Description | Export LiteLLM usage data to CloudZero AnyCost API for cost tracking and analysis |
| callback name | `cloudzero`|
| Supported Operations | • Automatic hourly data export<br/>• Manual data export<br/>• Dry run testing<br/>• Cost and token usage tracking |
| Data Format | CloudZero Billing Format (CBF) with proper resource tagging |
| Export Frequency | Hourly (configurable via `CLOUDZERO_EXPORT_INTERVAL_MINUTES`) |
## Environment Variables
| Variable | Required | Description | Example |
|----------|----------|-------------|---------|
| `CLOUDZERO_API_KEY` | Yes | Your CloudZero API key | `cz_api_xxxxxxxxxx` |
| `CLOUDZERO_CONNECTION_ID` | Yes | CloudZero connection ID for data submission | `conn_xxxxxxxxxx` |
| `CLOUDZERO_TIMEZONE` | No | Timezone for date handling (default: UTC) | `America/New_York` |
| `CLOUDZERO_EXPORT_INTERVAL_MINUTES` | No | Export frequency in minutes (default: 60) | `60` |
## Setup
### End to End Video Walkthrough
This video walks through the entire process of setting up LiteLLM with CloudZero integration and viewing LiteLLM exported usage data in CloudZero.
<iframe width="840" height="500" src="https://www.loom.com/embed/59b57593183f4cc3b1c05a2dd3277f92" frameborder="0" webkitallowfullscreen mozallowfullscreen allowfullscreen></iframe>
### Step 1: Configure Environment Variables
Set your CloudZero credentials in your environment:
```bash
export CLOUDZERO_API_KEY="cz_api_xxxxxxxxxx"
export CLOUDZERO_CONNECTION_ID="conn_xxxxxxxxxx"
export CLOUDZERO_TIMEZONE="UTC" # Optional, defaults to UTC
```
### Step 2: Enable CloudZero Integration
Add the CloudZero callback to your LiteLLM configuration YAML file:
```yaml
model_list:
- model_name: gpt-4o
litellm_params:
model: openai/gpt-4o
api_key: sk-xxxxxxx
litellm_settings:
callbacks: ["cloudzero"] # Enable CloudZero integration
```
### Step 3: Start LiteLLM Proxy
Start your LiteLLM proxy with the configuration:
```bash
litellm --config /path/to/config.yaml
```
## Testing Your Setup
### Dry Run Export
Call the dry run endpoint to test your CloudZero configuration without sending data to CloudZero. This endpoint will not send any data to CloudZero, but will return the data that would be exported.
```bash
curl -X POST "http://localhost:4000/cloudzero/dry-run" \
-H "Content-Type: application/json" \
-H "Authorization: Bearer sk-1234" \
-d '{
"limit": 10
}' | jq
```
**Expected Response:**
```json
{
"message": "CloudZero dry run export completed successfully.",
"status": "success",
"dry_run_data": {
"usage_data": [...],
"cbf_data": [...],
"summary": {
"total_cost": 0.05,
"total_tokens": 1250,
"total_records": 10
}
}
}
```
### Manual Export
Call the export endpoint to send data immediately to CloudZero. We suggest setting a small `limit` to test the export. This will only export the last 10 records to CloudZero. Note: Cloudzero can take up to 15 minutes to process the exported data.
```bash
curl -X POST "http://localhost:4000/cloudzero/export" \
-H "Content-Type: application/json" \
-H "Authorization: Bearer sk-1234" \
-d '{
"limit": 10
}' | jq
```
**Expected Response:**
```json
{
"message": "CloudZero export completed successfully",
"status": "success"
}
```
## Data Export Details
### Automatic Export Schedule
- **Frequency**: Every 60 minutes (configurable via `CLOUDZERO_EXPORT_INTERVAL_MINUTES`)
- **Data Processing**: LiteLLM automatically processes and exports usage data hourly
- **CloudZero Processing**: CloudZero typically takes 10-15 minutes to process data from LiteLLM
### Data Format
LiteLLM exports data in CloudZero Billing Format (CBF) with the following structure:
```json
{
"time/usage_start": "2024-01-15T14:00:00Z",
"cost/cost": 0.002,
"usage/amount": 150,
"usage/units": "tokens",
"resource/id": "czrn:litellm:openai:cross-region:team-123:llm-usage:gpt-4o",
"resource/service": "litellm",
"resource/account": "team-123",
"resource/region": "cross-region",
"resource/usage_family": "llm-usage",
"resource/tag:provider": "openai",
"resource/tag:model": "gpt-4o",
"resource/tag:prompt_tokens": "100",
"resource/tag:completion_tokens": "50"
}
```
### Resource Tagging
LiteLLM automatically creates comprehensive resource tags for cost attribution:
- **Provider Tags**: `openai`, `anthropic`, `azure`, etc.
- **Model Tags**: Specific model names like `gpt-4o`, `claude-3-sonnet`
- **Team/User Tags**: Team IDs and user IDs for cost allocation
- **Token Breakdown**: Separate tracking of prompt and completion tokens
- **Usage Metrics**: Total tokens consumed per request
## Advanced Configuration
### Custom Export Frequency
Change the export frequency (not recommended to go below 60 minutes):
```bash
export CLOUDZERO_EXPORT_INTERVAL_MINUTES=120 # Export every 2 hours
```
### Custom Time Range Export
Export data for a specific time range:
```bash
curl -X POST "http://localhost:4000/cloudzero/export" \
-H "Content-Type: application/json" \
-H "Authorization: Bearer sk-1234" \
-d '{
"start_time_utc": "2024-01-15T00:00:00Z",
"end_time_utc": "2024-01-15T23:59:59Z",
"operation": "replace_hourly"
}' | jq
```
## Troubleshooting
### Common Issues
1. **Missing Credentials Error**
```
CloudZero configuration missing. Please set CLOUDZERO_API_KEY and CLOUDZERO_CONNECTION_ID environment variables.
```
**Solution**: Ensure both environment variables are set with valid values.
2. **Connection Issues**
- Verify your CloudZero API key is valid
- Check that the connection ID exists in your CloudZero account
- Ensure your proxy has internet access to reach CloudZero's API
3. **No Data in CloudZero**
- CloudZero can take 10-15 minutes to process data
- Check that your LiteLLM proxy is generating usage data
- Use the dry-run endpoint to verify data is being formatted correctly
## Related Links
- [CloudZero Documentation](https://docs.cloudzero.com/)
- [CloudZero AnyCost API](https://docs.cloudzero.com/reference/anycost-api)

View file

@ -146,6 +146,7 @@ _custom_logger_compatible_callbacks_literal = Literal[
"aws_sqs",
"vector_store_pre_call_hook",
"dotprompt",
"cloudzero",
]
configured_cold_storage_logger: Optional[_custom_logger_compatible_callbacks_literal] = None
logged_real_time_event_types: Optional[Union[List[str], Literal["*"]]] = None

View file

@ -873,6 +873,9 @@ AZURE_STORAGE_MSFT_VERSION = "2019-07-07"
PROMETHEUS_BUDGET_METRICS_REFRESH_INTERVAL_MINUTES = int(
os.getenv("PROMETHEUS_BUDGET_METRICS_REFRESH_INTERVAL_MINUTES", 5)
)
CLOUDZERO_EXPORT_INTERVAL_MINUTES = int(
os.getenv("CLOUDZERO_EXPORT_INTERVAL_MINUTES", 60)
)
MCP_TOOL_NAME_PREFIX = "mcp_tool"
MAXIMUM_TRACEBACK_LINES_TO_LOG = int(os.getenv("MAXIMUM_TRACEBACK_LINES_TO_LOG", 100))
@ -927,6 +930,8 @@ LITELLM_CLI_SESSION_TOKEN_PREFIX = "litellm-session-token"
########################### DB CRON JOB NAMES ###########################
DB_SPEND_UPDATE_JOB_NAME = "db_spend_update_job"
PROMETHEUS_EMIT_BUDGET_METRICS_JOB_NAME = "prometheus_emit_budget_metrics"
CLOUDZERO_EXPORT_USAGE_DATA_JOB_NAME = "cloudzero_export_usage_data"
CLOUDZERO_MAX_FETCHED_DATA_RECORDS = int(os.getenv("CLOUDZERO_MAX_FETCHED_DATA_RECORDS", 50000))
SPEND_LOG_CLEANUP_JOB_NAME = "spend_log_cleanup"
SPEND_LOG_RUN_LOOPS = int(os.getenv("SPEND_LOG_RUN_LOOPS", 500))
SPEND_LOG_CLEANUP_BATCH_SIZE = int(os.getenv("SPEND_LOG_CLEANUP_BATCH_SIZE", 1000))

View file

@ -1,6 +1,8 @@
import os
from typing import Optional
from datetime import datetime
from typing import TYPE_CHECKING, Any, List, Optional, cast
import litellm
from litellm._logging import verbose_logger
from litellm.integrations.custom_logger import CustomLogger
@ -8,6 +10,11 @@ from .cz_stream_api import CloudZeroStreamer
from .database import LiteLLMDatabase
from .transform import CBFTransformer
if TYPE_CHECKING:
from apscheduler.schedulers.asyncio import AsyncIOScheduler
else:
AsyncIOScheduler = Any
class CloudZeroLogger(CustomLogger):
"""
@ -27,8 +34,66 @@ class CloudZeroLogger(CustomLogger):
self.api_key = api_key or os.getenv("CLOUDZERO_API_KEY")
self.connection_id = connection_id or os.getenv("CLOUDZERO_CONNECTION_ID")
self.timezone = timezone or os.getenv("CLOUDZERO_TIMEZONE", "UTC")
verbose_logger.debug(f"CloudZero Logger initialized with connection ID: {self.connection_id}, timezone: {self.timezone}")
async def export_usage_data(self, limit: Optional[int] = None, operation: str = "replace_hourly"):
async def initialize_cloudzero_export_job(self):
"""
Handler for initializing CloudZero export job.
Runs when CloudZero logger starts up.
- If redis cache is available, we use the pod lock manager to acquire a lock and export the data.
- Ensures only one pod exports the data at a time.
- If redis cache is not available, we export the data directly.
"""
from litellm.constants import (
CLOUDZERO_EXPORT_USAGE_DATA_JOB_NAME,
)
from litellm.proxy.proxy_server import proxy_logging_obj
pod_lock_manager = proxy_logging_obj.db_spend_update_writer.pod_lock_manager
# if using redis, ensure only one pod exports the data at a time
if pod_lock_manager and pod_lock_manager.redis_cache:
if await pod_lock_manager.acquire_lock(
cronjob_id=CLOUDZERO_EXPORT_USAGE_DATA_JOB_NAME
):
try:
await self._hourly_usage_data_export()
finally:
await pod_lock_manager.release_lock(
cronjob_id=CLOUDZERO_EXPORT_USAGE_DATA_JOB_NAME
)
else:
# if not using redis, export the data directly
await self._hourly_usage_data_export()
async def _hourly_usage_data_export(self):
"""
Exports the hourly usage data to CloudZero.
Start time: 1 hour ago
End time: current time
"""
from datetime import timedelta, timezone
from litellm.constants import CLOUDZERO_MAX_FETCHED_DATA_RECORDS
current_time_utc = datetime.now(timezone.utc)
one_hour_ago_utc = current_time_utc - timedelta(hours=1)
await self.export_usage_data(
limit=CLOUDZERO_MAX_FETCHED_DATA_RECORDS,
operation="replace_hourly",
start_time_utc=one_hour_ago_utc,
end_time_utc=current_time_utc
)
async def export_usage_data(
self,
limit: Optional[int] = None,
operation: str = "replace_hourly",
start_time_utc: Optional[datetime] = None,
end_time_utc: Optional[datetime] = None
):
"""
Exports the usage data to CloudZero.
@ -52,7 +117,11 @@ class CloudZeroLogger(CustomLogger):
# Initialize database connection and load data
database = LiteLLMDatabase()
verbose_logger.debug("CloudZero Logger: Loading usage data from database")
data = await database.get_usage_data(limit=limit)
data = await database.get_usage_data(
limit=limit,
start_time_utc=start_time_utc,
end_time_utc=end_time_utc
)
if data.is_empty():
verbose_logger.info("CloudZero Logger: No usage data found to export")
@ -86,10 +155,13 @@ class CloudZeroLogger(CustomLogger):
async def dry_run_export_usage_data(self, limit: Optional[int] = 10000):
"""
Only prints the data that would be exported to CloudZero.
Returns the data that would be exported to CloudZero without actually sending it.
Args:
limit: Limit number of records to display (default: 10000)
Returns:
dict: Contains usage_data, cbf_data, and summary statistics
"""
try:
verbose_logger.debug("CloudZero Logger: Starting dry run export")
@ -101,23 +173,64 @@ class CloudZeroLogger(CustomLogger):
if data.is_empty():
verbose_logger.warning("CloudZero Dry Run: No usage data found")
return
return {
"usage_data": [],
"cbf_data": [],
"summary": {
"total_records": 0,
"total_cost": 0,
"total_tokens": 0,
"unique_accounts": 0,
"unique_services": 0
}
}
verbose_logger.debug(f"CloudZero Dry Run: Processing {len(data)} records...")
# Convert usage data to dict format for response
usage_data_sample = data.head(50).to_dicts() # Return first 50 rows
# Transform data to CloudZero CBF format
transformer = CBFTransformer()
cbf_data = transformer.transform(data)
if cbf_data.is_empty():
verbose_logger.warning("CloudZero Dry Run: No valid data after transformation")
return
return {
"usage_data": usage_data_sample,
"cbf_data": [],
"summary": {
"total_records": len(usage_data_sample),
"total_cost": sum(row.get('spend', 0) for row in usage_data_sample),
"total_tokens": sum(row.get('prompt_tokens', 0) + row.get('completion_tokens', 0) for row in usage_data_sample),
"unique_accounts": 0,
"unique_services": 0
}
}
# Display the transformed data on screen
self._display_cbf_data_on_screen(cbf_data)
# Convert CBF data to dict format for response
cbf_data_dict = cbf_data.to_dicts()
# Calculate summary statistics
total_cost = sum(record.get('cost/cost', 0) for record in cbf_data_dict)
unique_accounts = len(set(record.get('resource/account', '') for record in cbf_data_dict if record.get('resource/account')))
unique_services = len(set(record.get('resource/service', '') for record in cbf_data_dict if record.get('resource/service')))
total_tokens = sum(record.get('usage/amount', 0) for record in cbf_data_dict)
verbose_logger.info(f"CloudZero Logger: Dry run completed for {len(cbf_data)} records")
return {
"usage_data": usage_data_sample,
"cbf_data": cbf_data_dict,
"summary": {
"total_records": len(cbf_data_dict),
"total_cost": total_cost,
"total_tokens": total_tokens,
"unique_accounts": unique_accounts,
"unique_services": unique_services
}
}
except Exception as e:
verbose_logger.error(f"CloudZero Logger: Error in dry run export: {str(e)}")
verbose_logger.error(f"CloudZero Dry Run Error: {str(e)}")
@ -144,6 +257,11 @@ class CloudZeroLogger(CustomLogger):
cbf_table = Table(show_header=True, header_style="bold cyan", box=SIMPLE, padding=(0, 1))
cbf_table.add_column("time/usage_start", style="blue", no_wrap=False)
cbf_table.add_column("cost/cost", style="green", justify="right", no_wrap=False)
cbf_table.add_column("entity_type", style="magenta", justify="right", no_wrap=False)
cbf_table.add_column("entity_id", style="magenta", justify="right", no_wrap=False)
cbf_table.add_column("team_id", style="cyan", no_wrap=False)
cbf_table.add_column("team_alias", style="cyan", no_wrap=False)
cbf_table.add_column("api_key_alias", style="yellow", no_wrap=False)
cbf_table.add_column("usage/amount", style="yellow", justify="right", no_wrap=False)
cbf_table.add_column("resource/id", style="magenta", no_wrap=False)
cbf_table.add_column("resource/service", style="cyan", no_wrap=False)
@ -159,10 +277,20 @@ class CloudZeroLogger(CustomLogger):
resource_service = str(record.get('resource/service', 'N/A'))
resource_account = str(record.get('resource/account', 'N/A'))
resource_region = str(record.get('resource/region', 'N/A'))
entity_type = str(record.get('entity_type', 'N/A'))
entity_id = str(record.get('entity_id', 'N/A'))
team_id = str(record.get('resource/tag:team_id', 'N/A'))
team_alias = str(record.get('resource/tag:team_alias', 'N/A'))
api_key_alias = str(record.get('resource/tag:api_key_alias', 'N/A'))
cbf_table.add_row(
time_usage_start,
cost_cost,
entity_type,
entity_id,
team_id,
team_alias,
api_key_alias,
usage_amount,
resource_id,
resource_service,
@ -187,4 +315,34 @@ class CloudZeroLogger(CustomLogger):
console.print(f" Unique Accounts: {unique_accounts}")
console.print(f" Unique Services: {unique_services}")
console.print("\n[dim]💡 This is the CloudZero CBF format ready for AnyCost ingestion[/dim]")
console.print("\n[dim]💡 This is the CloudZero CBF format ready for AnyCost ingestion[/dim]")
@staticmethod
async def init_cloudzero_background_job(scheduler: AsyncIOScheduler):
"""
Initialize the CloudZero background job.
Starts the background job that exports the usage data to CloudZero every hour.
"""
from litellm.constants import CLOUDZERO_EXPORT_INTERVAL_MINUTES
from litellm.integrations.custom_logger import CustomLogger
prometheus_loggers: List[CustomLogger] = (
litellm.logging_callback_manager.get_custom_loggers_for_type(
callback_type=CloudZeroLogger
)
)
# we need to get the initialized prometheus logger instance(s) and call logger.initialize_remaining_budget_metrics() on them
verbose_logger.debug("found %s cloudzero loggers", len(prometheus_loggers))
if len(prometheus_loggers) > 0:
cloudzero_logger = cast(CloudZeroLogger, prometheus_loggers[0])
verbose_logger.debug(
"Initializing remaining budget metrics as a cron job executing every %s minutes"
% CLOUDZERO_EXPORT_INTERVAL_MINUTES
)
scheduler.add_job(
cloudzero_logger.initialize_cloudzero_export_job,
"interval",
minutes=CLOUDZERO_EXPORT_INTERVAL_MINUTES
)

View file

@ -17,11 +17,16 @@
"""CloudZero Resource Names (CZRN) generation and validation for LiteLLM resources."""
import re
from enum import Enum
from typing import Any, cast
import litellm
class CZEntityType(str, Enum):
TEAM = "team"
class CZRNGenerator:
"""Generate CloudZero Resource Names (CZRNs) for LiteLLM resources."""
@ -49,8 +54,8 @@ class CZRNGenerator:
region = 'cross-region'
# Use the actual entity_id (team_id or user_id) as the owner account
entity_id = row.get('entity_id', 'unknown')
owner_account_id = self._normalize_component(entity_id)
team_id = row.get('team_id', 'unknown')
owner_account_id = self._normalize_component(team_id)
resource_type = 'llm-usage'

View file

@ -18,6 +18,7 @@
"""Database connection and data extraction for LiteLLM."""
from datetime import datetime
from typing import Any, Dict, Optional
import polars as pl
@ -35,85 +36,54 @@ class LiteLLMDatabase:
)
return prisma_client
async def get_usage_data(self, limit: Optional[int] = None) -> pl.DataFrame:
"""Retrieve consolidated usage data from LiteLLM daily spend tables."""
async def get_usage_data(
self,
limit: Optional[int] = None,
start_time_utc: Optional[datetime] = None,
end_time_utc: Optional[datetime] = None
) -> pl.DataFrame:
"""Retrieve usage data from LiteLLM daily user spend table."""
client = self._ensure_prisma_client()
# Union query to combine user, team, and tag spend data
query = """
WITH consolidated_spend AS (
-- User spend data
SELECT
id,
date,
user_id as entity_id,
'user' as entity_type,
api_key,
model,
model_group,
custom_llm_provider,
prompt_tokens,
completion_tokens,
spend,
api_requests,
successful_requests,
failed_requests,
cache_creation_input_tokens,
cache_read_input_tokens,
created_at,
updated_at
FROM "LiteLLM_DailyUserSpend"
UNION ALL
-- Team spend data
SELECT
id,
date,
team_id as entity_id,
'team' as entity_type,
api_key,
model,
model_group,
custom_llm_provider,
prompt_tokens,
completion_tokens,
spend,
api_requests,
successful_requests,
failed_requests,
cache_creation_input_tokens,
cache_read_input_tokens,
created_at,
updated_at
FROM "LiteLLM_DailyTeamSpend"
UNION ALL
-- Tag spend data
SELECT
id,
date,
tag as entity_id,
'tag' as entity_type,
api_key,
model,
model_group,
custom_llm_provider,
prompt_tokens,
completion_tokens,
spend,
api_requests,
successful_requests,
failed_requests,
cache_creation_input_tokens,
cache_read_input_tokens,
created_at,
updated_at
FROM "LiteLLM_DailyTagSpend"
)
SELECT * FROM consolidated_spend
ORDER BY date DESC, created_at DESC
# Build WHERE clause for time filtering
where_conditions = []
if start_time_utc:
where_conditions.append(f"dus.created_at >= '{start_time_utc.isoformat()}'")
if end_time_utc:
where_conditions.append(f"dus.created_at <= '{end_time_utc.isoformat()}'")
where_clause = ""
if where_conditions:
where_clause = "WHERE " + " AND ".join(where_conditions)
# Query to get user spend data with team information
query = f"""
SELECT
dus.id,
dus.date,
dus.user_id,
dus.api_key,
dus.model,
dus.model_group,
dus.custom_llm_provider,
dus.prompt_tokens,
dus.completion_tokens,
dus.spend,
dus.api_requests,
dus.successful_requests,
dus.failed_requests,
dus.cache_creation_input_tokens,
dus.cache_read_input_tokens,
dus.created_at,
dus.updated_at,
vt.team_id,
vt.key_alias as api_key_alias,
tt.team_alias
FROM "LiteLLM_DailyUserSpend" dus
LEFT JOIN "LiteLLM_VerificationToken" vt ON dus.api_key = vt.token
LEFT JOIN "LiteLLM_TeamTable" tt ON vt.team_id = tt.team_id
{where_clause}
ORDER BY dus.date DESC, dus.created_at DESC
"""
if limit:
@ -121,22 +91,21 @@ class LiteLLMDatabase:
try:
db_response = await client.db.query_raw(query)
# Convert the response to polars DataFrame
return pl.DataFrame(db_response)
# Convert the response to polars DataFrame with full schema inference
# This prevents schema mismatch errors when data types vary across rows
return pl.DataFrame(db_response, infer_schema_length=None)
except Exception as e:
raise Exception(f"Error retrieving usage data: {str(e)}")
async def get_table_info(self) -> Dict[str, Any]:
"""Get information about the consolidated daily spend tables."""
"""Get information about the daily user spend table."""
client = self._ensure_prisma_client()
try:
# Get combined row count from both tables
# Get row count from user spend table
user_count = await self._get_table_row_count('LiteLLM_DailyUserSpend')
team_count = await self._get_table_row_count('LiteLLM_DailyTeamSpend')
tag_count = await self._get_table_row_count('LiteLLM_DailyTagSpend')
# Get column structure from user spend table (representative)
# Get column structure from user spend table
query = """
SELECT column_name, data_type, is_nullable
FROM information_schema.columns
@ -147,12 +116,8 @@ class LiteLLMDatabase:
return {
'columns': columns_response,
'row_count': user_count + team_count + tag_count,
'table_breakdown': {
'user_spend': user_count,
'team_spend': team_count,
'tag_spend': tag_count
}
'row_count': user_count,
'table_name': 'LiteLLM_DailyUserSpend'
}
except Exception as e:
raise Exception(f"Error getting table info: {str(e)}")

View file

@ -24,7 +24,7 @@ from typing import Any, Optional
import polars as pl
from ...types.integrations.cloudzero import CBFRecord
from .cz_resource_names import CZRNGenerator
from .cz_resource_names import CZEntityType, CZRNGenerator
class CBFTransformer:
@ -92,17 +92,26 @@ class CBFTransformer:
resource_id = self.czrn_generator.create_from_litellm_data(row)
# Build dimensions for CloudZero
entity_id = str(row.get('entity_id', ''))
model = str(row.get('model', ''))
api_key_hash = str(row.get('api_key', ''))[:8] # First 8 chars for identification
# Handle team information with fallbacks
team_id = row.get('team_id')
team_alias = row.get('team_alias')
# Use team_alias if available, otherwise team_id, otherwise fallback to 'unknown'
entity_id = str(team_alias) if team_alias else (str(team_id) if team_id else 'unknown')
dimensions = {
'entity_type': str(row.get('entity_type', '')), # 'user' or 'team'
'entity_type': CZEntityType.TEAM.value,
'entity_id': entity_id,
'team_id': str(team_id) if team_id else 'unknown',
'team_alias': str(team_alias) if team_alias else 'unknown',
'model': model,
'model_group': str(row.get('model_group', '')),
'provider': str(row.get('custom_llm_provider', '')),
'api_key_prefix': api_key_hash,
'api_key_alias': str(row.get('api_key_alias', '')),
'api_requests': str(row.get('api_requests', 0)),
'successful_requests': str(row.get('successful_requests', 0)),
'failed_requests': str(row.get('failed_requests', 0)),
@ -138,10 +147,10 @@ class CBFTransformer:
# Add CZRN components that don't have direct CBF column mappings as resource tags
cbf_record['resource/tag:provider'] = provider # CZRN provider component
cbf_record['resource/tag:model'] = cloud_local_id # CZRN cloud-local-id component (model)
# Add resource tags for all dimensions (using resource/tag:<key> format)
for key, value in dimensions.items():
if value and value != 'N/A': # Only add non-empty tags
if value and value != 'N/A' and value != 'unknown': # Only add meaningful tags
cbf_record[f'resource/tag:{key}'] = str(value)
# Add token breakdown as resource tags for analysis

View file

@ -3369,7 +3369,14 @@ def _init_custom_logger_compatible_class( # noqa: PLR0915
galileo_logger = GalileoObserve()
_in_memory_loggers.append(galileo_logger)
return galileo_logger # type: ignore
elif logging_integration == "cloudzero":
from litellm.integrations.cloudzero.cloudzero import CloudZeroLogger
for callback in _in_memory_loggers:
if isinstance(callback, CloudZeroLogger):
return callback # type: ignore
cloudzero_logger = CloudZeroLogger()
_in_memory_loggers.append(cloudzero_logger)
return cloudzero_logger # type: ignore
elif logging_integration == "deepeval":
for callback in _in_memory_loggers:
if isinstance(callback, DeepEvalLogger):
@ -3589,6 +3596,11 @@ def get_custom_logger_compatible_class( # noqa: PLR0915
for callback in _in_memory_loggers:
if isinstance(callback, GalileoObserve):
return callback
elif logging_integration == "cloudzero":
from litellm.integrations.cloudzero.cloudzero import CloudZeroLogger
for callback in _in_memory_loggers:
if isinstance(callback, CloudZeroLogger):
return callback
elif logging_integration == "deepeval":
for callback in _in_memory_loggers:
if isinstance(callback, DeepEvalLogger):

View file

@ -3,3 +3,5 @@ model_list:
litellm_params:
model: openai/*
api_base: https://exampleopenaiendpoint-production-0ee2.up.railway.app/
litellm_settings:
callbacks: ["cloudzero"]

View file

@ -248,7 +248,9 @@ from litellm.proxy.management_endpoints.customer_endpoints import (
from litellm.proxy.management_endpoints.internal_user_endpoints import (
router as internal_user_router,
)
from litellm.proxy.management_endpoints.internal_user_endpoints import user_update
from litellm.proxy.management_endpoints.internal_user_endpoints import (
user_update,
)
from litellm.proxy.management_endpoints.key_management_endpoints import (
delete_verification_tokens,
duration_in_seconds,
@ -295,7 +297,9 @@ from litellm.proxy.middleware.prometheus_auth_middleware import PrometheusAuthMi
from litellm.proxy.openai_files_endpoints.files_endpoints import (
router as openai_files_router,
)
from litellm.proxy.openai_files_endpoints.files_endpoints import set_files_config
from litellm.proxy.openai_files_endpoints.files_endpoints import (
set_files_config,
)
from litellm.proxy.pass_through_endpoints.llm_passthrough_endpoints import (
passthrough_endpoint_router,
)
@ -3807,13 +3811,13 @@ class ProxyStartupEvent:
########################################################
# CloudZero Background Job
########################################################
from litellm.integrations.cloudzero.cloudzero import CloudZeroLogger
from litellm.proxy.spend_tracking.cloudzero_endpoints import (
init_cloudzero_background_job,
is_cloudzero_setup_in_db,
is_cloudzero_setup,
)
if await is_cloudzero_setup_in_db():
await init_cloudzero_background_job()
if await is_cloudzero_setup():
await CloudZeroLogger.init_cloudzero_background_job(scheduler=scheduler)
########################################################
# Prometheus Background Job

View file

@ -82,14 +82,8 @@ async def _get_cloudzero_settings():
cloudzero_config = await prisma_client.db.litellm_config.find_first(
where={"param_name": "cloudzero_settings"}
)
if not cloudzero_config or not cloudzero_config.param_value:
raise HTTPException(
status_code=400,
detail={
"error": "CloudZero settings not configured. Please run /cloudzero/init first."
},
)
if cloudzero_config is None:
return {}
settings = dict(cloudzero_config.param_value)
@ -257,62 +251,6 @@ async def update_cloudzero_settings(
_cloudzero_background_job_initialized = False
async def init_cloudzero_background_job():
"""
Initialize CloudZero background job if not already initialized.
This should be called from the proxy server startup.
"""
global _cloudzero_background_job_initialized
if _cloudzero_background_job_initialized:
verbose_proxy_logger.debug(
"CloudZero background job already initialized, skipping"
)
return
try:
from litellm.proxy.proxy_server import prisma_client
if prisma_client is None:
verbose_proxy_logger.warning(
"Prisma client not available, skipping CloudZero background job initialization"
)
return
# Get CloudZero settings from database
cloudzero_config = await prisma_client.db.litellm_config.find_first(
where={"param_name": "cloudzero_settings"}
)
if not cloudzero_config or not cloudzero_config.param_value:
verbose_proxy_logger.debug(
"CloudZero settings not configured, skipping background job initialization"
)
return
settings = dict(cloudzero_config.param_value)
# Initialize CloudZero logger with credentials
from litellm.integrations.cloudzero.cloudzero import CloudZeroLogger
logger = CloudZeroLogger(
api_key=settings["api_key"],
connection_id=settings["connection_id"],
timezone=settings["timezone"],
)
# Initialize the background job
#await logger.init_background_job()
_cloudzero_background_job_initialized = True
verbose_proxy_logger.info("CloudZero background job initialized successfully")
except Exception as e:
verbose_proxy_logger.error(
f"Error initializing CloudZero background job: {str(e)}"
)
async def is_cloudzero_setup_in_db() -> bool:
"""
Check if CloudZero is setup in the database.
@ -343,6 +281,47 @@ async def is_cloudzero_setup_in_db() -> bool:
return False
def is_cloudzero_setup_in_config() -> bool:
"""
Check if CloudZero is setup in config.yaml or environment variables.
CloudZero is considered setup in config if:
- "cloudzero" is in the callbacks list in config.yaml, OR
Returns:
bool: True if CloudZero is configured, False otherwise
"""
import litellm
return "cloudzero" in litellm.callbacks
async def is_cloudzero_setup() -> bool:
"""
Check if CloudZero is setup in either config.yaml/env vars OR database.
CloudZero is considered setup if:
- CloudZero is configured in config.yaml callbacks, OR
- CloudZero environment variables are set, OR
- CloudZero settings exist in the database
Returns:
bool: True if CloudZero is configured anywhere, False otherwise
"""
try:
# Check config.yaml/environment variables first
if is_cloudzero_setup_in_config():
return True
# Check database as fallback
if await is_cloudzero_setup_in_db():
return True
return False
except Exception as e:
verbose_proxy_logger.error(f"Error checking CloudZero setup: {str(e)}")
return False
@router.post(
"/cloudzero/init",
tags=["CloudZero"],
@ -383,9 +362,6 @@ async def init_cloudzero_settings(
verbose_proxy_logger.info("CloudZero settings initialized successfully")
# Initialize background job after settings are saved
await init_cloudzero_background_job()
return CloudZeroInitResponse(
message="CloudZero settings initialized successfully", status="success"
)
@ -412,15 +388,18 @@ async def cloudzero_dry_run_export(
Perform a dry run export using the CloudZero logger.
This endpoint uses the CloudZero logger to perform a dry run export,
which displays the data that would be exported without actually sending it to CloudZero.
which returns the data that would be exported without actually sending it to CloudZero.
Parameters:
- limit: Optional limit on number of records to process (default: 10000)
Returns:
- usage_data: Sample of the raw usage data (first 50 records)
- cbf_data: CloudZero CBF formatted data ready for export
- summary: Statistics including total cost, tokens, and record counts
Only admin users can perform CloudZero exports.
"""
from datetime import datetime
# Validation
if user_api_key_dict.user_role != LitellmUserRoles.PROXY_ADMIN:
raise HTTPException(
@ -430,19 +409,21 @@ async def cloudzero_dry_run_export(
try:
# Import and initialize CloudZero logger with credentials
from litellm.integrations.cloudzero.ll2cz.cloudzero import CloudZeroLogger
from litellm.integrations.cloudzero.cloudzero import CloudZeroLogger
# Initialize logger with credentials directly
logger = CloudZeroLogger()
await logger.dry_run_export_usage_data(
target_hour=datetime.utcnow(), limit=request.limit
dry_run_result = await logger.dry_run_export_usage_data(
limit=request.limit
)
verbose_proxy_logger.info("CloudZero dry run export completed successfully")
return CloudZeroExportResponse(
message="CloudZero dry run export completed successfully. Check logs for output.",
message="CloudZero dry run export completed successfully.",
status="success",
dry_run_data=dry_run_result,
summary=dry_run_result.get("summary") if dry_run_result else None,
)
except Exception as e:
@ -477,7 +458,6 @@ async def cloudzero_export(
Only admin users can perform CloudZero exports.
"""
from datetime import datetime
if user_api_key_dict.user_role != LitellmUserRoles.PROXY_ADMIN:
raise HTTPException(
@ -490,24 +470,28 @@ async def cloudzero_export(
settings = await _get_cloudzero_settings()
# Import and initialize CloudZero logger with credentials
from litellm.integrations.cloudzero.ll2cz.cloudzero import CloudZeroLogger
from litellm.integrations.cloudzero.cloudzero import CloudZeroLogger
# Initialize logger with credentials directly
logger = CloudZeroLogger(
api_key=settings["api_key"],
connection_id=settings["connection_id"],
timezone=settings["timezone"],
api_key=settings.get("api_key"),
connection_id=settings.get("connection_id"),
timezone=settings.get("timezone"),
)
await logger.export_usage_data(
target_hour=datetime.utcnow(),
limit=request.limit,
operation=request.operation,
start_time_utc=request.start_time_utc,
end_time_utc=request.end_time_utc,
)
verbose_proxy_logger.info("CloudZero export completed successfully")
return CloudZeroExportResponse(
message="CloudZero export completed successfully", status="success"
message="CloudZero export completed successfully",
status="success",
dry_run_data=None,
summary=None
)
except Exception as e:

View file

@ -2,7 +2,8 @@
CloudZero endpoint types for LiteLLM Proxy
"""
from typing import Optional
from datetime import datetime
from typing import Any, Dict, List, Optional
from pydantic import BaseModel, Field
@ -27,6 +28,8 @@ class CloudZeroExportRequest(BaseModel):
limit: Optional[int] = Field(None, description="Optional limit on number of records to export")
operation: str = Field(default="replace_hourly", description="CloudZero operation type (replace_hourly or sum)")
start_time_utc: Optional[datetime] = Field(None, description="Start time for data export in UTC")
end_time_utc: Optional[datetime] = Field(None, description="End time for data export in UTC")
class CloudZeroExportResponse(BaseModel):
@ -35,6 +38,8 @@ class CloudZeroExportResponse(BaseModel):
message: str
status: str
records_exported: Optional[int] = None
dry_run_data: Optional[Dict[str, Any]] = Field(None, description="Dry run data including usage data and CBF transformed data")
summary: Optional[Dict[str, Any]] = Field(None, description="Summary statistics for dry run")
class CloudZeroSettingsView(BaseModel):

View file

@ -0,0 +1,163 @@
"""
Test the CloudZero dry run endpoint functionality
"""
import os
import sys
from unittest.mock import AsyncMock, MagicMock, patch
import polars as pl
import pytest
sys.path.insert(0, os.path.abspath("../../../.."))
from litellm.integrations.cloudzero.cloudzero import CloudZeroLogger
class TestCloudZeroDryRunEndpoint:
"""Test suite for CloudZero dry run endpoint functionality."""
@pytest.mark.asyncio
async def test_dry_run_export_usage_data_returns_data(self):
"""
Test that dry_run_export_usage_data returns expected data structure
instead of just logging to console.
"""
logger = CloudZeroLogger()
# Mock database data
mock_usage_data = pl.DataFrame({
'date': ['2025-01-19', '2025-01-20'],
'model': ['gpt-4', 'gpt-3.5-turbo'],
'custom_llm_provider': ['openai', 'openai'],
'team_id': ['team1', 'team2'],
'team_alias': ['Team One', 'Team Two'],
'api_key_alias': ['key1', 'key2'],
'prompt_tokens': [100, 200],
'completion_tokens': [50, 100],
'spend': [0.01, 0.02],
'successful_requests': [1, 2]
})
# Mock CBF transformed data
mock_cbf_data = pl.DataFrame({
'time/usage_start': ['2025-01-19T00:00:00Z', '2025-01-20T00:00:00Z'],
'cost/cost': [0.01, 0.02],
'usage/amount': [150, 300],
'resource/service': ['openai', 'openai'],
'resource/account': ['litellm', 'litellm'],
'resource/region': ['us-east-1', 'us-east-1'],
'resource/id': ['gpt-4', 'gpt-3.5-turbo'],
'entity_type': ['user', 'user'],
'entity_id': ['team1', 'team2'],
'resource/tag:team_id': ['team1', 'team2'],
'resource/tag:team_alias': ['Team One', 'Team Two'],
'resource/tag:api_key_alias': ['key1', 'key2']
})
with patch('litellm.integrations.cloudzero.cloudzero.LiteLLMDatabase') as mock_db_class, \
patch('litellm.integrations.cloudzero.cloudzero.CBFTransformer') as mock_transformer_class:
# Setup mocks
mock_db = AsyncMock()
mock_db.get_usage_data.return_value = mock_usage_data
mock_db_class.return_value = mock_db
mock_transformer = MagicMock()
mock_transformer.transform.return_value = mock_cbf_data
mock_transformer_class.return_value = mock_transformer
# Call the method
result = await logger.dry_run_export_usage_data(limit=1000)
# Verify the result structure
assert isinstance(result, dict)
assert 'usage_data' in result
assert 'cbf_data' in result
assert 'summary' in result
# Verify usage_data
assert isinstance(result['usage_data'], list)
assert len(result['usage_data']) == 2
assert result['usage_data'][0]['model'] == 'gpt-4'
assert result['usage_data'][1]['model'] == 'gpt-3.5-turbo'
# Verify cbf_data
assert isinstance(result['cbf_data'], list)
assert len(result['cbf_data']) == 2
assert result['cbf_data'][0]['cost/cost'] == 0.01
assert result['cbf_data'][1]['cost/cost'] == 0.02
# Verify summary
summary = result['summary']
assert summary['total_records'] == 2
assert summary['total_cost'] == 0.03
assert summary['total_tokens'] == 450 # 150 + 300
assert summary['unique_accounts'] == 1
assert summary['unique_services'] == 1
@pytest.mark.asyncio
async def test_dry_run_export_usage_data_empty_data(self):
"""
Test that dry_run_export_usage_data handles empty data gracefully.
"""
logger = CloudZeroLogger()
# Mock empty database data
mock_empty_data = pl.DataFrame()
with patch('litellm.integrations.cloudzero.cloudzero.LiteLLMDatabase') as mock_db_class:
# Setup mocks
mock_db = AsyncMock()
mock_db.get_usage_data.return_value = mock_empty_data
mock_db_class.return_value = mock_db
# Call the method
result = await logger.dry_run_export_usage_data(limit=1000)
# Verify the result structure for empty data
assert isinstance(result, dict)
assert result['usage_data'] == []
assert result['cbf_data'] == []
assert result['summary']['total_records'] == 0
assert result['summary']['total_cost'] == 0
assert result['summary']['total_tokens'] == 0
@pytest.mark.asyncio
async def test_dry_run_export_usage_data_cbf_transformation_failure(self):
"""
Test that dry_run_export_usage_data handles CBF transformation failure gracefully.
"""
logger = CloudZeroLogger()
# Mock database data
mock_usage_data = pl.DataFrame({
'date': ['2025-01-19'],
'model': ['gpt-4'],
'spend': [0.01],
'successful_requests': [1]
})
# Mock empty CBF data (transformation failed)
mock_empty_cbf_data = pl.DataFrame()
with patch('litellm.integrations.cloudzero.cloudzero.LiteLLMDatabase') as mock_db_class, \
patch('litellm.integrations.cloudzero.cloudzero.CBFTransformer') as mock_transformer_class:
# Setup mocks
mock_db = AsyncMock()
mock_db.get_usage_data.return_value = mock_usage_data
mock_db_class.return_value = mock_db
mock_transformer = MagicMock()
mock_transformer.transform.return_value = mock_empty_cbf_data
mock_transformer_class.return_value = mock_transformer
# Call the method
result = await logger.dry_run_export_usage_data(limit=1000)
# Verify the result handles CBF transformation failure
assert isinstance(result, dict)
assert len(result['usage_data']) == 1 # Usage data should still be present
assert result['cbf_data'] == [] # CBF data should be empty
assert result['summary']['total_cost'] == 0.01 # Should calculate from usage data