Author SHA1 Message Date
clawbot a4337e949b Show each webhook's activity in the webhook list (closes #394)
check / check (push) Waiting to run
Each entry in the webhook list now shows when the webhook's last event
arrived, or "No events yet", and how many of its deliveries failed in
the last 24 hours, in red when that is not zero. The entrypoint and
target counts say how many are inactive.

The figures come from the statistics pane's own reads: the event totals
row, which now also gives the event count, shown as events within
retention, instead of counting every stored event, and the pane's query
over the deliveries finished in the last 24 hours. A webhook whose event
database cannot be read says so in its entry rather than showing zeros.

Model: opus-5-5
2026-10-02 07:34:37 +00:00
clawbot eb4c4cc849 Serve /metrics from a registry of its own (closes #227)
check / check (push) Waiting to run
A second metrics-enabled router in one process panicked on a duplicate collector registration, because every collector registered on Prometheus's global default registry. metrics.NewRegistry now builds one registry with the Go runtime and process collectors; fx provides it and the delivery metric set built on it. The middleware builds its HTTP recorder once on that registry (NewForTest on a fresh one), the engine and handlers take the metric set from fx, and nothing registers on the global default any more.

/metrics is served from the new registry with the same series names, labels and auth. A test builds two metrics-enabled routers in one process.

Model: opus-5-5
2026-10-02 09:06:20 +02:00
20 changed files with 688 additions and 92 deletions
+18 -8
View File
@@ -1768,7 +1768,7 @@ retries) is individually logged for full observability.
#### EventTotals and TargetTotals #### EventTotals and TargetTotals
Running counts in each event database, read by the statistics pane at the Running counts in each event database, read by the statistics pane at the
top of the webhook page. `EventTotals` is one row: top of the webhook page and by the webhook list. `EventTotals` is one row:
| Field | Type | Description | | Field | Type | Description |
| ---------------- | --------- | ----------- | | ---------------- | --------- | ----------- |
@@ -1800,6 +1800,14 @@ target. Its failure percentage for a window is the deliveries that became
`failed` in it out of all that became `delivered` or `failed` in it, and `failed` in it out of all that became `delivered` or `failed` in it, and
a dash when none did. a dash when none did.
The webhook list at `/hooks` shows three of the pane's figures for each
webhook: its events within retention and its last event, both from
`EventTotals`, and its deliveries that failed in the last 24 hours,
counted with the pane's query. It opens each webhook's event database once
(the handle stays open) and runs those two reads there, so its cost grows
with the number of webhooks and, for each, with the deliveries that
finished in the last 24 hours, never with the events stored.
#### Event-tier indexes #### Event-tier indexes
These indexes on the per-webhook event databases are declared in the model These indexes on the per-webhook event databases are declared in the model
@@ -1807,7 +1815,7 @@ tags, so `AutoMigrate` creates them on a fresh database:
| Table | Columns | Serves | | Table | Columns | Serves |
| ------------------ | --------------------------- | ------ | | ------------------ | --------------------------- | ------ |
| `deliveries` | `status`, `deleted_at`, `finished_at`, `target_id` | Startup recovery, the retry and pending sweeps every 60 seconds and the queue-depth sampler every 30 seconds, which select deliveries by status, and the webhook page's statistics, which count each target's deliveries by status and when they finished | | `deliveries` | `status`, `deleted_at`, `finished_at`, `target_id` | Startup recovery, the retry and pending sweeps every 60 seconds and the queue-depth sampler every 30 seconds, which select deliveries by status, and the webhook page's statistics and the webhook list, which count each target's deliveries by status and when they finished |
| `deliveries` | `event_id`, `deleted_at` | The event log, which loads each event's deliveries, and retention, which counts and deletes the deliveries of expired events | | `deliveries` | `event_id`, `deleted_at` | The event log, which loads each event's deliveries, and retention, which counts and deletes the deliveries of expired events |
| `delivery_results` | `delivery_id`, `deleted_at` | The event log, which loads the attempts of a page's deliveries, and retention, which deletes the attempts of expired events | | `delivery_results` | `delivery_id`, `deleted_at` | The event log, which loads the attempts of a page's deliveries, and retention, which deletes the attempts of expired events |
| `events` | `deleted_at`, `created_at` | The webhook page's statistics, which count recent events | | `events` | `deleted_at`, `created_at` | The webhook page's statistics, which count recent events |
@@ -2951,13 +2959,15 @@ Components are wired via Uber fx in this order:
7. `healthcheck.New` — Health check service 7. `healthcheck.New` — Health check service
8. `session.New` — Cookie-based session manager (key from database) 8. `session.New` — Cookie-based session manager (key from database)
9. `handlers.New` — HTTP handlers 9. `handlers.New` — HTTP handlers
10. `middleware.New` — HTTP middleware 10. `metrics.NewRegistry` — The registry `/metrics` serves
11. `delivery.New` — Event-driven delivery engine 11. `metrics.New` — The delivery collectors, registered on that registry
12. `delivery.NewArchiveSweeper` — Periodic pruning of idle archives 12. `middleware.New` — HTTP middleware
13. `delivery.Engine` → `delivery.Notifier` — interface bridge 13. `delivery.New` — Event-driven delivery engine
14. `delivery.Engine` → `delivery.WebhookEvictor` — interface bridge so 14. `delivery.NewArchiveSweeper` — Periodic pruning of idle archives
15. `delivery.Engine` → `delivery.Notifier` — interface bridge
16. `delivery.Engine` → `delivery.WebhookEvictor` — interface bridge so
deleting a webhook releases its archive writer deleting a webhook releases its archive writer
15. `server.New` — HTTP server and router 17. `server.New` — HTTP server and router
The server starts via `fx.Invoke(func(*server.Server, *delivery.Engine, The server starts via `fx.Invoke(func(*server.Server, *delivery.Engine,
*database.RetentionReaper, *delivery.ArchiveSweeper) {})`, which *database.RetentionReaper, *delivery.ArchiveSweeper) {})`, which
+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
+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"
@@ -399,7 +400,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)
@@ -414,7 +415,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)
@@ -437,7 +438,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)
@@ -445,8 +446,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 {
+4 -1
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"
@@ -64,6 +65,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
@@ -129,7 +132,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,
+21
View File
@@ -0,0 +1,21 @@
package handlers
import (
"net/http"
"github.com/prometheus/client_golang/prometheus/promhttp"
)
// HandleMetrics returns the Prometheus scrape handler for the
// registry built by metrics.NewRegistry, which the HTTP, delivery, Go
// runtime and process collectors register 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
}
+355
View File
@@ -0,0 +1,355 @@
package handlers_test
import (
"net/http"
"net/http/httptest"
"regexp"
"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/logger"
"sneak.berlin/go/webhooker/internal/session"
)
// failedHighlight is how the list marks a number of failed deliveries
// that is not zero.
const failedHighlight = `class="font-medium text-red-600"`
// listWebhook adds a webhook with the given name, owned by the test
// user.
func listWebhook(
t *testing.T, db *database.Database, name string,
) *database.Webhook {
t.Helper()
wh := &database.Webhook{UserID: deleteTestUserID, Name: name}
require.NoError(t, db.DB().Omit(clause.Associations).Create(wh).Error)
return wh
}
// renderWebhookList runs the real webhook list handler as the test user
// and returns the rendered page.
func renderWebhookList(
t *testing.T, h *handlers.Handlers, sess *session.Session,
) string {
t.Helper()
cookies := authenticatedCookies(
t, sess, deleteTestUserID, deleteTestUsername,
)
w := httptest.NewRecorder()
h.HandleSourceList().ServeHTTP(
w, getRequest(t, "/hooks", cookies, nil),
)
require.Equal(t, http.StatusOK, w.Code)
return w.Body.String()
}
// listCard returns one webhook's entry in a rendered webhook list, its
// markup as rendered and its text with the markup taken out and each
// run of space made one space.
func listCard(t *testing.T, page, webhookID string) (string, string) {
t.Helper()
_, card, found := strings.Cut(page, `href="/hook/`+webhookID+`"`)
require.True(t, found, "the list has no entry for %s", webhookID)
card, _, _ = strings.Cut(card, "</a>")
text := regexp.MustCompile(`<[^>]*>`).ReplaceAllString(card, " ")
return card, strings.Join(strings.Fields(text), " ")
}
// receiveEvents posts the given number of events to an entrypoint
// through the real receiver, and returns the webhook's event database
// and its events, oldest first.
func receiveEvents(
t *testing.T,
h *handlers.Handlers,
dbMgr *database.WebhookDBManager,
webhookID, path string,
count int,
) (*gorm.DB, []database.Event) {
t.Helper()
router := receiverRouter(h)
for range count {
require.Equal(t, http.StatusOK, postReceiver(t, router, path))
}
webhookDB, err := dbMgr.GetDB(webhookID)
require.NoError(t, err)
events := listEvents(t, webhookDB)
require.Len(t, events, count)
return webhookDB, events
}
// seedFailingWebhook adds a webhook with two entrypoints, one inactive,
// and four targets, one inactive. Three events each reach the three
// active targets. Two deliveries failed in the last 24 hours, one 30
// hours ago, and one was delivered. It returns the webhook and its
// newest event.
func seedFailingWebhook(
t *testing.T,
h *handlers.Handlers,
db *database.Database,
dbMgr *database.WebhookDBManager,
) (*database.Webhook, database.Event) {
t.Helper()
wh := listWebhook(t, db, "failing")
path := statsEntrypoint(t, db, wh.ID, true)
statsEntrypoint(t, db, wh.ID, false)
first := seedTarget(t, db, wh.ID, database.TargetTypeLog)
second := seedTarget(t, db, wh.ID, database.TargetTypeLog)
seedTarget(t, db, wh.ID, database.TargetTypeLog)
inactive := seedTarget(t, db, wh.ID, database.TargetTypeLog)
require.NoError(t, db.DB().Model(inactive).
Update("active", false).Error)
webhookDB, events := receiveEvents(t, h, dbMgr, wh.ID, path, 3)
now := time.Now()
statsFinish(t, webhookDB,
statsDelivery(t, webhookDB, events[0].ID, first.ID),
database.DeliveryStatusFailed, now.Add(-30*time.Hour))
statsFinish(t, webhookDB,
statsDelivery(t, webhookDB, events[1].ID, first.ID),
database.DeliveryStatusFailed, now.Add(-time.Hour))
statsFinish(t, webhookDB,
statsDelivery(t, webhookDB, events[2].ID, first.ID),
database.DeliveryStatusFailed, now.Add(-time.Minute))
statsFinish(t, webhookDB,
statsDelivery(t, webhookDB, events[2].ID, second.ID),
database.DeliveryStatusDelivered, now.Add(-time.Minute))
return wh, events[2]
}
// seedHealthyWebhook adds a webhook with one entrypoint and one target,
// both active, and two events, both delivered. It returns the webhook
// and its newest event.
func seedHealthyWebhook(
t *testing.T,
h *handlers.Handlers,
db *database.Database,
dbMgr *database.WebhookDBManager,
) (*database.Webhook, database.Event) {
t.Helper()
wh := listWebhook(t, db, "healthy")
path := statsEntrypoint(t, db, wh.ID, true)
target := seedTarget(t, db, wh.ID, database.TargetTypeLog)
webhookDB, events := receiveEvents(t, h, dbMgr, wh.ID, path, 2)
for _, ev := range events {
statsFinish(t, webhookDB,
statsDelivery(t, webhookDB, ev.ID, target.ID),
database.DeliveryStatusDelivered, time.Now())
}
return wh, events[1]
}
// lastEventText is how the list shows the arrival of an event.
func lastEventText(ev database.Event) string {
return ev.CreatedAt.UTC().Format("2006-01-02 15:04:05 UTC")
}
// TestSourceList_ShowsActivityOfEachWebhook checks the figures the list
// shows for a webhook with recent failures, a healthy one, a new one
// that has received no event, and one without an event database.
func TestSourceList_ShowsActivityOfEachWebhook(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
dbMgr *database.WebhookDBManager
)
app := newTestApp(t, &h, &sess, &db, &dbMgr)
app.RequireStart()
t.Cleanup(app.RequireStop)
failing, failingNewest := seedFailingWebhook(t, h, db, dbMgr)
healthy, healthyNewest := seedHealthyWebhook(t, h, db, dbMgr)
// Creating a webhook creates its event database.
fresh := listWebhook(t, db, "fresh")
require.NoError(t, dbMgr.CreateDB(fresh.ID))
quiet := listWebhook(t, db, "quiet")
page := renderWebhookList(t, h, sess)
card, text := listCard(t, page, failing.ID)
assert.Contains(t, text, "2 entrypoints, 1 inactive "+
"4 targets, 1 inactive "+
"3 events within retention "+
"Last event "+lastEventText(failingNewest)+" "+
"2 failed deliveries in the last 24 hours")
assert.Contains(t, card,
failedHighlight+">2 failed deliveries in the last 24 hours<")
card, text = listCard(t, page, healthy.ID)
assert.Contains(t, text, "1 entrypoint "+
"1 target "+
"2 events within retention "+
"Last event "+lastEventText(healthyNewest)+" "+
"0 failed deliveries in the last 24 hours")
assert.NotContains(t, text, "inactive")
assert.NotContains(t, card, failedHighlight)
card, text = listCard(t, page, fresh.ID)
assert.Contains(t, text, "0 entrypoints "+
"0 targets "+
"0 events within retention "+
"No events yet "+
"0 failed deliveries in the last 24 hours")
assert.NotContains(t, card, failedHighlight)
card, text = listCard(t, page, quiet.ID)
assert.Contains(t, text, "0 entrypoints "+
"0 targets "+
"0 events within retention "+
"No events yet "+
"0 failed deliveries in the last 24 hours")
assert.NotContains(t, card, failedHighlight)
assert.False(t, dbMgr.DBExists(quiet.ID),
"showing the list must not create an event database")
}
// TestSourceList_CountsOnlyEventsWithinRetention checks that once
// retention has removed one of a webhook's three events, the list
// counts the two still stored.
func TestSourceList_CountsOnlyEventsWithinRetention(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
dbMgr *database.WebhookDBManager
log *logger.Logger
)
app := newTestApp(t, &h, &sess, &db, &dbMgr, &log)
app.RequireStart()
t.Cleanup(app.RequireStop)
wh := &database.Webhook{
UserID: deleteTestUserID, Name: "pruned", RetentionDays: 14,
}
require.NoError(t, db.DB().Omit(clause.Associations).Create(wh).Error)
path := statsEntrypoint(t, db, wh.ID, true)
webhookDB, events := receiveEvents(t, h, dbMgr, wh.ID, path, 3)
statsAge(t, webhookDB, events[0].ID, time.Now().Add(-15*24*time.Hour))
statsPrune(t, db, dbMgr, log, webhookDB)
require.Len(t, listEvents(t, webhookDB), 2)
_, text := listCard(t, renderWebhookList(t, h, sess), wh.ID)
assert.Contains(t, text,
"1 entrypoint 0 targets 2 events within retention")
}
// TestSourceList_LastEventSurvivesPruningEveryEvent checks that once
// retention has removed every event of a webhook, the list still shows
// when the last one arrived rather than "No events yet".
func TestSourceList_LastEventSurvivesPruningEveryEvent(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
dbMgr *database.WebhookDBManager
log *logger.Logger
)
app := newTestApp(t, &h, &sess, &db, &dbMgr, &log)
app.RequireStart()
t.Cleanup(app.RequireStop)
wh := &database.Webhook{
UserID: deleteTestUserID, Name: "emptied", RetentionDays: 1,
}
require.NoError(t, db.DB().Omit(clause.Associations).Create(wh).Error)
path := statsEntrypoint(t, db, wh.ID, true)
webhookDB, events := receiveEvents(t, h, dbMgr, wh.ID, path, 1)
statsAge(t, webhookDB, events[0].ID, time.Now().Add(-50*time.Hour))
statsPrune(t, db, dbMgr, log, webhookDB)
require.Empty(t, listEvents(t, webhookDB))
_, text := listCard(t, renderWebhookList(t, h, sess), wh.ID)
assert.Contains(t, text,
"0 events within retention "+
"Last event "+lastEventText(events[0]))
assert.NotContains(t, text, "No events yet")
}
// TestSourceList_UnreadableEventDatabase checks that a webhook whose
// event database cannot be read says so in its entry instead of
// showing zeros, and that the rest of the list is still shown.
func TestSourceList_UnreadableEventDatabase(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
dbMgr *database.WebhookDBManager
)
app := newTestApp(t, &h, &sess, &db, &dbMgr)
app.RequireStart()
t.Cleanup(app.RequireStop)
broken := listWebhook(t, db, "broken")
statsEntrypoint(t, db, broken.ID, true)
brokenDB, err := dbMgr.GetDB(broken.ID)
require.NoError(t, err)
require.NoError(t,
brokenDB.Migrator().DropTable(&database.EventTotals{}))
quiet := listWebhook(t, db, "quiet")
page := renderWebhookList(t, h, sess)
_, text := listCard(t, page, broken.ID)
assert.Contains(t, text,
"1 entrypoint 0 targets The event figures could not be read.")
assert.NotContains(t, text, "events")
assert.NotContains(t, text, "failed")
_, text = listCard(t, page, quiet.ID)
assert.Contains(t, text, "No events yet")
}
+119 -22
View File
@@ -3,10 +3,12 @@ package handlers
import ( import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt"
"net/http" "net/http"
"slices" "slices"
"strconv" "strconv"
"strings" "strings"
"time"
"github.com/go-chi/chi" "github.com/go-chi/chi"
"github.com/google/uuid" "github.com/google/uuid"
@@ -20,9 +22,20 @@ import (
type WebhookListItem struct { type WebhookListItem struct {
database.Webhook database.Webhook
EntrypointCount int64 EntrypointCount int
TargetCount int64 InactiveEntrypointCount int
EventCount int64 TargetCount int
InactiveTargetCount int
// EventCount is how many events the webhook holds, LastEventAt
// when the newest arrived (nil before the first), and
// FailedLast24Hours how many of its deliveries failed in the last
// 24 hours. When the webhook's event database could not be read,
// EventsUnreadable is set and these three are not known.
EventCount int64
LastEventAt *time.Time
FailedLast24Hours int64
EventsUnreadable bool
} }
// errMissingURL signals that a required URL was not provided. // errMissingURL signals that a required URL was not provided.
@@ -154,7 +167,12 @@ func (h *Handlers) HandleSourceList() http.HandlerFunc {
return return
} }
items := h.buildWebhookListItems(webhooks) items, err := h.buildWebhookListItems(webhooks)
if err != nil {
h.serverError(w, r, "failed to list webhooks", err)
return
}
data := map[string]any{ data := map[string]any{
"Webhooks": items, "Webhooks": items,
@@ -164,36 +182,115 @@ func (h *Handlers) HandleSourceList() http.HandlerFunc {
} }
} }
// buildWebhookListItems builds list items with counts. // buildWebhookListItems builds the list's entry for each webhook. It
// fails when the main database cannot be read. A webhook whose event
// database cannot be read is marked on its own entry, and the error is
// logged.
func (h *Handlers) buildWebhookListItems( func (h *Handlers) buildWebhookListItems(
webhooks []database.Webhook, webhooks []database.Webhook,
) []WebhookListItem { ) ([]WebhookListItem, error) {
items := make([]WebhookListItem, len(webhooks)) items := make([]WebhookListItem, len(webhooks))
since := time.Now().Add(-longWindow)
for i := range webhooks { for i := range webhooks {
items[i].Webhook = webhooks[i] item := &items[i]
item.Webhook = webhooks[i]
h.db.DB().Model(&database.Entrypoint{}).Where( var err error
"webhook_id = ?", webhooks[i].ID,
).Count(&items[i].EntrypointCount)
h.db.DB().Model(&database.Target{}).Where( item.EntrypointCount, item.InactiveEntrypointCount, err =
"webhook_id = ?", webhooks[i].ID, h.countWithInactive(&database.Entrypoint{}, item.ID)
).Count(&items[i].TargetCount) if err != nil {
return nil, err
}
if h.dbMgr.DBExists(webhooks[i].ID) { item.TargetCount, item.InactiveTargetCount, err =
webhookDB, err := h.dbMgr.GetDB( h.countWithInactive(&database.Target{}, item.ID)
webhooks[i].ID, if err != nil {
return nil, err
}
// Opening an event database that does not exist would create
// it, and it would hold nothing to count.
if !h.dbMgr.DBExists(item.ID) {
continue
}
err = h.readListEventFigures(item, since)
if err != nil {
h.log.Error(
"failed to read webhook list figures",
"webhook_id", item.ID,
"error", err,
) )
if err == nil {
webhookDB.Model( item.EventsUnreadable = true
&database.Event{},
).Count(&items[i].EventCount)
}
} }
} }
return items return items, nil
}
// countWithInactive returns how many entrypoints or targets, as model
// says, a webhook has, and how many of them are inactive.
func (h *Handlers) countWithInactive(
model any, webhookID string,
) (int, int, error) {
var active []bool
err := h.db.DB().Model(model).
Where("webhook_id = ?", webhookID).
Pluck("active", &active).Error
if err != nil {
return 0, 0, fmt.Errorf(
"reading active flags of webhook %s: %w", webhookID, err,
)
}
inactive := 0
for _, a := range active {
if !a {
inactive++
}
}
return len(active), inactive, nil
}
// readListEventFigures fills in the figures the list shows from the
// webhook's event database, with the statistics pane's own queries:
// the event count and last arrival from the event totals row, and the
// deliveries that failed since the given time from the deliveries'
// status index.
func (h *Handlers) readListEventFigures(
item *WebhookListItem, since time.Time,
) error {
webhookDB, err := h.dbMgr.GetDB(item.ID)
if err != nil {
return err
}
var totals database.EventTotals
err = webhookDB.Take(&totals).Error
if err != nil {
return fmt.Errorf("reading event totals: %w", err)
}
item.EventCount = totals.Events - totals.EventsRemoved
item.LastEventAt = totals.LastEventAt
byTarget, err := finishedByTarget(webhookDB, since)
if err != nil {
return err
}
for _, f := range byTarget {
item.FailedLast24Hours += f.Failed
}
return nil
} }
// HandleSourceCreate shows the form to create a new webhook. // HandleSourceCreate shows the form to create a new webhook.
+27 -20
View File
@@ -3,17 +3,18 @@
// 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. In production, these collectors, the inbound HTTP metrics
// collectors register there too, so both surfaces are gathered by the // recorded in internal/middleware, and the Go runtime and process
// one promhttp handler mounted on the authenticated /metrics route. // collectors all 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 +58,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 // one on a registry of their own so they can gather what their own
// deliveries other tests are making concurrently. // deliveries recorded.
type Set struct { type Set struct {
eventsReceived prometheus.Counter eventsReceived prometheus.Counter
deliveryAttempts *prometheus.CounterVec deliveryAttempts *prometheus.CounterVec
@@ -93,7 +100,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 {
+7 -8
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"
) )
@@ -151,17 +150,17 @@ 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
// the default registry, which is the one the /metrics route gathers. // with the Middleware's one recorder, which New builds on the registry
// the /metrics route serves and NewForTest on a registry of its own.
// Every call reuses that recorder, 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 {
+28 -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,
@@ -455,3 +454,29 @@ func TestMetrics_StatusAndSizeStillRecorded(t *testing.T) {
"the interceptor must still count written bytes", "the interceptor must still count written bytes",
) )
} }
// TestMetrics_WorksOnNewForTestMiddleware pins that a Middleware built
// by NewForTest has a recorder of its own: its Metrics() serves a
// request instead of panicking, and a second one does not collide
// with the first.
func TestMetrics_WorksOnNewForTestMiddleware(t *testing.T) {
t.Parallel()
log := slog.New(slog.DiscardHandler)
cfg := &config.Config{Environment: "prod"}
ok := http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
_, _ = w.Write([]byte(okBody))
})
for range 2 {
h := middleware.NewForTest(log, cfg, nil).Metrics()(ok)
req := httptest.NewRequestWithContext(
t.Context(), http.MethodGet, okRoute, nil,
)
w := httptest.NewRecorder()
h.ServeHTTP(w, req)
assert.Equal(t, http.StatusOK, w.Code)
}
}
+19 -4
View File
@@ -14,6 +14,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"
@@ -149,10 +152,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
@@ -162,6 +166,14 @@ type Middleware struct {
params *MiddlewareParams params *MiddlewareParams
session *session.Session session *session.Session
// metricsRecorder records the inbound HTTP metrics. New builds
// it on the registry /metrics serves, NewForTest on a registry
// of its own. Either way it is built once per Middleware and
// Metrics reuses it, because building it registers its
// collectors, and a second registration on the same registry
// panics.
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().
@@ -180,6 +192,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
} }
+8
View File
@@ -3,12 +3,17 @@ package middleware
import ( import (
"log/slog" "log/slog"
"github.com/prometheus/client_golang/prometheus"
prommetrics "github.com/slok/go-http-metrics/metrics/prometheus"
"sneak.berlin/go/webhooker/internal/config" "sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/session" "sneak.berlin/go/webhooker/internal/session"
) )
// NewForTest creates a Middleware with the minimum dependencies // NewForTest creates a Middleware with the minimum dependencies
// needed for testing. This bypasses the fx lifecycle. // needed for testing. This bypasses the fx lifecycle.
//
// Its metrics recorder writes to a fresh registry of its own, so
// Metrics() works on it and two of them never collide.
func NewForTest( func NewForTest(
log *slog.Logger, log *slog.Logger,
cfg *config.Config, cfg *config.Config,
@@ -20,5 +25,8 @@ func NewForTest(
Config: cfg, Config: cfg,
}, },
session: sess, session: sess,
metricsRecorder: prommetrics.NewRecorder(
prommetrics.Config{Registry: prometheus.NewRegistry()},
),
} }
} }
+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"
) )
@@ -149,12 +148,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,
),
)
}) })
} }
+46
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/server" "sneak.berlin/go/webhooker/internal/server"
"sneak.berlin/go/webhooker/internal/session" "sneak.berlin/go/webhooker/internal/session"
@@ -113,6 +114,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,
@@ -1478,3 +1481,46 @@ 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, Go runtime and process series, and the series counting
// scrapes of /metrics itself.
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",
"process_start_time_seconds",
"promhttp_metric_handler_requests_total",
} {
assert.Contains(t, scrape.Body.String(), series)
}
}
}
+10 -4
View File
@@ -27,10 +27,16 @@
</div> </div>
<span class="badge-info">Retention: {{.RetentionLabel}}</span> <span class="badge-info">Retention: {{.RetentionLabel}}</span>
</div> </div>
<div class="flex gap-6 mt-4 text-sm text-gray-500"> <div class="flex flex-wrap gap-6 mt-4 text-sm text-gray-500">
<span>{{.EntrypointCount}} entrypoint{{if ne .EntrypointCount 1}}s{{end}}</span> <span>{{.EntrypointCount}} entrypoint{{if ne .EntrypointCount 1}}s{{end}}{{if .InactiveEntrypointCount}}, {{.InactiveEntrypointCount}} inactive{{end}}</span>
<span>{{.TargetCount}} target{{if ne .TargetCount 1}}s{{end}}</span> <span>{{.TargetCount}} target{{if ne .TargetCount 1}}s{{end}}{{if .InactiveTargetCount}}, {{.InactiveTargetCount}} inactive{{end}}</span>
<span>{{.EventCount}} event{{if ne .EventCount 1}}s{{end}}</span> {{if .EventsUnreadable}}
<span class="text-red-600">The event figures could not be read.</span>
{{else}}
<span>{{.EventCount}} event{{if ne .EventCount 1}}s{{end}} within retention</span>
<span>{{with .LastEventAt}}Last event {{.UTC.Format "2006-01-02 15:04:05 UTC"}}{{else}}No events yet{{end}}</span>
<span class="{{if .FailedLast24Hours}}font-medium text-red-600{{end}}">{{.FailedLast24Hours}} failed deliver{{if eq .FailedLast24Hours 1}}y{{else}}ies{{end}} in the last 24 hours</span>
{{end}}
</div> </div>
</a> </a>
{{end}} {{end}}