feat(columbo): give the agent fixed fleet-health probes and a bigger tool budget (#600)

This commit is contained in:
Antoine Lecompte
2026-08-31 16:10:02 -04:00
committed by GitHub
parent a1f5454e3e
commit bb0cd74bee
10 changed files with 284 additions and 20 deletions
+15 -3
View File
@@ -71,8 +71,9 @@ The split is **harness vs. model**, not "the agent service is trusted":
- **Harness (trusted, holds the secrets)**: the Go process. It owns the - **Harness (trusted, holds the secrets)**: the Go process. It owns the
OpenRouter key and the internal secret, executes every tool call itself, OpenRouter key and the internal secret, executes every tool call itself,
and posts the final note. The model only ever sees tool *results*. and posts the final note. The model only ever sees tool *results*.
- **Model (untrusted)**: fills the parameters of three typed tools — - **Model (untrusted)**: fills the parameters of four typed tools —
`query_metrics` (PromQL), `query_logs` (LogsQL), `jq` (in-process gojq over `query_metrics` (PromQL), `query_logs` (LogsQL), `query_health` (a probe
name from a fixed registry, see below), `jq` (in-process gojq over
stored results, no shell, no subprocess). No tool takes a URL, header, or stored results, no shell, no subprocess). No tool takes a URL, header, or
credential. There is no command execution and no filesystem access. credential. There is no command execution and no filesystem access.
@@ -92,6 +93,17 @@ filter is the only wall between the agent and other users' telemetry —
which is why it lives in `internal/o11y` with tests asserting a query that which is why it lives in `internal/o11y` with tests asserting a query that
names another user still comes back scoped. names another user still comes back scoped.
`query_health` is the one deliberate exception: fleet-wide platform health
(michael error rates and latency, storage-backend health, Ceph/RGW health,
pool capacity), so a user's 5xx errors can be correlated with a platform
incident. The model never composes the query — it picks a probe *name* from
the hand-maintained registry in `internal/agent/health.go` and a time range,
and the harness runs that probe's fixed PromQL unscoped
(`o11y.QueryFleetRange`). Cross-tenant leakage stays structurally
impossible: the registry must never include series carrying per-customer
labels (customerId, asn, repository ids), which is the review bar for adding
a probe.
Prompt injection is the main residual threat: ticket text and log lines are Prompt injection is the main residual threat: ticket text and log lines are
user-influenceable model input. The blast radius is bounded structurally — user-influenceable model input. The blast radius is bounded structurally —
read-only user-scoped tools, output only to the staff thread, note stamped read-only user-scoped tools, output only to the staff thread, note stamped
@@ -99,7 +111,7 @@ as AI-generated with the executed queries listed — so the worst case is a
misleading note that staff are told to verify. misleading note that staff are told to verify.
Hard limits per investigation: tool-call budget (`COLUMBO_MAX_TOOL_CALLS`, Hard limits per investigation: tool-call budget (`COLUMBO_MAX_TOOL_CALLS`,
16), wall clock (`COLUMBO_TIMEOUT_SECONDS`, 600), model calls retried on 20), wall clock (`COLUMBO_TIMEOUT_SECONDS`, 600), model calls retried on
transport errors/timeouts/5xx with per-attempt deadlines transport errors/timeouts/5xx with per-attempt deadlines
(`COLUMBO_MODEL_TIMEOUT_SECONDS` 120 × `COLUMBO_MODEL_ATTEMPTS` 3 — the (`COLUMBO_MODEL_TIMEOUT_SECONDS` 120 × `COLUMBO_MODEL_ATTEMPTS` 3 — the
response body is buffered per attempt so a mid-body stall retries instead of response body is buffered per attempt so a mid-body stall retries instead of
+6 -4
View File
@@ -38,10 +38,10 @@ type Investigation struct {
} }
type Config struct { type Config struct {
OpenRouterURL string OpenRouterURL string
APIKey string APIKey string
Model string Model string
TriageModel string TriageModel string
MetricsURL string MetricsURL string
LogsURL string LogsURL string
MaxToolCalls int MaxToolCalls int
@@ -130,6 +130,8 @@ Michael operation semantics (restic REST protocol):
- get_blob = blob read (restores, checks); check_blob = existence probe; delete_blob = cleanup (locks after every operation; data/index during prune); list_blobs = listing; save_config = repository initialization (happens once, before the first backup); create_repository = repository creation. - get_blob = blob read (restores, checks); check_blob = existence probe; delete_blob = cleanup (locks after every operation; data/index during prune); list_blobs = listing; save_config = repository initialization (happens once, before the first backup); create_repository = repository creation.
So: op:="save_blob" blob_type:="snapshots" = completed backups; op:="save_blob" blob_type:="data" = backup traffic; status:>=400 on michael = failing restic requests. So: op:="save_blob" blob_type:="snapshots" = completed backups; op:="save_blob" blob_type:="data" = backup traffic; status:>=400 on michael = failing restic requests.
Platform health (fleet-wide, not user-specific): the query_health tool runs fixed named probes over platform telemetry — michael error rates and latency, storage-backend health, Ceph/RGW health, pool capacity. When the user's telemetry shows server-side errors (5xx, timeouts), check whether a platform incident overlaps their error window; a healthy platform during that window is itself evidence. One or two probes over the incident window usually suffice — do not audit the whole platform.
Your final message becomes the staff note verbatim. Format: Your final message becomes the staff note verbatim. Format:
1. One-line verdict (e.g. "Backups from connection X have failed with 507 since 14:02 UTC"). 1. One-line verdict (e.g. "Backups from connection X have failed with 507 since 14:02 UTC").
2. Evidence: the specific log lines / metric numbers, with timestamps. 2. Evidence: the specific log lines / metric numbers, with timestamps.
+142
View File
@@ -0,0 +1,142 @@
package agent
import (
"context"
"fmt"
"strings"
"time"
)
// healthProbes is the complete set of fleet-wide queries the model can run.
// The model only ever picks a name from this list — probe queries are the one
// path that reaches the o11y backends without per-user scoping, so free-form
// query text must never be added here. Every query must stay free of series
// that carry per-customer labels (customerId, asn, repository ids).
type healthProbe struct {
name string
description string
query string
}
var healthProbes = []healthProbe{
{
"michael_requests_by_status",
"michael (restic backend) request rate by HTTP status",
`sum by (cluster, status) (rate({__name__="http.server.request.count"}[5m]))`,
},
{
"michael_error_ratio",
"michael server-error ratio (errors / all requests)",
`sum by (cluster) (rate({__name__="http.server.request.errors"}[5m])) / sum by (cluster) (rate({__name__="http.server.request.count"}[5m]))`,
},
{
"michael_latency_p99",
"michael p99 request duration in seconds",
`histogram_quantile(0.99, sum by (cluster, le) (rate({__name__="http.server.request.duration_bucket"}[5m])))`,
},
{
"backend_health",
"per-backend health as michael sees its storage backends (1 healthy, 0 unhealthy)",
`max by (cluster, backend) ({__name__="s3.backend.healthy"})`,
},
{
"backend_errors",
"error rate of michael's requests to each storage backend",
`sum by (cluster, backend) (rate({__name__="s3.backend.errors"}[5m]))`,
},
{
"backend_retries",
"michael's storage-backend retry rate by outcome (success/failure/denied)",
`sum by (cluster, outcome) (rate({__name__="s3.pool.retries"}[5m]))`,
},
{
"ceph_status",
"Ceph cluster health per storage cluster (0 OK, 1 WARN, 2 ERR)",
`ceph_health_status`,
},
{
"ceph_pg_problems",
"Ceph placement groups in degraded/undersized/inconsistent/backfilling/recovering/peering states",
`sum by (cluster, __name__) ({__name__=~"ceph_pg_(degraded|undersized|inconsistent|backfilling|recovering|peering)"})`,
},
{
"ceph_osds_down",
"number of Ceph OSDs down per storage cluster",
`count by (cluster) (ceph_osd_up) - sum by (cluster) (ceph_osd_up)`,
},
{
"rgw_errors",
"rate of failed requests at the Ceph RGW S3 gateways",
`sum by (cluster) (rate(ceph_rgw_failed_req[5m]))`,
},
{
"rgw_get_latency",
"average object GET latency at the RGW gateways in seconds",
`sum by (cluster) (rate(ceph_rgw_op_get_obj_lat_sum[5m])) / sum by (cluster) (rate(ceph_rgw_op_get_obj_lat_count[5m]))`,
},
{
"rgw_put_latency",
"average object PUT latency at the RGW gateways in seconds",
`sum by (cluster) (rate(ceph_rgw_op_put_obj_lat_sum[5m])) / sum by (cluster) (rate(ceph_rgw_op_put_obj_lat_count[5m]))`,
},
{
"rgw_queue",
"RGW request queue length per gateway daemon",
`sum by (cluster, ceph_daemon) (ceph_rgw_qlen)`,
},
{
"pool_capacity",
"Ceph pool fill percentage by pool name",
`ceph_pool_percent_used * on (cluster, pool_id) group_left(name) ceph_pool_metadata`,
},
}
func healthProbeByName(name string) (healthProbe, bool) {
for _, p := range healthProbes {
if p.name == name {
return p, true
}
}
return healthProbe{}, false
}
func healthDescription() string {
var b strings.Builder
b.WriteString("Check fleet-wide platform health (NOT user-specific). " +
"You pick a probe by name and a time range; a fixed query runs — there is no free-form query on this tool. " +
"Use it to test whether the user's symptoms coincide with a platform-side incident. Probes:")
for _, p := range healthProbes {
b.WriteString("\n- " + p.name + ": " + p.description)
}
return b.String()
}
type healthArgs struct {
Probe string `json:"probe" jsonschema:"description=Probe name from the list in the tool description"`
Start string `json:"start,omitempty" jsonschema:"description=Range start as RFC3339 or unix seconds; defaults to 24h ago, capped at 30 days back"`
End string `json:"end,omitempty" jsonschema:"description=Range end as RFC3339 or unix seconds; defaults to now"`
Step string `json:"step,omitempty" jsonschema:"description=Resolution step such as 5m; defaults to 5m, minimum 1m"`
}
func (t *toolbox) queryHealth(ctx context.Context, args healthArgs) (string, error) {
probe, ok := healthProbeByName(args.Probe)
if !ok {
return "", fmt.Errorf("unknown probe %q — pick one from the tool description", args.Probe)
}
if err := t.spend("health: " + probe.name); err != nil {
return "", err
}
start, end, err := resolveRange(args.Start, args.End, time.Now().UTC())
if err != nil {
return "", err
}
step, err := resolveStep(args.Step)
if err != nil {
return "", err
}
result, err := t.o11y.QueryFleetRange(ctx, probe.query, start, end, step)
if err != nil {
return "", err
}
return t.deliver(result), nil
}
@@ -0,0 +1,70 @@
package agent
import (
"context"
"net/http"
"net/http/httptest"
"net/url"
"strings"
"testing"
"columbo/internal/o11y"
)
func TestHealthProbeRegistryIsWellFormed(t *testing.T) {
seen := map[string]bool{}
for _, p := range healthProbes {
if p.name == "" || p.description == "" || p.query == "" {
t.Fatalf("incomplete probe %+v", p)
}
if seen[p.name] {
t.Fatalf("duplicate probe name %q", p.name)
}
seen[p.name] = true
if !strings.Contains(healthDescription(), p.name) {
t.Fatalf("probe %q missing from the tool description", p.name)
}
}
}
func TestQueryHealthRejectsUnknownProbe(t *testing.T) {
box := testBox(4, 1024)
_, err := box.queryHealth(context.Background(), healthArgs{Probe: "drop_tables"})
if err == nil || !strings.Contains(err.Error(), "unknown probe") {
t.Fatalf("expected an unknown-probe error, got %v", err)
}
if box.callsMade() != 0 {
t.Fatal("an unknown probe must not consume the tool budget")
}
}
func TestQueryHealthRunsTheRegistryQueryUnscoped(t *testing.T) {
var form url.Values
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if err := r.ParseForm(); err != nil {
t.Fatal(err)
}
form = r.Form
_, _ = w.Write([]byte(`{"status":"success"}`))
}))
t.Cleanup(srv.Close)
box := newToolbox(o11y.NewClient(srv.URL, srv.URL, "user-1"), NewResultStore(), 4, 1024)
out, err := box.queryHealth(context.Background(), healthArgs{Probe: "backend_health"})
if err != nil {
t.Fatal(err)
}
if out != `{"status":"success"}` {
t.Fatalf("out = %q", out)
}
probe, _ := healthProbeByName("backend_health")
if got := form.Get("query"); got != probe.query {
t.Fatalf("query = %q, want the registry query %q", got, probe.query)
}
if _, scoped := form["extra_label"]; scoped {
t.Fatal("health probes must not carry the per-user extra_label")
}
if got := box.queriesRun(); len(got) != 1 || got[0] != "health: backend_health" {
t.Fatalf("queriesRun = %v", got)
}
}
+9 -4
View File
@@ -28,8 +28,9 @@ const (
var errToolBudget = errors.New("tool budget exhausted — write your conclusion with what you have") var errToolBudget = errors.New("tool budget exhausted — write your conclusion with what you have")
// toolbox is the complete capability surface the model gets: two read-only // toolbox is the complete capability surface the model gets: two read-only
// queries pre-scoped to one user, and an in-process jq over stored results. // queries pre-scoped to one user, a fixed-probe fleet-health check, and an
// No tool takes a URL, a header, or a credential. // in-process jq over stored results. No tool takes a URL, a header, or a
// credential.
type toolbox struct { type toolbox struct {
o11y *o11y.Client o11y *o11y.Client
store *ResultStore store *ResultStore
@@ -54,6 +55,10 @@ func (t *toolbox) tools() ([]tool.BaseTool, error) {
if err != nil { if err != nil {
return nil, err return nil, err
} }
health, err := utils.InferTool("query_health", healthDescription(), audited("query_health", t.queryHealth))
if err != nil {
return nil, err
}
jq, err := utils.InferTool("jq", jqDescription, audited("jq", t.jq)) jq, err := utils.InferTool("jq", jqDescription, audited("jq", t.jq))
if err != nil { if err != nil {
return nil, err return nil, err
@@ -62,8 +67,8 @@ func (t *toolbox) tools() ([]tool.BaseTool, error) {
// not run failures: the model gets the error text and can correct itself // not run failures: the model gets the error text and can correct itself
// instead of the whole investigation dying on a syntax error. MaxStep // instead of the whole investigation dying on a syntax error. MaxStep
// still bounds a model that never recovers. // still bounds a model that never recovers.
wrapped := make([]tool.BaseTool, 0, 3) wrapped := make([]tool.BaseTool, 0, 4)
for _, t := range []tool.BaseTool{metrics, logs, jq} { for _, t := range []tool.BaseTool{metrics, logs, health, jq} {
wrapped = append(wrapped, utils.WrapToolWithErrorHandler(t, func(_ context.Context, err error) string { wrapped = append(wrapped, utils.WrapToolWithErrorHandler(t, func(_ context.Context, err error) string {
return "ERROR: " + err.Error() return "ERROR: " + err.Error()
})) }))
@@ -114,7 +114,7 @@ func TestToolErrorsBecomeToolResults(t *testing.T) {
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
jq := tools[2].(interface { jq := tools[3].(interface {
InvokableRun(ctx context.Context, argumentsInJSON string, opts ...tool.Option) (string, error) InvokableRun(ctx context.Context, argumentsInJSON string, opts ...tool.Option) (string, error)
}) })
out, err := jq.InvokableRun(context.Background(), `{"program":".","ref":"r99"}`) out, err := jq.InvokableRun(context.Background(), `{"program":".","ref":"r99"}`)
+1 -1
View File
@@ -83,7 +83,7 @@ func LoadConfig() Config {
LogsURL: envOr("O11Y_LOGS_URL", "http://localhost:9428"), LogsURL: envOr("O11Y_LOGS_URL", "http://localhost:9428"),
BotURL: envOr("FUTO_BACKUPS_BOT_URL", "http://localhost:3050"), BotURL: envOr("FUTO_BACKUPS_BOT_URL", "http://localhost:3050"),
GrafanaURL: envOr("GRAFANA_URL", "https://grafana.futostatus.com"), GrafanaURL: envOr("GRAFANA_URL", "https://grafana.futostatus.com"),
MaxToolCalls: envIntMin("COLUMBO_MAX_TOOL_CALLS", 16, 1), MaxToolCalls: envIntMin("COLUMBO_MAX_TOOL_CALLS", 20, 1),
InvestigationTimeout: time.Duration(envIntMin("COLUMBO_TIMEOUT_SECONDS", 600, 10)) * time.Second, InvestigationTimeout: time.Duration(envIntMin("COLUMBO_TIMEOUT_SECONDS", 600, 10)) * time.Second,
ModelCallTimeout: time.Duration(envIntMin("COLUMBO_MODEL_TIMEOUT_SECONDS", 120, 10)) * time.Second, ModelCallTimeout: time.Duration(envIntMin("COLUMBO_MODEL_TIMEOUT_SECONDS", 120, 10)) * time.Second,
ModelCallAttempts: envIntMin("COLUMBO_MODEL_ATTEMPTS", 3, 1), ModelCallAttempts: envIntMin("COLUMBO_MODEL_ATTEMPTS", 3, 1),
+16 -1
View File
@@ -6,7 +6,8 @@
// are unauthenticated from the cluster), so it must never depend on the // are unauthenticated from the cluster), so it must never depend on the
// model composing its queries correctly. Logs match either per-user field // model composing its queries correctly. Logs match either per-user field
// convention: michael writes `user`, the NestJS services `customerId` // convention: michael writes `user`, the NestJS services `customerId`
// (mirroring the yucca-per-user dashboard's scoping). // (mirroring the yucca-per-user dashboard's scoping). QueryFleetRange is the
// single unscoped exception; see its doc.
package o11y package o11y
import ( import (
@@ -61,6 +62,20 @@ func (c *Client) MetricNames(ctx context.Context, lookback time.Duration) ([]str
return parsed.Data, nil return parsed.Data, nil
} }
// QueryFleetRange runs a platform-health query with NO per-user scope. It is
// the one deliberate exception to this package's scoping rule and must only
// ever receive queries from the harness's fixed probe registry, never
// model-composed text — the registry is what keeps this path from becoming a
// cross-tenant hole.
func (c *Client) QueryFleetRange(ctx context.Context, query, start, end, step string) (string, error) {
params := url.Values{}
params.Set("query", query)
params.Set("start", start)
params.Set("end", end)
params.Set("step", step)
return c.do(ctx, c.MetricsURL+"/api/v1/query_range", params)
}
func (c *Client) QueryMetricsRange(ctx context.Context, query, start, end, step string) (string, error) { func (c *Client) QueryMetricsRange(ctx context.Context, query, start, end, step string) (string, error) {
params := url.Values{} params := url.Values{}
params.Set("query", query) params.Set("query", query)
@@ -111,6 +111,24 @@ func TestScopeLogsQLRejectsJoinAndUnion(t *testing.T) {
} }
} }
func TestQueryFleetRangeIsUnscoped(t *testing.T) {
srv, captured, form := recordingServer(t, http.StatusOK, `{"status":"success"}`)
client := NewClient(srv.URL, srv.URL, "user-1")
if _, err := client.QueryFleetRange(context.Background(), "ceph_health_status", "0", "1", "5m"); err != nil {
t.Fatal(err)
}
if captured.URL.Path != "/api/v1/query_range" {
t.Fatalf("unexpected path %q", captured.URL.Path)
}
if got := formValue(t, *form, "extra_label"); got != "" {
t.Fatalf("extra_label = %q, want none on the fleet path", got)
}
if got := formValue(t, *form, "query"); got != "ceph_health_status" {
t.Fatalf("query = %q", got)
}
}
func TestMetricNamesIsScopedToCustomer(t *testing.T) { func TestMetricNamesIsScopedToCustomer(t *testing.T) {
srv, captured, form := recordingServer(t, http.StatusOK, `{"status":"success","data":["api_request_count","blobs.uploaded_bytes"]}`) srv, captured, form := recordingServer(t, http.StatusOK, `{"status":"success","data":["api_request_count","blobs.uploaded_bytes"]}`)
client := NewClient(srv.URL, srv.URL, "user-1") client := NewClient(srv.URL, srv.URL, "user-1")
+6 -6
View File
@@ -49,12 +49,12 @@ func main() {
} }
runner := agent.NewRunner(agent.Config{ runner := agent.NewRunner(agent.Config{
OpenRouterURL: cfg.OpenRouterURL, OpenRouterURL: cfg.OpenRouterURL,
APIKey: cfg.OpenRouterAPIKey, APIKey: cfg.OpenRouterAPIKey,
Model: cfg.Model, Model: cfg.Model,
TriageModel: cfg.TriageModel, TriageModel: cfg.TriageModel,
MetricsURL: cfg.MetricsURL, MetricsURL: cfg.MetricsURL,
LogsURL: cfg.LogsURL, LogsURL: cfg.LogsURL,
MaxToolCalls: cfg.MaxToolCalls, MaxToolCalls: cfg.MaxToolCalls,
ToolResultBytes: cfg.ToolResultBytes, ToolResultBytes: cfg.ToolResultBytes,
ModelCallTimeout: cfg.ModelCallTimeout, ModelCallTimeout: cfg.ModelCallTimeout,