Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
182 changes: 182 additions & 0 deletions gateway/internal/adminapi/agentruns.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,182 @@
package adminapi

import (
"net/http"
"sort"
"time"
)

// ─── /_plugin/agents/:name/runs ──────────────────────────────────────
//
// The agent's runs in the window, one row per run-id, most recent
// activity first. Backs the "Recent runs" table on the AgentDetail
// page. Before this endpoint the page derived that table from
// `/histogram/cost?dimension=run-id`, which is not agent-scoped and
// orders by spend — so an operator selecting canvas-agent saw every
// agent's runs, ranked by cost, with the run they had just fired
// nowhere near the top.
//
// Same strategy as the other rollups: page the window out of
// Bifrost's /api/logs filtered by `metadata.agent-name`, then group
// by `metadata.run-id` in Go. Rows without a run-id are dropped (a
// bare call outside any run has nothing to link to). Same 200k-row
// ceiling as the rest of observability.go.

// AgentRunSummary is one row of /_plugin/agents/:name/runs.
type AgentRunSummary struct {
RunID string `json:"run_id"`
// UserID is `metadata.user-id` from the run's first row that
// carries one — the same key the People pages are keyed on, so
// the dashboard can link straight to /people/:id. Empty when
// no row was stamped.
UserID string `json:"user_id,omitempty"`
// Models the run called, most-used first (ties by name). A run
// usually has one; a "+N" affordance in the UI covers the rest.
Models []string `json:"models"`
TotalCost float64 `json:"total_cost"`
TotalTokens int64 `json:"total_tokens"`
RequestCount int64 `json:"request_count"`
FirstSeen string `json:"first_seen,omitempty"`
LastSeen string `json:"last_seen,omitempty"`
}

// AgentRunsResponse is the envelope for /_plugin/agents/:name/runs.
// `total` is the run count in the window before ?limit=/?offset=
// paging, so the UI can say "showing 50 of 120".
type AgentRunsResponse struct {
AgentName string `json:"agent_name"`
Window string `json:"window"`
Total int `json:"total"`
Runs []AgentRunSummary `json:"runs"`
}

func (h *observabilityHandlers) agentRuns(w http.ResponseWriter, r *http.Request, name string) {
if r.Method != http.MethodGet {
methodNotAllowed(w, http.MethodGet)
return
}
window, start, end, ok := parseWindow(w, r)
if !ok {
return
}
limit, offset, ok := parsePagination(w, r)
if !ok {
return
}
logs, err := h.logs.searchAll(r.Context(), searchOpts{
StartTime: &start,
EndTime: &end,
Metadata: map[string]string{"agent-name": name},
}, 1000, 200_000)
if err != nil {
writeUpstreamError(w, err, "agents.runs")
return
}

runs := summarizeRuns(logs)
total := len(runs)
if offset > len(runs) {
offset = len(runs)
}
runs = runs[offset:]
if len(runs) > limit {
runs = runs[:limit]
}
writeJSON(w, http.StatusOK, AgentRunsResponse{
AgentName: name,
Window: window,
Total: total,
Runs: runs,
})
}

