Author SHA1 Message Date
clawbot 3a4f3625e8 Document and test that unrouted methods get 405 without Allow on /s/*
check / check (push) Waiting to run
A method chi does not route, such as PROPFIND, is refused by the
top-level router before it reaches the /s group, so it gets 405
without an Allow header. The README row and the test's doc comment
now say so, and TestStaticServesOnlyGetAndHead checks PROPFIND.

Model: opus-5-5
2026-09-29 08:34:01 +00:00
clawbot 205539cde7 Restrict /s/* to GET and HEAD (closes #169)
The static file server was attached with Mount, which registers every
method, so POST, PUT and DELETE on an asset were answered 200 with the
file. It is now registered for GET and HEAD only, inside a /s group
whose method-not-allowed handler answers 405 with Allow: GET, HEAD
(chi's default 405 sends no Allow header).

TestStaticServesEveryMethod is inverted and renamed
TestStaticServesOnlyGetAndHead, and the README route table row for
/s/* now says the same.

Model: opus-5-5
2026-09-29 08:34:01 +00:00
clawbot f0adeafde3 Drop Set-Cookie from the recovered 500 (closes #193)
check / check (push) Successful in 3m35s
When a handler sets a cookie and then panics before sending
anything, the recover middleware now deletes Set-Cookie before
writing its 500, so a request that failed never hands the client
a credential. Every other header, Location included, is left as
http.Error leaves it, matching chi's Recoverer. A response that
was already sent is untouched.

Tests cover the uncommitted case (no cookie, Location kept) and
assert the cookie still reaches the client when the response was
committed before the panic.

Model: opus-5-5
2026-09-29 10:30:26 +02:00
clawbot f755c03110 Default-block Azure WireServer's public address (closes #245)
check / check (push) Successful in 4m34s
Add 168.63.129.16 (Azure WireServer) to blockedNetworks, the default
blocklist, not alwaysBlockedNetworks: it is public unicast, so an
operator who lists it in ALLOWED_EGRESS_CIDRS can reach it again. The
refusal message, the allowlist startup warning, the README and the
comments no longer call every blocked address private/reserved, and
no longer claim the allowlist cannot open any metadata endpoint.

Sources:
- https://learn.microsoft.com/en-us/azure/virtual-network/what-is-ip-address-168-63-129-16
- https://learn.microsoft.com/en-us/azure/virtual-machines/metadata-security-protocol/overview

Deviation: 147.75.207.243 (Equinix Metal) is not added; Equinix
documents only a hostname, and the service was sunset on 2026-06-30.

Model: opus-5-5
2026-09-29 10:22:07 +02:00
clawbot 4a724130ca Close archive writers when the delivery engine stops (closes #280)
check / check (push) Successful in 3m45s
The engine cached archive writers and never closed them at shutdown,
so after a clean stop an archive's rows could sit in its -wal while
the .db held no table. The engine's stop hook now evicts every cached
writer once its workers have returned, the same way deleting a webhook
does, so a clean stop leaves each archive as one file and a late write
is refused. If the workers do not return within the stop budget, the
writers are left open as a kill would leave them: closing would wait
on a write in progress, and a still-running worker would open new
ones.

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

Model: opus-5-5
2026-09-29 08:30:22 +02:00
clawbot d4f4ddf51f Send Content-Type once on a delivery (closes #246)
check / check (push) Successful in 3m20s
A delivery set Content-Type from the event's ContentType and then
added the inbound Content-Type from the event's stored headers, so a
target could receive two values. The inbound Content-Type is no
longer forwarded from the stored headers; the receiver already saves
it as the event's ContentType.

Which value is sent is now stated at applyRequestHeaders: a
Content-Type configured on the target, otherwise the event's
ContentType, otherwise none. A configured one still survives a
cross-origin 307/308 with its body.

Model: opus-5-5
2026-09-29 07:11:55 +02:00
clawbot 51580a2bc6 Fail a pending delivery whose target was deleted (closes #293)
check / check (push) Successful in 3m38s
Restart recovery and the pending sweep skipped a pending delivery
whose target was missing from the batch's target map, every minute,
for the life of the database. A miss now asks loadTarget: no row
fails the delivery terminally with a recorded reason; any other error
leaves it pending, since the map is also empty when its query failed;
a target found there is used.

The failure goes through the ownership-gated function the retrying
paths already used, now failMissingTarget. Once it owns the delivery
it re-reads the row and fails it only if the status is unchanged, so
a delivery sent and settled in between is left alone.

Model: opus-5-5
2026-09-29 06:48:19 +02:00
clawbot e0b211f960 Make deliveries refused while half-open wait a cooldown (closes #306)
check / check (push) Successful in 5m0s
While the breaker was half-open, Allow refused every delivery but the
probe and CooldownRemaining returned zero, so each queued task for the
target went straight back onto the retry channel and rewrote its status
on every pass until the probe finished.

CooldownRemaining now returns the whole cooldown while half-open, so a
refused delivery waits that long. A refused delivery already at
retrying is not written again, so the retry counter now moves only
when a refusal moves a delivery into retrying.

Model: opus-5-5
2026-09-29 05:48:20 +02:00
clawbot 978eb01b29 Open each event database once when callers race (closes #291)
check / check (push) Successful in 3m48s
GetDB opened the database on a cache miss and then tried to cache it,
so callers racing on a webhook's first use could each open the file,
and the losers closed their copies. On a new file the parallel opens
also create its tables at the same time, and one caller can fail with
"table already exists".

A mutex now covers the open: GetDB looks in the cache again under it,
then opens and caches. DeleteDB and CloseAll take the same mutex, so
neither runs while an open is under way. Reading an already cached
database takes no lock.

The new test starts many callers on one webhook at once and checks
that exactly one open happened.

Model: opus-5-5
2026-09-29 04:37:30 +02:00
clawbot 3cdab97930 Say make check needs make bootstrap on a fresh clone (closes #282)
check / check (push) Successful in 11s
The third-party browser assets are not committed, so on a fresh clone
make check fails in the tests until make bootstrap (or make assets) has
fetched them. The Entrypoints section now says so up front, and why the
check does not fetch them itself: it must not change files in the repo.

Model: opus-5-5
2026-09-29 04:30:14 +02:00
clawbot 83740b1de1 Do not resend a delivery that restart recovery already sent (closes #299)
check / check (push) Successful in 3m35s
Restart recovery could find a just-written delivery pending, send it
and release it before the receiver's Notify queued the same delivery.
Notify's claim then succeeded on the released id, and the worker sent
it again because the new-task path never read the delivery's row.

Before sending a new task the worker now reads the delivery's status
by primary key and skips the task unless the row still says pending,
as the retry path already does for retrying. Nothing else can change
the row while the worker owns the delivery. A row left pending by a
failed bookkeeping write is still sent again.

loadRetryDelivery is renamed loadDelivery now that both paths use it.

Model: opus-5-5
2026-09-29 04:12:14 +02:00
sneak aeeeca5ea1 Merge branch 'main' into next
check / check (push) Successful in 9s
2026-09-29 03:13:57 +02:00
clawbot 6ebac4fa71 Index the event-tier columns the sweeps, event log and retention scan (closes #314)
check / check (push) Successful in 3m4s
The per-webhook event databases had no secondary indexes, so startup
recovery, the retry and pending sweeps, the queue-depth sampler, the
event log and retention each read whole tables. Indexes declared in
the GORM model tags now serve them, and AutoMigrate adds them to new
and existing databases alike.

Each index also covers deleted_at: GORM adds deleted_at IS NULL to
these queries, and SQLite, with no table statistics, otherwise
prefers the existing deleted_at index. A test checks SQLite's plan
for each statement as GORM builds it.

Rule suppressed: lll on the three event-tier model structs, whose
struct tags cannot wrap.
The resubmitted_from_id scan is left to
#325.

Model: opus-4-8 (implementation); opus-5-5 (rework)
2026-09-28 14:13:22 +02:00
clawbot 237f131367 Default WEBHOOKER_ENVIRONMENT to prod (closes #307)
check / check (push) Successful in 3m14s
An unset WEBHOOKER_ENVIRONMENT now means prod, not dev. The only
thing dev still changes is CORS, which then answers every origin with
Access-Control-Allow-Origin: *, so an operator who forgets the
variable is no longer silently permissive; dev must be set
explicitly. Cookie Secure and CSRF strictness follow each request's
transport and are unaffected.

The README, comments and tests no longer describe dev as the default:
the deployment checklist asks only that the environment is not dev,
the Docker and nginx examples drop the now-redundant setting, and the
TRUSTED_PROXIES warning gives its real reason for firing in every
environment.

Model: opus-4-8 (implementation); opus-5-5 (rework)
2026-09-28 12:47:31 +02:00
clawbot 7ed1588443 Document running webhooker under upaas (closes #323)
check / check (push) Successful in 8s
Adds a "Running under upaas" section to the README: add no port
mapping, since upaas publishes mapped ports on every host interface,
and put the app on the reverse proxy's Docker network instead; one
data volume at /var/lib/webhooker, created owned by UID 1000 before
the first deploy; WEBHOOKER_ENVIRONMENT and TRUSTED_PROXIES; the
health check upaas reads 60 seconds after a deploy; and where the
first-run admin password appears and how to reset it.

upaas bind-mounts a host directory it never creates, and one made by
root stops the container at its data directory lock. The documented
creation step removes that; the image is unchanged.

Model: opus-5-5
2026-09-28 12:30:32 +02:00
clawbot b051821370 Say max_retries is the total attempt count, not a retry count (closes #316)
check / check (push) Successful in 3m54s
The help text under the field on both target forms and the max_retries rows in the README now say the number is the total number of delivery attempts: 0 is a single attempt with no retries and no circuit breaker, and N is N attempts in total. The delivery code already worked this way; only the wording was wrong, so an operator wanting one try plus two retries would have entered 2 instead of 3. A UI copy test renders both forms and pins the wording. Delivery behaviour is unchanged.

Model: opus-4-8 (implementation); fable-5-1 (merge)
2026-09-21 18:33:09 +02:00
clawbot 39afa69bfc Consolidate the data directory mode into one owner (closes #288)
check / check (push) Successful in 6m34s
Two packages each declared the 0o750 mode for DATA_DIR and both created the directory. internal/datadir now exports DirPerm as the single definition, and internal/database uses it in both places it creates the directory. The value is unchanged, so existing deployments see no permission change. datadir owns it because guarding and creating DATA_DIR is that package's whole purpose and it imports nothing that would form a cycle.

Model: opus-4-8 (implementation and review); fable-5-1 (merge)
2026-09-21 10:01:51 +02:00
clawbotandsneak 888eaf526b State the UUID-is-the-credential rule as a rule (closes #301) (#302)
check / check (push) Successful in 4m28s
Closes #301. Docs-only apart from one test comment; no behaviour change.

The receiver has authenticated on the entrypoint UUID alone since inbound signature verification was removed in #279. The README described that as the current state. It did not say it is the decision, which leaves a future contributor free to propose HMAC as an improvement rather than as a reversal.

What changed:

- `## The entrypoint URL is the authentication secret` now states the rule: the v4 UUID at `/webhook/{uuid}` is the credential and the only one; no shared secret, HMAC signature, bearer token or second factor will be added, including as defence in depth. It names the removal that settled it, and it says explicitly that signature headers a sender sends anyway are stored and forwarded but never checked — the previous text left that ambiguous.
- The same section carries the two consequences an operator has to act on: the URL is a capability, so keep it out of logs, tickets and screenshots; and rotation means minting a new entrypoint, not changing a key.
- It also handles the case the rule will next be argued from: a sender that only supports signed payloads to a well-known URL is a constraint on that integration, to be raised on its own terms, not grounds to reintroduce shared secrets.
- The rule is reachable without scrolling 1,100 lines: a pointer in the intro, a new first bullet under Authentication (which previously listed the web UI, the API and `/metrics` and said nothing about the receiver at all), and a sharpened bullet under Security.

Stale language found and corrected: one, in `internal/delivery/redirect_test.go`. Its comment justified same-origin header retention partly by "the inbound signature the receiver verifies" — in this repo's vocabulary "the receiver" is `/webhook/{uuid}`, which verifies nothing. The endpoint that verifies it is the delivery target's, and the comment now says so.

Two places that read like stale signing language were checked and left alone as accurate: `internal/delivery/redirect.go` and `internal/server/sentry.go` describe signature headers senders put on the receiver route, which do arrive and are forwarded — neither claims webhooker checks them.

`REPO_POLICIES.md` was deliberately not touched. It is the cross-project policy document synced from `sneak/prompts` and carries `last_modified` front matter for that purpose, so a webhooker-specific carve-out does not belong in it. Worth knowing: its hardening section ends "if a standard security hardening measure exists for HTTP services and is not listed here, it is still expected. When in doubt, harden" — that is the sentence a future HMAC proposal will cite, and only the README now answers it.

`TODO.md` is untouched per its own Workflow section (issue branches do not touch it).

Co-authored-by: sneak <sneak@sneak.berlin>
Reviewed-on: #302
Co-authored-by: clawbot <clawbot@noreply.example.org>
Co-committed-by: clawbot <clawbot@noreply.example.org>
2026-08-30 04:05:38 +02:00
38 changed files with 1631 additions and 241 deletions
+197 -71
View File
@@ -7,6 +7,13 @@ services, durably stores them, and delivers them to configured targets
with retry support, logging, and observability. Category: infrastructure with retry support, logging, and observability. Category: infrastructure
/ web service. License: MIT. / web service. License: MIT.
Each entrypoint is a version 4 UUID served at `/webhook/{uuid}`, and
that UUID is the entrypoint's only credential. webhooker does not use
shared secrets, HMAC signatures or token headers on the receiver, and
will not add them — read
[The entrypoint URL is the authentication secret](#the-entrypoint-url-is-the-authentication-secret)
before deploying one.
## Getting Started ## Getting Started
### Prerequisites ### Prerequisites
@@ -37,9 +44,9 @@ make bootstrap
# Run all checks (test, lint, format check) # Run all checks (test, lint, format check)
make check make check
# Run in development mode. DATA_DIR defaults to /var/lib/webhooker in # Run the server from the clone. DATA_DIR defaults to
# every environment, so set it (in .env or the shell) to a writable # /var/lib/webhooker in every environment, so set it (in .env or the
# directory when running from a clone. # shell) to a writable directory.
DATA_DIR=./data make dev DATA_DIR=./data make dev
# Build Docker image # Build Docker image
@@ -85,7 +92,8 @@ them at once. A variable already present in the real environment wins
over the file's value for the same name. over the file's value for the same name.
The environment is selected by setting `WEBHOOKER_ENVIRONMENT` to `dev` The environment is selected by setting `WEBHOOKER_ENVIRONMENT` to `dev`
or `prod` (default: `dev`). The setting controls exactly one behavior: or `prod` (default: `prod`; `dev` must be set explicitly). The setting
controls exactly one behavior:
| Behavior | `dev` | `prod` | | Behavior | `dev` | `prod` |
| -------- | ----------------------- | ---------------- | | -------- | ----------------------- | ---------------- |
@@ -127,7 +135,7 @@ TTY detection, and security headers are always applied.
| Variable | Description | Default | | Variable | Description | Default |
| ----------------------- | ----------------------------------- | -------- | | ----------------------- | ----------------------------------- | -------- |
| `WEBHOOKER_ENVIRONMENT` | `dev` or `prod` | `dev` | | `WEBHOOKER_ENVIRONMENT` | `dev` or `prod` | `prod` |
| `PORT` | HTTP listen port | `8080` | | `PORT` | HTTP listen port | `8080` |
| `BIND_ADDRESS` | IP address the HTTP listener binds. Loopback by default, so the cleartext listener is not published on every interface. The Docker image ships `0.0.0.0` instead. See [Bind address](#bind-address) | `127.0.0.1` (image: `0.0.0.0`) | | `BIND_ADDRESS` | IP address the HTTP listener binds. Loopback by default, so the cleartext listener is not published on every interface. The Docker image ships `0.0.0.0` instead. See [Bind address](#bind-address) | `127.0.0.1` (image: `0.0.0.0`) |
| `DATA_DIR` | Directory for all SQLite databases | `/var/lib/webhooker` | | `DATA_DIR` | Directory for all SQLite databases | `/var/lib/webhooker` |
@@ -149,6 +157,11 @@ private and reserved ranges — RFC 1918, loopback, CGNAT, link-local and
the rest — are refused, which stops a target from being used to make the rest — are refused, which stops a target from being used to make
webhooker probe the network it sits in. webhooker probe the network it sits in.
Besides the private and reserved ranges, the default blocklist refuses
public cloud metadata addresses: currently only `168.63.129.16`, Azure's
WireServer, which serves an Azure VM its credentials. Because it is a
public address, listing it in `ALLOWED_EGRESS_CIDRS` reopens it.
That default is also inconvenient for the thing webhooker is mostly That default is also inconvenient for the thing webhooker is mostly
for: taking a public webhook and forwarding it to something on your own for: taking a public webhook and forwarding it to something on your own
network. A container on the same Docker network, a box on `10.x`, a network. A container on the same Docker network, a box on `10.x`, a
@@ -187,15 +200,16 @@ Two things this setting cannot do:
the list is always an allowlist; an empty list (the default) means the list is always an allowlist; an empty list (the default) means
every private and reserved range stays refused. Note that every private and reserved range stays refused. Note that
`0.0.0.0/0` gets you most of the way there anyway, per above. `0.0.0.0/0` gets you most of the way there anyway, per above.
- **It cannot open link-local, or a cloud metadata endpoint that - **It cannot open link-local, or a cloud metadata endpoint at a
discloses credentials or user data.** An address is on the list below non-public address that discloses credentials or user data.** An
when both of these hold: the provider fixes it, so it cannot collide address is on the list below when it is not a public address and both
with anything you run; and reaching it hands out credentials, user of these hold: the provider fixes it, so it cannot collide with
data or bootstrap material. Those stay blocked no matter what you anything you run; and reaching it hands out credentials, user data or
list, including when you list them outright or list a supernet such bootstrap material. Those stay blocked no matter what you list,
as `0.0.0.0/0`, `::/0`, `fd00::/8` or `100.64.0.0/10`. Treat this as including when you list them outright or list a supernet such as
best effort rather than a guarantee — it is a hand-maintained list `0.0.0.0/0`, `::/0`, `fd00::/8` or `100.64.0.0/10`. Treat this as best
and the caveat below the table applies: effort rather than a guarantee — it is a hand-maintained list and the
caveat below the table applies:
| Blocked unconditionally | What it is | | Blocked unconditionally | What it is |
| ----------------------- | ---------- | | ----------------------- | ---------- |
@@ -234,7 +248,8 @@ Two things this setting cannot do:
encodings, which the default blocklist does not match. A publicly encodings, which the default blocklist does not match. A publicly
routable metadata address is not listed here, because nothing on this routable metadata address is not listed here, because nothing on this
list can be reopened and blocking one that way would leave you no list can be reopened and blocking one that way would leave you no
escape hatch at all. escape hatch at all; Azure's `168.63.129.16` is refused by the default
blocklist instead, as described above.
This list is not exhaustive of every cloud's metadata address — if This list is not exhaustive of every cloud's metadata address — if
yours is not here, do not allowlist the block that contains it. yours is not here, do not allowlist the block that contains it.
@@ -392,10 +407,9 @@ the bucket is. See [Rate Limiting](#rate-limiting).
The remedy is to set `TRUSTED_PROXIES` to your reverse proxy's The remedy is to set `TRUSTED_PROXIES` to your reverse proxy's
address, which restores per-client buckets. webhooker logs a warning address, which restores per-client buckets. webhooker logs a warning
at startup whenever `TRUSTED_PROXIES` is empty, in every environment — at startup whenever `TRUSTED_PROXIES` is empty, in every environment,
not only when `WEBHOOKER_ENVIRONMENT=prod`, because that variable because behind a proxy every client shares one bucket in `dev` and
defaults to `dev` and an operator who never set it is precisely the `prod` alike. The warning is informational when nothing proxies to the
one at risk. The warning is informational when nothing proxies to the
process: with no proxy in front, the peer address is the client's own process: with no proxy in front, the peer address is the client's own
and the buckets are already per-client. See and the buckets are already per-client. See
[Rate Limiting](#rate-limiting) for what each limit shares. [Rate Limiting](#rate-limiting) for what each limit shares.
@@ -631,7 +645,6 @@ decision:
docker run -d \ docker run -d \
-p 127.0.0.1:8080:8080 \ -p 127.0.0.1:8080:8080 \
-v /path/to/data:/var/lib/webhooker \ -v /path/to/data:/var/lib/webhooker \
-e WEBHOOKER_ENVIRONMENT=prod \
-e BIND_ADDRESS=0.0.0.0 \ -e BIND_ADDRESS=0.0.0.0 \
webhooker:latest webhooker:latest
``` ```
@@ -724,6 +737,66 @@ listing the directory and learning your webhook UUIDs from the
`events-{uuid}.db` filenames — not the barrier protecting the `events-{uuid}.db` filenames — not the barrier protecting the
credentials. credentials.
### Running under upaas
[upaas](https://git.eeqj.de/sneak/upaas) builds the image from this
repository's `Dockerfile` and runs it. The app needs:
- **Network and port:** add no port mapping in upaas. upaas publishes
every mapped port on all interfaces of the host
([upaas issue 113](https://git.eeqj.de/sneak/upaas/issues/113)),
which would put the plain-HTTP admin UI and receiver there. Instead,
set the app's Docker Network in upaas to your reverse proxy's Docker
network; the proxy then reaches the app at `upaas-` followed by the
app name, port `8080`. Leave `PORT` unset: the image's health check
probes `8080`.
- **Volume:** one host directory mounted at `/var/lib/webhooker`.
upaas bind-mounts the host path it is given and does not create it,
and the container does not start unless UID 1000 owns it (see
[Running with Docker](#running-with-docker)). Create it before the
first deploy:
```bash
mkdir -p /path/to/data
chown 1000:1000 /path/to/data
chmod 750 /path/to/data
```
- **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).
- 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`.
- Everything else is optional; see [Configuration](#configuration).
- **Health check:** the image's own, which requests
`/.well-known/healthcheck`. upaas reads the container's health 60
seconds after a deploy and marks the deploy failed unless it is
`healthy`.
- **First run:** the first start prints the `admin` password once, in
the banner described under [The admin account](#the-admin-account),
to the container's log. upaas names the container `upaas-` followed
by the app name, so for an app named `webhooker`:
```bash
docker logs upaas-webhooker
```
If the password is lost, stop the container, set a new password with
the app's own image and volume, and start it again (see
[Recovering a lost admin password](#recovering-a-lost-admin-password)):
```bash
docker stop upaas-webhooker
docker run --rm --volumes-from upaas-webhooker \
"$(docker inspect -f '{{.Image}}' upaas-webhooker)" \
/app/webhooker resetpw -generate admin
docker start upaas-webhooker
```
## Deployment behind a reverse proxy ## Deployment behind a reverse proxy
webhooker terminates no TLS of its own. It serves plaintext HTTP and webhooker terminates no TLS of its own. It serves plaintext HTTP and
@@ -745,17 +818,18 @@ reports.
serves the admin login form and the unauthenticated receiver with serves the admin login form and the unauthenticated receiver with
no TLS at all, and the proxy in front of it changes nothing about no TLS at all, and the proxy in front of it changes nothing about
that. that.
2. **Set `WEBHOOKER_ENVIRONMENT=prod`, and make sure the proxy sends 2. **Make sure the environment is not `dev` (leave
`X-Forwarded-Proto`.** These are two requirements, not one. The `WEBHOOKER_ENVIRONMENT` unset or set it to `prod`), and make sure
environment setting decides CORS and nothing else: the default the proxy sends `X-Forwarded-Proto`.** These are two requirements,
not one. The environment setting decides CORS and nothing else:
`dev` answers every origin with `Access-Control-Allow-Origin: *` `dev` answers every origin with `Access-Control-Allow-Origin: *`
(without credentials), which a server-rendered production (without credentials), which a server-rendered production
deployment has no use for. Cookie `Secure` and the strict deployment has no use for, and `prod` — the default — disables it.
Origin/Referer mode are **not** tied to it — they are decided per Cookie `Secure` and the strict Origin/Referer mode are **not** tied
request from the transport, which behind a proxy means the to it — they are decided per request from the transport, which
`X-Forwarded-Proto` header. The block below sets it; without it behind a proxy means the `X-Forwarded-Proto` header. The block below
every request is read as plaintext and cookies ship without sets it; without it every request is read as plaintext and cookies
`Secure`. See [Configuration](#configuration). ship without `Secure`. See [Configuration](#configuration).
3. **Set `TRUSTED_PROXIES` to the proxy's address.** Unset, every rate 3. **Set `TRUSTED_PROXIES` to the proxy's address.** Unset, every rate
limiter keys on the connecting peer, which behind a proxy is the limiter keys on the connecting peer, which behind a proxy is the
proxy on every request: all clients collapse into one global bucket proxy on every request: all clients collapse into one global bucket
@@ -847,7 +921,6 @@ sent — `$scheme` above does.
With that block, webhooker's environment is: With that block, webhooker's environment is:
```sh ```sh
WEBHOOKER_ENVIRONMENT=prod
BIND_ADDRESS=127.0.0.1 # the default; stated here to be explicit BIND_ADDRESS=127.0.0.1 # the default; stated here to be explicit
TRUSTED_PROXIES=127.0.0.1 TRUSTED_PROXIES=127.0.0.1
``` ```
@@ -902,15 +975,10 @@ scratch file**: it holds committed transactions that are not yet in the
have no readable schema at all. `-shm` is regenerable, but there is no have no readable schema at all. `-shm` is regenerable, but there is no
reason to separate the two — copy the directory and you have them. reason to separate the two — copy the directory and you have them.
A clean shutdown closes `webhooker.db` and every `events-*.db`, which A clean shutdown closes every database, which checkpoints and removes
checkpoints and removes their sidecars; a killed or crashed instance its sidecars; a killed or crashed instance leaves them, and they must be
leaves them, and they must be carried with the `.db`. **Archive carried with the `.db`. An archive the service has not opened since a
databases are different**: their handle is not closed at shutdown, so crash keeps that crash's sidecars, even across a later clean stop.
`archive-*.db-wal` and `-shm` normally survive a clean stop and the
`-wal` can hold every row the archive has. Measured on a stopped
instance: `archive-….db` 4096 bytes with no table, its `-wal` 157 KB
holding all 8 archived events. Copying `DATA_DIR` in full is what makes
this a non-issue; copying `.db` files out of it by name is not.
Configuration is **not** in `DATA_DIR` — it comes from the environment Configuration is **not** in `DATA_DIR` — it comes from the environment
and from a `.env` file read out of the process working directory. Back and from a `.env` file read out of the process working directory. Back
@@ -985,10 +1053,9 @@ The file becomes self-contained again when the handle closes, which
happens on the next write past the debounce window, when the connection happens on the next write past the debounce window, when the connection
pool retires the idle connection (about a minute after the last write), pool retires the idle connection (about a minute after the last write),
or at the idle archive sweep — measured, the same file was a complete or at the idle archive sweep — measured, the same file was a complete
20 KB `.db` with no sidecars about a minute after its last write. 20 KB `.db` with no sidecars about a minute after its last write. A
Shutdown is **not** on that list: the archive handle is not closed when clean stop closes it too. So either move `archive-{uuid}.db` together
the service stops. So either move `archive-{uuid}.db` together with any with any `-wal`/`-shm` beside it, or wait until there are none.
`-wal`/`-shm` beside it, or wait until there are none.
### Restore ### Restore
@@ -1007,12 +1074,10 @@ the service stops. So either move `archive-{uuid}.db` together with any
They are part of the database, and dropping a `-wal` silently They are part of the database, and dropping a `-wal` silently
discards every transaction it still holds. An `.backup` set will not discards every transaction it still holds. An `.backup` set will not
contain any: it writes a single consolidated file per database. A contain any: it writes a single consolidated file per database. A
stop-and-copy set has none for `webhooker.db` or the `events-*.db`, stop-and-copy set normally has none, because a clean stop closes
because a clean stop closes those and checkpoints their sidecars every database and checkpoints its sidecars away; the exception is an
away — but it will normally have them for `archive-*.db`, whose archive not opened since a crash. A copy salvaged from a crashed
handle stays open across shutdown, and those carry the archive's instance has them for everything, and needs all of them.
rows. A copy salvaged from a crashed instance has them for
everything, and needs all of them.
4. **Fix ownership.** The container runs as the non-root `webhooker` 4. **Fix ownership.** The container runs as the non-root `webhooker`
user, UID 1000 / GID 1000. Restored files must be owned by (or user, UID 1000 / GID 1000. Restored files must be owned by (or
@@ -1149,14 +1214,38 @@ backups at rest and restrict who can read them.
## The entrypoint URL is the authentication secret ## The entrypoint URL is the authentication secret
The receiver verifies nothing about an inbound request. The UUID in an **The entrypoint UUID is the credential, and it is the only one.**
entrypoint's URL is its credential: anyone who holds that URL can webhooker mints a version 4 UUID per entrypoint and serves it at
submit events to it, and the receiver checks nothing else about the `/webhook/{uuid}`. Possession of that URL is the authentication:
sender. Treat an entrypoint URL the way you would treat an API token. anyone who holds it can submit events to the entrypoint, and the
receiver verifies nothing else about the sender.
There is no way to rotate the UUID in place. To retire one, delete the There is no shared secret, no HMAC signature, no bearer token and no
entrypoint (or deactivate it, which answers `410`) and create a new second factor on the receiver, and none will be added. This was
one, then point the sender at the new URL. considered and rejected; the implementation that existed was removed
in [PR #279](https://git.eeqj.de/sneak/webhooker/pulls/279), closing
[issue #67](https://git.eeqj.de/sneak/webhooker/issues/67) and
[issue #241](https://git.eeqj.de/sneak/webhooker/issues/241). A
proposal to reintroduce any of them — including as "defence in depth"
alongside the UUID — is answered by this section. Inbound signature
headers a sender sends anyway (`X-Hub-Signature` and its
per-provider equivalents) are stored and forwarded as ordinary
headers; nothing checks them.
What that means for an operator:
- **The URL is a capability, so treat it as a secret.** Keep it out of
logs, ticket bodies, chat messages and screenshots. Anyone who reads
it anywhere can post events as that sender.
- **Rotating means minting a new entrypoint, not changing a key.**
There is no way to rotate the UUID in place. To retire one, delete
the entrypoint (or deactivate it, which answers `410`) and create a
new one, then point the sender at the new URL.
- **A sender that cannot be given a secret URL is a constraint on that
integration, not a reason to change this.** If a service only
supports signed payloads to a well-known URL, raise it as its own
problem — pick a different integration path, or accept that it
cannot be used. It is not grounds to reintroduce shared secrets.
## Entrypoints ## Entrypoints
@@ -1166,8 +1255,15 @@ standard: normalized scripts in `script/` are the entrypoints for the
development workflow. Ten of the Makefile's seventeen targets are thin development workflow. Ten of the Makefile's seventeen targets are thin
shims that call them; `build`, `run`, `dev`, `deps`, `clean`, `css` and shims that call them; `build`, `run`, `dev`, `deps`, `clean`, `css` and
`version` are inline commands with no script behind them, though `version` are inline commands with no script behind them, though
`build` and `version` both take their value from `script/version`. We `build` and `version` both take their value from `script/version`.
provide:
`make check` needs the third-party browser assets in `static/`, which
are not committed, so run `make bootstrap` (or just `make assets`) once
after cloning. Without them the tests fail with a message naming that
remedy. `make check` does not fetch them itself because it must not
change any files in the repo.
We provide:
- `script/bootstrap` — install all dependencies (idempotent) - `script/bootstrap` — install all dependencies (idempotent)
- `script/setup` — make a fresh clone ready for development - `script/setup` — make a fresh clone ready for development
@@ -1476,7 +1572,7 @@ events should be forwarded.
| `type` | TargetType | One of: `http`, `slack`, `database`, `log` | | `type` | TargetType | One of: `http`, `slack`, `database`, `log` |
| `active` | boolean | Whether deliveries are enabled (default: true) | | `active` | boolean | Whether deliveries are enabled (default: true) |
| `config` | JSON text | Type-specific configuration | | `config` | JSON text | Type-specific configuration |
| `max_retries` | integer | Maximum retry attempts for `http` and `slack` targets (0 = fire-and-forget, >0 = retries with backoff and a circuit breaker). Ignored by `database` and `log` targets | | `max_retries` | integer | Total delivery attempts for `http` and `slack` targets, not retries on top of the first: 0 is a single fire-and-forget attempt with no retries and no circuit breaker, and a value of N makes N attempts in all, with exponential backoff and a per-target circuit breaker. Ignored by `database` and `log` targets |
| `max_queue_size` | integer | Stored and shown on the target's detail view, but not enforced anywhere yet: nothing in the delivery engine consults it. Queue depth is set by the two fixed 10,000-entry channels | | `max_queue_size` | integer | Stored and shown on the target's detail view, but not enforced anywhere yet: nothing in the delivery engine consults it. Queue depth is set by the two fixed 10,000-entry channels |
**Relations:** Belongs to Webhook. Has many Deliveries. **Relations:** Belongs to Webhook. Has many Deliveries.
@@ -1484,12 +1580,12 @@ events should be forwarded.
**Target types:** **Target types:**
- **`http`** — Forward the event as an HTTP POST to a configured URL. - **`http`** — Forward the event as an HTTP POST to a configured URL.
Behavior depends on `max_retries`: when `max_retries` is 0 (the `max_retries` is the total number of delivery attempts, not retries on
default), the target operates in fire-and-forget mode — a single top of the first: when `max_retries` is 0 (the default), the target
attempt with no retries and no circuit breaker. When `max_retries` is operates in fire-and-forget mode, a single attempt with no retries and
greater than 0, failed deliveries are retried with exponential backoff no circuit breaker; a value of N makes up to N attempts in all,
up to `max_retries` attempts, protected by a per-target circuit retrying failed deliveries with exponential backoff and protecting them
breaker. with a per-target circuit breaker.
- **`slack`** — Post the event as a formatted message to a - **`slack`** — Post the event as a formatted message to a
Slack-compatible incoming webhook URL (`webhookUrl` in `config`). It Slack-compatible incoming webhook URL (`webhookUrl` in `config`). It
is built on the same HTTP core as `http` and honours `max_retries` is built on the same HTTP core as `http` and honours `max_retries`
@@ -1669,6 +1765,28 @@ retries) is individually logged for full observability.
**Relations:** Belongs to Delivery. **Relations:** Belongs to Delivery.
#### Event-tier indexes
These indexes on the per-webhook event databases are declared in the model
tags, so `AutoMigrate` creates them on a fresh and on an existing database:
| Table | Columns | Serves |
| ------------------ | --------------------------- | ------ |
| `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 |
| `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
deletes leave it out, but their lookups of expired rows keep it. SQLite keeps
no statistics on these tables, and without them it rates the `deleted_at`
index, which every live row matches, above an index on a column matched
against several values or compared with `<`. So every index but the last also
covers `deleted_at`. It comes second, so that retention's deletes can use the
index without it, except in `events`, where `created_at` is compared with `<`
and SQLite narrows by a `<` only on the last column it uses.
#### Common Fields #### Common Fields
Every entity except `Setting` includes these fields from `BaseModel`. Every entity except `Setting` includes these fields from `BaseModel`.
@@ -1932,7 +2050,7 @@ rescans the database anyway).
| ----------- | -------- | | ----------- | -------- |
| **Closed** | Normal operation. Deliveries flow through. Consecutive failures are counted. | | **Closed** | Normal operation. Deliveries flow through. Consecutive failures are counted. |
| **Open** | Target appears down. Deliveries are skipped and rescheduled for after the cooldown. | | **Open** | Target appears down. Deliveries are skipped and rescheduled for after the cooldown. |
| **Half-Open** | Cooldown expired. One probe delivery is allowed to test if the target has recovered. | | **Half-Open** | Cooldown expired. One probe delivery is allowed to test if the target has recovered. Other deliveries are rescheduled for one whole cooldown later. |
**Transitions:** **Transitions:**
@@ -1970,7 +2088,9 @@ operations), and log targets (stdout) do not use circuit breakers.
When a circuit is open and a new delivery arrives, the engine marks the When a circuit is open and a new delivery arrives, the engine marks the
delivery as `retrying` and schedules a retry timer for after the delivery as `retrying` and schedules a retry timer for after the
remaining cooldown period. This ensures no deliveries are lost — they're remaining cooldown period. This ensures no deliveries are lost — they're
just delayed until the target is healthy again. just delayed until the target is healthy again. A delivery already in
`retrying` keeps that status without another database write each time
the breaker turns it away.
### Metrics ### Metrics
@@ -1984,7 +2104,7 @@ arriving and being stored, they are just not getting anywhere.
| Metric | Type | Meaning | | Metric | Type | Meaning |
| ------ | ---- | ------- | | ------ | ---- | ------- |
| `webhooker_events_received_total` | counter | Events received and durably stored. Compare against the delivery counters on one dashboard | | `webhooker_events_received_total` | counter | Events received and durably stored. Compare against the delivery counters on one dashboard |
| `webhooker_delivery_attempts_total` | counter | Delivery attempts actually dispatched to a target. A delivery an open circuit breaker refused is not one: it is counted as a retry instead | | `webhooker_delivery_attempts_total` | counter | Delivery attempts actually dispatched to a target. A delivery a circuit breaker refused is not one: it is counted as a retry instead, but only when the refusal moves it into `retrying` |
| `webhooker_deliveries_succeeded_total` | counter | Deliveries that reached `delivered` | | `webhooker_deliveries_succeeded_total` | counter | Deliveries that reached `delivered` |
| `webhooker_deliveries_failed_total` | counter | Deliveries that failed terminally and will not be retried | | `webhooker_deliveries_failed_total` | counter | Deliveries that failed terminally and will not be retried |
| `webhooker_delivery_retries_total` | counter | Deliveries put back into `retrying` | | `webhooker_delivery_retries_total` | counter | Deliveries put back into `retrying` |
@@ -2601,7 +2721,7 @@ abuse limit later; they are tracked as future work.
| ------ | --------------------------- | ----------- | | ------ | --------------------------- | ----------- |
| `GET` | `/` | Root redirect, 303 (authenticated → `/sources`, unauthenticated → `/pages/login`) | | `GET` | `/` | Root redirect, 303 (authenticated → `/sources`, unauthenticated → `/pages/login`) |
| `GET` | `/.well-known/healthcheck` | Health check (JSON: `status`, `now`, `uptimeSeconds`, `uptimeHuman`, `version`, `appname`, `maintenanceMode`) | | `GET` | `/.well-known/healthcheck` | Health check (JSON: `status`, `now`, `uptimeSeconds`, `uptimeHuman`, `version`, `appname`, `maintenanceMode`) |
| any | `/s/*` | Static file serving (embedded CSS, JS). Mounted for every method, not just `GET`/`HEAD`: chi's `Mount` registers all methods and `http.FileServer` special-cases only `HEAD` (by omitting the body), so a `POST` or `DELETE` to an asset is answered `200` with the file. Pinned by `TestStaticServesEveryMethod` | | `GET`, `HEAD` | `/s/*` | Static file serving (embedded CSS, JS). `GET` and `HEAD` only — `POST`, `PUT`, `PATCH`, `DELETE`, `OPTIONS`, `TRACE` and `CONNECT` are answered `405 Method Not Allowed` with `Allow: GET, HEAD`. Any other method (such as `PROPFIND`) is refused by chi before it reaches this route, and gets `405` without an `Allow` header. Pinned by `TestStaticServesOnlyGetAndHead` |
| `POST` | `/webhook/{uuid}` | Webhook receiver endpoint. `POST` only — every other method is answered `405 Method Not Allowed` with `Allow: POST`. Rate limited (see [Rate Limiting](#rate-limiting)) | | `POST` | `/webhook/{uuid}` | Webhook receiver endpoint. `POST` only — every other method is answered `405 Method Not Allowed` with `Allow: POST`. Rate limited (see [Rate Limiting](#rate-limiting)) |
#### Authentication Endpoints #### Authentication Endpoints
@@ -2867,6 +2987,10 @@ check, see [The login endpoint](#the-login-endpoint).
### Authentication ### Authentication
- **Webhook receiver:** the entrypoint UUID in the URL, and nothing
else. No shared secret, no HMAC signature, no token header, and none
will be added — see
[The entrypoint URL is the authentication secret](#the-entrypoint-url-is-the-authentication-secret).
- **Web UI:** Cookie-based sessions using gorilla/sessions with - **Web UI:** Cookie-based sessions using gorilla/sessions with
encrypted cookies. Sessions are configured with HttpOnly, SameSite encrypted cookies. Sessions are configured with HttpOnly, SameSite
Lax, and Secure whenever the request is on TLS — the flag follows the Lax, and Secure whenever the request is on TLS — the flag follows the
@@ -2906,7 +3030,8 @@ check, see [The login endpoint](#the-login-endpoint).
mode mode
- **The entrypoint URL is the receiver's only credential.** Nothing - **The entrypoint URL is the receiver's only credential.** Nothing
about an inbound request is verified; possession of the UUID about an inbound request is verified; possession of the UUID
authorises submission (see authorises submission, and no shared secret or signature check will
be added alongside it (see
[The entrypoint URL is the authentication secret](#the-entrypoint-url-is-the-authentication-secret)) [The entrypoint URL is the authentication secret](#the-entrypoint-url-is-the-authentication-secret))
- **SSRF prevention** for HTTP delivery targets: private/reserved IP - **SSRF prevention** for HTTP delivery targets: private/reserved IP
ranges (RFC 1918, loopback, link-local, cloud metadata) are blocked ranges (RFC 1918, loopback, link-local, cloud metadata) are blocked
@@ -2962,7 +3087,8 @@ each hook. The order, read off the fx stop-hook log:
3. `server` — the HTTP drain, bounded separately by 3. `server` — the HTTP drain, bounded separately by
`server.ShutdownTimeout` (**3 seconds**), then a Sentry flush if `server.ShutdownTimeout` (**3 seconds**), then a Sentry flush if
`SENTRY_DSN` is set `SENTRY_DSN` is set
4. `delivery.Engine` 4. `delivery.Engine` — waits for its workers, then closes the archive
databases
5. `healthcheck` 5. `healthcheck`
6. `WebhookDBManager` 6. `WebhookDBManager`
7. the database close 7. the database close
+19 -16
View File
@@ -192,9 +192,10 @@ type Config struct {
// alwaysBlockedNetworks stays blocked no matter what is listed // alwaysBlockedNetworks stays blocked no matter what is listed
// here. That set is link-local plus the cloud metadata // here. That set is link-local plus the cloud metadata
// endpoints outside it that disclose credentials or user data // endpoints outside it that disclose credentials or user data
// at a provider-fixed address; it is not exhaustive of every // at a provider-fixed, non-public address; it is not
// cloud's metadata address. See alwaysBlockedNetworks for the // exhaustive of every cloud's metadata address. See
// authoritative list and the criterion it is built from. // alwaysBlockedNetworks for the authoritative list and the
// criterion it is built from.
AllowedEgressCIDRs []netip.Prefix AllowedEgressCIDRs []netip.Prefix
params *ConfigParams params *ConfigParams
@@ -585,12 +586,14 @@ func resolveMetricsAuth() (string, string, error) {
) )
} }
// resolveEnvironment reads WEBHOOKER_ENVIRONMENT, defaulting to // resolveEnvironment reads WEBHOOKER_ENVIRONMENT, defaulting to prod
// dev, and rejects unrecognised values. // when it is unset so a deployment that forgets the variable is not
// silently permissive; dev must be set explicitly. It rejects
// unrecognised values.
func resolveEnvironment() (string, error) { func resolveEnvironment() (string, error) {
environment := os.Getenv("WEBHOOKER_ENVIRONMENT") environment := os.Getenv("WEBHOOKER_ENVIRONMENT")
if environment == "" { if environment == "" {
environment = EnvironmentDev environment = EnvironmentProd
} }
if environment != EnvironmentDev && if environment != EnvironmentDev &&
@@ -744,12 +747,14 @@ func (c *Config) warnEgressAllowlist(log *slog.Logger) {
log.Warn( log.Warn(
"ALLOWED_EGRESS_CIDRS lets delivery targets reach these "+ "ALLOWED_EGRESS_CIDRS lets delivery targets reach these "+
"otherwise-blocked private/reserved networks. Anyone "+ "otherwise-blocked networks. Anyone who can create a "+
"who can create a delivery target can now make this "+ "delivery target can now make this process issue "+
"process issue requests into them, and read back the "+ "requests into them, and read back the response. Only "+
"response. Link-local and the known cloud instance "+ "the addresses the README lists as blocked "+
"metadata endpoints outside it stay blocked "+ "unconditionally stay blocked regardless of what is "+
"regardless of what is listed here.", "listed here; a public cloud metadata address such as "+
"168.63.129.16 is reachable once it, or a block "+
"covering it, is listed.",
"allowedEgressCIDRs", "allowedEgressCIDRs",
strings.Join(PrefixStrings(c.AllowedEgressCIDRs), ","), strings.Join(PrefixStrings(c.AllowedEgressCIDRs), ","),
) )
@@ -772,10 +777,8 @@ func (c *Config) warnEgressAllowlist(log *slog.Logger) {
// everyone else's wrong passwords, and the receiver's limits become // everyone else's wrong passwords, and the receiver's limits become
// service-wide ceilings. // service-wide ceilings.
// //
// The warning is deliberately not gated on WEBHOOKER_ENVIRONMENT. That // The warning is deliberately not gated on WEBHOOKER_ENVIRONMENT:
// variable defaults to dev, so gating on it would silence the warning // behind a proxy every client shares one bucket in dev and prod alike.
// for exactly the operator who forgot to configure the deployment —
// the case it exists to catch.
// //
// The default of trusting nobody is deliberate — trusting forwarded // The default of trusting nobody is deliberate — trusting forwarded
// headers from arbitrary peers lets any client choose its own bucket — // headers from arbitrary peers lets any client choose its own bucket —
+13 -17
View File
@@ -44,9 +44,9 @@ func TestEnvironmentConfig(t *testing.T) {
isProd bool isProd bool
}{ }{
{ {
name: "default is dev", name: "default is prod",
isDev: true, isDev: false,
isProd: false, isProd: true,
}, },
{ {
name: "explicit dev", name: "explicit dev",
@@ -834,12 +834,13 @@ func TestEgressAllowlistWarning(t *testing.T) {
// to be able to read back which networks are open. // to be able to read back which networks are open.
assert.Contains(t, logged, "10.0.0.0/8") assert.Contains(t, logged, "10.0.0.0/8")
assert.Contains(t, logged, "127.0.0.0/8") assert.Contains(t, logged, "127.0.0.0/8")
// What stays shut. Asserted on the clause naming the // What stays shut is the whole unconditional set, not
// wider set rather than on "Link-local" alone, so the // link-local alone; a public metadata address is not in
// string cannot narrow back to link-local only while // it, so a listed block covering it opens it.
// the always-blocked set covers ULA, CGNAT and two assert.Contains(t, logged, "blocked unconditionally")
// public metadata addresses as well. assert.Contains(t, logged, "168.63.129.16 is reachable")
assert.Contains(t, logged, "metadata endpoints outside it") // The listed blocks need not be private or reserved.
assert.NotContains(t, logged, "private/reserved")
}) })
} }
} }
@@ -848,10 +849,9 @@ func TestEgressAllowlistWarning(t *testing.T) {
// tells an operator a deployment behind a reverse proxy shares one // tells an operator a deployment behind a reverse proxy shares one
// rate-limit bucket between every client, which turns the receiver // rate-limit bucket between every client, which turns the receiver
// limits into service-wide ceilings and collapses login failure // limits into service-wide ceilings and collapses login failure
// counting. It must fire whenever TRUSTED_PROXIES is empty, // counting. It must fire whenever TRUSTED_PROXIES is empty, in any
// in any environment: WEBHOOKER_ENVIRONMENT defaults to dev, so gating // environment, because behind a proxy every client shares one bucket
// on it would silence the warning for exactly the operator who never // in dev and prod alike. It stays quiet once proxies are named.
// configured the deployment. It stays quiet once proxies are named.
func TestSharedRateLimitBucketWarning(t *testing.T) { func TestSharedRateLimitBucketWarning(t *testing.T) {
tests := []struct { tests := []struct {
name string name string
@@ -871,10 +871,6 @@ func TestSharedRateLimitBucketWarning(t *testing.T) {
expectWarning: false, expectWarning: false,
}, },
{ {
// The default environment. An internet-exposed
// deployment whose operator never set
// WEBHOOKER_ENVIRONMENT lands here and has exactly
// the exposure the warning announces.
name: "dev without trusted proxies warns", name: "dev without trusted proxies warns",
environment: config.EnvironmentDev, environment: config.EnvironmentDev,
expectWarning: true, expectWarning: true,
+4 -2
View File
@@ -17,12 +17,12 @@ import (
"gorm.io/gorm" "gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/banner" "sneak.berlin/go/webhooker/internal/banner"
"sneak.berlin/go/webhooker/internal/config" "sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/datadir"
"sneak.berlin/go/webhooker/internal/gormlog" "sneak.berlin/go/webhooker/internal/gormlog"
"sneak.berlin/go/webhooker/internal/logger" "sneak.berlin/go/webhooker/internal/logger"
) )
const ( const (
dataDirPerm = 0750
randomPasswordLen = 16 randomPasswordLen = 16
sessionKeyLen = 32 sessionKeyLen = 32
) )
@@ -185,7 +185,9 @@ func (d *Database) connect() error {
// caller's decision. // caller's decision.
func (d *Database) connectTo(dataDir string) error { func (d *Database) connectTo(dataDir string) error {
// Ensure the data directory exists before opening the database. // Ensure the data directory exists before opening the database.
err := os.MkdirAll(dataDir, dataDirPerm) // datadir.DirPerm is the single source of the directory mode; this
// package creates the directory too, since either may run first.
err := os.MkdirAll(dataDir, datadir.DirPerm)
if err != nil { if err != nil {
return fmt.Errorf( return fmt.Errorf(
"creating data directory %s: %w", "creating data directory %s: %w",
@@ -0,0 +1,166 @@
package database_test
import (
"context"
"fmt"
"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"
)
// TestWebhookDBManager_OpenAddsEventTierIndexes verifies that opening a
// per-webhook database that predates these indexes creates them. It
// stands in for an older database file by dropping the indexes
// AutoMigrate just created, then reopening the same file.
func TestWebhookDBManager_OpenAddsEventTierIndexes(t *testing.T) {
t.Parallel()
indexes := []struct {
model any
name string
}{
{&database.Delivery{}, "idx_deliveries_status"},
{&database.Delivery{}, "idx_deliveries_event_id"},
{&database.DeliveryResult{}, "idx_delivery_results_delivery_id"},
{&database.Event{}, "idx_events_deleted_at_created_at"},
{&database.Event{}, "idx_events_created_at"},
}
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)
// A fresh database has them.
for _, ix := range indexes {
require.True(t, db.Migrator().HasIndex(ix.model, ix.name))
}
// Stand in for a database file created before the indexes existed.
for _, ix := range indexes {
require.NoError(t, db.Migrator().DropIndex(ix.model, ix.name))
require.False(t, db.Migrator().HasIndex(ix.model, ix.name))
}
// Drop the cached connection so the next open reopens the file and
// runs AutoMigrate against it, as a restart would.
require.NoError(t, mgr.CloseAll())
db, err = mgr.GetDB(webhookID)
require.NoError(t, err)
for _, ix := range indexes {
assert.True(t, db.Migrator().HasIndex(ix.model, ix.name),
"opening the existing database should create %s", ix.name)
}
}
// TestEventTierQueriesUseTheirIndexes verifies that the statements the
// indexes are for use them. GORM builds each statement in a dry run as
// the code named above it does, soft-delete condition included, and
// SQLite, which keeps no statistics on these tables, must plan to seek
// on each index listed by the columns in parentheses.
func TestEventTierQueriesUseTheirIndexes(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})
ids := []string{
uuid.New().String(), uuid.New().String(), uuid.New().String(),
}
cutoff := time.Now()
var (
deliveries []database.Delivery
results []database.DeliveryResult
depths []struct{ Depth int }
)
byStatus := "idx_deliveries_status (status=? AND deleted_at=?)"
byEvent := "idx_deliveries_event_id (event_id=? AND deleted_at=?)"
byAge := "idx_events_deleted_at_created_at (deleted_at=? AND created_at<?)"
// The delivery engine: recovery and the retry sweep, the sweep for
// stranded pending deliveries, and the queue depth count.
assertPlanUses(t, db, dry.Where(
"status = ?", database.DeliveryStatusRetrying,
).Find(&deliveries), byStatus)
assertPlanUses(t, db, dry.Where(
"status = ? AND updated_at < ?",
database.DeliveryStatusPending, cutoff,
).Limit(500).Find(&deliveries), byStatus)
assertPlanUses(t, db, dry.Model(&database.Delivery{}).
Select("target_id", "status", "count(*) as depth").
Where("status IN ?", []database.DeliveryStatus{
database.DeliveryStatusPending,
database.DeliveryStatusRetrying,
}).Group("target_id, status").Find(&depths), byStatus)
// The event log: each event's deliveries, then their attempts
// (loadEventsWithDeliveries, loadDeliveryResults).
assertPlanUses(t, db, dry.Where("event_id = ?", ids[0]).
Find(&deliveries), byEvent)
assertPlanUses(t, db, dry.Where("delivery_id IN ?", ids).
Order("attempt_num ASC").Find(&results),
"idx_delivery_results_delivery_id (delivery_id=? AND deleted_at=?)")
// 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").
Where("created_at < ?", cutoff)
}
assertPlanUses(t, db, dry.Unscoped().Where(
"delivery_id IN (?)", dry.Model(&database.Delivery{}).
Select("id").Where("event_id IN (?)", expiredEventIDs()),
).Delete(&database.DeliveryResult{}),
"idx_delivery_results_delivery_id (delivery_id=?)", byEvent, byAge)
assertPlanUses(t, db, dry.Unscoped().Where(
"event_id IN (?)", expiredEventIDs(),
).Delete(&database.Delivery{}),
"idx_deliveries_event_id (event_id=?)", byAge)
assertPlanUses(t, db, dry.Unscoped().Where(
"created_at < ?", cutoff,
).Delete(&database.Event{}), "idx_events_created_at (created_at<?)")
}
// 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.
func assertPlanUses(
t *testing.T, db, built *gorm.DB, indexes ...string,
) {
t.Helper()
var plan []struct{ Detail string }
require.NoError(t, db.Raw(
"EXPLAIN QUERY PLAN "+built.Statement.SQL.String(),
built.Statement.Vars...,
).Scan(&plan).Error)
for _, index := range indexes {
assert.Contains(t, fmt.Sprint(plan), index,
built.Statement.SQL.String())
}
}
+12 -3
View File
@@ -1,5 +1,7 @@
package database package database
import "gorm.io/gorm"
// DeliveryStatus represents the status of a delivery // DeliveryStatus represents the status of a delivery
type DeliveryStatus string type DeliveryStatus string
@@ -29,12 +31,19 @@ func (s DeliveryStatus) Terminal() bool {
} }
// Delivery represents a delivery attempt for an event to a target // Delivery represents a delivery attempt for an event to a target
//
//nolint:lll // a struct tag cannot wrap
type Delivery struct { type Delivery struct {
BaseModel BaseModel
EventID string `gorm:"type:uuid;not null" json:"eventId"` EventID string `gorm:"type:uuid;not null;index:idx_deliveries_event_id,priority:1" json:"eventId"`
TargetID string `gorm:"type:uuid;not null" json:"targetId"` TargetID string `gorm:"type:uuid;not null" json:"targetId"`
Status DeliveryStatus `gorm:"not null;default:'pending'" json:"status"` Status DeliveryStatus `gorm:"not null;default:'pending';index:idx_deliveries_status,priority:1" json:"status"`
// DeletedAt repeats the BaseModel field only to be the second column
// of the event_id and status indexes, for the reason DeliveryResult
// gives.
DeletedAt gorm.DeletedAt `gorm:"index:idx_deliveries_event_id,priority:2;index:idx_deliveries_status,priority:2" json:"deletedAt,omitzero"`
// Relations // Relations
Event Event `json:"event,omitzero"` Event Event `json:"event,omitzero"`
+12 -1
View File
@@ -1,10 +1,21 @@
package database package database
import "gorm.io/gorm"
// DeliveryResult represents the result of a delivery attempt // DeliveryResult represents the result of a delivery attempt
//
//nolint:lll // a struct tag cannot wrap
type DeliveryResult struct { type DeliveryResult struct {
BaseModel BaseModel
DeliveryID string `gorm:"type:uuid;not null" json:"deliveryId"` // DeliveryID and DeletedAt make up one index, in that order.
// DeletedAt repeats the BaseModel field only to join it: GORM adds
// "deleted_at IS NULL" to almost every query, and where a column is
// matched against several values SQLite otherwise reads through the
// deleted_at index, which every live row matches.
DeliveryID string `gorm:"type:uuid;not null;index:idx_delivery_results_delivery_id,priority:1" json:"deliveryId"`
DeletedAt gorm.DeletedAt `gorm:"index:idx_delivery_results_delivery_id,priority:2" json:"deletedAt,omitzero"`
AttemptNum int `gorm:"not null" json:"attemptNum"` AttemptNum int `gorm:"not null" json:"attemptNum"`
Success bool `json:"success"` Success bool `json:"success"`
StatusCode int `json:"statusCode,omitempty"` StatusCode int `json:"statusCode,omitempty"`
+18
View File
@@ -1,9 +1,27 @@
package database package database
import (
"time"
"gorm.io/gorm"
)
// Event represents a captured webhook event // Event represents a captured webhook event
//
//nolint:lll // a struct tag cannot wrap
type Event struct { type Event struct {
BaseModel BaseModel
// CreatedAt and DeletedAt repeat the BaseModel fields only to index
// them for retention, which finds events by age. Its lookups carry
// GORM's "deleted_at IS NULL" (see DeliveryResult) and compare
// created_at with <, so their index has deleted_at first: SQLite
// narrows by a < only on the last column it uses. Its final delete
// has no deleted_at condition and uses the index on created_at
// alone. The other tables keep the unindexed BaseModel created_at.
CreatedAt time.Time `gorm:"index;index:idx_events_deleted_at_created_at,priority:2" json:"createdAt"`
DeletedAt gorm.DeletedAt `gorm:"index:idx_events_deleted_at_created_at,priority:1" json:"deletedAt,omitzero"`
WebhookID string `gorm:"type:uuid;not null" json:"webhookId"` WebhookID string `gorm:"type:uuid;not null" json:"webhookId"`
EntrypointID string `gorm:"type:uuid;not null" json:"entrypointId"` EntrypointID string `gorm:"type:uuid;not null" json:"entrypointId"`
+42 -29
View File
@@ -13,6 +13,7 @@ import (
"gorm.io/driver/sqlite" "gorm.io/driver/sqlite"
"gorm.io/gorm" "gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/config" "sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/datadir"
"sneak.berlin/go/webhooker/internal/gormlog" "sneak.berlin/go/webhooker/internal/gormlog"
"sneak.berlin/go/webhooker/internal/logger" "sneak.berlin/go/webhooker/internal/logger"
) )
@@ -40,6 +41,11 @@ type WebhookDBManager struct {
dataDir string dataDir string
dbs sync.Map // map[webhookID]*gorm.DB dbs sync.Map // map[webhookID]*gorm.DB
log *slog.Logger log *slog.Logger
// mu is held while a database is opened, deleted, or closed, so
// each file has at most one open handle. Reading an already cached
// handle does not take it.
mu sync.Mutex
} }
// NewWebhookDBManager creates a new WebhookDBManager and // NewWebhookDBManager creates a new WebhookDBManager and
@@ -53,8 +59,9 @@ func NewWebhookDBManager(
log: params.Logger.Get(), log: params.Logger.Get(),
} }
// Create data directory if it doesn't exist // Create data directory if it doesn't exist. datadir.DirPerm is the
err := os.MkdirAll(m.dataDir, dataDirPerm) // single source of the directory mode; either package may run first.
err := os.MkdirAll(m.dataDir, datadir.DirPerm)
if err != nil { if err != nil {
return nil, fmt.Errorf( return nil, fmt.Errorf(
"creating data directory %s: %w", "creating data directory %s: %w",
@@ -84,43 +91,39 @@ func (m *WebhookDBManager) GetDB(
) (*gorm.DB, error) { ) (*gorm.DB, error) {
// Fast path: already open // Fast path: already open
if val, ok := m.dbs.Load(webhookID); ok { if val, ok := m.dbs.Load(webhookID); ok {
cachedDB, castOK := val.(*gorm.DB) return asGormDB(val, webhookID)
if !castOK { }
return nil, fmt.Errorf(
"%w for webhook %s",
errInvalidCachedDBType,
webhookID,
)
}
return cachedDB, nil // Slow path: open the database under the lock, looking in the
// cache again first. A caller that raced another one here then
// waits for its handle instead of opening a second one.
m.mu.Lock()
defer m.mu.Unlock()
if val, ok := m.dbs.Load(webhookID); ok {
return asGormDB(val, webhookID)
} }
// Slow path: open/create the database
db, err := m.openDB(webhookID) db, err := m.openDB(webhookID)
if err != nil { if err != nil {
return nil, err return nil, err
} }
// Store it; if another goroutine beat us, close ours m.dbs.Store(webhookID, db)
actual, loaded := m.dbs.LoadOrStore(webhookID, db)
if loaded {
// Another goroutine created it first; close our duplicate
sqlDB, closeErr := db.DB()
if closeErr == nil {
_ = sqlDB.Close()
}
existingDB, castOK := actual.(*gorm.DB) return db, nil
if !castOK { }
return nil, fmt.Errorf(
"%w for webhook %s",
errInvalidCachedDBType,
webhookID,
)
}
return existingDB, nil // asGormDB returns a value read from the cache as the database
// handle it is.
func asGormDB(val any, webhookID string) (*gorm.DB, error) {
db, ok := val.(*gorm.DB)
if !ok {
return nil, fmt.Errorf(
"%w for webhook %s",
errInvalidCachedDBType,
webhookID,
)
} }
return db, nil return db, nil
@@ -151,6 +154,11 @@ func (m *WebhookDBManager) DBExists(
func (m *WebhookDBManager) DeleteDB( func (m *WebhookDBManager) DeleteDB(
webhookID string, webhookID string,
) error { ) error {
// Held until the files are gone, so GetDB cannot open the file
// again between the close and the removal.
m.mu.Lock()
defer m.mu.Unlock()
// Close and remove from cache // Close and remove from cache
if val, ok := m.dbs.LoadAndDelete(webhookID); ok { if val, ok := m.dbs.LoadAndDelete(webhookID); ok {
if gormDB, castOK := val.(*gorm.DB); castOK { if gormDB, castOK := val.(*gorm.DB); castOK {
@@ -184,6 +192,11 @@ func (m *WebhookDBManager) DeleteDB(
// CloseAll closes all open per-webhook database connections. // CloseAll closes all open per-webhook database connections.
// Called during application shutdown. // Called during application shutdown.
func (m *WebhookDBManager) CloseAll() error { func (m *WebhookDBManager) CloseAll() error {
// An open already under way finishes and is cached first, so it
// is closed here rather than cached after this loop has passed.
m.mu.Lock()
defer m.mu.Unlock()
var lastErr error var lastErr error
m.dbs.Range(func(key, value any) bool { m.dbs.Range(func(key, value any) bool {
@@ -1,10 +1,14 @@
package database_test package database_test
import ( import (
"bytes"
"context" "context"
"log/slog"
"net/http" "net/http"
"os" "os"
"path/filepath" "path/filepath"
"strings"
"sync"
"testing" "testing"
"github.com/google/uuid" "github.com/google/uuid"
@@ -104,6 +108,54 @@ func TestWebhookDBManager_CreateAndGetDB(t *testing.T) {
assert.Equal(t, `{"test": true}`, readEvent.Body) assert.Equal(t, `{"test": true}`, readEvent.Body)
} }
// Many callers ask for one webhook's database at the same moment,
// before it is cached. Only one of them may open the file; the others
// must wait for its handle. openDB logs one "opened per-webhook
// database" line per open, and those lines are what is counted.
func TestWebhookDBManager_ConcurrentFirstTouchOpensOnce(t *testing.T) {
t.Parallel()
var logs bytes.Buffer
mgr := database.NewTestWebhookDBManagerWithLogger(
t.TempDir(),
slog.New(slog.NewTextHandler(&logs, nil)),
)
t.Cleanup(func() { assert.NoError(t, mgr.CloseAll()) })
webhookID := uuid.New().String()
const callers = 16
start := make(chan struct{})
handles := make([]*gorm.DB, callers)
errs := make([]error, callers)
var wg sync.WaitGroup
for i := range callers {
wg.Go(func() {
<-start
handles[i], errs[i] = mgr.GetDB(webhookID)
})
}
close(start)
wg.Wait()
for i := range callers {
require.NoError(t, errs[i])
assert.Same(t, handles[0], handles[i])
}
assert.Equal(
t, 1,
strings.Count(logs.String(), "opened per-webhook database"),
)
}
func TestWebhookDBManager_DeleteDB(t *testing.T) { func TestWebhookDBManager_DeleteDB(t *testing.T) {
t.Parallel() t.Parallel()
+6 -4
View File
@@ -29,9 +29,11 @@ import (
// process that was killed with SIGKILL blocks nothing. // process that was killed with SIGKILL blocks nothing.
const LockFileName = "webhooker.lock" const LockFileName = "webhooker.lock"
// dirPerm is the mode Acquire creates DATA_DIR with. It matches what // DirPerm is the mode DATA_DIR is created with. It is the single
// internal/database uses, since whichever runs first creates it. // source of that mode: internal/database consumes it rather than
const dirPerm = 0o750 // keeping its own copy, so the two packages that both create the
// directory cannot drift into disagreeing about its permissions.
const DirPerm = 0o750
// ErrLocked reports that another live process holds the data // ErrLocked reports that another live process holds the data
// directory. Callers that need to know whether a deployment is running // directory. Callers that need to know whether a deployment is running
@@ -64,7 +66,7 @@ func Acquire(dir string) (*Lock, error) {
return nil, ErrNoDir return nil, ErrNoDir
} }
err := os.MkdirAll(dir, dirPerm) err := os.MkdirAll(dir, DirPerm)
if err != nil { if err != nil {
return nil, fmt.Errorf( return nil, fmt.Errorf(
"creating data directory %s: %w", dir, err, "creating data directory %s: %w", dir, err,
+10 -2
View File
@@ -76,12 +76,20 @@ func (cb *CircuitBreaker) Allow() bool {
} }
} }
// CooldownRemaining returns how much time is left before // CooldownRemaining returns how long a delivery that Allow refused
// an open circuit transitions to half-open. // should wait before it is tried again. Closed, it returns zero.
// Open, it returns what is left of the cooldown, or zero once that
// has passed. Half-open, it returns the whole cooldown: the one
// probe delivery is still in flight, and if it fails the circuit
// reopens for that long.
func (cb *CircuitBreaker) CooldownRemaining() time.Duration { func (cb *CircuitBreaker) CooldownRemaining() time.Duration {
cb.mu.Lock() cb.mu.Lock()
defer cb.mu.Unlock() defer cb.mu.Unlock()
if cb.state == CircuitHalfOpen {
return cb.cooldown
}
if cb.state != CircuitOpen { if cb.state != CircuitOpen {
return 0 return 0
} }
+5 -3
View File
@@ -267,7 +267,7 @@ func TestCircuitBreaker_CooldownRemaining_ClosedReturnsZero(
) )
} }
func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsZero( func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsCooldown(
t *testing.T, t *testing.T,
) { ) {
t.Parallel() t.Parallel()
@@ -282,9 +282,11 @@ func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsZero(
require.True(t, cb.Allow()) require.True(t, cb.Allow())
assert.Equal(t, time.Duration(0), // The cooldown newShortCooldownCB gives the breaker.
assert.Equal(t, 50*time.Millisecond,
cb.CooldownRemaining(), cb.CooldownRemaining(),
"half-open circuit should have zero cooldown remaining", "a delivery refused while half-open should wait "+
"a whole cooldown",
) )
} }
+97 -25
View File
@@ -362,6 +362,15 @@ func (e *Engine) start() {
// stop cancels the worker pool's context and waits for the pool // stop cancels the worker pool's context and waits for the pool
// to drain, bounded by the stop hook's context: a wedged worker // to drain, bounded by the stop hook's context: a wedged worker
// must not hang the process past fx's stop timeout. // must not hang the process past fx's stop timeout.
//
// Once the pool has drained it closes the archive writers, so a
// clean stop leaves no archive -wal behind. Nothing else holds a
// writer for long by then: the archive sweeper stops before the
// engine, and deleting a webhook only closes one. If the pool did
// not drain in time, the writers are left open, as a kill would
// leave them. Closing them would wait for any write in progress,
// and a worker still running would then open new writers that
// nothing closes, so it gains nothing over a kill.
func (e *Engine) stop(ctx context.Context) error { func (e *Engine) stop(ctx context.Context) error {
e.log.Info("delivery engine stopping") e.log.Info("delivery engine stopping")
@@ -376,6 +385,8 @@ func (e *Engine) stop(ctx context.Context) error {
return err return err
} }
e.dbTarget.evictAll()
e.log.Info("delivery engine stopped") e.log.Info("delivery engine stopped")
return nil return nil
@@ -438,6 +449,31 @@ func (e *Engine) processNewTask(
return return
} }
// Restart recovery can send and release this delivery before the
// receiver's Notify queues it. Ownership cannot refuse a delivery
// nobody holds, so the row decides whether it still needs sending.
row, err := e.loadDelivery(webhookDB, task.DeliveryID)
if err != nil {
e.log.Error(
"failed to load delivery",
"delivery_id", task.DeliveryID,
"error", err,
)
return
}
if row.Status != database.DeliveryStatusPending {
e.log.Info(
"delivery already handled, not sent again",
"delivery_id", task.DeliveryID,
"event_id", task.EventID,
"status", row.Status,
)
return
}
event := buildEventFromTask(task) event := buildEventFromTask(task)
event, err = e.hydrateEvent( event, err = e.hydrateEvent(
@@ -482,7 +518,7 @@ func (e *Engine) processRetryTask(
return return
} }
d, err := e.loadRetryDelivery( d, err := e.loadDelivery(
webhookDB, task.DeliveryID, webhookDB, task.DeliveryID,
) )
if err != nil { if err != nil {
@@ -703,9 +739,7 @@ func (e *Engine) recoverSingleRetry(
// webhook on one bad read would be a far larger fault than // webhook on one bad read would be a far larger fault than
// the strand it is meant to clear. // the strand it is meant to clear.
if errors.Is(err, gorm.ErrRecordNotFound) { if errors.Is(err, gorm.ErrRecordNotFound) {
e.failMissingTargetRetry( e.failMissingTarget(webhookDB, webhookID, d)
webhookDB, webhookID, d,
)
return return
} }
@@ -1108,9 +1142,7 @@ func (e *Engine) sweepSingleRetry(
// Deleted is terminal, unreadable is not; see // Deleted is terminal, unreadable is not; see
// recoverSingleRetry. // recoverSingleRetry.
if errors.Is(err, gorm.ErrRecordNotFound) { if errors.Is(err, gorm.ErrRecordNotFound) {
e.failMissingTargetRetry( e.failMissingTarget(webhookDB, webhookID, d)
webhookDB, webhookID, d,
)
return return
} }
@@ -1224,19 +1256,19 @@ func (e *Engine) failUnretryableRetry(
e.failDelivery(webhookDB, d, target.Type, reason) e.failDelivery(webhookDB, d, target.Type, reason)
} }
// failMissingTargetRetry terminally fails an orphaned retrying // failMissingTarget terminally fails a recovered delivery, pending or
// delivery whose target row is gone. Both restart recovery and the // retrying, whose target row is gone. Restart recovery and the periodic
// periodic sweep call it, so the transition exists once. // sweep call it for both statuses, so the transition exists once.
// //
// Until it existed both paths logged the failed lookup and returned, // Until it existed those paths logged the failed lookup and moved on,
// which left the delivery retrying for the life of the database and // which left the delivery where it was for the life of the database and
// the sweep repeating the same error every minute forever. Failing it // the sweep repeating the same error every minute forever. Failing it
// with a recorded reason is the treatment the other orphaned-retry // with a recorded reason is the treatment the other orphaned-retry
// cases already get, so all of them read alike in the event log. // cases already get, so all of them read alike in the event log.
// //
// Logged at warn rather than error: a deleted target is an operator // Logged at warn rather than error: a deleted target is an operator
// action, not a system fault. // action, not a system fault.
func (e *Engine) failMissingTargetRetry( func (e *Engine) failMissingTarget(
webhookDB *gorm.DB, webhookDB *gorm.DB,
webhookID string, webhookID string,
d *database.Delivery, d *database.Delivery,
@@ -1249,13 +1281,37 @@ func (e *Engine) failMissingTargetRetry(
defer e.inflight.release(d.ID) defer e.inflight.release(d.ID)
// The batch was read before ownership was taken, and a worker may
// have settled the delivery and let it go in between. Only a row
// still in the status the batch read is failed.
row, err := e.loadDelivery(webhookDB, d.ID)
if err != nil {
e.log.Error(
"failed to load delivery",
"delivery_id", d.ID,
"error", err,
)
return
}
if row.Status != d.Status {
e.log.Debug(
"delivery already handled, not failed",
"delivery_id", d.ID,
"status", row.Status,
)
return
}
targetType, reason := e.missingTargetReason(d.TargetID) targetType, reason := e.missingTargetReason(d.TargetID)
e.log.Warn( e.log.Warn(
"failing orphaned retrying delivery: "+ "failing recovered delivery: its target no longer exists",
"its target no longer exists",
"webhook_id", webhookID, "webhook_id", webhookID,
"delivery_id", d.ID, "delivery_id", d.ID,
"status", d.Status,
"target_id", d.TargetID, "target_id", d.TargetID,
"target_type", targetType, "target_type", targetType,
) )
@@ -1289,15 +1345,14 @@ func (e *Engine) missingTargetReason(
if err != nil { if err != nil {
return "", fmt.Sprintf( return "", fmt.Sprintf(
"target %s no longer exists; the delivery "+ "target %s no longer exists; the delivery "+
"cannot be retried and has been failed "+ "has been failed terminally",
"terminally",
targetID, targetID,
) )
} }
return target.Type, fmt.Sprintf( return target.Type, fmt.Sprintf(
"target %q (type %s) was deleted; the delivery "+ "target %q (type %s) was deleted; the delivery "+
"cannot be retried and has been failed terminally", "has been failed terminally",
target.Name, target.Type, target.Name, target.Type,
) )
} }
@@ -1643,7 +1698,7 @@ func (e *Engine) hydrateEvent(
return event, nil return event, nil
} }
func (e *Engine) loadRetryDelivery( func (e *Engine) loadDelivery(
webhookDB *gorm.DB, deliveryID string, webhookDB *gorm.DB, deliveryID string,
) (*database.Delivery, error) { ) (*database.Delivery, error) {
var d database.Delivery var d database.Delivery
@@ -1996,13 +2051,30 @@ func (e *Engine) sendRecoveredDeliveries(
target, ok := targetMap[deliveries[i].TargetID] target, ok := targetMap[deliveries[i].TargetID]
if !ok { if !ok {
e.log.Error( // A missing entry does not mean the target is gone: the
"target not found for delivery", // map is also empty when its query failed. Only a lookup
"delivery_id", deliveries[i].ID, // that finds no row ends the delivery; any other error
"target_id", deliveries[i].TargetID, // leaves it pending for the next sweep. See
) // recoverSingleRetry.
var err error
continue target, err = e.loadTarget(deliveries[i].TargetID)
if errors.Is(err, gorm.ErrRecordNotFound) {
e.failMissingTarget(webhookDB, webhookID, &deliveries[i])
continue
}
if err != nil {
e.log.Error(
"failed to load target for recovered delivery",
"delivery_id", deliveries[i].ID,
"target_id", deliveries[i].TargetID,
"error", err,
)
continue
}
} }
if !e.takeForRedispatch( if !e.takeForRedispatch(
@@ -2,6 +2,8 @@ package delivery_test
import ( import (
"context" "context"
"fmt"
"path/filepath"
"testing" "testing"
"time" "time"
@@ -269,3 +271,88 @@ func TestEngine_StopHookHonoursStopTimeout(t *testing.T) {
requireStopHookExpires(t, lc.hooks[0], "delivery engine") requireStopHookExpires(t, lc.hooks[0], "delivery engine")
} }
// deliverToArchive runs one delivery to a database target through
// the running engine and returns the webhook's archive file path.
// The archive writer holds the file open afterwards.
func deliverToArchive(t *testing.T, s iSetup) string {
t.Helper()
deliveryID, task := seedLogTask(t, s)
task.TargetType = database.TargetTypeDatabase
s.Engine.Notify([]delivery.Task{task})
iWaitForDelivered(t, s.WebhookDB, deliveryID)
return filepath.Join(
filepath.Dir(s.DBMgr.DBPath(s.WebhookID)),
fmt.Sprintf("archive-%s.db", s.WebhookID),
)
}
// TestEngine_StopHookClosesArchives is the regression test for an
// archive split across two files by a clean stop. The engine never
// closed its archive writers, so after a stop the archived rows
// could sit in archive-{id}.db-wal while archive-{id}.db held no
// table at all, and copying the .db on its own gave an empty
// database.
func TestEngine_StopHookClosesArchives(t *testing.T) {
t.Parallel()
s := newISetup(t)
lc := startEngineViaHook(t, s.Engine)
path := deliverToArchive(t, s)
require.FileExists(
t, path+"-wal",
"an open archive should have a -wal for the stop to remove",
)
require.NoError(t, lc.hooks[0].OnStop(context.Background()))
wals, err := filepath.Glob(
filepath.Join(filepath.Dir(path), "archive-*.db-wal"),
)
require.NoError(t, err)
require.Empty(
t, wals, "a clean stop must leave no archive -wal behind",
)
// With no -wal beside it, the row can only be in the .db.
count, err := countArchivedRows(path)
require.NoError(t, err)
require.Equal(t, int64(1), count)
}
// TestEngine_StopHookTimeoutLeavesArchivesOpen covers a stop whose
// budget runs out while a worker is still running. The archive
// writers are left open, as a kill would leave them: closing them
// would wait for any write in progress, and that worker would then
// open new writers that nothing closes.
func TestEngine_StopHookTimeoutLeavesArchivesOpen(t *testing.T) {
t.Parallel()
s := newISetup(t)
lc := startEngineViaHook(t, s.Engine)
deliverToArchive(t, s)
release := make(chan struct{})
t.Cleanup(func() {
close(release)
s.Engine.EvictWebhook(s.WebhookID)
})
s.Engine.ExportWedgeWorker(release)
requireStopHookExpires(t, lc.hooks[0], "delivery engine")
require.True(
t, s.Engine.ExportArchiveHandleOpen(s.WebhookID),
"a stop that timed out must not close archive writers",
)
}
+175
View File
@@ -17,6 +17,7 @@ import (
"time" "time"
"github.com/google/uuid" "github.com/google/uuid"
"github.com/prometheus/client_golang/prometheus"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
"gorm.io/driver/sqlite" "gorm.io/driver/sqlite"
@@ -24,6 +25,7 @@ import (
_ "modernc.org/sqlite" _ "modernc.org/sqlite"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/delivery"
"sneak.berlin/go/webhooker/internal/metrics"
) )
// testContentType is the event content type used in tests. // testContentType is the event content type used in tests.
@@ -894,6 +896,100 @@ func TestDeliverHTTP_CircuitBreakerBlocks(t *testing.T) {
) )
} }
// recordingScheduler keeps the delay of every retry it is asked to
// schedule, and schedules nothing.
type recordingScheduler struct {
delays []time.Duration
}
func (s *recordingScheduler) ScheduleRetry(
_ delivery.Task, delay time.Duration,
) {
s.delays = append(s.delays, delay)
}
// TestDeliverHTTP_HalfOpenBreakerDelaysQueuedTasks proves that while a
// half-open breaker's one probe delivery is in flight, every other task
// for the target is put back with a whole cooldown as its delay rather
// than none, and that its status is written the first time the breaker
// turns it away and not on each pass after that.
func TestDeliverHTTP_HalfOpenBreakerDelaysQueuedTasks(t *testing.T) {
t.Parallel()
db := testWebhookDB(t)
e := testEngine(t, 1)
// Every write of retrying moves the retry counter, so on a registry
// this test owns the counter is the number of those writes.
reg := prometheus.NewRegistry()
e.ExportSetMetrics(metrics.New(reg))
targetID := uuid.New().String()
cb := newShortCooldownCB(t)
e.ExportSetCircuitBreaker(targetID, cb)
for range delivery.ExportDefaultFailureThreshold {
cb.RecordFailure()
}
time.Sleep(60 * time.Millisecond)
require.True(t, cb.Allow(), "the probe delivery should go through")
require.Equal(t, delivery.CircuitHalfOpen, cb.State())
cfg := newHTTPTargetConfig(
"http://will-not-be-called.invalid",
)
sched := &recordingScheduler{}
const queued, passes = 3, 4
for range queued {
event := seedEvent(t, db, `{"cb":"half-open"}`)
dlv := seedDelivery(
t, db, event.ID, targetID,
database.DeliveryStatusPending,
)
for range passes {
// Each pass starts from the stored row, as a retry does.
var row database.Delivery
require.NoError(t, db.First(
&row, "id = ?", dlv.ID,
).Error)
fix := buildHTTPFixture(
row, event, targetID,
"test-cb-half-open", cfg, 5, 1,
)
e.ExportDeliverHTTPWithScheduler(
context.TODO(), db, fix.Delivery, fix.Task, sched,
)
}
assertDeliveryStatus(t, db, dlv.ID,
database.DeliveryStatusRetrying,
)
}
require.Len(t, sched.delays, queued*passes)
for _, delay := range sched.delays {
// The cooldown newShortCooldownCB gives the breaker.
assert.Equal(t, 50*time.Millisecond, delay,
"a task turned away while half-open should wait "+
"a whole cooldown",
)
}
assert.InDelta(t, float64(queued),
mCounter(t, reg, mRetries, mTypeHTTP), 0,
"status should be written once per task, not once per pass",
)
}
func TestGetCircuitBreaker_CreatesOnDemand(t *testing.T) { func TestGetCircuitBreaker_CreatesOnDemand(t *testing.T) {
t.Parallel() t.Parallel()
@@ -1070,6 +1166,10 @@ func TestIsForwardableHeader(t *testing.T) {
assert.False(t, assert.False(t,
delivery.ExportIsForwardableHeader("Content-Length"), delivery.ExportIsForwardableHeader("Content-Length"),
) )
assert.False(t,
delivery.ExportIsForwardableHeader("Content-Type"),
)
} }
func TestTruncate(t *testing.T) { func TestTruncate(t *testing.T) {
@@ -1151,6 +1251,81 @@ func TestDoHTTPRequest_ForwardsHeaders(t *testing.T) {
) )
} }
// The event's stored inbound headers carry the same Content-Type the
// receiver saved as the event's ContentType, so a delivery could send
// it twice. It must go out exactly once, with a Content-Type configured
// on the target winning, then the event's ContentType.
func TestApplyRequestHeaders_SendsOneContentType(t *testing.T) {
t.Parallel()
cases := map[string]struct {
inbound string
event string
configured string
want []string
}{
"inbound and event agree": {
inbound: testContentType,
event: testContentType,
want: []string{testContentType},
},
"inbound and event disagree": {
inbound: "text/plain",
event: testContentType,
want: []string{testContentType},
},
"event has none": {
inbound: testContentType,
want: nil,
},
"target configures its own": {
inbound: testContentType,
event: testContentType,
configured: "application/xml",
want: []string{"application/xml"},
},
}
for name, tc := range cases {
t.Run(name, func(t *testing.T) {
t.Parallel()
inbound, err := json.Marshal(map[string][]string{
headerContentType: {tc.inbound},
})
require.NoError(t, err)
cfg := &delivery.HTTPTargetConfig{}
if tc.configured != "" {
cfg.Headers = map[string]string{
headerContentType: tc.configured,
}
}
req, err := http.NewRequestWithContext(
context.Background(),
http.MethodPost,
"https://target.example.com/hook",
http.NoBody,
)
require.NoError(t, err)
delivery.ExportApplyRequestHeaders(
req,
&database.Event{
Headers: string(inbound),
ContentType: tc.event,
},
cfg,
)
assert.Equal(t,
tc.want, req.Header.Values(headerContentType),
)
})
}
}
func TestProcessDelivery_RoutesToCorrectHandler( func TestProcessDelivery_RoutesToCorrectHandler(
t *testing.T, t *testing.T,
) { ) {
+46
View File
@@ -101,6 +101,19 @@ func (e *Engine) ExportDeliverHTTP(
e.httpTarget.Deliver(ctx, webhookDB, d, task, e) e.httpTarget.Deliver(ctx, webhookDB, d, task, e)
} }
// ExportDeliverHTTPWithScheduler delivers via the http target, handing
// any retry to sched instead of the engine, so a test can see the
// delay each retry is given.
func (e *Engine) ExportDeliverHTTPWithScheduler(
ctx context.Context,
webhookDB *gorm.DB,
d *database.Delivery,
task *Task,
sched Scheduler,
) {
e.httpTarget.Deliver(ctx, webhookDB, d, task, sched)
}
// ExportDeliverDatabase delivers via the database target. // ExportDeliverDatabase delivers via the database target.
func (e *Engine) ExportDeliverDatabase( func (e *Engine) ExportDeliverDatabase(
webhookDB *gorm.DB, d *database.Delivery, webhookDB *gorm.DB, d *database.Delivery,
@@ -179,6 +192,14 @@ func (e *Engine) ExportGetCircuitBreaker(
return e.httpTarget.getCircuitBreaker(targetID) return e.httpTarget.getCircuitBreaker(targetID)
} }
// ExportSetCircuitBreaker makes cb the http target's circuit breaker
// for targetID, so a test can use one with a short cooldown.
func (e *Engine) ExportSetCircuitBreaker(
targetID string, cb *CircuitBreaker,
) {
e.httpTarget.circuitBreakers.Store(targetID, cb)
}
// ExportParseHTTPConfig exposes parseHTTPConfig. // ExportParseHTTPConfig exposes parseHTTPConfig.
func (e *Engine) ExportParseHTTPConfig( func (e *Engine) ExportParseHTTPConfig(
configJSON string, configJSON string,
@@ -321,6 +342,31 @@ func (e *Engine) ExportRecoverRetryingDeliveries(
e.recoverRetryingDeliveries(webhookDB, webhookID) e.recoverRetryingDeliveries(webhookDB, webhookID)
} }
// ExportFailMissingTarget exposes failMissingTarget, so a test can hand
// it a delivery as a batch read it earlier.
func (e *Engine) ExportFailMissingTarget(
webhookDB *gorm.DB,
webhookID string,
d *database.Delivery,
) {
e.failMissingTarget(webhookDB, webhookID, d)
}
// ExportSendRecoveredDeliveries exposes sendRecoveredDeliveries, so a
// test can hand it a target map that lacks a delivery's target.
func (e *Engine) ExportSendRecoveredDeliveries(
ctx context.Context,
webhookDB *gorm.DB,
deliveries []database.Delivery,
webhookID string,
targetMap map[string]database.Target,
settled map[string]struct{},
) {
e.sendRecoveredDeliveries(
ctx, webhookDB, deliveries, webhookID, targetMap, settled,
)
}
// ExportDeliveryCh returns the delivery channel. // ExportDeliveryCh returns the delivery channel.
func (e *Engine) ExportDeliveryCh() chan Task { func (e *Engine) ExportDeliveryCh() chan Task {
return e.deliveryCh return e.deliveryCh
+67
View File
@@ -273,6 +273,73 @@ func TestOwnershipIsReleasedAfterDelivery(t *testing.T) {
) )
} }
// TestNotifyAfterRecoveryDoesNotSendAgain is the startup race of
// https://git.eeqj.de/sneak/webhooker/issues/299. The receiver has
// written a delivery, restart recovery finds it pending, sends it and
// releases it, and only then does the receiver's Notify for it arrive.
// Nothing owns the delivery by then, so Notify takes it.
func TestNotifyAfterRecoveryDoesNotSendAgain(t *testing.T) {
t.Parallel()
targetID := uuid.New().String()
s := fSweepSetup(t, targetID, "recovered")
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"recovered":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
s.Engine.ExportStart()
defer func() {
require.NoError(
t, s.Engine.ExportStop(context.Background()),
)
}()
// Restart recovery sends the delivery and lets it go.
iWaitForDelivered(t, s.WebhookDB, d.ID)
require.Eventually(
t,
func() bool {
return s.Engine.ExportInflightHeld() == 0
},
5*time.Second, 20*time.Millisecond,
)
body := event.Body
s.Engine.Notify([]delivery.Task{{
DeliveryID: d.ID,
EventID: event.ID,
WebhookID: s.WebhookID,
TargetID: targetID,
TargetName: "recovered",
TargetType: database.TargetTypeLog,
Body: &body,
EntrypointID: event.EntrypointID,
}})
// Notify took the delivery, and a worker releases it once it has
// run the task.
require.Eventually(
t,
func() bool {
return s.Engine.ExportInflightHeld() == 0
},
5*time.Second, 20*time.Millisecond,
)
assert.Len(
t, iResults(t, s.WebhookDB, d.ID), 1,
"the delivery was sent a second time",
)
}
// TestRetryingRecoverySkipsASuccessfulResult is the retrying-side twin // TestRetryingRecoverySkipsASuccessfulResult is the retrying-side twin
// of the pending reconcile. A second attempt that reached the receiver // of the pending reconcile. A second attempt that reached the receiver
// and whose status write then failed sits at retrying holding a // and whose status write then failed sits at retrying holding a
+4 -3
View File
@@ -412,9 +412,10 @@ func TestDeliveryMetrics_BreakerBlockedIsNotAnAttempt(
s.Engine.ExportProcessRetryTask(context.TODO(), &blocked) s.Engine.ExportProcessRetryTask(context.TODO(), &blocked)
// The breaker refused it: rescheduled, so the retry counter // The breaker refused it: rescheduled without rewriting the
// moved, but nothing was attempted or timed. // retrying status it already had, so the retry counter did not
assert.InDelta(t, retriesBefore+1, // move, and nothing was attempted or timed.
assert.InDelta(t, retriesBefore,
mCounter(t, reg, mRetries, mTypeHTTP), 0) mCounter(t, reg, mRetries, mTypeHTTP), 0)
assert.InDelta(t, threshold, assert.InDelta(t, threshold,
mCounter(t, reg, mAttempts, mTypeHTTP), 0) mCounter(t, reg, mAttempts, mTypeHTTP), 0)
+14 -7
View File
@@ -170,7 +170,8 @@ func TestDelivery_CrossOriginRedirectDropsOriginScopedHeaders(
// Stripping must not fire within the configured origin, or every // Stripping must not fire within the configured origin, or every
// destination that redirects its own path would lose its // destination that redirects its own path would lose its
// credential and start answering 401 — and would lose the inbound // credential and start answering 401 — and would lose the inbound
// signature the receiver verifies. // signature header the target endpoint verifies. webhooker's own
// receiver verifies no signature; it only forwards the header.
func TestDelivery_SameOriginRedirectKeepsOriginScopedHeaders( func TestDelivery_SameOriginRedirectKeepsOriginScopedHeaders(
t *testing.T, t *testing.T,
) { ) {
@@ -338,10 +339,11 @@ func TestRedirectPolicy_StopsAtHopCap(t *testing.T) {
// The set the redirect policy strips is whatever the delivery path // The set the redirect policy strips is whatever the delivery path
// actually put on the wire, so a header added to the forward set is // actually put on the wire, so a header added to the forward set is
// covered without a second edit. A header the event never carried // covered without a second edit. A header the event never carried
// is not in the set, and the delivery path's own two are deliberately // is not in the set, and neither is the inbound Content-Type, because
// excluded: Content-Type describes the body, which a 307 carries // it is not forwarded. Two more are deliberately excluded: a
// across hosts, and the inbound User-Agent every real sender supplies // Content-Type configured on the target describes the body, which a
// is overwritten before the request goes out. // 307 carries across hosts, and the inbound User-Agent every real
// sender supplies is overwritten before the request goes out.
func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) { func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
t.Parallel() t.Parallel()
@@ -370,6 +372,7 @@ func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
&delivery.HTTPTargetConfig{ &delivery.HTTPTargetConfig{
Headers: map[string]string{ Headers: map[string]string{
probeHeaderName: probeHeaderValue, probeHeaderName: probeHeaderValue,
"Content-Type": testContentType,
}, },
}, },
) )
@@ -377,7 +380,11 @@ func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
assert.Equal(t, assert.Equal(t,
[]string{probeHeaderName, inboundHeaderName}, names, []string{probeHeaderName, inboundHeaderName}, names,
"both header classes are reported, and only those: "+ "both header classes are reported, and only those: "+
"Host is never forwarded, Content-Type and "+ "Host and the inbound Content-Type are never "+
"User-Agent are the delivery path's own", "forwarded, User-Agent is the delivery path's own",
)
assert.NotContains(t, names, "Content-Type",
"a Content-Type configured on the target must survive "+
"a cross-origin 307/308 with the body it describes",
) )
} }
+10 -7
View File
@@ -26,7 +26,7 @@ var (
"hostname resolved to no IP addresses", "hostname resolved to no IP addresses",
) )
errBlockedIP = errors.New( errBlockedIP = errors.New(
"blocked private/reserved IP range", "blocked private, reserved or cloud metadata address",
) )
errBlockedMetadata = errors.New( errBlockedMetadata = errors.New(
"blocked link-local or cloud instance metadata " + "blocked link-local or cloud instance metadata " +
@@ -37,9 +37,10 @@ var (
) )
) )
// blockedNetworks contains all private/reserved IP ranges // blockedNetworks is the default blocklist: the private and
// that should be blocked to prevent SSRF attacks. An operator // reserved IP ranges, plus the public cloud metadata addresses,
// can permit specific blocks out of this set with // that are blocked to prevent SSRF attacks. An operator can
// permit specific blocks out of this set with
// ALLOWED_EGRESS_CIDRS; see Guard. // ALLOWED_EGRESS_CIDRS; see Guard.
// //
//nolint:gochecknoglobals // package-level network list is appropriate here //nolint:gochecknoglobals // package-level network list is appropriate here
@@ -122,6 +123,8 @@ func init() {
"::1/128", "::1/128",
"fc00::/7", "fc00::/7",
"fe80::/10", "fe80::/10",
// Azure WireServer, a public address that serves VM credentials.
"168.63.129.16/32",
}) })
// Every entry is named. The set must not grow or shrink // Every entry is named. The set must not grow or shrink
@@ -216,8 +219,8 @@ func matchesAny(networks []*net.IPNet, ip net.IP) bool {
} }
// isBlockedIP checks whether an IP address falls within // isBlockedIP checks whether an IP address falls within
// any blocked private/reserved network range, before any // the default blocklist, before any operator allowlist is
// operator allowlist is considered. // considered.
func isBlockedIP(ip net.IP) bool { func isBlockedIP(ip net.IP) bool {
return matchesAny(blockedNetworks, ip) return matchesAny(blockedNetworks, ip)
} }
@@ -320,7 +323,7 @@ func (g *Guard) allows(ip net.IP) bool {
// //
// 1. alwaysBlockedNetworks is refused before the allowlist is // 1. alwaysBlockedNetworks is refused before the allowlist is
// consulted, so no configured CIDR reaches link-local or a // consulted, so no configured CIDR reaches link-local or a
// cloud instance metadata endpoint. // cloud metadata endpoint at a non-public address.
// 2. The allowlist is consulted next, so a listed private // 2. The allowlist is consulted next, so a listed private
// network becomes reachable. // network becomes reachable.
// 3. Everything else keeps the default blocklist's answer. // 3. Everything else keeps the default blocklist's answer.
+35
View File
@@ -390,6 +390,41 @@ func TestGuardAllowlist_PublicUnaffected(t *testing.T) {
} }
} }
// TestGuardAllowlist_AzureWireServerReopenable covers Azure's
// WireServer, a public address that serves VM credentials. The
// default guard refuses it, but because it is public it sits in
// the default blocklist rather than the unconditional set, so an
// operator who lists it can reach it.
func TestGuardAllowlist_AzureWireServerReopenable(t *testing.T) {
t.Parallel()
const wireServerIP = "168.63.129.16"
target := "http://" + wireServerIP + "/?comp=versions"
defaultGuard := delivery.NewTestGuard()
err := defaultGuard.ValidateTargetURL(context.Background(), target)
require.Error(t, err,
"WireServer must be refused with no allowlist set",
)
assert.NotContains(t, err.Error(), metadataRefusalClause,
"WireServer must be refused by the default blocklist, "+
"which an allowlist can override",
)
assertDialRefused(t, defaultGuard, target)
listed := delivery.NewTestGuard(
netip.MustParsePrefix(wireServerIP + "/32"),
)
assert.NoError(t,
listed.ValidateTargetURL(context.Background(), target),
"an operator who lists WireServer must be able to reach it",
)
}
// TestGuardCheckIP_BothPathsShareOneDecision asserts that the // TestGuardCheckIP_BothPathsShareOneDecision asserts that the
// validator and the dialer are not two policies that happen to // validator and the dialer are not two policies that happen to
// agree: both are defined in terms of checkIP, so the exported // agree: both are defined in terms of checkIP, so the exported
+18
View File
@@ -277,6 +277,24 @@ func (t *databaseTarget) evict(webhookID string) {
) )
} }
// evictAll evicts every cached archive writer, exactly as evict
// does for one webhook. The engine calls it at shutdown, once its
// workers have returned. Closing the last handle on an archive
// moves the contents of its -wal into the .db and removes the
// -wal, so a clean stop leaves each archive as a single file.
func (t *databaseTarget) evictAll() {
t.mu.Lock()
writers := t.writers
t.writers = nil
t.mu.Unlock()
for _, w := range writers {
w.evict()
}
}
// sweepWebhook prunes one webhook's archive of rows older than // sweepWebhook prunes one webhook's archive of rows older than
// expiry, without requiring a write. It returns nil (nothing to // expiry, without requiring a write. It returns nil (nothing to
// do) when the archive file does not exist, so a sweep never // do) when the archive file does not exist, so a sweep never
@@ -1,6 +1,7 @@
package delivery_test package delivery_test
import ( import (
"context"
"errors" "errors"
"fmt" "fmt"
"net/http" "net/http"
@@ -361,3 +362,46 @@ func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) {
"a later delivery should recreate the writer", "a later delivery should recreate the writer",
) )
} }
// TestEngineStop_WriteAfterStopIsRefused proves the engine's stop
// closes each archive writer the way deleting its webhook does: a
// write that reaches a writer after the stop is refused, reopens
// nothing and adds no row.
func TestEngineStop_WriteAfterStopIsRefused(t *testing.T) {
t.Parallel()
eng, _ := evictTestEngine(t)
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":true}`)
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
eng.ExportDeliverDatabase(webhookDB, d)
w := eng.ExportArchiveWriterFor(event.WebhookID)
require.NotNil(t, w)
require.True(t, w.HandleOpen())
require.NoError(t, eng.ExportStop(context.Background()))
err := w.Write(evictTestRow("ev-after-stop"), 0)
require.ErrorIs(
t, err, delivery.ErrExportArchiveWriterEvicted,
"a write after the stop must be refused",
)
assert.False(
t, w.HandleOpen(),
"a refused write must not reopen the archive",
)
assert.False(
t, eng.ExportHasArchiveWriter(event.WebhookID),
"the stop should empty the registry",
)
count, err := countArchivedRows(w.Path())
require.NoError(t, err)
assert.Equal(
t, int64(1), count, "the refused row must not be written",
)
}
+2 -1
View File
@@ -11,10 +11,11 @@ import (
"sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/delivery"
) )
// Literals these tests repeat, named so that the header name and the // Literals these tests repeat, named so that the header names and the
// keep-forever archive config each have one definition. // keep-forever archive config each have one definition.
const ( const (
headerAuthorization = "Authorization" headerAuthorization = "Authorization"
headerContentType = "Content-Type"
bearerValue = "Bearer abc" bearerValue = "Bearer abc"
archiveConfigNever = "{\"expiry\":\"never\"}" archiveConfigNever = "{\"expiry\":\"never\"}"
) )
+21 -8
View File
@@ -197,10 +197,14 @@ func (c *httpCore) circuitBreakerBlock(
"cooldown_remaining", remaining, "cooldown_remaining", remaining,
) )
c.eng.settleStatus( // A delivery already at retrying is left as it is, so a task
webhookDB, d, d.Target.Type, // the breaker keeps turning away writes nothing each time.
database.DeliveryStatusRetrying, if d.Status != database.DeliveryStatusRetrying {
) c.eng.settleStatus(
webhookDB, d, d.Target.Type,
database.DeliveryStatusRetrying,
)
}
retryTask := *task retryTask := *task
sched.ScheduleRetry(retryTask, remaining) sched.ScheduleRetry(retryTask, remaining)
@@ -537,6 +541,11 @@ func isForwardableHeader(name string) bool {
"Upgrade", "Proxy-Authorization", "Upgrade", "Proxy-Authorization",
"Proxy-Connection", "Content-Length": "Proxy-Connection", "Content-Length":
return false return false
case "Content-Type":
// applyRequestHeaders sets Content-Type itself. The receiver
// already stored this inbound value as the event's
// ContentType, so forwarding it too would send it twice.
return false
default: default:
return true return true
} }
@@ -549,6 +558,10 @@ func isForwardableHeader(name string) bool {
// policy strips exactly that set on a hop that leaves the origin, // policy strips exactly that set on a hop that leaves the origin,
// so the forward set is decided here and only here — a header added // so the forward set is decided here and only here — a header added
// to it is covered off-origin without a second edit elsewhere. // to it is covered off-origin without a second edit elsewhere.
//
// Content-Type goes out once: a Content-Type configured on the target
// wins, otherwise the event's ContentType, otherwise none. The inbound
// Content-Type in the event's headers is never forwarded.
func applyRequestHeaders( func applyRequestHeaders(
req *http.Request, req *http.Request,
event *database.Event, event *database.Event,
@@ -569,10 +582,10 @@ func applyRequestHeaders(
req.Header.Set("User-Agent", "webhooker/1.0") req.Header.Set("User-Agent", "webhooker/1.0")
// Content-Type describes the body being sent rather than the // A Content-Type configured on the target describes the body
// sender, and the delivery path sets it from the event itself. // being sent rather than the sender. A 307/308 preserves the
// A 307/308 preserves the body across hosts, so stripping it // body across hosts, so stripping it would send that body
// would send that body untyped. // untyped.
delete(originScoped, "Content-Type") delete(originScoped, "Content-Type")
// User-Agent is overwritten just above, so an inbound one never // User-Agent is overwritten just above, so an inbound one never
+270 -9
View File
@@ -18,16 +18,17 @@ import (
// https://git.eeqj.de/sneak/webhooker/issues/107: a delivery failed // https://git.eeqj.de/sneak/webhooker/issues/107: a delivery failed
// with nothing in its event log to say why, and a retrying delivery // with nothing in its event log to say why, and a retrying delivery
// whose target was deleted, which used to keep sending and then never // whose target was deleted, which used to keep sending and then never
// terminalise. // terminalise. Section 4 is the same deleted-target gap for a pending
// delivery: https://git.eeqj.de/sneak/webhooker/issues/293.
// tUnknownType is a target type no build implements. It stands in for // tUnknownType is a target type no build implements. It stands in for
// a target whose type was written by a build that knew a type this one // a target whose type was written by a build that knew a type this one
// does not. // does not.
const tUnknownType = database.TargetType("pubsub") const tUnknownType = database.TargetType("pubsub")
// tSeedDeletedTarget creates a target, a retrying delivery against it // tSeedDeletedTarget creates a target, a delivery against it at the
// with one recorded failed attempt, and then deletes the target the // given status with one recorded failed attempt, and then deletes the
// way the source page does. // target the way the source page does.
// //
// It asserts the delete is soft, because that is the whole reason the // It asserts the delete is soft, because that is the whole reason the
// engine could not tell a deleted target from a target id that never // engine could not tell a deleted target from a target id that never
@@ -36,6 +37,7 @@ func tSeedDeletedTarget(
t *testing.T, t *testing.T,
s iSetup, s iSetup,
name, url string, name, url string,
status database.DeliveryStatus,
) string { ) string {
t.Helper() t.Helper()
@@ -51,8 +53,7 @@ func tSeedDeletedTarget(
) )
d := iSeedDelivery( d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID, t, s.WebhookDB, event.ID, targetID, status,
database.DeliveryStatusRetrying,
) )
iSeedFailedResult(t, s.WebhookDB, d.ID) iSeedFailedResult(t, s.WebhookDB, d.ID)
@@ -173,6 +174,7 @@ func TestRecoverSingleRetry_TargetDeleted(t *testing.T) {
deliveryID := tSeedDeletedTarget( deliveryID := tSeedDeletedTarget(
t, s, "gone-on-recovery", "http://example.com/hook", t, s, "gone-on-recovery", "http://example.com/hook",
database.DeliveryStatusRetrying,
) )
s.Engine.ExportRecoverWebhookDeliveries( s.Engine.ExportRecoverWebhookDeliveries(
@@ -210,6 +212,7 @@ func TestSweepSingleRetry_TargetDeleted(t *testing.T) {
deliveryID := tSeedDeletedTarget( deliveryID := tSeedDeletedTarget(
t, s, "gone-on-sweep", "http://example.com/hook", t, s, "gone-on-sweep", "http://example.com/hook",
database.DeliveryStatusRetrying,
) )
// Twice, because the bug was an error the sweep repeated every // Twice, because the bug was an error the sweep repeated every
@@ -278,11 +281,11 @@ func TestSweepSingleRetry_TargetNeverExisted(t *testing.T) {
) )
} }
// TestFailMissingTargetRetry_WritesNoTargetRow holds the new terminal // TestFailMissingTarget_WritesNoTargetRow holds the new terminal path
// path to the same rule as the existing one: no target row, and so no // to the same rule as the existing one: no target row, and so no
// plaintext target config, may be written into the per-webhook event // plaintext target config, may be written into the per-webhook event
// database. See https://git.eeqj.de/sneak/webhooker/issues/206. // database. See https://git.eeqj.de/sneak/webhooker/issues/206.
func TestFailMissingTargetRetry_WritesNoTargetRow( func TestFailMissingTarget_WritesNoTargetRow(
t *testing.T, t *testing.T,
) { ) {
t.Parallel() t.Parallel()
@@ -297,6 +300,7 @@ func TestFailMissingTargetRetry_WritesNoTargetRow(
deliveryID := tSeedDeletedTarget( deliveryID := tSeedDeletedTarget(
t, s, "credential-bearing", hookURL, t, s, "credential-bearing", hookURL,
database.DeliveryStatusRetrying,
) )
s.Engine.ExportSweepWebhookRetries( s.Engine.ExportSweepWebhookRetries(
@@ -529,3 +533,260 @@ func TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone(
assert.Zero(t, s.Engine.ExportInflightHeld()) assert.Zero(t, s.Engine.ExportInflightHeld())
} }
// --- 4. A pending delivery whose target is gone ---
func TestRecoverPending_TargetDeleted(t *testing.T) {
t.Parallel()
s := newISetup(t)
deliveryID := tSeedDeletedTarget(
t, s, "gone-while-pending", "http://example.com/hook",
database.DeliveryStatusPending,
)
s.Engine.ExportRecoverWebhookDeliveries(
context.Background(), s.WebhookID,
)
iAssertStatus(
t, s.WebhookDB, deliveryID,
database.DeliveryStatusFailed,
)
last := tLastResult(t, s, deliveryID, 2)
assert.False(t, last.Success)
assert.Equal(t, 2, last.AttemptNum)
assert.Contains(t, last.Error, "gone-while-pending")
assert.Contains(t, last.Error, "was deleted")
assert.Empty(t, fDrain(s.Engine),
"a delivery whose target is gone was sent",
)
assert.Zero(t, s.Engine.ExportInflightHeld(),
"the terminal path leaked its ownership reference",
)
}
// TestRecoverPending_TargetDeleted_LeavesAnOwnedDeliveryAlone: the
// terminal write takes ownership like every other recovery write, so a
// delivery the engine still holds is not failed underneath its worker.
func TestRecoverPending_TargetDeleted_LeavesAnOwnedDeliveryAlone(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
deliveryID := tSeedDeletedTarget(
t, s, "gone-but-owned", "http://example.com/hook",
database.DeliveryStatusPending,
)
require.True(t, s.Engine.ExportRetainDelivery(deliveryID))
s.Engine.ExportRecoverWebhookDeliveries(
context.Background(), s.WebhookID,
)
iAssertStatus(
t, s.WebhookDB, deliveryID,
database.DeliveryStatusPending,
)
assert.Len(t, iResults(t, s.WebhookDB, deliveryID), 1,
"a delivery the engine owns was failed underneath it",
)
}
// TestFailMissingTarget_LeavesASettledDeliveryAlone: the recovery paths
// read their batch before taking ownership, and a worker may send a
// delivery and let it go in between. The terminal write goes by the row
// as it is now, not as the batch read it.
func TestFailMissingTarget_LeavesASettledDeliveryAlone(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
deliveryID := tSeedDeletedTarget(
t, s, "gone-after-sending", "http://example.com/hook",
database.DeliveryStatusPending,
)
var batch database.Delivery
require.NoError(t, s.WebhookDB.First(
&batch, "id = ?", deliveryID,
).Error)
// A worker settles the delivery after the batch was read.
require.NoError(t, s.WebhookDB.Model(&database.Delivery{}).
Where("id = ?", deliveryID).
Update("status", database.DeliveryStatusDelivered).Error)
s.Engine.ExportFailMissingTarget(
s.WebhookDB, s.WebhookID, &batch,
)
iAssertStatus(
t, s.WebhookDB, deliveryID,
database.DeliveryStatusDelivered,
)
assert.Len(t, iResults(t, s.WebhookDB, deliveryID), 1,
"a delivery settled after the batch read was then failed",
)
assert.Zero(t, s.Engine.ExportInflightHeld())
}
// TestSweepPending_TargetDeleted sweeps twice over a batch that also
// holds a healthy stranded delivery. The one whose target is gone is
// failed once and then left alone; the healthy one is queued by the
// first sweep and not again by the second.
func TestSweepPending_TargetDeleted(t *testing.T) {
t.Parallel()
liveTargetID := uuid.New().String()
s := fSweepSetup(t, liveTargetID, "still-there")
deliveryID := tSeedDeletedTarget(
t, s, "gone-on-pending-sweep", "http://example.com/hook",
database.DeliveryStatusPending,
)
rAgePending(t, s.WebhookDB, deliveryID)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"target":"live"}`,
)
healthy := iSeedDelivery(
t, s.WebhookDB, event.ID, liveTargetID,
database.DeliveryStatusPending,
)
rAgePending(t, s.WebhookDB, healthy.ID)
ctx := context.Background()
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
tasks := fDrain(s.Engine)
require.Len(t, tasks, 1,
"the first sweep did not queue the healthy delivery",
)
assert.Equal(t, healthy.ID, tasks[0].DeliveryID)
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
assert.Empty(t, fDrain(s.Engine),
"the second sweep queued a delivery again",
)
iAssertStatus(
t, s.WebhookDB, deliveryID,
database.DeliveryStatusFailed,
)
last := tLastResult(t, s, deliveryID, 2)
assert.Contains(t, last.Error, "gone-on-pending-sweep")
assert.Contains(t, last.Error, "was deleted")
}
// TestSendRecoveredDeliveries_TargetMissingFromMap: the batch's target
// map is empty when its query failed, so every delivery in the batch is
// looked up on its own. A healthy one is sent to the target that lookup
// finds.
func TestSendRecoveredDeliveries_TargetMissingFromMap(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, "found-on-lookup",
database.TargetTypeLog, "", 0,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"map":"empty"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
s.Engine.ExportSendRecoveredDeliveries(
context.Background(), s.WebhookDB,
[]database.Delivery{d}, s.WebhookID,
map[string]database.Target{}, nil,
)
tasks := fDrain(s.Engine)
require.Len(t, tasks, 1,
"the healthy delivery was not queued exactly once",
)
assert.Equal(t, d.ID, tasks[0].DeliveryID)
assert.Equal(t, targetID, tasks[0].TargetID)
assert.Equal(t, database.TargetTypeLog, tasks[0].TargetType)
iAssertStatus(
t, s.WebhookDB, d.ID,
database.DeliveryStatusPending,
)
}
// TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone: a failed
// read of the main database is not a deleted target. Restart recovery
// holds every pending delivery of the webhook in one batch, so failing
// on this would fail all of them.
func TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, "healthy",
database.TargetTypeLog, "", 0,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"still":"pending"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
sqlDB, err := s.MainDB.DB()
require.NoError(t, err)
require.NoError(t, sqlDB.Close())
s.Engine.ExportRecoverPendingDeliveries(
context.Background(), s.WebhookDB, s.WebhookID,
)
iAssertStatus(
t, s.WebhookDB, d.ID,
database.DeliveryStatusPending,
)
assert.Empty(t, iResults(t, s.WebhookDB, d.ID),
"an unreadable main database produced a terminal "+
"failure row",
)
assert.Empty(t, fDrain(s.Engine))
assert.Zero(t, s.Engine.ExportInflightHeld())
}
+77
View File
@@ -300,3 +300,80 @@ func TestEntrypointCopyButtonIsProgressiveEnhancement(t *testing.T) {
"the page must render to completion, not abort partway", "the page must render to completion, not abort partway",
) )
} }
// maxRetriesHelp is the wording both target forms must carry. The
// delivery core makes max_retries attempts in total, not that many
// retries on top of a first try (a fresh delivery starts at attempt 1
// and target_http gives up once the attempt number reaches
// max_retries), and 0 is special-cased to a single fire-and-forget
// attempt with no circuit breaker.
const maxRetriesHelp = "This is the total number of delivery attempts, " +
"not retries on top of the first: a value of 3 makes three attempts " +
"in all. 0 means a single attempt with no retries and no circuit " +
"breaker."
// TestTargetFormMaxRetriesCopyMatchesBehaviour pins the max_retries
// help text on both the create form (the add-target form on the webhook
// detail page) and the edit form, so the copy cannot drift back to
// calling the number a retry count.
func TestTargetFormMaxRetriesCopyMatchesBehaviour(t *testing.T) {
t.Parallel()
var h *handlers.Handlers
var sess *session.Session
app := newTestApp(t, &h, &sess)
app.RequireStart()
t.Cleanup(app.RequireStop)
webhook := &database.Webhook{Name: "wh", RetentionDays: 14}
webhook.ID = testWebhookID
entrypoint := database.Entrypoint{Path: "abc123"}
entrypoint.ID = "ep-1"
createBody := renderPage(
t, h, sess, "source_detail.html", map[string]any{
dataKeyWebhook: webhook,
"Entrypoints": handlers.NewEntrypointViews(
[]database.Entrypoint{entrypoint},
),
"Targets": delivery.NewTargetViews(nil),
"Events": []database.Event{},
"BaseURL": "https://hooks.example.com",
},
)
assert.Contains(
t, createBody, maxRetriesHelp,
"the add-target form must explain max_retries as total attempts",
)
// A slack target exercises the same max_retries field while needing
// only Config.URL from the edit template, so the test data stays
// minimal. The Target key mirrors the field names the template reads
// off the handler's view value.
editBody := renderPage(
t, h, sess, "target_edit.html", map[string]any{
dataKeyWebhook: webhook,
"Target": map[string]any{
"ID": "tg-1",
"Name": "t",
"Type": "slack",
"Active": true,
"MaxRetries": 3,
"Config": map[string]any{
"URL": "https://hooks.slack.com/services/x",
},
},
dataKeyError: "",
},
)
assert.Contains(
t, editBody, maxRetriesHelp,
"the target edit form must explain max_retries as total attempts",
)
}
+2 -3
View File
@@ -380,9 +380,8 @@ func csrfTookStrictPath(
// TestCSRF_ForwardedProtoSpellingsTakeStrictPath runs the header // TestCSRF_ForwardedProtoSpellingsTakeStrictPath runs the header
// spellings a real proxy emits through the middleware. The environment // spellings a real proxy emits through the middleware. The environment
// is dev -- the DEFAULT when WEBHOOKER_ENVIRONMENT is unset -- to pin // is set to dev -- the permissive setting -- to pin that the routing is
// that the routing is a per-request transport decision and owes // a per-request transport decision and owes nothing to configuration.
// nothing to configuration.
func TestCSRF_ForwardedProtoSpellingsTakeStrictPath(t *testing.T) { func TestCSRF_ForwardedProtoSpellingsTakeStrictPath(t *testing.T) {
t.Parallel() t.Parallel()
+7
View File
@@ -133,6 +133,11 @@ func (w *recoverResponseWriter) Unwrap() http.ResponseWriter {
// what the access log records and the metrics count, and outside the // what the access log records and the metrics count, and outside the
// sentryhttp handler, whose Repanic option depends on something // sentryhttp handler, whose Repanic option depends on something
// further out recovering what it re-raises. // further out recovering what it re-raises.
//
// Unlike http.Error on its own, it deletes any Set-Cookie the handler
// set before panicking, because a request that failed must not hand
// the client a credential; every other header is left to http.Error.
// See https://git.eeqj.de/sneak/webhooker/issues/193.
func (s *Middleware) Recoverer() func(http.Handler) http.Handler { func (s *Middleware) Recoverer() func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler { return func(next http.Handler) http.Handler {
return http.HandlerFunc(func( return http.HandlerFunc(func(
@@ -164,6 +169,8 @@ func (s *Middleware) Recoverer() func(http.Handler) http.Handler {
return return
} }
rw.Header().Del("Set-Cookie")
http.Error( http.Error(
rw, rw,
http.StatusText( http.StatusText(
+31 -2
View File
@@ -304,16 +304,44 @@ func TestRecovererRepanicsErrAbortHandler(t *testing.T) {
) )
} }
// TestRecovererDropsSetCookieFromTheRecovered500 covers a handler that
// sets a cookie and a redirect target and then panics before sending
// anything. A request that failed must not hand the client a
// credential, so the 500 carries no cookie; Location is left alone.
func TestRecovererDropsSetCookieFromTheRecovered500(t *testing.T) {
t.Parallel()
probe := newRecovererProbe(
t, false,
func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Set-Cookie", "session=x")
w.Header().Set("Location", "/after")
panic(panicMarker)
},
)
resp, err := probe.get(t)
require.NoError(t, err)
require.NoError(t, resp.Body.Close())
assert.Equal(t, http.StatusInternalServerError, resp.StatusCode)
assert.Empty(t, resp.Cookies())
assert.Equal(t, "/after", resp.Header.Get("Location"))
}
// TestRecovererKeepsAnAlreadyCommittedResponse covers a handler that // TestRecovererKeepsAnAlreadyCommittedResponse covers a handler that
// panics after sending its status. The bytes are already on the wire, // panics after sending its status. The bytes are already on the wire,
// so a second WriteHeader would change nothing the client sees and // cookie included, so a second WriteHeader would change nothing the
// would draw net/http's "superfluous response.WriteHeader" report. // client sees and would draw net/http's "superfluous
// response.WriteHeader" report.
func TestRecovererKeepsAnAlreadyCommittedResponse(t *testing.T) { func TestRecovererKeepsAnAlreadyCommittedResponse(t *testing.T) {
t.Parallel() t.Parallel()
probe := newRecovererProbe( probe := newRecovererProbe(
t, false, t, false,
func(w http.ResponseWriter, _ *http.Request) { func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Set-Cookie", "session=x")
w.WriteHeader(committedStatus) w.WriteHeader(committedStatus)
_, _ = w.Write([]byte("partial")) _, _ = w.Write([]byte("partial"))
@@ -331,6 +359,7 @@ func TestRecovererKeepsAnAlreadyCommittedResponse(t *testing.T) {
assert.Equal(t, committedStatus, resp.StatusCode) assert.Equal(t, committedStatus, resp.StatusCode)
assert.Equal(t, "partial", string(body)) assert.Equal(t, "partial", string(body))
assert.Len(t, resp.Cookies(), 1)
record := probe.panicRecord(t) record := probe.panicRecord(t)
assert.Equal(t, panicMarker, record["panic"]) assert.Equal(t, panicMarker, record["panic"])
+1 -1
View File
@@ -5,7 +5,7 @@
// several packages, by hand, and the answers disagreed. The session // several packages, by hand, and the answers disagreed. The session
// cookie's Secure attribute was decided at startup from the configured // cookie's Secure attribute was decided at startup from the configured
// environment while the CSRF cookie's was decided per-request, so a // environment while the CSRF cookie's was decided per-request, so a
// deployment behind a TLS proxy in the default environment emitted one // deployment behind a TLS proxy in the dev environment emitted one
// Secure cookie and one non-Secure cookie on the same response. // Secure cookie and one non-Secure cookie on the same response.
// Everything kept working, which is exactly why nobody noticed. // Everything kept working, which is exactly why nobody noticed.
// //
+17 -3
View File
@@ -92,11 +92,25 @@ func (s *Server) setupGlobalMiddleware() {
func (s *Server) setupRoutes() { func (s *Server) setupRoutes() {
s.router.Get("/", s.h.HandleIndex()) s.router.Get("/", s.h.HandleIndex())
s.router.Mount( // Static assets answer GET and HEAD only. chi's default 405
"/s", // carries no Allow header, so this group supplies its own.
http.StripPrefix("/s", http.FileServer(http.FS(static.Static))), staticFiles := http.StripPrefix(
"/s", http.FileServer(http.FS(static.Static)),
) )
s.router.Route("/s", func(r chi.Router) {
r.MethodNotAllowed(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Allow", "GET, HEAD")
http.Error(
w,
"Method Not Allowed",
http.StatusMethodNotAllowed,
)
})
r.Method(http.MethodGet, "/*", staticFiles)
r.Method(http.MethodHead, "/*", staticFiles)
})
s.router.Route("/api/v1", func(_ chi.Router) { s.router.Route("/api/v1", func(_ chi.Router) {
// API routes will be added here. // API routes will be added here.
}) })
+39 -16
View File
@@ -396,13 +396,15 @@ func (e *testEnv) storedHash(t *testing.T, username string) string {
// --- /s static group --- // --- /s static group ---
// TestStaticServesEveryMethod pins what the static mount actually // TestStaticServesOnlyGetAndHead pins the methods the static group
// answers. chi's Mount registers the handler for all methods and // answers: GET and HEAD are served the asset, and the other methods
// http.FileServer only special-cases HEAD (by suppressing the body), // chi routes (POST, PUT, DELETE and the rest) are refused with 405
// so a POST or a DELETE to an asset is served the file rather than // and an Allow header naming those two. A method chi does not route,
// refused. The README documents this; the test is what keeps the two // such as PROPFIND, is refused with 405 by the top-level router
// from drifting. // before it reaches the static group, so it gets no Allow header.
func TestStaticServesEveryMethod(t *testing.T) { // The README documents this; the test is what keeps the two from
// drifting.
func TestStaticServesOnlyGetAndHead(t *testing.T) {
t.Parallel() t.Parallel()
env := newTestEnv(t) env := newTestEnv(t)
@@ -417,6 +419,7 @@ func TestStaticServesEveryMethod(t *testing.T) {
http.MethodPost, http.MethodPost,
http.MethodPut, http.MethodPut,
http.MethodDelete, http.MethodDelete,
"PROPFIND",
} { } {
t.Run(method, func(t *testing.T) { t.Run(method, func(t *testing.T) {
t.Parallel() t.Parallel()
@@ -428,18 +431,38 @@ func TestStaticServesEveryMethod(t *testing.T) {
w := httptest.NewRecorder() w := httptest.NewRecorder()
env.router.ServeHTTP(w, req) env.router.ServeHTTP(w, req)
assert.Equal(t, http.StatusOK, w.Code, switch method {
"static mount answers every method") case http.MethodGet:
assert.Equal(t, http.StatusOK, w.Code)
if method == http.MethodHead { assert.Equal(t, body, w.Body.Bytes(),
"the asset itself is returned")
case http.MethodHead:
assert.Equal(t, http.StatusOK, w.Code)
assert.Empty(t, w.Body.Bytes(), assert.Empty(t, w.Body.Bytes(),
"HEAD must not carry a body") "HEAD must not carry a body")
case "PROPFIND":
return assert.Equal(
t, http.StatusMethodNotAllowed, w.Code,
)
assert.Empty(t, w.Header().Get("Allow"),
"chi refuses a method it does not route "+
"before the static group runs")
assert.NotContains(
t, w.Body.String(), string(body),
"a refused method must not get the asset",
)
default:
assert.Equal(
t, http.StatusMethodNotAllowed, w.Code,
)
assert.Equal(
t, "GET, HEAD", w.Header().Get("Allow"),
)
assert.NotContains(
t, w.Body.String(), string(body),
"a refused method must not get the asset",
)
} }
assert.Equal(t, body, w.Body.Bytes(),
"the asset itself is returned")
}) })
} }
} }
+2 -2
View File
@@ -146,8 +146,8 @@ func newStore(key []byte) *sessions.CookieStore {
// //
// This is decided per-request, not once at startup. Deciding it at // This is decided per-request, not once at startup. Deciding it at
// startup from the configured environment is what this replaces, and // startup from the configured environment is what this replaces, and
// it got the DEFAULT posture wrong: "dev" is the environment when // it got the DEFAULT posture wrong: "dev" was then the environment when
// WEBHOOKER_ENVIRONMENT is unset, so a deployment terminating TLS at a // WEBHOOKER_ENVIRONMENT was unset, so a deployment terminating TLS at a
// proxy without also setting the environment emitted the // proxy without also setting the environment emitted the
// authentication cookie with no Secure attribute -- silently, and on // authentication cookie with no Secure attribute -- silently, and on
// the same response as a CSRF cookie that did have one. // the same response as a CSRF cookie that did have one.
+2 -2
View File
@@ -990,8 +990,8 @@ func sessionCookieFrom(
// TestSave_SecureFollowsRequestTransport is the regression test for // TestSave_SecureFollowsRequestTransport is the regression test for
// the defect this replaces: Secure was fixed at startup from the // the defect this replaces: Secure was fixed at startup from the
// configured environment, and "dev" is the environment when // configured environment, and "dev" was then the environment when
// WEBHOOKER_ENVIRONMENT is unset. A deployment behind a TLS proxy in // WEBHOOKER_ENVIRONMENT was unset. A deployment behind a TLS proxy in
// that DEFAULT posture shipped the authentication cookie with no // that DEFAULT posture shipped the authentication cookie with no
// Secure attribute and said nothing about it. // Secure attribute and said nothing about it.
// //
+6 -3
View File
@@ -120,9 +120,12 @@
<label class="text-sm text-gray-700">Timeout (seconds, blank = default):</label> <label class="text-sm text-gray-700">Timeout (seconds, blank = default):</label>
<input type="number" name="timeout" min="0" max="300" :disabled="targetType !== 'http'" class="input text-sm w-24"> <input type="number" name="timeout" min="0" max="300" :disabled="targetType !== 'http'" class="input text-sm w-24">
</div> </div>
<div x-show="targetType === 'http'" class="flex gap-2 items-center"> <div x-show="targetType === 'http'">
<label class="text-sm text-gray-700">Max retries (0 = fire-and-forget):</label> <div class="flex gap-2 items-center">
<input type="number" name="max_retries" value="0" min="0" max="20" class="input text-sm w-24"> <label class="text-sm text-gray-700">Max retries:</label>
<input type="number" name="max_retries" value="0" min="0" max="20" class="input text-sm w-24">
</div>
<p class="text-xs text-gray-500 mt-1">This is the total number of delivery attempts, not retries on top of the first: a value of 3 makes three attempts in all. 0 means a single attempt with no retries and no circuit breaker.</p>
</div> </div>
<div x-show="targetType === 'slack'"> <div x-show="targetType === 'slack'">
<input type="url" name="url" placeholder="https://hooks.slack.com/services/..." :disabled="targetType !== 'slack'" class="input text-sm"> <input type="url" name="url" placeholder="https://hooks.slack.com/services/..." :disabled="targetType !== 'slack'" class="input text-sm">
+1 -1
View File
@@ -69,7 +69,7 @@
<div class="form-group"> <div class="form-group">
<label for="max_retries" class="label">Max retries</label> <label for="max_retries" class="label">Max retries</label>
<input type="number" id="max_retries" name="max_retries" value="{{.Target.MaxRetries}}" min="0" max="20" class="input"> <input type="number" id="max_retries" name="max_retries" value="{{.Target.MaxRetries}}" min="0" max="20" class="input">
<p class="text-xs text-gray-500 mt-1">0 is fire-and-forget: one attempt, no circuit breaker.</p> <p class="text-xs text-gray-500 mt-1">This is the total number of delivery attempts, not retries on top of the first: a value of 3 makes three attempts in all. 0 means a single attempt with no retries and no circuit breaker.</p>
</div> </div>
{{end}} {{end}}