updates on the ai and agent and all thse things awith onboarding
This commit is contained in:
228
internal/ai/telemetry/insights.go
Normal file
228
internal/ai/telemetry/insights.go
Normal file
@@ -0,0 +1,228 @@
|
||||
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
|
||||
}
|
||||
146
internal/ai/telemetry/parse.go
Normal file
146
internal/ai/telemetry/parse.go
Normal file
@@ -0,0 +1,146 @@
|
||||
// Package telemetry records what AI_engine's agents actually do.
|
||||
//
|
||||
// AI_engine publishes two fire-and-forget subjects on plain NATS
|
||||
// (core/message_bus.py publish_telemetry):
|
||||
//
|
||||
// telemetry.task — after every task: agent, task id/type, status, error, duration
|
||||
// telemetry.agent — every ~5 s per agent: status, current task, counters
|
||||
//
|
||||
// Tasks become rows in aiagentruns (Postgres); agent state goes to Redis with a
|
||||
// short TTL. Both are observability only: nothing here may slow a booking or
|
||||
// fail a request, and a malformed event is dropped, never guessed at.
|
||||
package telemetry
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"strings"
|
||||
"time"
|
||||
"unicode/utf8"
|
||||
|
||||
"doormile/models"
|
||||
)
|
||||
|
||||
const (
|
||||
maxIDLen = 64
|
||||
maxErrorLen = 2000
|
||||
)
|
||||
|
||||
// ErrInvalid is returned for an event that cannot be recorded as-is.
|
||||
var ErrInvalid = errors.New("invalid telemetry event")
|
||||
|
||||
type taskEvent struct {
|
||||
TS string `json:"ts"`
|
||||
AgentID string `json:"agent_id"`
|
||||
TaskID string `json:"task_id"`
|
||||
TaskType string `json:"task_type"`
|
||||
Status string `json:"status"`
|
||||
Error *string `json:"error"`
|
||||
DurationMS float64 `json:"duration_ms"`
|
||||
}
|
||||
|
||||
// engineTimeLayouts are what Python's datetime.now().isoformat() produces:
|
||||
// a naive local time, with or without microseconds.
|
||||
var engineTimeLayouts = []string{"2006-01-02T15:04:05.999999", "2006-01-02T15:04:05"}
|
||||
|
||||
func parseEngineTime(s string) *time.Time {
|
||||
for _, layout := range engineTimeLayouts {
|
||||
if t, err := time.Parse(layout, s); err == nil {
|
||||
return &t
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// truncate cuts s to at most n bytes without splitting a UTF-8 rune.
|
||||
func truncate(s string, n int) string {
|
||||
if len(s) <= n {
|
||||
return s
|
||||
}
|
||||
s = s[:n]
|
||||
for !utf8.ValidString(s) {
|
||||
s = s[:len(s)-1]
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
// ParseTask turns a telemetry.task payload into a run row stamped with
|
||||
// receivedAt. It refuses — rather than repairs — an event with no agent id or
|
||||
// no status, or an id too long to be one: a run attributed to the wrong agent
|
||||
// is worse than a missing one.
|
||||
func ParseTask(body []byte, receivedAt time.Time) (models.AIAgentRun, error) {
|
||||
var ev taskEvent
|
||||
if err := json.Unmarshal(body, &ev); err != nil {
|
||||
return models.AIAgentRun{}, ErrInvalid
|
||||
}
|
||||
agent := strings.TrimSpace(ev.AgentID)
|
||||
status := strings.ToLower(strings.TrimSpace(ev.Status))
|
||||
if agent == "" || len(agent) > maxIDLen || status == "" || len(status) > 20 {
|
||||
return models.AIAgentRun{}, ErrInvalid
|
||||
}
|
||||
|
||||
run := models.AIAgentRun{
|
||||
Agentid: agent,
|
||||
Tasktype: truncate(strings.TrimSpace(ev.TaskType), maxIDLen),
|
||||
Status: status,
|
||||
Durationms: int(ev.DurationMS),
|
||||
Occurredat: parseEngineTime(ev.TS),
|
||||
Receivedat: receivedAt,
|
||||
}
|
||||
if run.Durationms < 0 {
|
||||
run.Durationms = 0
|
||||
}
|
||||
if id := strings.TrimSpace(ev.TaskID); id != "" {
|
||||
id = truncate(id, maxIDLen)
|
||||
run.Taskid = &id
|
||||
}
|
||||
if ev.Error != nil {
|
||||
run.Error = truncate(*ev.Error, maxErrorLen)
|
||||
}
|
||||
return run, nil
|
||||
}
|
||||
|
||||
// AgentState is an agent's latest telemetry.agent heartbeat, as kept in Redis.
|
||||
type AgentState struct {
|
||||
AgentID string `json:"agentid"`
|
||||
Status string `json:"status"`
|
||||
CurrentTask *string `json:"currenttask"`
|
||||
TasksCompleted int64 `json:"taskscompleted"`
|
||||
TasksFailed int64 `json:"tasksfailed"`
|
||||
LastSeenAt time.Time `json:"lastseenat"`
|
||||
}
|
||||
|
||||
type agentEvent struct {
|
||||
AgentID string `json:"agent_id"`
|
||||
Status string `json:"status"`
|
||||
CurrentTask *string `json:"current_task"`
|
||||
TasksCompleted int64 `json:"tasks_completed"`
|
||||
TasksFailed int64 `json:"tasks_failed"`
|
||||
}
|
||||
|
||||
// ParseAgent turns a telemetry.agent payload into the state kept for it.
|
||||
func ParseAgent(body []byte, seenAt time.Time) (AgentState, error) {
|
||||
var ev agentEvent
|
||||
if err := json.Unmarshal(body, &ev); err != nil {
|
||||
return AgentState{}, ErrInvalid
|
||||
}
|
||||
agent := strings.TrimSpace(ev.AgentID)
|
||||
if agent == "" || len(agent) > maxIDLen {
|
||||
return AgentState{}, ErrInvalid
|
||||
}
|
||||
st := AgentState{
|
||||
AgentID: agent,
|
||||
Status: truncate(strings.TrimSpace(ev.Status), 20),
|
||||
TasksCompleted: ev.TasksCompleted,
|
||||
TasksFailed: ev.TasksFailed,
|
||||
LastSeenAt: seenAt,
|
||||
}
|
||||
if ev.CurrentTask != nil {
|
||||
t := truncate(*ev.CurrentTask, maxIDLen)
|
||||
st.CurrentTask = &t
|
||||
}
|
||||
return st, nil
|
||||
}
|
||||
|
||||
// StateKey is the Redis key an agent's latest state lives under.
|
||||
func StateKey(agentID string) string { return "ai:agent:state:" + agentID }
|
||||
169
internal/ai/telemetry/recorder.go
Normal file
169
internal/ai/telemetry/recorder.go
Normal file
@@ -0,0 +1,169 @@
|
||||
package telemetry
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"doormile/models"
|
||||
"doormile/utils"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
"github.com/redis/go-redis/v9"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
const (
|
||||
// QueueGroup makes each event land on exactly one backend replica, so a
|
||||
// run is written once however many pods subscribe.
|
||||
QueueGroup = "doormile-backend-telemetry"
|
||||
|
||||
bufferSize = 2000
|
||||
flushEvery = 2 * time.Second
|
||||
flushAt = 200
|
||||
stateTTL = 5 * time.Minute
|
||||
RetentionDays = 30
|
||||
)
|
||||
|
||||
// Recorder buffers runs and writes them in batches. A NATS callback must never
|
||||
// block on Postgres: a slow database would back up the subscription and, in
|
||||
// nats.go, eventually mark it a slow consumer. So the callback only enqueues,
|
||||
// and a full buffer drops (counted, logged) rather than waits.
|
||||
type Recorder struct {
|
||||
db *gorm.DB
|
||||
rdb *redis.Client
|
||||
runs chan models.AIAgentRun
|
||||
now func() time.Time
|
||||
dropped atomic.Int64
|
||||
written atomic.Int64
|
||||
once sync.Once
|
||||
}
|
||||
|
||||
// NewRecorder builds a recorder. rdb may be nil: agent state is then not kept.
|
||||
func NewRecorder(db *gorm.DB, rdb *redis.Client) *Recorder {
|
||||
// time.Now, NOT utils.DBNow. DBNow returns IST digits labelled UTC, which
|
||||
// is right only for the legacy timestamp-WITHOUT-time-zone columns. This
|
||||
// table is created by AutoMigrate, so its columns are timestamptz, and a
|
||||
// DBNow value lands 5h30m in the future — caught by the Phase 4 end-to-end
|
||||
// run, where a run received at 20:57 IST read back as 02:27 next day.
|
||||
return &Recorder{db: db, rdb: rdb, runs: make(chan models.AIAgentRun, bufferSize), now: time.Now}
|
||||
}
|
||||
|
||||
// Enqueue offers a run to the writer without blocking. It reports whether the
|
||||
// run was accepted.
|
||||
func (r *Recorder) Enqueue(run models.AIAgentRun) bool {
|
||||
select {
|
||||
case r.runs <- run:
|
||||
return true
|
||||
default:
|
||||
if n := r.dropped.Add(1); n == 1 || n%500 == 0 {
|
||||
utils.Error("ai telemetry: run buffer full, dropping", "dropped_total", n)
|
||||
}
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// HandleTask is the telemetry.task callback.
|
||||
func (r *Recorder) HandleTask(body []byte) {
|
||||
run, err := ParseTask(body, r.now())
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
r.Enqueue(run)
|
||||
}
|
||||
|
||||
// HandleAgent is the telemetry.agent callback. Best effort: Redis down means
|
||||
// the Insights page shows no live state, never that a request fails.
|
||||
func (r *Recorder) HandleAgent(body []byte) {
|
||||
if r.rdb == nil {
|
||||
return
|
||||
}
|
||||
st, err := ParseAgent(body, r.now())
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
b, _ := json.Marshal(st)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
_ = r.rdb.Set(ctx, StateKey(st.AgentID), b, stateTTL).Err()
|
||||
}
|
||||
|
||||
// Flush writes whatever is buffered. Duplicates of an (agent, task id) pair —
|
||||
// a redelivery — are ignored by the unique index.
|
||||
func (r *Recorder) Flush(batch []models.AIAgentRun) {
|
||||
if len(batch) == 0 {
|
||||
return
|
||||
}
|
||||
if err := r.db.Clauses(clause.OnConflict{DoNothing: true}).CreateInBatches(batch, 200).Error; err != nil {
|
||||
utils.Error("ai telemetry: writing runs failed", "count", len(batch), "error", err.Error())
|
||||
return
|
||||
}
|
||||
r.written.Add(int64(len(batch)))
|
||||
}
|
||||
|
||||
func (r *Recorder) writeLoop() {
|
||||
ticker := time.NewTicker(flushEvery)
|
||||
defer ticker.Stop()
|
||||
batch := make([]models.AIAgentRun, 0, flushAt)
|
||||
for {
|
||||
select {
|
||||
case run := <-r.runs:
|
||||
batch = append(batch, run)
|
||||
if len(batch) >= flushAt {
|
||||
r.Flush(batch)
|
||||
batch = batch[:0]
|
||||
}
|
||||
case <-ticker.C:
|
||||
r.Flush(batch)
|
||||
batch = batch[:0]
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Prune deletes runs older than the retention window. Returns rows removed.
|
||||
func (r *Recorder) Prune() (int64, error) {
|
||||
cutoff := r.now().AddDate(0, 0, -RetentionDays)
|
||||
res := r.db.Where("receivedat < ?", cutoff).Delete(&models.AIAgentRun{})
|
||||
return res.RowsAffected, res.Error
|
||||
}
|
||||
|
||||
func (r *Recorder) pruneLoop() {
|
||||
for {
|
||||
if n, err := r.Prune(); err != nil {
|
||||
utils.Error("ai telemetry: pruning old runs failed", "error", err.Error())
|
||||
} else if n > 0 {
|
||||
utils.Info("ai telemetry: pruned old runs", "count", n)
|
||||
}
|
||||
time.Sleep(24 * time.Hour)
|
||||
}
|
||||
}
|
||||
|
||||
// Start subscribes to AI_engine's telemetry and starts the writer. A nil NATS
|
||||
// connection (NATS down, or not configured) logs and returns: the API serves
|
||||
// without it, and the Insights page says no telemetry is being received.
|
||||
func (r *Recorder) Start(nc *nats.Conn) {
|
||||
if nc == nil {
|
||||
utils.Info("ai telemetry: NATS not connected, agent runs will not be recorded")
|
||||
return
|
||||
}
|
||||
r.once.Do(func() {
|
||||
go r.writeLoop()
|
||||
go r.pruneLoop()
|
||||
if _, err := nc.QueueSubscribe("telemetry.task", QueueGroup, func(m *nats.Msg) { r.HandleTask(m.Data) }); err != nil {
|
||||
utils.Error("ai telemetry: subscribe telemetry.task failed", "error", err.Error())
|
||||
}
|
||||
if _, err := nc.QueueSubscribe("telemetry.agent", QueueGroup, func(m *nats.Msg) { r.HandleAgent(m.Data) }); err != nil {
|
||||
utils.Error("ai telemetry: subscribe telemetry.agent failed", "error", err.Error())
|
||||
}
|
||||
Receiving.Store(true)
|
||||
utils.Info("ai telemetry: recording agent runs", "queue_group", QueueGroup)
|
||||
})
|
||||
}
|
||||
|
||||
// Receiving reports whether this process subscribed to telemetry at boot. The
|
||||
// Insights endpoint returns it, so an empty page can say "not connected"
|
||||
// rather than "no runs".
|
||||
var Receiving atomic.Bool
|
||||
135
internal/ai/telemetry/store_integration_test.go
Normal file
135
internal/ai/telemetry/store_integration_test.go
Normal file
@@ -0,0 +1,135 @@
|
||||
package telemetry
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"doormile/internal/testpg"
|
||||
"doormile/models"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// Real SQL against a THROWAWAY Postgres (REGISTRY_TEST_DSN); skipped otherwise.
|
||||
// Drops and recreates aiagentruns and agent_decisions in its own schema.
|
||||
|
||||
func pgDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
dsn := os.Getenv("REGISTRY_TEST_DSN")
|
||||
if dsn == "" {
|
||||
t.Skip("REGISTRY_TEST_DSN not set; skipping Postgres integration test")
|
||||
}
|
||||
db := testpg.Open(t, dsn, "aitelemetry_test")
|
||||
all := []any{&models.AIAgentRun{}, &models.AgentDecision{}}
|
||||
if err := db.Migrator().DropTable(all...); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := db.AutoMigrate(all...); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return db
|
||||
}
|
||||
|
||||
func strp(s string) *string { return &s }
|
||||
|
||||
func TestPGFlushStoresRunsOnceAndAggregates(t *testing.T) {
|
||||
db := pgDB(t)
|
||||
r := &Recorder{db: db, now: func() time.Time { return now }}
|
||||
|
||||
batch := []models.AIAgentRun{
|
||||
{Agentid: "EXCEPTION_AGENT", Taskid: strp("t1"), Status: "completed", Durationms: 100, Receivedat: now},
|
||||
{Agentid: "EXCEPTION_AGENT", Taskid: strp("t2"), Status: "failed", Error: "boom", Durationms: 300, Receivedat: now},
|
||||
{Agentid: "DISPATCH_AGENT", Taskid: nil, Status: "completed", Durationms: 50, Receivedat: now},
|
||||
{Agentid: "DISPATCH_AGENT", Taskid: nil, Status: "completed", Durationms: 70, Receivedat: now},
|
||||
}
|
||||
r.Flush(batch)
|
||||
// A redelivered event (same agent + task id) is ignored, not double-counted.
|
||||
r.Flush([]models.AIAgentRun{{Agentid: "EXCEPTION_AGENT", Taskid: strp("t1"), Status: "completed", Durationms: 999, Receivedat: now}})
|
||||
|
||||
var n int64
|
||||
db.Model(&models.AIAgentRun{}).Count(&n)
|
||||
if n != 4 {
|
||||
t.Fatalf("stored %d runs, want 4 (redelivery ignored, null task ids both kept)", n)
|
||||
}
|
||||
|
||||
stats, err := RunStats(db, now.Add(-time.Hour))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s := SummariseRuns(stats)
|
||||
if s.Total != 4 || s.Failed != 1 || len(s.PerAgent) != 2 {
|
||||
t.Fatalf("summary %+v", s)
|
||||
}
|
||||
for _, a := range s.PerAgent {
|
||||
if a.Agentid == "EXCEPTION_AGENT" && (a.Runs != 2 || a.Failed != 1 || a.Avgdurationms != 200 || a.Lastrunat == nil) {
|
||||
t.Errorf("EXCEPTION_AGENT stats %+v", a)
|
||||
}
|
||||
}
|
||||
|
||||
// Outside the window, nothing.
|
||||
if stats, _ := RunStats(db, now.Add(time.Hour)); len(stats) != 0 {
|
||||
t.Errorf("runs before the window were counted: %+v", stats)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPGPruneRemovesOnlyExpiredRuns(t *testing.T) {
|
||||
db := pgDB(t)
|
||||
r := &Recorder{db: db, now: func() time.Time { return now }}
|
||||
r.Flush([]models.AIAgentRun{
|
||||
{Agentid: "A", Status: "completed", Receivedat: now.AddDate(0, 0, -(RetentionDays + 1))},
|
||||
{Agentid: "A", Status: "completed", Receivedat: now.AddDate(0, 0, -1)},
|
||||
})
|
||||
removed, err := r.Prune()
|
||||
if err != nil || removed != 1 {
|
||||
t.Fatalf("pruned %d, %v; want 1", removed, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPGDecisionsCountAndPage(t *testing.T) {
|
||||
db := pgDB(t)
|
||||
ok, pend := "success", (*string)(nil)
|
||||
b1 := uint64(501)
|
||||
for i, d := range []models.AgentDecision{
|
||||
{DecisionType: "miler_assignment", BookingID: &b1, Context: `{"rider":"secret"}`, Decision: `{"miler_id":8}`, Reasoning: "nearest", Outcome: &ok},
|
||||
{DecisionType: "miler_assignment", Context: `{}`, Decision: `{"miler_id":9}`, Reasoning: "load", Outcome: pend},
|
||||
{DecisionType: "stall_response", Context: `{}`, Decision: `{"action":"alert"}`, Reasoning: "stalled 12m", Outcome: pend},
|
||||
} {
|
||||
d.CreatedAt = now.Add(time.Duration(i) * time.Minute)
|
||||
if err := db.Create(&d).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
counts, err := DecisionCounts(db, now.Add(-time.Hour))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s := SummariseDecisions(counts)
|
||||
if s.Total != 3 || s.ByType[0].Decisiontype != "miler_assignment" || s.ByType[0].Outcomes["success"] != 1 || s.ByType[0].Outcomes["pending"] != 1 {
|
||||
t.Fatalf("decision summary %+v", s)
|
||||
}
|
||||
|
||||
page, err := RecentDecisions(db, "", 0, 2)
|
||||
if err != nil || len(page) != 2 || page[0].Decisiontype != "stall_response" {
|
||||
t.Fatalf("first page %+v, %v", page, err)
|
||||
}
|
||||
var dec map[string]any
|
||||
if json.Unmarshal(page[0].Decision, &dec) != nil || dec["action"] != "alert" {
|
||||
t.Errorf("decision jsonb not returned as JSON: %s", page[0].Decision)
|
||||
}
|
||||
next, err := RecentDecisions(db, "", page[1].ID, 2)
|
||||
if err != nil || len(next) != 1 || next[0].Bookingid == nil || *next[0].Bookingid != 501 {
|
||||
t.Fatalf("second page %+v, %v", next, err)
|
||||
}
|
||||
only, _ := RecentDecisions(db, "stall_response", 0, 10)
|
||||
if len(only) != 1 {
|
||||
t.Errorf("type filter returned %d rows", len(only))
|
||||
}
|
||||
b, _ := json.Marshal(next)
|
||||
if strings.Contains(string(b), "secret") {
|
||||
t.Error("the context column (rider data) leaked into the decisions list")
|
||||
}
|
||||
}
|
||||
177
internal/ai/telemetry/telemetry_test.go
Normal file
177
internal/ai/telemetry/telemetry_test.go
Normal file
@@ -0,0 +1,177 @@
|
||||
package telemetry
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
"unicode/utf8"
|
||||
|
||||
"doormile/models"
|
||||
)
|
||||
|
||||
var now = time.Date(2026, 9, 29, 18, 0, 0, 0, time.UTC)
|
||||
|
||||
// The exact shape core/agent.py publishes after a task.
|
||||
const engineTask = `{"kind":"task","ts":"2026-09-29T17:59:58.123456","agent_id":"EXCEPTION_AGENT",
|
||||
"task_id":"3f2c9a","task_type":"handle_stall","status":"completed","error":null,"duration_ms":412}`
|
||||
|
||||
func TestParseTaskReadsTheEngineShape(t *testing.T) {
|
||||
run, err := ParseTask([]byte(engineTask), now)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if run.Agentid != "EXCEPTION_AGENT" || run.Tasktype != "handle_stall" || run.Status != "completed" || run.Durationms != 412 {
|
||||
t.Errorf("parsed %+v", run)
|
||||
}
|
||||
if run.Taskid == nil || *run.Taskid != "3f2c9a" {
|
||||
t.Errorf("task id = %v", run.Taskid)
|
||||
}
|
||||
if !run.Receivedat.Equal(now) {
|
||||
t.Errorf("receivedat = %v, want the backend's clock", run.Receivedat)
|
||||
}
|
||||
if run.Occurredat == nil || run.Occurredat.Format("15:04:05") != "17:59:58" {
|
||||
t.Errorf("occurredat = %v", run.Occurredat)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseTaskKeepsAFailureAndItsError(t *testing.T) {
|
||||
run, err := ParseTask([]byte(`{"agent_id":"DISPATCH_AGENT","task_id":"x","status":"FAILED","error":"boom","duration_ms":5}`), now)
|
||||
if err != nil || run.Status != "failed" || run.Error != "boom" {
|
||||
t.Fatalf("got %+v, %v", run, err)
|
||||
}
|
||||
}
|
||||
|
||||
// A run attributed to the wrong agent is worse than a missing one.
|
||||
func TestParseTaskRefusesWhatItCannotAttribute(t *testing.T) {
|
||||
cases := map[string]string{
|
||||
"not json": `{nope`,
|
||||
"no agent": `{"status":"completed"}`,
|
||||
"blank agent": `{"agent_id":" ","status":"completed"}`,
|
||||
"no status": `{"agent_id":"A"}`,
|
||||
"agent id too long": `{"agent_id":"` + strings.Repeat("A", 65) + `","status":"completed"}`,
|
||||
}
|
||||
for name, body := range cases {
|
||||
if _, err := ParseTask([]byte(body), now); err != ErrInvalid {
|
||||
t.Errorf("%s: want ErrInvalid, got %v", name, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseTaskCleansUpEdgeValues(t *testing.T) {
|
||||
long := strings.Repeat("é", 1500) // 3000 bytes of two-byte runes
|
||||
run, err := ParseTask([]byte(`{"agent_id":"A","status":"completed","duration_ms":-40,"ts":"garbage","error":"`+long+`"}`), now)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if run.Durationms != 0 {
|
||||
t.Errorf("negative duration kept: %d", run.Durationms)
|
||||
}
|
||||
if run.Taskid != nil {
|
||||
t.Error("an absent task id must be stored as null, not an empty string")
|
||||
}
|
||||
if run.Occurredat != nil {
|
||||
t.Error("an unparseable engine time must be null, not guessed")
|
||||
}
|
||||
if len(run.Error) > maxErrorLen || !utf8.ValidString(run.Error) {
|
||||
t.Errorf("error not truncated safely: %d bytes, valid=%v", len(run.Error), utf8.ValidString(run.Error))
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseAgent(t *testing.T) {
|
||||
st, err := ParseAgent([]byte(`{"kind":"agent","agent_id":"JARVIS","status":"idle","current_task":null,"tasks_completed":12,"tasks_failed":1}`), now)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if st.AgentID != "JARVIS" || st.Status != "idle" || st.TasksCompleted != 12 || st.TasksFailed != 1 || !st.LastSeenAt.Equal(now) {
|
||||
t.Errorf("parsed %+v", st)
|
||||
}
|
||||
if _, err := ParseAgent([]byte(`{"status":"idle"}`), now); err != ErrInvalid {
|
||||
t.Error("an agent event with no id was accepted")
|
||||
}
|
||||
}
|
||||
|
||||
// The NATS callback must never block: a full buffer drops and counts.
|
||||
func TestEnqueueDropsInsteadOfBlocking(t *testing.T) {
|
||||
r := &Recorder{runs: make(chan models.AIAgentRun, 1), now: func() time.Time { return now }}
|
||||
if !r.Enqueue(models.AIAgentRun{Agentid: "A"}) {
|
||||
t.Fatal("first run refused")
|
||||
}
|
||||
done := make(chan bool)
|
||||
go func() { done <- r.Enqueue(models.AIAgentRun{Agentid: "B"}) }()
|
||||
select {
|
||||
case ok := <-done:
|
||||
if ok || r.dropped.Load() != 1 {
|
||||
t.Errorf("full buffer: accepted=%v dropped=%d", ok, r.dropped.Load())
|
||||
}
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("Enqueue blocked on a full buffer")
|
||||
}
|
||||
}
|
||||
|
||||
func TestHandlersIgnoreBadInputAndMissingRedis(t *testing.T) {
|
||||
r := &Recorder{runs: make(chan models.AIAgentRun, 4), now: func() time.Time { return now }}
|
||||
r.HandleTask([]byte(`{bad`))
|
||||
r.HandleAgent([]byte(`{"agent_id":"A","status":"idle"}`)) // rdb nil: must not panic
|
||||
if len(r.runs) != 0 {
|
||||
t.Error("an invalid task event was enqueued")
|
||||
}
|
||||
r.HandleTask([]byte(engineTask))
|
||||
if len(r.runs) != 1 {
|
||||
t.Error("a valid task event was not enqueued")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSummariseRuns(t *testing.T) {
|
||||
s := SummariseRuns([]AgentRunStats{
|
||||
{Agentid: "B", Runs: 3, Failed: 1},
|
||||
{Agentid: "A", Runs: 10, Failed: 0},
|
||||
{Agentid: "C", Runs: 3, Failed: 2},
|
||||
})
|
||||
if s.Total != 16 || s.Failed != 3 {
|
||||
t.Errorf("totals %d/%d", s.Total, s.Failed)
|
||||
}
|
||||
var order []string
|
||||
for _, a := range s.PerAgent {
|
||||
order = append(order, a.Agentid)
|
||||
}
|
||||
if strings.Join(order, ",") != "A,B,C" {
|
||||
t.Errorf("order %v, want busiest first then by id", order)
|
||||
}
|
||||
if empty := SummariseRuns(nil); empty.PerAgent == nil || empty.Total != 0 {
|
||||
t.Error("no runs must summarise to an empty list, not null")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSummariseDecisions(t *testing.T) {
|
||||
s := SummariseDecisions([]DecisionCount{
|
||||
{Decisiontype: "miler_assignment", Outcome: "success", Count: 7},
|
||||
{Decisiontype: "stall_response", Outcome: "pending", Count: 2},
|
||||
{Decisiontype: "miler_assignment", Outcome: "pending", Count: 3},
|
||||
})
|
||||
if s.Total != 12 || len(s.ByType) != 2 {
|
||||
t.Fatalf("summary %+v", s)
|
||||
}
|
||||
first := s.ByType[0]
|
||||
if first.Decisiontype != "miler_assignment" || first.Total != 10 || first.Outcomes["success"] != 7 || first.Outcomes["pending"] != 3 {
|
||||
t.Errorf("miler_assignment rolled up wrong: %+v", first)
|
||||
}
|
||||
}
|
||||
|
||||
func TestClampDays(t *testing.T) {
|
||||
for in, want := range map[int]int{0: 7, -3: 7, 1: 1, 7: 7, 30: 30, 90: RetentionDays} {
|
||||
if got := ClampDays(in); got != want {
|
||||
t.Errorf("ClampDays(%d) = %d, want %d", in, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// aiagentruns is created by AutoMigrate, so its columns are timestamptz. The
|
||||
// recorder must stamp a true instant; utils.DBNow (IST digits labelled UTC) is
|
||||
// 5h30m off as an instant and made the Phase 4 end-to-end run read a 20:57 IST
|
||||
// run back as 02:27 the next day.
|
||||
func TestRecorderStampsARealInstant(t *testing.T) {
|
||||
r := NewRecorder(nil, nil)
|
||||
if d := r.now().Sub(time.Now()); d > time.Minute || d < -time.Minute {
|
||||
t.Fatalf("recorder clock is %v off real time; it must not use utils.DBNow", d)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user