Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a4337e949b | ||
|
|
eb4c4cc849 |
@@ -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
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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{
|
||||||
|
|||||||
@@ -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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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 {
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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")
|
||||||
|
}
|
||||||
@@ -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
@@ -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{
|
||||||
|
|||||||
@@ -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,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 {
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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 = ¶ms
|
s.params = ¶ms
|
||||||
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,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()},
|
||||||
|
),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
@@ -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,
|
|
||||||
),
|
|
||||||
)
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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}}
|
||||||
|
|||||||
Reference in New Issue
Block a user