// 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 }