// summarizeRuns groups log rows by `metadata.run-id` and returns one
// summary per run, sorted by last activity, newest first (ties by
// run-id so the order is stable across polls). Rows with no run-id
// are skipped.
func summarizeRuns(logs []logstoreLog) []AgentRunSummary {
type agg struct {
user string
models map[string]int64
cost float64
tokens int64
count int64
first time.Time
firstSeen string
last time.Time
lastSeen string
}
byRun := map[string]*agg{}
for _, l := range logs {
runID := l.Metadata["run-id"]
if runID == "" {
continue
}
a, ok := byRun[runID]
if !ok {
a = &agg{models: map[string]int64{}}
byRun[runID] = a
}
if a.user == "" {
a.user = l.Metadata["user-id"]
}
if l.Model != "" {
a.models[l.Model]++
}
a.cost += l.Cost
a.tokens += l.tokens()
a.count++
// Compare as times, not strings: RFC3339Nano trims trailing
// zeros, so "…:00Z" sorts after "…:00.5Z" lexicographically.
ts := parseLogTimestamp(l.Timestamp)
if a.firstSeen == "" || ts.Before(a.first) {
a.first, a.firstSeen = ts, l.Timestamp
}
if a.lastSeen == "" || ts.After(a.last) {
a.last, a.lastSeen = ts, l.Timestamp
}
}

out := make([]AgentRunSummary, 0, len(byRun))
for id, a := range byRun {
models := make([]string, 0, len(a.models))
for m := range a.models {
models = append(models, m)
}
sort.Slice(models, func(i, j int) bool {
if a.models[models[i]] != a.models[models[j]] {
return a.models[models[i]] > a.models[models[j]]
}
return models[i] < models[j]
})
out = append(out, AgentRunSummary{
RunID: id,
UserID: a.user,
Models: models,
TotalCost: a.cost,
TotalTokens: a.tokens,
RequestCount: a.count,
FirstSeen: a.firstSeen,
LastSeen: a.lastSeen,
})
}
sort.Slice(out, func(i, j int) bool {
ti, tj := parseLogTimestamp(out[i].LastSeen), parseLogTimestamp(out[j].LastSeen)
if !ti.Equal(tj) {
return ti.After(tj)
}
return out[i].RunID < out[j].RunID
})
return out
}

// parseLogTimestamp reads a Bifrost row timestamp. Unparseable
// values sort as the zero time — oldest — rather than being dropped,
// so a malformed row still counts toward the run's totals.
func parseLogTimestamp(s string) time.Time {
t, err := time.Parse(time.RFC3339Nano, s)
if err != nil {
return time.Time{}
}
return t
}
160 changes: 160 additions & 0 deletions gateway/internal/adminapi/agentruns_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,160 @@
package adminapi

import (
"net/http"
"testing"
"time"
)

// /_plugin/agents/:name/runs — the per-run rollup behind the
// AgentDetail "Recent runs" table. Same fakeBifrost harness as
// observability_test.go; phase7Logs gives coder two runs (r1 by
// u_alice on haiku, r3 by u_bob on gpt-4o-mini) and web-search one
// (r2), which must not leak into coder's list.

func TestAgentRuns_ScopedAndNewestFirst(t *testing.T) {
now := time.Now().UTC()
srv := newObservabilityTestServer(t, newFakeBifrost(t, phase7Logs(now)))

out := decodeOK[AgentRunsResponse](t, bearerGet(t, srv, "/_plugin/agents/coder/runs?window=1h"))
if out.AgentName != "coder" || out.Window != "1h" || out.Total != 2 || len(out.Runs) != 2 {
t.Fatalf("envelope: %+v", out)
}
// r3's last call (4m ago) is newer than r1's (28m ago).
r3, r1 := out.Runs[0], out.Runs[1]
if r3.RunID != "r3" || r1.RunID != "r1" {
t.Fatalf("order: %s, %s", r3.RunID, r1.RunID)
}
if r3.UserID != "u_bob" || r3.RequestCount != 2 || r3.TotalTokens != 100 ||
len(r3.Models) != 1 || r3.Models[0] != "gpt-4o-mini" {
t.Errorf("r3: %+v", r3)
}
if r3.TotalCost < 0.039 || r3.TotalCost > 0.041 {
t.Errorf("r3 cost: %v", r3.TotalCost)
}
if r1.UserID != "u_alice" || r1.RequestCount != 2 || r1.TotalTokens != 360 ||
len(r1.Models) != 1 || r1.Models[0] != "claude-3-5-haiku" {
t.Errorf("r1: %+v", r1)
}
if r1.TotalCost < 0.149 || r1.TotalCost > 0.151 {
t.Errorf("r1 cost: %v", r1.TotalCost)
}
base := now.Truncate(10 * time.Minute)
if r1.FirstSeen != base.Add(-29*time.Minute).Format(time.RFC3339Nano) ||
r1.LastSeen != base.Add(-28*time.Minute).Format(time.RFC3339Nano) {
t.Errorf("r1 seen: %s .. %s", r1.FirstSeen, r1.LastSeen)
}
for _, r := range out.Runs {
if r.RunID == "r2" {
t.Fatal("web-search's run leaked into coder's list")
}
}

none := decodeOK[AgentRunsResponse](t, bearerGet(t, srv, "/_plugin/agents/ghost/runs"))
if none.Total != 0 || len(none.Runs) != 0 || none.Window != "24h" {
t.Errorf("unknown agent: %+v", none)
}
}

