diff --git a/cmd/gateway_http_wiring.go b/cmd/gateway_http_wiring.go index 3479472e..35eb9dfd 100644 --- a/cmd/gateway_http_wiring.go +++ b/cmd/gateway_http_wiring.go @@ -135,7 +135,7 @@ func (d *gatewayDeps) wireHTTPHandlersOnServer( // Usage analytics API if d.pgStores.Snapshots != nil { - d.server.SetUsageHandler(httpapi.NewUsageHandler(d.pgStores.Snapshots, d.pgStores.DB)) + d.server.SetUsageHandler(httpapi.NewUsageHandler(d.pgStores.Snapshots, d.pgStores.UsageEvents, d.pgStores.DB)) } if d.pgStores.UsageCaps != nil { d.server.SetUsageCapsHandler(httpapi.NewUsageCapsHandler(d.pgStores.UsageCaps, d.pgStores.Tenants)) diff --git a/cmd/gateway_managed.go b/cmd/gateway_managed.go index 179ce8da..e1c6b5f6 100644 --- a/cmd/gateway_managed.go +++ b/cmd/gateway_managed.go @@ -246,6 +246,7 @@ func wireExtras( ModelPricing: appCfg.Telemetry.ModelPricing, TracingStore: stores.Tracing, UsageCaps: usageCapSvc, + UsageEvents: stores.UsageEvents, MemoryStore: stores.Memory, ContactStore: stores.Contacts, TenantStore: stores.Tenants, diff --git a/cmd/gateway_setup.go b/cmd/gateway_setup.go index aba394db..3f6ac2a4 100644 --- a/cmd/gateway_setup.go +++ b/cmd/gateway_setup.go @@ -322,7 +322,7 @@ func wireTracingAndCron( // Start snapshot worker for hourly usage aggregation var snapshotWorker *tracing.SnapshotWorker if stores.Snapshots != nil { - snapshotWorker = tracing.NewSnapshotWorker(stores.DB, stores.Snapshots) + snapshotWorker = tracing.NewSnapshotWorker(stores.DB, stores.Snapshots, stores.UsageEvents) snapshotWorker.Start() // Backfill historical data in background diff --git a/internal/agent/loop_context.go b/internal/agent/loop_context.go index 80d166cd..e8804344 100644 --- a/internal/agent/loop_context.go +++ b/internal/agent/loop_context.go @@ -378,6 +378,8 @@ func (l *Loop) injectContext(ctx context.Context, req *RunRequest) (contextSetup AgentKey: l.id, TenantID: l.tenantID, UserID: req.UserID, + RunID: req.RunID, + SessionKey: req.SessionKey, CredentialUserID: credUserID, AgentType: l.agentType, SenderID: req.SenderID, @@ -388,6 +390,7 @@ func (l *Loop) injectContext(ctx context.Context, req *RunRequest) (contextSetup SharedContext: store.IsSharedContext(ctx), RestrictToWorkspace: l.restrictToWs != nil && *l.restrictToWs, BuiltinToolSettings: l.builtinToolSettings, + Channel: req.Channel, ChannelType: req.ChannelType, SubagentsCfg: l.subagentsCfg, ParentModel: l.model, diff --git a/internal/agent/loop_pipeline_tool_callbacks.go b/internal/agent/loop_pipeline_tool_callbacks.go index 90373ab1..7a74155a 100644 --- a/internal/agent/loop_pipeline_tool_callbacks.go +++ b/internal/agent/loop_pipeline_tool_callbacks.go @@ -21,7 +21,7 @@ func (l *Loop) makeExecuteToolCall(req *RunRequest, bridgeRS *runState) func(ctx emitRun := makeToolEmitRun(l, req) return func(ctx context.Context, state *pipeline.RunState, tc providers.ToolCall) ([]providers.Message, error) { tc = l.normalizeToolCall(tc) - registryName := l.resolveToolCallName(tc.Name) + registryName := l.canonicalToolName(l.resolveToolCallName(tc.Name)) argsJSON, _ := json.Marshal(tc.Arguments) slog.Info("tool call", "agent", l.id, "tool", tc.Name, "args_len", len(argsJSON)) @@ -34,7 +34,7 @@ func (l *Loop) makeExecuteToolCall(req *RunRequest, bridgeRS *runState) func(ctx // Emit tool span start for tracing. toolStart := time.Now().UTC() - toolSpanID := l.emitToolSpanStart(ctx, toolStart, tc.Name, tc.ID, string(argsJSON)) + toolSpanID := l.emitToolSpanStart(ctx, toolStart, registryName, tc.ID, string(argsJSON)) // Inject agent audio snapshot so TTS tool (and any future audio consumers) // can read agent-level voice/model config without an extra DB lookup. @@ -55,6 +55,7 @@ func (l *Loop) makeExecuteToolCall(req *RunRequest, bridgeRS *runState) func(ctx toolDuration := time.Since(toolStart) l.emitToolSpanEnd(ctx, toolSpanID, toolStart, result) + l.recordToolUsageEvent(ctx, req, registryName, tc.Name, tc.ID, tc.Arguments, toolStart, result, toolSpanID) // v3 evolution metrics: record tool execution non-blocking (best-effort). l.recordToolMetric(ctx, req.SessionKey, registryName, !result.IsError, toolDuration) @@ -73,6 +74,10 @@ func (l *Loop) makeExecuteToolCall(req *RunRequest, bridgeRS *runState) func(ctx type toolRawResult struct { result *tools.Result duration time.Duration + start time.Time + spanID uuid.UUID + toolName string + rawName string } // makeExecuteToolRaw wraps tool I/O only (parallel-safe, no state mutation). @@ -81,7 +86,7 @@ func (l *Loop) makeExecuteToolRaw(req *RunRequest) func(ctx context.Context, tc emitRun := makeToolEmitRun(l, req) return func(ctx context.Context, tc providers.ToolCall) (providers.Message, any, error) { tc = l.normalizeToolCall(tc) - registryName := l.resolveToolCallName(tc.Name) + registryName := l.canonicalToolName(l.resolveToolCallName(tc.Name)) argsJSON, _ := json.Marshal(tc.Arguments) slog.Info("tool call", "agent", l.id, "tool", tc.Name, "args_len", len(argsJSON)) @@ -98,7 +103,7 @@ func (l *Loop) makeExecuteToolRaw(req *RunRequest) func(ctx context.Context, tc // Emit tool span start (goroutine-safe: channel send only). start := time.Now().UTC() - spanID := l.emitToolSpanStart(ctx, start, tc.Name, tc.ID, string(argsJSON)) + spanID := l.emitToolSpanStart(ctx, start, registryName, tc.ID, string(argsJSON)) // Inject agent audio snapshot (parallel path — same as sequential makeExecuteToolCall). if l.agentUUID != uuid.Nil { @@ -125,7 +130,7 @@ func (l *Loop) makeExecuteToolRaw(req *RunRequest) func(ctx context.Context, tc ToolCallID: tc.ID, IsError: result.IsError, } - return msg, &toolRawResult{result: result, duration: dur}, nil + return msg, &toolRawResult{result: result, duration: dur, start: start, spanID: spanID, toolName: registryName, rawName: tc.Name}, nil } } @@ -151,20 +156,33 @@ func (l *Loop) makeProcessToolResult(req *RunRequest, bridgeRS *runState) func(c emitRun := makeToolEmitRun(l, req) return func(ctx context.Context, state *pipeline.RunState, tc providers.ToolCall, rawMsg providers.Message, rawData any) []providers.Message { tc = l.normalizeToolCall(tc) - registryName := l.resolveToolCallName(tc.Name) + registryName := l.canonicalToolName(l.resolveToolCallName(tc.Name)) // Extract result and timing from toolRawResult wrapper. var result *tools.Result var dur time.Duration + var start time.Time + var spanID uuid.UUID + var rawName string if raw, ok := rawData.(*toolRawResult); ok && raw != nil { result = raw.result dur = raw.duration + start = raw.start + spanID = raw.spanID + registryName = raw.toolName + rawName = raw.rawName } else if r, ok := rawData.(*tools.Result); ok { result = r // backward compat } if result == nil { return []providers.Message{rawMsg} } + if rawName == "" { + rawName = tc.Name + } + if !start.IsZero() { + l.recordToolUsageEvent(ctx, req, registryName, rawName, tc.ID, tc.Arguments, start, result, spanID) + } // Record tool metrics (non-blocking, best-effort). l.recordToolMetric(ctx, req.SessionKey, registryName, !result.IsError, dur) diff --git a/internal/agent/loop_types.go b/internal/agent/loop_types.go index 816e06a1..d98cd06e 100644 --- a/internal/agent/loop_types.go +++ b/internal/agent/loop_types.go @@ -243,6 +243,7 @@ type Loop struct { budgetMonthlyCents int tracingStore store.TracingStore usageCaps *usagecaps.Service + usageEvents store.UsageEventStore // Memory store for extractive memory fallback (writes directly when LLM flush fails) memStore store.MemoryStore @@ -437,6 +438,7 @@ type LoopConfig struct { BudgetMonthlyCents int TracingStore store.TracingStore UsageCaps *usagecaps.Service + UsageEvents store.UsageEventStore // Memory store for extractive memory fallback (writes directly when LLM flush fails) MemoryStore store.MemoryStore @@ -578,6 +580,7 @@ func NewLoop(cfg LoopConfig) *Loop { budgetMonthlyCents: cfg.BudgetMonthlyCents, tracingStore: cfg.TracingStore, usageCaps: cfg.UsageCaps, + usageEvents: cfg.UsageEvents, memStore: cfg.MemoryStore, mcpStore: cfg.MCPStore, mcpPool: cfg.MCPPool, diff --git a/internal/agent/resolver.go b/internal/agent/resolver.go index 823e17fc..e92e4ff1 100644 --- a/internal/agent/resolver.go +++ b/internal/agent/resolver.go @@ -100,6 +100,7 @@ type ResolverDeps struct { // Tracing store for budget enforcement queries TracingStore store.TracingStore UsageCaps *usagecaps.Service + UsageEvents store.UsageEventStore // Memory store for extractive memory fallback MemoryStore store.MemoryStore @@ -534,6 +535,7 @@ func NewManagedResolver(deps ResolverDeps) ResolverFunc { BudgetMonthlyCents: derefInt(ag.BudgetMonthlyCents), TracingStore: deps.TracingStore, UsageCaps: deps.UsageCaps, + UsageEvents: deps.UsageEvents, MemoryStore: deps.MemoryStore, MCPStore: deps.MCPStore, MCPPool: deps.MCPPool, diff --git a/internal/agent/skill_slash_commands.go b/internal/agent/skill_slash_commands.go index c9b1f34e..66be2c73 100644 --- a/internal/agent/skill_slash_commands.go +++ b/internal/agent/skill_slash_commands.go @@ -41,6 +41,7 @@ func (l *Loop) applySkillSlashCommand(ctx context.Context, message, extraPrompt message = result.RemainingPrompt } skillFilter = []string{result.Skill.Slug} + l.recordSkillSlashUsageEvent(ctx, result.Skill.Slug) case skillSlashCommandList: message = "List the available skills shown in the system instructions." case skillSlashCommandHelp: diff --git a/internal/agent/usage_events.go b/internal/agent/usage_events.go new file mode 100644 index 00000000..854f389b --- /dev/null +++ b/internal/agent/usage_events.go @@ -0,0 +1,200 @@ +package agent + +import ( + "context" + "encoding/json" + "log/slog" + "slices" + "time" + + "github.com/google/uuid" + + "github.com/nextlevelbuilder/goclaw/internal/store" + "github.com/nextlevelbuilder/goclaw/internal/tools" + "github.com/nextlevelbuilder/goclaw/internal/tracing" +) + +type mcpUsageTool interface { + ServerName() string + OriginalName() string +} + +func (l *Loop) canonicalToolName(name string) string { + if l.registry == nil { + return name + } + if tool, ok := l.registry.Get(name); ok && tool != nil { + return tool.Name() + } + return name +} + +func (l *Loop) recordToolUsageEvent(ctx context.Context, req *RunRequest, canonicalName, rawName, toolCallID string, args map[string]any, start time.Time, result *tools.Result, spanID uuid.UUID) { + if l.usageEvents == nil || result == nil { + return + } + if tracing.TraceIDFromContext(ctx) == uuid.Nil { + return + } + + resourceName := canonicalName + resourceID := canonicalName + eventType := store.UsageEventTypeToolCall + resourceType := store.UsageResourceTypeTool + source := store.UsageSourceToolCall + metadata := map[string]any{} + + if canonicalName == "use_skill" { + eventType = store.UsageEventTypeSkillActivation + resourceType = store.UsageResourceTypeSkill + source = store.UsageSourceUseSkill + if skill, _ := args["name"].(string); skill != "" { + resourceName = skill + resourceID = skill + } + } else if l.registry != nil { + tool, ok := l.registry.Get(canonicalName) + if !ok { + tool = nil + } + if mcpTool, ok := tool.(mcpUsageTool); ok { + eventType = store.UsageEventTypeMCPToolCall + resourceType = store.UsageResourceTypeMCPTool + resourceName = mcpTool.ServerName() + "/" + mcpTool.OriginalName() + resourceID = mcpTool.OriginalName() + metadata["server"] = mcpTool.ServerName() + metadata["tool"] = mcpTool.OriginalName() + } else if l.isRuntimeTool(canonicalName) { + eventType = store.UsageEventTypeRuntimeToolCall + resourceType = store.UsageResourceTypeRuntimeTool + } + } + + if rawName != "" && rawName != canonicalName { + metadata["raw_tool_name"] = rawName + } + + event := l.baseUsageEvent(ctx, req, start, eventType, resourceType, resourceName, resourceID, source) + event.SpanID = uuidPtr(spanID) + event.Status = "completed" + if result.IsError { + event.Status = "error" + event.ErrorCount = 1 + } + event.DurationMS = int(time.Since(start).Milliseconds()) + if result.Usage != nil { + event.InputTokens = int64(result.Usage.PromptTokens) + event.OutputTokens = int64(result.Usage.CompletionTokens) + event.TotalTokens = int64(result.Usage.TotalTokens) + if event.TotalTokens == 0 { + event.TotalTokens = event.InputTokens + event.OutputTokens + } + } + event.Provider = result.Provider + event.Model = result.Model + event.Metadata = usageMetadata(metadata) + l.insertUsageEventBestEffort(ctx, event) +} + +func (l *Loop) recordSkillSlashUsageEvent(ctx context.Context, skillSlug string) { + if l.usageEvents == nil || skillSlug == "" { + return + } + traceID := tracing.TraceIDFromContext(ctx) + if traceID == uuid.Nil { + return + } + rc := store.RunContextFromCtx(ctx) + event := l.baseUsageEvent(ctx, nil, time.Now().UTC(), + store.UsageEventTypeSkillActivation, + store.UsageResourceTypeSkill, + skillSlug, + skillSlug, + store.UsageSourceSlashCommand, + ) + event.TraceID = uuidPtr(traceID) + if rc != nil { + event.RunID = rc.RunID + event.SessionKey = rc.SessionKey + event.Channel = rc.Channel + if rc.TeamID != "" { + if teamID, err := uuid.Parse(rc.TeamID); err == nil { + event.TeamID = &teamID + } + } + } + event.Metadata = usageMetadata(map[string]any{"activation_source": store.UsageSourceSlashCommand}) + l.insertUsageEventBestEffort(ctx, event) +} + +func (l *Loop) baseUsageEvent(ctx context.Context, req *RunRequest, eventTime time.Time, eventType, resourceType, resourceName, resourceID, source string) store.UsageEvent { + tenantID := store.TenantIDFromContext(ctx) + if tenantID == uuid.Nil { + tenantID = l.tenantID + } + traceID := tracing.TraceIDFromContext(ctx) + event := store.UsageEvent{ + ID: uuid.New(), + TenantID: tenantID, + EventTime: eventTime.UTC(), + BucketHour: eventTime.UTC().Truncate(time.Hour), + EventType: eventType, + ResourceType: resourceType, + ResourceName: resourceName, + ResourceID: resourceID, + Source: source, + AgentID: uuidPtr(l.agentUUID), + TeamID: tracing.TraceTeamIDPtrFromContext(ctx), + TraceID: uuidPtr(traceID), + Status: "completed", + CallCount: 1, + } + if req != nil { + event.RunID = req.RunID + event.SessionKey = req.SessionKey + event.Channel = req.Channel + if req.TeamID != "" && event.TeamID == nil { + if teamID, err := uuid.Parse(req.TeamID); err == nil { + event.TeamID = &teamID + } + } + } + return event +} + +func (l *Loop) isRuntimeTool(toolName string) bool { + if l.registry == nil { + return false + } + members, ok := l.registry.GetToolGroup("runtime") + return ok && slices.Contains(members, toolName) +} + +func (l *Loop) insertUsageEventBestEffort(ctx context.Context, event store.UsageEvent) { + tenantID := event.TenantID + go func() { + bgCtx, cancel := context.WithTimeout(store.WithTenantID(context.Background(), tenantID), 5*time.Second) + defer cancel() + if err := l.usageEvents.InsertEvent(bgCtx, &event); err != nil { + slog.Debug("usage.event.record_failed", "resource", event.ResourceName, "error", err) + } + }() +} + +func uuidPtr(id uuid.UUID) *uuid.UUID { + if id == uuid.Nil { + return nil + } + return &id +} + +func usageMetadata(values map[string]any) json.RawMessage { + if len(values) == 0 { + return nil + } + data, err := json.Marshal(values) + if err != nil { + return nil + } + return data +} diff --git a/internal/agent/usage_events_test.go b/internal/agent/usage_events_test.go new file mode 100644 index 00000000..59cff8da --- /dev/null +++ b/internal/agent/usage_events_test.go @@ -0,0 +1,164 @@ +package agent + +import ( + "context" + "encoding/json" + "testing" + "time" + + "github.com/google/uuid" + + "github.com/nextlevelbuilder/goclaw/internal/store" + "github.com/nextlevelbuilder/goclaw/internal/tools" + "github.com/nextlevelbuilder/goclaw/internal/tracing" +) + +type fakeUsageEventStore struct { + events chan store.UsageEvent +} + +func newFakeUsageEventStore() *fakeUsageEventStore { + return &fakeUsageEventStore{events: make(chan store.UsageEvent, 8)} +} + +func (s *fakeUsageEventStore) InsertEvent(_ context.Context, event *store.UsageEvent) error { + if event != nil { + s.events <- *event + } + return nil +} + +func (s *fakeUsageEventStore) InsertEvents(ctx context.Context, events []store.UsageEvent) error { + for i := range events { + if err := s.InsertEvent(ctx, &events[i]); err != nil { + return err + } + } + return nil +} + +func (s *fakeUsageEventStore) GetEventTimeSeries(context.Context, store.UsageEventQuery) ([]store.UsageEventTimeSeries, error) { + return nil, nil +} + +func (s *fakeUsageEventStore) RefreshEventRollupHour(context.Context, time.Time) error { + return nil +} + +func (s *fakeUsageEventStore) GetLatestEventRollupBucket(context.Context) (*time.Time, error) { + return nil, nil +} + +func (s *fakeUsageEventStore) GetEventBreakdown(context.Context, store.UsageEventQuery) ([]store.UsageEventBreakdown, error) { + return nil, nil +} + +func (s *fakeUsageEventStore) GetEventSummary(context.Context, store.UsageEventQuery) (*store.UsageEventSummary, error) { + return nil, nil +} + +type usageFakeTool struct { + name string +} + +func (t usageFakeTool) Name() string { return t.name } +func (t usageFakeTool) Description() string { return "" } +func (t usageFakeTool) Parameters() map[string]any { return nil } +func (t usageFakeTool) Execute(context.Context, map[string]any) *tools.Result { + return tools.NewResult("ok") +} + +func TestRecordToolUsageEvent_RuntimeAliasUsesCanonicalName(t *testing.T) { + storeSpy := newFakeUsageEventStore() + registry := tools.NewRegistry() + registry.Register(usageFakeTool{name: "exec"}) + registry.RegisterAlias("Bash", "exec") + loop := &Loop{registry: registry, usageEvents: storeSpy, agentUUID: uuid.New(), tenantID: uuid.New()} + ctx := tracing.WithTraceID(store.WithTenantID(t.Context(), loop.tenantID), uuid.New()) + + canonical := loop.canonicalToolName("Bash") + loop.recordToolUsageEvent(ctx, &RunRequest{RunID: "run-1", SessionKey: "session-1", Channel: "telegram"}, canonical, "Bash", "call-1", + map[string]any{"command": "echo secret"}, time.Now().Add(-time.Millisecond), tools.NewResult("ok"), uuid.New()) + + event := waitUsageEvent(t, storeSpy) + if event.EventType != store.UsageEventTypeRuntimeToolCall { + t.Fatalf("event type = %q, want %q", event.EventType, store.UsageEventTypeRuntimeToolCall) + } + if event.ResourceName != "exec" { + t.Fatalf("resource = %q, want canonical exec", event.ResourceName) + } + if string(event.Metadata) == "" || json.Valid(event.Metadata) == false { + t.Fatalf("metadata should contain valid alias metadata, got %q", string(event.Metadata)) + } + if containsJSONKey(event.Metadata, "command") { + t.Fatalf("metadata leaked raw command args: %s", string(event.Metadata)) + } +} + +func TestRecordToolUsageEvent_UseSkillCountsSkillName(t *testing.T) { + storeSpy := newFakeUsageEventStore() + registry := tools.NewRegistry() + registry.Register(tools.NewUseSkillTool()) + loop := &Loop{registry: registry, usageEvents: storeSpy, agentUUID: uuid.New(), tenantID: uuid.New()} + ctx := tracing.WithTraceID(store.WithTenantID(t.Context(), loop.tenantID), uuid.New()) + + loop.recordToolUsageEvent(ctx, &RunRequest{RunID: "run-1", SessionKey: "session-1"}, "use_skill", "use_skill", "call-1", + map[string]any{"name": "ck:plan"}, time.Now().Add(-time.Millisecond), tools.NewResult("ok"), uuid.New()) + + event := waitUsageEvent(t, storeSpy) + if event.EventType != store.UsageEventTypeSkillActivation { + t.Fatalf("event type = %q, want skill activation", event.EventType) + } + if event.ResourceType != store.UsageResourceTypeSkill || event.ResourceName != "ck:plan" { + t.Fatalf("resource = %s/%s, want skill ck:plan", event.ResourceType, event.ResourceName) + } + if event.Source != store.UsageSourceUseSkill { + t.Fatalf("source = %q, want use_skill", event.Source) + } +} + +func TestRecordSkillSlashUsageEvent_RequiresTraceContext(t *testing.T) { + storeSpy := newFakeUsageEventStore() + loop := &Loop{usageEvents: storeSpy, agentUUID: uuid.New(), tenantID: uuid.New()} + ctx := store.WithRunContext(store.WithTenantID(t.Context(), loop.tenantID), &store.RunContext{ + RunID: "run-1", + SessionKey: "session-1", + Channel: "web", + }) + + loop.recordSkillSlashUsageEvent(ctx, "ck:plan") + select { + case event := <-storeSpy.events: + t.Fatalf("unexpected event without trace: %+v", event) + case <-time.After(50 * time.Millisecond): + } + + loop.recordSkillSlashUsageEvent(tracing.WithTraceID(ctx, uuid.New()), "ck:plan") + event := waitUsageEvent(t, storeSpy) + if event.Source != store.UsageSourceSlashCommand || event.ResourceName != "ck:plan" { + t.Fatalf("slash event = %s/%s, want slash-command ck:plan", event.Source, event.ResourceName) + } + if event.RunID != "run-1" || event.SessionKey != "session-1" || event.Channel != "web" { + t.Fatalf("run context not copied: run=%q session=%q channel=%q", event.RunID, event.SessionKey, event.Channel) + } +} + +func waitUsageEvent(t *testing.T, storeSpy *fakeUsageEventStore) store.UsageEvent { + t.Helper() + select { + case event := <-storeSpy.events: + return event + case <-time.After(time.Second): + t.Fatal("timed out waiting for usage event") + return store.UsageEvent{} + } +} + +func containsJSONKey(raw json.RawMessage, key string) bool { + var values map[string]any + if err := json.Unmarshal(raw, &values); err != nil { + return false + } + _, ok := values[key] + return ok +} diff --git a/internal/http/usage.go b/internal/http/usage.go index d07108b6..5b29e275 100644 --- a/internal/http/usage.go +++ b/internal/http/usage.go @@ -5,6 +5,7 @@ import ( "fmt" "log/slog" "net/http" + "strconv" "time" "github.com/google/uuid" @@ -14,18 +15,22 @@ import ( // UsageHandler serves pre-computed usage analytics from snapshots. type UsageHandler struct { - snapshots store.SnapshotStore - db *sql.DB + snapshots store.SnapshotStore + usageEvents store.UsageEventStore + db *sql.DB } -func NewUsageHandler(snapshots store.SnapshotStore, db *sql.DB) *UsageHandler { - return &UsageHandler{snapshots: snapshots, db: db} +func NewUsageHandler(snapshots store.SnapshotStore, usageEvents store.UsageEventStore, db *sql.DB) *UsageHandler { + return &UsageHandler{snapshots: snapshots, usageEvents: usageEvents, db: db} } func (h *UsageHandler) RegisterRoutes(mux *http.ServeMux) { mux.HandleFunc("GET /v1/usage/timeseries", h.authMiddleware(h.handleTimeSeries)) mux.HandleFunc("GET /v1/usage/breakdown", h.authMiddleware(h.handleBreakdown)) mux.HandleFunc("GET /v1/usage/summary", h.authMiddleware(h.handleSummary)) + mux.HandleFunc("GET /v1/usage/events/timeseries", h.authMiddleware(h.handleEventTimeSeries)) + mux.HandleFunc("GET /v1/usage/events/breakdown", h.authMiddleware(h.handleEventBreakdown)) + mux.HandleFunc("GET /v1/usage/events/summary", h.authMiddleware(h.handleEventSummary)) } func (h *UsageHandler) authMiddleware(next http.HandlerFunc) http.HandlerFunc { @@ -128,6 +133,81 @@ func (h *UsageHandler) handleSummary(w http.ResponseWriter, r *http.Request) { }) } +func (h *UsageHandler) handleEventTimeSeries(w http.ResponseWriter, r *http.Request) { + if h.usageEvents == nil { + writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "usage event analytics unavailable"}) + return + } + q, err := parseUsageEventFilters(r) + if err != nil { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()}) + return + } + if q.From.IsZero() || q.To.IsZero() { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": "from and to are required"}) + return + } + if q.GroupBy == "" { + q.GroupBy = "hour" + } + points, err := h.usageEvents.GetEventTimeSeries(r.Context(), q) + if err != nil { + slog.Error("usage.events.timeseries query failed", "error", err) + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "internal server error"}) + return + } + writeJSON(w, http.StatusOK, map[string]any{"points": points}) +} + +func (h *UsageHandler) handleEventBreakdown(w http.ResponseWriter, r *http.Request) { + if h.usageEvents == nil { + writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "usage event analytics unavailable"}) + return + } + q, err := parseUsageEventFilters(r) + if err != nil { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()}) + return + } + if q.From.IsZero() || q.To.IsZero() { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": "from and to are required"}) + return + } + if q.GroupBy == "" { + q.GroupBy = "resource" + } + rows, err := h.usageEvents.GetEventBreakdown(r.Context(), q) + if err != nil { + slog.Error("usage.events.breakdown query failed", "error", err) + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "internal server error"}) + return + } + writeJSON(w, http.StatusOK, map[string]any{"rows": rows}) +} + +func (h *UsageHandler) handleEventSummary(w http.ResponseWriter, r *http.Request) { + if h.usageEvents == nil { + writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "usage event analytics unavailable"}) + return + } + q, err := parseUsageEventFilters(r) + if err != nil { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()}) + return + } + if q.From.IsZero() || q.To.IsZero() { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": "from and to are required"}) + return + } + summary, err := h.usageEvents.GetEventSummary(r.Context(), q) + if err != nil { + slog.Error("usage.events.summary query failed", "error", err) + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "internal server error"}) + return + } + writeJSON(w, http.StatusOK, map[string]any{"summary": summary}) +} + // usageSummary is the response shape for summary endpoint. type usageSummary struct { Requests int `json:"requests"` @@ -246,3 +326,51 @@ func parseSnapshotFilters(r *http.Request) store.SnapshotQuery { return q } +func parseUsageEventFilters(r *http.Request) (store.UsageEventQuery, error) { + q := store.UsageEventQuery{} + values := r.URL.Query() + if values.Get("user_id") != "" { + return q, fmt.Errorf("user_id filter is not supported") + } + if v := values.Get("from"); v != "" { + t, err := time.Parse(time.RFC3339, v) + if err != nil { + return q, fmt.Errorf("invalid from") + } + q.From = t + } + if v := values.Get("to"); v != "" { + t, err := time.Parse(time.RFC3339, v) + if err != nil { + return q, fmt.Errorf("invalid to") + } + q.To = t + } + if !q.From.IsZero() && !q.To.IsZero() && !q.To.After(q.From) { + return q, fmt.Errorf("to must be after from") + } + if v := values.Get("agent_id"); v != "" { + id, err := uuid.Parse(v) + if err != nil { + return q, fmt.Errorf("invalid agent_id") + } + q.AgentID = &id + } + q.Channel = values.Get("channel") + q.EventType = values.Get("event_type") + q.ResourceType = values.Get("resource_type") + q.ResourceName = values.Get("resource_name") + q.Provider = values.Get("provider") + q.Model = values.Get("model") + q.Status = values.Get("status") + q.Source = values.Get("source") + q.GroupBy = values.Get("group_by") + if v := values.Get("limit"); v != "" { + limit, err := strconv.Atoi(v) + if err != nil || limit < 0 { + return q, fmt.Errorf("invalid limit") + } + q.Limit = limit + } + return q, nil +} diff --git a/internal/store/pg/factory.go b/internal/store/pg/factory.go index ea2bc565..33cff7f5 100644 --- a/internal/store/pg/factory.go +++ b/internal/store/pg/factory.go @@ -46,6 +46,7 @@ func NewPGStores(cfg store.StoreConfig) (*store.Stores, error) { Contacts: NewPGContactStore(db), Activity: NewPGActivityStore(db), Snapshots: NewPGSnapshotStore(db), + UsageEvents: NewPGUsageEventStore(db), BrowserCookies: NewPGBrowserCookieStore(db, cfg.EncryptionKey), SecureCLI: NewPGSecureCLIStore(db, cfg.EncryptionKey), SecureCLIGrants: NewPGSecureCLIAgentGrantStore(db, cfg.EncryptionKey), diff --git a/internal/store/pg/usage_event.go b/internal/store/pg/usage_event.go new file mode 100644 index 00000000..bde6e3b2 --- /dev/null +++ b/internal/store/pg/usage_event.go @@ -0,0 +1,421 @@ +package pg + +import ( + "context" + "database/sql" + "fmt" + "strings" + "time" + + "github.com/google/uuid" + + "github.com/nextlevelbuilder/goclaw/internal/store" +) + +type PGUsageEventStore struct { + db *sql.DB +} + +func NewPGUsageEventStore(db *sql.DB) *PGUsageEventStore { + return &PGUsageEventStore{db: db} +} + +const usageEventFieldCount = 28 +const usageRollupFieldCount = 21 + +func (s *PGUsageEventStore) InsertEvent(ctx context.Context, event *store.UsageEvent) error { + if event == nil { + return nil + } + return s.InsertEvents(ctx, []store.UsageEvent{*event}) +} + +func (s *PGUsageEventStore) InsertEvents(ctx context.Context, events []store.UsageEvent) error { + if len(events) == 0 { + return nil + } + for i := range events { + prepareUsageEvent(ctx, &events[i]) + } + + vals := make([]string, len(events)) + args := make([]any, 0, len(events)*usageEventFieldCount) + for i, event := range events { + base := i * usageEventFieldCount + placeholders := make([]string, usageEventFieldCount) + for j := range usageEventFieldCount { + placeholders[j] = fmt.Sprintf("$%d", base+j+1) + } + vals[i] = "(" + strings.Join(placeholders, ", ") + ")" + args = append(args, + event.ID, event.TenantID, event.EventTime, event.BucketHour, + 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.DurationMS, event.CallCount, event.ErrorCount, jsonOrNull(event.Metadata), event.CreatedAt, + ) + } + + query := `INSERT INTO usage_events ( + id, tenant_id, event_time, bucket_hour, + 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, + duration_ms, call_count, error_count, metadata, created_at + ) VALUES ` + strings.Join(vals, ", ") + ` + ON CONFLICT DO NOTHING` + _, err := s.db.ExecContext(ctx, query, args...) + return err +} + +func (s *PGUsageEventStore) RefreshEventRollupHour(ctx context.Context, bucketHour time.Time) error { + start := bucketHour.UTC().Truncate(time.Hour) + end := start.Add(time.Hour) + rows, err := s.db.QueryContext(ctx, `SELECT + tenant_id, + bucket_hour, + event_type, + resource_type, + resource_name, + source, + agent_id, + channel, + provider, + model, + status, + COALESCE(SUM(input_tokens), 0), + COALESCE(SUM(output_tokens), 0), + COALESCE(SUM(total_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) + ELSE 0 END, + COALESCE(SUM(call_count), 0), + COALESCE(SUM(error_count), 0) + FROM usage_events + WHERE event_time >= $1 AND event_time < $2 + GROUP BY tenant_id, bucket_hour, event_type, resource_type, resource_name, source, agent_id, channel, provider, model, status`, + start, end) + if err != nil { + return fmt.Errorf("aggregate usage event rollup: %w", err) + } + defer rows.Close() + + now := time.Now().UTC() + var rollups []store.UsageEventRollup + for rows.Next() { + rollup := store.UsageEventRollup{ID: uuid.New(), CreatedAt: now, UpdatedAt: now} + if err := rows.Scan( + &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.DurationMS, &rollup.CallCount, &rollup.ErrorCount, + ); err != nil { + return fmt.Errorf("scan usage event rollup: %w", err) + } + rollups = append(rollups, rollup) + } + if err := rows.Err(); err != nil { + return err + } + return s.upsertEventRollups(ctx, rollups) +} + +func (s *PGUsageEventStore) GetLatestEventRollupBucket(ctx context.Context) (*time.Time, error) { + var t sql.NullTime + err := s.db.QueryRowContext(ctx, `SELECT MAX(bucket_hour) FROM usage_event_rollups`).Scan(&t) + if err != nil { + return nil, fmt.Errorf("get latest event rollup bucket: %w", err) + } + if !t.Valid { + return nil, nil + } + return &t.Time, nil +} + +func (s *PGUsageEventStore) upsertEventRollups(ctx context.Context, rollups []store.UsageEventRollup) error { + if len(rollups) == 0 { + return nil + } + vals := make([]string, len(rollups)) + args := make([]any, 0, len(rollups)*usageRollupFieldCount) + for i, rollup := range rollups { + base := i * usageRollupFieldCount + placeholders := make([]string, usageRollupFieldCount) + for j := range usageRollupFieldCount { + placeholders[j] = fmt.Sprintf("$%d", base+j+1) + } + vals[i] = "(" + strings.Join(placeholders, ", ") + ")" + args = append(args, + 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.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, + duration_ms, call_count, error_count, created_at, updated_at + ) VALUES ` + strings.Join(vals, ", ") + ` + ON CONFLICT ( + tenant_id, + bucket_hour, + event_type, + resource_type, + resource_name, + source, + COALESCE(agent_id, '00000000-0000-0000-0000-000000000000'::uuid), + channel, + provider, + model, + status + ) DO UPDATE SET + input_tokens = EXCLUDED.input_tokens, + output_tokens = EXCLUDED.output_tokens, + total_tokens = EXCLUDED.total_tokens, + cost_usd = EXCLUDED.cost_usd, + duration_ms = EXCLUDED.duration_ms, + call_count = EXCLUDED.call_count, + error_count = EXCLUDED.error_count, + updated_at = EXCLUDED.updated_at` + _, err := s.db.ExecContext(ctx, query, args...) + return err +} + +func (s *PGUsageEventStore) GetEventTimeSeries(ctx context.Context, q store.UsageEventQuery) ([]store.UsageEventTimeSeries, error) { + bucketExpr := "bucket_hour" + if q.GroupBy == "day" { + bucketExpr = "date_trunc('day', bucket_hour)" + } + where, args := buildUsageEventWhere(ctx, q, "bucket_hour") + query := fmt.Sprintf(`SELECT + %s AS bucket_time, + COALESCE(SUM(call_count), 0), + COALESCE(SUM(error_count), 0), + COALESCE(SUM(input_tokens), 0), + COALESCE(SUM(output_tokens), 0), + COALESCE(SUM(total_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) + ELSE 0 END + FROM usage_event_rollups + %s + GROUP BY bucket_time + ORDER BY bucket_time`, bucketExpr, where) + + rows, err := s.db.QueryContext(ctx, query, args...) + if err != nil { + return nil, fmt.Errorf("get usage event timeseries: %w", err) + } + defer rows.Close() + + var result []store.UsageEventTimeSeries + for rows.Next() { + var point store.UsageEventTimeSeries + if err := rows.Scan( + &point.BucketTime, &point.Calls, &point.Errors, + &point.InputTokens, &point.OutputTokens, &point.TotalTokens, + &point.CostUSD, &point.AvgDurationMS, + ); err != nil { + return nil, fmt.Errorf("scan usage event timeseries: %w", err) + } + result = append(result, point) + } + return result, rows.Err() +} + +func (s *PGUsageEventStore) GetEventBreakdown(ctx context.Context, q store.UsageEventQuery) ([]store.UsageEventBreakdown, error) { + groupCol := usageEventGroupColumn(q.GroupBy) + where, args := buildUsageEventWhere(ctx, q, "bucket_hour") + if where == "" { + where = " WHERE 1=1" + } + limit := q.Limit + if limit <= 0 || limit > 100 { + limit = 25 + } + args = append(args, limit) + limitPlaceholder := fmt.Sprintf("$%d", len(args)) + + query := fmt.Sprintf(`SELECT + %s AS key, + MIN(event_type), + MIN(resource_type), + MIN(resource_name), + MIN(source), + COALESCE(SUM(call_count), 0), + COALESCE(SUM(error_count), 0), + COALESCE(SUM(input_tokens), 0), + COALESCE(SUM(output_tokens), 0), + COALESCE(SUM(total_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) + ELSE 0 END + FROM usage_event_rollups + %s + GROUP BY %s + ORDER BY SUM(call_count) DESC, key ASC + LIMIT %s`, groupCol, where, groupCol, limitPlaceholder) + + rows, err := s.db.QueryContext(ctx, query, args...) + if err != nil { + return nil, fmt.Errorf("get usage event breakdown: %w", err) + } + defer rows.Close() + + var result []store.UsageEventBreakdown + for rows.Next() { + var row store.UsageEventBreakdown + 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.CostUSD, &row.AvgDurationMS, + ); err != nil { + return nil, fmt.Errorf("scan usage event breakdown: %w", err) + } + result = append(result, row) + } + return result, rows.Err() +} + +func (s *PGUsageEventStore) GetEventSummary(ctx context.Context, q store.UsageEventQuery) (*store.UsageEventSummary, error) { + where, args := buildUsageEventWhere(ctx, q, "bucket_hour") + query := `SELECT + COALESCE(SUM(call_count), 0), + COALESCE(SUM(error_count), 0), + COALESCE(SUM(input_tokens), 0), + COALESCE(SUM(output_tokens), 0), + COALESCE(SUM(total_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) + ELSE 0 END + FROM usage_event_rollups` + where + 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, + ); err != nil { + return nil, fmt.Errorf("get usage event summary: %w", err) + } + return &summary, nil +} + +func prepareUsageEvent(ctx context.Context, event *store.UsageEvent) { + if event.ID == uuid.Nil { + event.ID = uuid.New() + } + if event.TenantID == uuid.Nil { + event.TenantID = store.TenantIDFromContext(ctx) + } + if event.TenantID == uuid.Nil { + event.TenantID = store.MasterTenantID + } + if event.EventTime.IsZero() { + event.EventTime = time.Now().UTC() + } + event.EventTime = event.EventTime.UTC() + if event.BucketHour.IsZero() { + event.BucketHour = event.EventTime.Truncate(time.Hour) + } + if event.CallCount <= 0 { + event.CallCount = 1 + } + if event.Status == "" { + event.Status = "completed" + } + if event.CreatedAt.IsZero() { + event.CreatedAt = time.Now().UTC() + } +} + +func buildUsageEventWhere(ctx context.Context, q store.UsageEventQuery, timeColumn string) (string, []any) { + var conds []string + var args []any + idx := 1 + + if !store.IsCrossTenant(ctx) { + if tenantID := store.TenantIDFromContext(ctx); tenantID != uuid.Nil { + conds = append(conds, fmt.Sprintf("tenant_id = $%d", idx)) + args = append(args, tenantID) + idx++ + } + } + add := func(col string, value any) { + conds = append(conds, fmt.Sprintf("%s = $%d", col, idx)) + args = append(args, value) + idx++ + } + if !q.From.IsZero() { + conds = append(conds, fmt.Sprintf("%s >= $%d", timeColumn, idx)) + args = append(args, q.From.UTC()) + idx++ + } + if !q.To.IsZero() { + conds = append(conds, fmt.Sprintf("%s < $%d", timeColumn, idx)) + args = append(args, q.To.UTC()) + idx++ + } + if q.AgentID != nil { + add("agent_id", *q.AgentID) + } + if q.Channel != "" { + add("channel", q.Channel) + } + if q.EventType != "" { + add("event_type", q.EventType) + } + if q.ResourceType != "" { + add("resource_type", q.ResourceType) + } + if q.ResourceName != "" { + add("resource_name", q.ResourceName) + } + if q.Provider != "" { + add("provider", q.Provider) + } + if q.Model != "" { + add("model", q.Model) + } + if q.Status != "" { + add("status", q.Status) + } + if q.Source != "" { + add("source", q.Source) + } + if len(conds) == 0 { + return "", nil + } + return " WHERE " + strings.Join(conds, " AND "), args +} + +func usageEventGroupColumn(groupBy string) string { + switch groupBy { + case "event_type": + return "event_type" + case "resource_type": + return "resource_type" + case "source": + return "source" + case "status": + return "status" + case "agent": + return "COALESCE(agent_id::TEXT, '')" + case "channel": + return "channel" + case "provider": + return "provider" + case "model": + return "model" + default: + return "resource_name" + } +} diff --git a/internal/store/run_context.go b/internal/store/run_context.go index 2c825571..4645b90c 100644 --- a/internal/store/run_context.go +++ b/internal/store/run_context.go @@ -25,6 +25,8 @@ type RunContext struct { AgentKey string TenantID uuid.UUID UserID string + RunID string + SessionKey string CredentialUserID string // resolved tenant user for credential lookups (empty = use UserID) AgentType string SenderID string @@ -39,6 +41,7 @@ type RunContext struct { // Tool configuration BuiltinToolSettings map[string][]byte + Channel string ChannelType string ChannelContextScope ChannelContextScope SubagentsCfg *config.SubagentsConfig diff --git a/internal/store/sqlitestore/factory.go b/internal/store/sqlitestore/factory.go index 284c4c1a..e2d8820f 100644 --- a/internal/store/sqlitestore/factory.go +++ b/internal/store/sqlitestore/factory.go @@ -50,6 +50,7 @@ func NewSQLiteStores(cfg store.StoreConfig) (*store.Stores, error) { SkillTenantCfgs: NewSQLiteSkillTenantConfigStore(db), SystemConfigs: NewSQLiteSystemConfigStore(db), Snapshots: NewSQLiteSnapshotStore(db), + UsageEvents: NewSQLiteUsageEventStore(db), Cron: NewSQLiteCronStore(db), ChannelInstances: NewSQLiteChannelInstanceStore(db, cfg.EncryptionKey), Pairing: NewSQLitePairingStore(db), diff --git a/internal/store/sqlitestore/schema.go b/internal/store/sqlitestore/schema.go index af8816e0..fd793772 100644 --- a/internal/store/sqlitestore/schema.go +++ b/internal/store/sqlitestore/schema.go @@ -16,7 +16,7 @@ var schemaSQL string // SchemaVersion is the current SQLite schema version. // Bump this when adding new migration steps below. -const SchemaVersion = 47 +const SchemaVersion = 48 // migrations maps version → SQL to apply when upgrading FROM that version. // schema.sql always represents the LATEST full schema (for fresh DBs). @@ -848,6 +848,91 @@ DROP TABLE skill_user_grants; ALTER TABLE skill_user_grants_new RENAME TO skill_user_grants; CREATE INDEX IF NOT EXISTS idx_skill_user_grants_user ON skill_user_grants(user_id); CREATE INDEX IF NOT EXISTS idx_skill_user_grants_tenant ON skill_user_grants(tenant_id);`, + // Version 47 → 48: append-only usage event analytics. + 47: `CREATE TABLE IF NOT EXISTS usage_events ( + id TEXT NOT NULL PRIMARY KEY, + tenant_id TEXT NOT NULL REFERENCES tenants(id), + event_time TEXT NOT NULL, + bucket_hour TEXT NOT NULL, + event_type TEXT NOT NULL, + resource_type TEXT NOT NULL, + resource_name TEXT NOT NULL, + resource_id TEXT NOT NULL DEFAULT '', + source TEXT NOT NULL DEFAULT '', + agent_id TEXT REFERENCES agents(id) ON DELETE SET NULL, + team_id TEXT, + trace_id TEXT REFERENCES traces(id) ON DELETE SET NULL, + span_id TEXT REFERENCES spans(id) ON DELETE SET NULL, + run_id TEXT NOT NULL DEFAULT '', + session_key TEXT NOT NULL DEFAULT '', + channel TEXT NOT NULL DEFAULT '', + provider TEXT NOT NULL DEFAULT '', + model TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT '', + input_tokens BIGINT NOT NULL DEFAULT 0, + output_tokens BIGINT NOT NULL DEFAULT 0, + total_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, + error_count INTEGER NOT NULL DEFAULT 0, + metadata TEXT, + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')) +); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_time + ON usage_events(tenant_id, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_resource_time + ON usage_events(tenant_id, resource_type, resource_name, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_type_time + ON usage_events(tenant_id, event_type, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_agent_time + ON usage_events(tenant_id, agent_id, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_channel_time + ON usage_events(tenant_id, channel, event_time DESC) WHERE channel != ''; +CREATE UNIQUE INDEX IF NOT EXISTS idx_usage_events_trace_span_type_source + ON usage_events(trace_id, span_id, event_type, source) + WHERE trace_id IS NOT NULL AND span_id IS NOT NULL; +CREATE TABLE IF NOT EXISTS usage_event_rollups ( + id TEXT NOT NULL PRIMARY KEY, + tenant_id TEXT NOT NULL REFERENCES tenants(id), + bucket_hour TEXT NOT NULL, + event_type TEXT NOT NULL, + resource_type TEXT NOT NULL, + resource_name TEXT NOT NULL, + source TEXT NOT NULL DEFAULT '', + agent_id TEXT REFERENCES agents(id) ON DELETE SET NULL, + channel TEXT NOT NULL DEFAULT '', + provider TEXT NOT NULL DEFAULT '', + model TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT '', + input_tokens BIGINT NOT NULL DEFAULT 0, + output_tokens BIGINT NOT NULL DEFAULT 0, + total_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, + error_count INTEGER NOT NULL DEFAULT 0, + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + updated_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')) +); +CREATE UNIQUE INDEX IF NOT EXISTS idx_usage_event_rollups_unique + ON usage_event_rollups( + tenant_id, + bucket_hour, + event_type, + resource_type, + resource_name, + source, + COALESCE(agent_id, '00000000-0000-0000-0000-000000000000'), + channel, + provider, + model, + status + ); +CREATE INDEX IF NOT EXISTS idx_usage_event_rollups_tenant_hour + ON usage_event_rollups(tenant_id, bucket_hour DESC); +CREATE INDEX IF NOT EXISTS idx_usage_event_rollups_resource_hour + ON usage_event_rollups(tenant_id, resource_type, resource_name, bucket_hour DESC);`, } const addChannelMemoryExtractionTables = ` diff --git a/internal/store/sqlitestore/schema.sql b/internal/store/sqlitestore/schema.sql index df8eb703..06304e0e 100644 --- a/internal/store/sqlitestore/schema.sql +++ b/internal/store/sqlitestore/schema.sql @@ -1348,6 +1348,102 @@ CREATE UNIQUE INDEX IF NOT EXISTS idx_usage_snapshots_unique ON usage_snapshots( ); CREATE INDEX IF NOT EXISTS idx_usage_snapshots_tenant ON usage_snapshots(tenant_id); +-- ============================================================ +-- Table: usage_events +-- ============================================================ + +CREATE TABLE IF NOT EXISTS usage_events ( + id TEXT NOT NULL PRIMARY KEY, + tenant_id TEXT NOT NULL REFERENCES tenants(id), + event_time TEXT NOT NULL, + bucket_hour TEXT NOT NULL, + event_type TEXT NOT NULL, + resource_type TEXT NOT NULL, + resource_name TEXT NOT NULL, + resource_id TEXT NOT NULL DEFAULT '', + source TEXT NOT NULL DEFAULT '', + agent_id TEXT REFERENCES agents(id) ON DELETE SET NULL, + team_id TEXT, + trace_id TEXT REFERENCES traces(id) ON DELETE SET NULL, + span_id TEXT REFERENCES spans(id) ON DELETE SET NULL, + run_id TEXT NOT NULL DEFAULT '', + session_key TEXT NOT NULL DEFAULT '', + channel TEXT NOT NULL DEFAULT '', + provider TEXT NOT NULL DEFAULT '', + model TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT '', + input_tokens BIGINT NOT NULL DEFAULT 0, + output_tokens BIGINT NOT NULL DEFAULT 0, + total_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, + error_count INTEGER NOT NULL DEFAULT 0, + metadata TEXT, + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')) +); + +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_time + ON usage_events(tenant_id, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_resource_time + ON usage_events(tenant_id, resource_type, resource_name, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_type_time + ON usage_events(tenant_id, event_type, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_agent_time + ON usage_events(tenant_id, agent_id, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_channel_time + ON usage_events(tenant_id, channel, event_time DESC) WHERE channel != ''; +CREATE UNIQUE INDEX IF NOT EXISTS idx_usage_events_trace_span_type_source + ON usage_events(trace_id, span_id, event_type, source) + WHERE trace_id IS NOT NULL AND span_id IS NOT NULL; + +-- ============================================================ +-- Table: usage_event_rollups +-- ============================================================ + +CREATE TABLE IF NOT EXISTS usage_event_rollups ( + id TEXT NOT NULL PRIMARY KEY, + tenant_id TEXT NOT NULL REFERENCES tenants(id), + bucket_hour TEXT NOT NULL, + event_type TEXT NOT NULL, + resource_type TEXT NOT NULL, + resource_name TEXT NOT NULL, + source TEXT NOT NULL DEFAULT '', + agent_id TEXT REFERENCES agents(id) ON DELETE SET NULL, + channel TEXT NOT NULL DEFAULT '', + provider TEXT NOT NULL DEFAULT '', + model TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT '', + input_tokens BIGINT NOT NULL DEFAULT 0, + output_tokens BIGINT NOT NULL DEFAULT 0, + total_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, + error_count INTEGER NOT NULL DEFAULT 0, + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + updated_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')) +); + +CREATE UNIQUE INDEX IF NOT EXISTS idx_usage_event_rollups_unique + ON usage_event_rollups( + tenant_id, + bucket_hour, + event_type, + resource_type, + resource_name, + source, + COALESCE(agent_id, '00000000-0000-0000-0000-000000000000'), + channel, + provider, + model, + status + ); +CREATE INDEX IF NOT EXISTS idx_usage_event_rollups_tenant_hour + ON usage_event_rollups(tenant_id, bucket_hour DESC); +CREATE INDEX IF NOT EXISTS idx_usage_event_rollups_resource_hour + ON usage_event_rollups(tenant_id, resource_type, resource_name, bucket_hour DESC); + -- ============================================================ -- Table: builtin_tools -- ============================================================ diff --git a/internal/store/sqlitestore/usage_events.go b/internal/store/sqlitestore/usage_events.go new file mode 100644 index 00000000..c2c276a9 --- /dev/null +++ b/internal/store/sqlitestore/usage_events.go @@ -0,0 +1,423 @@ +//go:build sqlite || sqliteonly + +package sqlitestore + +import ( + "context" + "database/sql" + "fmt" + "strings" + "time" + + "github.com/google/uuid" + + "github.com/nextlevelbuilder/goclaw/internal/store" +) + +type SQLiteUsageEventStore struct { + db *sql.DB +} + +func NewSQLiteUsageEventStore(db *sql.DB) *SQLiteUsageEventStore { + return &SQLiteUsageEventStore{db: db} +} + +const sqliteUsageEventFieldCount = 28 +const sqliteUsageEventBatchSize = 30 +const sqliteUsageRollupFieldCount = 21 + +func (s *SQLiteUsageEventStore) InsertEvent(ctx context.Context, event *store.UsageEvent) error { + if event == nil { + return nil + } + return s.InsertEvents(ctx, []store.UsageEvent{*event}) +} + +func (s *SQLiteUsageEventStore) InsertEvents(ctx context.Context, events []store.UsageEvent) error { + if len(events) == 0 { + return nil + } + for start := 0; start < len(events); start += sqliteUsageEventBatchSize { + end := start + sqliteUsageEventBatchSize + if end > len(events) { + end = len(events) + } + if err := s.insertBatch(ctx, events[start:end]); err != nil { + return err + } + } + return nil +} + +func (s *SQLiteUsageEventStore) insertBatch(ctx context.Context, events []store.UsageEvent) error { + placeholderRow := "(" + strings.Repeat("?, ", sqliteUsageEventFieldCount-1) + "?)" + vals := make([]string, len(events)) + args := make([]any, 0, len(events)*sqliteUsageEventFieldCount) + for i := range events { + prepareSQLiteUsageEvent(ctx, &events[i]) + event := events[i] + vals[i] = placeholderRow + args = append(args, + event.ID, event.TenantID, event.EventTime, event.BucketHour, + 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.DurationMS, event.CallCount, event.ErrorCount, jsonOrNull(event.Metadata), event.CreatedAt, + ) + } + query := `INSERT INTO usage_events ( + id, tenant_id, event_time, bucket_hour, + 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, + duration_ms, call_count, error_count, metadata, created_at + ) VALUES ` + strings.Join(vals, ", ") + ` + ON CONFLICT DO NOTHING` + _, err := s.db.ExecContext(ctx, query, args...) + return err +} + +func (s *SQLiteUsageEventStore) RefreshEventRollupHour(ctx context.Context, bucketHour time.Time) error { + start := bucketHour.UTC().Truncate(time.Hour) + end := start.Add(time.Hour) + rows, err := s.db.QueryContext(ctx, `SELECT + tenant_id, + bucket_hour, + event_type, + resource_type, + resource_name, + source, + agent_id, + channel, + provider, + model, + status, + COALESCE(SUM(input_tokens), 0), + COALESCE(SUM(output_tokens), 0), + COALESCE(SUM(total_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) + ELSE 0 END, + COALESCE(SUM(call_count), 0), + COALESCE(SUM(error_count), 0) + FROM usage_events + WHERE event_time >= ? AND event_time < ? + GROUP BY tenant_id, bucket_hour, event_type, resource_type, resource_name, source, agent_id, channel, provider, model, status`, + start, end) + if err != nil { + return fmt.Errorf("aggregate usage event rollup: %w", err) + } + defer rows.Close() + + now := time.Now().UTC() + var rollups []store.UsageEventRollup + for rows.Next() { + rollup := store.UsageEventRollup{ID: uuid.New(), CreatedAt: now, UpdatedAt: now} + var bucketTime sqliteTime + if err := rows.Scan( + &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.DurationMS, &rollup.CallCount, &rollup.ErrorCount, + ); err != nil { + return fmt.Errorf("scan usage event rollup: %w", err) + } + rollup.BucketHour = bucketTime.Time + rollups = append(rollups, rollup) + } + if err := rows.Err(); err != nil { + return err + } + return s.upsertEventRollups(ctx, rollups) +} + +func (s *SQLiteUsageEventStore) GetLatestEventRollupBucket(ctx context.Context) (*time.Time, error) { + var nt nullSqliteTime + err := s.db.QueryRowContext(ctx, `SELECT MAX(bucket_hour) FROM usage_event_rollups`).Scan(&nt) + if err != nil { + return nil, fmt.Errorf("get latest event rollup bucket: %w", err) + } + if !nt.Valid { + return nil, nil + } + return &nt.Time, nil +} + +func (s *SQLiteUsageEventStore) upsertEventRollups(ctx context.Context, rollups []store.UsageEventRollup) error { + if len(rollups) == 0 { + return nil + } + placeholderRow := "(" + strings.Repeat("?, ", sqliteUsageRollupFieldCount-1) + "?)" + vals := make([]string, len(rollups)) + args := make([]any, 0, len(rollups)*sqliteUsageRollupFieldCount) + for i, rollup := range rollups { + vals[i] = placeholderRow + args = append(args, + 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.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, + duration_ms, call_count, error_count, created_at, updated_at + ) VALUES ` + strings.Join(vals, ", ") + ` + ON CONFLICT ( + tenant_id, + bucket_hour, + event_type, + resource_type, + resource_name, + source, + COALESCE(agent_id, '00000000-0000-0000-0000-000000000000'), + channel, + provider, + model, + status + ) DO UPDATE SET + input_tokens = excluded.input_tokens, + output_tokens = excluded.output_tokens, + total_tokens = excluded.total_tokens, + cost_usd = excluded.cost_usd, + duration_ms = excluded.duration_ms, + call_count = excluded.call_count, + error_count = excluded.error_count, + updated_at = excluded.updated_at` + _, err := s.db.ExecContext(ctx, query, args...) + return err +} + +func (s *SQLiteUsageEventStore) GetEventTimeSeries(ctx context.Context, q store.UsageEventQuery) ([]store.UsageEventTimeSeries, error) { + bucketExpr := "bucket_hour" + if q.GroupBy == "day" { + bucketExpr = "strftime('%Y-%m-%d 00:00:00', bucket_hour)" + } + where, args := buildSQLiteUsageEventWhere(ctx, q, "bucket_hour") + query := fmt.Sprintf(`SELECT + %s AS bucket_time, + COALESCE(SUM(call_count), 0), + COALESCE(SUM(error_count), 0), + COALESCE(SUM(input_tokens), 0), + COALESCE(SUM(output_tokens), 0), + COALESCE(SUM(total_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) + ELSE 0 END + FROM usage_event_rollups + %s + GROUP BY bucket_time + ORDER BY bucket_time`, bucketExpr, where) + + rows, err := s.db.QueryContext(ctx, query, args...) + if err != nil { + return nil, fmt.Errorf("get usage event timeseries: %w", err) + } + defer rows.Close() + + var result []store.UsageEventTimeSeries + for rows.Next() { + var point store.UsageEventTimeSeries + var bucketTime sqliteTime + if err := rows.Scan( + &bucketTime, &point.Calls, &point.Errors, + &point.InputTokens, &point.OutputTokens, &point.TotalTokens, + &point.CostUSD, &point.AvgDurationMS, + ); err != nil { + return nil, fmt.Errorf("scan usage event timeseries: %w", err) + } + point.BucketTime = bucketTime.Time + result = append(result, point) + } + return result, rows.Err() +} + +func (s *SQLiteUsageEventStore) GetEventBreakdown(ctx context.Context, q store.UsageEventQuery) ([]store.UsageEventBreakdown, error) { + groupCol := sqliteUsageEventGroupColumn(q.GroupBy) + where, args := buildSQLiteUsageEventWhere(ctx, q, "bucket_hour") + if where == "" { + where = " WHERE 1=1" + } + limit := q.Limit + if limit <= 0 || limit > 100 { + limit = 25 + } + args = append(args, limit) + query := fmt.Sprintf(`SELECT + %s AS key, + MIN(event_type), + MIN(resource_type), + MIN(resource_name), + MIN(source), + COALESCE(SUM(call_count), 0), + COALESCE(SUM(error_count), 0), + COALESCE(SUM(input_tokens), 0), + COALESCE(SUM(output_tokens), 0), + COALESCE(SUM(total_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) + ELSE 0 END + FROM usage_event_rollups + %s + GROUP BY %s + ORDER BY SUM(call_count) DESC, key ASC + LIMIT ?`, groupCol, where, groupCol) + + rows, err := s.db.QueryContext(ctx, query, args...) + if err != nil { + return nil, fmt.Errorf("get usage event breakdown: %w", err) + } + defer rows.Close() + + var result []store.UsageEventBreakdown + for rows.Next() { + var row store.UsageEventBreakdown + 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.CostUSD, &row.AvgDurationMS, + ); err != nil { + return nil, fmt.Errorf("scan usage event breakdown: %w", err) + } + result = append(result, row) + } + return result, rows.Err() +} + +func (s *SQLiteUsageEventStore) GetEventSummary(ctx context.Context, q store.UsageEventQuery) (*store.UsageEventSummary, error) { + where, args := buildSQLiteUsageEventWhere(ctx, q, "bucket_hour") + query := `SELECT + COALESCE(SUM(call_count), 0), + COALESCE(SUM(error_count), 0), + COALESCE(SUM(input_tokens), 0), + COALESCE(SUM(output_tokens), 0), + COALESCE(SUM(total_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) + ELSE 0 END + FROM usage_event_rollups` + where + 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, + ); err != nil { + return nil, fmt.Errorf("get usage event summary: %w", err) + } + return &summary, nil +} + +func prepareSQLiteUsageEvent(ctx context.Context, event *store.UsageEvent) { + if event.ID == uuid.Nil { + event.ID = uuid.New() + } + if event.TenantID == uuid.Nil { + event.TenantID = store.TenantIDFromContext(ctx) + } + if event.TenantID == uuid.Nil { + event.TenantID = store.MasterTenantID + } + if event.EventTime.IsZero() { + event.EventTime = time.Now().UTC() + } + event.EventTime = event.EventTime.UTC() + if event.BucketHour.IsZero() { + event.BucketHour = event.EventTime.Truncate(time.Hour) + } + if event.CallCount <= 0 { + event.CallCount = 1 + } + if event.Status == "" { + event.Status = "completed" + } + if event.CreatedAt.IsZero() { + event.CreatedAt = time.Now().UTC() + } +} + +func buildSQLiteUsageEventWhere(ctx context.Context, q store.UsageEventQuery, timeColumn string) (string, []any) { + var conds []string + var args []any + + if !store.IsCrossTenant(ctx) { + if tenantID := store.TenantIDFromContext(ctx); tenantID != uuid.Nil { + conds = append(conds, "tenant_id = ?") + args = append(args, tenantID) + } + } + add := func(col string, value any) { + conds = append(conds, col+" = ?") + args = append(args, value) + } + if !q.From.IsZero() { + conds = append(conds, timeColumn+" >= ?") + args = append(args, q.From.UTC()) + } + if !q.To.IsZero() { + conds = append(conds, timeColumn+" < ?") + args = append(args, q.To.UTC()) + } + if q.AgentID != nil { + add("agent_id", *q.AgentID) + } + if q.Channel != "" { + add("channel", q.Channel) + } + if q.EventType != "" { + add("event_type", q.EventType) + } + if q.ResourceType != "" { + add("resource_type", q.ResourceType) + } + if q.ResourceName != "" { + add("resource_name", q.ResourceName) + } + if q.Provider != "" { + add("provider", q.Provider) + } + if q.Model != "" { + add("model", q.Model) + } + if q.Status != "" { + add("status", q.Status) + } + if q.Source != "" { + add("source", q.Source) + } + if len(conds) == 0 { + return "", nil + } + return " WHERE " + strings.Join(conds, " AND "), args +} + +func sqliteUsageEventGroupColumn(groupBy string) string { + switch groupBy { + case "event_type": + return "event_type" + case "resource_type": + return "resource_type" + case "source": + return "source" + case "status": + return "status" + case "agent": + return "COALESCE(agent_id, '')" + case "channel": + return "channel" + case "provider": + return "provider" + case "model": + return "model" + default: + return "resource_name" + } +} diff --git a/internal/store/sqlitestore/usage_events_test.go b/internal/store/sqlitestore/usage_events_test.go new file mode 100644 index 00000000..c6c94063 --- /dev/null +++ b/internal/store/sqlitestore/usage_events_test.go @@ -0,0 +1,88 @@ +//go:build sqlite || sqliteonly + +package sqlitestore + +import ( + "context" + "testing" + "time" + + "github.com/google/uuid" + + "github.com/nextlevelbuilder/goclaw/internal/store" +) + +func TestSQLiteUsageEventStoreRefreshRollupHourIsIdempotent(t *testing.T) { + db := openTestDB(t) + if err := EnsureSchema(db); err != nil { + t.Fatalf("EnsureSchema: %v", err) + } + + eventStore := NewSQLiteUsageEventStore(db) + ctx := store.WithTenantID(context.Background(), store.MasterTenantID) + bucket := time.Date(2026, 6, 12, 8, 0, 0, 0, time.UTC) + + if err := eventStore.InsertEvent(ctx, &store.UsageEvent{ + ID: uuid.New(), + TenantID: store.MasterTenantID, + EventTime: bucket.Add(10 * time.Minute), + EventType: store.UsageEventTypeRuntimeToolCall, + ResourceType: store.UsageResourceTypeRuntimeTool, + ResourceName: "exec", + ResourceID: "exec", + Source: store.UsageSourceToolCall, + Status: "completed", + InputTokens: 10, + TotalTokens: 10, + DurationMS: 50, + CallCount: 1, + }); err != nil { + t.Fatalf("InsertEvent first: %v", err) + } + if err := eventStore.RefreshEventRollupHour(ctx, bucket); err != nil { + t.Fatalf("Refresh first: %v", err) + } + + if err := eventStore.InsertEvent(ctx, &store.UsageEvent{ + ID: uuid.New(), + TenantID: store.MasterTenantID, + EventTime: bucket.Add(20 * time.Minute), + EventType: store.UsageEventTypeRuntimeToolCall, + ResourceType: store.UsageResourceTypeRuntimeTool, + ResourceName: "exec", + ResourceID: "exec", + Source: store.UsageSourceToolCall, + Status: "error", + OutputTokens: 5, + TotalTokens: 5, + DurationMS: 150, + CallCount: 1, + ErrorCount: 1, + }); err != nil { + t.Fatalf("InsertEvent second: %v", err) + } + if err := eventStore.RefreshEventRollupHour(ctx, bucket); err != nil { + t.Fatalf("Refresh second: %v", err) + } + if err := eventStore.RefreshEventRollupHour(ctx, bucket); err != nil { + t.Fatalf("Refresh idempotent: %v", err) + } + + summary, err := eventStore.GetEventSummary(ctx, store.UsageEventQuery{ + From: bucket, + To: bucket.Add(time.Hour), + ResourceType: store.UsageResourceTypeRuntimeTool, + }) + if err != nil { + t.Fatalf("GetEventSummary: %v", err) + } + if summary.Calls != 2 || summary.Errors != 1 { + t.Fatalf("summary calls/errors = %d/%d, want 2/1", summary.Calls, summary.Errors) + } + if summary.TotalTokens != 15 { + t.Fatalf("summary tokens = %d, want 15", summary.TotalTokens) + } + if summary.AvgDurationMS != 100 { + t.Fatalf("summary avg duration = %d, want 100", summary.AvgDurationMS) + } +} diff --git a/internal/store/stores.go b/internal/store/stores.go index 2b3b07b6..13eff72f 100644 --- a/internal/store/stores.go +++ b/internal/store/stores.go @@ -26,6 +26,7 @@ type Stores struct { Contacts ContactStore Activity ActivityStore Snapshots SnapshotStore + UsageEvents UsageEventStore BrowserCookies BrowserCookieStore SecureCLI SecureCLIStore SecureCLIGrants SecureCLIAgentGrantStore diff --git a/internal/store/usage_event_store.go b/internal/store/usage_event_store.go new file mode 100644 index 00000000..19b14730 --- /dev/null +++ b/internal/store/usage_event_store.go @@ -0,0 +1,145 @@ +package store + +import ( + "context" + "encoding/json" + "time" + + "github.com/google/uuid" +) + +const ( + UsageEventTypeToolCall = "tool_call" + UsageEventTypeSkillActivation = "skill_activation" + UsageEventTypeMCPToolCall = "mcp_tool_call" + UsageEventTypeRuntimeToolCall = "runtime_tool_call" + + UsageResourceTypeTool = "tool" + UsageResourceTypeSkill = "skill" + UsageResourceTypeMCPTool = "mcp_tool" + UsageResourceTypeRuntimeTool = "runtime_tool" + + UsageSourceToolCall = "tool_call" + UsageSourceUseSkill = "use_skill" + UsageSourceSlashCommand = "slash-command" +) + +// 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"` +} + +// UsageEventQuery filters usage event analytics. +type UsageEventQuery struct { + From time.Time + To time.Time + AgentID *uuid.UUID + Channel string + EventType string + ResourceType string + ResourceName string + Provider string + Model string + Status string + Source string + GroupBy string + Limit int +} + +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"` +} + +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"` +} + +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"` +} + +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"` +} + +type UsageEventStore interface { + InsertEvent(ctx context.Context, event *UsageEvent) error + InsertEvents(ctx context.Context, events []UsageEvent) error + RefreshEventRollupHour(ctx context.Context, bucketHour time.Time) error + GetLatestEventRollupBucket(ctx context.Context) (*time.Time, error) + GetEventTimeSeries(ctx context.Context, q UsageEventQuery) ([]UsageEventTimeSeries, error) + GetEventBreakdown(ctx context.Context, q UsageEventQuery) ([]UsageEventBreakdown, error) + GetEventSummary(ctx context.Context, q UsageEventQuery) (*UsageEventSummary, error) +} diff --git a/internal/tracing/snapshot_worker.go b/internal/tracing/snapshot_worker.go index a5c14791..672c364e 100644 --- a/internal/tracing/snapshot_worker.go +++ b/internal/tracing/snapshot_worker.go @@ -14,17 +14,19 @@ import ( // SnapshotWorker periodically aggregates trace/span data into usage_snapshots. type SnapshotWorker struct { - db *sql.DB - snapshots store.SnapshotStore - stopCh chan struct{} - wg sync.WaitGroup + db *sql.DB + snapshots store.SnapshotStore + usageEvents store.UsageEventStore + stopCh chan struct{} + wg sync.WaitGroup } -func NewSnapshotWorker(db *sql.DB, snapshots store.SnapshotStore) *SnapshotWorker { +func NewSnapshotWorker(db *sql.DB, snapshots store.SnapshotStore, usageEvents store.UsageEventStore) *SnapshotWorker { return &SnapshotWorker{ - db: db, - snapshots: snapshots, - stopCh: make(chan struct{}), + db: db, + snapshots: snapshots, + usageEvents: usageEvents, + stopCh: make(chan struct{}), } } @@ -73,6 +75,11 @@ func (w *SnapshotWorker) loop() { // catchUp computes snapshots for all missed hours between latest bucket and current hour. func (w *SnapshotWorker) catchUp() { ctx := context.Background() + w.catchUpSnapshots(ctx) + w.catchUpUsageEvents(ctx) +} + +func (w *SnapshotWorker) catchUpSnapshots(ctx context.Context) { now := time.Now().UTC() targetHour := now.Truncate(time.Hour).Add(-time.Hour) // previous complete hour @@ -100,6 +107,34 @@ func (w *SnapshotWorker) catchUp() { } } +func (w *SnapshotWorker) catchUpUsageEvents(ctx context.Context) { + if w.usageEvents == nil { + return + } + now := time.Now().UTC() + targetHour := now.Truncate(time.Hour).Add(-time.Hour) + + latest, err := w.usageEvents.GetLatestEventRollupBucket(ctx) + if err != nil { + slog.Warn("usage_event_rollup: get latest bucket", "error", err) + return + } + + startHour := targetHour + if latest != nil { + startHour = latest.Add(time.Hour) + } + + for h := startHour; !h.After(targetHour); h = h.Add(time.Hour) { + start := time.Now() + if err := w.usageEvents.RefreshEventRollupHour(ctx, h); err != nil { + slog.Warn("usage_event_rollup: aggregate hour failed", "hour", h.Format(time.RFC3339), "error", err) + return + } + slog.Info("usage event rollup computed", "hour", h.Format(time.RFC3339), "duration_ms", time.Since(start).Milliseconds()) + } +} + // Backfill populates usage_snapshots from historical trace/span data. // Returns the number of hours processed. func (w *SnapshotWorker) Backfill(ctx context.Context) (int, error) { @@ -131,6 +166,41 @@ func (w *SnapshotWorker) Backfill(ctx context.Context) (int, error) { } count++ } + eventCount, err := w.backfillUsageEvents(ctx) + if err != nil { + return count, err + } + return count + eventCount, nil +} + +func (w *SnapshotWorker) backfillUsageEvents(ctx context.Context) (int, error) { + if w.usageEvents == nil { + return 0, nil + } + latest, err := w.usageEvents.GetLatestEventRollupBucket(ctx) + if err != nil { + return 0, fmt.Errorf("get latest event rollup bucket: %w", err) + } + + var earliest sql.NullTime + if err := w.db.QueryRowContext(ctx, `SELECT MIN(event_time) FROM usage_events`).Scan(&earliest); err != nil || !earliest.Valid { + return 0, nil + } + + startHour := earliest.Time.UTC().Truncate(time.Hour) + if latest != nil { + startHour = latest.Add(time.Hour) + } + + endHour := time.Now().UTC().Truncate(time.Hour) + count := 0 + for h := startHour; h.Before(endHour); h = h.Add(time.Hour) { + if err := w.usageEvents.RefreshEventRollupHour(ctx, h); err != nil { + slog.Warn("usage_event_rollup backfill: aggregate hour failed", "hour", h.Format(time.RFC3339), "error", err) + continue + } + count++ + } return count, nil } @@ -288,8 +358,8 @@ type agentMemoryCounts struct { // agentKGCounts holds point-in-time KG counts for one agent. type agentKGCounts struct { - AgentID uuid.UUID - Entities int + AgentID uuid.UUID + Entities int Relations int } diff --git a/internal/upgrade/version.go b/internal/upgrade/version.go index 67b4ce6e..04f1a23d 100644 --- a/internal/upgrade/version.go +++ b/internal/upgrade/version.go @@ -2,4 +2,4 @@ package upgrade // RequiredSchemaVersion is the schema migration version this binary requires. // Bump this whenever adding a new SQL migration file. -const RequiredSchemaVersion uint = 78 +const RequiredSchemaVersion uint = 79 diff --git a/migrations/000079_usage_event_analytics.down.sql b/migrations/000079_usage_event_analytics.down.sql new file mode 100644 index 00000000..9b387991 --- /dev/null +++ b/migrations/000079_usage_event_analytics.down.sql @@ -0,0 +1,2 @@ +DROP TABLE IF EXISTS usage_event_rollups; +DROP TABLE IF EXISTS usage_events; diff --git a/migrations/000079_usage_event_analytics.up.sql b/migrations/000079_usage_event_analytics.up.sql new file mode 100644 index 00000000..2fb7f6e3 --- /dev/null +++ b/migrations/000079_usage_event_analytics.up.sql @@ -0,0 +1,88 @@ +CREATE TABLE IF NOT EXISTS usage_events ( + id UUID PRIMARY KEY, + tenant_id UUID NOT NULL REFERENCES tenants(id), + event_time TIMESTAMPTZ NOT NULL, + bucket_hour TIMESTAMPTZ NOT NULL, + event_type TEXT NOT NULL, + resource_type TEXT NOT NULL, + resource_name TEXT NOT NULL, + resource_id TEXT NOT NULL DEFAULT '', + source TEXT NOT NULL DEFAULT '', + agent_id UUID REFERENCES agents(id) ON DELETE SET NULL, + team_id UUID REFERENCES teams(id) ON DELETE SET NULL, + trace_id UUID REFERENCES traces(id) ON DELETE SET NULL, + span_id UUID REFERENCES spans(id) ON DELETE SET NULL, + run_id TEXT NOT NULL DEFAULT '', + session_key TEXT NOT NULL DEFAULT '', + channel TEXT NOT NULL DEFAULT '', + provider TEXT NOT NULL DEFAULT '', + model TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT '', + input_tokens BIGINT NOT NULL DEFAULT 0, + output_tokens BIGINT NOT NULL DEFAULT 0, + total_tokens BIGINT NOT NULL DEFAULT 0, + cost_usd DOUBLE PRECISION NOT NULL DEFAULT 0, + duration_ms INTEGER NOT NULL DEFAULT 0, + call_count INTEGER NOT NULL DEFAULT 1, + error_count INTEGER NOT NULL DEFAULT 0, + metadata JSONB, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); + +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_time + ON usage_events(tenant_id, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_resource_time + ON usage_events(tenant_id, resource_type, resource_name, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_type_time + ON usage_events(tenant_id, event_type, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_agent_time + ON usage_events(tenant_id, agent_id, event_time DESC); +CREATE INDEX IF NOT EXISTS idx_usage_events_tenant_channel_time + ON usage_events(tenant_id, channel, event_time DESC) + WHERE channel != ''; +CREATE UNIQUE INDEX IF NOT EXISTS idx_usage_events_trace_span_type_source + ON usage_events(trace_id, span_id, event_type, source) + WHERE trace_id IS NOT NULL AND span_id IS NOT NULL; + +CREATE TABLE IF NOT EXISTS usage_event_rollups ( + id UUID PRIMARY KEY, + tenant_id UUID NOT NULL REFERENCES tenants(id), + bucket_hour TIMESTAMPTZ NOT NULL, + event_type TEXT NOT NULL, + resource_type TEXT NOT NULL, + resource_name TEXT NOT NULL, + source TEXT NOT NULL DEFAULT '', + agent_id UUID REFERENCES agents(id) ON DELETE SET NULL, + channel TEXT NOT NULL DEFAULT '', + provider TEXT NOT NULL DEFAULT '', + model TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT '', + input_tokens BIGINT NOT NULL DEFAULT 0, + output_tokens BIGINT NOT NULL DEFAULT 0, + total_tokens BIGINT NOT NULL DEFAULT 0, + cost_usd DOUBLE PRECISION NOT NULL DEFAULT 0, + duration_ms INTEGER NOT NULL DEFAULT 0, + call_count INTEGER NOT NULL DEFAULT 0, + error_count INTEGER NOT NULL DEFAULT 0, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); + +CREATE UNIQUE INDEX IF NOT EXISTS idx_usage_event_rollups_unique + ON usage_event_rollups ( + tenant_id, + bucket_hour, + event_type, + resource_type, + resource_name, + source, + COALESCE(agent_id, '00000000-0000-0000-0000-000000000000'::uuid), + channel, + provider, + model, + status + ); +CREATE INDEX IF NOT EXISTS idx_usage_event_rollups_tenant_hour + ON usage_event_rollups(tenant_id, bucket_hour DESC); +CREATE INDEX IF NOT EXISTS idx_usage_event_rollups_resource_hour + ON usage_event_rollups(tenant_id, resource_type, resource_name, bucket_hour DESC); diff --git a/ui/web/src/i18n/locales/en/usage.json b/ui/web/src/i18n/locales/en/usage.json index 477b8608..45d35540 100644 --- a/ui/web/src/i18n/locales/en/usage.json +++ b/ui/web/src/i18n/locales/en/usage.json @@ -133,6 +133,36 @@ "avgDuration": "Avg Duration", "cost": "Cost" }, + "events": { + "title": "Resource Event Analytics", + "description": "Counts tool, skill, MCP, and runtime activity from traced executions.", + "tabs": { + "tool": "Tools", + "skill": "Skills", + "mcp_tool": "MCP", + "runtime_tool": "Runtime" + }, + "metrics": { + "calls": "Calls", + "errorRate": "Error rate", + "tokens": "Tokens", + "avgDuration": "Avg duration", + "cost": "Cost" + }, + "activeBuckets": "{{count}} active buckets", + "unknownSource": "Unknown", + "empty": "No event data for this selection.", + "table": { + "resource": "Resource", + "source": "Source", + "calls": "Calls", + "errors": "Errors", + "errorRate": "Error rate", + "avgDuration": "Avg duration", + "tokens": "Tokens", + "cost": "Cost" + } + }, "tooltip": { "date": "Date", "total": "Total", diff --git a/ui/web/src/i18n/locales/vi/usage.json b/ui/web/src/i18n/locales/vi/usage.json index dd2a9f45..fe01398d 100644 --- a/ui/web/src/i18n/locales/vi/usage.json +++ b/ui/web/src/i18n/locales/vi/usage.json @@ -133,6 +133,36 @@ "avgDuration": "TB thời lượng", "cost": "Chi phí" }, + "events": { + "title": "Phân tích sự kiện tài nguyên", + "description": "Đếm hoạt động tool, skill, MCP và runtime từ các execution có trace.", + "tabs": { + "tool": "Tool", + "skill": "Skill", + "mcp_tool": "MCP", + "runtime_tool": "Runtime" + }, + "metrics": { + "calls": "Lượt gọi", + "errorRate": "Tỷ lệ lỗi", + "tokens": "Token", + "avgDuration": "TB thời lượng", + "cost": "Chi phí" + }, + "activeBuckets": "{{count}} bucket có dữ liệu", + "unknownSource": "Không rõ", + "empty": "Không có dữ liệu event cho lựa chọn này.", + "table": { + "resource": "Tài nguyên", + "source": "Nguồn", + "calls": "Lượt gọi", + "errors": "Lỗi", + "errorRate": "Tỷ lệ lỗi", + "avgDuration": "TB thời lượng", + "tokens": "Token", + "cost": "Chi phí" + } + }, "tooltip": { "date": "Ngày", "total": "Tổng", diff --git a/ui/web/src/i18n/locales/zh/usage.json b/ui/web/src/i18n/locales/zh/usage.json index b200c31f..27d56a31 100644 --- a/ui/web/src/i18n/locales/zh/usage.json +++ b/ui/web/src/i18n/locales/zh/usage.json @@ -133,6 +133,36 @@ "avgDuration": "平均时长", "cost": "费用" }, + "events": { + "title": "资源事件分析", + "description": "统计带 trace 的 tool、skill、MCP 和 runtime 活动。", + "tabs": { + "tool": "工具", + "skill": "技能", + "mcp_tool": "MCP", + "runtime_tool": "Runtime" + }, + "metrics": { + "calls": "调用", + "errorRate": "错误率", + "tokens": "令牌", + "avgDuration": "平均时长", + "cost": "费用" + }, + "activeBuckets": "{{count}} 个活跃时间桶", + "unknownSource": "未知", + "empty": "当前选择没有事件数据。", + "table": { + "resource": "资源", + "source": "来源", + "calls": "调用", + "errors": "错误", + "errorRate": "错误率", + "avgDuration": "平均时长", + "tokens": "令牌", + "cost": "费用" + } + }, "tooltip": { "date": "日期", "total": "合计", diff --git a/ui/web/src/pages/usage/components/usage-event-analytics-panel.tsx b/ui/web/src/pages/usage/components/usage-event-analytics-panel.tsx new file mode 100644 index 00000000..186325c0 --- /dev/null +++ b/ui/web/src/pages/usage/components/usage-event-analytics-panel.tsx @@ -0,0 +1,138 @@ +import { useMemo, useState } from "react"; +import { useTranslation } from "react-i18next"; +import { Cable, Sparkles, Terminal, Wrench } from "lucide-react"; +import { Badge } from "@/components/ui/badge"; +import { Button } from "@/components/ui/button"; +import { formatCost, formatDuration, formatTokens } from "@/lib/format"; +import { cn } from "@/lib/utils"; +import { useUsageFilterContext } from "../context/usage-filter-context"; +import { useUsageEventAnalytics, type UsageEventResourceType } from "../hooks/use-usage-event-analytics"; + +const RESOURCE_TABS: Array<{ value: UsageEventResourceType; icon: typeof Wrench }> = [ + { value: "tool", icon: Wrench }, + { value: "skill", icon: Sparkles }, + { value: "mcp_tool", icon: Cable }, + { value: "runtime_tool", icon: Terminal }, +]; + +const EMPTY_SUMMARY = { calls: 0, errors: 0, input_tokens: 0, output_tokens: 0, total_tokens: 0, cost_usd: 0, avg_duration_ms: 0 }; + +export function UsageEventAnalyticsPanel() { + const { t } = useTranslation("usage"); + const { filters } = useUsageFilterContext(); + const [resourceType, setResourceType] = useState("tool"); + const { summary, rows, sourceRows, points, loading, error } = useUsageEventAnalytics(filters, resourceType); + const current = summary ?? EMPTY_SUMMARY; + const errorRate = current.calls > 0 ? (current.errors / current.calls) * 100 : 0; + const activeBuckets = useMemo(() => points.filter((point) => point.calls > 0).length, [points]); + const apiError = error instanceof Error ? error.message : error ? String(error) : null; + + return ( +
+
+
+

{t("analytics.events.title")}

+

{t("analytics.events.description")}

+
+
+ {RESOURCE_TABS.map((tab) => { + const Icon = tab.icon; + const active = resourceType === tab.value; + return ( + + ); + })} +
+
+ + {apiError ? ( +
+ {t("common:error", "Error")}: {apiError} +
+ ) : null} + +
+ + + + + +
+ +
+ {t("analytics.events.activeBuckets", { count: activeBuckets })} + {sourceRows.length > 0 ? · : null} + {sourceRows.map((row) => ( + + {row.key || t("analytics.events.unknownSource")}: {row.calls.toLocaleString()} + + ))} +
+ +
+ + + + + + + + + + + + + + + {loading ? ( + Array.from({ length: 4 }).map((_, idx) => ( + + + + )) + ) : rows.length === 0 ? ( + + + + ) : ( + rows.map((row) => { + const rowErrorRate = row.calls > 0 ? (row.errors / row.calls) * 100 : 0; + return ( + 0 && "bg-destructive/5")}> + + + + + + + + + + ); + }) + )} + +
{t("analytics.events.table.resource")}{t("analytics.events.table.source")}{t("analytics.events.table.calls")}{t("analytics.events.table.errors")}{t("analytics.events.table.errorRate")}{t("analytics.events.table.avgDuration")}{t("analytics.events.table.tokens")}{t("analytics.events.table.cost")}
{t("analytics.events.empty")}
{row.key || row.resource_name || "—"}{row.source || "—"}{row.calls.toLocaleString()}{row.errors.toLocaleString()}{rowErrorRate.toFixed(1)}%{formatDuration(row.avg_duration_ms)}{formatTokens(row.total_tokens)}{formatCost(row.cost_usd)}
+
+
+ ); +} + +function Metric({ label, value, loading }: { label: string; value: string; loading: boolean }) { + return ( +
+
{label}
+ {loading ?
:
{value}
} +
+ ); +} diff --git a/ui/web/src/pages/usage/hooks/use-usage-event-analytics.ts b/ui/web/src/pages/usage/hooks/use-usage-event-analytics.ts new file mode 100644 index 00000000..eaf676ad --- /dev/null +++ b/ui/web/src/pages/usage/hooks/use-usage-event-analytics.ts @@ -0,0 +1,115 @@ +import { useQuery } from "@tanstack/react-query"; +import { useHttp } from "@/hooks/use-ws"; +import type { UsageFilters } from "../context/usage-filter-context"; + +export type UsageEventResourceType = "tool" | "skill" | "mcp_tool" | "runtime_tool"; + +export interface UsageEventSummary { + calls: number; + errors: number; + input_tokens: number; + output_tokens: number; + total_tokens: number; + cost_usd: number; + avg_duration_ms: number; +} + +export interface UsageEventBreakdown { + key: string; + event_type: string; + resource_type: string; + resource_name: string; + source: string; + calls: number; + errors: number; + input_tokens: number; + output_tokens: number; + total_tokens: number; + cost_usd: number; + avg_duration_ms: number; +} + +export interface UsageEventTimeSeries { + bucket_time: string; + calls: number; + errors: number; + input_tokens: number; + output_tokens: number; + total_tokens: number; + cost_usd: number; + avg_duration_ms: number; +} + +function buildParams(filters: UsageFilters, resourceType: UsageEventResourceType, extra?: Record): Record { + const p: Record = { + from: filters.from, + to: filters.to, + resource_type: resourceType, + }; + if (filters.agentId) p.agent_id = filters.agentId; + if (filters.provider) p.provider = filters.provider; + if (filters.model) p.model = filters.model; + if (filters.channel) p.channel = filters.channel; + return { ...p, ...extra }; +} + +function filterKey(f: UsageFilters, resourceType: UsageEventResourceType) { + return [resourceType, f.from, f.to, f.agentId, f.provider, f.model, f.channel, f.granularity] as const; +} + +const QUERY_OPTS = { staleTime: 60_000, refetchOnWindowFocus: false } as const; + +export function useUsageEventAnalytics(filters: UsageFilters, resourceType: UsageEventResourceType) { + const http = useHttp(); + const fk = filterKey(filters, resourceType); + + const summaryQuery = useQuery({ + queryKey: ["usage", "events", "summary", ...fk], + queryFn: () => + http.get<{ summary: UsageEventSummary }>("/v1/usage/events/summary", buildParams(filters, resourceType)), + placeholderData: (prev) => prev, + ...QUERY_OPTS, + }); + + const breakdownQuery = useQuery({ + queryKey: ["usage", "events", "breakdown", "resource", ...fk], + queryFn: () => + http.get<{ rows: UsageEventBreakdown[] }>( + "/v1/usage/events/breakdown", + buildParams(filters, resourceType, { group_by: "resource", limit: "25" }), + ), + placeholderData: (prev) => prev, + ...QUERY_OPTS, + }); + + const sourceQuery = useQuery({ + queryKey: ["usage", "events", "breakdown", "source", ...fk], + queryFn: () => + http.get<{ rows: UsageEventBreakdown[] }>( + "/v1/usage/events/breakdown", + buildParams(filters, resourceType, { group_by: "source", limit: "10" }), + ), + placeholderData: (prev) => prev, + ...QUERY_OPTS, + }); + + const timeseriesQuery = useQuery({ + queryKey: ["usage", "events", "timeseries", ...fk], + queryFn: () => + http.get<{ points: UsageEventTimeSeries[] }>( + "/v1/usage/events/timeseries", + buildParams(filters, resourceType, { group_by: filters.granularity }), + ), + placeholderData: (prev) => prev, + ...QUERY_OPTS, + }); + + return { + summary: summaryQuery.data?.summary ?? null, + rows: breakdownQuery.data?.rows ?? [], + sourceRows: sourceQuery.data?.rows ?? [], + points: timeseriesQuery.data?.points ?? [], + loading: summaryQuery.isLoading || breakdownQuery.isLoading || sourceQuery.isLoading || timeseriesQuery.isLoading, + error: summaryQuery.error ?? breakdownQuery.error ?? sourceQuery.error ?? timeseriesQuery.error, + }; +} diff --git a/ui/web/src/pages/usage/usage-page.tsx b/ui/web/src/pages/usage/usage-page.tsx index 71a36584..b23b2b3e 100644 --- a/ui/web/src/pages/usage/usage-page.tsx +++ b/ui/web/src/pages/usage/usage-page.tsx @@ -22,6 +22,7 @@ import { DurationChart } from "./components/duration-chart"; import { KnowledgeChart } from "./components/knowledge-chart"; import { TopModelsTable } from "./components/top-models-table"; import { UsageCapsPanel } from "./components/usage-caps-panel"; +import { UsageEventAnalyticsPanel } from "./components/usage-event-analytics-panel"; const EMPTY_SUMMARY = { requests: 0, input_tokens: 0, output_tokens: 0, cost: 0, errors: 0, unique_users: 0, llm_calls: 0, tool_calls: 0, avg_duration_ms: 0 }; @@ -113,6 +114,10 @@ function AnalyticsDashboard() { + + + +