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
54 changes: 36 additions & 18 deletions internal/http/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -129,24 +129,7 @@ func (s *serverImpl) Start(ctx *stopper.Context) error {
if err != nil {
return err
}
handler := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
defer r.Body.Close()
if s.metricsWriter != nil {
s.metricsWriter.Copy(ctx, w)
}
metrics, err := s.registry.Gather()
if err != nil {
s.errorResponse(w, "Error gathering metrics", err)
return
}
for _, m := range metrics {
_, err = expfmt.MetricFamilyToText(w, m)
if err != nil {
s.errorResponse(w, "Error gathering metrics", err)
}
}

})
handler := http.HandlerFunc(s.metricsHandler(ctx))

http.Handle(s.config.Endpoint, gziphandler.GzipHandler(handler))
s.debugInfo()
Expand Down Expand Up @@ -180,6 +163,41 @@ func (s *serverImpl) Start(ctx *stopper.Context) error {
return nil
}

// metricsHandler returns the handler that serves the metrics endpoint. It first
// copies the metrics fetched from the upstream source (applying any histogram
// translators) and then appends the metrics gathered from the local registry.
func (s *serverImpl) metricsHandler(ctx *stopper.Context) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
defer r.Body.Close()
if s.metricsWriter != nil {
// A Copy failure is usually the upstream source being unreachable
// or returning malformed metrics, which is actionable, so it is
// logged at error level. It can also be the client disconnecting
// mid-response, in which case the stream is already broken and we
// cannot report anything back either way.
if err := s.metricsWriter.Copy(ctx, w); err != nil {
log.Errorf("Error copying metrics from source: %s", err.Error())
return
}
}
metrics, err := s.registry.Gather()
if err != nil {
s.errorResponse(w, "Error gathering metrics", err)
return
}
for _, m := range metrics {
// A write failure here means the client has disconnected, which is
// an expected event (for example a scrape timeout). We stop writing
// and log at debug level rather than trying to send an error to a
// client that is no longer listening.
if _, err := expfmt.MetricFamilyToText(w, m); err != nil {
log.Debugf("Stopped writing metrics to client: %s", err.Error())
return
}
}
}
}

func (s *serverImpl) errorResponse(w http.ResponseWriter, msg string, err error) {
log.Errorf("%s: %s", msg, err.Error())
w.WriteHeader(http.StatusInternalServerError)
Expand Down
103 changes: 103 additions & 0 deletions internal/http/server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@ import (
"context"
"crypto/tls"
"crypto/x509"
"net/http"
"net/http/httptest"
_ "net/http/pprof"
"os"
"testing"
Expand All @@ -29,10 +31,40 @@ import (
"github.com/cockroachlabs/visus/internal/metric"
"github.com/cockroachlabs/visus/internal/server"
"github.com/cockroachlabs/visus/internal/store"
"github.com/prometheus/client_golang/prometheus"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// disconnectedResponseWriter simulates a client that has disconnected: every
// write fails, mimicking the "client disconnected" error returned by the HTTP/2
// server once the peer is gone. It records whether WriteHeader was called so
// tests can assert the handler does not try to send an error status back to a
// client that is no longer listening.
type disconnectedResponseWriter struct {
header http.Header
wroteHeader bool
statusCode int
}

var _ http.ResponseWriter = &disconnectedResponseWriter{}

func (d *disconnectedResponseWriter) Header() http.Header {
if d.header == nil {
d.header = http.Header{}
}
return d.header
}

func (d *disconnectedResponseWriter) Write([]byte) (int, error) {
return 0, errors.New("client disconnected")
}

func (d *disconnectedResponseWriter) WriteHeader(statusCode int) {
d.wroteHeader = true
d.statusCode = statusCode
}

// TestRefreshHistograms verifies we reload the histograms from the store.
func TestRefreshHistograms(t *testing.T) {
r := require.New(t)
Expand Down Expand Up @@ -198,3 +230,74 @@ func TestRefreshTLSConfig(t *testing.T) {
r.NoError(err)
a.Equal(expectedCert.Certificate, cert.Certificate)
}

// newTestRegistry returns a registry with a single registered gauge so the
// metrics handler has something to gather.
func newTestRegistry(t *testing.T) *prometheus.Registry {
t.Helper()
registry := prometheus.NewRegistry()
gauge := prometheus.NewGauge(prometheus.GaugeOpts{
Name: "visus_test_gauge",
Help: "gauge used by the metrics handler tests",
})
gauge.Set(1)
require.NoError(t, registry.Register(gauge))
return registry
}

// TestMetricsHandler verifies that a normal scrape returns the gathered metrics.
func TestMetricsHandler(t *testing.T) {
a := assert.New(t)
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
stop := stopper.WithContext(ctx)
server := &serverImpl{
registry: newTestRegistry(t),
}
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/metrics", nil)
server.metricsHandler(stop)(rec, req)
a.Equal(http.StatusOK, rec.Code)
a.Contains(rec.Body.String(), "visus_test_gauge")
}

// TestMetricsHandlerClientDisconnect verifies that when the client disconnects
// mid-response the handler stops writing and does not try to send an error
// status back to the client that is no longer listening.
func TestMetricsHandlerClientDisconnect(t *testing.T) {
a := assert.New(t)
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
stop := stopper.WithContext(ctx)
server := &serverImpl{
registry: newTestRegistry(t),
}
w := &disconnectedResponseWriter{}
req := httptest.NewRequest(http.MethodGet, "/metrics", nil)
// The handler must not panic and must not attempt to write an error status
// back to the disconnected client.
server.metricsHandler(stop)(w, req)
a.False(w.wroteHeader, "handler should not write a status back to a disconnected client")
}

// TestMetricsHandlerCopyError verifies that when copying from the upstream
// source fails the handler returns early without writing gathered metrics.
func TestMetricsHandlerCopyError(t *testing.T) {
r := require.New(t)
a := assert.New(t)
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
stop := stopper.WithContext(ctx)
// A file source that does not exist makes Copy fail on open.
writer, err := metric.NewWriter("file:///nonexistent-source.txt", nil, nil)
r.NoError(err)
server := &serverImpl{
registry: newTestRegistry(t),
metricsWriter: writer,
}
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/metrics", nil)
server.metricsHandler(stop)(rec, req)
// Copy failed, so we return before gathering registry metrics.
a.NotContains(rec.Body.String(), "visus_test_gauge")
}
4 changes: 3 additions & 1 deletion internal/metric/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -115,7 +115,9 @@ func (w *Writer) Copy(ctx context.Context, out io.Writer) error {
for _, mf := range metricFamilies {
if mf.GetType() == dto.MetricType_HISTOGRAM && translators != nil {
for _, translator := range translators {
translator.Translate(ctx, mf, out)
if err := translator.Translate(ctx, mf, out); err != nil {
return err
}
}
} else {
_, err := expfmt.MetricFamilyToText(out, mf)
Expand Down