// Models are ordered by call count; the user-id comes from the first
// row that carries one; rows with no run-id are skipped rather than
// crashing the rollup; sub-second timestamps compare as times.
func TestAgentRuns_ModelsUserAndOrphans(t *testing.T) {
now := time.Now().UTC()
ts := func(d time.Duration) string { return now.Add(-d).Format(time.RFC3339Nano) }
md := func(run, user string) map[string]string {
m := map[string]string{"agent-name": "coder"}
if run != "" {
m["run-id"] = run
}
if user != "" {
m["user-id"] = user
}
return m
}
logs := []fakeLog{
{ID: "1", Timestamp: ts(3 * time.Minute), Model: "big", Cost: 1, Metadata: md("rx", "")},
{ID: "2", Timestamp: ts(2 * time.Minute), Model: "small", Cost: 1, Metadata: md("rx", "u_carol")},
{ID: "3", Timestamp: ts(1 * time.Minute), Model: "small", Cost: 1, Metadata: md("rx", "u_dave")},
{ID: "4", Timestamp: ts(30 * time.Second), Model: "", Cost: 1, Metadata: md("rx", "")},
// No run-id: excluded from every run, must not 500.
{ID: "5", Timestamp: ts(10 * time.Second), Model: "big", Cost: 9, Metadata: md("", "u_carol")},
// A second run whose only call is a whole second older than
// rx's newest but has a *lexicographically* larger timestamp
// ("…:SSZ" vs "…:SS.5Z"). Must sort after rx.
{ID: "6", Timestamp: now.Add(-31 * time.Second).Truncate(time.Second).Format(time.RFC3339Nano),
Model: "big", Cost: 1, Metadata: md("ry", "u_erin")},
}
srv := newObservabilityTestServer(t, newFakeBifrost(t, logs))

out := decodeOK[AgentRunsResponse](t, bearerGet(t, srv, "/_plugin/agents/coder/runs?window=1h"))
if out.Total != 2 || len(out.Runs) != 2 {
t.Fatalf("want 2 runs, got %+v", out)
}
rx := out.Runs[0]
if rx.RunID != "rx" || out.Runs[1].RunID != "ry" {
t.Fatalf("order: %s, %s", out.Runs[0].RunID, out.Runs[1].RunID)
}
if rx.UserID != "u_carol" {
t.Errorf("user: %q", rx.UserID)
}
if len(rx.Models) != 2 || rx.Models[0] != "small" || rx.Models[1] != "big" {
t.Errorf("models: %v", rx.Models)
}
if rx.RequestCount != 4 || rx.TotalCost != 4 {
t.Errorf("totals: %+v", rx)
}
if rx.FirstSeen != ts(3*time.Minute) || rx.LastSeen != ts(30*time.Second) {
t.Errorf("seen: %s .. %s", rx.FirstSeen, rx.LastSeen)
}
}

