mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-10-03 07:12:50 +00:00
feat(usage): preserve cache tokens in usage events (#1311)
This commit is contained in:
1 parent
89b83af931
commit
606c8e60e1
10 files changed
+230
-93
No files matched your search
@@ -86,6 +86,9 @@ func (l *Loop) recordToolUsageEvent(ctx context.Context, req *RunRequest, canoni
|
||||
event.InputTokens = int64(result.Usage.PromptTokens)
|
||||
event.OutputTokens = int64(result.Usage.CompletionTokens)
|
||||
event.TotalTokens = int64(result.Usage.TotalTokens)
|
||||
event.CacheReadTokens = int64(result.Usage.CacheReadTokens)
|
||||
event.CacheCreateTokens = int64(result.Usage.CacheCreationTokens)
|
||||
event.ThinkingTokens = int64(result.Usage.ThinkingTokens)
|
||||
if event.TotalTokens == 0 {
|
||||
event.TotalTokens = event.InputTokens + event.OutputTokens
|
||||
}
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"github.com/nextlevelbuilder/goclaw/internal/providers"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/store"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/tools"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/tracing"
|
||||
@@ -117,6 +118,34 @@ func TestRecordToolUsageEvent_UseSkillCountsSkillName(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestRecordToolUsageEvent_PreservesProviderCacheUsage(t *testing.T) {
|
||||
storeSpy := newFakeUsageEventStore()
|
||||
loop := &Loop{usageEvents: storeSpy, agentUUID: uuid.New(), tenantID: uuid.New()}
|
||||
ctx := tracing.WithTraceID(store.WithTenantID(t.Context(), loop.tenantID), uuid.New())
|
||||
result := tools.NewResult("ok")
|
||||
result.Provider = "openai-codex"
|
||||
result.Model = "gpt-5.5"
|
||||
result.Usage = &providers.Usage{
|
||||
PromptTokens: 100,
|
||||
CompletionTokens: 20,
|
||||
TotalTokens: 120,
|
||||
CacheReadTokens: 80,
|
||||
CacheCreationTokens: 10,
|
||||
ThinkingTokens: 5,
|
||||
}
|
||||
|
||||
loop.recordToolUsageEvent(ctx, &RunRequest{RunID: "run-1", SessionKey: "session-1"}, "read_image", "read_image", "call-1",
|
||||
nil, time.Now().Add(-time.Millisecond), result, uuid.New())
|
||||
|
||||
event := waitUsageEvent(t, storeSpy)
|
||||
if event.InputTokens != 100 || event.OutputTokens != 20 || event.TotalTokens != 120 {
|
||||
t.Fatalf("tokens = input %d output %d total %d", event.InputTokens, event.OutputTokens, event.TotalTokens)
|
||||
}
|
||||
if event.CacheReadTokens != 80 || event.CacheCreateTokens != 10 || event.ThinkingTokens != 5 {
|
||||
t.Fatalf("cache/thinking tokens = read %d create %d thinking %d", event.CacheReadTokens, event.CacheCreateTokens, event.ThinkingTokens)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRecordSkillSlashUsageEvent_RequiresTraceContext(t *testing.T) {
|
||||
storeSpy := newFakeUsageEventStore()
|
||||
loop := &Loop{usageEvents: storeSpy, agentUUID: uuid.New(), tenantID: uuid.New()}
|
||||
|
||||
@@ -20,8 +20,8 @@ func NewPGUsageEventStore(db *sql.DB) *PGUsageEventStore {
|
||||
return &PGUsageEventStore{db: db}
|
||||
}
|
||||
|
||||
const usageEventFieldCount = 28
|
||||
const usageRollupFieldCount = 21
|
||||
const usageEventFieldCount = 31
|
||||
const usageRollupFieldCount = 24
|
||||
|
||||
func (s *PGUsageEventStore) InsertEvent(ctx context.Context, event *store.UsageEvent) error {
|
||||
if event == nil {
|
||||
@@ -52,7 +52,8 @@ func (s *PGUsageEventStore) InsertEvents(ctx context.Context, events []store.Usa
|
||||
event.EventType, event.ResourceType, event.ResourceName, event.ResourceID, event.Source,
|
||||
nilUUID(event.AgentID), nilUUID(event.TeamID), nilUUID(event.TraceID), nilUUID(event.SpanID),
|
||||
event.RunID, event.SessionKey, event.Channel, event.Provider, event.Model, event.Status,
|
||||
event.InputTokens, event.OutputTokens, event.TotalTokens, event.CostUSD,
|
||||
event.InputTokens, event.OutputTokens, event.TotalTokens,
|
||||
event.CacheReadTokens, event.CacheCreateTokens, event.ThinkingTokens, event.CostUSD,
|
||||
event.DurationMS, event.CallCount, event.ErrorCount, jsonOrNull(event.Metadata), event.CreatedAt,
|
||||
)
|
||||
}
|
||||
@@ -62,7 +63,8 @@ func (s *PGUsageEventStore) InsertEvents(ctx context.Context, events []store.Usa
|
||||
event_type, resource_type, resource_name, resource_id, source,
|
||||
agent_id, team_id, trace_id, span_id,
|
||||
run_id, session_key, channel, provider, model, status,
|
||||
input_tokens, output_tokens, total_tokens, cost_usd,
|
||||
input_tokens, output_tokens, total_tokens,
|
||||
cache_read_tokens, cache_create_tokens, thinking_tokens, cost_usd,
|
||||
duration_ms, call_count, error_count, metadata, created_at
|
||||
) VALUES ` + strings.Join(vals, ", ") + `
|
||||
ON CONFLICT DO NOTHING`
|
||||
@@ -88,6 +90,9 @@ func (s *PGUsageEventStore) RefreshEventRollupHour(ctx context.Context, bucketHo
|
||||
COALESCE(SUM(input_tokens), 0),
|
||||
COALESCE(SUM(output_tokens), 0),
|
||||
COALESCE(SUM(total_tokens), 0),
|
||||
COALESCE(SUM(cache_read_tokens), 0),
|
||||
COALESCE(SUM(cache_create_tokens), 0),
|
||||
COALESCE(SUM(thinking_tokens), 0),
|
||||
COALESCE(SUM(cost_usd), 0),
|
||||
CASE WHEN COALESCE(SUM(call_count), 0) > 0
|
||||
THEN COALESCE(SUM(duration_ms * call_count), 0) / SUM(call_count)
|
||||
@@ -111,7 +116,8 @@ func (s *PGUsageEventStore) RefreshEventRollupHour(ctx context.Context, bucketHo
|
||||
&rollup.TenantID, &rollup.BucketHour, &rollup.EventType, &rollup.ResourceType,
|
||||
&rollup.ResourceName, &rollup.Source, &rollup.AgentID, &rollup.Channel,
|
||||
&rollup.Provider, &rollup.Model, &rollup.Status,
|
||||
&rollup.InputTokens, &rollup.OutputTokens, &rollup.TotalTokens, &rollup.CostUSD,
|
||||
&rollup.InputTokens, &rollup.OutputTokens, &rollup.TotalTokens,
|
||||
&rollup.CacheReadTokens, &rollup.CacheCreateTokens, &rollup.ThinkingTokens, &rollup.CostUSD,
|
||||
&rollup.DurationMS, &rollup.CallCount, &rollup.ErrorCount,
|
||||
); err != nil {
|
||||
return fmt.Errorf("scan usage event rollup: %w", err)
|
||||
@@ -153,14 +159,16 @@ func (s *PGUsageEventStore) upsertEventRollups(ctx context.Context, rollups []st
|
||||
rollup.ID, rollup.TenantID, rollup.BucketHour, rollup.EventType, rollup.ResourceType,
|
||||
rollup.ResourceName, rollup.Source, nilUUID(rollup.AgentID), rollup.Channel,
|
||||
rollup.Provider, rollup.Model, rollup.Status, rollup.InputTokens, rollup.OutputTokens,
|
||||
rollup.TotalTokens, rollup.CostUSD, rollup.DurationMS, rollup.CallCount,
|
||||
rollup.TotalTokens, rollup.CacheReadTokens, rollup.CacheCreateTokens, rollup.ThinkingTokens,
|
||||
rollup.CostUSD, rollup.DurationMS, rollup.CallCount,
|
||||
rollup.ErrorCount, rollup.CreatedAt, rollup.UpdatedAt,
|
||||
)
|
||||
}
|
||||
query := `INSERT INTO usage_event_rollups (
|
||||
id, tenant_id, bucket_hour, event_type, resource_type, resource_name, source,
|
||||
agent_id, channel, provider, model, status,
|
||||
input_tokens, output_tokens, total_tokens, cost_usd,
|
||||
input_tokens, output_tokens, total_tokens,
|
||||
cache_read_tokens, cache_create_tokens, thinking_tokens, cost_usd,
|
||||
duration_ms, call_count, error_count, created_at, updated_at
|
||||
) VALUES ` + strings.Join(vals, ", ") + `
|
||||
ON CONFLICT (
|
||||
@@ -179,6 +187,9 @@ func (s *PGUsageEventStore) upsertEventRollups(ctx context.Context, rollups []st
|
||||
input_tokens = EXCLUDED.input_tokens,
|
||||
output_tokens = EXCLUDED.output_tokens,
|
||||
total_tokens = EXCLUDED.total_tokens,
|
||||
cache_read_tokens = EXCLUDED.cache_read_tokens,
|
||||
cache_create_tokens = EXCLUDED.cache_create_tokens,
|
||||
thinking_tokens = EXCLUDED.thinking_tokens,
|
||||
cost_usd = EXCLUDED.cost_usd,
|
||||
duration_ms = EXCLUDED.duration_ms,
|
||||
call_count = EXCLUDED.call_count,
|
||||
@@ -201,6 +212,9 @@ func (s *PGUsageEventStore) GetEventTimeSeries(ctx context.Context, q store.Usag
|
||||
COALESCE(SUM(input_tokens), 0),
|
||||
COALESCE(SUM(output_tokens), 0),
|
||||
COALESCE(SUM(total_tokens), 0),
|
||||
COALESCE(SUM(cache_read_tokens), 0),
|
||||
COALESCE(SUM(cache_create_tokens), 0),
|
||||
COALESCE(SUM(thinking_tokens), 0),
|
||||
COALESCE(SUM(cost_usd), 0),
|
||||
CASE WHEN COALESCE(SUM(call_count), 0) > 0
|
||||
THEN COALESCE(SUM(duration_ms * call_count), 0) / SUM(call_count)
|
||||
@@ -222,6 +236,7 @@ func (s *PGUsageEventStore) GetEventTimeSeries(ctx context.Context, q store.Usag
|
||||
if err := rows.Scan(
|
||||
&point.BucketTime, &point.Calls, &point.Errors,
|
||||
&point.InputTokens, &point.OutputTokens, &point.TotalTokens,
|
||||
&point.CacheReadTokens, &point.CacheCreateTokens, &point.ThinkingTokens,
|
||||
&point.CostUSD, &point.AvgDurationMS,
|
||||
); err != nil {
|
||||
return nil, fmt.Errorf("scan usage event timeseries: %w", err)
|
||||
@@ -255,6 +270,9 @@ func (s *PGUsageEventStore) GetEventBreakdown(ctx context.Context, q store.Usage
|
||||
COALESCE(SUM(input_tokens), 0),
|
||||
COALESCE(SUM(output_tokens), 0),
|
||||
COALESCE(SUM(total_tokens), 0),
|
||||
COALESCE(SUM(cache_read_tokens), 0),
|
||||
COALESCE(SUM(cache_create_tokens), 0),
|
||||
COALESCE(SUM(thinking_tokens), 0),
|
||||
COALESCE(SUM(cost_usd), 0),
|
||||
CASE WHEN COALESCE(SUM(call_count), 0) > 0
|
||||
THEN COALESCE(SUM(duration_ms * call_count), 0) / SUM(call_count)
|
||||
@@ -277,6 +295,7 @@ func (s *PGUsageEventStore) GetEventBreakdown(ctx context.Context, q store.Usage
|
||||
if err := rows.Scan(
|
||||
&row.Key, &row.EventType, &row.ResourceType, &row.ResourceName, &row.Source,
|
||||
&row.Calls, &row.Errors, &row.InputTokens, &row.OutputTokens, &row.TotalTokens,
|
||||
&row.CacheReadTokens, &row.CacheCreateTokens, &row.ThinkingTokens,
|
||||
&row.CostUSD, &row.AvgDurationMS,
|
||||
); err != nil {
|
||||
return nil, fmt.Errorf("scan usage event breakdown: %w", err)
|
||||
@@ -294,6 +313,9 @@ func (s *PGUsageEventStore) GetEventSummary(ctx context.Context, q store.UsageEv
|
||||
COALESCE(SUM(input_tokens), 0),
|
||||
COALESCE(SUM(output_tokens), 0),
|
||||
COALESCE(SUM(total_tokens), 0),
|
||||
COALESCE(SUM(cache_read_tokens), 0),
|
||||
COALESCE(SUM(cache_create_tokens), 0),
|
||||
COALESCE(SUM(thinking_tokens), 0),
|
||||
COALESCE(SUM(cost_usd), 0),
|
||||
CASE WHEN COALESCE(SUM(call_count), 0) > 0
|
||||
THEN COALESCE(SUM(duration_ms * call_count), 0) / SUM(call_count)
|
||||
@@ -302,7 +324,8 @@ func (s *PGUsageEventStore) GetEventSummary(ctx context.Context, q store.UsageEv
|
||||
var summary store.UsageEventSummary
|
||||
if err := s.db.QueryRowContext(ctx, query, args...).Scan(
|
||||
&summary.Calls, &summary.Errors, &summary.InputTokens, &summary.OutputTokens,
|
||||
&summary.TotalTokens, &summary.CostUSD, &summary.AvgDurationMS,
|
||||
&summary.TotalTokens, &summary.CacheReadTokens, &summary.CacheCreateTokens, &summary.ThinkingTokens,
|
||||
&summary.CostUSD, &summary.AvgDurationMS,
|
||||
); err != nil {
|
||||
return nil, fmt.Errorf("get usage event summary: %w", err)
|
||||
}
|
||||
|
||||
@@ -16,7 +16,7 @@ var schemaSQL string
|
||||
|
||||
// SchemaVersion is the current SQLite schema version.
|
||||
// Bump this when adding new migration steps below.
|
||||
const SchemaVersion = 52
|
||||
const SchemaVersion = 53
|
||||
|
||||
// migrations maps version → SQL to apply when upgrading FROM that version.
|
||||
// schema.sql always represents the LATEST full schema (for fresh DBs).
|
||||
@@ -884,6 +884,13 @@ CREATE INDEX IF NOT EXISTS idx_mcp_oauth_tokens_server_tenant ON mcp_oauth_token
|
||||
// Version 51 → 52: add last_heartbeat_at to webhook_calls for lease heartbeat.
|
||||
// Mirrors PG migration 000085. Idempotent-guarded via idempotentColumnMigration(51).
|
||||
51: `ALTER TABLE webhook_calls ADD COLUMN last_heartbeat_at TEXT;`,
|
||||
// Version 52 → 53: preserve provider cache/thinking token dimensions in usage event analytics.
|
||||
52: `ALTER TABLE usage_events ADD COLUMN cache_read_tokens BIGINT NOT NULL DEFAULT 0;
|
||||
ALTER TABLE usage_events ADD COLUMN cache_create_tokens BIGINT NOT NULL DEFAULT 0;
|
||||
ALTER TABLE usage_events ADD COLUMN thinking_tokens BIGINT NOT NULL DEFAULT 0;
|
||||
ALTER TABLE usage_event_rollups ADD COLUMN cache_read_tokens BIGINT NOT NULL DEFAULT 0;
|
||||
ALTER TABLE usage_event_rollups ADD COLUMN cache_create_tokens BIGINT NOT NULL DEFAULT 0;
|
||||
ALTER TABLE usage_event_rollups ADD COLUMN thinking_tokens BIGINT NOT NULL DEFAULT 0;`,
|
||||
}
|
||||
|
||||
const addUsageEventAnalyticsTables = `
|
||||
@@ -910,6 +917,9 @@ CREATE TABLE IF NOT EXISTS usage_events (
|
||||
input_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
output_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
total_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cache_read_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cache_create_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
thinking_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cost_usd NUMERIC(12,6) NOT NULL DEFAULT 0,
|
||||
duration_ms INTEGER NOT NULL DEFAULT 0,
|
||||
call_count INTEGER NOT NULL DEFAULT 1,
|
||||
@@ -946,6 +956,9 @@ CREATE TABLE IF NOT EXISTS usage_event_rollups (
|
||||
input_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
output_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
total_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cache_read_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cache_create_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
thinking_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cost_usd NUMERIC(12,6) NOT NULL DEFAULT 0,
|
||||
duration_ms INTEGER NOT NULL DEFAULT 0,
|
||||
call_count INTEGER NOT NULL DEFAULT 0,
|
||||
|
||||
@@ -1377,6 +1377,9 @@ CREATE TABLE IF NOT EXISTS usage_events (
|
||||
input_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
output_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
total_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cache_read_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cache_create_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
thinking_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cost_usd NUMERIC(12,6) NOT NULL DEFAULT 0,
|
||||
duration_ms INTEGER NOT NULL DEFAULT 0,
|
||||
call_count INTEGER NOT NULL DEFAULT 1,
|
||||
@@ -1419,6 +1422,9 @@ CREATE TABLE IF NOT EXISTS usage_event_rollups (
|
||||
input_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
output_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
total_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cache_read_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cache_create_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
thinking_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
cost_usd NUMERIC(12,6) NOT NULL DEFAULT 0,
|
||||
duration_ms INTEGER NOT NULL DEFAULT 0,
|
||||
call_count INTEGER NOT NULL DEFAULT 0,
|
||||
|
||||
@@ -22,9 +22,9 @@ func NewSQLiteUsageEventStore(db *sql.DB) *SQLiteUsageEventStore {
|
||||
return &SQLiteUsageEventStore{db: db}
|
||||
}
|
||||
|
||||
const sqliteUsageEventFieldCount = 28
|
||||
const sqliteUsageEventFieldCount = 31
|
||||
const sqliteUsageEventBatchSize = 30
|
||||
const sqliteUsageRollupFieldCount = 21
|
||||
const sqliteUsageRollupFieldCount = 24
|
||||
|
||||
func (s *SQLiteUsageEventStore) InsertEvent(ctx context.Context, event *store.UsageEvent) error {
|
||||
if event == nil {
|
||||
@@ -62,7 +62,8 @@ func (s *SQLiteUsageEventStore) insertBatch(ctx context.Context, events []store.
|
||||
event.EventType, event.ResourceType, event.ResourceName, event.ResourceID, event.Source,
|
||||
nilUUID(event.AgentID), nilUUID(event.TeamID), nilUUID(event.TraceID), nilUUID(event.SpanID),
|
||||
event.RunID, event.SessionKey, event.Channel, event.Provider, event.Model, event.Status,
|
||||
event.InputTokens, event.OutputTokens, event.TotalTokens, event.CostUSD,
|
||||
event.InputTokens, event.OutputTokens, event.TotalTokens,
|
||||
event.CacheReadTokens, event.CacheCreateTokens, event.ThinkingTokens, event.CostUSD,
|
||||
event.DurationMS, event.CallCount, event.ErrorCount, jsonOrNull(event.Metadata), event.CreatedAt,
|
||||
)
|
||||
}
|
||||
@@ -71,7 +72,8 @@ func (s *SQLiteUsageEventStore) insertBatch(ctx context.Context, events []store.
|
||||
event_type, resource_type, resource_name, resource_id, source,
|
||||
agent_id, team_id, trace_id, span_id,
|
||||
run_id, session_key, channel, provider, model, status,
|
||||
input_tokens, output_tokens, total_tokens, cost_usd,
|
||||
input_tokens, output_tokens, total_tokens,
|
||||
cache_read_tokens, cache_create_tokens, thinking_tokens, cost_usd,
|
||||
duration_ms, call_count, error_count, metadata, created_at
|
||||
) VALUES ` + strings.Join(vals, ", ") + `
|
||||
ON CONFLICT DO NOTHING`
|
||||
@@ -97,6 +99,9 @@ func (s *SQLiteUsageEventStore) RefreshEventRollupHour(ctx context.Context, buck
|
||||
COALESCE(SUM(input_tokens), 0),
|
||||
COALESCE(SUM(output_tokens), 0),
|
||||
COALESCE(SUM(total_tokens), 0),
|
||||
COALESCE(SUM(cache_read_tokens), 0),
|
||||
COALESCE(SUM(cache_create_tokens), 0),
|
||||
COALESCE(SUM(thinking_tokens), 0),
|
||||
COALESCE(SUM(cost_usd), 0),
|
||||
CASE WHEN COALESCE(SUM(call_count), 0) > 0
|
||||
THEN COALESCE(SUM(duration_ms * call_count), 0) / SUM(call_count)
|
||||
@@ -121,7 +126,8 @@ func (s *SQLiteUsageEventStore) RefreshEventRollupHour(ctx context.Context, buck
|
||||
&rollup.TenantID, &bucketTime, &rollup.EventType, &rollup.ResourceType,
|
||||
&rollup.ResourceName, &rollup.Source, &rollup.AgentID, &rollup.Channel,
|
||||
&rollup.Provider, &rollup.Model, &rollup.Status,
|
||||
&rollup.InputTokens, &rollup.OutputTokens, &rollup.TotalTokens, &rollup.CostUSD,
|
||||
&rollup.InputTokens, &rollup.OutputTokens, &rollup.TotalTokens,
|
||||
&rollup.CacheReadTokens, &rollup.CacheCreateTokens, &rollup.ThinkingTokens, &rollup.CostUSD,
|
||||
&rollup.DurationMS, &rollup.CallCount, &rollup.ErrorCount,
|
||||
); err != nil {
|
||||
return fmt.Errorf("scan usage event rollup: %w", err)
|
||||
@@ -160,14 +166,16 @@ func (s *SQLiteUsageEventStore) upsertEventRollups(ctx context.Context, rollups
|
||||
rollup.ID, rollup.TenantID, rollup.BucketHour, rollup.EventType, rollup.ResourceType,
|
||||
rollup.ResourceName, rollup.Source, nilUUID(rollup.AgentID), rollup.Channel,
|
||||
rollup.Provider, rollup.Model, rollup.Status, rollup.InputTokens, rollup.OutputTokens,
|
||||
rollup.TotalTokens, rollup.CostUSD, rollup.DurationMS, rollup.CallCount,
|
||||
rollup.TotalTokens, rollup.CacheReadTokens, rollup.CacheCreateTokens, rollup.ThinkingTokens,
|
||||
rollup.CostUSD, rollup.DurationMS, rollup.CallCount,
|
||||
rollup.ErrorCount, rollup.CreatedAt, rollup.UpdatedAt,
|
||||
)
|
||||
}
|
||||
query := `INSERT INTO usage_event_rollups (
|
||||
id, tenant_id, bucket_hour, event_type, resource_type, resource_name, source,
|
||||
agent_id, channel, provider, model, status,
|
||||
input_tokens, output_tokens, total_tokens, cost_usd,
|
||||
input_tokens, output_tokens, total_tokens,
|
||||
cache_read_tokens, cache_create_tokens, thinking_tokens, cost_usd,
|
||||
duration_ms, call_count, error_count, created_at, updated_at
|
||||
) VALUES ` + strings.Join(vals, ", ") + `
|
||||
ON CONFLICT (
|
||||
@@ -186,6 +194,9 @@ func (s *SQLiteUsageEventStore) upsertEventRollups(ctx context.Context, rollups
|
||||
input_tokens = excluded.input_tokens,
|
||||
output_tokens = excluded.output_tokens,
|
||||
total_tokens = excluded.total_tokens,
|
||||
cache_read_tokens = excluded.cache_read_tokens,
|
||||
cache_create_tokens = excluded.cache_create_tokens,
|
||||
thinking_tokens = excluded.thinking_tokens,
|
||||
cost_usd = excluded.cost_usd,
|
||||
duration_ms = excluded.duration_ms,
|
||||
call_count = excluded.call_count,
|
||||
@@ -208,6 +219,9 @@ func (s *SQLiteUsageEventStore) GetEventTimeSeries(ctx context.Context, q store.
|
||||
COALESCE(SUM(input_tokens), 0),
|
||||
COALESCE(SUM(output_tokens), 0),
|
||||
COALESCE(SUM(total_tokens), 0),
|
||||
COALESCE(SUM(cache_read_tokens), 0),
|
||||
COALESCE(SUM(cache_create_tokens), 0),
|
||||
COALESCE(SUM(thinking_tokens), 0),
|
||||
COALESCE(SUM(cost_usd), 0),
|
||||
CASE WHEN COALESCE(SUM(call_count), 0) > 0
|
||||
THEN COALESCE(SUM(duration_ms * call_count), 0) / SUM(call_count)
|
||||
@@ -230,6 +244,7 @@ func (s *SQLiteUsageEventStore) GetEventTimeSeries(ctx context.Context, q store.
|
||||
if err := rows.Scan(
|
||||
&bucketTime, &point.Calls, &point.Errors,
|
||||
&point.InputTokens, &point.OutputTokens, &point.TotalTokens,
|
||||
&point.CacheReadTokens, &point.CacheCreateTokens, &point.ThinkingTokens,
|
||||
&point.CostUSD, &point.AvgDurationMS,
|
||||
); err != nil {
|
||||
return nil, fmt.Errorf("scan usage event timeseries: %w", err)
|
||||
@@ -262,6 +277,9 @@ func (s *SQLiteUsageEventStore) GetEventBreakdown(ctx context.Context, q store.U
|
||||
COALESCE(SUM(input_tokens), 0),
|
||||
COALESCE(SUM(output_tokens), 0),
|
||||
COALESCE(SUM(total_tokens), 0),
|
||||
COALESCE(SUM(cache_read_tokens), 0),
|
||||
COALESCE(SUM(cache_create_tokens), 0),
|
||||
COALESCE(SUM(thinking_tokens), 0),
|
||||
COALESCE(SUM(cost_usd), 0),
|
||||
CASE WHEN COALESCE(SUM(call_count), 0) > 0
|
||||
THEN COALESCE(SUM(duration_ms * call_count), 0) / SUM(call_count)
|
||||
@@ -284,6 +302,7 @@ func (s *SQLiteUsageEventStore) GetEventBreakdown(ctx context.Context, q store.U
|
||||
if err := rows.Scan(
|
||||
&row.Key, &row.EventType, &row.ResourceType, &row.ResourceName, &row.Source,
|
||||
&row.Calls, &row.Errors, &row.InputTokens, &row.OutputTokens, &row.TotalTokens,
|
||||
&row.CacheReadTokens, &row.CacheCreateTokens, &row.ThinkingTokens,
|
||||
&row.CostUSD, &row.AvgDurationMS,
|
||||
); err != nil {
|
||||
return nil, fmt.Errorf("scan usage event breakdown: %w", err)
|
||||
@@ -301,6 +320,9 @@ func (s *SQLiteUsageEventStore) GetEventSummary(ctx context.Context, q store.Usa
|
||||
COALESCE(SUM(input_tokens), 0),
|
||||
COALESCE(SUM(output_tokens), 0),
|
||||
COALESCE(SUM(total_tokens), 0),
|
||||
COALESCE(SUM(cache_read_tokens), 0),
|
||||
COALESCE(SUM(cache_create_tokens), 0),
|
||||
COALESCE(SUM(thinking_tokens), 0),
|
||||
COALESCE(SUM(cost_usd), 0),
|
||||
CASE WHEN COALESCE(SUM(call_count), 0) > 0
|
||||
THEN COALESCE(SUM(duration_ms * call_count), 0) / SUM(call_count)
|
||||
@@ -309,7 +331,8 @@ func (s *SQLiteUsageEventStore) GetEventSummary(ctx context.Context, q store.Usa
|
||||
var summary store.UsageEventSummary
|
||||
if err := s.db.QueryRowContext(ctx, query, args...).Scan(
|
||||
&summary.Calls, &summary.Errors, &summary.InputTokens, &summary.OutputTokens,
|
||||
&summary.TotalTokens, &summary.CostUSD, &summary.AvgDurationMS,
|
||||
&summary.TotalTokens, &summary.CacheReadTokens, &summary.CacheCreateTokens, &summary.ThinkingTokens,
|
||||
&summary.CostUSD, &summary.AvgDurationMS,
|
||||
); err != nil {
|
||||
return nil, fmt.Errorf("get usage event summary: %w", err)
|
||||
}
|
||||
|
||||
@@ -27,34 +27,37 @@ const (
|
||||
// UsageEvent is an append-only analytics row for resource usage.
|
||||
// It intentionally excludes user IDs, raw prompts, tool args, outputs, and shell commands.
|
||||
type UsageEvent struct {
|
||||
ID uuid.UUID `json:"id" db:"id"`
|
||||
TenantID uuid.UUID `json:"tenant_id" db:"tenant_id"`
|
||||
EventTime time.Time `json:"event_time" db:"event_time"`
|
||||
BucketHour time.Time `json:"bucket_hour" db:"bucket_hour"`
|
||||
EventType string `json:"event_type" db:"event_type"`
|
||||
ResourceType string `json:"resource_type" db:"resource_type"`
|
||||
ResourceName string `json:"resource_name" db:"resource_name"`
|
||||
ResourceID string `json:"resource_id" db:"resource_id"`
|
||||
Source string `json:"source" db:"source"`
|
||||
AgentID *uuid.UUID `json:"agent_id,omitempty" db:"agent_id"`
|
||||
TeamID *uuid.UUID `json:"team_id,omitempty" db:"team_id"`
|
||||
TraceID *uuid.UUID `json:"trace_id,omitempty" db:"trace_id"`
|
||||
SpanID *uuid.UUID `json:"span_id,omitempty" db:"span_id"`
|
||||
RunID string `json:"run_id" db:"run_id"`
|
||||
SessionKey string `json:"session_key" db:"session_key"`
|
||||
Channel string `json:"channel" db:"channel"`
|
||||
Provider string `json:"provider" db:"provider"`
|
||||
Model string `json:"model" db:"model"`
|
||||
Status string `json:"status" db:"status"`
|
||||
InputTokens int64 `json:"input_tokens" db:"input_tokens"`
|
||||
OutputTokens int64 `json:"output_tokens" db:"output_tokens"`
|
||||
TotalTokens int64 `json:"total_tokens" db:"total_tokens"`
|
||||
CostUSD float64 `json:"cost_usd" db:"cost_usd"`
|
||||
DurationMS int `json:"duration_ms" db:"duration_ms"`
|
||||
CallCount int `json:"call_count" db:"call_count"`
|
||||
ErrorCount int `json:"error_count" db:"error_count"`
|
||||
Metadata json.RawMessage `json:"metadata,omitempty" db:"metadata"`
|
||||
CreatedAt time.Time `json:"created_at" db:"created_at"`
|
||||
ID uuid.UUID `json:"id" db:"id"`
|
||||
TenantID uuid.UUID `json:"tenant_id" db:"tenant_id"`
|
||||
EventTime time.Time `json:"event_time" db:"event_time"`
|
||||
BucketHour time.Time `json:"bucket_hour" db:"bucket_hour"`
|
||||
EventType string `json:"event_type" db:"event_type"`
|
||||
ResourceType string `json:"resource_type" db:"resource_type"`
|
||||
ResourceName string `json:"resource_name" db:"resource_name"`
|
||||
ResourceID string `json:"resource_id" db:"resource_id"`
|
||||
Source string `json:"source" db:"source"`
|
||||
AgentID *uuid.UUID `json:"agent_id,omitempty" db:"agent_id"`
|
||||
TeamID *uuid.UUID `json:"team_id,omitempty" db:"team_id"`
|
||||
TraceID *uuid.UUID `json:"trace_id,omitempty" db:"trace_id"`
|
||||
SpanID *uuid.UUID `json:"span_id,omitempty" db:"span_id"`
|
||||
RunID string `json:"run_id" db:"run_id"`
|
||||
SessionKey string `json:"session_key" db:"session_key"`
|
||||
Channel string `json:"channel" db:"channel"`
|
||||
Provider string `json:"provider" db:"provider"`
|
||||
Model string `json:"model" db:"model"`
|
||||
Status string `json:"status" db:"status"`
|
||||
InputTokens int64 `json:"input_tokens" db:"input_tokens"`
|
||||
OutputTokens int64 `json:"output_tokens" db:"output_tokens"`
|
||||
TotalTokens int64 `json:"total_tokens" db:"total_tokens"`
|
||||
CacheReadTokens int64 `json:"cache_read_tokens" db:"cache_read_tokens"`
|
||||
CacheCreateTokens int64 `json:"cache_create_tokens" db:"cache_create_tokens"`
|
||||
ThinkingTokens int64 `json:"thinking_tokens" db:"thinking_tokens"`
|
||||
CostUSD float64 `json:"cost_usd" db:"cost_usd"`
|
||||
DurationMS int `json:"duration_ms" db:"duration_ms"`
|
||||
CallCount int `json:"call_count" db:"call_count"`
|
||||
ErrorCount int `json:"error_count" db:"error_count"`
|
||||
Metadata json.RawMessage `json:"metadata,omitempty" db:"metadata"`
|
||||
CreatedAt time.Time `json:"created_at" db:"created_at"`
|
||||
}
|
||||
|
||||
// UsageEventQuery filters usage event analytics.
|
||||
@@ -75,63 +78,75 @@ type UsageEventQuery struct {
|
||||
}
|
||||
|
||||
type UsageEventSummary struct {
|
||||
Calls int `json:"calls" db:"calls"`
|
||||
Errors int `json:"errors" db:"errors"`
|
||||
InputTokens int64 `json:"input_tokens" db:"input_tokens"`
|
||||
OutputTokens int64 `json:"output_tokens" db:"output_tokens"`
|
||||
TotalTokens int64 `json:"total_tokens" db:"total_tokens"`
|
||||
CostUSD float64 `json:"cost_usd" db:"cost_usd"`
|
||||
AvgDurationMS int `json:"avg_duration_ms" db:"avg_duration_ms"`
|
||||
Calls int `json:"calls" db:"calls"`
|
||||
Errors int `json:"errors" db:"errors"`
|
||||
InputTokens int64 `json:"input_tokens" db:"input_tokens"`
|
||||
OutputTokens int64 `json:"output_tokens" db:"output_tokens"`
|
||||
TotalTokens int64 `json:"total_tokens" db:"total_tokens"`
|
||||
CacheReadTokens int64 `json:"cache_read_tokens" db:"cache_read_tokens"`
|
||||
CacheCreateTokens int64 `json:"cache_create_tokens" db:"cache_create_tokens"`
|
||||
ThinkingTokens int64 `json:"thinking_tokens" db:"thinking_tokens"`
|
||||
CostUSD float64 `json:"cost_usd" db:"cost_usd"`
|
||||
AvgDurationMS int `json:"avg_duration_ms" db:"avg_duration_ms"`
|
||||
}
|
||||
|
||||
type UsageEventTimeSeries struct {
|
||||
BucketTime time.Time `json:"bucket_time" db:"bucket_time"`
|
||||
Calls int `json:"calls" db:"calls"`
|
||||
Errors int `json:"errors" db:"errors"`
|
||||
InputTokens int64 `json:"input_tokens" db:"input_tokens"`
|
||||
OutputTokens int64 `json:"output_tokens" db:"output_tokens"`
|
||||
TotalTokens int64 `json:"total_tokens" db:"total_tokens"`
|
||||
CostUSD float64 `json:"cost_usd" db:"cost_usd"`
|
||||
AvgDurationMS int `json:"avg_duration_ms" db:"avg_duration_ms"`
|
||||
BucketTime time.Time `json:"bucket_time" db:"bucket_time"`
|
||||
Calls int `json:"calls" db:"calls"`
|
||||
Errors int `json:"errors" db:"errors"`
|
||||
InputTokens int64 `json:"input_tokens" db:"input_tokens"`
|
||||
OutputTokens int64 `json:"output_tokens" db:"output_tokens"`
|
||||
TotalTokens int64 `json:"total_tokens" db:"total_tokens"`
|
||||
CacheReadTokens int64 `json:"cache_read_tokens" db:"cache_read_tokens"`
|
||||
CacheCreateTokens int64 `json:"cache_create_tokens" db:"cache_create_tokens"`
|
||||
ThinkingTokens int64 `json:"thinking_tokens" db:"thinking_tokens"`
|
||||
CostUSD float64 `json:"cost_usd" db:"cost_usd"`
|
||||
AvgDurationMS int `json:"avg_duration_ms" db:"avg_duration_ms"`
|
||||
}
|
||||
|
||||
type UsageEventBreakdown struct {
|
||||
Key string `json:"key" db:"key"`
|
||||
EventType string `json:"event_type" db:"event_type"`
|
||||
ResourceType string `json:"resource_type" db:"resource_type"`
|
||||
ResourceName string `json:"resource_name" db:"resource_name"`
|
||||
Source string `json:"source" db:"source"`
|
||||
Calls int `json:"calls" db:"calls"`
|
||||
Errors int `json:"errors" db:"errors"`
|
||||
InputTokens int64 `json:"input_tokens" db:"input_tokens"`
|
||||
OutputTokens int64 `json:"output_tokens" db:"output_tokens"`
|
||||
TotalTokens int64 `json:"total_tokens" db:"total_tokens"`
|
||||
CostUSD float64 `json:"cost_usd" db:"cost_usd"`
|
||||
AvgDurationMS int `json:"avg_duration_ms" db:"avg_duration_ms"`
|
||||
Key string `json:"key" db:"key"`
|
||||
EventType string `json:"event_type" db:"event_type"`
|
||||
ResourceType string `json:"resource_type" db:"resource_type"`
|
||||
ResourceName string `json:"resource_name" db:"resource_name"`
|
||||
Source string `json:"source" db:"source"`
|
||||
Calls int `json:"calls" db:"calls"`
|
||||
Errors int `json:"errors" db:"errors"`
|
||||
InputTokens int64 `json:"input_tokens" db:"input_tokens"`
|
||||
OutputTokens int64 `json:"output_tokens" db:"output_tokens"`
|
||||
TotalTokens int64 `json:"total_tokens" db:"total_tokens"`
|
||||
CacheReadTokens int64 `json:"cache_read_tokens" db:"cache_read_tokens"`
|
||||
CacheCreateTokens int64 `json:"cache_create_tokens" db:"cache_create_tokens"`
|
||||
ThinkingTokens int64 `json:"thinking_tokens" db:"thinking_tokens"`
|
||||
CostUSD float64 `json:"cost_usd" db:"cost_usd"`
|
||||
AvgDurationMS int `json:"avg_duration_ms" db:"avg_duration_ms"`
|
||||
}
|
||||
|
||||
type UsageEventRollup struct {
|
||||
ID uuid.UUID `json:"id" db:"id"`
|
||||
TenantID uuid.UUID `json:"tenant_id" db:"tenant_id"`
|
||||
BucketHour time.Time `json:"bucket_hour" db:"bucket_hour"`
|
||||
EventType string `json:"event_type" db:"event_type"`
|
||||
ResourceType string `json:"resource_type" db:"resource_type"`
|
||||
ResourceName string `json:"resource_name" db:"resource_name"`
|
||||
Source string `json:"source" db:"source"`
|
||||
AgentID *uuid.UUID `json:"agent_id,omitempty" db:"agent_id"`
|
||||
Channel string `json:"channel" db:"channel"`
|
||||
Provider string `json:"provider" db:"provider"`
|
||||
Model string `json:"model" db:"model"`
|
||||
Status string `json:"status" db:"status"`
|
||||
InputTokens int64 `json:"input_tokens" db:"input_tokens"`
|
||||
OutputTokens int64 `json:"output_tokens" db:"output_tokens"`
|
||||
TotalTokens int64 `json:"total_tokens" db:"total_tokens"`
|
||||
CostUSD float64 `json:"cost_usd" db:"cost_usd"`
|
||||
DurationMS int `json:"duration_ms" db:"duration_ms"`
|
||||
CallCount int `json:"call_count" db:"call_count"`
|
||||
ErrorCount int `json:"error_count" db:"error_count"`
|
||||
CreatedAt time.Time `json:"created_at" db:"created_at"`
|
||||
UpdatedAt time.Time `json:"updated_at" db:"updated_at"`
|
||||
ID uuid.UUID `json:"id" db:"id"`
|
||||
TenantID uuid.UUID `json:"tenant_id" db:"tenant_id"`
|
||||
BucketHour time.Time `json:"bucket_hour" db:"bucket_hour"`
|
||||
EventType string `json:"event_type" db:"event_type"`
|
||||
ResourceType string `json:"resource_type" db:"resource_type"`
|
||||
ResourceName string `json:"resource_name" db:"resource_name"`
|
||||
Source string `json:"source" db:"source"`
|
||||
AgentID *uuid.UUID `json:"agent_id,omitempty" db:"agent_id"`
|
||||
Channel string `json:"channel" db:"channel"`
|
||||
Provider string `json:"provider" db:"provider"`
|
||||
Model string `json:"model" db:"model"`
|
||||
Status string `json:"status" db:"status"`
|
||||
InputTokens int64 `json:"input_tokens" db:"input_tokens"`
|
||||
OutputTokens int64 `json:"output_tokens" db:"output_tokens"`
|
||||
TotalTokens int64 `json:"total_tokens" db:"total_tokens"`
|
||||
CacheReadTokens int64 `json:"cache_read_tokens" db:"cache_read_tokens"`
|
||||
CacheCreateTokens int64 `json:"cache_create_tokens" db:"cache_create_tokens"`
|
||||
ThinkingTokens int64 `json:"thinking_tokens" db:"thinking_tokens"`
|
||||
CostUSD float64 `json:"cost_usd" db:"cost_usd"`
|
||||
DurationMS int `json:"duration_ms" db:"duration_ms"`
|
||||
CallCount int `json:"call_count" db:"call_count"`
|
||||
ErrorCount int `json:"error_count" db:"error_count"`
|
||||
CreatedAt time.Time `json:"created_at" db:"created_at"`
|
||||
UpdatedAt time.Time `json:"updated_at" db:"updated_at"`
|
||||
}
|
||||
|
||||
type UsageEventStore interface {
|
||||
|
||||
@@ -0,0 +1,7 @@
|
||||
ALTER TABLE usage_event_rollups DROP COLUMN IF EXISTS thinking_tokens;
|
||||
ALTER TABLE usage_event_rollups DROP COLUMN IF EXISTS cache_create_tokens;
|
||||
ALTER TABLE usage_event_rollups DROP COLUMN IF EXISTS cache_read_tokens;
|
||||
|
||||
ALTER TABLE usage_events DROP COLUMN IF EXISTS thinking_tokens;
|
||||
ALTER TABLE usage_events DROP COLUMN IF EXISTS cache_create_tokens;
|
||||
ALTER TABLE usage_events DROP COLUMN IF EXISTS cache_read_tokens;
|
||||
@@ -0,0 +1,9 @@
|
||||
ALTER TABLE usage_events
|
||||
ADD COLUMN IF NOT EXISTS cache_read_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
ADD COLUMN IF NOT EXISTS cache_create_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
ADD COLUMN IF NOT EXISTS thinking_tokens BIGINT NOT NULL DEFAULT 0;
|
||||
|
||||
ALTER TABLE usage_event_rollups
|
||||
ADD COLUMN IF NOT EXISTS cache_read_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
ADD COLUMN IF NOT EXISTS cache_create_tokens BIGINT NOT NULL DEFAULT 0,
|
||||
ADD COLUMN IF NOT EXISTS thinking_tokens BIGINT NOT NULL DEFAULT 0;
|
||||
@@ -10,6 +10,9 @@ export interface UsageEventSummary {
|
||||
input_tokens: number;
|
||||
output_tokens: number;
|
||||
total_tokens: number;
|
||||
cache_read_tokens: number;
|
||||
cache_create_tokens: number;
|
||||
thinking_tokens: number;
|
||||
cost_usd: number;
|
||||
avg_duration_ms: number;
|
||||
}
|
||||
@@ -25,6 +28,9 @@ export interface UsageEventBreakdown {
|
||||
input_tokens: number;
|
||||
output_tokens: number;
|
||||
total_tokens: number;
|
||||
cache_read_tokens: number;
|
||||
cache_create_tokens: number;
|
||||
thinking_tokens: number;
|
||||
cost_usd: number;
|
||||
avg_duration_ms: number;
|
||||
}
|
||||
@@ -36,6 +42,9 @@ export interface UsageEventTimeSeries {
|
||||
input_tokens: number;
|
||||
output_tokens: number;
|
||||
total_tokens: number;
|
||||
cache_read_tokens: number;
|
||||
cache_create_tokens: number;
|
||||
thinking_tokens: number;
|
||||
cost_usd: number;
|
||||
avg_duration_ms: number;
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user