Skip to content
Open
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
2 changes: 1 addition & 1 deletion charts/nudgebee-agent/values.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ runnerServiceAccount:
runner:
image:
repository: ghcr.io/nudgebee/nudgebee-agent
tag: 2026-08-12T06-01-54_7ec58f7701909a3ce172ad2a9235f8b15255e363
tag: 2026-08-12T08-10-04_e13ce92bc8c238e883d050b1276c6128b2c8fb44
# Image template the pod_profiler action launches debugger pods from.
# The agent substitutes `{}` for the variant (bpf, jvm, python, perf, ruby).
# Surfaces as PROFILER_IMAGE; leave empty to fall back to the binary default.
Expand Down
63 changes: 45 additions & 18 deletions runner/cmd/agent/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (
"encoding/json"
"errors"
"fmt"
"io"
"log/slog"
"maps"
"net/http"
Expand Down Expand Up @@ -991,9 +992,10 @@ func run(ctx context.Context, logger *slog.Logger, cfg *config.Config) error {
probeCtx, probeCancel := context.WithTimeout(gctx, 30*time.Second)
defer probeCancel()
probeClient := &http.Client{Timeout: 5 * time.Second}
logsProvider, logsURL, logsOK, logCfg := probeLogsProvider(probeCtx, cfg)
logsProvider, logsURL, logsOK, logsErr, logCfg := probeLogsProvider(probeCtx, cfg)
as := telemetry.DetectAutoScaler(probeCtx, typedKube, providerInfo.Provider, logger)
clickhouseStatus, clickhouseErr := probeClickhouse(probeCtx, probeClient, clickhouseHost, clickhousePort)
promConnected, promErr := prometheusConnected(probeCtx, promClient, logger)
return telemetry.Datasources{
PrometheusURL: cfg.PrometheusURL,
AlertManagerURL: cfg.AlertManagerURL,
Expand All @@ -1002,8 +1004,10 @@ func run(ctx context.Context, logger *slog.Logger, cfg *config.Config) error {
LogsProvider: logsProvider,
LogsProviderURL: logsURL,
LogsProviderStatus: logsOK,
LogsProviderError: logsErr,
LogProviderConfig: logCfg,
PrometheusConnected: prometheusConnected(probeCtx, promClient, logger),
PrometheusConnected: promConnected,
PrometheusConnectedError: promErr,
NodeAgentCount: queryNodeAgentCount(probeCtx, promClient, logger),
PrometheusRetentionTime: telemetry.PrometheusRetention(probeCtx, promClient, logger),
PrometheusAdditionalLabels: promExtraLabels,
Expand Down Expand Up @@ -1169,24 +1173,27 @@ func (a *grafanaAdapter) HandlePrometheus(ctx context.Context, r *dispatch.Grafa
// auth — Chronosphere, Thanos Query, Grafana Mimir, Amazon Managed Prometheus —
// are reported Connected when metric queries work. Returns false on any error
// so a broken backend shows Disconnected rather than panicking the tick.
func prometheusConnected(ctx context.Context, c *prometheus.Client, logger *slog.Logger) bool {
func prometheusConnected(ctx context.Context, c *prometheus.Client, logger *slog.Logger) (ok bool, reason string) {
if c == nil || c.BaseURL == "" {
return false
return false, ""
}
cctx, cancel := context.WithTimeout(ctx, 5*time.Second)
defer cancel()
raw, err := c.Query(cctx, "vector(1)", "", "")
if err != nil {
logger.Debug("prometheus health query failed", "err", err)
return false
return false, err.Error()
}
var resp struct {
Status string `json:"status"`
}
if err := json.Unmarshal(raw, &resp); err != nil {
return false
return false, err.Error()
}
if resp.Status == "success" {
return true, ""
}
return resp.Status == "success"
return false, fmt.Sprintf("prometheus query returned status %q", resp.Status)
}

// queryNodeAgentCount
Expand Down Expand Up @@ -1249,46 +1256,55 @@ func selectedLogsProvider(cfg *config.Config) (provider, url string) {
// configured provider's own client at action-handler time. Fail-closed: any
// non-2xx → status=false, URL stays in payload so the UI can show "URL
// configured but unhealthy".
func probeLogsProvider(ctx context.Context, cfg *config.Config) (provider, url string, ok bool, providerCfg map[string]any) {
func probeLogsProvider(ctx context.Context, cfg *config.Config) (provider, url string, ok bool, reason string, providerCfg map[string]any) {
httpClient := &http.Client{Timeout: 5 * time.Second}
switch {
case cfg.PinotURL != "":
ok = httpProbe(ctx, httpClient, cfg.PinotURL+"/health")
return "pinot", cfg.PinotURL, ok, map[string]any{}
err := httpProbeErr(ctx, httpClient, cfg.PinotURL+"/health")
return "pinot", cfg.PinotURL, err == nil, errString(err), map[string]any{}
case cfg.ElasticsearchEnabled && cfg.ElasticsearchURL != "":
// ES exposes a `_cluster/health` endpoint; we treat 200 as healthy.
// Probe with the configured credentials so the badge reflects whether
// queries will actually succeed — a secured OpenSearch/ES otherwise 401s
// on an unauthenticated probe even when the configured creds work fine.
ok = httpProbe(ctx, httpClient, cfg.ElasticsearchURL+"/_cluster/health", esAuthHeader(cfg))
err := httpProbeErr(ctx, httpClient, cfg.ElasticsearchURL+"/_cluster/health", esAuthHeader(cfg))
providerCfg = map[string]any{}
if v := os.Getenv("ELASTICSEARCH_LOG_INDEX"); v != "" {
providerCfg["default_index"] = v
}
return "ES", cfg.ElasticsearchURL, ok, providerCfg
return "ES", cfg.ElasticsearchURL, err == nil, errString(err), providerCfg
case cfg.SignozURL != "":
// Signoz health endpoint: /api/v1/health.
ok = httpProbe(ctx, httpClient, cfg.SignozURL+"/api/v1/health")
err := httpProbeErr(ctx, httpClient, cfg.SignozURL+"/api/v1/health")
providerCfg = map[string]any{}
// Report the Signoz server version so the backend/UI can surface it
// and gate version-specific behaviour. /api/v1/version is unauthed.
if v := fetchSignozVersion(ctx, httpClient, cfg.SignozURL); v != "" {
providerCfg["version"] = v
}
return "signoz", cfg.SignozURL, ok, providerCfg
return "signoz", cfg.SignozURL, err == nil, errString(err), providerCfg
case cfg.LokiURL != "":
// LOKI_URL points at the loki gateway, whose nginx only proxies the
// `/loki/...` API paths — the backend `/ready` is not exposed there and
// 404s. Probe a gateway-served API endpoint instead so the badge
// reflects query reachability.
ok = httpProbe(ctx, httpClient, cfg.LokiURL+"/loki/api/v1/status/buildinfo")
err := httpProbeErr(ctx, httpClient, cfg.LokiURL+"/loki/api/v1/status/buildinfo")
providerCfg = map[string]any{"url": cfg.LokiURL}
return "loki", cfg.LokiURL, ok, providerCfg
return "loki", cfg.LokiURL, err == nil, errString(err), providerCfg
default:
return "", "", false, map[string]any{}
return "", "", false, "", map[string]any{}
}
}

// errString renders a probe failure for the health UI: empty when healthy so
// the wire field clears, the error text otherwise.
func errString(err error) string {
if err == nil {
return ""
}
return err.Error()
}

// probeClickhouse mirrors the legacy _check_clickhouse → db.health() probe.
// Returns false (without probing) when CLICKHOUSE_HOST is unset — the Helm
// chart only wires the host when clickhouse/otel-collector is enabled, so
Expand Down Expand Up @@ -1394,7 +1410,18 @@ func httpProbeErr(ctx context.Context, c *http.Client, url string, headers ...ma
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("HTTP %d", resp.StatusCode)
// Include a compact body snippet (whitespace collapsed, truncated) —
// backends put the useful detail ("token is expired", CORS/auth pages)
// in the body, and the UI renders this string verbatim.
body, _ := io.ReadAll(io.LimitReader(resp.Body, 512))
msg := strings.Join(strings.Fields(string(body)), " ")
if len(msg) > 200 {
msg = msg[:200] + "…"
}
if msg == "" {
return fmt.Errorf("HTTP %d", resp.StatusCode)
}
return fmt.Errorf("HTTP %d: %s", resp.StatusCode, msg)
}
return nil
}
Expand Down
16 changes: 9 additions & 7 deletions runner/cmd/agent/probe_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ func TestProbeLogsProvider_ESDisabledFallsThroughToLoki(t *testing.T) {
ElasticsearchEnabled: false,
LokiURL: loki.URL,
}
provider, url, ok, _ := probeLogsProvider(context.Background(), cfg)
provider, url, ok, _, _ := probeLogsProvider(context.Background(), cfg)
if provider != "loki" {
t.Fatalf("provider = %q, want loki", provider)
}
Expand Down Expand Up @@ -67,7 +67,7 @@ func TestProbeLogsProvider_StrayESURLDoesNotMaskSignoz(t *testing.T) {
ElasticsearchEnabled: false,
SignozURL: signoz.URL,
}
provider, url, ok, _ := probeLogsProvider(context.Background(), cfg)
provider, url, ok, _, _ := probeLogsProvider(context.Background(), cfg)
if provider != "signoz" {
t.Fatalf("provider = %q, want signoz (stray ES URL must not mask SigNoz)", provider)
}
Expand Down Expand Up @@ -97,7 +97,7 @@ func TestProbeLogsProvider_ESProbeSendsAuth(t *testing.T) {
ElasticsearchUser: "admin",
ElasticsearchPassword: "pw",
}
provider, _, ok, _ := probeLogsProvider(context.Background(), cfg)
provider, _, ok, _, _ := probeLogsProvider(context.Background(), cfg)
if provider != "ES" {
t.Fatalf("provider = %q, want ES", provider)
}
Expand Down Expand Up @@ -134,7 +134,7 @@ func TestPrometheusConnected_ChronosphereStyleBackend(t *testing.T) {

c := prometheus.New(srv.URL, nil)
c.ExtraHeaders = config.ParseHeaders("Authorization: Bearer tok")
if !prometheusConnected(context.Background(), c, slog.Default()) {
if ok, _ := prometheusConnected(context.Background(), c, slog.Default()); !ok {
t.Error("expected connected=true for query-only backend serving /api/v1/query")
}
if !queried {
Expand All @@ -151,14 +151,16 @@ func TestPrometheusConnected_FailuresReportDisconnected(t *testing.T) {
w.WriteHeader(http.StatusInternalServerError)
}))
defer down.Close()
if prometheusConnected(context.Background(), prometheus.New(down.URL, nil), slog.Default()) {
if ok, reason := prometheusConnected(context.Background(), prometheus.New(down.URL, nil), slog.Default()); ok {
t.Error("expected connected=false when backend returns 500")
} else if reason == "" {
t.Error("expected a non-empty failure reason when backend returns 500")
}
// Nil / unconfigured client → not connected, no panic.
if prometheusConnected(context.Background(), nil, slog.Default()) {
if ok, _ := prometheusConnected(context.Background(), nil, slog.Default()); ok {
t.Error("expected connected=false for nil client")
}
if prometheusConnected(context.Background(), prometheus.New("", nil), slog.Default()) {
if ok, _ := prometheusConnected(context.Background(), prometheus.New("", nil), slog.Default()); ok {
t.Error("expected connected=false for empty base URL")
}
}
Expand Down
Loading
Loading