3 Commits

Author SHA1 Message Date
e83eb2977e Bound shutdown hooks by their stop context (closes #102)
All checks were successful
check / check (push) Successful in 3m45s
fx hands OnStop a context carrying the application's stop timeout,
and the delivery engine, the retention reaper, and the archive
sweeper all discarded it and called wg.Wait() bare. A worker wedged
inside a delivery target that never returns, or a sweep blocked on a
locked SQLite database, hung the process forever instead of letting
it exit when the timeout expired.

All three now wait through internal/lifecycle.WaitForShutdown, which
selects the drained WaitGroup against the stop context and, on
timeout, logs at error naming the component and returns an error
rather than reporting a clean stop.

Engine.stop also gains the cancel != nil guard its two mirrored
components already had.
2026-08-12 09:45:25 +00:00
d19e33671c Gate forwarded-header trust behind trusted-proxy config (closes #88)
All checks were successful
check / check (push) Successful in 6s
All three rate limiters (receiver, login, password change) now key on the connection's own address unless the direct peer is inside the new TRUSTED_PROXIES CIDR list, in which case X-Forwarded-For is walked right to left for the first non-proxy hop. Default is the empty list, which trusts nothing. A set-but-unparseable value aborts startup.
2026-08-12 11:36:10 +02:00
aab448b076 Clarify web UI terminology, copy, and the entrypoint URL (closes #57)
All checks were successful
check / check (push) Successful in 4s
Unifies user-visible copy on "Webhook" (routes and URLs unchanged), drops
the placeholder Profile settings section, and adds a copy-to-clipboard
affordance for the entrypoint URL as progressive enhancement — the button
stays hidden unless both the target element and the Clipboard API resolve,
so no dead control appears without JavaScript and the URL stays selectable.

Retention copy now matches what the code does: deletion is permanent, 0
retains forever, and a blank field means the default on create or the
current value on edit. The permanent-deletion sentence is suppressed for a
retain-forever webhook, which the reaper exempts before computing a cutoff.

Template tests gained a render-completed assertion. Without it, a page that
aborted mid-render still satisfied assertions matching the already-flushed
prefix, because renderTemplate streams to the ResponseWriter (#123).
2026-08-11 15:42:08 +02:00
25 changed files with 1526 additions and 134 deletions

View File

@@ -95,6 +95,54 @@ TTY detection, and security headers are always applied.
| `SENTRY_DSN` | Sentry error reporting DSN | `""` |
| `SESSION_IDLE_TIMEOUT` | Idle session timeout (Go duration) | `24h` |
| `RECEIVER_RATE_LIMIT` | Receiver requests/minute per IP per entrypoint | `120` |
| `TRUSTED_PROXIES` | CIDRs whose forwarded headers are trusted | `""` (none) |
#### Trusted proxies
`TRUSTED_PROXIES` is a comma-separated list of CIDR blocks (a bare
address such as `192.168.1.7` is accepted and treated as a single
host), for example `192.168.1.7, 2001:db8::5`. It decides whose
`X-Forwarded-For` header the rate limiters believe, so it should name
the addresses of your reverse proxies and nothing else.
`X-Forwarded-For` is honoured **only** when the connecting peer is
inside one of these blocks; for every other peer the client identity is
the connection's own address and the header is ignored. The default is
the empty list, which trusts nobody — anything else would let any
client pick its own rate limit bucket, minting a fresh one per request
or draining someone else's. Set it to the address of your reverse
proxy, and to nothing wider. A set but unparseable value aborts
startup.
`X-Real-IP` and `True-Client-IP` are **never** read, from any peer.
Reverse proxies append to `X-Forwarded-For` but forward other client
headers verbatim, so a single-valued header is client-controlled even
behind a trusted proxy.
Within a trusted request, `X-Forwarded-For` is read right to left,
because the rightmost entry is the one the nearest proxy appended and
everything left of it may have been written by the client. The first
hop that is not itself a trusted proxy is taken as the client. A hop
that is not a bare IP address — `ip:port`, a bracketed IPv6 literal,
the token `unknown` — ends the walk and the peer address is used
instead, since past such an entry the chain is not the shape assumed
here. The peer address is likewise used when the header is absent or
every hop in it is a trusted proxy.
Two operator requirements follow:
- Your proxy must **append** the peer address to `X-Forwarded-For`
(nginx `$proxy_add_x_forwarded_for`, HAProxy `option forwardfor`,
Caddy and AWS ALB by default), and must append a bare address with
no port.
- List proxy hosts **only**. Any address inside `TRUSTED_PROXIES`
chooses its own rate-limit key: its `X-Forwarded-For` is walked, so
it can name a different address on every request to get a fresh
bucket each time, or name another client's address to drain that
client's bucket. Never list a block that also covers clients — a
broad `10.0.0.0/8` on a network where clients live in the same range
makes all three limits, including the unauthenticated webhook
receiver, silently bypassable by every client in the block.
Sessions are bounded by two independent clocks, and end at whichever
one runs out first:
@@ -124,8 +172,9 @@ fatal configuration error: webhooker logs the offending variable and
its value and refuses to start, rather than silently running with a
substituted default. `PORT=eighty`, `DEBUG=ture`, and
`RETENTION_SWEEP_INTERVAL=1 hour` all abort startup. `PORT` must
additionally be a number in the range 165535, and
`RECEIVER_RATE_LIMIT` must be at least 1.
additionally be a number in the range 165535,
`RECEIVER_RATE_LIMIT` must be at least 1, and every entry in
`TRUSTED_PROXIES` must be a CIDR block or a bare IP address.
Boolean variables (`DEBUG`, `MAINTENANCE_MODE`) accept exactly the
spellings Go's `strconv.ParseBool` accepts — `1`, `t`, `T`, `TRUE`,
@@ -802,6 +851,16 @@ legitimate webhook senders). Requests over the limit receive HTTP 429
with a `Retry-After` header. A set-but-invalid `RECEIVER_RATE_LIMIT`
value aborts startup rather than silently falling back to the default.
Every limiter here — receiver, login, and password change — identifies
the client the same way, through one shared key function: the
connection's own address, unless the peer is listed in
`TRUSTED_PROXIES`, in which case the forwarded client address is used
instead. See [Trusted proxies](#trusted-proxies). Deployed without that
variable set, a client behind a reverse proxy shares one bucket with
every other client behind the same proxy, which is the safe direction
to be wrong in: set `TRUSTED_PROXIES` to the proxy's address to get
per-client limits back.
Finer-grained per-webhook rate limits (configured in the web UI and
enforced in the webhook handler) can layer on top of this env-level
abuse limit later; they are tracked as future work.

View File

@@ -26,6 +26,10 @@ capability in the README rationale).
# Completed Steps
- 2026-08-11 Web UI cleanup: nav terminology unified on Webhooks, the
Profile settings placeholder removed, a progressive-enhancement copy
button for the entrypoint URL, and retention form copy that states the
actual policy (deletion by the reaper, 0 retains forever) (#57)
- 2026-08-09 Inactivity-based session timeout: sliding idle expiry
(`SESSION_IDLE_TIMEOUT`, default `24h`) refreshed on authenticated
requests, with the 7-day absolute cap kept as an independent

View File

@@ -5,8 +5,10 @@ import (
"errors"
"fmt"
"log/slog"
"net/netip"
"os"
"strconv"
"strings"
"time"
"go.uber.org/fx"
@@ -45,6 +47,11 @@ const (
// maxPort is the highest valid TCP port number. The lower
// bound (at least 1) is enforced by envPositiveInt.
maxPort = 65535
// mappedV4Offset is the number of leading bits an IPv4-mapped
// IPv6 prefix spends on the ::ffff:0:0/96 wrapper, so a /104
// covers the same addresses as an IPv4 /8.
mappedV4Offset = 96
)
// ErrInvalidEnvironment is returned when WEBHOOKER_ENVIRONMENT
@@ -59,6 +66,11 @@ var ErrNonPositiveValue = errors.New("value must be positive")
// TCP port number is set above the valid port range.
var ErrInvalidPort = errors.New("invalid port")
// ErrInvalidCIDR is returned when an environment variable holding a
// list of CIDR blocks contains an entry that is neither a CIDR block
// nor a bare IP address.
var ErrInvalidCIDR = errors.New("invalid CIDR")
//nolint:revive // ConfigParams is a standard fx naming convention.
type ConfigParams struct {
fx.In
@@ -90,6 +102,17 @@ type Config struct {
// client IP may send to a single webhook receiver entrypoint.
ReceiverRateLimit int
// TrustedProxies is the set of networks whose members are
// allowed to speak for the client with X-Forwarded-For, the
// only forwarded header read. It is empty unless
// TRUSTED_PROXIES is set, and empty means no peer is
// trusted: forwarded headers are then ignored entirely and
// clients are identified by the connection's own address.
// Members can choose their own rate-limit key, so this must
// name proxy hosts only, never a block that also covers
// clients.
TrustedProxies []netip.Prefix
params *ConfigParams
log *slog.Logger
}
@@ -212,6 +235,71 @@ func envDuration(
return d, nil
}
// parseCIDR parses one trusted-proxy list entry, which may be a
// CIDR block ("10.0.0.0/8") or a bare address ("10.0.0.1", treated
// as a single-host block).
//
// Both forms are unmapped, because peer addresses are unmapped
// before they are matched against the list: an IPv4-mapped prefix
// left in that form would silently never match.
func parseCIDR(entry string) (netip.Prefix, error) {
if strings.Contains(entry, "/") {
prefix, err := netip.ParsePrefix(entry)
if err != nil {
return netip.Prefix{}, err //nolint:wrapcheck // wrapped by caller
}
if addr := prefix.Addr(); addr.Is4In6() &&
prefix.Bits() >= mappedV4Offset {
prefix = netip.PrefixFrom(
addr.Unmap(), prefix.Bits()-mappedV4Offset,
)
}
return prefix.Masked(), nil
}
addr, err := netip.ParseAddr(entry)
if err != nil {
return netip.Prefix{}, err //nolint:wrapcheck // wrapped by caller
}
return netip.PrefixFrom(addr.Unmap(), addr.Unmap().BitLen()), nil
}
// envPrefixList returns the value of the named environment variable
// parsed as a comma-separated list of CIDR blocks (bare addresses
// allowed). An unset, empty, or blank value yields an empty list. A
// set value containing an unparseable entry is a hard error naming
// the key and the bad entry, so startup fails loudly rather than
// silently running with a list the operator did not intend.
func envPrefixList(key string) ([]netip.Prefix, error) {
v := strings.TrimSpace(os.Getenv(key))
if v == "" {
return nil, nil
}
var prefixes []netip.Prefix
for entry := range strings.SplitSeq(v, ",") {
entry = strings.TrimSpace(entry)
if entry == "" {
continue
}
prefix, err := parseCIDR(entry)
if err != nil {
return nil, fmt.Errorf(
"%w: %s: %q: %w", ErrInvalidCIDR, key, entry, err,
)
}
prefixes = append(prefixes, prefix)
}
return prefixes, nil
}
// resolveEnvironment reads WEBHOOKER_ENVIRONMENT, defaulting to
// dev, and rejects unrecognised values.
func resolveEnvironment() (string, error) {
@@ -282,6 +370,11 @@ func loadFromEnv() (*Config, error) {
return nil, err
}
trustedProxies, err := envPrefixList("TRUSTED_PROXIES")
if err != nil {
return nil, err
}
return &Config{
DataDir: envString("DATA_DIR"),
Debug: debug,
@@ -294,6 +387,7 @@ func loadFromEnv() (*Config, error) {
RetentionSweepInterval: retentionSweepInterval,
SessionIdleTimeout: sessionIdleTimeout,
ReceiverRateLimit: receiverRateLimit,
TrustedProxies: trustedProxies,
}, nil
}
@@ -335,6 +429,7 @@ func New(lc fx.Lifecycle, params ConfigParams) (*Config, error) {
"dataDir", s.DataDir,
"retentionSweepInterval", s.RetentionSweepInterval.String(),
"receiverRateLimit", s.ReceiverRateLimit,
"trustedProxies", len(s.TrustedProxies),
"hasSentryDSN", s.SentryDSN != "",
"hasMetricsAuth",
s.MetricsUsername != "" && s.MetricsPassword != "",

View File

@@ -20,6 +20,10 @@ const (
caseUnsetUsesDefault = "unset uses default"
caseValidValueParsed = "valid value is parsed"
caseUnparseableFails = "unparseable value fails startup"
// cidrPrivateV4 is the sample trusted-proxy block the
// TRUSTED_PROXIES cases are built from.
cidrPrivateV4 = "10.0.0.0/8"
)
func TestEnvironmentConfig(t *testing.T) {
@@ -179,9 +183,10 @@ func TestRetentionSweepInterval(t *testing.T) {
}
}
// expectStartupError asserts that fx refuses to build the app,
// which is what a set-but-invalid environment value must cause.
func expectStartupError(t *testing.T) {
// startupError builds the app config.New belongs to and returns
// the error fx reports, which is non-nil whenever an environment
// value is set but invalid.
func startupError(t *testing.T) error {
t.Helper()
var cfg *config.Config
@@ -196,7 +201,33 @@ func expectStartupError(t *testing.T) {
fx.Populate(&cfg),
)
assert.Error(t, app.Err())
return app.Err()
}
// expectStartupError asserts that fx refuses to build the app,
// which is what a set-but-invalid environment value must cause.
func expectStartupError(t *testing.T) {
t.Helper()
assert.Error(t, startupError(t))
}
// expectStartupErrorFor asserts that startup fails, that the error
// names the offending variable so an operator can find it, and,
// when sentinel is non-nil, that it wraps that sentinel.
func expectStartupErrorFor(
t *testing.T,
key string,
sentinel error,
) {
t.Helper()
err := startupError(t)
require.ErrorContains(t, err, key)
if sentinel != nil {
require.ErrorIs(t, err, sentinel)
}
}
func testRetentionSweepIntervalSuccess(
@@ -351,7 +382,11 @@ func TestReceiverRateLimit(t *testing.T) {
set bool
value string
expectError bool
expected int
// sentinel, when set, must be wrapped by the startup
// error; every error case must additionally name the
// variable in its message.
sentinel error
expected int
}{
{
name: caseUnsetUsesDefault,
@@ -375,12 +410,14 @@ func TestReceiverRateLimit(t *testing.T) {
set: true,
value: "0",
expectError: true,
sentinel: config.ErrNonPositiveValue,
},
{
name: "negative fails startup",
set: true,
value: "-5",
expectError: true,
sentinel: config.ErrNonPositiveValue,
},
}
@@ -399,7 +436,9 @@ func TestReceiverRateLimit(t *testing.T) {
}
if tt.expectError {
expectStartupError(t)
expectStartupErrorFor(
t, "RECEIVER_RATE_LIMIT", tt.sentinel,
)
} else {
testReceiverRateLimitSuccess(t, tt.expected)
}
@@ -432,3 +471,116 @@ func testReceiverRateLimitSuccess(
assert.Equal(t, expected, cfg.ReceiverRateLimit)
}
func TestTrustedProxies(t *testing.T) {
tests := []struct {
name string
set bool
value string
expectError bool
expected []string
}{
{
// The default must be "trust nobody": an empty list
// means forwarded headers are ignored, never that
// every peer may speak for the client.
name: caseUnsetUsesDefault,
set: false,
expected: []string{},
},
{
name: "blank value trusts nothing",
set: true,
value: " ",
expected: []string{},
},
{
name: caseValidValueParsed,
set: true,
value: cidrPrivateV4 + ", 192.168.1.7 ,2001:db8::/32",
expected: []string{
cidrPrivateV4, "192.168.1.7/32", "2001:db8::/32",
},
},
{
name: "host bits are masked off",
set: true,
value: "10.1.2.3/8",
expected: []string{cidrPrivateV4},
},
{
// Peer addresses are unmapped before they are
// matched, so an IPv4-mapped prefix kept in that
// form could never match anything.
name: "IPv4-mapped prefix is unmapped",
set: true,
value: "::ffff:10.0.0.0/104",
expected: []string{cidrPrivateV4},
},
{
name: caseUnparseableFails,
set: true,
value: cidrPrivateV4 + ",not-an-address",
expectError: true,
},
{
name: "out-of-range prefix length fails startup",
set: true,
value: "10.0.0.0/33",
expectError: true,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests.
t.Setenv("WEBHOOKER_ENVIRONMENT", "dev")
if tt.set {
t.Setenv("TRUSTED_PROXIES", tt.value)
} else {
require.NoError(t, os.Unsetenv("TRUSTED_PROXIES"))
}
if tt.expectError {
expectStartupErrorFor(
t, "TRUSTED_PROXIES", config.ErrInvalidCIDR,
)
} else {
testTrustedProxiesSuccess(t, tt.expected)
}
})
}
}
func testTrustedProxiesSuccess(
t *testing.T,
expected []string,
) {
t.Helper()
var cfg *config.Config
app := fxtest.New(
t,
fx.Provide(
globals.New,
logger.New,
config.New,
),
fx.Populate(&cfg),
)
require.NoError(t, app.Err())
app.RequireStart()
defer app.RequireStop()
got := make([]string, 0, len(cfg.TrustedProxies))
for _, prefix := range cfg.TrustedProxies {
got = append(got, prefix.String())
}
assert.Equal(t, expected, got)
}

View File

@@ -46,8 +46,19 @@ func (r *RetentionReaper) ExportStart() {
}
// ExportStop stops the reaper's background loop for tests.
func (r *RetentionReaper) ExportStop() {
r.stop()
func (r *RetentionReaper) ExportStop(ctx context.Context) error {
return r.stop(ctx)
}
// ExportWedgeLoop adds a goroutine to the reaper's WaitGroup that
// never observes cancellation and returns only when release is
// closed. It stands in for a sweep stuck on a locked database.
func (r *RetentionReaper) ExportWedgeLoop(
release <-chan struct{},
) {
r.wg.Go(func() {
<-release
})
}
// ExportSetInterval overrides the sweep interval for tests.

View File

@@ -10,6 +10,7 @@ import (
"go.uber.org/fx"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/lifecycle"
"sneak.berlin/go/webhooker/internal/logger"
)
@@ -62,8 +63,9 @@ func NewRetentionReaper(
}
// registerHooks wires the reaper's start and stop into the fx
// lifecycle. The start hook's context is deliberately ignored: see
// start for why the sweep loop must not inherit it.
// lifecycle. The start hook's context is deliberately ignored (see
// start for why the sweep loop must not inherit it); the stop hook's
// context is honoured (see stop).
func (r *RetentionReaper) registerHooks(lc fx.Lifecycle) {
lc.Append(fx.Hook{
//nolint:contextcheck // Not inheriting the hook context is
@@ -73,10 +75,8 @@ func (r *RetentionReaper) registerHooks(lc fx.Lifecycle) {
return nil
},
OnStop: func(_ context.Context) error {
r.stop()
return nil
OnStop: func(ctx context.Context) error {
return r.stop(ctx)
},
})
}
@@ -105,15 +105,27 @@ func (r *RetentionReaper) start() {
)
}
func (r *RetentionReaper) stop() {
// stop cancels the sweep loop's context and waits for it to
// exit, bounded by the stop hook's context: a sweep wedged on a
// locked database must not hang the process past fx's stop
// timeout.
func (r *RetentionReaper) stop(ctx context.Context) error {
r.log.Info("retention reaper stopping")
if r.cancel != nil {
r.cancel()
}
r.wg.Wait()
err := lifecycle.WaitForShutdown(
ctx, r.log, "retention reaper", &r.wg,
)
if err != nil {
return err
}
r.log.Info("retention reaper stopped")
return nil
}
func (r *RetentionReaper) run(ctx context.Context) {

View File

@@ -26,6 +26,13 @@ const (
// reaperTestRetentionDays is the retention policy the lifecycle
// tests give their webhook.
reaperTestRetentionDays = 30
// reaperWedgeStopTimeout is the stop timeout the wedged-shutdown
// test hands OnStop, standing in for fx's StopTimeout. The test
// asserts only that the hook returns at all, and allows it
// reaperStopTimeout — forty times this budget — to do so, so no
// assertion races the wall clock.
reaperWedgeStopTimeout = 250 * time.Millisecond
)
// recordingLifecycle is a minimal fx.Lifecycle that records the
@@ -207,3 +214,59 @@ func TestRetentionReaper_StopHookStopsLoop(t *testing.T) {
"a stopped reaper must not sweep anything",
)
}
// TestRetentionReaper_StopHookHonoursStopTimeout is the
// regression test for a shutdown that could never complete. fx
// hands OnStop a context carrying the application's stop timeout;
// an OnStop that discards it and calls wg.Wait() bare hangs the
// process forever on a sweep blocked on a locked SQLite database
// — precisely when a bounded shutdown matters most.
//
// The wedged goroutine here never observes cancellation, so the
// hook can only return by honouring its context, and it must say
// so rather than reporting a clean stop.
func TestRetentionReaper_StopHookHonoursStopTimeout(
t *testing.T,
) {
t.Parallel()
env := setupRetentionTest(t)
env.reaper.ExportSetInterval(reaperTestInterval)
lc := startReaperViaHook(t, env.reaper)
release := make(chan struct{})
t.Cleanup(func() { close(release) })
env.reaper.ExportWedgeLoop(release)
stopCtx, cancel := context.WithTimeout(
context.Background(), reaperWedgeStopTimeout,
)
defer cancel()
var stopErr error
stopped := make(chan struct{})
go func() {
defer close(stopped)
stopErr = lc.hooks[0].OnStop(stopCtx)
}()
select {
case <-stopped:
case <-time.After(reaperStopTimeout):
t.Fatal(
"OnStop did not return: it discarded the stop " +
"context and is waiting on a wedged goroutine " +
"that will never observe cancellation",
)
}
require.ErrorIs(t, stopErr, context.DeadlineExceeded)
require.ErrorContains(t, stopErr, "retention reaper")
}

View File

@@ -10,6 +10,7 @@ import (
"go.uber.org/fx"
"sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/lifecycle"
"sneak.berlin/go/webhooker/internal/logger"
)
@@ -67,10 +68,9 @@ func NewArchiveSweeper(
}
// registerHooks wires the sweeper's start and stop into the fx
// lifecycle. Both hook contexts are deliberately ignored: see
// start for why the background loop must not inherit the start
// hook's context, and stop for why shutdown blocks on the loop
// rather than on the stop hook's deadline.
// lifecycle. The start hook's context is deliberately ignored
// (see start for why the background loop must not inherit it);
// the stop hook's context is honoured (see stop).
func (s *ArchiveSweeper) registerHooks(lc fx.Lifecycle) {
lc.Append(fx.Hook{
//nolint:contextcheck // Not passing the hook context is
@@ -80,10 +80,8 @@ func (s *ArchiveSweeper) registerHooks(lc fx.Lifecycle) {
return nil
},
OnStop: func(_ context.Context) error {
s.stop()
return nil
OnStop: func(ctx context.Context) error {
return s.stop(ctx)
},
})
}
@@ -113,15 +111,27 @@ func (s *ArchiveSweeper) start() {
)
}
func (s *ArchiveSweeper) stop() {
// stop cancels the sweep loop's context and waits for it to
// exit, bounded by the stop hook's context: a prune wedged on a
// locked archive must not hang the process past fx's stop
// timeout.
func (s *ArchiveSweeper) stop(ctx context.Context) error {
s.log.Info("archive sweeper stopping")
if s.cancel != nil {
s.cancel()
}
s.wg.Wait()
err := lifecycle.WaitForShutdown(
ctx, s.log, "archive sweeper", &s.wg,
)
if err != nil {
return err
}
s.log.Info("archive sweeper stopped")
return nil
}
func (s *ArchiveSweeper) run(ctx context.Context) {

View File

@@ -14,7 +14,6 @@ import (
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.uber.org/fx"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"gorm.io/gorm/clause"
@@ -226,17 +225,6 @@ func countArchivedRows(path string) (int64, error) {
return count, nil
}
// captureLifecycle is a minimal fx.Lifecycle that records the
// hooks a component registers, so a test can invoke the real
// OnStart/OnStop functions with a context of its choosing.
type captureLifecycle struct {
hooks []fx.Hook
}
func (l *captureLifecycle) Append(h fx.Hook) {
l.hooks = append(l.hooks, h)
}
// TestArchiveSweeper_LoopOutlivesStartHookContext is the
// regression test for a sweeper that never swept. fx calls
// OnStart with a context carrying the application's start
@@ -270,7 +258,7 @@ func TestArchiveSweeper_LoopOutlivesStartHookContext(
// Drive the genuine fx hooks the application registers,
// rather than a test-only entry point.
lc := &captureLifecycle{}
lc := &recordingLifecycle{}
env.sweeper.ExportRegisterHooks(lc)
require.Len(t, lc.hooks, 1)
@@ -924,7 +912,36 @@ func TestArchiveSweeper_StopsCleanly(t *testing.T) {
env.sweeper.ExportSetInterval(time.Millisecond)
env.sweeper.ExportStart()
// stop blocks on the loop's WaitGroup, so returning at all
// proves the loop observed the cancellation and exited.
env.sweeper.ExportStop()
// stop blocks on the loop's WaitGroup, so returning without
// error proves the loop observed the cancellation and exited
// well inside the stop context.
require.NoError(
t, env.sweeper.ExportStop(context.Background()),
)
}
// TestArchiveSweeper_StopHookHonoursStopTimeout is the sweeper's
// half of the same shutdown defect the engine and the retention
// reaper carried: an OnStop that discards its context and waits
// on the WaitGroup bare hangs the process forever on a prune
// wedged inside a locked archive.
func TestArchiveSweeper_StopHookHonoursStopTimeout(
t *testing.T,
) {
t.Parallel()
env := setupSweeperTest(t)
lc := &recordingLifecycle{}
env.sweeper.ExportRegisterHooks(lc)
require.Len(t, lc.hooks, 1)
require.NoError(t, lc.hooks[0].OnStart(context.Background()))
release := make(chan struct{})
t.Cleanup(func() { close(release) })
env.sweeper.ExportWedgeLoop(release)
requireStopHookExpires(t, lc.hooks[0], "archive sweeper")
}

View File

@@ -13,6 +13,7 @@ import (
"go.uber.org/fx"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/lifecycle"
"sneak.berlin/go/webhooker/internal/logger"
)
@@ -234,8 +235,9 @@ func (e *Engine) ScheduleRetry(
}
// registerHooks wires the engine's start and stop into the fx
// lifecycle. The start hook's context is deliberately ignored:
// see start for why the worker pool must not inherit it.
// lifecycle. The start hook's context is deliberately ignored
// (see start for why the worker pool must not inherit it); the
// stop hook's context is honoured (see stop).
func (e *Engine) registerHooks(lc fx.Lifecycle) {
lc.Append(fx.Hook{
//nolint:contextcheck // Not inheriting the hook context
@@ -245,10 +247,8 @@ func (e *Engine) registerHooks(lc fx.Lifecycle) {
return nil
},
OnStop: func(_ context.Context) error {
e.stop()
return nil
OnStop: func(ctx context.Context) error {
return e.stop(ctx)
},
})
}
@@ -289,11 +289,26 @@ func (e *Engine) start() {
)
}
func (e *Engine) stop() {
// stop cancels the worker pool's context and waits for the pool
// to drain, bounded by the stop hook's context: a wedged worker
// must not hang the process past fx's stop timeout.
func (e *Engine) stop(ctx context.Context) error {
e.log.Info("delivery engine stopping")
e.cancel()
e.wg.Wait()
if e.cancel != nil {
e.cancel()
}
err := lifecycle.WaitForShutdown(
ctx, e.log, "delivery engine", &e.wg,
)
if err != nil {
return err
}
e.log.Info("delivery engine stopped")
return nil
}
func (e *Engine) worker(ctx context.Context) {

View File

@@ -501,7 +501,7 @@ func TestWorkerLifecycle_StartStop(t *testing.T) {
iWaitForDelivered(t, s.WebhookDB, d.ID)
s.Engine.ExportStop()
require.NoError(t, s.Engine.ExportStop(context.Background()))
}
// iWaitForDelivered polls until the delivery reaches the
@@ -567,7 +567,7 @@ func TestWorkerLifecycle_ProcessesRetryChannel(
iWaitForDelivered(t, s.WebhookDB, d.ID)
s.Engine.ExportStop()
require.NoError(t, s.Engine.ExportStop(context.Background()))
}
// --- processDelivery: unknown target type ---

View File

@@ -27,6 +27,13 @@ const (
// and a ready deliveryCh are chosen between at random and a
// doomed pool still delivers.
hookSettleDelay = 250 * time.Millisecond
// wedgeStopTimeout is the stop timeout a wedged-shutdown test
// hands OnStop, standing in for fx's StopTimeout. The test
// asserts only that the hook returns at all, and allows it
// hookStopTimeout — forty times this budget — to do so, so no
// assertion here races the wall clock.
wedgeStopTimeout = 250 * time.Millisecond
)
// recordingLifecycle is a minimal fx.Lifecycle that records the
@@ -40,6 +47,44 @@ func (l *recordingLifecycle) Append(h fx.Hook) {
l.hooks = append(l.hooks, h)
}
// requireStopHookExpires drives hook.OnStop with a stop context
// that expires while a wedged goroutine is still running, and
// requires the hook to return the deadline error naming
// component instead of blocking on the WaitGroup forever.
func requireStopHookExpires(
t *testing.T, hook fx.Hook, component string,
) {
t.Helper()
stopCtx, cancel := context.WithTimeout(
context.Background(), wedgeStopTimeout,
)
defer cancel()
var stopErr error
stopped := make(chan struct{})
go func() {
defer close(stopped)
stopErr = hook.OnStop(stopCtx)
}()
select {
case <-stopped:
case <-time.After(hookStopTimeout):
t.Fatal(
"OnStop did not return: it discarded the stop " +
"context and is waiting on a wedged goroutine " +
"that will never observe cancellation",
)
}
require.ErrorIs(t, stopErr, context.DeadlineExceeded)
require.ErrorContains(t, stopErr, component)
}
// startEngineViaHook drives the genuine fx hooks the application
// registers for the engine, handing OnStart a context that is
// already done, and returns only once a pool that inherited that
@@ -197,3 +242,30 @@ func TestEngine_StopHookStopsWorkers(t *testing.T) {
"a stopped engine must not deliver anything",
)
}
// TestEngine_StopHookHonoursStopTimeout is the regression test
// for a shutdown that could never complete. fx hands OnStop a
// context carrying the application's stop timeout; an OnStop
// that discards it and calls wg.Wait() bare hangs the process
// forever on a single worker stuck inside a delivery target that
// never returns — precisely when a bounded shutdown matters
// most.
//
// The wedged goroutine here never observes cancellation, so the
// hook can only return by honouring its context, and it must say
// so rather than reporting a clean stop.
func TestEngine_StopHookHonoursStopTimeout(t *testing.T) {
t.Parallel()
s := newISetup(t)
lc := startEngineViaHook(t, s.Engine)
release := make(chan struct{})
t.Cleanup(func() { close(release) })
s.Engine.ExportWedgeWorker(release)
requireStopHookExpires(t, lc.hooks[0], "delivery engine")
}

View File

@@ -216,8 +216,19 @@ func (e *Engine) ExportRegisterHooks(lc fx.Lifecycle) {
}
// ExportStop exposes stop for testing.
func (e *Engine) ExportStop() {
e.stop()
func (e *Engine) ExportStop(ctx context.Context) error {
return e.stop(ctx)
}
// ExportWedgeWorker adds a goroutine to the engine's WaitGroup
// that never observes cancellation and returns only when release
// is closed. It stands in for a worker stuck inside a delivery
// target that never returns, which is the only way stop can be
// made to outlast its context.
func (e *Engine) ExportWedgeWorker(release <-chan struct{}) {
e.wg.Go(func() {
<-release
})
}
// ExportDeliveryCh returns the delivery channel.
@@ -518,8 +529,19 @@ func (s *ArchiveSweeper) ExportRegisterHooks(lc fx.Lifecycle) {
}
// ExportStop stops the sweeper's background loop for tests.
func (s *ArchiveSweeper) ExportStop() {
s.stop()
func (s *ArchiveSweeper) ExportStop(ctx context.Context) error {
return s.stop(ctx)
}
// ExportWedgeLoop adds a goroutine to the sweeper's WaitGroup
// that never observes cancellation and returns only when release
// is closed. It stands in for a prune stuck on a locked archive.
func (s *ArchiveSweeper) ExportWedgeLoop(
release <-chan struct{},
) {
s.wg.Go(func() {
<-release
})
}
// ExportSetInterval overrides the sweep interval for tests.

View File

@@ -0,0 +1,299 @@
package handlers_test
import (
"context"
"net/http"
"net/http/httptest"
"strconv"
"testing"
"time"
"github.com/stretchr/testify/assert"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
"sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/session"
)
// Template data keys the page templates read. The handlers package has
// its own unexported constants for these; this is the external test
// package, so it needs its own.
const (
dataKeyWebhook = "Webhook"
dataKeyError = "Error"
)
// testWebhookID is the identifier given to the webhook under test on
// pages that render one.
const testWebhookID = "wh-1"
// renderPage renders a page template through the real template set as
// an authenticated user and returns the resulting HTML.
func renderPage(
t *testing.T,
h *handlers.Handlers,
sess *session.Session,
page string,
data map[string]any,
) string {
t.Helper()
cookies := authenticatedCookies(t, sess, "test-user-id", "testuser")
req := httptest.NewRequestWithContext(
context.Background(), http.MethodGet, "/", nil,
)
for _, c := range cookies {
req.AddCookie(c)
}
w := httptest.NewRecorder()
h.RenderTemplateForTest(w, req, page, data)
return w.Body.String()
}
// TestNavbarUsesWebhookTerminology pins the user-visible navigation
// label to "Webhooks". The /sources route is deliberately unchanged, so
// the assertion targets the link text rather than the href.
func TestNavbarUsesWebhookTerminology(t *testing.T) {
t.Parallel()
var h *handlers.Handlers
var sess *session.Session
app := newTestApp(t, &h, &sess)
app.RequireStart()
t.Cleanup(app.RequireStop)
// One item, so the list body renders too: it calls
// WebhookListItem.RetentionLabel, promoted from the embedded
// Webhook and therefore a pointer method. An empty list would
// skip that call and hide a template error behind the
// navigation assertions below.
item := handlers.WebhookListItem{}
item.Name = "wh"
item.ID = testWebhookID
item.RetentionDays = 14
body := renderPage(t, h, sess, "sources_list.html", map[string]any{
"Webhooks": []handlers.WebhookListItem{item},
})
assert.Contains(t, body, "Retention: 14 days")
assert.Contains(t, body, `class="btn-text">Webhooks</a>`)
assert.Contains(
t, body, `class="btn-text w-full text-left">Webhooks</a>`,
)
assert.Contains(
t, body,
`<h1 class="text-2xl font-medium text-gray-900">Webhooks</h1>`,
)
assert.NotContains(
t, body, ">Sources<",
"no user-visible element may still be labelled Sources",
)
assert.Contains(
t, body, `href="/sources"`,
"the /sources route itself must not change",
)
}
// TestEditPageUsesWebhookTerminology pins the edit page's heading and
// its back link. The link's href still points at /source/{id}, which is
// intentional: only user-visible copy changes.
func TestEditPageUsesWebhookTerminology(t *testing.T) {
t.Parallel()
var h *handlers.Handlers
var sess *session.Session
app := newTestApp(t, &h, &sess)
app.RequireStart()
t.Cleanup(app.RequireStop)
// The webhook goes in as a pointer because source_edit.html calls
// Webhook.RetentionLabel, a pointer method: a map element is not
// addressable, so a value here renders an error instead of the
// page.
webhook := &database.Webhook{Name: "wh", RetentionDays: 14}
webhook.ID = testWebhookID
body := renderPage(t, h, sess, "source_edit.html", map[string]any{
dataKeyWebhook: webhook,
dataKeyError: "",
})
assert.Contains(t, body, "Edit Webhook")
assert.NotContains(t, body, ">Sources<")
assert.Contains(t, body, `href="/source/wh-1"`)
}
// TestCreateFormRetentionCopyMatchesBehaviour pins the create form's
// retention copy to what the code does: the reaper permanently deletes
// events past the cutoff, an empty field falls back to
// DefaultRetentionDays, and 0 is rewritten to the retain-forever
// sentinel by Webhook.BeforeSave.
func TestCreateFormRetentionCopyMatchesBehaviour(t *testing.T) {
t.Parallel()
var h *handlers.Handlers
var sess *session.Session
app := newTestApp(t, &h, &sess)
app.RequireStart()
t.Cleanup(app.RequireStop)
body := renderPage(t, h, sess, "sources_new.html", map[string]any{
"Name": "",
"Description": "",
"DefaultRetentionDays": database.DefaultRetentionDays,
dataKeyError: "",
})
assert.Contains(
t, body,
"permanently deletes events older than this",
"the form must say retention is enforced by deletion",
)
assert.Contains(t, body, "Enter 0 to retain events forever")
assert.Contains(
t, body,
"leave blank to use the default of "+
strconv.Itoa(database.DefaultRetentionDays)+" days",
"blank means the default, not forever",
)
}
// TestEditFormRetentionCopyMatchesBehaviour pins the edit form's
// retention copy, including that it states the stored policy via
// RetentionLabel and that an empty field leaves that policy unchanged
// rather than meaning forever.
func TestEditFormRetentionCopyMatchesBehaviour(t *testing.T) {
t.Parallel()
var h *handlers.Handlers
var sess *session.Session
app := newTestApp(t, &h, &sess)
app.RequireStart()
t.Cleanup(app.RequireStop)
finite := &database.Webhook{Name: "wh", RetentionDays: 14}
finite.ID = testWebhookID
body := renderPage(t, h, sess, "source_edit.html", map[string]any{
dataKeyWebhook: finite,
dataKeyError: "",
})
assert.Contains(t, body, "Currently 14 days.")
assert.Contains(
t, body,
"permanently deletes events older than this",
)
assert.Contains(t, body, "Enter 0 to retain events forever")
assert.Contains(
t, body,
"leave blank to keep the current setting",
"blank means unchanged, not forever",
)
forever := &database.Webhook{
Name: "wh",
RetentionDays: database.RetentionForeverDays,
}
forever.ID = "wh-2"
foreverBody := renderPage(
t, h, sess, "source_edit.html", map[string]any{
dataKeyWebhook: forever,
dataKeyError: "",
},
)
assert.Contains(
t, foreverBody, "Currently forever.",
"a retain-forever webhook must not read as a day count",
)
assert.Contains(
t, foreverBody,
"No events are deleted while retention is set to forever",
)
assert.NotContains(
t, foreverBody,
"permanently deletes events older than this",
"the reaper skips retain-forever webhooks, so the form "+
"must not claim it deletes their events",
)
}
// TestEntrypointCopyButtonIsProgressiveEnhancement proves the copy
// affordance degrades: the button ships with the hidden attribute, so a
// browser that never runs app.js shows no dead control, and the URL is
// rendered as ordinary selectable text either way.
func TestEntrypointCopyButtonIsProgressiveEnhancement(t *testing.T) {
t.Parallel()
var h *handlers.Handlers
var sess *session.Session
app := newTestApp(t, &h, &sess)
app.RequireStart()
t.Cleanup(app.RequireStop)
entrypoint := database.Entrypoint{Path: "abc123"}
entrypoint.ID = "ep-1"
// The webhook goes in as a pointer because source_detail.html
// calls Webhook.RetentionLabel, a pointer method: a map element
// is not addressable, so a value here aborts execution partway
// down the page, after the copy button has already been flushed
// to the response.
webhook := &database.Webhook{Name: "wh", RetentionDays: 14}
webhook.ID = testWebhookID
webhook.CreatedAt = time.Date(
2026, time.January, 2, 3, 4, 5, 0, time.UTC,
)
body := renderPage(t, h, sess, "source_detail.html", map[string]any{
dataKeyWebhook: webhook,
"Entrypoints": []database.Entrypoint{entrypoint},
// The handler passes delivery.NewTargetViews(targets), never
// raw targets, so the test data has to have that same shape.
"Targets": delivery.NewTargetViews(nil),
"Events": []database.Event{},
"BaseURL": "https://hooks.example.com",
})
assert.Contains(
t, body,
`<code id="entrypoint-url-ep-1"`,
)
assert.Contains(t, body, "https://hooks.example.com/webhook/abc123")
assert.Contains(
t, body,
`hidden data-copy-target="entrypoint-url-ep-1"`,
"the button must start hidden and be revealed by script",
)
// renderTemplate streams to the ResponseWriter, so an abort
// midway still leaves everything above it in the body. This pins
// content from the last line of the template, which is below the
// assertions above: without it, a page that renders the copy
// button and then 500s passes.
assert.Contains(
t, body, "Retention: 14 days",
"the page must render to completion, not abort partway",
)
}

View File

@@ -0,0 +1,57 @@
// Package lifecycle holds helpers shared by the components that
// register fx start and stop hooks.
package lifecycle
import (
"context"
"fmt"
"log/slog"
"sync"
)
// WaitForShutdown waits for wg to drain, bounded by ctx.
//
// fx hands OnStop a context carrying the application's stop
// timeout. A bare wg.Wait() discards that deadline, so a single
// goroutine that never observes cancellation — a delivery target
// that never returns, a SQLite operation blocked on a lock —
// hangs the process forever instead of letting it exit when the
// timeout expires, which is exactly when a clean shutdown matters
// most.
//
// On timeout it logs at error naming component and returns an
// error: the goroutines are still running, and reporting success
// would hide an unclean shutdown from the operator. The waiting
// goroutine outlives this call and exits when (if) wg drains; it
// holds nothing but the channel it closes.
func WaitForShutdown(
ctx context.Context,
log *slog.Logger,
component string,
wg *sync.WaitGroup,
) error {
done := make(chan struct{})
go func() {
defer close(done)
wg.Wait()
}()
select {
case <-done:
return nil
case <-ctx.Done():
log.Error(
"shutdown timed out, goroutines still running",
"component", component,
"error", ctx.Err(),
)
return fmt.Errorf(
"%s: shutdown timed out, "+
"goroutines still running: %w",
component, ctx.Err(),
)
}
}

View File

@@ -0,0 +1,62 @@
package lifecycle_test
import (
"context"
"log/slog"
"sync"
"testing"
"time"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/lifecycle"
)
// waitTimeout is the stop budget the timeout case gives a
// goroutine that never returns. The test's own patience is the
// go test deadline, so the only thing this value affects is how
// long the case takes.
const waitTimeout = 100 * time.Millisecond
func discardLogger() *slog.Logger {
return slog.New(slog.DiscardHandler)
}
func TestWaitForShutdown_DrainedGroup(t *testing.T) {
t.Parallel()
var wg sync.WaitGroup
wg.Go(func() {})
require.NoError(
t,
lifecycle.WaitForShutdown(
context.Background(), discardLogger(),
"test component", &wg,
),
)
}
func TestWaitForShutdown_ContextExpires(t *testing.T) {
t.Parallel()
release := make(chan struct{})
t.Cleanup(func() { close(release) })
var wg sync.WaitGroup
wg.Go(func() { <-release })
ctx, cancel := context.WithTimeout(
context.Background(), waitTimeout,
)
defer cancel()
err := lifecycle.WaitForShutdown(
ctx, discardLogger(), "test component", &wg,
)
require.ErrorIs(t, err, context.DeadlineExceeded)
require.ErrorContains(t, err, "test component")
}

View File

@@ -2,6 +2,9 @@ package middleware
import (
"net/http"
"net/netip"
"slices"
"strings"
"time"
"github.com/go-chi/httprate"
@@ -31,13 +34,120 @@ const (
receiverRateInterval = 1 * time.Minute
)
// normalizeAddr strips the IPv4-in-IPv6 wrapper and any zone from
// addr so that comparisons and bucket keys are canonical.
func normalizeAddr(addr netip.Addr) netip.Addr {
return addr.Unmap().WithZone("")
}
// isTrustedProxy reports whether addr belongs to a network the
// operator listed in TRUSTED_PROXIES. The list is empty by default,
// so by default nothing is trusted.
func (m *Middleware) isTrustedProxy(addr netip.Addr) bool {
for _, prefix := range m.params.Config.TrustedProxies {
if prefix.Contains(addr) {
return true
}
}
return false
}
// forwardedClientAddr returns the client address named by this
// request's X-Forwarded-For chain. It is consulted only for requests
// whose direct peer is a trusted proxy.
//
// X-Forwarded-For is the only header read. X-Real-IP and
// True-Client-IP are deliberately ignored: the reverse proxies in
// common use append to X-Forwarded-For and pass any other header the
// client sent through untouched, so believing a single-valued header
// would let a client behind the trusted proxy name its own bucket —
// the very bypass this gating exists to close.
//
// The chain is walked right to left, because the rightmost entry is
// the one the nearest proxy appended and everything to its left may
// have been written by the client. The first hop that is not itself
// a trusted proxy is the client. A hop that cannot be read as a bare
// address ends the walk: past it the chain is not the shape assumed
// here, so the caller falls back to the peer address.
func (m *Middleware) forwardedClientAddr(
r *http.Request,
) (netip.Addr, bool) {
hops := strings.Split(
strings.Join(r.Header.Values("X-Forwarded-For"), ","), ",",
)
for _, hop := range slices.Backward(hops) {
hop = strings.TrimSpace(hop)
if hop == "" {
continue
}
addr, err := netip.ParseAddr(hop)
if err != nil {
return netip.Addr{}, false
}
if addr = normalizeAddr(addr); !m.isTrustedProxy(addr) {
return addr, true
}
}
return netip.Addr{}, false
}
// rateLimitKey is the client identity every rate limiter in this
// package buckets on. Forwarded headers are honoured only when the
// direct peer (RemoteAddr) is inside the configured trusted-proxy
// set; otherwise the peer address itself is the key. Without that
// gate any client could mint a fresh bucket per request, or starve
// another client's bucket, by picking an X-Forwarded-For value —
// which makes every limit here decorative against a deliberate
// attacker.
func (m *Middleware) rateLimitKey(r *http.Request) (string, error) {
return m.clientKey(r), nil
}
// clientKey computes the bucket key described on rateLimitKey.
func (m *Middleware) clientKey(r *http.Request) string {
peer, err := netip.ParseAddr(ipFromHostPort(r.RemoteAddr))
if err != nil {
// Not an address we can reason about; key on the raw
// value rather than collapsing such peers into one
// shared bucket.
return r.RemoteAddr
}
peer = normalizeAddr(peer)
if !m.isTrustedProxy(peer) {
return peer.String()
}
if addr, ok := m.forwardedClientAddr(r); ok {
return addr.String()
}
return peer.String()
}
// tooManyRequests returns the 429 handler shared by every limiter:
// it logs the rejection with logMessage and answers with
// responseMessage. httprate adds the Retry-After header (RFC 6585).
func (m *Middleware) tooManyRequests(
logMessage, responseMessage string,
) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
m.log.Warn(logMessage, "path", r.URL.Path)
http.Error(w, responseMessage, http.StatusTooManyRequests)
}
}
// LoginRateLimit returns middleware that enforces per-IP rate
// limiting on login attempts using go-chi/httprate. Only POST
// requests are rate-limited; GET requests (rendering the login
// form) pass through unaffected. When the rate limit is exceeded,
// a 429 Too Many Requests response is returned. IP extraction
// honours X-Forwarded-For, X-Real-IP, and True-Client-IP headers
// for reverse-proxy setups.
// a 429 Too Many Requests response is returned. Clients are
// identified by rateLimitKey.
func (m *Middleware) LoginRateLimit() func(http.Handler) http.Handler {
return m.postRateLimit(
loginRateLimit,
@@ -66,9 +176,7 @@ func (m *Middleware) PasswordChangeRateLimit() func(http.Handler) http.Handler {
// limit on POST requests only; all other methods pass through
// unaffected. Requests over the limit receive a 429 with the
// given response message, and each rejection is logged with the
// given log message. IP extraction honours X-Forwarded-For,
// X-Real-IP, and True-Client-IP headers for reverse-proxy
// setups.
// given log message. Clients are identified by rateLimitKey.
func (m *Middleware) postRateLimit(
limit int,
interval time.Duration,
@@ -77,19 +185,10 @@ func (m *Middleware) postRateLimit(
limiter := httprate.Limit(
limit,
interval,
httprate.WithKeyFuncs(httprate.KeyByRealIP),
httprate.WithLimitHandler(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
m.log.Warn(logMessage,
"path", r.URL.Path,
)
http.Error(
w,
responseMessage,
http.StatusTooManyRequests,
)
},
)),
httprate.WithKeyFuncs(m.rateLimitKey),
httprate.WithLimitHandler(
m.tooManyRequests(logMessage, responseMessage),
),
)
return func(next http.Handler) http.Handler {
@@ -116,31 +215,19 @@ func (m *Middleware) postRateLimit(
// path (the path contains the entrypoint UUID, so each sender
// is limited per entrypoint without affecting other senders or
// other entrypoints). The limit is Config.ReceiverRateLimit
// requests per minute. Requests over the limit receive a 429;
// httprate adds the Retry-After header (RFC 6585). IP
// extraction honours X-Forwarded-For, X-Real-IP, and
// True-Client-IP headers for reverse-proxy setups.
// requests per minute. Requests over the limit receive a 429.
// Clients are identified by rateLimitKey.
func (m *Middleware) ReceiverRateLimit() func(http.Handler) http.Handler {
return httprate.Limit(
m.params.Config.ReceiverRateLimit,
receiverRateInterval,
httprate.WithKeyFuncs(
httprate.KeyByRealIP,
m.rateLimitKey,
httprate.KeyByEndpoint,
),
httprate.WithLimitHandler(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
m.log.Warn(
"webhook receiver rate limit exceeded",
"path", r.URL.Path,
)
http.Error(
w,
"Too many requests. "+
"Please slow down.",
http.StatusTooManyRequests,
)
},
httprate.WithLimitHandler(m.tooManyRequests(
"webhook receiver rate limit exceeded",
"Too many requests. Please slow down.",
)),
)
}

View File

@@ -2,9 +2,11 @@ package middleware_test
import (
"context"
"fmt"
"log/slog"
"net/http"
"net/http/httptest"
"net/netip"
"os"
"testing"
@@ -182,11 +184,22 @@ func TestLoginRateLimit_IndependentPerIP(t *testing.T) {
)
}
// receiverLimitedHandler builds a ReceiverRateLimit-wrapped
// handler with the given per-minute limit.
func receiverLimitedHandler(
t *testing.T, limit int,
) http.Handler {
// okHandler is the terminal handler the limiter middleware wraps
// in these tests: it answers 200 to anything that reaches it.
func okHandler() http.Handler {
return http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
},
)
}
// rateLimitMiddleware builds a Middleware around cfg, whose
// TrustedProxies field is what the rate limit key function gates
// forwarded-header trust on.
func rateLimitMiddleware(
t *testing.T, cfg *config.Config,
) *middleware.Middleware {
t.Helper()
log := slog.New(slog.NewTextHandler(
@@ -194,17 +207,53 @@ func receiverLimitedHandler(
&slog.HandlerOptions{Level: slog.LevelDebug},
))
m := middleware.NewForTest(
log,
&config.Config{ReceiverRateLimit: limit},
nil,
return middleware.NewForTest(log, cfg, nil)
}
// trustedProxies parses CIDR strings for a test Config.
func trustedProxies(cidrs ...string) []netip.Prefix {
prefixes := make([]netip.Prefix, 0, len(cidrs))
for _, cidr := range cidrs {
prefixes = append(prefixes, netip.MustParsePrefix(cidr))
}
return prefixes
}
// postWithHeaders sends one POST to the handler from peer with the
// given headers set and returns the recorder.
func postWithHeaders(
handler http.Handler,
peer, path string,
headers map[string]string,
) *httptest.ResponseRecorder {
req := httptest.NewRequestWithContext(
context.Background(), http.MethodPost, path, nil,
)
req.RemoteAddr = peer
for name, value := range headers {
req.Header.Set(name, value)
}
w := httptest.NewRecorder()
handler.ServeHTTP(w, req)
return w
}
// receiverLimitedHandler builds a ReceiverRateLimit-wrapped
// handler with the given per-minute limit and no trusted proxies.
func receiverLimitedHandler(
t *testing.T, limit int,
) http.Handler {
t.Helper()
m := rateLimitMiddleware(
t, &config.Config{ReceiverRateLimit: limit},
)
return m.ReceiverRateLimit()(http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
},
))
return m.ReceiverRateLimit()(okHandler())
}
// receiverPost sends one POST to the handler from the given IP
@@ -311,3 +360,252 @@ func TestReceiverRateLimit_CountsEveryMethod(t *testing.T) {
"a GET over the limit must be rate-limited",
)
}
const (
loginPath = "/pages/login"
headerXFF = "X-Forwarded-For"
headerReal = "X-Real-IP"
headerTrue = "True-Client-IP"
)
// assertSharedBucket drives the login limiter from peer with the
// trusted-proxy set proxies, sending one more request than the limit
// allows and varying the headers on each with headers(i). Every
// request must land in the same bucket, so the last one is rejected:
// if any of the varying header values reached the key, the run would
// have minted fresh buckets and nothing would be rejected.
func assertSharedBucket(
t *testing.T,
proxies []netip.Prefix,
peer string,
headers func(i int) map[string]string,
msg string,
) {
t.Helper()
m := rateLimitMiddleware(
t, &config.Config{TrustedProxies: proxies},
)
handler := m.LoginRateLimit()(okHandler())
for i := range middleware.LoginRateLimitConst {
w := postWithHeaders(handler, peer, loginPath, headers(i))
assert.Equal(
t, http.StatusOK, w.Code,
"request %d should pass", i,
)
}
w := postWithHeaders(
handler, peer, loginPath,
headers(middleware.LoginRateLimitConst),
)
assert.Equal(t, http.StatusTooManyRequests, w.Code, msg)
}
// TestRateLimitKey_SpoofedForwardedFromUntrustedPeer is the test
// this gating exists for: with no trusted proxies configured (the
// default), a client that rotates a forwarded header on every
// request must stay in one bucket. If forwarded headers were
// trusted unconditionally, each spoofed value would mint a fresh
// bucket and the limit would stop no one.
func TestRateLimitKey_SpoofedForwardedFromUntrustedPeer(
t *testing.T,
) {
t.Parallel()
for _, header := range []string{
headerXFF, headerReal, headerTrue,
} {
t.Run(header, func(t *testing.T) {
t.Parallel()
assertSharedBucket(
t, nil, "203.0.113.9:44444",
func(i int) map[string]string {
return map[string]string{
header: fmt.Sprintf(
"198.51.100.%d", i+1,
),
}
},
"a spoofed "+header+" from an untrusted peer "+
"must not mint a fresh bucket",
)
})
}
}
// TestRateLimitKey_SingleValuedHeadersIgnoredFromTrustedPeer is the
// regression test for the bypass hiding inside the trusted case.
// Real reverse proxies (nginx, HAProxy, Caddy, ALB) set only
// X-Forwarded-For and pass every other client header through
// verbatim, so a client behind the configured proxy can send its own
// X-Real-IP or True-Client-IP. Reading either would hand that client
// a fresh bucket per request from inside exactly the deployment
// TRUSTED_PROXIES exists to serve, so neither header is read at all.
func TestRateLimitKey_SingleValuedHeadersIgnoredFromTrustedPeer(
t *testing.T,
) {
t.Parallel()
for _, header := range []string{headerReal, headerTrue} {
t.Run(header, func(t *testing.T) {
t.Parallel()
assertSharedBucket(
t, trustedProxies("10.0.0.0/8"),
"10.0.0.1:44444",
func(i int) map[string]string {
return map[string]string{
header: fmt.Sprintf(
"198.51.100.%d", i+1,
),
}
},
header+" from a trusted peer must not mint a "+
"fresh bucket: only X-Forwarded-For is read",
)
})
}
}
// TestRateLimitKey_MalformedRightmostHopFallsBackToPeer covers the
// other end of the chain walk. The rightmost X-Forwarded-For entry
// is the one the trusted proxy appended; if it cannot be read as an
// address the chain is not the shape the walk assumes, and every
// entry to its left may have come from the client. The walk must
// stop and fall back to the peer rather than select one of them.
func TestRateLimitKey_MalformedRightmostHopFallsBackToPeer(
t *testing.T,
) {
t.Parallel()
// Forms seen in the wild: host:port (Azure Application
// Gateway, IIS ARR), a bracketed IPv6 literal, and the
// RFC 7239 placeholder token.
for _, tail := range []string{
"198.51.100.7:1234", "[2001:db8::1]", "unknown",
} {
t.Run(tail, func(t *testing.T) {
t.Parallel()
assertSharedBucket(
t, trustedProxies("10.0.0.0/8"),
"10.0.0.1:44444",
func(i int) map[string]string {
return map[string]string{
headerXFF: fmt.Sprintf(
"9.9.9.%d, %s", i+1, tail,
),
}
},
"an unparseable rightmost hop must fall back "+
"to the peer address, not select a "+
"client-controlled entry",
)
})
}
}
// TestRateLimitKey_ForwardedHonouredFromTrustedPeer checks the
// other half: when the direct peer is a configured trusted proxy,
// the forwarded client address is what buckets are keyed on, so
// one sender behind the proxy cannot exhaust another's limit.
func TestRateLimitKey_ForwardedHonouredFromTrustedPeer(
t *testing.T,
) {
t.Parallel()
m := rateLimitMiddleware(t, &config.Config{
TrustedProxies: trustedProxies("10.0.0.0/8"),
})
handler := m.LoginRateLimit()(okHandler())
const peer = "10.0.0.1:44444"
first := map[string]string{headerXFF: "198.51.100.7"}
for range middleware.LoginRateLimitConst {
postWithHeaders(handler, peer, loginPath, first)
}
w := postWithHeaders(handler, peer, loginPath, first)
assert.Equal(
t, http.StatusTooManyRequests, w.Code,
"the forwarded client's own bucket must fill up",
)
w = postWithHeaders(
handler, peer, loginPath,
map[string]string{headerXFF: "198.51.100.8"},
)
assert.Equal(
t, http.StatusOK, w.Code,
"a forwarded header from a trusted peer must be honoured",
)
}
// TestRateLimitKey_ChainWalkSkipsClientPrepended covers the
// residual spoofing route behind a trusted proxy: the client
// controls the leftmost X-Forwarded-For entries, so the key is the
// rightmost hop that is not itself trusted. Rotating the prepended
// entry must not create new buckets.
func TestRateLimitKey_ChainWalkSkipsClientPrepended(t *testing.T) {
t.Parallel()
assertSharedBucket(
t, trustedProxies("10.0.0.0/8"), "10.0.0.1:44444",
func(i int) map[string]string {
return map[string]string{
headerXFF: fmt.Sprintf(
"9.9.9.%d, 198.51.100.7, 10.0.0.2", i+1,
),
}
},
"a client-prepended X-Forwarded-For entry must not "+
"mint a fresh bucket",
)
}
// TestReceiverRateLimit_IgnoresForwardedFromUntrustedPeer proves
// the receiver limiter uses the same gated key function as the
// POST limiters.
func TestReceiverRateLimit_IgnoresForwardedFromUntrustedPeer(
t *testing.T,
) {
t.Parallel()
const (
limit = 3
peer = "203.0.113.10:44444"
path = "/webhook/uuid-d"
)
handler := receiverLimitedHandler(t, limit)
for i := range limit {
w := postWithHeaders(
handler, peer, path,
map[string]string{
headerXFF: fmt.Sprintf(
"198.51.100.%d", i+1,
),
},
)
assert.Equal(
t, http.StatusOK, w.Code,
"request %d should pass", i,
)
}
w := postWithHeaders(
handler, peer, path,
map[string]string{headerXFF: "198.51.100.200"},
)
assert.Equal(
t, http.StatusTooManyRequests, w.Code,
"a spoofed X-Forwarded-For from an untrusted peer must "+
"not mint a fresh receiver bucket",
)
}

View File

@@ -1,2 +1,60 @@
// Webhooker client-side JavaScript
console.log("Webhooker loaded");
// Copy-to-clipboard, as progressive enhancement.
//
// Markup renders each copy button with the `hidden` attribute and a
// `data-copy-target` pointing at the id of the element holding the
// text. This script reveals a button only once it has both a resolvable
// target and a usable Clipboard API, so a browser without either shows
// no button at all and the text stays selectable.
(function () {
"use strict";
const revertDelayMs = 2000;
function flash(button, message) {
const original = button.getAttribute("data-copy-label");
button.textContent = message;
window.setTimeout(function () {
button.textContent = original;
}, revertDelayMs);
}
function wire(button) {
const target = document.getElementById(
button.getAttribute("data-copy-target")
);
if (!target) {
return;
}
button.setAttribute("data-copy-label", button.textContent);
button.addEventListener("click", function () {
navigator.clipboard.writeText(target.textContent.trim()).then(
function () {
flash(button, "Copied");
},
function () {
flash(button, "Copy failed");
}
);
});
button.removeAttribute("hidden");
}
function init() {
if (!navigator.clipboard || !navigator.clipboard.writeText) {
return;
}
const buttons = document.querySelectorAll("[data-copy-target]");
buttons.forEach(wire);
}
if (document.readyState === "loading") {
document.addEventListener("DOMContentLoaded", init);
} else {
init();
}
})();

View File

@@ -16,7 +16,7 @@
<!-- Desktop navigation -->
<div class="hidden md:flex items-center gap-4">
{{if .User}}
<a href="/sources" class="btn-text">Sources</a>
<a href="/sources" class="btn-text">Webhooks</a>
<a href="/user/{{.User.Username}}" class="btn-text">
<svg class="w-5 h-5 mr-1" fill="currentColor" viewBox="0 0 16 16">
<path d="M11 6a3 3 0 1 1-6 0 3 3 0 0 1 6 0z"/>
@@ -38,7 +38,7 @@
<div x-show="open" x-cloak x-transition class="md:hidden mt-4 pt-4 border-t border-gray-200">
<div class="flex flex-col gap-2">
{{if .User}}
<a href="/sources" class="btn-text w-full text-left">Sources</a>
<a href="/sources" class="btn-text w-full text-left">Webhooks</a>
<a href="/user/{{.User.Username}}" class="btn-text w-full text-left">Profile</a>
<form method="POST" action="/pages/logout">
<input type="hidden" name="csrf_token" value="{{.CSRFToken}}">

View File

@@ -34,24 +34,18 @@
<hr class="border-gray-200 mb-6">
<div class="grid grid-cols-1 md:grid-cols-2 gap-8">
<div>
<h3 class="text-lg font-medium text-gray-900 mb-3">Account Information</h3>
<dl class="space-y-3">
<div class="flex">
<dt class="w-32 text-sm font-medium text-gray-500">Username</dt>
<dd class="text-sm text-gray-900">{{.User.Username}}</dd>
</div>
<div class="flex">
<dt class="w-32 text-sm font-medium text-gray-500">Account Type</dt>
<dd class="text-sm text-gray-900">Standard User</dd>
</div>
</dl>
</div>
<div>
<h3 class="text-lg font-medium text-gray-900 mb-3">Settings</h3>
<p class="text-sm text-gray-500">Profile settings and preferences will be available here.</p>
</div>
<div>
<h3 class="text-lg font-medium text-gray-900 mb-3">Account Information</h3>
<dl class="space-y-3">
<div class="flex">
<dt class="w-32 text-sm font-medium text-gray-500">Username</dt>
<dd class="text-sm text-gray-900">{{.User.Username}}</dd>
</div>
<div class="flex">
<dt class="w-32 text-sm font-medium text-gray-500">Account Type</dt>
<dd class="text-sm text-gray-900">Standard User</dd>
</div>
</dl>
</div>
</div>

View File

@@ -69,7 +69,12 @@
</form>
</div>
</div>
<code class="text-xs text-gray-500 break-all block mt-1">{{$.BaseURL}}/webhook/{{.Path}}</code>
<div class="flex items-start gap-2 mt-1">
<code id="entrypoint-url-{{.ID}}" class="text-xs text-gray-500 break-all block flex-1">{{$.BaseURL}}/webhook/{{.Path}}</code>
<!-- Hidden until app.js reveals it; without the
script the URL above stays selectable. -->
<button type="button" hidden data-copy-target="entrypoint-url-{{.ID}}" class="text-xs text-gray-500 hover:text-primary-600">Copy</button>
</div>
</div>
{{else}}
<div class="p-4 text-sm text-gray-500">No entrypoints configured.</div>

View File

@@ -29,7 +29,7 @@
<div class="form-group">
<label for="retention_days" class="label">Retention (days)</label>
<input type="number" id="retention_days" name="retention_days" value="{{.Webhook.RetentionDays}}" min="0" class="input">
<p class="text-xs text-gray-500 mt-1">Currently {{.Webhook.RetentionLabel}}. Enter 0 to retain events forever.</p>
<p class="text-xs text-gray-500 mt-1">Currently {{.Webhook.RetentionLabel}}.{{if .Webhook.RetainsForever}} No events are deleted while retention is set to forever.{{else}} A periodic cleanup permanently deletes events older than this, along with their delivery records.{{end}} Enter 0 to retain events forever; leave blank to keep the current setting.</p>
</div>
<div class="flex gap-3">

View File

@@ -1,6 +1,6 @@
{{template "base" .}}
{{define "title"}}Sources - Webhooker{{end}}
{{define "title"}}Webhooks - Webhooker{{end}}
{{define "content"}}
<div class="max-w-6xl mx-auto px-6 py-8">

View File

@@ -29,7 +29,7 @@
<div class="form-group">
<label for="retention_days" class="label">Retention (days)</label>
<input type="number" id="retention_days" name="retention_days" value="{{.DefaultRetentionDays}}" min="0" class="input">
<p class="text-xs text-gray-500 mt-1">How long to keep event data. Enter 0 to retain events forever.</p>
<p class="text-xs text-gray-500 mt-1">A periodic cleanup permanently deletes events older than this, along with their delivery records. Enter 0 to retain events forever; leave blank to use the default of {{.DefaultRetentionDays}} days.</p>
</div>
<div class="flex gap-3">