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