package telemetry import ( "context" "encoding/json" "sort" "time" "github.com/redis/go-redis/v9" "gorm.io/gorm" ) // AgentRunStats is one agent's runs over the window. type AgentRunStats struct { Agentid string `json:"agentid"` Runs int64 `json:"runs"` Failed int64 `json:"failed"` Avgdurationms float64 `json:"avgdurationms"` Lastrunat *time.Time `json:"lastrunat"` } // DecisionCount is decisions of one type with one outcome over the window. // Outcome is "pending" while none has been recorded. type DecisionCount struct { Decisiontype string `json:"decisiontype"` Outcome string `json:"outcome"` Count int64 `json:"count"` } // DecisionTypeStats rolls DecisionCount rows up per type. type DecisionTypeStats struct { Decisiontype string `json:"decisiontype"` Total int64 `json:"total"` Outcomes map[string]int64 `json:"outcomes"` } // Insights is everything the Insights tab shows. type Insights struct { Days int `json:"days"` Since time.Time `json:"since"` // Receiving is whether this backend is subscribed to AI_engine telemetry. // False means an empty run list is "not connected", not "no activity". Receiving bool `json:"receiving"` Runs RunsSummary `json:"runs"` Decisions DecisionSummary `json:"decisions"` Live []AgentState `json:"live"` } type RunsSummary struct { Total int64 `json:"total"` Failed int64 `json:"failed"` PerAgent []AgentRunStats `json:"peragent"` } type DecisionSummary struct { Total int64 `json:"total"` ByType []DecisionTypeStats `json:"bytype"` } // ClampDays keeps the window to what the table retains. func ClampDays(days int) int { switch { case days < 1: return 7 case days > RetentionDays: return RetentionDays default: return days } } // SummariseRuns totals per-agent rows, busiest agent first. func SummariseRuns(rows []AgentRunStats) RunsSummary { out := RunsSummary{PerAgent: append([]AgentRunStats{}, rows...)} for _, r := range rows { out.Total += r.Runs out.Failed += r.Failed } sort.SliceStable(out.PerAgent, func(i, j int) bool { if out.PerAgent[i].Runs != out.PerAgent[j].Runs { return out.PerAgent[i].Runs > out.PerAgent[j].Runs } return out.PerAgent[i].Agentid < out.PerAgent[j].Agentid }) return out } // SummariseDecisions rolls (type, outcome) counts up per type, largest first. func SummariseDecisions(rows []DecisionCount) DecisionSummary { byType := map[string]*DecisionTypeStats{} var order []string var total int64 for _, r := range rows { st, ok := byType[r.Decisiontype] if !ok { st = &DecisionTypeStats{Decisiontype: r.Decisiontype, Outcomes: map[string]int64{}} byType[r.Decisiontype] = st order = append(order, r.Decisiontype) } st.Total += r.Count st.Outcomes[r.Outcome] += r.Count total += r.Count } out := DecisionSummary{Total: total, ByType: make([]DecisionTypeStats, 0, len(order))} for _, k := range order { out.ByType = append(out.ByType, *byType[k]) } sort.SliceStable(out.ByType, func(i, j int) bool { if out.ByType[i].Total != out.ByType[j].Total { return out.ByType[i].Total > out.ByType[j].Total } return out.ByType[i].Decisiontype < out.ByType[j].Decisiontype }) return out } // RunStats reads per-agent run statistics since the cutoff. A fixed single // grouped query, whatever the volume. func RunStats(db *gorm.DB, since time.Time) ([]AgentRunStats, error) { var rows []AgentRunStats err := db.Table("aiagentruns"). Select(`agentid, COUNT(*) AS runs, COUNT(*) FILTER (WHERE status <> 'completed') AS failed, COALESCE(AVG(durationms), 0) AS avgdurationms, MAX(receivedat) AS lastrunat`). Where("receivedat >= ?", since). Group("agentid"). Scan(&rows).Error return rows, err } // DecisionCounts reads agent_decisions grouped by type and outcome. The table // is written by the decision engine through POST /internal/agent-decisions. func DecisionCounts(db *gorm.DB, since time.Time) ([]DecisionCount, error) { var rows []DecisionCount err := db.Table("agent_decisions"). Select("decision_type AS decisiontype, COALESCE(outcome, 'pending') AS outcome, COUNT(*) AS count"). Where("created_at >= ?", since). Group("decision_type, COALESCE(outcome, 'pending')"). Scan(&rows).Error return rows, err } // LiveStates reads the latest heartbeat of each agent that has one. Missing // keys (an agent silent for over five minutes) are simply absent. func LiveStates(rdb *redis.Client, agentIDs []string) []AgentState { out := []AgentState{} if rdb == nil || len(agentIDs) == 0 { return out } keys := make([]string, len(agentIDs)) for i, id := range agentIDs { keys[i] = StateKey(id) } ctx, cancel := context.WithTimeout(context.Background(), time.Second) defer cancel() vals, err := rdb.MGet(ctx, keys...).Result() if err != nil { return out } for _, v := range vals { s, ok := v.(string) if !ok { continue } var st AgentState if json.Unmarshal([]byte(s), &st) == nil && st.AgentID != "" { out = append(out, st) } } return out } // DecisionRow is one recent decision, as the Insights list shows it. type DecisionRow struct { ID uint64 `json:"id"` Decisiontype string `json:"decisiontype"` Bookingid *uint64 `json:"bookingid"` Decision json.RawMessage `json:"decision"` Reasoning string `json:"reasoning"` Outcome *string `json:"outcome"` Createdat time.Time `json:"createdat"` } type decisionScan struct { ID uint64 Decisiontype string Bookingid *uint64 Decision string Reasoning string Outcome *string Createdat time.Time } // RecentDecisions pages agent_decisions newest first. beforeID (0 = start) // is a keyset cursor, so a page is stable while new rows arrive. The context // column is deliberately not returned: it can be large and holds rider data. func RecentDecisions(db *gorm.DB, decisionType string, beforeID uint64, limit int) ([]DecisionRow, error) { if limit <= 0 || limit > 100 { limit = 25 } q := db.Table("agent_decisions"). Select("id, decision_type AS decisiontype, booking_id AS bookingid, COALESCE(decision::text, 'null') AS decision, reasoning, outcome, created_at AS createdat"). Order("id DESC").Limit(limit) if decisionType != "" { q = q.Where("decision_type = ?", decisionType) } if beforeID > 0 { q = q.Where("id < ?", beforeID) } var rows []decisionScan if err := q.Scan(&rows).Error; err != nil { return nil, err } out := make([]DecisionRow, 0, len(rows)) for _, r := range rows { raw := json.RawMessage(r.Decision) if !json.Valid(raw) { raw = json.RawMessage("null") } out = append(out, DecisionRow{ ID: r.ID, Decisiontype: r.Decisiontype, Bookingid: r.Bookingid, Decision: raw, Reasoning: truncate(r.Reasoning, 1000), Outcome: r.Outcome, Createdat: r.Createdat, }) } return out, nil }