Autor SHA1 Zpráva Datum
clawbot fa67d883e9 Close archive writers when the delivery engine stops (closes #280)
check / check (push) Successful in 3m36s
The engine cached archive writers and never closed them at shutdown,
so after a clean stop an archive's rows could sit in its -wal while
the .db held no table. The engine's stop hook now evicts every cached
writer once its workers have returned, the same way deleting a webhook
does, so a clean stop leaves each archive as one file and a late write
is refused. If the workers do not return within the stop budget, the
writers are left open as a kill would leave them.

The README no longer says archives keep their sidecars across a clean
stop.

Model: opus-5-5
2026-09-29 05:00:30 +00:00
27 změnil soubory, kde provedl 129 přidání a 431 odebrání
+11 -18
Zobrazit soubor
@@ -157,11 +157,6 @@ private and reserved ranges — RFC 1918, loopback, CGNAT, link-local and
the rest — are refused, which stops a target from being used to make
webhooker probe the network it sits in.
Besides the private and reserved ranges, the default blocklist refuses
public cloud metadata addresses: currently only `168.63.129.16`, Azure's
WireServer, which serves an Azure VM its credentials. Because it is a
public address, listing it in `ALLOWED_EGRESS_CIDRS` reopens it.
That default is also inconvenient for the thing webhooker is mostly
for: taking a public webhook and forwarding it to something on your own
network. A container on the same Docker network, a box on `10.x`, a
@@ -200,16 +195,15 @@ Two things this setting cannot do:
the list is always an allowlist; an empty list (the default) means
every private and reserved range stays refused. Note that
`0.0.0.0/0` gets you most of the way there anyway, per above.
- **It cannot open link-local, or a cloud metadata endpoint at a
non-public address that discloses credentials or user data.** An
address is on the list below when it is not a public address and both
of these hold: the provider fixes it, so it cannot collide with
anything you run; and reaching it hands out credentials, user data or
bootstrap material. Those stay blocked no matter what you list,
including when you list them outright or list a supernet such as
`0.0.0.0/0`, `::/0`, `fd00::/8` or `100.64.0.0/10`. Treat this as best
effort rather than a guarantee — it is a hand-maintained list and the
caveat below the table applies:
- **It cannot open link-local, or a cloud metadata endpoint that
discloses credentials or user data.** An address is on the list below
when both of these hold: the provider fixes it, so it cannot collide
with anything you run; and reaching it hands out credentials, user
data or bootstrap material. Those stay blocked no matter what you
list, including when you list them outright or list a supernet such
as `0.0.0.0/0`, `::/0`, `fd00::/8` or `100.64.0.0/10`. Treat this as
best effort rather than a guarantee — it is a hand-maintained list
and the caveat below the table applies:
| Blocked unconditionally | What it is |
| ----------------------- | ---------- |
@@ -248,8 +242,7 @@ Two things this setting cannot do:
encodings, which the default blocklist does not match. A publicly
routable metadata address is not listed here, because nothing on this
list can be reopened and blocking one that way would leave you no
escape hatch at all; Azure's `168.63.129.16` is refused by the default
blocklist instead, as described above.
escape hatch at all.
This list is not exhaustive of every cloud's metadata address — if
yours is not here, do not allowlist the block that contains it.
@@ -2721,7 +2714,7 @@ abuse limit later; they are tracked as future work.
| ------ | --------------------------- | ----------- |
| `GET` | `/` | Root redirect, 303 (authenticated → `/sources`, unauthenticated → `/pages/login`) |
| `GET` | `/.well-known/healthcheck` | Health check (JSON: `status`, `now`, `uptimeSeconds`, `uptimeHuman`, `version`, `appname`, `maintenanceMode`) |
| `GET`, `HEAD` | `/s/*` | Static file serving (embedded CSS, JS). `GET` and `HEAD` only — `POST`, `PUT`, `PATCH`, `DELETE`, `OPTIONS`, `TRACE` and `CONNECT` are answered `405 Method Not Allowed` with `Allow: GET, HEAD`. Any other method (such as `PROPFIND`) is refused by chi before it reaches this route, and gets `405` without an `Allow` header. Pinned by `TestStaticServesOnlyGetAndHead` |
| any | `/s/*` | Static file serving (embedded CSS, JS). Mounted for every method, not just `GET`/`HEAD`: chi's `Mount` registers all methods and `http.FileServer` special-cases only `HEAD` (by omitting the body), so a `POST` or `DELETE` to an asset is answered `200` with the file. Pinned by `TestStaticServesEveryMethod` |
| `POST` | `/webhook/{uuid}` | Webhook receiver endpoint. `POST` only — every other method is answered `405 Method Not Allowed` with `Allow: POST`. Rate limited (see [Rate Limiting](#rate-limiting)) |
#### Authentication Endpoints
-5
Zobrazit soubor
@@ -16,7 +16,6 @@ import (
"sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/resetpw"
"sneak.berlin/go/webhooker/internal/server"
@@ -178,10 +177,6 @@ func newApp() *fx.App {
healthcheck.New,
session.New,
handlers.New,
// The registry /metrics serves, and the delivery
// collectors registered on it.
metrics.NewRegistry,
metrics.New,
middleware.New,
// The one SSRF guard both target-creation validation
// and the delivery dialer consult, so they cannot
+9 -12
Zobrazit soubor
@@ -192,10 +192,9 @@ type Config struct {
// alwaysBlockedNetworks stays blocked no matter what is listed
// here. That set is link-local plus the cloud metadata
// endpoints outside it that disclose credentials or user data
// at a provider-fixed, non-public address; it is not
// exhaustive of every cloud's metadata address. See
// alwaysBlockedNetworks for the authoritative list and the
// criterion it is built from.
// at a provider-fixed address; it is not exhaustive of every
// cloud's metadata address. See alwaysBlockedNetworks for the
// authoritative list and the criterion it is built from.
AllowedEgressCIDRs []netip.Prefix
params *ConfigParams
@@ -747,14 +746,12 @@ func (c *Config) warnEgressAllowlist(log *slog.Logger) {
log.Warn(
"ALLOWED_EGRESS_CIDRS lets delivery targets reach these "+
"otherwise-blocked networks. Anyone who can create a "+
"delivery target can now make this process issue "+
"requests into them, and read back the response. Only "+
"the addresses the README lists as blocked "+
"unconditionally stay blocked regardless of what is "+
"listed here; a public cloud metadata address such as "+
"168.63.129.16 is reachable once it, or a block "+
"covering it, is listed.",
"otherwise-blocked private/reserved networks. Anyone "+
"who can create a delivery target can now make this "+
"process issue requests into them, and read back the "+
"response. Link-local and the known cloud instance "+
"metadata endpoints outside it stay blocked "+
"regardless of what is listed here.",
"allowedEgressCIDRs",
strings.Join(PrefixStrings(c.AllowedEgressCIDRs), ","),
)
+6 -7
Zobrazit soubor
@@ -834,13 +834,12 @@ func TestEgressAllowlistWarning(t *testing.T) {
// to be able to read back which networks are open.
assert.Contains(t, logged, "10.0.0.0/8")
assert.Contains(t, logged, "127.0.0.0/8")
// What stays shut is the whole unconditional set, not
// link-local alone; a public metadata address is not in
// it, so a listed block covering it opens it.
assert.Contains(t, logged, "blocked unconditionally")
assert.Contains(t, logged, "168.63.129.16 is reachable")
// The listed blocks need not be private or reserved.
assert.NotContains(t, logged, "private/reserved")
// What stays shut. Asserted on the clause naming the
// wider set rather than on "Link-local" alone, so the
// string cannot narrow back to link-local only while
// the always-blocked set covers ULA, CGNAT and two
// public metadata addresses as well.
assert.Contains(t, logged, "metadata endpoints outside it")
})
}
}
+8 -9
Zobrazit soubor
@@ -148,7 +148,6 @@ type EngineParams struct {
DBManager *database.WebhookDBManager
Logger *logger.Logger
SSRFGuard *Guard
Metrics *metrics.Set
}
// Engine processes queued deliveries in the background
@@ -168,10 +167,10 @@ type Engine struct {
retryCh chan Task
workers int
// mtr is the delivery metric set. Production wires the one
// registered on the registry /metrics serves; a test can
// substitute a set registered on a registry it holds, so it can
// gather what its own deliveries recorded.
// mtr is the delivery metric set. Production wires the
// process-wide one; a test can substitute a set registered on
// a private registry so its assertions are not disturbed by
// deliveries other tests are making at the same time.
mtr *metrics.Set
// targets maps each target type to its implementation.
@@ -205,7 +204,7 @@ func New(
deliveryCh: make(chan Task, deliveryChannelSize),
retryCh: make(chan Task, retryChannelSize),
workers: defaultWorkers,
mtr: params.Metrics,
mtr: metrics.Default(),
}
e.initTargets(&http.Client{
@@ -369,9 +368,9 @@ func (e *Engine) start() {
// writer for long by then: the archive sweeper stops before the
// engine, and deleting a webhook only closes one. If the pool did
// not drain in time, the writers are left open, as a kill would
// leave them. Closing them would wait for any write in progress,
// and a worker still running would then open new writers that
// nothing closes, so it gains nothing over a kill.
// leave them: a worker still running may be mid-write, and
// closing its writer would wait on that write and then fail the
// next delivery the worker archives.
func (e *Engine) stop(ctx context.Context) error {
e.log.Info("delivery engine stopping")
+4 -4
Zobrazit soubor
@@ -327,10 +327,10 @@ func TestEngine_StopHookClosesArchives(t *testing.T) {
}
// TestEngine_StopHookTimeoutLeavesArchivesOpen covers a stop whose
// budget runs out while a worker is still running. The archive
// writers are left open, as a kill would leave them: closing them
// would wait for any write in progress, and that worker would then
// open new writers that nothing closes.
// budget runs out while a worker is still running. That worker may
// be in the middle of an archive write, so the archive writers are
// left open, as a kill would leave them, rather than closed
// underneath it.
func TestEngine_StopHookTimeoutLeavesArchivesOpen(t *testing.T) {
t.Parallel()
-79
Zobrazit soubor
@@ -1166,10 +1166,6 @@ func TestIsForwardableHeader(t *testing.T) {
assert.False(t,
delivery.ExportIsForwardableHeader("Content-Length"),
)
assert.False(t,
delivery.ExportIsForwardableHeader("Content-Type"),
)
}
func TestTruncate(t *testing.T) {
@@ -1251,81 +1247,6 @@ func TestDoHTTPRequest_ForwardsHeaders(t *testing.T) {
)
}
// The event's stored inbound headers carry the same Content-Type the
// receiver saved as the event's ContentType, so a delivery could send
// it twice. It must go out exactly once, with a Content-Type configured
// on the target winning, then the event's ContentType.
func TestApplyRequestHeaders_SendsOneContentType(t *testing.T) {
t.Parallel()
cases := map[string]struct {
inbound string
event string
configured string
want []string
}{
"inbound and event agree": {
inbound: testContentType,
event: testContentType,
want: []string{testContentType},
},
"inbound and event disagree": {
inbound: "text/plain",
event: testContentType,
want: []string{testContentType},
},
"event has none": {
inbound: testContentType,
want: nil,
},
"target configures its own": {
inbound: testContentType,
event: testContentType,
configured: "application/xml",
want: []string{"application/xml"},
},
}
for name, tc := range cases {
t.Run(name, func(t *testing.T) {
t.Parallel()
inbound, err := json.Marshal(map[string][]string{
headerContentType: {tc.inbound},
})
require.NoError(t, err)
cfg := &delivery.HTTPTargetConfig{}
if tc.configured != "" {
cfg.Headers = map[string]string{
headerContentType: tc.configured,
}
}
req, err := http.NewRequestWithContext(
context.Background(),
http.MethodPost,
"https://target.example.com/hook",
http.NoBody,
)
require.NoError(t, err)
delivery.ExportApplyRequestHeaders(
req,
&database.Event{
Headers: string(inbound),
ContentType: tc.event,
},
cfg,
)
assert.Equal(t,
tc.want, req.Header.Values(headerContentType),
)
})
}
}
func TestProcessDelivery_RoutesToCorrectHandler(
t *testing.T,
) {
+5 -5
Zobrazit soubor
@@ -9,7 +9,6 @@ import (
"net/url"
"time"
"github.com/prometheus/client_golang/prometheus"
"go.uber.org/fx"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
@@ -390,7 +389,7 @@ func NewTestEngine(
deliveryCh: make(chan Task, deliveryChannelSize),
retryCh: make(chan Task, retryChannelSize),
workers: workers,
mtr: metrics.New(prometheus.NewRegistry()),
mtr: metrics.Default(),
}
e.initTargets(client)
@@ -405,7 +404,7 @@ func NewTestEngineSmallRetry(
e := &Engine{
log: log,
retryCh: make(chan Task, 1),
mtr: metrics.New(prometheus.NewRegistry()),
mtr: metrics.Default(),
}
e.initTargets(nil)
@@ -428,7 +427,7 @@ func NewTestEngineWithDB(
deliveryCh: make(chan Task, deliveryChannelSize),
retryCh: make(chan Task, retryChannelSize),
workers: workers,
mtr: metrics.New(prometheus.NewRegistry()),
mtr: metrics.Default(),
}
e.initTargets(client)
@@ -436,7 +435,8 @@ func NewTestEngineWithDB(
}
// ExportSetMetrics substitutes the engine's metric set, so a test can
// assert on collectors registered on a registry it holds.
// assert on collectors registered on a private registry instead of
// the process-wide ones every other test is also moving.
func (e *Engine) ExportSetMetrics(mtr *metrics.Set) {
e.mtr = mtr
}
+3 -2
Zobrazit soubor
@@ -35,8 +35,9 @@ const (
)
// mIsolate gives the setup's engine a metric set registered on a
// registry this test holds, so its exact assertions can gather from
// it.
// private registry. The process-wide collectors are moved by every
// other delivery test running in parallel, so exact assertions are
// only possible against a registry this test owns.
func mIsolate(
t *testing.T, s iSetup,
) *prometheus.Registry {
+6 -12
Zobrazit soubor
@@ -339,11 +339,10 @@ func TestRedirectPolicy_StopsAtHopCap(t *testing.T) {
// The set the redirect policy strips is whatever the delivery path
// actually put on the wire, so a header added to the forward set is
// covered without a second edit. A header the event never carried
// is not in the set, and neither is the inbound Content-Type, because
// it is not forwarded. Two more are deliberately excluded: a
// Content-Type configured on the target describes the body, which a
// 307 carries across hosts, and the inbound User-Agent every real
// sender supplies is overwritten before the request goes out.
// is not in the set, and the delivery path's own two are deliberately
// excluded: Content-Type describes the body, which a 307 carries
// across hosts, and the inbound User-Agent every real sender supplies
// is overwritten before the request goes out.
func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
t.Parallel()
@@ -372,7 +371,6 @@ func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
&delivery.HTTPTargetConfig{
Headers: map[string]string{
probeHeaderName: probeHeaderValue,
"Content-Type": testContentType,
},
},
)
@@ -380,11 +378,7 @@ func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
assert.Equal(t,
[]string{probeHeaderName, inboundHeaderName}, names,
"both header classes are reported, and only those: "+
"Host and the inbound Content-Type are never "+
"forwarded, User-Agent is the delivery path's own",
)
assert.NotContains(t, names, "Content-Type",
"a Content-Type configured on the target must survive "+
"a cross-origin 307/308 with the body it describes",
"Host is never forwarded, Content-Type and "+
"User-Agent are the delivery path's own",
)
}
+7 -10
Zobrazit soubor
@@ -26,7 +26,7 @@ var (
"hostname resolved to no IP addresses",
)
errBlockedIP = errors.New(
"blocked private, reserved or cloud metadata address",
"blocked private/reserved IP range",
)
errBlockedMetadata = errors.New(
"blocked link-local or cloud instance metadata " +
@@ -37,10 +37,9 @@ var (
)
)
// blockedNetworks is the default blocklist: the private and
// reserved IP ranges, plus the public cloud metadata addresses,
// that are blocked to prevent SSRF attacks. An operator can
// permit specific blocks out of this set with
// blockedNetworks contains all private/reserved IP ranges
// that should be blocked to prevent SSRF attacks. An operator
// can permit specific blocks out of this set with
// ALLOWED_EGRESS_CIDRS; see Guard.
//
//nolint:gochecknoglobals // package-level network list is appropriate here
@@ -123,8 +122,6 @@ func init() {
"::1/128",
"fc00::/7",
"fe80::/10",
// Azure WireServer, a public address that serves VM credentials.
"168.63.129.16/32",
})
// Every entry is named. The set must not grow or shrink
@@ -219,8 +216,8 @@ func matchesAny(networks []*net.IPNet, ip net.IP) bool {
}
// isBlockedIP checks whether an IP address falls within
// the default blocklist, before any operator allowlist is
// considered.
// any blocked private/reserved network range, before any
// operator allowlist is considered.
func isBlockedIP(ip net.IP) bool {
return matchesAny(blockedNetworks, ip)
}
@@ -323,7 +320,7 @@ func (g *Guard) allows(ip net.IP) bool {
//
// 1. alwaysBlockedNetworks is refused before the allowlist is
// consulted, so no configured CIDR reaches link-local or a
// cloud metadata endpoint at a non-public address.
// cloud instance metadata endpoint.
// 2. The allowlist is consulted next, so a listed private
// network becomes reachable.
// 3. Everything else keeps the default blocklist's answer.
-35
Zobrazit soubor
@@ -390,41 +390,6 @@ func TestGuardAllowlist_PublicUnaffected(t *testing.T) {
}
}
// TestGuardAllowlist_AzureWireServerReopenable covers Azure's
// WireServer, a public address that serves VM credentials. The
// default guard refuses it, but because it is public it sits in
// the default blocklist rather than the unconditional set, so an
// operator who lists it can reach it.
func TestGuardAllowlist_AzureWireServerReopenable(t *testing.T) {
t.Parallel()
const wireServerIP = "168.63.129.16"
target := "http://" + wireServerIP + "/?comp=versions"
defaultGuard := delivery.NewTestGuard()
err := defaultGuard.ValidateTargetURL(context.Background(), target)
require.Error(t, err,
"WireServer must be refused with no allowlist set",
)
assert.NotContains(t, err.Error(), metadataRefusalClause,
"WireServer must be refused by the default blocklist, "+
"which an allowlist can override",
)
assertDialRefused(t, defaultGuard, target)
listed := delivery.NewTestGuard(
netip.MustParsePrefix(wireServerIP + "/32"),
)
assert.NoError(t,
listed.ValidateTargetURL(context.Background(), target),
"an operator who lists WireServer must be able to reach it",
)
}
// TestGuardCheckIP_BothPathsShareOneDecision asserts that the
// validator and the dialer are not two policies that happen to
// agree: both are defined in terms of checkIP, so the exported
+1 -2
Zobrazit soubor
@@ -11,11 +11,10 @@ import (
"sneak.berlin/go/webhooker/internal/delivery"
)
// Literals these tests repeat, named so that the header names and the
// Literals these tests repeat, named so that the header name and the
// keep-forever archive config each have one definition.
const (
headerAuthorization = "Authorization"
headerContentType = "Content-Type"
bearerValue = "Bearer abc"
archiveConfigNever = "{\"expiry\":\"never\"}"
)
+4 -13
Zobrazit soubor
@@ -541,11 +541,6 @@ func isForwardableHeader(name string) bool {
"Upgrade", "Proxy-Authorization",
"Proxy-Connection", "Content-Length":
return false
case "Content-Type":
// applyRequestHeaders sets Content-Type itself. The receiver
// already stored this inbound value as the event's
// ContentType, so forwarding it too would send it twice.
return false
default:
return true
}
@@ -558,10 +553,6 @@ func isForwardableHeader(name string) bool {
// policy strips exactly that set on a hop that leaves the origin,
// so the forward set is decided here and only here — a header added
// to it is covered off-origin without a second edit elsewhere.
//
// Content-Type goes out once: a Content-Type configured on the target
// wins, otherwise the event's ContentType, otherwise none. The inbound
// Content-Type in the event's headers is never forwarded.
func applyRequestHeaders(
req *http.Request,
event *database.Event,
@@ -582,10 +573,10 @@ func applyRequestHeaders(
req.Header.Set("User-Agent", "webhooker/1.0")
// A Content-Type configured on the target describes the body
// being sent rather than the sender. A 307/308 preserves the
// body across hosts, so stripping it would send that body
// untyped.
// Content-Type describes the body being sent rather than the
// sender, and the delivery path sets it from the event itself.
// A 307/308 preserves the body across hosts, so stripping it
// would send that body untyped.
delete(originScoped, "Content-Type")
// User-Agent is overwritten just above, so an inbound one never
+1 -4
Zobrazit soubor
@@ -12,7 +12,6 @@ import (
"net/http"
"sync/atomic"
"github.com/prometheus/client_golang/prometheus"
"go.uber.org/fx"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
@@ -62,8 +61,6 @@ type HandlersParams struct {
Notifier delivery.Notifier
Evictor delivery.WebhookEvictor
SSRFGuard *delivery.Guard
Metrics *metrics.Set
Registry *prometheus.Registry
}
// Handlers provides HTTP handler methods for all application
@@ -125,7 +122,7 @@ func New(
s.mw = params.Middleware
s.notifier = params.Notifier
s.evictor = params.Evictor
s.mtr = params.Metrics
s.mtr = metrics.Default()
s.ssrf = params.SSRFGuard
// Parse all page templates once at startup
-3
Zobrazit soubor
@@ -20,7 +20,6 @@ import (
"sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/session"
)
@@ -110,8 +109,6 @@ func newTestApp(
func(r *recordingEvictor) delivery.WebhookEvictor {
return r
},
metrics.NewRegistry,
metrics.New,
middleware.New,
delivery.NewGuard,
handlers.New,
-20
Zobrazit soubor
@@ -1,20 +0,0 @@
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
}
+20 -26
Zobrazit soubor
@@ -3,17 +3,17 @@
// deliveries are attempted, how they end, how long they take, how
// deep the queues are, and how many circuit breakers are open.
//
// It also builds the registry the authenticated /metrics route
// serves. These collectors, the inbound HTTP metrics recorded in
// internal/middleware, and the Go runtime and process collectors all
// register on that one registry, never on Prometheus's global default.
// The inbound HTTP metrics come from the go-http-metrics recorder in
// internal/middleware and land on prometheus.DefaultRegisterer. These
// collectors register there too, so both surfaces are gathered by the
// one promhttp handler mounted on the authenticated /metrics route.
package metrics
import (
"sync"
"time"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/collectors"
"github.com/prometheus/client_golang/prometheus/promauto"
"sneak.berlin/go/webhooker/internal/database"
)
@@ -57,31 +57,25 @@ var knownTargetTypes = []database.TargetType{
database.TargetTypeSlack,
}
// NewRegistry returns the registry /metrics serves, carrying the Go
// runtime and process collectors that Prometheus's global default
// registry carries, so the go_* and process_* series stay in the
// scrape.
// defaultSet is the process-wide metric set, registered on the same
// registry the HTTP middleware and the /metrics handler already use.
// It is built on first use rather than in an init so that a test
// binary that never touches metrics never registers them.
//
// A registry of its own, rather than the global default, is what lets
// two dependency graphs in one process — two tests, say — each
// register their collectors without the second registration
// panicking.
func NewRegistry() *prometheus.Registry {
reg := prometheus.NewRegistry()
reg.MustRegister(
collectors.NewGoCollector(),
collectors.NewProcessCollector(
collectors.ProcessCollectorOpts{},
),
)
//nolint:gochecknoglobals // one process-wide registration, by design
var defaultSet = sync.OnceValue(func() *Set {
return New(prometheus.DefaultRegisterer)
})
return reg
// Default returns the process-wide metric set.
func Default() *Set {
return defaultSet()
}
// Set is one registered group of webhooker's delivery collectors.
// Production builds one on the registry /metrics serves; tests build
// their own against a private registry so assertions are not
// disturbed by deliveries other tests are making concurrently.
// Production uses the single Default set; tests build their own
// against a private registry so assertions are not disturbed by
// deliveries other tests are making concurrently.
type Set struct {
eventsReceived prometheus.Counter
deliveryAttempts *prometheus.CounterVec
@@ -99,7 +93,7 @@ type Set struct {
// New registers a full set of delivery collectors on reg and returns
// it. It panics if reg already holds them, which is the intended
// behaviour for a duplicate registration.
func New(reg *prometheus.Registry) *Set {
func New(reg prometheus.Registerer) *Set {
factory := promauto.With(reg)
s := &Set{
+2 -1
Zobrazit soubor
@@ -10,7 +10,8 @@ import (
// MetricsMiddlewareForTest builds the metrics recording middleware
// against a caller-supplied recorder, so a test can gather from its
// own Prometheus registry without building a whole Middleware.
// own Prometheus registry rather than the process-wide default one
// that Middleware.Metrics uses.
func MetricsMiddlewareForTest(
rec httpmetrics.Recorder,
) func(http.Handler) http.Handler {
+7 -4
Zobrazit soubor
@@ -7,6 +7,7 @@ import (
"github.com/go-chi/chi"
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"
"github.com/slok/go-http-metrics/middleware/std"
)
@@ -151,14 +152,16 @@ func (r boundedLabelRecorder) AddInflightRequests(
var _ httpmetrics.Recorder = boundedLabelRecorder{}
// Metrics returns middleware that records Prometheus HTTP metrics on
// the registry the /metrics route serves. Every call shares the one
// recorder New built, so any number of routers can install it.
// the default registry, which is the one the /metrics route gathers.
func (s *Middleware) Metrics() func(http.Handler) http.Handler {
return metricsMiddleware(s.metricsRecorder)
return metricsMiddleware(
prommetrics.NewRecorder(prommetrics.Config{}),
)
}
// metricsMiddleware builds the recording middleware against a given
// recorder, so tests can gather from a registry of their own.
// recorder, so tests can gather from a registry of their own instead
// of the process-wide default.
func metricsMiddleware(
rec httpmetrics.Recorder,
) func(http.Handler) http.Handler {
+3 -2
Zobrazit soubor
@@ -57,8 +57,9 @@ const (
// Server.setupWebhookRoutes inside it. That ordering is the whole
// defect, so a test that flattens it would prove nothing.
//
// The recorder writes to a registry of the test's own, so each test
// observes only its own traffic.
// The recorder writes to a registry of the test's own rather than the
// process-wide default one, so each test observes only its own
// traffic.
func metricsTestRouter(
t *testing.T,
receiverLimit int,
+4 -17
Zobrazit soubor
@@ -13,9 +13,6 @@ import (
"github.com/go-chi/chi"
"github.com/go-chi/chi/middleware"
"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"
"sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/globals"
@@ -151,11 +148,10 @@ const (
type MiddlewareParams struct {
fx.In
Logger *logger.Logger
Globals *globals.Globals
Config *config.Config
Session *session.Session
Registry *prometheus.Registry
Logger *logger.Logger
Globals *globals.Globals
Config *config.Config
Session *session.Session
}
// Middleware provides HTTP middleware for logging, CORS, auth, and
@@ -165,12 +161,6 @@ type Middleware struct {
params *MiddlewareParams
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
// concurrent password hashing. It is built on first use so that
// every construction path gets one; see guard().
@@ -189,9 +179,6 @@ func New(
s.params = &params
s.log = params.Logger.Get()
s.session = params.Session
s.metricsRecorder = prommetrics.NewRecorder(
prommetrics.Config{Registry: params.Registry},
)
return s, nil
}
-7
Zobrazit soubor
@@ -133,11 +133,6 @@ func (w *recoverResponseWriter) Unwrap() http.ResponseWriter {
// what the access log records and the metrics count, and outside the
// sentryhttp handler, whose Repanic option depends on something
// further out recovering what it re-raises.
//
// Unlike http.Error on its own, it deletes any Set-Cookie the handler
// set before panicking, because a request that failed must not hand
// the client a credential; every other header is left to http.Error.
// See https://git.eeqj.de/sneak/webhooker/issues/193.
func (s *Middleware) Recoverer() func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(
@@ -169,8 +164,6 @@ func (s *Middleware) Recoverer() func(http.Handler) http.Handler {
return
}
rw.Header().Del("Set-Cookie")
http.Error(
rw,
http.StatusText(
+2 -31
Zobrazit soubor
@@ -304,44 +304,16 @@ func TestRecovererRepanicsErrAbortHandler(t *testing.T) {
)
}
// TestRecovererDropsSetCookieFromTheRecovered500 covers a handler that
// sets a cookie and a redirect target and then panics before sending
// anything. A request that failed must not hand the client a
// credential, so the 500 carries no cookie; Location is left alone.
func TestRecovererDropsSetCookieFromTheRecovered500(t *testing.T) {
t.Parallel()
probe := newRecovererProbe(
t, false,
func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Set-Cookie", "session=x")
w.Header().Set("Location", "/after")
panic(panicMarker)
},
)
resp, err := probe.get(t)
require.NoError(t, err)
require.NoError(t, resp.Body.Close())
assert.Equal(t, http.StatusInternalServerError, resp.StatusCode)
assert.Empty(t, resp.Cookies())
assert.Equal(t, "/after", resp.Header.Get("Location"))
}
// TestRecovererKeepsAnAlreadyCommittedResponse covers a handler that
// panics after sending its status. The bytes are already on the wire,
// cookie included, so a second WriteHeader would change nothing the
// client sees and would draw net/http's "superfluous
// response.WriteHeader" report.
// so a second WriteHeader would change nothing the client sees and
// would draw net/http's "superfluous response.WriteHeader" report.
func TestRecovererKeepsAnAlreadyCommittedResponse(t *testing.T) {
t.Parallel()
probe := newRecovererProbe(
t, false,
func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Set-Cookie", "session=x")
w.WriteHeader(committedStatus)
_, _ = w.Write([]byte("partial"))
@@ -359,7 +331,6 @@ func TestRecovererKeepsAnAlreadyCommittedResponse(t *testing.T) {
assert.Equal(t, committedStatus, resp.StatusCode)
assert.Equal(t, "partial", string(body))
assert.Len(t, resp.Cookies(), 1)
record := probe.panicRecord(t)
assert.Equal(t, panicMarker, record["panic"])
-3
Zobrazit soubor
@@ -24,7 +24,6 @@ import (
"sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/resetpw"
"sneak.berlin/go/webhooker/internal/session"
@@ -164,8 +163,6 @@ func newServerApp(
session.New,
func() delivery.Notifier { return &noopNotifier{} },
func() delivery.WebhookEvictor { return &noopEvictor{} },
metrics.NewRegistry,
metrics.New,
middleware.New,
delivery.NewGuard,
handlers.New,
+10 -18
Zobrazit soubor
@@ -7,6 +7,7 @@ import (
sentryhttp "github.com/getsentry/sentry-go/http"
"github.com/go-chi/chi"
"github.com/go-chi/chi/middleware"
"github.com/prometheus/client_golang/prometheus/promhttp"
"sneak.berlin/go/webhooker/static"
)
@@ -91,25 +92,11 @@ func (s *Server) setupGlobalMiddleware() {
func (s *Server) setupRoutes() {
s.router.Get("/", s.h.HandleIndex())
// Static assets answer GET and HEAD only. chi's default 405
// carries no Allow header, so this group supplies its own.
staticFiles := http.StripPrefix(
"/s", http.FileServer(http.FS(static.Static)),
s.router.Mount(
"/s",
http.StripPrefix("/s", http.FileServer(http.FS(static.Static))),
)
s.router.Route("/s", func(r chi.Router) {
r.MethodNotAllowed(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Allow", "GET, HEAD")
http.Error(
w,
"Method Not Allowed",
http.StatusMethodNotAllowed,
)
})
r.Method(http.MethodGet, "/*", staticFiles)
r.Method(http.MethodHead, "/*", staticFiles)
})
s.router.Route("/api/v1", func(_ chi.Router) {
// API routes will be added here.
})
@@ -129,7 +116,12 @@ func (s *Server) setupRoutes() {
if s.params.Config.MetricsAuthEnabled() {
s.router.Group(func(r chi.Router) {
r.Use(s.mw.MetricsAuth())
r.Get("/metrics", s.h.HandleMetrics())
r.Get(
"/metrics",
http.HandlerFunc(
promhttp.Handler().ServeHTTP,
),
)
})
}
+16 -82
Zobrazit soubor
@@ -23,7 +23,6 @@ import (
"sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/server"
"sneak.berlin/go/webhooker/internal/session"
@@ -113,8 +112,6 @@ func newTestEnvWithConfig(
session.New,
func() delivery.Notifier { return &noopNotifier{} },
func() delivery.WebhookEvictor { return &noopEvictor{} },
metrics.NewRegistry,
metrics.New,
middleware.New,
delivery.NewGuard,
handlers.New,
@@ -399,15 +396,13 @@ func (e *testEnv) storedHash(t *testing.T, username string) string {
// --- /s static group ---
// TestStaticServesOnlyGetAndHead pins the methods the static group
// answers: GET and HEAD are served the asset, and the other methods
// chi routes (POST, PUT, DELETE and the rest) are refused with 405
// and an Allow header naming those two. A method chi does not route,
// such as PROPFIND, is refused with 405 by the top-level router
// before it reaches the static group, so it gets no Allow header.
// The README documents this; the test is what keeps the two from
// drifting.
func TestStaticServesOnlyGetAndHead(t *testing.T) {
// TestStaticServesEveryMethod pins what the static mount actually
// answers. chi's Mount registers the handler for all methods and
// http.FileServer only special-cases HEAD (by suppressing the body),
// so a POST or a DELETE to an asset is served the file rather than
// refused. The README documents this; the test is what keeps the two
// from drifting.
func TestStaticServesEveryMethod(t *testing.T) {
t.Parallel()
env := newTestEnv(t)
@@ -422,7 +417,6 @@ func TestStaticServesOnlyGetAndHead(t *testing.T) {
http.MethodPost,
http.MethodPut,
http.MethodDelete,
"PROPFIND",
} {
t.Run(method, func(t *testing.T) {
t.Parallel()
@@ -434,38 +428,18 @@ func TestStaticServesOnlyGetAndHead(t *testing.T) {
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
switch method {
case http.MethodGet:
assert.Equal(t, http.StatusOK, w.Code)
assert.Equal(t, body, w.Body.Bytes(),
"the asset itself is returned")
case http.MethodHead:
assert.Equal(t, http.StatusOK, w.Code)
assert.Equal(t, http.StatusOK, w.Code,
"static mount answers every method")
if method == http.MethodHead {
assert.Empty(t, w.Body.Bytes(),
"HEAD must not carry a body")
case "PROPFIND":
assert.Equal(
t, http.StatusMethodNotAllowed, w.Code,
)
assert.Empty(t, w.Header().Get("Allow"),
"chi refuses a method it does not route "+
"before the static group runs")
assert.NotContains(
t, w.Body.String(), string(body),
"a refused method must not get the asset",
)
default:
assert.Equal(
t, http.StatusMethodNotAllowed, w.Code,
)
assert.Equal(
t, "GET, HEAD", w.Header().Get("Allow"),
)
assert.NotContains(
t, w.Body.String(), string(body),
"a refused method must not get the asset",
)
return
}
assert.Equal(t, body, w.Body.Bytes(),
"the asset itself is returned")
})
}
}
@@ -964,43 +938,3 @@ 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)
}
}
}