mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-08-21 08:26:20 +00:00
feat(usage): add event analytics dashboard
This commit is contained in:
@@ -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))
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
+132
-4
@@ -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
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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 = `
|
||||
|
||||
@@ -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
|
||||
-- ============================================================
|
||||
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -26,6 +26,7 @@ type Stores struct {
|
||||
Contacts ContactStore
|
||||
Activity ActivityStore
|
||||
Snapshots SnapshotStore
|
||||
UsageEvents UsageEventStore
|
||||
BrowserCookies BrowserCookieStore
|
||||
SecureCLI SecureCLIStore
|
||||
SecureCLIGrants SecureCLIAgentGrantStore
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -0,0 +1,2 @@
|
||||
DROP TABLE IF EXISTS usage_event_rollups;
|
||||
DROP TABLE IF EXISTS usage_events;
|
||||
@@ -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);
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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": "合计",
|
||||
|
||||
@@ -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<UsageEventResourceType>("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 (
|
||||
<section className="space-y-4 rounded-md border p-3 sm:p-4">
|
||||
<div className="flex flex-col gap-3 lg:flex-row lg:items-center lg:justify-between">
|
||||
<div>
|
||||
<h3 className="text-sm font-semibold">{t("analytics.events.title")}</h3>
|
||||
<p className="text-xs text-muted-foreground">{t("analytics.events.description")}</p>
|
||||
</div>
|
||||
<div className="grid grid-cols-2 gap-2 sm:flex">
|
||||
{RESOURCE_TABS.map((tab) => {
|
||||
const Icon = tab.icon;
|
||||
const active = resourceType === tab.value;
|
||||
return (
|
||||
<Button
|
||||
key={tab.value}
|
||||
type="button"
|
||||
variant={active ? "default" : "outline"}
|
||||
size="sm"
|
||||
className="gap-1.5 justify-start sm:justify-center"
|
||||
onClick={() => setResourceType(tab.value)}
|
||||
>
|
||||
<Icon className="h-3.5 w-3.5" />
|
||||
{t(`analytics.events.tabs.${tab.value}`)}
|
||||
</Button>
|
||||
);
|
||||
})}
|
||||
</div>
|
||||
</div>
|
||||
|
||||
{apiError ? (
|
||||
<div className="rounded-md border border-destructive/50 bg-destructive/10 p-3 text-sm text-destructive">
|
||||
{t("common:error", "Error")}: {apiError}
|
||||
</div>
|
||||
) : null}
|
||||
|
||||
<div className="grid gap-3 sm:grid-cols-2 lg:grid-cols-5">
|
||||
<Metric label={t("analytics.events.metrics.calls")} value={current.calls.toLocaleString()} loading={loading} />
|
||||
<Metric label={t("analytics.events.metrics.errorRate")} value={`${errorRate.toFixed(1)}%`} loading={loading} />
|
||||
<Metric label={t("analytics.events.metrics.tokens")} value={formatTokens(current.total_tokens)} loading={loading} />
|
||||
<Metric label={t("analytics.events.metrics.avgDuration")} value={formatDuration(current.avg_duration_ms)} loading={loading} />
|
||||
<Metric label={t("analytics.events.metrics.cost")} value={formatCost(current.cost_usd)} loading={loading} />
|
||||
</div>
|
||||
|
||||
<div className="flex flex-wrap items-center gap-2 text-xs text-muted-foreground">
|
||||
<span>{t("analytics.events.activeBuckets", { count: activeBuckets })}</span>
|
||||
{sourceRows.length > 0 ? <span>·</span> : null}
|
||||
{sourceRows.map((row) => (
|
||||
<Badge key={row.key} variant="outline" className="font-normal">
|
||||
{row.key || t("analytics.events.unknownSource")}: {row.calls.toLocaleString()}
|
||||
</Badge>
|
||||
))}
|
||||
</div>
|
||||
|
||||
<div className="overflow-x-auto">
|
||||
<table className="w-full min-w-[720px] text-sm">
|
||||
<thead>
|
||||
<tr className="border-b bg-muted/50 text-xs text-muted-foreground">
|
||||
<th className="px-3 py-2 text-left font-medium">{t("analytics.events.table.resource")}</th>
|
||||
<th className="px-3 py-2 text-left font-medium">{t("analytics.events.table.source")}</th>
|
||||
<th className="px-3 py-2 text-right font-medium">{t("analytics.events.table.calls")}</th>
|
||||
<th className="px-3 py-2 text-right font-medium">{t("analytics.events.table.errors")}</th>
|
||||
<th className="px-3 py-2 text-right font-medium">{t("analytics.events.table.errorRate")}</th>
|
||||
<th className="px-3 py-2 text-right font-medium">{t("analytics.events.table.avgDuration")}</th>
|
||||
<th className="px-3 py-2 text-right font-medium">{t("analytics.events.table.tokens")}</th>
|
||||
<th className="px-3 py-2 text-right font-medium">{t("analytics.events.table.cost")}</th>
|
||||
</tr>
|
||||
</thead>
|
||||
<tbody>
|
||||
{loading ? (
|
||||
Array.from({ length: 4 }).map((_, idx) => (
|
||||
<tr key={idx} className="border-b last:border-0">
|
||||
<td colSpan={8} className="px-3 py-2"><div className="h-7 animate-pulse rounded bg-muted" /></td>
|
||||
</tr>
|
||||
))
|
||||
) : rows.length === 0 ? (
|
||||
<tr>
|
||||
<td colSpan={8} className="px-3 py-8 text-center text-muted-foreground">{t("analytics.events.empty")}</td>
|
||||
</tr>
|
||||
) : (
|
||||
rows.map((row) => {
|
||||
const rowErrorRate = row.calls > 0 ? (row.errors / row.calls) * 100 : 0;
|
||||
return (
|
||||
<tr key={`${row.resource_type}:${row.key}:${row.source}`} className={cn("border-b last:border-0", row.errors > 0 && "bg-destructive/5")}>
|
||||
<td className="px-3 py-2 font-medium">{row.key || row.resource_name || "—"}</td>
|
||||
<td className="px-3 py-2 text-muted-foreground">{row.source || "—"}</td>
|
||||
<td className="px-3 py-2 text-right">{row.calls.toLocaleString()}</td>
|
||||
<td className="px-3 py-2 text-right text-muted-foreground">{row.errors.toLocaleString()}</td>
|
||||
<td className="px-3 py-2 text-right text-muted-foreground">{rowErrorRate.toFixed(1)}%</td>
|
||||
<td className="px-3 py-2 text-right text-muted-foreground">{formatDuration(row.avg_duration_ms)}</td>
|
||||
<td className="px-3 py-2 text-right text-muted-foreground">{formatTokens(row.total_tokens)}</td>
|
||||
<td className="px-3 py-2 text-right">{formatCost(row.cost_usd)}</td>
|
||||
</tr>
|
||||
);
|
||||
})
|
||||
)}
|
||||
</tbody>
|
||||
</table>
|
||||
</div>
|
||||
</section>
|
||||
);
|
||||
}
|
||||
|
||||
function Metric({ label, value, loading }: { label: string; value: string; loading: boolean }) {
|
||||
return (
|
||||
<div className="rounded-md border bg-card p-3">
|
||||
<div className="text-xs text-muted-foreground">{label}</div>
|
||||
{loading ? <div className="mt-2 h-6 animate-pulse rounded bg-muted" /> : <div className="mt-1 text-lg font-semibold">{value}</div>}
|
||||
</div>
|
||||
);
|
||||
}
|
||||
@@ -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<string, string>): Record<string, string> {
|
||||
const p: Record<string, string> = {
|
||||
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,
|
||||
};
|
||||
}
|
||||
@@ -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() {
|
||||
<UsageCapsPanel />
|
||||
</ErrorBoundary>
|
||||
|
||||
<ErrorBoundary>
|
||||
<UsageEventAnalyticsPanel />
|
||||
</ErrorBoundary>
|
||||
|
||||
<ErrorBoundary>
|
||||
<TokenAreaChart data={timeseries} loading={loading} granularity={filters.granularity} />
|
||||
</ErrorBoundary>
|
||||
|
||||
Reference in New Issue
Block a user