func TestAgentRuns_Pagination(t *testing.T) {
now := time.Now().UTC()
srv := newObservabilityTestServer(t, newFakeBifrost(t, phase7Logs(now)))

page := decodeOK[AgentRunsResponse](t, bearerGet(t, srv, "/_plugin/agents/coder/runs?window=1h&limit=1"))
if page.Total != 2 || len(page.Runs) != 1 || page.Runs[0].RunID != "r3" {
t.Errorf("limit=1: %+v", page)
}
next := decodeOK[AgentRunsResponse](t, bearerGet(t, srv, "/_plugin/agents/coder/runs?window=1h&limit=1&offset=1"))
if next.Total != 2 || len(next.Runs) != 1 || next.Runs[0].RunID != "r1" {
t.Errorf("offset=1: %+v", next)
}
past := decodeOK[AgentRunsResponse](t, bearerGet(t, srv, "/_plugin/agents/coder/runs?window=1h&offset=9"))
if past.Total != 2 || len(past.Runs) != 0 {
t.Errorf("offset past end: %+v", past)
}

resp := bearerGet(t, srv, "/_plugin/agents/coder/runs?limit=0")
resp.Body.Close()
if resp.StatusCode != http.StatusBadRequest {
t.Errorf("limit=0: want 400, got %d", resp.StatusCode)
}
resp = bearerGet(t, srv, "/_plugin/agents/coder/runs?window=bogus")
resp.Body.Close()
if resp.StatusCode != http.StatusBadRequest {
t.Errorf("bad window: want 400, got %d", resp.StatusCode)
}
}

func TestAgentRuns_404WithoutLogstore(t *testing.T) {
srv, _ := newBudgetTestServer(t, nil) // no logstore in routeDeps
resp := bearerDo(t, srv, http.MethodGet, "/_plugin/agents/coder/runs", "")
resp.Body.Close()
if resp.StatusCode != http.StatusNotFound {
t.Fatalf("want 404, got %d", resp.StatusCode)
}
}

func TestAgentRuns_Upstream502(t *testing.T) {
now := time.Now().UTC()
bf := newFakeBifrost(t, phase7Logs(now))
srv := newObservabilityTestServer(t, bf)
bf.failNextWith = http.StatusInternalServerError
resp := bearerGet(t, srv, "/_plugin/agents/coder/runs")
resp.Body.Close()
if resp.StatusCode != http.StatusBadGateway {
t.Fatalf("want 502, got %d", resp.StatusCode)
}
}
11 changes: 11 additions & 0 deletions gateway/internal/adminapi/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -360,6 +360,17 @@ func registerRoutes(mux *http.ServeMux, deps routeDeps) {
obs.agentSpend(w, r, parts[0])
return
}
// `<name>/runs` (GET): the agent's runs in the window, newest
// activity first, with user / models / spend per run. Backs
// the AgentDetail "Recent runs" table. Same logstore gate.
if len(parts) == 2 && parts[0] != "" && parts[1] == "runs" {
if obs == nil {
http.NotFound(w, r)
return
}
obs.agentRuns(w, r, parts[0])
return
}
// `/_plugin/agents/catalog` (single segment) is the catalog
// list — every registry agent, traffic or not. Distinct from
// `<name>/catalog` (two segments) which is one agent's detail.
Expand Down
23 changes: 23 additions & 0 deletions gateway/internal/adminapi/ui/src/api/queries.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import { apiFetch, ApiCallError, getErrorMessage } from "./client";
import type {
AgentBudgetResponse,
AgentCatalogResponse,
AgentRunsResponse,
AgentEvalsResponse,
AgentStateResponse,
KillAgentResponse,
Expand Down Expand Up @@ -196,6 +197,28 @@ export function useAgentBudgets(names: string[]) {
return out;
}

// ─── /agents/:name/runs ─────────────────────────────────────────────
//
// One row per run of this agent in the window, newest activity
// first, with the user, model(s), spend and call count. Scoped
// server-side by `metadata.agent-name`, so unlike a run-id
// histogram it never shows another agent's runs. 30s poll, the
// by-agent cadence: an operator fires a run and expects it to land
// at the top of the table on the next tick.

export function useAgentRuns(name: string | undefined, window: Window) {
return useQuery({
queryKey: ["agents", name, "runs", window],
queryFn: () =>
apiFetch<AgentRunsResponse>(
`/agents/${encodeURIComponent(name!)}/runs?window=${encodeURIComponent(window)}`
),
enabled: !!name,
refetchInterval: 30_000,
staleTime: 10_000,
});
}

// ─── /agents/catalog (list) ─────────────────────────────────────────
//
// The whole registry — every catalog agent, traffic or not. The Agents
Expand Down
Loading
Loading