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