147 lines
4.3 KiB
Go
147 lines
4.3 KiB
Go
// 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 }
|