Author SHA1 Message Date
clawbot ef8e8cf398 Serve /metrics from a registry of its own (closes #227)
check / check (push) Successful in 4m2s
The HTTP metrics recorder, the delivery collectors and the Go and
process collectors now register on one prometheus.Registry that fx
provides, instead of Prometheus's global default registry, and
/metrics serves that registry. A second metrics-enabled router in one
process, or the server tests run with -count=2, no longer panics on a
duplicate registration.

The middleware builds its recorder once, in New, so installing
Metrics() on more than one router over the same graph is also safe.
The scrape keeps the same series and labels, including go_*,
process_* and promhttp_metric_handler_*.

Model: opus-5-5
2026-09-29 09:20:29 +00:00
22 changed files with 155 additions and 638 deletions
-2
View File
@@ -3280,5 +3280,3 @@ MIT
## Author ## Author
[@sneak](https://sneak.berlin) [@sneak](https://sneak.berlin)
+5
View File
@@ -16,6 +16,7 @@ import (
"sneak.berlin/go/webhooker/internal/handlers" "sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck" "sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger" "sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware" "sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/resetpw" "sneak.berlin/go/webhooker/internal/resetpw"
"sneak.berlin/go/webhooker/internal/server" "sneak.berlin/go/webhooker/internal/server"
@@ -177,6 +178,10 @@ func newApp() *fx.App {
healthcheck.New, healthcheck.New,
session.New, session.New,
handlers.New, handlers.New,
// The registry /metrics serves, and the delivery
// collectors registered on it.
metrics.NewRegistry,
metrics.New,
middleware.New, middleware.New,
// The one SSRF guard both target-creation validation // The one SSRF guard both target-creation validation
// and the delivery dialer consult, so they cannot // and the delivery dialer consult, so they cannot
+1 -1
View File
@@ -4,7 +4,6 @@ go 1.26.1
require ( require (
github.com/99designs/basicauth-go v0.0.0-20230316000542-bf6f9cbbf0f8 github.com/99designs/basicauth-go v0.0.0-20230316000542-bf6f9cbbf0f8
github.com/dustin/go-humanize v1.0.1
github.com/getsentry/sentry-go v0.25.0 github.com/getsentry/sentry-go v0.25.0
github.com/go-chi/chi v1.5.5 github.com/go-chi/chi v1.5.5
github.com/go-chi/cors v1.2.1 github.com/go-chi/cors v1.2.1
@@ -30,6 +29,7 @@ require (
github.com/beorn7/perks v1.0.1 // indirect github.com/beorn7/perks v1.0.1 // indirect
github.com/cespare/xxhash/v2 v2.2.0 // indirect github.com/cespare/xxhash/v2 v2.2.0 // indirect
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect
github.com/dustin/go-humanize v1.0.1 // indirect
github.com/gorilla/securecookie v1.1.2 // indirect github.com/gorilla/securecookie v1.1.2 // indirect
github.com/jinzhu/inflection v1.0.0 // indirect github.com/jinzhu/inflection v1.0.0 // indirect
github.com/jinzhu/now v1.1.5 // indirect github.com/jinzhu/now v1.1.5 // indirect
+6 -5
View File
@@ -148,6 +148,7 @@ type EngineParams struct {
DBManager *database.WebhookDBManager DBManager *database.WebhookDBManager
Logger *logger.Logger Logger *logger.Logger
SSRFGuard *Guard SSRFGuard *Guard
Metrics *metrics.Set
} }
// Engine processes queued deliveries in the background // Engine processes queued deliveries in the background
@@ -167,10 +168,10 @@ type Engine struct {
retryCh chan Task retryCh chan Task
workers int workers int
// mtr is the delivery metric set. Production wires the // mtr is the delivery metric set. Production wires the one
// process-wide one; a test can substitute a set registered on // registered on the registry /metrics serves; a test can
// a private registry so its assertions are not disturbed by // substitute a set registered on a registry it holds, so it can
// deliveries other tests are making at the same time. // gather what its own deliveries recorded.
mtr *metrics.Set mtr *metrics.Set
// targets maps each target type to its implementation. // targets maps each target type to its implementation.
@@ -204,7 +205,7 @@ func New(
deliveryCh: make(chan Task, deliveryChannelSize), deliveryCh: make(chan Task, deliveryChannelSize),
retryCh: make(chan Task, retryChannelSize), retryCh: make(chan Task, retryChannelSize),
workers: defaultWorkers, workers: defaultWorkers,
mtr: metrics.Default(), mtr: params.Metrics,
} }
e.initTargets(&http.Client{ e.initTargets(&http.Client{
+5 -5
View File
@@ -9,6 +9,7 @@ import (
"net/url" "net/url"
"time" "time"
"github.com/prometheus/client_golang/prometheus"
"go.uber.org/fx" "go.uber.org/fx"
"gorm.io/gorm" "gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
@@ -389,7 +390,7 @@ func NewTestEngine(
deliveryCh: make(chan Task, deliveryChannelSize), deliveryCh: make(chan Task, deliveryChannelSize),
retryCh: make(chan Task, retryChannelSize), retryCh: make(chan Task, retryChannelSize),
workers: workers, workers: workers,
mtr: metrics.Default(), mtr: metrics.New(prometheus.NewRegistry()),
} }
e.initTargets(client) e.initTargets(client)
@@ -404,7 +405,7 @@ func NewTestEngineSmallRetry(
e := &Engine{ e := &Engine{
log: log, log: log,
retryCh: make(chan Task, 1), retryCh: make(chan Task, 1),
mtr: metrics.Default(), mtr: metrics.New(prometheus.NewRegistry()),
} }
e.initTargets(nil) e.initTargets(nil)
@@ -427,7 +428,7 @@ func NewTestEngineWithDB(
deliveryCh: make(chan Task, deliveryChannelSize), deliveryCh: make(chan Task, deliveryChannelSize),
retryCh: make(chan Task, retryChannelSize), retryCh: make(chan Task, retryChannelSize),
workers: workers, workers: workers,
mtr: metrics.Default(), mtr: metrics.New(prometheus.NewRegistry()),
} }
e.initTargets(client) e.initTargets(client)
@@ -435,8 +436,7 @@ func NewTestEngineWithDB(
} }
// ExportSetMetrics substitutes the engine's metric set, so a test can // ExportSetMetrics substitutes the engine's metric set, so a test can
// assert on collectors registered on a private registry instead of // assert on collectors registered on a registry it holds.
// the process-wide ones every other test is also moving.
func (e *Engine) ExportSetMetrics(mtr *metrics.Set) { func (e *Engine) ExportSetMetrics(mtr *metrics.Set) {
e.mtr = mtr e.mtr = mtr
} }
+2 -3
View File
@@ -35,9 +35,8 @@ const (
) )
// mIsolate gives the setup's engine a metric set registered on a // mIsolate gives the setup's engine a metric set registered on a
// private registry. The process-wide collectors are moved by every // registry this test holds, so its exact assertions can gather from
// other delivery test running in parallel, so exact assertions are // it.
// only possible against a registry this test owns.
func mIsolate( func mIsolate(
t *testing.T, s iSetup, t *testing.T, s iSetup,
) *prometheus.Registry { ) *prometheus.Registry {
+1 -7
View File
@@ -1,8 +1,6 @@
package handlers package handlers
import ( import (
"time"
"sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/delivery"
) )
@@ -27,7 +25,7 @@ const maxRenderedResponseBytes = 4096
// cut, so an oversized stored response never becomes a Go // cut, so an oversized stored response never becomes a Go
// string at all. // string at all.
const deliveryResultColumns = "delivery_id, attempt_num, success, " + const deliveryResultColumns = "delivery_id, attempt_num, success, " +
"status_code, error, duration, created_at, " + "status_code, error, duration, " +
"substr(cast(response_body as blob), 1, ?) AS response_body, " + "substr(cast(response_body as blob), 1, ?) AS response_body, " +
"length(cast(response_body as blob)) AS response_bytes" "length(cast(response_body as blob)) AS response_bytes"
@@ -108,10 +106,6 @@ type deliveryResultRow struct {
Duration int64 Duration int64
ResponseBody []byte ResponseBody []byte
ResponseBytes int64 ResponseBytes int64
// CreatedAt is when the attempt's result was recorded, which
// is when the attempt finished.
CreatedAt time.Time
} }
// view projects a loaded row for rendering, stripping the // view projects a loaded row for rendering, stripping the
+5 -2
View File
@@ -12,6 +12,7 @@ import (
"net/http" "net/http"
"sync/atomic" "sync/atomic"
"github.com/prometheus/client_golang/prometheus"
"go.uber.org/fx" "go.uber.org/fx"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/delivery"
@@ -28,7 +29,7 @@ const (
// maxBodyShift is the bit shift for 1 MB body limit. // maxBodyShift is the bit shift for 1 MB body limit.
maxBodyShift = 20 maxBodyShift = 20
// recentEventLimit is the number of recent events to show. // recentEventLimit is the number of recent events to show.
recentEventLimit = 50 recentEventLimit = 20
// paginationPerPage is the number of items per page. // paginationPerPage is the number of items per page.
paginationPerPage = 25 paginationPerPage = 25
@@ -61,6 +62,8 @@ type HandlersParams struct {
Notifier delivery.Notifier Notifier delivery.Notifier
Evictor delivery.WebhookEvictor Evictor delivery.WebhookEvictor
SSRFGuard *delivery.Guard SSRFGuard *delivery.Guard
Metrics *metrics.Set
Registry *prometheus.Registry
} }
// Handlers provides HTTP handler methods for all application // Handlers provides HTTP handler methods for all application
@@ -122,7 +125,7 @@ func New(
s.mw = params.Middleware s.mw = params.Middleware
s.notifier = params.Notifier s.notifier = params.Notifier
s.evictor = params.Evictor s.evictor = params.Evictor
s.mtr = metrics.Default() s.mtr = params.Metrics
s.ssrf = params.SSRFGuard s.ssrf = params.SSRFGuard
// Parse all page templates once at startup // Parse all page templates once at startup
+3
View File
@@ -20,6 +20,7 @@ import (
"sneak.berlin/go/webhooker/internal/handlers" "sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck" "sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger" "sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware" "sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/session" "sneak.berlin/go/webhooker/internal/session"
) )
@@ -109,6 +110,8 @@ func newTestApp(
func(r *recordingEvictor) delivery.WebhookEvictor { func(r *recordingEvictor) delivery.WebhookEvictor {
return r return r
}, },
metrics.NewRegistry,
metrics.New,
middleware.New, middleware.New,
delivery.NewGuard, delivery.NewGuard,
handlers.New, handlers.New,
+20
View File
@@ -0,0 +1,20 @@
package handlers
import (
"net/http"
"github.com/prometheus/client_golang/prometheus/promhttp"
)
// HandleMetrics returns the Prometheus scrape handler for the
// registry every collector in this process registers on. It is what
// promhttp.Handler builds for the global default registry, including
// the promhttp_metric_handler_* series that count scrapes, pointed at
// that registry instead.
func (s *Handlers) HandleMetrics() http.HandlerFunc {
reg := s.params.Registry
return promhttp.InstrumentMetricHandler(
reg, promhttp.HandlerFor(reg, promhttp.HandlerOpts{}),
).ServeHTTP
}
-248
View File
@@ -1,248 +0,0 @@
package handlers
import (
"net/http"
"strconv"
"time"
"github.com/dustin/go-humanize"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
)
// recentEventColumns is the recent events list's projection. It
// reads the body's size and never the body itself, for the reason
// maxRenderedBodyBytes gives; the cast to blob makes length count
// bytes rather than characters.
const recentEventColumns = "id, created_at, method, content_type, " +
"resubmitted_from_id, length(cast(body as blob)) AS body_bytes"
// RecentEventView is one row of the recent events list on a
// webhook's page.
type RecentEventView struct {
Method string
ContentType string
// ResubmittedFromID names the event this one was copied from,
// empty for an event that arrived on the receiver.
ResubmittedFromID string
// Received is how long ago the event arrived, and ReceivedUTC
// the full timestamp the page shows on hover.
Received string
ReceivedUTC string
// Size is the size of the stored body.
Size string
// ProcessingTime is how long the event's slowest delivery
// took; see processingTime.
ProcessingTime string
// Status is what the webhook's HTTP target answered, and
// StatusClass its colour; see targetStatus. Both are empty
// unless the webhook has exactly one HTTP target.
Status string
StatusClass string
}
// recentEventRow is one row of recentEventColumns.
type recentEventRow struct {
ID string
CreatedAt time.Time
Method string
ContentType string
ResubmittedFromID *string
BodyBytes uint64
}
// singleHTTPTargetID returns the ID of the webhook's HTTP target
// when it has exactly one, and "" when it has none or several.
func singleHTTPTargetID(targets []database.Target) string {
id := ""
count := 0
for i := range targets {
if targets[i].Type == database.TargetTypeHTTP {
id = targets[i].ID
count++
}
}
if count != 1 {
return ""
}
return id
}
// loadRecentEvents loads the webhook's recentEventLimit newest
// events for its page, newest first. statusTargetID is the
// webhook's only HTTP target, or "" when the list shows no status.
func (h *Handlers) loadRecentEvents(
webhookDB *gorm.DB, webhookID, statusTargetID string,
) ([]RecentEventView, error) {
var rows []recentEventRow
err := webhookDB.Model(&database.Event{}).
Select(recentEventColumns).
Where("webhook_id = ?", webhookID).
Order("created_at DESC").
Limit(recentEventLimit).
Find(&rows).Error
if err != nil {
return nil, err
}
eventIDs := make([]string, len(rows))
for i := range rows {
eventIDs[i] = rows[i].ID
}
// Oldest first, so an event's last delivery to a target is its
// newest: a replay adds a delivery rather than changing the
// earlier one.
var deliveries []database.Delivery
err = webhookDB.
Select("id, event_id, target_id, status, created_at").
Where("event_id IN ?", eventIDs).
Order("created_at ASC").
Find(&deliveries).Error
if err != nil {
return nil, err
}
byEvent := make(map[string][]database.Delivery, len(rows))
deliveryIDs := make([]string, len(deliveries))
for i := range deliveries {
eventID := deliveries[i].EventID
byEvent[eventID] = append(byEvent[eventID], deliveries[i])
deliveryIDs[i] = deliveries[i].ID
}
attempts, err := h.loadDeliveryResults(webhookDB, deliveryIDs)
if err != nil {
return nil, err
}
views := make([]RecentEventView, len(rows))
for i := range rows {
views[i] = rows[i].view(
byEvent[rows[i].ID], attempts, statusTargetID,
)
}
return views, nil
}
// view projects a loaded row for rendering. deliveries is the
// event's deliveries, oldest first, and attempts their recorded
// attempts keyed by delivery ID.
func (r *recentEventRow) view(
deliveries []database.Delivery,
attempts map[string][]deliveryResultRow,
statusTargetID string,
) RecentEventView {
v := RecentEventView{
Method: r.Method,
ContentType: r.ContentType,
Received: humanize.Time(r.CreatedAt),
ReceivedUTC: r.CreatedAt.UTC().Format(time.DateTime) + " UTC",
Size: humanize.Bytes(r.BodyBytes),
ProcessingTime: processingTime(deliveries, attempts),
}
if r.ResubmittedFromID != nil {
v.ResubmittedFromID = *r.ResubmittedFromID
}
if statusTargetID != "" {
v.Status, v.StatusClass = targetStatus(
deliveries, attempts, statusTargetID,
)
}
return v
}
// processingTime is how long the event's slowest delivery took,
// from being queued to its last recorded attempt, time spent
// waiting between retries included. A delivery is queued when its
// event is received, or when an operator replays it, so a replay
// is timed from the replay rather than from the event's arrival.
// It is "in progress" while any delivery is pending or retrying,
// and empty for an event with no deliveries.
func processingTime(
deliveries []database.Delivery,
attempts map[string][]deliveryResultRow,
) string {
if len(deliveries) == 0 {
return ""
}
var slowest time.Duration
for i := range deliveries {
if !deliveries[i].Status.Terminal() {
return "in progress"
}
tries := attempts[deliveries[i].ID]
if len(tries) == 0 {
continue
}
last := tries[len(tries)-1].CreatedAt
slowest = max(slowest, last.Sub(deliveries[i].CreatedAt))
}
return slowest.Round(time.Millisecond).String()
}
// targetStatus is what the target answered for the event, and the
// colour to show it in: the HTTP status code of the last attempt of
// the event's newest delivery to the target. Without a code it is
// "no response" when that attempt failed before a response
// arrived, the delivery's status ("pending") before any attempt,
// and "not sent" when the event has no delivery to the target.
func targetStatus(
deliveries []database.Delivery,
attempts map[string][]deliveryResultRow,
targetID string,
) (string, string) {
newest := -1
for i := range deliveries {
if deliveries[i].TargetID == targetID {
newest = i
}
}
if newest < 0 {
return "not sent", "text-gray-400"
}
tries := attempts[deliveries[newest].ID]
if len(tries) == 0 {
return string(deliveries[newest].Status), "text-gray-400"
}
code := tries[len(tries)-1].StatusCode
switch {
case code == 0:
return "no response", "text-red-600"
case code >= http.StatusInternalServerError:
return strconv.Itoa(code), "text-red-600"
case code >= http.StatusBadRequest:
return strconv.Itoa(code), "text-yellow-600"
case code >= http.StatusMultipleChoices:
return strconv.Itoa(code), "text-gray-500"
case code >= http.StatusOK:
return strconv.Itoa(code), "text-green-600"
default:
return strconv.Itoa(code), "text-gray-500"
}
}
-295
View File
@@ -1,295 +0,0 @@
package handlers_test
import (
"fmt"
"net/http"
"strings"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/session"
)
// statusTitle marks the status column's cell in a recent events
// row; it is absent from the page when the column is not shown.
const statusTitle = `title="HTTP status from the HTTP target"`
// recentEventsFixture is one started app and a webhook whose
// recent events list a test fills.
type recentEventsFixture struct {
h *handlers.Handlers
sess *session.Session
db *database.Database
webhook *database.Webhook
webhookDB *gorm.DB
}
func newRecentEventsFixture(t *testing.T) *recentEventsFixture {
t.Helper()
f := &recentEventsFixture{}
var dbMgr *database.WebhookDBManager
app := newTestApp(t, &f.h, &f.sess, &f.db, &dbMgr)
app.RequireStart()
t.Cleanup(app.RequireStop)
f.webhook = seedWebhook(t, f.db)
webhookDB, err := dbMgr.GetDB(f.webhook.ID)
require.NoError(t, err)
f.webhookDB = webhookDB
return f
}
func (f *recentEventsFixture) render(t *testing.T) string {
t.Helper()
return renderSourceDetailPage(t, f.h, f.sess, f.webhook.ID)
}
// event records an event received at receivedAt.
func (f *recentEventsFixture) event(
t *testing.T, contentType, body string, receivedAt time.Time,
) *database.Event {
t.Helper()
event := &database.Event{
WebhookID: f.webhook.ID,
Method: http.MethodPost,
Body: body,
ContentType: contentType,
}
event.CreatedAt = receivedAt
require.NoError(t, f.webhookDB.Omit(
clause.Associations,
).Create(event).Error)
return event
}
// delivery records a delivery of the event to the target, queued
// when the event was received.
func (f *recentEventsFixture) delivery(
t *testing.T,
event *database.Event,
targetID string,
status database.DeliveryStatus,
) *database.Delivery {
t.Helper()
return f.deliveryQueuedAt(
t, event, targetID, status, event.CreatedAt,
)
}
// deliveryQueuedAt records a delivery of the event to the target,
// queued at queuedAt, as a replay is.
func (f *recentEventsFixture) deliveryQueuedAt(
t *testing.T,
event *database.Event,
targetID string,
status database.DeliveryStatus,
queuedAt time.Time,
) *database.Delivery {
t.Helper()
dlv := &database.Delivery{
EventID: event.ID,
TargetID: targetID,
Status: status,
}
dlv.CreatedAt = queuedAt
require.NoError(t, f.webhookDB.Omit(
clause.Associations,
).Create(dlv).Error)
return dlv
}
// attempt records one attempt of the delivery that finished took
// after the delivery was queued, with HTTP status code (0 for no
// response).
func (f *recentEventsFixture) attempt(
t *testing.T, dlv *database.Delivery, code int, took time.Duration,
) {
t.Helper()
result := &database.DeliveryResult{
DeliveryID: dlv.ID,
AttemptNum: 1,
StatusCode: code,
}
result.CreatedAt = dlv.CreatedAt.Add(took)
require.NoError(t, f.webhookDB.Omit(
clause.Associations,
).Create(result).Error)
}
// statusCell is the status column's cell as the page renders it.
func statusCell(class, text string) string {
return `<span class="font-medium ` + class + `" ` + statusTitle +
`>` + text + `</span>`
}
// TestHandleSourceDetail_ShowsFiftyNewestEvents proves the list
// holds the 50 newest events, newest first, and not one more.
func TestHandleSourceDetail_ShowsFiftyNewestEvents(t *testing.T) {
t.Parallel()
f := newRecentEventsFixture(t)
base := time.Now().Add(-time.Hour)
for i := range 51 {
f.event(
t, fmt.Sprintf("application/x-recent-%02d", i), "{}",
base.Add(time.Duration(i)*time.Second),
)
}
body := f.render(t)
assert.Equal(t, 50, strings.Count(body, `title="Body size"`))
assert.NotContains(t, body, "application/x-recent-00")
assert.Contains(t, body, "application/x-recent-01")
assert.Less(
t,
strings.Index(body, "application/x-recent-50"),
strings.Index(body, "application/x-recent-49"),
)
}
// TestHandleSourceDetail_RecentEventColumns proves a row shows its
// time relative with the UTC timestamp on hover, its body size,
// and its processing time once every delivery has finished.
func TestHandleSourceDetail_RecentEventColumns(t *testing.T) {
t.Parallel()
f := newRecentEventsFixture(t)
logTarget := seedTarget(t, f.db, f.webhook.ID, database.TargetTypeLog)
receivedAt := time.Now().Add(-210 * time.Second).
UTC().Truncate(time.Second)
done := f.event(
t, contentTypeJSON, strings.Repeat("x", 2048), receivedAt,
)
f.attempt(
t,
f.delivery(t, done, logTarget.ID, database.DeliveryStatusDelivered),
0, 1500*time.Millisecond,
)
waiting := f.event(t, "text/plain", "{}", receivedAt)
f.delivery(t, waiting, logTarget.ID, database.DeliveryStatusPending)
body := f.render(t)
assert.Contains(
t, body,
`<span title="`+receivedAt.Format(time.DateTime)+
` UTC">3 minutes ago</span>`,
)
assert.Contains(t, body, `<span title="Body size">2.0 kB</span>`)
assert.Contains(t, body, ">1.5s</span>")
assert.Contains(t, body, ">in progress</span>")
}
// TestHandleSourceDetail_StatusWithSingleHTTPTarget proves that a
// webhook with exactly one HTTP target shows, colour-coded, what
// that target answered for each event. The log target beside it
// does not count against "exactly one".
func TestHandleSourceDetail_StatusWithSingleHTTPTarget(t *testing.T) {
t.Parallel()
f := newRecentEventsFixture(t)
target := seedTarget(t, f.db, f.webhook.ID, database.TargetTypeHTTP)
seedTarget(t, f.db, f.webhook.ID, database.TargetTypeLog)
now := time.Now()
for _, code := range []int{204, 302, 404, 503, 0} {
dlv := f.delivery(
t, f.event(t, contentTypeJSON, "{}", now), target.ID,
database.DeliveryStatusDelivered,
)
f.attempt(t, dlv, code, time.Second)
}
f.delivery(
t, f.event(t, contentTypeJSON, "{}", now), target.ID,
database.DeliveryStatusPending,
)
f.event(t, contentTypeJSON, "{}", now)
// A replay is a newer delivery, and its answer is the one shown.
replayed := f.event(t, contentTypeJSON, "{}", now)
f.attempt(t, f.delivery(
t, replayed, target.ID, database.DeliveryStatusFailed,
), 502, time.Second)
f.attempt(t, f.deliveryQueuedAt(
t, replayed, target.ID, database.DeliveryStatusDelivered,
now.Add(time.Minute),
), 200, time.Second)
body := f.render(t)
assert.Contains(t, body, statusCell("text-green-600", "204"))
assert.Contains(t, body, statusCell("text-gray-500", "302"))
assert.Contains(t, body, statusCell("text-yellow-600", "404"))
assert.Contains(t, body, statusCell("text-red-600", "503"))
assert.Contains(t, body, statusCell("text-red-600", "no response"))
assert.Contains(t, body, statusCell("text-gray-400", "pending"))
assert.Contains(t, body, statusCell("text-gray-400", "not sent"))
assert.Contains(t, body, statusCell("text-green-600", "200"))
assert.NotContains(t, body, ">502<")
}
// TestHandleSourceDetail_NoStatusWithoutSingleHTTPTarget proves the
// status column is absent when the webhook has no HTTP target or
// more than one.
func TestHandleSourceDetail_NoStatusWithoutSingleHTTPTarget(
t *testing.T,
) {
t.Parallel()
cases := map[string][]database.TargetType{
"none": {database.TargetTypeLog},
"several": {database.TargetTypeHTTP, database.TargetTypeHTTP},
}
for name, types := range cases {
t.Run(name, func(t *testing.T) {
t.Parallel()
f := newRecentEventsFixture(t)
event := f.event(t, contentTypeJSON, "{}", time.Now())
for _, tt := range types {
target := seedTarget(t, f.db, f.webhook.ID, tt)
f.attempt(t, f.delivery(
t, event, target.ID,
database.DeliveryStatusDelivered,
), 200, time.Second)
}
body := f.render(t)
assert.Contains(t, body, `title="Body size"`)
assert.NotContains(t, body, statusTitle)
})
}
}
+6 -11
View File
@@ -415,21 +415,16 @@ func (h *Handlers) renderSourceDetail(
"webhook_id = ?", webhook.ID, "webhook_id = ?", webhook.ID,
).Find(&targets) ).Find(&targets)
var events []RecentEventView var events []database.Event
if h.dbMgr.DBExists(webhook.ID) { if h.dbMgr.DBExists(webhook.ID) {
webhookDB, dbErr := h.dbMgr.GetDB(webhook.ID) webhookDB, dbErr := h.dbMgr.GetDB(webhook.ID)
if dbErr == nil { if dbErr == nil {
events, dbErr = h.loadRecentEvents( webhookDB.Where(
webhookDB, webhook.ID, singleHTTPTargetID(targets), "webhook_id = ?", webhook.ID,
) ).Order("created_at DESC").Limit(
} recentEventLimit,
).Find(&events)
if dbErr != nil {
h.log.Error(
"failed to load recent events",
"webhook_id", webhook.ID, "error", dbErr,
)
} }
} }
+26 -20
View File
@@ -3,17 +3,17 @@
// deliveries are attempted, how they end, how long they take, how // deliveries are attempted, how they end, how long they take, how
// deep the queues are, and how many circuit breakers are open. // deep the queues are, and how many circuit breakers are open.
// //
// The inbound HTTP metrics come from the go-http-metrics recorder in // It also builds the registry the authenticated /metrics route
// internal/middleware and land on prometheus.DefaultRegisterer. These // serves. These collectors, the inbound HTTP metrics recorded in
// collectors register there too, so both surfaces are gathered by the // internal/middleware, and the Go runtime and process collectors all
// one promhttp handler mounted on the authenticated /metrics route. // register on that one registry, never on Prometheus's global default.
package metrics package metrics
import ( import (
"sync"
"time" "time"
"github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/collectors"
"github.com/prometheus/client_golang/prometheus/promauto" "github.com/prometheus/client_golang/prometheus/promauto"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
) )
@@ -57,25 +57,31 @@ var knownTargetTypes = []database.TargetType{
database.TargetTypeSlack, database.TargetTypeSlack,
} }
// defaultSet is the process-wide metric set, registered on the same // NewRegistry returns the registry /metrics serves, carrying the Go
// registry the HTTP middleware and the /metrics handler already use. // runtime and process collectors that Prometheus's global default
// It is built on first use rather than in an init so that a test // registry carries, so the go_* and process_* series stay in the
// binary that never touches metrics never registers them. // scrape.
// //
//nolint:gochecknoglobals // one process-wide registration, by design // A registry of its own, rather than the global default, is what lets
var defaultSet = sync.OnceValue(func() *Set { // two dependency graphs in one process — two tests, say — each
return New(prometheus.DefaultRegisterer) // register their collectors without the second registration
}) // panicking.
func NewRegistry() *prometheus.Registry {
reg := prometheus.NewRegistry()
reg.MustRegister(
collectors.NewGoCollector(),
collectors.NewProcessCollector(
collectors.ProcessCollectorOpts{},
),
)
// Default returns the process-wide metric set. return reg
func Default() *Set {
return defaultSet()
} }
// Set is one registered group of webhooker's delivery collectors. // Set is one registered group of webhooker's delivery collectors.
// Production uses the single Default set; tests build their own // Production builds one on the registry /metrics serves; tests build
// against a private registry so assertions are not disturbed by // their own against a private registry so assertions are not
// deliveries other tests are making concurrently. // disturbed by deliveries other tests are making concurrently.
type Set struct { type Set struct {
eventsReceived prometheus.Counter eventsReceived prometheus.Counter
deliveryAttempts *prometheus.CounterVec deliveryAttempts *prometheus.CounterVec
@@ -93,7 +99,7 @@ type Set struct {
// New registers a full set of delivery collectors on reg and returns // New registers a full set of delivery collectors on reg and returns
// it. It panics if reg already holds them, which is the intended // it. It panics if reg already holds them, which is the intended
// behaviour for a duplicate registration. // behaviour for a duplicate registration.
func New(reg prometheus.Registerer) *Set { func New(reg *prometheus.Registry) *Set {
factory := promauto.With(reg) factory := promauto.With(reg)
s := &Set{ s := &Set{
+1 -2
View File
@@ -10,8 +10,7 @@ import (
// MetricsMiddlewareForTest builds the metrics recording middleware // MetricsMiddlewareForTest builds the metrics recording middleware
// against a caller-supplied recorder, so a test can gather from its // against a caller-supplied recorder, so a test can gather from its
// own Prometheus registry rather than the process-wide default one // own Prometheus registry without building a whole Middleware.
// that Middleware.Metrics uses.
func MetricsMiddlewareForTest( func MetricsMiddlewareForTest(
rec httpmetrics.Recorder, rec httpmetrics.Recorder,
) func(http.Handler) http.Handler { ) func(http.Handler) http.Handler {
+4 -7
View File
@@ -7,7 +7,6 @@ import (
"github.com/go-chi/chi" "github.com/go-chi/chi"
httpmetrics "github.com/slok/go-http-metrics/metrics" httpmetrics "github.com/slok/go-http-metrics/metrics"
prommetrics "github.com/slok/go-http-metrics/metrics/prometheus"
ghmm "github.com/slok/go-http-metrics/middleware" ghmm "github.com/slok/go-http-metrics/middleware"
"github.com/slok/go-http-metrics/middleware/std" "github.com/slok/go-http-metrics/middleware/std"
) )
@@ -152,16 +151,14 @@ func (r boundedLabelRecorder) AddInflightRequests(
var _ httpmetrics.Recorder = boundedLabelRecorder{} var _ httpmetrics.Recorder = boundedLabelRecorder{}
// Metrics returns middleware that records Prometheus HTTP metrics on // Metrics returns middleware that records Prometheus HTTP metrics on
// the default registry, which is the one the /metrics route gathers. // the registry the /metrics route serves. Every call shares the one
// recorder New built, so any number of routers can install it.
func (s *Middleware) Metrics() func(http.Handler) http.Handler { func (s *Middleware) Metrics() func(http.Handler) http.Handler {
return metricsMiddleware( return metricsMiddleware(s.metricsRecorder)
prommetrics.NewRecorder(prommetrics.Config{}),
)
} }
// metricsMiddleware builds the recording middleware against a given // metricsMiddleware builds the recording middleware against a given
// recorder, so tests can gather from a registry of their own instead // recorder, so tests can gather from a registry of their own.
// of the process-wide default.
func metricsMiddleware( func metricsMiddleware(
rec httpmetrics.Recorder, rec httpmetrics.Recorder,
) func(http.Handler) http.Handler { ) func(http.Handler) http.Handler {
+2 -3
View File
@@ -57,9 +57,8 @@ const (
// Server.setupWebhookRoutes inside it. That ordering is the whole // Server.setupWebhookRoutes inside it. That ordering is the whole
// defect, so a test that flattens it would prove nothing. // defect, so a test that flattens it would prove nothing.
// //
// The recorder writes to a registry of the test's own rather than the // The recorder writes to a registry of the test's own, so each test
// process-wide default one, so each test observes only its own // observes only its own traffic.
// traffic.
func metricsTestRouter( func metricsTestRouter(
t *testing.T, t *testing.T,
receiverLimit int, receiverLimit int,
+17 -4
View File
@@ -13,6 +13,9 @@ import (
"github.com/go-chi/chi" "github.com/go-chi/chi"
"github.com/go-chi/chi/middleware" "github.com/go-chi/chi/middleware"
"github.com/go-chi/cors" "github.com/go-chi/cors"
"github.com/prometheus/client_golang/prometheus"
httpmetrics "github.com/slok/go-http-metrics/metrics"
prommetrics "github.com/slok/go-http-metrics/metrics/prometheus"
"go.uber.org/fx" "go.uber.org/fx"
"sneak.berlin/go/webhooker/internal/config" "sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/globals" "sneak.berlin/go/webhooker/internal/globals"
@@ -148,10 +151,11 @@ const (
type MiddlewareParams struct { type MiddlewareParams struct {
fx.In fx.In
Logger *logger.Logger Logger *logger.Logger
Globals *globals.Globals Globals *globals.Globals
Config *config.Config Config *config.Config
Session *session.Session Session *session.Session
Registry *prometheus.Registry
} }
// Middleware provides HTTP middleware for logging, CORS, auth, and // Middleware provides HTTP middleware for logging, CORS, auth, and
@@ -161,6 +165,12 @@ type Middleware struct {
params *MiddlewareParams params *MiddlewareParams
session *session.Session session *session.Session
// metricsRecorder records the inbound HTTP metrics on the
// registry /metrics serves. It is built once, in New, because
// building it registers its collectors, and a second
// registration on the same registry panics; see Metrics.
metricsRecorder httpmetrics.Recorder
// loginGuard counts failed credential verifications and bounds // loginGuard counts failed credential verifications and bounds
// concurrent password hashing. It is built on first use so that // concurrent password hashing. It is built on first use so that
// every construction path gets one; see guard(). // every construction path gets one; see guard().
@@ -179,6 +189,9 @@ func New(
s.params = &params s.params = &params
s.log = params.Logger.Get() s.log = params.Logger.Get()
s.session = params.Session s.session = params.Session
s.metricsRecorder = prommetrics.NewRecorder(
prommetrics.Config{Registry: params.Registry},
)
return s, nil return s, nil
} }
+3
View File
@@ -24,6 +24,7 @@ import (
"sneak.berlin/go/webhooker/internal/handlers" "sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck" "sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger" "sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware" "sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/resetpw" "sneak.berlin/go/webhooker/internal/resetpw"
"sneak.berlin/go/webhooker/internal/session" "sneak.berlin/go/webhooker/internal/session"
@@ -163,6 +164,8 @@ func newServerApp(
session.New, session.New,
func() delivery.Notifier { return &noopNotifier{} }, func() delivery.Notifier { return &noopNotifier{} },
func() delivery.WebhookEvictor { return &noopEvictor{} }, func() delivery.WebhookEvictor { return &noopEvictor{} },
metrics.NewRegistry,
metrics.New,
middleware.New, middleware.New,
delivery.NewGuard, delivery.NewGuard,
handlers.New, handlers.New,
+1 -7
View File
@@ -7,7 +7,6 @@ import (
sentryhttp "github.com/getsentry/sentry-go/http" sentryhttp "github.com/getsentry/sentry-go/http"
"github.com/go-chi/chi" "github.com/go-chi/chi"
"github.com/go-chi/chi/middleware" "github.com/go-chi/chi/middleware"
"github.com/prometheus/client_golang/prometheus/promhttp"
"sneak.berlin/go/webhooker/static" "sneak.berlin/go/webhooker/static"
) )
@@ -130,12 +129,7 @@ func (s *Server) setupRoutes() {
if s.params.Config.MetricsAuthEnabled() { if s.params.Config.MetricsAuthEnabled() {
s.router.Group(func(r chi.Router) { s.router.Group(func(r chi.Router) {
r.Use(s.mw.MetricsAuth()) r.Use(s.mw.MetricsAuth())
r.Get( r.Get("/metrics", s.h.HandleMetrics())
"/metrics",
http.HandlerFunc(
promhttp.Handler().ServeHTTP,
),
)
}) })
} }
+43
View File
@@ -23,6 +23,7 @@ import (
"sneak.berlin/go/webhooker/internal/handlers" "sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck" "sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger" "sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware" "sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/server" "sneak.berlin/go/webhooker/internal/server"
"sneak.berlin/go/webhooker/internal/session" "sneak.berlin/go/webhooker/internal/session"
@@ -112,6 +113,8 @@ func newTestEnvWithConfig(
session.New, session.New,
func() delivery.Notifier { return &noopNotifier{} }, func() delivery.Notifier { return &noopNotifier{} },
func() delivery.WebhookEvictor { return &noopEvictor{} }, func() delivery.WebhookEvictor { return &noopEvictor{} },
metrics.NewRegistry,
metrics.New,
middleware.New, middleware.New,
delivery.NewGuard, delivery.NewGuard,
handlers.New, handlers.New,
@@ -961,3 +964,43 @@ func TestMetricsRouteUnmountedOnHalfSetConfig(t *testing.T) {
}) })
} }
} }
// TestTwoMetricsRoutersInOneProcess pins
// https://git.eeqj.de/sneak/webhooker/issues/227: a second
// metrics-enabled router in one process used to panic, because the
// HTTP metrics registered on Prometheus's global default registry.
// Two routers are built over separate dependency graphs and a third
// over the first graph again, and each must still serve the HTTP,
// delivery and Go runtime series.
func TestTwoMetricsRoutersInOneProcess(t *testing.T) {
t.Parallel()
first := newTestEnvWithConfig(
t, metricsConfig(t, metricsUser, metricsAuthValue),
)
second := newTestEnvWithConfig(
t, metricsConfig(t, metricsUser, metricsAuthValue),
)
third := &testEnv{
router: server.NewRouterForTest(
first.log.Get(), first.cfg, first.mw, first.hnd,
),
}
for _, env := range []*testEnv{first, second, third} {
env.get("/", nil)
scrape := env.metricsRequest(metricsUser, metricsAuthValue)
require.Equal(t, http.StatusOK, scrape.Code)
for _, series := range []string{
"http_request_duration_seconds",
"http_response_size_bytes",
"http_requests_inflight",
"webhooker_events_received_total",
"go_goroutines",
} {
assert.Contains(t, scrape.Body.String(), series)
}
}
}
+4 -16
View File
@@ -187,24 +187,12 @@
<div class="divide-y divide-gray-100"> <div class="divide-y divide-gray-100">
{{range .Events}} {{range .Events}}
<div class="p-4"> <div class="p-4">
<div class="flex flex-wrap items-center justify-between gap-3"> <div class="flex items-center justify-between">
<div class="flex flex-wrap items-center gap-3"> <div class="flex items-center gap-3">
<span class="badge-info">{{.Method}}</span> <span class="badge-info">{{.Method}}</span>
<span class="text-sm text-gray-500 break-all">{{.ContentType}}</span> <span class="text-sm text-gray-500">{{.ContentType}}</span>
{{if .ResubmittedFromID}}
<span class="text-xs text-gray-500" title="This event is a copy of {{.ResubmittedFromID}}">resubmitted copy</span>
{{end}}
</div>
<div class="flex flex-wrap items-center gap-3 text-xs text-gray-400">
<span title="Body size">{{.Size}}</span>
{{if .ProcessingTime}}
<span title="Processing time: how long the slowest delivery took, from being queued to its last attempt">{{.ProcessingTime}}</span>
{{end}}
{{if .Status}}
<span class="font-medium {{.StatusClass}}" title="HTTP status from the HTTP target">{{.Status}}</span>
{{end}}
<span title="{{.ReceivedUTC}}">{{.Received}}</span>
</div> </div>
<span class="text-xs text-gray-400">{{.CreatedAt.Format "2006-01-02 15:04:05 UTC"}}</span>
</div> </div>
</div> </div>
{{else}} {{else}}