Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8b0854ccb7 | ||
|
|
70e708ee4f | ||
|
|
815260b405 | ||
|
|
6af979420b | ||
|
|
cfd043ed24 |
@@ -147,7 +147,7 @@ TTY detection, and security headers are always applied.
|
||||
| `RETENTION_SWEEP_INTERVAL` | How often the retention reaper and archive sweeper run (Go duration, must be positive) | `1h` |
|
||||
| `SESSION_IDLE_TIMEOUT` | Idle session timeout (Go duration) | `24h` |
|
||||
| `RECEIVER_RATE_LIMIT` | Receiver requests/minute per IP per entrypoint (10x that per IP across the route) | `120` |
|
||||
| `TRUSTED_PROXIES` | CIDRs whose forwarded headers are trusted (unset: all clients behind a proxy share one rate-limit bucket; a correct login password is never throttled either way) | `""` (none) |
|
||||
| `TRUSTED_PROXIES` | CIDRs whose forwarded headers are trusted. A set value replaces the default. Under the default, any client with a private address, whether it connects directly or through a trusted proxy, can choose its own rate-limit key by sending its own `X-Forwarded-For`; if any clients have private addresses, set it to the proxy's address alone. See [Trusted proxies](#trusted-proxies) | `10.0.0.0/8,172.16.0.0/12,192.168.0.0/16` (RFC 1918) |
|
||||
| `ALLOWED_EGRESS_CIDRS` | CIDRs that delivery targets may reach despite the SSRF blocklist. Read [Allowing egress to your own network](#allowing-egress-to-your-own-network) before setting it | `""` (none) |
|
||||
|
||||
#### Allowing egress to your own network
|
||||
@@ -379,41 +379,48 @@ unlocked.
|
||||
`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` header the rate limiters believe, so it should cover
|
||||
the addresses of your reverse proxies.
|
||||
|
||||
`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.
|
||||
the connection's own address and the header is ignored. Unset (or
|
||||
empty), the list is the RFC 1918 private ranges: `10.0.0.0/8`,
|
||||
`172.16.0.0/12` and `192.168.0.0/16`. That covers a reverse proxy
|
||||
reaching webhooker over a Docker network or a private LAN without
|
||||
anything set. A set value replaces the default entirely. A set but
|
||||
unparseable value aborts startup.
|
||||
|
||||
That default is safe against forged headers, but leaving it unset in
|
||||
production has a cost you must know about. Production runs behind a
|
||||
TLS-terminating reverse proxy, so with `TRUSTED_PROXIES` unset every
|
||||
request keys on the proxy's own address and all clients share a single
|
||||
Trusting those ranges has two consequences for clients with private
|
||||
addresses:
|
||||
|
||||
- Any such client, whether it connects directly or through the proxy,
|
||||
can choose its own rate-limit key by sending its own
|
||||
`X-Forwarded-For`. A direct client's header is walked because the
|
||||
client is itself trusted; behind the proxy, the client's own address
|
||||
is skipped as a trusted hop when the chain is walked (below), so the
|
||||
entry it wrote is taken as the client. If any of your clients have
|
||||
private addresses, you must set `TRUSTED_PROXIES` to the proxy's
|
||||
address alone.
|
||||
- A client behind the proxy that sends no `X-Forwarded-For` of its own
|
||||
shares the proxy's bucket, because its own address is skipped too.
|
||||
Setting the list to the proxy's address alone gives each its own
|
||||
bucket.
|
||||
|
||||
A proxy the list does not cover, such as nginx on the same host
|
||||
reaching webhooker over loopback, is not trusted: every request through
|
||||
it keys on the proxy's own address and all clients share a single
|
||||
bucket per limit. The receiver limits become service-wide ceilings,
|
||||
and the login endpoint's failure counting collapses onto one key, so a
|
||||
stranger's wrong passwords throttle every other client's wrong
|
||||
passwords.
|
||||
passwords. Set `TRUSTED_PROXIES` to that proxy's address to restore
|
||||
per-client buckets.
|
||||
|
||||
What it cannot do is lock the operator out. The login endpoint
|
||||
verifies credentials **before** it consults any limit and charges only
|
||||
failures, so a correct password is never throttled no matter how full
|
||||
the bucket is. See [Rate Limiting](#rate-limiting).
|
||||
|
||||
The remedy is to set `TRUSTED_PROXIES` to your reverse proxy's
|
||||
address, which restores per-client buckets. webhooker logs a warning
|
||||
at startup whenever `TRUSTED_PROXIES` is empty, in every environment,
|
||||
because behind a proxy every client shares one bucket in `dev` and
|
||||
`prod` alike. The warning is informational when nothing proxies to the
|
||||
process: with no proxy in front, the peer address is the client's own
|
||||
and the buckets are already per-client. See
|
||||
[Rate Limiting](#rate-limiting) for what each limit shares.
|
||||
|
||||
`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
|
||||
@@ -435,14 +442,14 @@ Two operator requirements follow:
|
||||
(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`
|
||||
- Keep clients out of the list. 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.
|
||||
client's bucket. A block that also covers clients — the default, on
|
||||
a network where clients have private addresses — makes every rate
|
||||
limit, including the unauthenticated webhook receiver's, silently
|
||||
bypassable by every client in the block.
|
||||
|
||||
#### Sessions
|
||||
|
||||
@@ -741,10 +748,16 @@ repository's `Dockerfile` and runs it. The app needs:
|
||||
- **Volume:** one host directory mounted at `/var/lib/webhooker`.
|
||||
- **Environment variables:**
|
||||
- `WEBHOOKER_ENVIRONMENT=prod`
|
||||
- `TRUSTED_PROXIES`: your reverse proxy's address on that Docker
|
||||
network. The `remoteIP` field of the `http request` log line for a
|
||||
request that came through the proxy shows it; the health check's
|
||||
own lines show `::1`. See [Trusted proxies](#trusted-proxies).
|
||||
- `TRUSTED_PROXIES`: Docker networks use private addresses, so the
|
||||
default covers your reverse proxy on that network. Under the
|
||||
default, any client with a private address, whether it connects
|
||||
directly or through the proxy, can choose its own rate-limit key
|
||||
by sending its own `X-Forwarded-For`. If any clients have private
|
||||
addresses, or the network's addresses are outside the RFC 1918
|
||||
ranges, set it to the proxy's address there. The `remoteIP` field
|
||||
of the `http request` log line for a request that came through
|
||||
the proxy shows it; the health check's own lines show `::1`. See
|
||||
[Trusted proxies](#trusted-proxies).
|
||||
- Leave `BIND_ADDRESS` and `DATA_DIR` unset: the image sets
|
||||
`BIND_ADDRESS` to `0.0.0.0`, and `DATA_DIR` defaults to
|
||||
`/var/lib/webhooker`.
|
||||
@@ -807,12 +820,14 @@ reports.
|
||||
behind a proxy means the `X-Forwarded-Proto` header. The block below
|
||||
sets it; without it every request is read as plaintext and cookies
|
||||
ship without `Secure`. See [Configuration](#configuration).
|
||||
3. **Set `TRUSTED_PROXIES` to the proxy's address.** Unset, every rate
|
||||
limiter keys on the connecting peer, which behind a proxy is the
|
||||
proxy on every request: all clients collapse into one global bucket
|
||||
per limit and the receiver's per-IP limits become service-wide
|
||||
ceilings. See [Trusted proxies](#trusted-proxies). List the proxy
|
||||
and nothing else.
|
||||
3. **Make sure `TRUSTED_PROXIES` covers the proxy's address.** Unset,
|
||||
it covers the RFC 1918 private ranges, so a proxy on a Docker
|
||||
network or a private LAN is covered and one on loopback is not. For
|
||||
a proxy it does not cover, every rate limiter keys on the connecting
|
||||
peer, which is the proxy on every request: all clients collapse into
|
||||
one global bucket per limit and the receiver's per-IP limits become
|
||||
service-wide ceilings. See [Trusted proxies](#trusted-proxies). If
|
||||
any clients have private addresses, list the proxy and nothing else.
|
||||
4. **Send `Host` as `$http_host`, not `$host`.** `$host` strips the
|
||||
port. webhooker's Origin/Referer check compares against the host it
|
||||
was given, so on any port other than 443 `$host` makes every form
|
||||
@@ -1070,7 +1085,7 @@ unconditionally against whatever files it finds:
|
||||
- the main database on connect — `Setting`, `User`, `APIKey`, `Webhook`,
|
||||
`Entrypoint`, `Target`
|
||||
- each event database when it is lazily opened — `Event`, `Delivery`,
|
||||
`DeliveryResult`, `Totals`
|
||||
`DeliveryResult`
|
||||
- each archive database on every open and reopen
|
||||
|
||||
There is no schema version table, no migration ledger, and no down
|
||||
@@ -1363,10 +1378,11 @@ It uses:
|
||||
- **[go-chi/httprate](https://github.com/go-chi/httprate)** for
|
||||
sliding-window rate limiting of the password-change and webhook
|
||||
receiver endpoints. The bucket is per client IP only when
|
||||
`TRUSTED_PROXIES` names the reverse proxy; unset, every client
|
||||
behind that proxy shares one bucket per limit. The login endpoint
|
||||
counts failed attempts itself instead, so that a correct password is
|
||||
never throttled (see [Rate Limiting](#rate-limiting))
|
||||
`TRUSTED_PROXIES` covers the reverse proxy (by default it covers the
|
||||
RFC 1918 private ranges); otherwise every client behind that proxy
|
||||
shares one bucket per limit. The login endpoint counts failed
|
||||
attempts itself instead, so that a correct password is never
|
||||
throttled (see [Rate Limiting](#rate-limiting))
|
||||
- **[Prometheus](https://prometheus.io)** for metrics, served at
|
||||
`/metrics` behind basic auth
|
||||
- **[Sentry](https://sentry.io)** for optional error reporting
|
||||
@@ -1384,7 +1400,7 @@ The codebase uses consistent naming throughout (rename completed in
|
||||
|
||||
### Data Model
|
||||
|
||||
webhooker's data model has ten entities organized into two tiers: the
|
||||
webhooker's data model has nine entities organized into two tiers: the
|
||||
**application tier** (user and webhook configuration) and the **event
|
||||
tier** (event ingestion, delivery, and logging).
|
||||
|
||||
@@ -1413,10 +1429,6 @@ tier** (event ingestion, delivery, and logging).
|
||||
│ ┌──────────┐ ┌──────────┐ ┌─────────────────┐ │
|
||||
│ │ Event │──1:N──│ Delivery │──1:N──│ DeliveryResult │ │
|
||||
│ └──────────┘ └──────────┘ └─────────────────┘ │
|
||||
│ │
|
||||
│ ┌──────────┐ │
|
||||
│ │ Totals │ (one row of running counts) │
|
||||
│ └──────────┘ │
|
||||
└─────────────────────────────────────────────────────────────┘
|
||||
```
|
||||
|
||||
@@ -1669,7 +1681,6 @@ status across potentially multiple attempts.
|
||||
| `event_id` | UUID | Foreign key → Event |
|
||||
| `target_id`| UUID | Foreign key → Target |
|
||||
| `status` | DeliveryStatus | One of: `pending`, `delivered`, `failed`, `retrying` |
|
||||
| `finished_at` | timestamp | When the delivery became `delivered` or `failed` (nullable; empty while `pending` or `retrying`) |
|
||||
|
||||
**Relations:** Belongs to Event. Belongs to Target. Has many
|
||||
DeliveryResults.
|
||||
@@ -1737,29 +1748,6 @@ retries) is individually logged for full observability.
|
||||
|
||||
**Relations:** Belongs to Delivery.
|
||||
|
||||
#### Totals
|
||||
|
||||
The one row of running counts in each event database, read by the
|
||||
statistics pane at the top of the webhook page.
|
||||
|
||||
| Field | Type | Description |
|
||||
| -------------------- | ------- | ----------- |
|
||||
| `events` | integer | Events ever stored, resubmitted copies included |
|
||||
| `deliveries` | integer | Deliveries ever created, replays included |
|
||||
| `failures` | integer | Deliveries that ever became `failed` |
|
||||
| `events_removed` | integer | Events retention has deleted |
|
||||
| `deliveries_removed` | integer | Deliveries retention has deleted |
|
||||
| `failures_removed` | integer | Failed deliveries retention has deleted |
|
||||
|
||||
Each count changes in the transaction that writes or deletes the rows it
|
||||
counts. The pane shows each of the first three as a lifetime figure, and
|
||||
less what retention removed as the figure within retention, so neither
|
||||
needs the rows themselves. Its last-10-minutes and last-24-hours figures
|
||||
are counted from the `events` and `deliveries` indexes over just that
|
||||
window. Its failure percentage for a window is the deliveries that became
|
||||
`failed` in it out of all that became `delivered` or `failed` in it, and
|
||||
a dash when none did.
|
||||
|
||||
#### Event-tier indexes
|
||||
|
||||
These indexes on the per-webhook event databases are declared in the model
|
||||
@@ -1767,10 +1755,10 @@ tags, so `AutoMigrate` creates them on a fresh and on an existing database:
|
||||
|
||||
| Table | Columns | Serves |
|
||||
| ------------------ | --------------------------- | ------ |
|
||||
| `deliveries` | `status`, `deleted_at`, `finished_at` | Startup recovery, the retry and pending sweeps every 60 seconds and the queue-depth sampler every 30 seconds, which select deliveries by status, and the webhook page's statistics, which count deliveries by status and when they finished |
|
||||
| `deliveries` | `status`, `deleted_at` | Startup recovery, the retry and pending sweeps every 60 seconds and the queue-depth sampler every 30 seconds, which select deliveries by status |
|
||||
| `deliveries` | `event_id`, `deleted_at` | The event log, which loads each event's deliveries, and retention, which selects and deletes the deliveries of expired events |
|
||||
| `delivery_results` | `delivery_id`, `deleted_at` | The event log, which loads the attempts of a page's deliveries, and retention, which deletes the attempts of expired events |
|
||||
| `events` | `deleted_at`, `created_at` | Retention, which selects expired events by age, and the webhook page's statistics, which count recent events and find the newest |
|
||||
| `events` | `deleted_at`, `created_at` | Retention, which selects expired events by age |
|
||||
| `events` | `created_at` | Retention's delete of the expired events themselves |
|
||||
|
||||
GORM's soft delete adds `deleted_at IS NULL` to these queries; retention's
|
||||
@@ -1784,10 +1772,9 @@ and SQLite narrows by a `<` only on the last column it uses.
|
||||
|
||||
#### Common Fields
|
||||
|
||||
Every entity except `Setting` and `Totals` includes these fields from
|
||||
`BaseModel`. `Setting` is a bare key-value row with no `id`, no
|
||||
timestamps and no soft delete, and `Totals` is a single row of counts
|
||||
with only a numeric `id`:
|
||||
Every entity except `Setting` includes these fields from `BaseModel`.
|
||||
`Setting` is a bare key-value row with no `id`, no timestamps and no
|
||||
soft delete:
|
||||
|
||||
| Field | Type | Description |
|
||||
| ------------ | --------- | ----------- |
|
||||
@@ -1829,7 +1816,6 @@ encryption key is generated and stored, and an `admin` user is created.
|
||||
- **Events** — captured incoming webhook payloads
|
||||
- **Deliveries** — event-to-target pairings and their status
|
||||
- **DeliveryResults** — individual delivery attempt logs
|
||||
- **Totals** — running counts of the above, kept through retention
|
||||
|
||||
Per-webhook databases are created automatically when a webhook is
|
||||
created (and lazily on first access for webhooks that predate this
|
||||
@@ -2566,47 +2552,48 @@ the tree is checked out: four checkouts have reported 3,959, 3,961,
|
||||
client-supplied field was cut, and that the shipped chain's stack
|
||||
arrived uncut — never the numbers.
|
||||
|
||||
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. That address becomes a bucket by family: IPv4 keys on the full
|
||||
address, IPv6 on its `/64` prefix. A routed `/64` is the normal
|
||||
Every limiter here — receiver, login, password change, delivery replay
|
||||
and event resubmit — identifies the client the same way, through one
|
||||
shared key function: the connection's own address, unless the peer is
|
||||
inside `TRUSTED_PROXIES`, in which case the forwarded client address is
|
||||
used instead. That address becomes a bucket by family: IPv4 keys on
|
||||
the full address, IPv6 on its `/64` prefix. A routed `/64` is the normal
|
||||
residential and mobile IPv6 allocation, so keying IPv6 per address would
|
||||
let one subscriber rotate source addresses and mint a fresh bucket per
|
||||
request, evading these limits at the network layer without spoofing
|
||||
anything; the cost is that distinct clients inside one `/64` share a
|
||||
bucket. IPv4-mapped addresses (`::ffff:1.2.3.4`) key as the IPv4 address
|
||||
they carry. 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. Set `TRUSTED_PROXIES` to the
|
||||
proxy's address to get per-client limits back. What the shared bucket
|
||||
they carry. See [Trusted proxies](#trusted-proxies). When that variable
|
||||
does not cover the reverse proxy, a client behind it shares one bucket
|
||||
with every other client behind the same proxy. Set `TRUSTED_PROXIES` to
|
||||
the proxy's address to get per-client limits back. What the shared bucket
|
||||
costs is not the same for every limiter, and the two cases pull in
|
||||
opposite directions:
|
||||
|
||||
- For the **receiver** limits it costs throughput, which is the safe
|
||||
direction to be wrong in: sharing can only make a limit bind sooner,
|
||||
never let a sender past it. It matters more for the aggregate limit
|
||||
than for the per-entrypoint one: with `TRUSTED_PROXIES` unset behind
|
||||
the reverse proxy a production deployment is required to run behind,
|
||||
every request keys on the proxy, so the aggregate limit becomes a
|
||||
service-wide ceiling of 1200 requests per minute across all senders
|
||||
and all entrypoints, where the per-entrypoint limit's capacity still
|
||||
grows with the number of entrypoints. Any deployment with more than a
|
||||
handful of busy entrypoints must set `TRUSTED_PROXIES`.
|
||||
than for the per-entrypoint one: when `TRUSTED_PROXIES` does not
|
||||
cover the reverse proxy a production deployment is required to run
|
||||
behind, every request keys on the proxy, so the aggregate limit
|
||||
becomes a service-wide ceiling of 1200 requests per minute across all
|
||||
senders and all entrypoints, where the per-entrypoint limit's
|
||||
capacity still grows with the number of entrypoints. Any deployment
|
||||
with more than a handful of busy entrypoints must make sure
|
||||
`TRUSTED_PROXIES` covers its proxy.
|
||||
- For the **login and password-change** limits it costs precision, not
|
||||
availability. Login failures from every client land in one counter,
|
||||
so a stranger's wrong passwords make the operator's own wrong
|
||||
passwords answer `429` sooner; the operator's _correct_ password is
|
||||
never affected, because it is never counted. Production deployments
|
||||
should still set `TRUSTED_PROXIES`; webhooker warns at startup
|
||||
whenever it is empty, in any environment.
|
||||
should still make sure `TRUSTED_PROXIES` covers their proxy.
|
||||
|
||||
#### The login endpoint
|
||||
|
||||
The login `POST` is the one endpoint with no pre-emptive limiter in
|
||||
front of it, and that is deliberate. A limiter that spends budget on
|
||||
arrival is a lockout in this deployment shape: sharing one bucket, a
|
||||
arrival is a lockout wherever clients share one bucket, as they do
|
||||
behind a reverse proxy that `TRUSTED_PROXIES` does not cover: a
|
||||
stranger sending five POSTs a minute — about 0.08 requests per second,
|
||||
from anywhere — keeps it permanently full, and the operator has no
|
||||
second administrative path. So the handler inverts the order:
|
||||
@@ -2703,8 +2690,10 @@ re-fills both verification slots on its first two requests. The
|
||||
remedies are to block the source at the reverse proxy, or to
|
||||
rate-limit `POST /pages/login` there — the one place a limit can be
|
||||
applied without reintroducing the lockout, because the proxy sees the
|
||||
real client address. Setting `TRUSTED_PROXIES` does not stop the
|
||||
saturation, but it makes the source visible in the failure logs.
|
||||
real client address. `TRUSTED_PROXIES` does not stop the saturation.
|
||||
The flood's source is in the proxy's access log: webhooker's own logs
|
||||
record the proxy's address, not the client's (see
|
||||
[Deployment behind a reverse proxy](#deployment-behind-a-reverse-proxy)).
|
||||
|
||||
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
|
||||
@@ -3061,10 +3050,9 @@ check, see [The login endpoint](#the-login-endpoint).
|
||||
It runs behind session auth, so only a client already holding a
|
||||
valid session reaches it, and an operator throttled out of changing
|
||||
a password can still log in. The bucket is per client IP only when
|
||||
`TRUSTED_PROXIES` names the reverse proxy; unset, every client
|
||||
`TRUSTED_PROXIES` covers the reverse proxy; otherwise every client
|
||||
shares one bucket, which costs precision rather than availability
|
||||
(see [Rate Limiting](#rate-limiting)). webhooker warns at startup
|
||||
whenever `TRUSTED_PROXIES` is empty
|
||||
(see [Rate Limiting](#rate-limiting))
|
||||
- Prometheus metrics behind basic auth
|
||||
- Static assets embedded in binary (no filesystem access needed at
|
||||
runtime)
|
||||
|
||||
+22
-60
@@ -75,6 +75,11 @@ const (
|
||||
// internet-exposed endpoint.
|
||||
defaultReceiverRateLimit = 120
|
||||
|
||||
// defaultTrustedProxies is TRUSTED_PROXIES when it is unset: the
|
||||
// RFC 1918 private ranges, which a reverse proxy reaching the
|
||||
// process over a Docker network or a private LAN connects from.
|
||||
defaultTrustedProxies = "10.0.0.0/8,172.16.0.0/12,192.168.0.0/16"
|
||||
|
||||
// maxPort is the highest valid TCP port number. The lower
|
||||
// bound (at least 1) is enforced by envPositiveInt.
|
||||
maxPort = 65535
|
||||
@@ -172,13 +177,14 @@ type Config struct {
|
||||
|
||||
// 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.
|
||||
// only forwarded header read. Unless TRUSTED_PROXIES is set it
|
||||
// is the RFC 1918 private ranges (defaultTrustedProxies).
|
||||
// Other peers' forwarded headers are ignored and they are
|
||||
// identified by the connection's own address. Under the
|
||||
// default any client with a private address, directly or
|
||||
// through a trusted proxy, can choose its own rate-limit
|
||||
// key, so where any clients have private addresses this must
|
||||
// be set to the proxy hosts alone.
|
||||
TrustedProxies []netip.Prefix
|
||||
|
||||
// AllowedEgressCIDRs is the set of networks a delivery target
|
||||
@@ -460,14 +466,15 @@ func parseCIDR(entry string) (netip.Prefix, error) {
|
||||
|
||||
// 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) {
|
||||
// allowed). An unset, empty, or blank value is read as defaultValue
|
||||
// instead. 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, defaultValue string) ([]netip.Prefix, error) {
|
||||
v := strings.TrimSpace(os.Getenv(key))
|
||||
if v == "" {
|
||||
return nil, nil
|
||||
v = defaultValue
|
||||
}
|
||||
|
||||
var prefixes []netip.Prefix
|
||||
@@ -681,12 +688,12 @@ func loadFromEnv() (*Config, error) {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
trustedProxies, err := envPrefixList("TRUSTED_PROXIES")
|
||||
trustedProxies, err := envPrefixList("TRUSTED_PROXIES", defaultTrustedProxies)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
allowedEgressCIDRs, err := envPrefixList("ALLOWED_EGRESS_CIDRS")
|
||||
allowedEgressCIDRs, err := envPrefixList("ALLOWED_EGRESS_CIDRS", "")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -760,50 +767,6 @@ func (c *Config) warnEgressAllowlist(log *slog.Logger) {
|
||||
)
|
||||
}
|
||||
|
||||
// warnSharedRateLimitBucket logs a startup warning whenever
|
||||
// TRUSTED_PROXIES is empty, in any environment.
|
||||
//
|
||||
// With no trusted proxies every rate limiter keys on the connecting
|
||||
// peer's address. Whether that is harmless or dangerous depends on
|
||||
// what is in front of the process, which this code cannot observe:
|
||||
// with nothing in front, the peer is the client and the limits are
|
||||
// per-client as intended; behind a reverse proxy the peer is the proxy
|
||||
// for every request, so all clients share one bucket per limiter.
|
||||
//
|
||||
// The login endpoint no longer spends budget on arrival — it verifies
|
||||
// credentials first and charges only failures — so a shared bucket
|
||||
// cannot deny the operator a correct password. What it does collapse
|
||||
// is the failure counting: one client's wrong passwords throttle
|
||||
// everyone else's wrong passwords, and the receiver's limits become
|
||||
// service-wide ceilings.
|
||||
//
|
||||
// The warning is deliberately not gated on WEBHOOKER_ENVIRONMENT:
|
||||
// behind a proxy every client shares one bucket in dev and prod alike.
|
||||
//
|
||||
// The default of trusting nobody is deliberate — trusting forwarded
|
||||
// headers from arbitrary peers lets any client choose its own bucket —
|
||||
// so this warns rather than failing startup or changing the key.
|
||||
func (c *Config) warnSharedRateLimitBucket(log *slog.Logger) {
|
||||
if len(c.TrustedProxies) > 0 {
|
||||
return
|
||||
}
|
||||
|
||||
log.Warn(
|
||||
"TRUSTED_PROXIES is empty: every rate limit keys on the "+
|
||||
"connecting peer's address. With nothing proxying to "+
|
||||
"this process that is the client itself and the limits "+
|
||||
"are per-client as intended. Behind a reverse proxy the "+
|
||||
"peer is the proxy on every request, so all clients "+
|
||||
"share one bucket per limit: the receiver limits become "+
|
||||
"service-wide ceilings, and one client's failed logins "+
|
||||
"throttle every other client's failed logins — a "+
|
||||
"correct password still gets in. If anything proxies to "+
|
||||
"this process, set TRUSTED_PROXIES to its address.",
|
||||
"environment", c.Environment,
|
||||
"trustedProxies", len(c.TrustedProxies),
|
||||
)
|
||||
}
|
||||
|
||||
// New creates a Config by reading environment variables.
|
||||
//
|
||||
//nolint:revive // lc parameter is required by fx even if unused.
|
||||
@@ -849,7 +812,6 @@ func New(lc fx.Lifecycle, params ConfigParams) (*Config, error) {
|
||||
"hasMetricsAuth", s.MetricsAuthEnabled(),
|
||||
)
|
||||
|
||||
s.warnSharedRateLimitBucket(log)
|
||||
s.warnEgressAllowlist(log)
|
||||
|
||||
return s, nil
|
||||
|
||||
+14
-101
@@ -551,6 +551,11 @@ func testReceiverRateLimitSuccess(
|
||||
}
|
||||
|
||||
func TestTrustedProxies(t *testing.T) {
|
||||
// Unset, the RFC 1918 private ranges are trusted, so a reverse
|
||||
// proxy on a Docker network or a private LAN is covered without
|
||||
// configuration.
|
||||
defaultProxies := []string{cidrPrivateV4, "172.16.0.0/12", "192.168.0.0/16"}
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
set bool
|
||||
@@ -559,18 +564,21 @@ func TestTrustedProxies(t *testing.T) {
|
||||
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{},
|
||||
expected: defaultProxies,
|
||||
},
|
||||
{
|
||||
name: "blank value trusts nothing",
|
||||
name: "blank value uses default",
|
||||
set: true,
|
||||
value: " ",
|
||||
expected: []string{},
|
||||
expected: defaultProxies,
|
||||
},
|
||||
{
|
||||
name: "set value replaces the default entirely",
|
||||
set: true,
|
||||
value: "203.0.113.7",
|
||||
expected: []string{"203.0.113.7/32"},
|
||||
},
|
||||
{
|
||||
name: caseValidValueParsed,
|
||||
@@ -845,101 +853,6 @@ func TestEgressAllowlistWarning(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestSharedRateLimitBucketWarning covers the startup warning that
|
||||
// tells an operator a deployment behind a reverse proxy shares one
|
||||
// rate-limit bucket between every client, which turns the receiver
|
||||
// limits into service-wide ceilings and collapses login failure
|
||||
// counting. It must fire whenever TRUSTED_PROXIES is empty, in any
|
||||
// environment, because behind a proxy every client shares one bucket
|
||||
// in dev and prod alike. It stays quiet once proxies are named.
|
||||
func TestSharedRateLimitBucketWarning(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
environment string
|
||||
trustedProxies string
|
||||
expectWarning bool
|
||||
}{
|
||||
{
|
||||
name: "prod without trusted proxies warns",
|
||||
environment: config.EnvironmentProd,
|
||||
expectWarning: true,
|
||||
},
|
||||
{
|
||||
name: "prod with trusted proxies is quiet",
|
||||
environment: config.EnvironmentProd,
|
||||
trustedProxies: cidrPrivateV4,
|
||||
expectWarning: false,
|
||||
},
|
||||
{
|
||||
name: "dev without trusted proxies warns",
|
||||
environment: config.EnvironmentDev,
|
||||
expectWarning: true,
|
||||
},
|
||||
{
|
||||
name: "dev with trusted proxies is quiet",
|
||||
environment: config.EnvironmentDev,
|
||||
trustedProxies: cidrPrivateV4,
|
||||
expectWarning: false,
|
||||
},
|
||||
}
|
||||
|
||||
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", tt.environment)
|
||||
|
||||
if tt.trustedProxies == "" {
|
||||
require.NoError(
|
||||
t, os.Unsetenv("TRUSTED_PROXIES"),
|
||||
)
|
||||
} else {
|
||||
t.Setenv("TRUSTED_PROXIES", tt.trustedProxies)
|
||||
}
|
||||
|
||||
var buf bytes.Buffer
|
||||
|
||||
log := slog.New(slog.NewJSONHandler(
|
||||
&buf, &slog.HandlerOptions{
|
||||
Level: slog.LevelDebug,
|
||||
},
|
||||
))
|
||||
|
||||
require.NoError(
|
||||
t,
|
||||
config.WarnSharedRateLimitBucketForTest(log),
|
||||
)
|
||||
|
||||
if !tt.expectWarning {
|
||||
assert.Empty(t, buf.String())
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
logged := buf.String()
|
||||
|
||||
assert.Contains(t, logged, `"level":"WARN"`)
|
||||
assert.Contains(t, logged, "TRUSTED_PROXIES")
|
||||
assert.Contains(t, logged, "share one bucket")
|
||||
assert.Contains(
|
||||
t, logged, "throttle every other client's failed logins",
|
||||
)
|
||||
// The warning must not claim a lockout the login
|
||||
// endpoint no longer permits: credentials are verified
|
||||
// before any budget is spent.
|
||||
assert.Contains(
|
||||
t, logged, "a correct password still gets in",
|
||||
)
|
||||
// The text must stay accurate for a developer with
|
||||
// nothing in front of the process, where an empty
|
||||
// list costs nothing.
|
||||
assert.Contains(
|
||||
t, logged, "nothing proxying to this process",
|
||||
)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// metricsEnv describes what one subtest below puts in the
|
||||
// environment for a single METRICS_ variable. A variable that is
|
||||
// set to the empty string and one that is not set at all are
|
||||
|
||||
@@ -6,21 +6,6 @@ import "log/slog"
|
||||
// the external config_test package so each helper can be covered by
|
||||
// its own table-driven test without weakening the package API.
|
||||
|
||||
// WarnSharedRateLimitBucketForTest loads a Config from the current
|
||||
// environment and emits its startup warnings to log. The real logger
|
||||
// writes to stdout, so this lets the warning's firing condition be
|
||||
// asserted against a handler the test controls.
|
||||
func WarnSharedRateLimitBucketForTest(log *slog.Logger) error {
|
||||
c, err := loadFromEnv()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
c.warnSharedRateLimitBucket(log)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// WarnEgressAllowlistForTest loads a Config from the current
|
||||
// environment and emits its egress-allowlist startup warning to
|
||||
// log, so a test can assert both that the warning fires only when
|
||||
|
||||
@@ -93,7 +93,6 @@ func TestEventTierQueriesUseTheirIndexes(t *testing.T) {
|
||||
deliveries []database.Delivery
|
||||
results []database.DeliveryResult
|
||||
depths []struct{ Depth int }
|
||||
failed struct{ Count int64 }
|
||||
)
|
||||
|
||||
byStatus := "idx_deliveries_status (status=? AND deleted_at=?)"
|
||||
@@ -124,7 +123,7 @@ func TestEventTierQueriesUseTheirIndexes(t *testing.T) {
|
||||
Order("attempt_num ASC").Find(&results),
|
||||
"idx_delivery_results_delivery_id (delivery_id=? AND deleted_at=?)")
|
||||
|
||||
// Retention's deletes (deleteExpired), whose subqueries are built
|
||||
// Retention's three deletes (reapExpired), whose subqueries are built
|
||||
// afresh for each statement as it builds them.
|
||||
expiredEventIDs := func() *gorm.DB {
|
||||
return dry.Model(&database.Event{}).Select("id").
|
||||
@@ -136,12 +135,6 @@ func TestEventTierQueriesUseTheirIndexes(t *testing.T) {
|
||||
Select("id").Where("event_id IN (?)", expiredEventIDs()),
|
||||
).Delete(&database.DeliveryResult{}),
|
||||
"idx_delivery_results_delivery_id (delivery_id=?)", byEvent, byAge)
|
||||
assertPlanUses(t, db, dry.Unscoped().Model(&database.Delivery{}).
|
||||
Select("count(CASE WHEN status = ? THEN 1 END) AS count",
|
||||
database.DeliveryStatusFailed).
|
||||
Where("event_id IN (?)", expiredEventIDs()).
|
||||
Take(&failed),
|
||||
"idx_deliveries_event_id (event_id=?)", byAge)
|
||||
assertPlanUses(t, db, dry.Unscoped().Where(
|
||||
"event_id IN (?)", expiredEventIDs(),
|
||||
).Delete(&database.Delivery{}),
|
||||
@@ -151,54 +144,6 @@ func TestEventTierQueriesUseTheirIndexes(t *testing.T) {
|
||||
).Delete(&database.Event{}), "idx_events_created_at (created_at<?)")
|
||||
}
|
||||
|
||||
// TestStatisticsQueriesUseTheirIndexes does the same for the webhook
|
||||
// page's statistics (readEventStats in the handlers): deliveries in
|
||||
// progress, deliveries finished and events received since a time, and
|
||||
// the newest event, which must come straight off an index rather than
|
||||
// from sorting every event.
|
||||
func TestStatisticsQueriesUseTheirIndexes(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
mgr, lc := setupTestWebhookDBManager(t)
|
||||
ctx := context.Background()
|
||||
require.NoError(t, lc.Start(ctx))
|
||||
|
||||
defer func() { require.NoError(t, lc.Stop(ctx)) }()
|
||||
|
||||
db, err := mgr.GetDB(uuid.New().String())
|
||||
require.NoError(t, err)
|
||||
|
||||
dry := db.Session(&gorm.Session{DryRun: true})
|
||||
since := time.Now()
|
||||
|
||||
var (
|
||||
count int64
|
||||
newest []time.Time
|
||||
)
|
||||
|
||||
assertPlanUses(t, db, dry.Model(&database.Delivery{}).
|
||||
Where("status IN ?", []database.DeliveryStatus{
|
||||
database.DeliveryStatusPending,
|
||||
database.DeliveryStatusRetrying,
|
||||
}).Count(&count),
|
||||
"idx_deliveries_status (status=? AND deleted_at=?)")
|
||||
assertPlanUses(t, db, dry.Model(&database.Delivery{}).
|
||||
Where("status = ? AND finished_at >= ?",
|
||||
database.DeliveryStatusFailed, since).Count(&count),
|
||||
"idx_deliveries_status "+
|
||||
"(status=? AND deleted_at=? AND finished_at>?)")
|
||||
assertPlanUses(t, db, dry.Model(&database.Event{}).
|
||||
Where("created_at >= ?", since).Count(&count),
|
||||
"idx_events_deleted_at_created_at "+
|
||||
"(deleted_at=? AND created_at>?)")
|
||||
|
||||
newestEvent := dry.Model(&database.Event{}).
|
||||
Order("created_at DESC").Limit(1).Pluck("created_at", &newest)
|
||||
assertPlanUses(t, db, newestEvent,
|
||||
"idx_events_deleted_at_created_at (deleted_at=?)")
|
||||
assert.NotContains(t, queryPlan(t, db, newestEvent), "TEMP B-TREE")
|
||||
}
|
||||
|
||||
// assertPlanUses asserts that SQLite's plan for a statement GORM built
|
||||
// in a dry run, run with the same SQL and arguments GORM would send,
|
||||
// names each of the given indexes.
|
||||
@@ -207,18 +152,6 @@ func assertPlanUses(
|
||||
) {
|
||||
t.Helper()
|
||||
|
||||
plan := queryPlan(t, db, built)
|
||||
|
||||
for _, index := range indexes {
|
||||
assert.Contains(t, plan, index, built.Statement.SQL.String())
|
||||
}
|
||||
}
|
||||
|
||||
// queryPlan returns SQLite's plan for a statement GORM built in a dry
|
||||
// run, run with the same SQL and arguments GORM would send.
|
||||
func queryPlan(t *testing.T, db, built *gorm.DB) string {
|
||||
t.Helper()
|
||||
|
||||
var plan []struct{ Detail string }
|
||||
|
||||
require.NoError(t, db.Raw(
|
||||
@@ -226,5 +159,8 @@ func queryPlan(t *testing.T, db, built *gorm.DB) string {
|
||||
built.Statement.Vars...,
|
||||
).Scan(&plan).Error)
|
||||
|
||||
return fmt.Sprint(plan)
|
||||
for _, index := range indexes {
|
||||
assert.Contains(t, fmt.Sprint(plan), index,
|
||||
built.Statement.SQL.String())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,10 +1,6 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
import "gorm.io/gorm"
|
||||
|
||||
// DeliveryStatus represents the status of a delivery
|
||||
type DeliveryStatus string
|
||||
@@ -49,12 +45,6 @@ type Delivery struct {
|
||||
// gives.
|
||||
DeletedAt gorm.DeletedAt `gorm:"index:idx_deliveries_event_id,priority:2;index:idx_deliveries_status,priority:2" json:"deletedAt,omitzero"`
|
||||
|
||||
// FinishedAt is when the delivery became delivered or failed, and
|
||||
// nil while it is pending or retrying. It ends the status index,
|
||||
// so the webhook page counts the deliveries that finished in a
|
||||
// recent window by reading that window from the index.
|
||||
FinishedAt *time.Time `gorm:"index:idx_deliveries_status,priority:3" json:"finishedAt,omitempty"`
|
||||
|
||||
// Relations
|
||||
Event Event `json:"event,omitzero"`
|
||||
Target Target `json:"target,omitzero"`
|
||||
|
||||
@@ -1,73 +0,0 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// Totals is the single row of running totals in a webhook's event
|
||||
// database. It is what keeps the webhook page's lifetime figures right
|
||||
// after retention has removed the rows they count, and what lets the
|
||||
// page show them without counting every row.
|
||||
//
|
||||
// Storing an event, creating a delivery and failing a delivery each
|
||||
// add one, and retention adds what it deletes to the Removed columns.
|
||||
// Every addition goes through AddTotals, in the transaction that
|
||||
// writes or deletes the rows it counts.
|
||||
type Totals struct {
|
||||
ID int64 `gorm:"primaryKey"`
|
||||
|
||||
Events int64 `gorm:"not null"`
|
||||
Deliveries int64 `gorm:"not null"`
|
||||
Failures int64 `gorm:"not null"`
|
||||
|
||||
EventsRemoved int64 `gorm:"not null"`
|
||||
DeliveriesRemoved int64 `gorm:"not null"`
|
||||
FailuresRemoved int64 `gorm:"not null"`
|
||||
}
|
||||
|
||||
// TableName names the table AddTotals updates.
|
||||
func (Totals) TableName() string {
|
||||
return "totals"
|
||||
}
|
||||
|
||||
// EventsWithinRetention is how many of the webhook's events are still
|
||||
// stored.
|
||||
func (t Totals) EventsWithinRetention() int64 {
|
||||
return t.Events - t.EventsRemoved
|
||||
}
|
||||
|
||||
// DeliveriesWithinRetention is how many of the webhook's deliveries
|
||||
// are still stored.
|
||||
func (t Totals) DeliveriesWithinRetention() int64 {
|
||||
return t.Deliveries - t.DeliveriesRemoved
|
||||
}
|
||||
|
||||
// FailuresWithinRetention is how many of the webhook's failed
|
||||
// deliveries are still stored.
|
||||
func (t Totals) FailuresWithinRetention() int64 {
|
||||
return t.Failures - t.FailuresRemoved
|
||||
}
|
||||
|
||||
// AddTotals adds each count in add to the webhook's running totals.
|
||||
// Call it on the transaction that writes or deletes the rows it
|
||||
// counts, so the totals change exactly when those rows do.
|
||||
func AddTotals(tx *gorm.DB, add Totals) error {
|
||||
err := tx.Exec(
|
||||
`UPDATE totals SET
|
||||
events = events + ?,
|
||||
deliveries = deliveries + ?,
|
||||
failures = failures + ?,
|
||||
events_removed = events_removed + ?,
|
||||
deliveries_removed = deliveries_removed + ?,
|
||||
failures_removed = failures_removed + ?`,
|
||||
add.Events, add.Deliveries, add.Failures,
|
||||
add.EventsRemoved, add.DeliveriesRemoved, add.FailuresRemoved,
|
||||
).Error
|
||||
if err != nil {
|
||||
return fmt.Errorf("adding to running totals: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -2,7 +2,7 @@ package database
|
||||
|
||||
// Migrate runs database migrations for the main application database.
|
||||
// Only configuration-tier models are stored in the main database.
|
||||
// Event-tier models (Event, Delivery, DeliveryResult, Totals) live in
|
||||
// Event-tier models (Event, Delivery, DeliveryResult) live in
|
||||
// per-webhook dedicated databases managed by WebhookDBManager.
|
||||
func (d *Database) Migrate() error {
|
||||
return d.db.AutoMigrate(
|
||||
|
||||
@@ -267,101 +267,55 @@ func retentionCutoff(
|
||||
|
||||
// reapExpired hard-deletes, in foreign-key-safe order, the delivery
|
||||
// results, deliveries, and events associated with events older than
|
||||
// cutoff, and adds what it deleted to the running totals, all in one
|
||||
// transaction. Deletes are unscoped so rows are physically removed
|
||||
// rather than soft-deleted, reclaiming disk. It returns the number of
|
||||
// events deleted.
|
||||
// cutoff. Deletes are unscoped so rows are physically removed rather
|
||||
// than soft-deleted, reclaiming disk. It returns the number of events
|
||||
// deleted.
|
||||
func reapExpired(db *gorm.DB, cutoff time.Time) (int64, error) {
|
||||
var removed Totals
|
||||
|
||||
err := db.Transaction(func(tx *gorm.DB) error {
|
||||
var err error
|
||||
|
||||
removed, err = deleteExpired(tx, cutoff)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return AddTotals(tx, removed)
|
||||
})
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
return removed.EventsRemoved, nil
|
||||
}
|
||||
|
||||
// deleteExpired runs reapExpired's deletes and returns how many
|
||||
// events, deliveries and failed deliveries they removed.
|
||||
func deleteExpired(tx *gorm.DB, cutoff time.Time) (Totals, error) {
|
||||
var removed Totals
|
||||
|
||||
// Fresh subqueries are built per statement to avoid reusing a
|
||||
// mutated builder across executions.
|
||||
expiredEventIDs := func() *gorm.DB {
|
||||
return tx.Model(&Event{}).
|
||||
return db.Model(&Event{}).
|
||||
Select("id").
|
||||
Where("created_at < ?", cutoff)
|
||||
}
|
||||
expiredDeliveryIDs := func() *gorm.DB {
|
||||
return tx.Model(&Delivery{}).
|
||||
return db.Model(&Delivery{}).
|
||||
Select("id").
|
||||
Where("event_id IN (?)", expiredEventIDs())
|
||||
}
|
||||
|
||||
// 1. Delivery results whose delivery belongs to an expired event.
|
||||
res := tx.Unscoped().
|
||||
res := db.Unscoped().
|
||||
Where("delivery_id IN (?)", expiredDeliveryIDs()).
|
||||
Delete(&DeliveryResult{})
|
||||
if res.Error != nil {
|
||||
return removed, fmt.Errorf(
|
||||
return 0, fmt.Errorf(
|
||||
"deleting expired delivery results: %w",
|
||||
res.Error,
|
||||
)
|
||||
}
|
||||
|
||||
// 2. Deliveries belonging to an expired event, after counting the
|
||||
// failed ones among them. The status is tested in the select list
|
||||
// rather than the WHERE clause: there, SQLite would read every
|
||||
// failed delivery the webhook has through the status index,
|
||||
// instead of only the expired ones through the event_id index.
|
||||
var failed struct{ Count int64 }
|
||||
|
||||
err := tx.Unscoped().Model(&Delivery{}).
|
||||
Select("count(CASE WHEN status = ? THEN 1 END) AS count",
|
||||
DeliveryStatusFailed).
|
||||
Where("event_id IN (?)", expiredEventIDs()).
|
||||
Take(&failed).Error
|
||||
if err != nil {
|
||||
return removed, fmt.Errorf(
|
||||
"counting expired failed deliveries: %w", err,
|
||||
)
|
||||
}
|
||||
|
||||
del := tx.Unscoped().
|
||||
// 2. Deliveries belonging to an expired event.
|
||||
del := db.Unscoped().
|
||||
Where("event_id IN (?)", expiredEventIDs()).
|
||||
Delete(&Delivery{})
|
||||
if del.Error != nil {
|
||||
return removed, fmt.Errorf(
|
||||
return 0, fmt.Errorf(
|
||||
"deleting expired deliveries: %w",
|
||||
del.Error,
|
||||
)
|
||||
}
|
||||
|
||||
// 3. The expired events themselves.
|
||||
ev := tx.Unscoped().
|
||||
ev := db.Unscoped().
|
||||
Where("created_at < ?", cutoff).
|
||||
Delete(&Event{})
|
||||
if ev.Error != nil {
|
||||
return removed, fmt.Errorf(
|
||||
return 0, fmt.Errorf(
|
||||
"deleting expired events: %w",
|
||||
ev.Error,
|
||||
)
|
||||
}
|
||||
|
||||
removed.EventsRemoved = ev.RowsAffected
|
||||
removed.DeliveriesRemoved = del.RowsAffected
|
||||
removed.FailuresRemoved = failed.Count
|
||||
|
||||
return removed, nil
|
||||
return ev.RowsAffected, nil
|
||||
}
|
||||
|
||||
@@ -1,126 +0,0 @@
|
||||
package database_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gorm.io/gorm"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
)
|
||||
|
||||
// readTotals reads a webhook database's row of running totals,
|
||||
// asserting that it has exactly one.
|
||||
func readTotals(t *testing.T, db *gorm.DB) database.Totals {
|
||||
t.Helper()
|
||||
|
||||
var rows []database.Totals
|
||||
|
||||
require.NoError(t, db.Find(&rows).Error)
|
||||
require.Len(t, rows, 1)
|
||||
|
||||
return rows[0]
|
||||
}
|
||||
|
||||
// TestWebhookDBManager_TotalsRowSurvivesReopen verifies that a new
|
||||
// event database starts with one row of zero totals, and that opening
|
||||
// it again keeps that row and what was added to it.
|
||||
func TestWebhookDBManager_TotalsRowSurvivesReopen(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
mgr, lc := setupTestWebhookDBManager(t)
|
||||
ctx := context.Background()
|
||||
require.NoError(t, lc.Start(ctx))
|
||||
|
||||
defer func() { require.NoError(t, lc.Stop(ctx)) }()
|
||||
|
||||
webhookID := uuid.New().String()
|
||||
|
||||
db, err := mgr.GetDB(webhookID)
|
||||
require.NoError(t, err)
|
||||
|
||||
fresh := readTotals(t, db)
|
||||
assert.Equal(t, database.Totals{ID: fresh.ID}, fresh)
|
||||
|
||||
require.NoError(t, database.AddTotals(db, database.Totals{
|
||||
Events: 2, Deliveries: 3, Failures: 1,
|
||||
}))
|
||||
|
||||
// Drop the cached connection so the next open reopens the file,
|
||||
// as a restart would.
|
||||
require.NoError(t, mgr.CloseAll())
|
||||
|
||||
db, err = mgr.GetDB(webhookID)
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.Equal(t, database.Totals{
|
||||
ID: fresh.ID, Events: 2, Deliveries: 3, Failures: 1,
|
||||
}, readTotals(t, db))
|
||||
}
|
||||
|
||||
// TestRetentionReaper_AddsWhatItRemovesToTotals verifies that a sweep
|
||||
// leaves the lifetime totals alone and adds the events, deliveries and
|
||||
// failed deliveries it deletes to the removed totals, so the totals
|
||||
// within retention match the rows still stored.
|
||||
func TestRetentionReaper_AddsWhatItRemovesToTotals(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
env := setupRetentionTest(t)
|
||||
|
||||
webhookID := createWebhook(t, env.mainDB.DB(), 30)
|
||||
|
||||
db, err := env.mgr.GetDB(webhookID)
|
||||
require.NoError(t, err)
|
||||
|
||||
now := time.Now()
|
||||
expired := now.Add(-40 * 24 * time.Hour)
|
||||
|
||||
seedEventChain(t, db, webhookID, expired)
|
||||
expiredFailure := seedEventChain(t, db, webhookID, expired)
|
||||
recentFailure := seedEventChain(
|
||||
t, db, webhookID, now.Add(-24*time.Hour),
|
||||
)
|
||||
|
||||
for _, id := range []string{
|
||||
expiredFailure.deliveryID, recentFailure.deliveryID,
|
||||
} {
|
||||
require.NoError(t, db.Model(&database.Delivery{}).
|
||||
Where("id = ?", id).
|
||||
Update("status", database.DeliveryStatusFailed).Error)
|
||||
}
|
||||
|
||||
// The totals storing those rows would have left.
|
||||
require.NoError(t, database.AddTotals(db, database.Totals{
|
||||
Events: 3, Deliveries: 3, Failures: 2,
|
||||
}))
|
||||
|
||||
env.reaper.ExportSweep(context.Background())
|
||||
|
||||
totals := readTotals(t, db)
|
||||
assert.Equal(t, database.Totals{
|
||||
ID: totals.ID,
|
||||
Events: 3, Deliveries: 3, Failures: 2,
|
||||
EventsRemoved: 2, DeliveriesRemoved: 2, FailuresRemoved: 1,
|
||||
}, totals)
|
||||
|
||||
var events, deliveries, failures int64
|
||||
|
||||
require.NoError(t, db.Model(&database.Event{}).Count(&events).Error)
|
||||
require.NoError(t, db.Model(&database.Delivery{}).
|
||||
Count(&deliveries).Error)
|
||||
require.NoError(t, db.Model(&database.Delivery{}).
|
||||
Where("status = ?", database.DeliveryStatusFailed).
|
||||
Count(&failures).Error)
|
||||
|
||||
assert.Equal(t, events, totals.EventsWithinRetention())
|
||||
assert.Equal(t, deliveries, totals.DeliveriesWithinRetention())
|
||||
assert.Equal(t, failures, totals.FailuresWithinRetention())
|
||||
|
||||
// A sweep with nothing left to remove changes nothing.
|
||||
env.reaper.ExportSweep(context.Background())
|
||||
|
||||
assert.Equal(t, totals, readTotals(t, db))
|
||||
}
|
||||
@@ -35,8 +35,7 @@ var errInvalidCachedDBType = errors.New(
|
||||
|
||||
// WebhookDBManager manages per-webhook SQLite database files
|
||||
// for event storage. Each webhook gets its own dedicated
|
||||
// database containing Events, Deliveries, DeliveryResults and the
|
||||
// running Totals of them.
|
||||
// database containing Events, Deliveries, and DeliveryResults.
|
||||
// Database connections are opened lazily and cached.
|
||||
type WebhookDBManager struct {
|
||||
dataDir string
|
||||
@@ -295,7 +294,7 @@ func (m *WebhookDBManager) openDB(
|
||||
|
||||
// Run migrations for event-tier models only
|
||||
err = db.AutoMigrate(
|
||||
&Event{}, &Delivery{}, &DeliveryResult{}, &Totals{},
|
||||
&Event{}, &Delivery{}, &DeliveryResult{},
|
||||
)
|
||||
if err != nil {
|
||||
_ = sqlDB.Close()
|
||||
@@ -306,17 +305,6 @@ func (m *WebhookDBManager) openDB(
|
||||
)
|
||||
}
|
||||
|
||||
// A new database gets its row of running totals, all zero.
|
||||
err = db.FirstOrCreate(&Totals{}).Error
|
||||
if err != nil {
|
||||
_ = sqlDB.Close()
|
||||
|
||||
return nil, fmt.Errorf(
|
||||
"creating running totals for webhook database %s: %w",
|
||||
webhookID, err,
|
||||
)
|
||||
}
|
||||
|
||||
m.log.Info(
|
||||
"opened per-webhook database",
|
||||
"webhook_id", webhookID,
|
||||
|
||||
@@ -1,100 +0,0 @@
|
||||
package delivery_test
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gorm.io/gorm"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
)
|
||||
|
||||
// failureTotal reads the running failure total of a webhook database.
|
||||
func failureTotal(t *testing.T, db *gorm.DB) int64 {
|
||||
t.Helper()
|
||||
|
||||
var totals database.Totals
|
||||
|
||||
require.NoError(t, db.Take(&totals).Error)
|
||||
|
||||
return totals.Failures
|
||||
}
|
||||
|
||||
// TestUpdateDeliveryStatus_FinishTimeAndFailureTotal pins what a status
|
||||
// write records for the webhook page's statistics: the time a delivery
|
||||
// finished, set only when it becomes delivered or failed, and one more
|
||||
// on the failure total when it fails.
|
||||
func TestUpdateDeliveryStatus_FinishTimeAndFailureTotal(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
tests := []struct {
|
||||
status database.DeliveryStatus
|
||||
finished bool
|
||||
failures int64
|
||||
}{
|
||||
{database.DeliveryStatusRetrying, false, 0},
|
||||
{database.DeliveryStatusDelivered, true, 0},
|
||||
{database.DeliveryStatusFailed, true, 1},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(string(tt.status), func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db := testWebhookDB(t)
|
||||
e := testEngine(t, 1)
|
||||
event := seedEvent(t, db, `{}`)
|
||||
d := seedDelivery(
|
||||
t, db, event.ID, uuid.New().String(),
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
before := time.Now()
|
||||
|
||||
require.NoError(t, e.ExportUpdateDeliveryStatus(
|
||||
db, &d, tt.status,
|
||||
))
|
||||
|
||||
var stored database.Delivery
|
||||
|
||||
require.NoError(t, db.First(&stored, "id = ?", d.ID).Error)
|
||||
assert.Equal(t, tt.status, stored.Status)
|
||||
|
||||
if tt.finished {
|
||||
require.NotNil(t, stored.FinishedAt)
|
||||
assert.False(t, stored.FinishedAt.Before(before))
|
||||
} else {
|
||||
assert.Nil(t, stored.FinishedAt)
|
||||
}
|
||||
|
||||
assert.Equal(t, tt.failures, failureTotal(t, db))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestUpdateDeliveryStatus_DeletedDeliveryIsNotCounted covers a
|
||||
// delivery retention deleted while the engine still held it. Failing
|
||||
// it afterwards writes no row, so it adds no failure either: retention
|
||||
// has already counted what it removed.
|
||||
func TestUpdateDeliveryStatus_DeletedDeliveryIsNotCounted(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db := testWebhookDB(t)
|
||||
e := testEngine(t, 1)
|
||||
event := seedEvent(t, db, `{}`)
|
||||
d := seedDelivery(
|
||||
t, db, event.ID, uuid.New().String(),
|
||||
database.DeliveryStatusRetrying,
|
||||
)
|
||||
|
||||
require.NoError(t, db.Unscoped().
|
||||
Delete(&database.Delivery{}, "id = ?", d.ID).Error)
|
||||
|
||||
require.NoError(t, e.ExportUpdateDeliveryStatus(
|
||||
db, &d, database.DeliveryStatusFailed,
|
||||
))
|
||||
|
||||
assert.Zero(t, failureTotal(t, db))
|
||||
}
|
||||
@@ -1554,9 +1554,8 @@ func (e *Engine) updateDeliveryStatus(
|
||||
targetType database.TargetType,
|
||||
status database.DeliveryStatus,
|
||||
) error {
|
||||
err := webhookDB.Transaction(func(tx *gorm.DB) error {
|
||||
return writeDeliveryStatus(tx, d, status)
|
||||
})
|
||||
err := webhookDB.Model(d).
|
||||
Update("status", status).Error
|
||||
if err != nil {
|
||||
return fmt.Errorf(
|
||||
"updating delivery %s to status %s: %w",
|
||||
@@ -1575,33 +1574,6 @@ func (e *Engine) updateDeliveryStatus(
|
||||
return nil
|
||||
}
|
||||
|
||||
// writeDeliveryStatus writes a delivery's new status. A delivery that
|
||||
// becomes delivered or failed also gets the time it finished, and a
|
||||
// failed one is added to the webhook's running failure total. The
|
||||
// failure is counted only if the row was still there to update:
|
||||
// retention may have deleted it while the engine was working on it.
|
||||
func writeDeliveryStatus(
|
||||
tx *gorm.DB,
|
||||
d *database.Delivery,
|
||||
status database.DeliveryStatus,
|
||||
) error {
|
||||
columns := map[string]any{"status": status}
|
||||
if status.Terminal() {
|
||||
columns["finished_at"] = time.Now()
|
||||
}
|
||||
|
||||
res := tx.Model(d).Updates(columns)
|
||||
if res.Error != nil {
|
||||
return res.Error
|
||||
}
|
||||
|
||||
if status == database.DeliveryStatusFailed && res.RowsAffected > 0 {
|
||||
return database.AddTotals(tx, database.Totals{Failures: 1})
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// settleStatus moves a delivery to its outcome status and reports a
|
||||
// failed write through bookkeepingFailed, which leaves the row
|
||||
// recoverable. It exists so the target call sites read as one
|
||||
|
||||
@@ -57,9 +57,7 @@ func testWebhookDB(t *testing.T) *gorm.DB {
|
||||
&database.Event{},
|
||||
&database.Delivery{},
|
||||
&database.DeliveryResult{},
|
||||
&database.Totals{},
|
||||
))
|
||||
require.NoError(t, db.Create(&database.Totals{}).Error)
|
||||
|
||||
return db
|
||||
}
|
||||
|
||||
@@ -150,16 +150,6 @@ func (e *Engine) ExportDeliverSlack(
|
||||
)
|
||||
}
|
||||
|
||||
// ExportUpdateDeliveryStatus exposes updateDeliveryStatus. It passes no
|
||||
// target type, so no metric moves.
|
||||
func (e *Engine) ExportUpdateDeliveryStatus(
|
||||
webhookDB *gorm.DB,
|
||||
d *database.Delivery,
|
||||
status database.DeliveryStatus,
|
||||
) error {
|
||||
return e.updateDeliveryStatus(webhookDB, d, "", status)
|
||||
}
|
||||
|
||||
// ExportProcessNewTask exposes processNewTask.
|
||||
func (e *Engine) ExportProcessNewTask(
|
||||
ctx context.Context, task *Task,
|
||||
|
||||
@@ -103,9 +103,10 @@ func (h *Handlers) renderLoginError(
|
||||
// The credential check runs BEFORE any rate-limit budget is
|
||||
// consulted, and only a failed check spends budget. That is what
|
||||
// keeps the single administrative path reachable: behind the reverse
|
||||
// proxy this deployment requires, with TRUSTED_PROXIES unset, every
|
||||
// client shares one bucket, so a limiter spent on arrival lets any
|
||||
// stranger deny the operator's own correct password indefinitely.
|
||||
// proxy this deployment requires, when TRUSTED_PROXIES does not cover
|
||||
// it, every client shares one bucket, so a limiter spent on arrival
|
||||
// lets any stranger deny the operator's own correct password
|
||||
// indefinitely.
|
||||
//
|
||||
// Verifying first means every login POST costs an Argon2id hash, so
|
||||
// the work is taken under a bounded number of verification slots.
|
||||
|
||||
@@ -25,7 +25,7 @@ const (
|
||||
|
||||
// sharedProxyPeer is the whole point of this file. Production is
|
||||
// required to run behind a TLS-terminating reverse proxy, and
|
||||
// TRUSTED_PROXIES defaults to empty, so every client — attacker
|
||||
// when TRUSTED_PROXIES does not cover it every client — attacker
|
||||
// and operator alike — reaches the process from the proxy's
|
||||
// address and shares one rate-limit bucket. Both parties in
|
||||
// these tests therefore use the same RemoteAddr.
|
||||
@@ -115,11 +115,11 @@ func floodFailures(
|
||||
// done-criterion of https://git.eeqj.de/sneak/webhooker/issues/150.
|
||||
//
|
||||
// The attacker and the operator share one rate-limit bucket, because
|
||||
// behind the mandated reverse proxy with TRUSTED_PROXIES unset every
|
||||
// client keys on the proxy's address. The attacker floods the
|
||||
// operator's own username — a single-admin product has a predictable
|
||||
// one — far past the failure limit. The operator must still be able
|
||||
// to log in with the correct password.
|
||||
// behind the mandated reverse proxy, when TRUSTED_PROXIES does not
|
||||
// cover it, every client keys on the proxy's address. The attacker
|
||||
// floods the operator's own username — a single-admin product has a
|
||||
// predictable one — far past the failure limit. The operator must
|
||||
// still be able to log in with the correct password.
|
||||
//
|
||||
// This fails if credentials stop being verified ahead of the limiter.
|
||||
func TestLogin_StrangersFloodCannotLockOutTheOperator(t *testing.T) {
|
||||
|
||||
@@ -299,8 +299,7 @@ func countInFlightDeliveries(
|
||||
return count, err
|
||||
}
|
||||
|
||||
// createReplayDelivery writes the new pending delivery row, adds it to
|
||||
// the webhook's running totals in the same transaction, and returns
|
||||
// createReplayDelivery writes the new pending delivery row and returns
|
||||
// the task that carries it to the delivery engine.
|
||||
//
|
||||
// The row is written with associations omitted, and neither Event nor
|
||||
@@ -320,14 +319,7 @@ func createReplayDelivery(
|
||||
Status: database.DeliveryStatusPending,
|
||||
}
|
||||
|
||||
err := webhookDB.Transaction(func(tx *gorm.DB) error {
|
||||
err := tx.Omit(clause.Associations).Create(dlv).Error
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return database.AddTotals(tx, database.Totals{Deliveries: 1})
|
||||
})
|
||||
err := webhookDB.Omit(clause.Associations).Create(dlv).Error
|
||||
if err != nil {
|
||||
return delivery.Task{}, err
|
||||
}
|
||||
|
||||
@@ -69,21 +69,6 @@ func (s *Handlers) LoadEventLogViewsForTest(
|
||||
return views
|
||||
}
|
||||
|
||||
// WebhookStatsForTest returns the figures the statistics pane on a
|
||||
// webhook's page shows, from the webhook's entrypoints and targets
|
||||
// loaded as that page loads them.
|
||||
func (s *Handlers) WebhookStatsForTest(webhookID string) *WebhookStats {
|
||||
var entrypoints []database.Entrypoint
|
||||
|
||||
s.db.DB().Where("webhook_id = ?", webhookID).Find(&entrypoints)
|
||||
|
||||
var targets []database.Target
|
||||
|
||||
s.db.DB().Where("webhook_id = ?", webhookID).Find(&targets)
|
||||
|
||||
return s.loadWebhookStats(webhookID, entrypoints, targets)
|
||||
}
|
||||
|
||||
// AddTemplateForTest registers a template under a page name so that
|
||||
// the handlers_test package can drive the render path with a
|
||||
// template of its own.
|
||||
|
||||
@@ -91,22 +91,18 @@ type Handlers struct {
|
||||
|
||||
// parsePageTemplate parses a page-specific template set from the
|
||||
// embedded FS. Each page template is combined with the shared
|
||||
// base, htmlheader, and navbar templates, and with any further files
|
||||
// the page includes. The page file must be listed first so that its
|
||||
// root action ({{template "base" .}}) becomes the template set's entry
|
||||
// point.
|
||||
func parsePageTemplate(
|
||||
pageFile string, included ...string,
|
||||
) *template.Template {
|
||||
files := append([]string{
|
||||
pageFile,
|
||||
"base.html",
|
||||
"htmlheader.html",
|
||||
"navbar.html",
|
||||
}, included...)
|
||||
|
||||
// base, htmlheader, and navbar templates. The page file must be
|
||||
// listed first so that its root action ({{template "base" .}})
|
||||
// becomes the template set's entry point.
|
||||
func parsePageTemplate(pageFile string) *template.Template {
|
||||
return template.Must(
|
||||
template.ParseFS(templates.Templates, files...),
|
||||
template.ParseFS(
|
||||
templates.Templates,
|
||||
pageFile,
|
||||
"base.html",
|
||||
"htmlheader.html",
|
||||
"navbar.html",
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -135,7 +131,7 @@ func New(
|
||||
"profile.html": parsePageTemplate("profile.html"),
|
||||
"sources_list.html": parsePageTemplate("sources_list.html"),
|
||||
"sources_new.html": parsePageTemplate("sources_new.html"),
|
||||
"source_detail.html": parsePageTemplate("source_detail.html", "webhook_stats.html"),
|
||||
"source_detail.html": parsePageTemplate("source_detail.html"),
|
||||
"source_edit.html": parsePageTemplate("source_edit.html"),
|
||||
"source_logs.html": parsePageTemplate("source_logs.html"),
|
||||
"target_edit.html": parsePageTemplate("target_edit.html"),
|
||||
|
||||
@@ -450,7 +450,6 @@ func (h *Handlers) renderSourceDetail(
|
||||
"Targets": delivery.NewTargetViews(targets),
|
||||
"Events": events,
|
||||
"BaseURL": baseURL,
|
||||
"Stats": h.loadWebhookStats(webhook.ID, entrypoints, targets),
|
||||
}
|
||||
|
||||
h.renderTemplate(w, r, "source_detail.html", data)
|
||||
|
||||
@@ -252,12 +252,11 @@ func requestEventSource(
|
||||
}
|
||||
}
|
||||
|
||||
// createAndFanOut writes the event and one pending delivery per target,
|
||||
// and adds them to the webhook's running totals, in a single
|
||||
// transaction, then hands the tasks to the delivery engine. It is the
|
||||
// only path by which an event and its deliveries are created, so a
|
||||
// resubmitted event is retried, SSRF-guarded and circuit-broken
|
||||
// exactly as a received one is.
|
||||
// createAndFanOut writes the event and one pending delivery per target
|
||||
// in a single transaction, then hands the tasks to the delivery
|
||||
// engine. It is the only path by which an event and its deliveries are
|
||||
// created, so a resubmitted event is retried, SSRF-guarded and
|
||||
// circuit-broken exactly as a received one is.
|
||||
//
|
||||
// The tasks are returned as well as queued, so a caller can report how
|
||||
// many targets the event went to.
|
||||
@@ -297,16 +296,6 @@ func (h *Handlers) createAndFanOut(
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
err = database.AddTotals(tx, database.Totals{
|
||||
Events: 1,
|
||||
Deliveries: int64(len(tasks)),
|
||||
})
|
||||
if err != nil {
|
||||
tx.Rollback()
|
||||
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
err = tx.Commit().Error
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf(
|
||||
|
||||
@@ -1,212 +0,0 @@
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"gorm.io/gorm"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
)
|
||||
|
||||
// The spans of the two recent windows the statistics pane reports on:
|
||||
// the last 10 minutes and the last 24 hours.
|
||||
const (
|
||||
shortWindow = 10 * time.Minute
|
||||
longWindow = 24 * time.Hour
|
||||
)
|
||||
|
||||
// percent turns a fraction into a percentage.
|
||||
const percent = 100
|
||||
|
||||
// WebhookStats holds the figures in the statistics pane at the top of
|
||||
// the webhook page.
|
||||
type WebhookStats struct {
|
||||
Entrypoints int
|
||||
ActiveEntrypoints int
|
||||
Targets int
|
||||
ActiveTargets int
|
||||
|
||||
// Totals holds the lifetime counts of events, deliveries and
|
||||
// failures, and how many of each retention has removed.
|
||||
Totals database.Totals
|
||||
|
||||
// InProgress counts the deliveries still pending or retrying.
|
||||
InProgress int64
|
||||
|
||||
// LastEventAt is when the newest stored event arrived, or nil when
|
||||
// none is stored.
|
||||
LastEventAt *time.Time
|
||||
|
||||
Last10Minutes RecentWindow
|
||||
Last24Hours RecentWindow
|
||||
}
|
||||
|
||||
// RecentWindow holds what happened in one recent window: the events
|
||||
// received in it, and the deliveries that became delivered or failed in
|
||||
// it.
|
||||
type RecentWindow struct {
|
||||
Events int64
|
||||
Delivered int64
|
||||
Failed int64
|
||||
}
|
||||
|
||||
// FailurePercent is the share of the deliveries finished in the window
|
||||
// that failed, or a dash when none finished. Deliveries still pending
|
||||
// or retrying are not counted either way.
|
||||
func (w RecentWindow) FailurePercent() string {
|
||||
finished := w.Delivered + w.Failed
|
||||
if finished == 0 {
|
||||
return "—"
|
||||
}
|
||||
|
||||
return fmt.Sprintf(
|
||||
"%.1f%%", percent*float64(w.Failed)/float64(finished),
|
||||
)
|
||||
}
|
||||
|
||||
// loadWebhookStats gathers the figures for the statistics pane from the
|
||||
// webhook's entrypoints and targets, as the page has already loaded
|
||||
// them, and from its event database. It returns nil, and logs why, when
|
||||
// the event database cannot be read.
|
||||
func (h *Handlers) loadWebhookStats(
|
||||
webhookID string,
|
||||
entrypoints []database.Entrypoint,
|
||||
targets []database.Target,
|
||||
) *WebhookStats {
|
||||
stats := &WebhookStats{
|
||||
Entrypoints: len(entrypoints),
|
||||
Targets: len(targets),
|
||||
}
|
||||
|
||||
for i := range entrypoints {
|
||||
if entrypoints[i].Active {
|
||||
stats.ActiveEntrypoints++
|
||||
}
|
||||
}
|
||||
|
||||
for i := range targets {
|
||||
if targets[i].Active {
|
||||
stats.ActiveTargets++
|
||||
}
|
||||
}
|
||||
|
||||
// Opening an event database that does not exist would create it,
|
||||
// and it would hold nothing to count.
|
||||
if !h.dbMgr.DBExists(webhookID) {
|
||||
return stats
|
||||
}
|
||||
|
||||
webhookDB, err := h.dbMgr.GetDB(webhookID)
|
||||
if err == nil {
|
||||
err = readEventStats(webhookDB, time.Now(), stats)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
h.log.Error(
|
||||
"failed to read webhook statistics",
|
||||
"webhook_id", webhookID,
|
||||
"error", err,
|
||||
)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
return stats
|
||||
}
|
||||
|
||||
// readEventStats fills in the figures that come from the webhook's
|
||||
// event database. None of them reads every stored row: the totals are
|
||||
// one row, and every other figure is read from an index, over only the
|
||||
// rows it counts.
|
||||
func readEventStats(
|
||||
db *gorm.DB, now time.Time, stats *WebhookStats,
|
||||
) error {
|
||||
err := db.Take(&stats.Totals).Error
|
||||
if err != nil {
|
||||
return fmt.Errorf("reading running totals: %w", err)
|
||||
}
|
||||
|
||||
err = db.Model(&database.Delivery{}).
|
||||
Where("status IN ?", []database.DeliveryStatus{
|
||||
database.DeliveryStatusPending,
|
||||
database.DeliveryStatusRetrying,
|
||||
}).
|
||||
Count(&stats.InProgress).Error
|
||||
if err != nil {
|
||||
return fmt.Errorf("counting deliveries in progress: %w", err)
|
||||
}
|
||||
|
||||
var newest []time.Time
|
||||
|
||||
err = db.Model(&database.Event{}).
|
||||
Order("created_at DESC").
|
||||
Limit(1).
|
||||
Pluck("created_at", &newest).Error
|
||||
if err != nil {
|
||||
return fmt.Errorf("reading newest event time: %w", err)
|
||||
}
|
||||
|
||||
if len(newest) > 0 {
|
||||
stats.LastEventAt = &newest[0]
|
||||
}
|
||||
|
||||
stats.Last10Minutes, err = readRecentWindow(
|
||||
db, now.Add(-shortWindow),
|
||||
)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
stats.Last24Hours, err = readRecentWindow(
|
||||
db, now.Add(-longWindow),
|
||||
)
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
// readRecentWindow counts the events received, and the deliveries that
|
||||
// became delivered or failed, since the given time.
|
||||
func readRecentWindow(
|
||||
db *gorm.DB, since time.Time,
|
||||
) (RecentWindow, error) {
|
||||
var w RecentWindow
|
||||
|
||||
err := db.Model(&database.Event{}).
|
||||
Where("created_at >= ?", since).
|
||||
Count(&w.Events).Error
|
||||
if err != nil {
|
||||
return w, fmt.Errorf("counting recent events: %w", err)
|
||||
}
|
||||
|
||||
w.Delivered, err = countFinishedSince(
|
||||
db, database.DeliveryStatusDelivered, since,
|
||||
)
|
||||
if err != nil {
|
||||
return w, err
|
||||
}
|
||||
|
||||
w.Failed, err = countFinishedSince(
|
||||
db, database.DeliveryStatusFailed, since,
|
||||
)
|
||||
|
||||
return w, err
|
||||
}
|
||||
|
||||
// countFinishedSince counts the deliveries that reached the given
|
||||
// final status since the given time.
|
||||
func countFinishedSince(
|
||||
db *gorm.DB, status database.DeliveryStatus, since time.Time,
|
||||
) (int64, error) {
|
||||
var n int64
|
||||
|
||||
err := db.Model(&database.Delivery{}).
|
||||
Where("status = ? AND finished_at >= ?", status, since).
|
||||
Count(&n).Error
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf(
|
||||
"counting deliveries %s recently: %w", status, err,
|
||||
)
|
||||
}
|
||||
|
||||
return n, nil
|
||||
}
|
||||
@@ -1,328 +0,0 @@
|
||||
package handlers_test
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"go.uber.org/fx/fxtest"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
"sneak.berlin/go/webhooker/internal/config"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
"sneak.berlin/go/webhooker/internal/handlers"
|
||||
"sneak.berlin/go/webhooker/internal/logger"
|
||||
"sneak.berlin/go/webhooker/internal/session"
|
||||
)
|
||||
|
||||
// statsEntrypoint adds an entrypoint to a webhook and returns its path.
|
||||
func statsEntrypoint(
|
||||
t *testing.T, db *database.Database, webhookID string, active bool,
|
||||
) string {
|
||||
t.Helper()
|
||||
|
||||
ep := &database.Entrypoint{
|
||||
WebhookID: webhookID,
|
||||
Path: uuid.New().String(),
|
||||
}
|
||||
|
||||
require.NoError(t, db.DB().Omit(clause.Associations).Create(ep).Error)
|
||||
require.NoError(t, db.DB().Model(ep).Update("active", active).Error)
|
||||
|
||||
return ep.Path
|
||||
}
|
||||
|
||||
// statsDelivery returns the id of an event's delivery to a target.
|
||||
func statsDelivery(
|
||||
t *testing.T, webhookDB *gorm.DB, eventID, targetID string,
|
||||
) string {
|
||||
t.Helper()
|
||||
|
||||
var d database.Delivery
|
||||
|
||||
require.NoError(t, webhookDB.Where(
|
||||
"event_id = ? AND target_id = ?", eventID, targetID,
|
||||
).First(&d).Error)
|
||||
|
||||
return d.ID
|
||||
}
|
||||
|
||||
// statsFinish settles a delivery as the delivery engine does: its
|
||||
// final status and the time it finished, and for a failure one more on
|
||||
// the webhook's failure total, in one transaction.
|
||||
func statsFinish(
|
||||
t *testing.T,
|
||||
webhookDB *gorm.DB,
|
||||
deliveryID string,
|
||||
status database.DeliveryStatus,
|
||||
at time.Time,
|
||||
) {
|
||||
t.Helper()
|
||||
|
||||
require.NoError(t, webhookDB.Transaction(func(tx *gorm.DB) error {
|
||||
err := tx.Model(&database.Delivery{}).
|
||||
Where("id = ?", deliveryID).
|
||||
Updates(map[string]any{"status": status, "finished_at": at}).
|
||||
Error
|
||||
if err != nil || status != database.DeliveryStatusFailed {
|
||||
return err
|
||||
}
|
||||
|
||||
return database.AddTotals(tx, database.Totals{Failures: 1})
|
||||
}))
|
||||
}
|
||||
|
||||
// statsAge moves an event's arrival back to the given time.
|
||||
func statsAge(
|
||||
t *testing.T, webhookDB *gorm.DB, eventID string, at time.Time,
|
||||
) {
|
||||
t.Helper()
|
||||
|
||||
require.NoError(t, webhookDB.Model(&database.Event{}).
|
||||
Where("id = ?", eventID).
|
||||
Update("created_at", at).Error)
|
||||
}
|
||||
|
||||
// seedStatsHistory builds the webhook the statistics test checks: one
|
||||
// day of retention, two entrypoints (one inactive) and three targets
|
||||
// (one inactive). Three events arrive through the receiver, and so
|
||||
// each has a delivery to the two active targets. The oldest event is
|
||||
// past retention, the middle one six hours old, the newest just in.
|
||||
// Their deliveries are settled as the delivery engine would, and a
|
||||
// replay adds a pending delivery to the oldest event. It returns the
|
||||
// webhook, its event database and the newest event.
|
||||
func seedStatsHistory(
|
||||
t *testing.T,
|
||||
h *handlers.Handlers,
|
||||
sess *session.Session,
|
||||
db *database.Database,
|
||||
dbMgr *database.WebhookDBManager,
|
||||
) (*database.Webhook, *gorm.DB, database.Event) {
|
||||
t.Helper()
|
||||
|
||||
wh := &database.Webhook{
|
||||
UserID: deleteTestUserID, Name: "stats", RetentionDays: 1,
|
||||
}
|
||||
require.NoError(t, db.DB().Omit(clause.Associations).Create(wh).Error)
|
||||
|
||||
path := statsEntrypoint(t, db, wh.ID, true)
|
||||
statsEntrypoint(t, db, wh.ID, false)
|
||||
|
||||
first := seedConfiguredTarget(
|
||||
t, db, wh.ID, database.TargetTypeHTTP,
|
||||
`{"url":"`+replayTargetURL+`"}`,
|
||||
)
|
||||
second := seedTarget(t, db, wh.ID, database.TargetTypeLog)
|
||||
inactive := seedTarget(t, db, wh.ID, database.TargetTypeLog)
|
||||
require.NoError(t, db.DB().Model(inactive).
|
||||
Update("active", false).Error)
|
||||
|
||||
router := receiverRouter(h)
|
||||
|
||||
for range 3 {
|
||||
require.Equal(t, http.StatusOK, postReceiver(t, router, path))
|
||||
}
|
||||
|
||||
webhookDB, err := dbMgr.GetDB(wh.ID)
|
||||
require.NoError(t, err)
|
||||
|
||||
events := listEvents(t, webhookDB)
|
||||
require.Len(t, events, 3)
|
||||
|
||||
oldest, middle, newest := events[0], events[1], events[2]
|
||||
now := time.Now()
|
||||
|
||||
statsAge(t, webhookDB, oldest.ID, now.Add(-50*time.Hour))
|
||||
statsAge(t, webhookDB, middle.ID, now.Add(-6*time.Hour))
|
||||
|
||||
oldestFailure := statsDelivery(t, webhookDB, oldest.ID, first.ID)
|
||||
statsFinish(t, webhookDB, oldestFailure,
|
||||
database.DeliveryStatusFailed, now.Add(-49*time.Hour))
|
||||
statsFinish(t, webhookDB,
|
||||
statsDelivery(t, webhookDB, oldest.ID, second.ID),
|
||||
database.DeliveryStatusDelivered, now.Add(-49*time.Hour))
|
||||
statsFinish(t, webhookDB,
|
||||
statsDelivery(t, webhookDB, middle.ID, first.ID),
|
||||
database.DeliveryStatusFailed, now.Add(-5*time.Hour))
|
||||
statsFinish(t, webhookDB,
|
||||
statsDelivery(t, webhookDB, middle.ID, second.ID),
|
||||
database.DeliveryStatusFailed, now.Add(-time.Minute))
|
||||
statsFinish(t, webhookDB,
|
||||
statsDelivery(t, webhookDB, newest.ID, first.ID),
|
||||
database.DeliveryStatusDelivered, now.Add(-2*time.Minute))
|
||||
|
||||
require.Equal(t, http.StatusSeeOther,
|
||||
postReplay(t, h, sess, wh.ID, oldestFailure).Code)
|
||||
|
||||
return wh, webhookDB, newest
|
||||
}
|
||||
|
||||
// statsPrune runs the real retention reaper until it has removed one
|
||||
// event from the webhook's database, then stops it.
|
||||
func statsPrune(
|
||||
t *testing.T,
|
||||
db *database.Database,
|
||||
dbMgr *database.WebhookDBManager,
|
||||
log *logger.Logger,
|
||||
webhookDB *gorm.DB,
|
||||
) {
|
||||
t.Helper()
|
||||
|
||||
lc := fxtest.NewLifecycle(t)
|
||||
database.NewRetentionReaper(lc, database.RetentionReaperParams{
|
||||
Config: &config.Config{
|
||||
RetentionSweepInterval: 10 * time.Millisecond,
|
||||
},
|
||||
Database: db,
|
||||
DBManager: dbMgr,
|
||||
Logger: log,
|
||||
})
|
||||
|
||||
lc.RequireStart()
|
||||
|
||||
require.Eventually(t, func() bool {
|
||||
var totals database.Totals
|
||||
|
||||
err := webhookDB.Take(&totals).Error
|
||||
|
||||
return err == nil && totals.EventsRemoved == 1
|
||||
}, 10*time.Second, 10*time.Millisecond)
|
||||
|
||||
lc.RequireStop()
|
||||
}
|
||||
|
||||
// assertStatsTotals checks the lifetime events, deliveries and
|
||||
// failures, and those within retention.
|
||||
func assertStatsTotals(
|
||||
t *testing.T, totals database.Totals, lifetime, within [3]int64,
|
||||
) {
|
||||
t.Helper()
|
||||
|
||||
gotLifetime := [3]int64{
|
||||
totals.Events, totals.Deliveries, totals.Failures,
|
||||
}
|
||||
gotWithin := [3]int64{
|
||||
totals.EventsWithinRetention(),
|
||||
totals.DeliveriesWithinRetention(),
|
||||
totals.FailuresWithinRetention(),
|
||||
}
|
||||
|
||||
assert.Equal(t, lifetime, gotLifetime,
|
||||
"lifetime events, deliveries, failures")
|
||||
assert.Equal(t, within, gotWithin,
|
||||
"events, deliveries, failures within retention")
|
||||
}
|
||||
|
||||
// TestWebhookStats_EveryFigureAcrossRetentionPrune checks every figure
|
||||
// the statistics pane shows for the history seedStatsHistory builds,
|
||||
// before and after the real retention reaper removes the oldest event.
|
||||
func TestWebhookStats_EveryFigureAcrossRetentionPrune(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
var (
|
||||
h *handlers.Handlers
|
||||
sess *session.Session
|
||||
db *database.Database
|
||||
dbMgr *database.WebhookDBManager
|
||||
log *logger.Logger
|
||||
)
|
||||
|
||||
app := newTestApp(t, &h, &sess, &db, &dbMgr, &log)
|
||||
app.RequireStart()
|
||||
|
||||
t.Cleanup(app.RequireStop)
|
||||
|
||||
wh, webhookDB, newest := seedStatsHistory(t, h, sess, db, dbMgr)
|
||||
|
||||
stats := h.WebhookStatsForTest(wh.ID)
|
||||
require.NotNil(t, stats)
|
||||
|
||||
assert.Equal(t, 2, stats.Entrypoints)
|
||||
assert.Equal(t, 1, stats.ActiveEntrypoints)
|
||||
assert.Equal(t, 3, stats.Targets)
|
||||
assert.Equal(t, 2, stats.ActiveTargets)
|
||||
assertStatsTotals(t, stats.Totals, [3]int64{3, 7, 3}, [3]int64{3, 7, 3})
|
||||
assert.Equal(t, int64(2), stats.InProgress)
|
||||
require.NotNil(t, stats.LastEventAt)
|
||||
assert.True(t, newest.CreatedAt.Equal(*stats.LastEventAt))
|
||||
assert.Equal(t, handlers.RecentWindow{
|
||||
Events: 1, Delivered: 1, Failed: 1,
|
||||
}, stats.Last10Minutes)
|
||||
assert.Equal(t, handlers.RecentWindow{
|
||||
Events: 2, Delivered: 1, Failed: 2,
|
||||
}, stats.Last24Hours)
|
||||
assert.Equal(t, "50.0%", stats.Last10Minutes.FailurePercent())
|
||||
assert.Equal(t, "66.7%", stats.Last24Hours.FailurePercent())
|
||||
|
||||
// Retention removes the oldest event with its three deliveries,
|
||||
// one of them failed and one the pending replay.
|
||||
statsPrune(t, db, dbMgr, log, webhookDB)
|
||||
|
||||
after := h.WebhookStatsForTest(wh.ID)
|
||||
require.NotNil(t, after)
|
||||
|
||||
assertStatsTotals(t, after.Totals, [3]int64{3, 7, 3}, [3]int64{2, 4, 2})
|
||||
assert.Equal(t, int64(1), after.InProgress)
|
||||
assert.Equal(t, stats.LastEventAt, after.LastEventAt)
|
||||
assert.Equal(t, stats.Last10Minutes, after.Last10Minutes)
|
||||
assert.Equal(t, stats.Last24Hours, after.Last24Hours)
|
||||
|
||||
body := renderSourceDetailPage(t, h, sess, wh.ID)
|
||||
assert.Contains(t, body, "Statistics")
|
||||
assert.Contains(t, body, "Within retention")
|
||||
assert.Contains(t, body, "50.0%")
|
||||
assert.Contains(t, body, "66.7%")
|
||||
}
|
||||
|
||||
// TestWebhookStats_WebhookWithNoEvents covers a webhook whose event
|
||||
// database has never been opened: every count is zero, the
|
||||
// percentages are a dash, and showing the page does not create the
|
||||
// database.
|
||||
func TestWebhookStats_WebhookWithNoEvents(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
var (
|
||||
h *handlers.Handlers
|
||||
sess *session.Session
|
||||
db *database.Database
|
||||
dbMgr *database.WebhookDBManager
|
||||
)
|
||||
|
||||
app := newTestApp(t, &h, &sess, &db, &dbMgr)
|
||||
app.RequireStart()
|
||||
|
||||
t.Cleanup(app.RequireStop)
|
||||
|
||||
wh := seedWebhook(t, db)
|
||||
|
||||
assert.Equal(t, &handlers.WebhookStats{}, h.WebhookStatsForTest(wh.ID))
|
||||
assert.Equal(t, "—", handlers.RecentWindow{}.FailurePercent())
|
||||
|
||||
body := renderSourceDetailPage(t, h, sess, wh.ID)
|
||||
assert.Contains(t, body, "Statistics")
|
||||
assert.False(t, dbMgr.DBExists(wh.ID))
|
||||
}
|
||||
|
||||
// TestRecentWindow_FailurePercent pins the percentage: failed
|
||||
// deliveries out of all that finished in the window.
|
||||
func TestRecentWindow_FailurePercent(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
tests := []struct {
|
||||
window handlers.RecentWindow
|
||||
want string
|
||||
}{
|
||||
{handlers.RecentWindow{}, "—"},
|
||||
{handlers.RecentWindow{Events: 4}, "—"},
|
||||
{handlers.RecentWindow{Delivered: 3, Failed: 1}, "25.0%"},
|
||||
{handlers.RecentWindow{Failed: 2}, "100.0%"},
|
||||
{handlers.RecentWindow{Delivered: 2}, "0.0%"},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
assert.Equal(t, tt.want, tt.window.FailurePercent(), tt.window)
|
||||
}
|
||||
}
|
||||
@@ -108,10 +108,10 @@ type failureWindow struct {
|
||||
//
|
||||
// A limiter that spends budget on arrival cannot protect a
|
||||
// single-admin product: behind the reverse proxy the deployment
|
||||
// requires, with TRUSTED_PROXIES unset, every client keys on the
|
||||
// proxy, so a stranger trickling five POSTs a minute keeps the one
|
||||
// bucket full and the operator's own correct password is answered 429
|
||||
// forever. There is no second administrative path.
|
||||
// requires, when TRUSTED_PROXIES does not cover it, every client
|
||||
// keys on the proxy, so a stranger trickling five POSTs a minute
|
||||
// keeps the one bucket full and the operator's own correct password
|
||||
// is answered 429 forever. There is no second administrative path.
|
||||
//
|
||||
// So budget is spent only by a FAILED verification. A correct
|
||||
// password is never throttled, whatever the counters say, which is
|
||||
|
||||
@@ -123,9 +123,8 @@ func bucketKey(addr netip.Addr) string {
|
||||
return prefix.String()
|
||||
}
|
||||
|
||||
// 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.
|
||||
// isTrustedProxy reports whether addr belongs to a network in
|
||||
// TRUSTED_PROXIES, which by default is the RFC 1918 private ranges.
|
||||
func (m *Middleware) isTrustedProxy(addr netip.Addr) bool {
|
||||
for _, prefix := range m.params.Config.TrustedProxies {
|
||||
if prefix.Contains(addr) {
|
||||
|
||||
@@ -384,8 +384,8 @@ const (
|
||||
// trustedProxyCIDR is the proxy network the forwarded-path
|
||||
// tests configure, and trustedPeer an address inside it. A
|
||||
// production deployment is required to run behind a reverse
|
||||
// proxy with TRUSTED_PROXIES set, so this is the shape the
|
||||
// bucketing has to hold in.
|
||||
// proxy that TRUSTED_PROXIES covers, either by the default or by
|
||||
// a set value, so this is the shape the bucketing has to hold in.
|
||||
trustedProxyCIDR = "10.0.0.0/8"
|
||||
trustedPeer = "10.0.0.1:44444"
|
||||
)
|
||||
@@ -426,8 +426,8 @@ func assertSharedBucket(
|
||||
}
|
||||
|
||||
// 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
|
||||
// this gating exists for: from a peer that is not a trusted
|
||||
// proxy, 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.
|
||||
@@ -1097,8 +1097,9 @@ func TestPostRateLimit_IPv4IndependentPerAddress(t *testing.T) {
|
||||
// that arrives from trustedPeer — a configured trusted proxy — and
|
||||
// names forwarded as its client in X-Forwarded-For. That is the
|
||||
// production path: a deployment is required to run behind a reverse
|
||||
// proxy with TRUSTED_PROXIES set, so the forwarded address, not the
|
||||
// peer, is what the limiters bucket on there.
|
||||
// proxy that TRUSTED_PROXIES covers, either by the default or by a
|
||||
// set value, so the forwarded address, not the peer, is what the
|
||||
// limiters bucket on there.
|
||||
func forwardedKeyFor(
|
||||
t *testing.T, m *middleware.Middleware, forwarded string,
|
||||
) string {
|
||||
@@ -1178,9 +1179,9 @@ func TestRateLimitKey_ForwardedIPv6BucketsByPrefix(t *testing.T) {
|
||||
//
|
||||
// Every existing test of this fallback uses an IPv4 proxy, where
|
||||
// bucketKey is the identity function, so replacing the call with
|
||||
// peer.String() leaves the whole suite green. Only operator-listed
|
||||
// addresses reach this line and the fallback is fail-closed, so this
|
||||
// pins behaviour rather than fixing a defect.
|
||||
// peer.String() leaves the whole suite green. Only addresses inside
|
||||
// TRUSTED_PROXIES reach this line and the fallback is fail-closed, so
|
||||
// this pins behaviour rather than fixing a defect.
|
||||
func TestRateLimitKey_TrustedPeerUnusableForwardedMasksPeer(
|
||||
t *testing.T,
|
||||
) {
|
||||
|
||||
@@ -154,11 +154,12 @@ func (s *Server) setupPageRoutes() {
|
||||
r.Use(s.mw.NoCache())
|
||||
|
||||
// The login POST carries no pre-emptive rate limiter. Behind
|
||||
// the reverse proxy production requires, with TRUSTED_PROXIES
|
||||
// unset, every client shares one bucket, so a limiter spent
|
||||
// on arrival lets any stranger deny the operator the only
|
||||
// administrative path. The handler verifies credentials first
|
||||
// and charges only failures; see Handlers.authenticateUser.
|
||||
// the reverse proxy production requires, when TRUSTED_PROXIES
|
||||
// does not cover it, every client shares one bucket, so a
|
||||
// limiter spent on arrival lets any stranger deny the operator
|
||||
// the only administrative path. The handler verifies
|
||||
// credentials first and charges only failures; see
|
||||
// Handlers.authenticateUser.
|
||||
r.Get("/login", s.h.HandleLoginPage())
|
||||
r.Post("/login", s.h.HandleLoginSubmit())
|
||||
|
||||
|
||||
@@ -24,8 +24,6 @@
|
||||
</div>
|
||||
</div>
|
||||
|
||||
{{template "webhook_stats" .}}
|
||||
|
||||
<div class="grid grid-cols-1 lg:grid-cols-2 gap-6">
|
||||
<!-- Entrypoints -->
|
||||
<div class="card">
|
||||
|
||||
@@ -1,87 +0,0 @@
|
||||
{{define "webhook_stats"}}
|
||||
<!-- Statistics pane at the top of the webhook page. -->
|
||||
<div class="card mb-6">
|
||||
<div class="p-4 border-b border-gray-200">
|
||||
<h2 class="text-lg font-medium text-gray-900">Statistics</h2>
|
||||
</div>
|
||||
{{with .Stats}}
|
||||
<div class="p-4 flex flex-wrap gap-6 text-sm border-b border-gray-200">
|
||||
<div>
|
||||
<span class="text-gray-500">Entrypoints</span>
|
||||
<span class="font-medium text-gray-900">{{.Entrypoints}}</span>
|
||||
<span class="text-gray-500">({{.ActiveEntrypoints}} active)</span>
|
||||
</div>
|
||||
<div>
|
||||
<span class="text-gray-500">Targets</span>
|
||||
<span class="font-medium text-gray-900">{{.Targets}}</span>
|
||||
<span class="text-gray-500">({{.ActiveTargets}} active)</span>
|
||||
</div>
|
||||
<div>
|
||||
<span class="text-gray-500">Deliveries in progress</span>
|
||||
<span class="font-medium text-gray-900">{{.InProgress}}</span>
|
||||
</div>
|
||||
<div>
|
||||
<span class="text-gray-500">Last event</span>
|
||||
<span class="font-medium text-gray-900">{{with .LastEventAt}}{{.Format "2006-01-02 15:04:05 UTC"}}{{else}}none{{end}}</span>
|
||||
</div>
|
||||
<div>
|
||||
<span class="text-gray-500">Retention</span>
|
||||
<span class="font-medium text-gray-900">{{$.Webhook.RetentionLabel}}</span>
|
||||
</div>
|
||||
</div>
|
||||
<div class="p-4 grid grid-cols-1 lg:grid-cols-2 gap-6 text-sm">
|
||||
<div>
|
||||
<div class="flex py-2 border-b border-gray-200 text-xs text-gray-500 uppercase tracking-wide">
|
||||
<span class="flex-1"></span>
|
||||
<span class="w-32 text-center">Lifetime</span>
|
||||
<span class="w-32 text-center">Within retention</span>
|
||||
</div>
|
||||
<div class="divide-y divide-gray-100">
|
||||
<div class="flex py-2">
|
||||
<span class="flex-1 text-gray-600">Events</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Totals.Events}}</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Totals.EventsWithinRetention}}</span>
|
||||
</div>
|
||||
<div class="flex py-2">
|
||||
<span class="flex-1 text-gray-600">Deliveries</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Totals.Deliveries}}</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Totals.DeliveriesWithinRetention}}</span>
|
||||
</div>
|
||||
<div class="flex py-2">
|
||||
<span class="flex-1 text-gray-600">Failures</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Totals.Failures}}</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Totals.FailuresWithinRetention}}</span>
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
<div>
|
||||
<div class="flex py-2 border-b border-gray-200 text-xs text-gray-500 uppercase tracking-wide">
|
||||
<span class="flex-1"></span>
|
||||
<span class="w-32 text-center">Last 10 minutes</span>
|
||||
<span class="w-32 text-center">Last 24 hours</span>
|
||||
</div>
|
||||
<div class="divide-y divide-gray-100">
|
||||
<div class="flex py-2">
|
||||
<span class="flex-1 text-gray-600">Events</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Last10Minutes.Events}}</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Last24Hours.Events}}</span>
|
||||
</div>
|
||||
<div class="flex py-2">
|
||||
<span class="flex-1 text-gray-600">Failures</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Last10Minutes.Failed}}</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Last24Hours.Failed}}</span>
|
||||
</div>
|
||||
<div class="flex py-2">
|
||||
<span class="flex-1 text-gray-600">Failure percentage</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Last10Minutes.FailurePercent}}</span>
|
||||
<span class="w-32 text-center text-gray-900">{{.Last24Hours.FailurePercent}}</span>
|
||||
</div>
|
||||
</div>
|
||||
<p class="mt-2 text-xs text-gray-500">Failure percentage is the failed deliveries out of all deliveries that finished in the window. Deliveries still pending or retrying are not counted.</p>
|
||||
</div>
|
||||
</div>
|
||||
{{else}}
|
||||
<div class="p-4 text-sm text-gray-500">The statistics could not be read.</div>
|
||||
{{end}}
|
||||
</div>
|
||||
{{end}}
|
||||
Reference in New Issue
Block a user