From a201ecd13296b718e86129bd5026842c7d87d33f Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 11:46:50 +0200 Subject: [PATCH] Deploy: main into prod (#343) Brings `prod`, which upaas deploys, up to `main` at `9cf9cdd`, the merge of https://git.eeqj.de/sneak/webhooker/pulls/321. `prod` was cut from `main` at `251cb3d` (1.0.0b1). What it deploys is everything listed in https://git.eeqj.de/sneak/webhooker/pulls/321. For running it: - With `WEBHOOKER_ENVIRONMENT` unset, the instance runs as `prod` and sends no `Access-Control-Allow-Origin: *`. - Each event database gains its new indexes the first time it is opened after the upgrade. - `webhooker_delivery_retries_total` no longer counts a circuit breaker holding back a delivery that is already `retrying`. Not in this PR yet: https://git.eeqj.de/sneak/webhooker/issues/340, in which the container sets its own data directory owner and mode before start. It is in progress on `next`. Once it reaches `main`, this PR carries it, because the PR follows `main`. Model: opus-5-5 Co-authored-by: Jeffrey Paul <1+sneak@noreply.example.org> Reviewed-on: https://git.eeqj.de/sneak/webhooker/pulls/343 Co-authored-by: clawbot <35+clawbot@noreply.example.org> --- README.md | 266 ++++++++++++----- internal/config/config.go | 35 ++- internal/config/config_test.go | 30 +- internal/database/database.go | 6 +- internal/database/event_tier_indexes_test.go | 166 +++++++++++ internal/database/model_delivery.go | 15 +- internal/database/model_delivery_result.go | 13 +- internal/database/model_event.go | 18 ++ internal/database/webhook_db_manager.go | 71 +++-- internal/database/webhook_db_manager_test.go | 52 ++++ internal/datadir/lock.go | 10 +- internal/delivery/circuit_breaker.go | 12 +- internal/delivery/circuit_breaker_test.go | 8 +- internal/delivery/engine.go | 122 ++++++-- internal/delivery/engine_lifecycle_test.go | 87 ++++++ internal/delivery/engine_test.go | 175 +++++++++++ internal/delivery/export_test.go | 46 +++ internal/delivery/inflight_test.go | 67 +++++ internal/delivery/metrics_test.go | 7 +- internal/delivery/redirect_test.go | 21 +- internal/delivery/ssrf.go | 17 +- internal/delivery/ssrf_allowlist_test.go | 35 +++ internal/delivery/target_database.go | 18 ++ .../delivery/target_database_evict_test.go | 44 +++ internal/delivery/target_headers_test.go | 3 +- internal/delivery/target_http.go | 29 +- internal/delivery/terminal_state_test.go | 279 +++++++++++++++++- internal/handlers/ui_copy_test.go | 77 +++++ internal/middleware/csrf_test.go | 5 +- internal/middleware/recoverer.go | 7 + internal/middleware/recoverer_test.go | 33 ++- internal/reqtls/reqtls.go | 2 +- internal/session/session.go | 4 +- internal/session/session_test.go | 4 +- templates/source_detail.html | 9 +- templates/target_edit.html | 2 +- 36 files changed, 1574 insertions(+), 221 deletions(-) create mode 100644 internal/database/event_tier_indexes_test.go diff --git a/README.md b/README.md index c66e5ca..5729606 100644 --- a/README.md +++ b/README.md @@ -7,6 +7,13 @@ services, durably stores them, and delivers them to configured targets with retry support, logging, and observability. Category: infrastructure / 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 ### Prerequisites @@ -37,9 +44,9 @@ make bootstrap # Run all checks (test, lint, format check) make check -# Run in development mode. DATA_DIR defaults to /var/lib/webhooker in -# every environment, so set it (in .env or the shell) to a writable -# directory when running from a clone. +# Run the server from the clone. DATA_DIR defaults to +# /var/lib/webhooker in every environment, so set it (in .env or the +# shell) to a writable directory. DATA_DIR=./data make dev # 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. 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` | | -------- | ----------------------- | ---------------- | @@ -127,7 +135,7 @@ TTY detection, and security headers are always applied. | Variable | Description | Default | | ----------------------- | ----------------------------------- | -------- | -| `WEBHOOKER_ENVIRONMENT` | `dev` or `prod` | `dev` | +| `WEBHOOKER_ENVIRONMENT` | `dev` or `prod` | `prod` | | `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`) | | `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 webhooker probe the network it sits in. +Besides the private and reserved ranges, the default blocklist refuses +public cloud metadata addresses: currently only `168.63.129.16`, Azure's +WireServer, which serves an Azure VM its credentials. Because it is a +public address, listing it in `ALLOWED_EGRESS_CIDRS` reopens it. + That default is also inconvenient for the thing webhooker is mostly for: taking a public webhook and forwarding it to something on your own network. A container on the same Docker network, a box on `10.x`, a @@ -187,15 +200,16 @@ Two things this setting cannot do: the list is always an allowlist; an empty list (the default) means every private and reserved range stays refused. Note that `0.0.0.0/0` gets you most of the way there anyway, per above. -- **It cannot open link-local, or a cloud metadata endpoint that - discloses credentials or user data.** An address is on the list below - when both of these hold: the provider fixes it, so it cannot collide - with anything you run; and reaching it hands out credentials, user - data or bootstrap material. Those stay blocked no matter what you - list, including when you list them outright or list a supernet such - as `0.0.0.0/0`, `::/0`, `fd00::/8` or `100.64.0.0/10`. Treat this as - best effort rather than a guarantee — it is a hand-maintained list - and the caveat below the table applies: +- **It cannot open link-local, or a cloud metadata endpoint at a + non-public address that discloses credentials or user data.** An + address is on the list below when it is not a public address and both + of these hold: the provider fixes it, so it cannot collide with + anything you run; and reaching it hands out credentials, user data or + bootstrap material. Those stay blocked no matter what you list, + including when you list them outright or list a supernet such as + `0.0.0.0/0`, `::/0`, `fd00::/8` or `100.64.0.0/10`. Treat this as best + effort rather than a guarantee — it is a hand-maintained list and the + caveat below the table applies: | 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 routable metadata address is not listed here, because nothing on this list can be reopened and blocking one that way would leave you no - escape hatch at all. + 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 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 address, which restores per-client buckets. webhooker logs a warning -at startup whenever `TRUSTED_PROXIES` is empty, in every environment — -not only when `WEBHOOKER_ENVIRONMENT=prod`, because that variable -defaults to `dev` and an operator who never set it is precisely the -one at risk. The warning is informational when nothing proxies to the +at startup whenever `TRUSTED_PROXIES` is empty, in every environment, +because behind a proxy every client shares one bucket in `dev` and +`prod` alike. The warning is informational when nothing proxies to the process: with no proxy in front, the peer address is the client's own and the buckets are already per-client. See [Rate Limiting](#rate-limiting) for what each limit shares. @@ -631,7 +645,6 @@ decision: docker run -d \ -p 127.0.0.1:8080:8080 \ -v /path/to/data:/var/lib/webhooker \ - -e WEBHOOKER_ENVIRONMENT=prod \ -e BIND_ADDRESS=0.0.0.0 \ 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 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 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 no TLS at all, and the proxy in front of it changes nothing about that. -2. **Set `WEBHOOKER_ENVIRONMENT=prod`, and make sure the proxy sends - `X-Forwarded-Proto`.** These are two requirements, not one. The - environment setting decides CORS and nothing else: the default +2. **Make sure the environment is not `dev` (leave + `WEBHOOKER_ENVIRONMENT` unset or set it to `prod`), and make sure + 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: *` (without credentials), which a server-rendered production - deployment has no use for. Cookie `Secure` and the strict - Origin/Referer mode are **not** tied to it — they are decided per - request from the transport, which behind a proxy means the - `X-Forwarded-Proto` header. The block below sets it; without it - every request is read as plaintext and cookies ship without - `Secure`. See [Configuration](#configuration). + deployment has no use for, and `prod` — the default — disables it. + Cookie `Secure` and the strict Origin/Referer mode are **not** tied + to it — they are decided per request from the transport, which + behind a proxy means the `X-Forwarded-Proto` header. The block below + sets it; without it every request is read as plaintext and cookies + ship without `Secure`. See [Configuration](#configuration). 3. **Set `TRUSTED_PROXIES` to the proxy's address.** Unset, every rate limiter keys on the connecting peer, which behind a proxy is the proxy on every request: all clients collapse into one global bucket @@ -847,7 +921,6 @@ sent — `$scheme` above does. With that block, webhooker's environment is: ```sh -WEBHOOKER_ENVIRONMENT=prod BIND_ADDRESS=127.0.0.1 # the default; stated here to be explicit 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 reason to separate the two — copy the directory and you have them. -A clean shutdown closes `webhooker.db` and every `events-*.db`, which -checkpoints and removes their sidecars; a killed or crashed instance -leaves them, and they must be carried with the `.db`. **Archive -databases are different**: their handle is not closed at shutdown, so -`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. +A clean shutdown closes every database, which checkpoints and removes +its sidecars; a killed or crashed instance leaves them, and they must be +carried with the `.db`. An archive the service has not opened since a +crash keeps that crash's sidecars, even across a later clean stop. Configuration is **not** in `DATA_DIR` — it comes from the environment 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 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 -20 KB `.db` with no sidecars about a minute after its last write. -Shutdown is **not** on that list: the archive handle is not closed when -the service stops. So either move `archive-{uuid}.db` together with any -`-wal`/`-shm` beside it, or wait until there are none. +20 KB `.db` with no sidecars about a minute after its last write. A +clean stop closes it too. So either move `archive-{uuid}.db` together +with any `-wal`/`-shm` beside it, or wait until there are none. ### 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 discards every transaction it still holds. An `.backup` set will not contain any: it writes a single consolidated file per database. A - stop-and-copy set has none for `webhooker.db` or the `events-*.db`, - because a clean stop closes those and checkpoints their sidecars - away — but it will normally have them for `archive-*.db`, whose - handle stays open across shutdown, and those carry the archive's - rows. A copy salvaged from a crashed instance has them for - everything, and needs all of them. + stop-and-copy set normally has none, because a clean stop closes + every database and checkpoints its sidecars away; the exception is an + archive not opened since a crash. 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` 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 receiver verifies nothing about an inbound request. The UUID in an -entrypoint's URL is its credential: anyone who holds that URL can -submit events to it, and the receiver checks nothing else about the -sender. Treat an entrypoint URL the way you would treat an API token. +**The entrypoint UUID is the credential, and it is the only one.** +webhooker mints a version 4 UUID per entrypoint and serves it at +`/webhook/{uuid}`. Possession of that URL is the authentication: +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 -entrypoint (or deactivate it, which answers `410`) and create a new -one, then point the sender at the new URL. +There is no shared secret, no HMAC signature, no bearer token and no +second factor on the receiver, and none will be added. This was +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 @@ -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 shims that call them; `build`, `run`, `dev`, `deps`, `clean`, `css` and `version` are inline commands with no script behind them, though -`build` and `version` both take their value from `script/version`. We -provide: +`build` and `version` both take their value from `script/version`. + +`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/setup` — make a fresh clone ready for development @@ -1476,7 +1572,7 @@ events should be forwarded. | `type` | TargetType | One of: `http`, `slack`, `database`, `log` | | `active` | boolean | Whether deliveries are enabled (default: true) | | `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 | **Relations:** Belongs to Webhook. Has many Deliveries. @@ -1484,12 +1580,12 @@ events should be forwarded. **Target types:** - **`http`** — Forward the event as an HTTP POST to a configured URL. - Behavior depends on `max_retries`: when `max_retries` is 0 (the - default), the target operates in fire-and-forget mode — a single - attempt with no retries and no circuit breaker. When `max_retries` is - greater than 0, failed deliveries are retried with exponential backoff - up to `max_retries` attempts, protected by a per-target circuit - breaker. + `max_retries` is the total number of delivery attempts, not retries on + top of the first: when `max_retries` is 0 (the default), the target + operates in fire-and-forget mode, a single attempt with no retries and + no circuit breaker; a value of N makes up to N attempts in all, + retrying failed deliveries with exponential backoff and protecting them + with a per-target circuit breaker. - **`slack`** — Post the event as a formatted message to a Slack-compatible incoming webhook URL (`webhookUrl` in `config`). It 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. +#### 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 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. | | **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:** @@ -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 delivery as `retrying` and schedules a retry timer for after the 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 @@ -1984,7 +2104,7 @@ arriving and being stored, they are just not getting anywhere. | Metric | Type | Meaning | | ------ | ---- | ------- | | `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_failed_total` | counter | Deliveries that failed terminally and will not be retried | | `webhooker_delivery_retries_total` | counter | Deliveries put back into `retrying` | @@ -2867,6 +2987,10 @@ check, see [The login endpoint](#the-login-endpoint). ### 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 encrypted cookies. Sessions are configured with HttpOnly, SameSite 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 - **The entrypoint URL is the receiver's only credential.** Nothing 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)) - **SSRF prevention** for HTTP delivery targets: private/reserved IP 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 `server.ShutdownTimeout` (**3 seconds**), then a Sentry flush if `SENTRY_DSN` is set -4. `delivery.Engine` +4. `delivery.Engine` — waits for its workers, then closes the archive + databases 5. `healthcheck` 6. `WebhookDBManager` 7. the database close diff --git a/internal/config/config.go b/internal/config/config.go index 6a7a628..1b9a4e7 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -192,9 +192,10 @@ type Config struct { // alwaysBlockedNetworks stays blocked no matter what is listed // here. That set is link-local plus the cloud metadata // endpoints outside it that disclose credentials or user data - // at a provider-fixed address; it is not exhaustive of every - // cloud's metadata address. See alwaysBlockedNetworks for the - // authoritative list and the criterion it is built from. + // at a provider-fixed, non-public address; it is not + // exhaustive of every cloud's metadata address. See + // alwaysBlockedNetworks for the authoritative list and the + // criterion it is built from. AllowedEgressCIDRs []netip.Prefix params *ConfigParams @@ -585,12 +586,14 @@ func resolveMetricsAuth() (string, string, error) { ) } -// resolveEnvironment reads WEBHOOKER_ENVIRONMENT, defaulting to -// dev, and rejects unrecognised values. +// resolveEnvironment reads WEBHOOKER_ENVIRONMENT, defaulting to prod +// 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) { environment := os.Getenv("WEBHOOKER_ENVIRONMENT") if environment == "" { - environment = EnvironmentDev + environment = EnvironmentProd } if environment != EnvironmentDev && @@ -744,12 +747,14 @@ func (c *Config) warnEgressAllowlist(log *slog.Logger) { log.Warn( "ALLOWED_EGRESS_CIDRS lets delivery targets reach these "+ - "otherwise-blocked private/reserved networks. Anyone "+ - "who can create a delivery target can now make this "+ - "process issue requests into them, and read back the "+ - "response. Link-local and the known cloud instance "+ - "metadata endpoints outside it stay blocked "+ - "regardless of what is listed here.", + "otherwise-blocked networks. Anyone who can create a "+ + "delivery target can now make this process issue "+ + "requests into them, and read back the response. Only "+ + "the addresses the README lists as blocked "+ + "unconditionally stay blocked regardless of what is "+ + "listed here; a public cloud metadata address such as "+ + "168.63.129.16 is reachable once it, or a block "+ + "covering it, is listed.", "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 // service-wide ceilings. // -// The warning is deliberately not gated on WEBHOOKER_ENVIRONMENT. That -// variable defaults to dev, so gating on it would silence the warning -// for exactly the operator who forgot to configure the deployment — -// the case it exists to catch. +// The warning is deliberately not gated on WEBHOOKER_ENVIRONMENT: +// behind a proxy every client shares one bucket in dev and prod alike. // // The default of trusting nobody is deliberate — trusting forwarded // headers from arbitrary peers lets any client choose its own bucket — diff --git a/internal/config/config_test.go b/internal/config/config_test.go index f38f7fd..a7cc76e 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -44,9 +44,9 @@ func TestEnvironmentConfig(t *testing.T) { isProd bool }{ { - name: "default is dev", - isDev: true, - isProd: false, + name: "default is prod", + isDev: false, + isProd: true, }, { name: "explicit dev", @@ -834,12 +834,13 @@ func TestEgressAllowlistWarning(t *testing.T) { // to be able to read back which networks are open. assert.Contains(t, logged, "10.0.0.0/8") assert.Contains(t, logged, "127.0.0.0/8") - // What stays shut. Asserted on the clause naming the - // wider set rather than on "Link-local" alone, so the - // string cannot narrow back to link-local only while - // the always-blocked set covers ULA, CGNAT and two - // public metadata addresses as well. - assert.Contains(t, logged, "metadata endpoints outside it") + // What stays shut is the whole unconditional set, not + // link-local alone; a public metadata address is not in + // it, so a listed block covering it opens it. + assert.Contains(t, logged, "blocked unconditionally") + assert.Contains(t, logged, "168.63.129.16 is reachable") + // The listed blocks need not be private or reserved. + assert.NotContains(t, logged, "private/reserved") }) } } @@ -848,10 +849,9 @@ func TestEgressAllowlistWarning(t *testing.T) { // tells an operator a deployment behind a reverse proxy shares one // rate-limit bucket between every client, which turns the receiver // limits into service-wide ceilings and collapses login failure -// counting. It must fire whenever TRUSTED_PROXIES is empty, -// in any environment: WEBHOOKER_ENVIRONMENT defaults to dev, so gating -// on it would silence the warning for exactly the operator who never -// configured the deployment. It stays quiet once proxies are named. +// counting. It must fire whenever TRUSTED_PROXIES is empty, in any +// environment, because behind a proxy every client shares one bucket +// in dev and prod alike. It stays quiet once proxies are named. func TestSharedRateLimitBucketWarning(t *testing.T) { tests := []struct { name string @@ -871,10 +871,6 @@ func TestSharedRateLimitBucketWarning(t *testing.T) { 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", environment: config.EnvironmentDev, expectWarning: true, diff --git a/internal/database/database.go b/internal/database/database.go index f2880e8..ba28bae 100644 --- a/internal/database/database.go +++ b/internal/database/database.go @@ -17,12 +17,12 @@ import ( "gorm.io/gorm" "sneak.berlin/go/webhooker/internal/banner" "sneak.berlin/go/webhooker/internal/config" + "sneak.berlin/go/webhooker/internal/datadir" "sneak.berlin/go/webhooker/internal/gormlog" "sneak.berlin/go/webhooker/internal/logger" ) const ( - dataDirPerm = 0750 randomPasswordLen = 16 sessionKeyLen = 32 ) @@ -185,7 +185,9 @@ func (d *Database) connect() error { // caller's decision. func (d *Database) connectTo(dataDir string) error { // 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 { return fmt.Errorf( "creating data directory %s: %w", diff --git a/internal/database/event_tier_indexes_test.go b/internal/database/event_tier_indexes_test.go new file mode 100644 index 0000000..d25dca5 --- /dev/null +++ b/internal/database/event_tier_indexes_test.go @@ -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()) + } +} diff --git a/internal/database/model_delivery.go b/internal/database/model_delivery.go index 71f6b6d..0ebd2bf 100644 --- a/internal/database/model_delivery.go +++ b/internal/database/model_delivery.go @@ -1,5 +1,7 @@ package database +import "gorm.io/gorm" + // DeliveryStatus represents the status of a delivery type DeliveryStatus string @@ -29,12 +31,19 @@ func (s DeliveryStatus) Terminal() bool { } // Delivery represents a delivery attempt for an event to a target +// +//nolint:lll // a struct tag cannot wrap type Delivery struct { BaseModel - EventID string `gorm:"type:uuid;not null" json:"eventId"` - TargetID string `gorm:"type:uuid;not null" json:"targetId"` - Status DeliveryStatus `gorm:"not null;default:'pending'" json:"status"` + 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"` + 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 Event Event `json:"event,omitzero"` diff --git a/internal/database/model_delivery_result.go b/internal/database/model_delivery_result.go index 56c9cc2..7bf1b88 100644 --- a/internal/database/model_delivery_result.go +++ b/internal/database/model_delivery_result.go @@ -1,10 +1,21 @@ package database +import "gorm.io/gorm" + // DeliveryResult represents the result of a delivery attempt +// +//nolint:lll // a struct tag cannot wrap type DeliveryResult struct { 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"` Success bool `json:"success"` StatusCode int `json:"statusCode,omitempty"` diff --git a/internal/database/model_event.go b/internal/database/model_event.go index dafc235..347eaac 100644 --- a/internal/database/model_event.go +++ b/internal/database/model_event.go @@ -1,9 +1,27 @@ package database +import ( + "time" + + "gorm.io/gorm" +) + // Event represents a captured webhook event +// +//nolint:lll // a struct tag cannot wrap type Event struct { 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"` EntrypointID string `gorm:"type:uuid;not null" json:"entrypointId"` diff --git a/internal/database/webhook_db_manager.go b/internal/database/webhook_db_manager.go index 81ca427..a628b56 100644 --- a/internal/database/webhook_db_manager.go +++ b/internal/database/webhook_db_manager.go @@ -13,6 +13,7 @@ import ( "gorm.io/driver/sqlite" "gorm.io/gorm" "sneak.berlin/go/webhooker/internal/config" + "sneak.berlin/go/webhooker/internal/datadir" "sneak.berlin/go/webhooker/internal/gormlog" "sneak.berlin/go/webhooker/internal/logger" ) @@ -40,6 +41,11 @@ type WebhookDBManager struct { dataDir string dbs sync.Map // map[webhookID]*gorm.DB 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 @@ -53,8 +59,9 @@ func NewWebhookDBManager( log: params.Logger.Get(), } - // Create data directory if it doesn't exist - err := os.MkdirAll(m.dataDir, dataDirPerm) + // Create data directory if it doesn't exist. datadir.DirPerm is the + // single source of the directory mode; either package may run first. + err := os.MkdirAll(m.dataDir, datadir.DirPerm) if err != nil { return nil, fmt.Errorf( "creating data directory %s: %w", @@ -84,43 +91,39 @@ func (m *WebhookDBManager) GetDB( ) (*gorm.DB, error) { // Fast path: already open if val, ok := m.dbs.Load(webhookID); ok { - cachedDB, castOK := val.(*gorm.DB) - if !castOK { - return nil, fmt.Errorf( - "%w for webhook %s", - errInvalidCachedDBType, - webhookID, - ) - } + return asGormDB(val, 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) if err != nil { return nil, err } - // Store it; if another goroutine beat us, close ours - 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() - } + m.dbs.Store(webhookID, db) - existingDB, castOK := actual.(*gorm.DB) - if !castOK { - return nil, fmt.Errorf( - "%w for webhook %s", - errInvalidCachedDBType, - webhookID, - ) - } + return db, nil +} - 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 @@ -151,6 +154,11 @@ func (m *WebhookDBManager) DBExists( func (m *WebhookDBManager) DeleteDB( webhookID string, ) 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 if val, ok := m.dbs.LoadAndDelete(webhookID); ok { if gormDB, castOK := val.(*gorm.DB); castOK { @@ -184,6 +192,11 @@ func (m *WebhookDBManager) DeleteDB( // CloseAll closes all open per-webhook database connections. // Called during application shutdown. 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 m.dbs.Range(func(key, value any) bool { diff --git a/internal/database/webhook_db_manager_test.go b/internal/database/webhook_db_manager_test.go index 771f9e4..29a788c 100644 --- a/internal/database/webhook_db_manager_test.go +++ b/internal/database/webhook_db_manager_test.go @@ -1,10 +1,14 @@ package database_test import ( + "bytes" "context" + "log/slog" "net/http" "os" "path/filepath" + "strings" + "sync" "testing" "github.com/google/uuid" @@ -104,6 +108,54 @@ func TestWebhookDBManager_CreateAndGetDB(t *testing.T) { 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) { t.Parallel() diff --git a/internal/datadir/lock.go b/internal/datadir/lock.go index a50929c..bb605cd 100644 --- a/internal/datadir/lock.go +++ b/internal/datadir/lock.go @@ -29,9 +29,11 @@ import ( // process that was killed with SIGKILL blocks nothing. const LockFileName = "webhooker.lock" -// dirPerm is the mode Acquire creates DATA_DIR with. It matches what -// internal/database uses, since whichever runs first creates it. -const dirPerm = 0o750 +// DirPerm is the mode DATA_DIR is created with. It is the single +// source of that mode: internal/database consumes it rather than +// 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 // directory. Callers that need to know whether a deployment is running @@ -64,7 +66,7 @@ func Acquire(dir string) (*Lock, error) { return nil, ErrNoDir } - err := os.MkdirAll(dir, dirPerm) + err := os.MkdirAll(dir, DirPerm) if err != nil { return nil, fmt.Errorf( "creating data directory %s: %w", dir, err, diff --git a/internal/delivery/circuit_breaker.go b/internal/delivery/circuit_breaker.go index 0d01c70..afb77d4 100644 --- a/internal/delivery/circuit_breaker.go +++ b/internal/delivery/circuit_breaker.go @@ -76,12 +76,20 @@ func (cb *CircuitBreaker) Allow() bool { } } -// CooldownRemaining returns how much time is left before -// an open circuit transitions to half-open. +// CooldownRemaining returns how long a delivery that Allow refused +// 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 { cb.mu.Lock() defer cb.mu.Unlock() + if cb.state == CircuitHalfOpen { + return cb.cooldown + } + if cb.state != CircuitOpen { return 0 } diff --git a/internal/delivery/circuit_breaker_test.go b/internal/delivery/circuit_breaker_test.go index 53829ae..8b85084 100644 --- a/internal/delivery/circuit_breaker_test.go +++ b/internal/delivery/circuit_breaker_test.go @@ -267,7 +267,7 @@ func TestCircuitBreaker_CooldownRemaining_ClosedReturnsZero( ) } -func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsZero( +func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsCooldown( t *testing.T, ) { t.Parallel() @@ -282,9 +282,11 @@ func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsZero( 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(), - "half-open circuit should have zero cooldown remaining", + "a delivery refused while half-open should wait "+ + "a whole cooldown", ) } diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index e023918..4786f08 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -362,6 +362,15 @@ func (e *Engine) start() { // stop cancels the worker pool's context and waits for the pool // to drain, bounded by the stop hook's context: a wedged worker // must not hang the process past fx's stop timeout. +// +// 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 { e.log.Info("delivery engine stopping") @@ -376,6 +385,8 @@ func (e *Engine) stop(ctx context.Context) error { return err } + e.dbTarget.evictAll() + e.log.Info("delivery engine stopped") return nil @@ -438,6 +449,31 @@ func (e *Engine) processNewTask( 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, err = e.hydrateEvent( @@ -482,7 +518,7 @@ func (e *Engine) processRetryTask( return } - d, err := e.loadRetryDelivery( + d, err := e.loadDelivery( webhookDB, task.DeliveryID, ) if err != nil { @@ -703,9 +739,7 @@ func (e *Engine) recoverSingleRetry( // webhook on one bad read would be a far larger fault than // the strand it is meant to clear. if errors.Is(err, gorm.ErrRecordNotFound) { - e.failMissingTargetRetry( - webhookDB, webhookID, d, - ) + e.failMissingTarget(webhookDB, webhookID, d) return } @@ -1108,9 +1142,7 @@ func (e *Engine) sweepSingleRetry( // Deleted is terminal, unreadable is not; see // recoverSingleRetry. if errors.Is(err, gorm.ErrRecordNotFound) { - e.failMissingTargetRetry( - webhookDB, webhookID, d, - ) + e.failMissingTarget(webhookDB, webhookID, d) return } @@ -1224,19 +1256,19 @@ func (e *Engine) failUnretryableRetry( e.failDelivery(webhookDB, d, target.Type, reason) } -// failMissingTargetRetry terminally fails an orphaned retrying -// delivery whose target row is gone. Both restart recovery and the -// periodic sweep call it, so the transition exists once. +// failMissingTarget terminally fails a recovered delivery, pending or +// retrying, whose target row is gone. Restart recovery and the periodic +// sweep call it for both statuses, so the transition exists once. // -// Until it existed both paths logged the failed lookup and returned, -// which left the delivery retrying for the life of the database and +// Until it existed those paths logged the failed lookup and moved on, +// 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 // with a recorded reason is the treatment the other orphaned-retry // 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 // action, not a system fault. -func (e *Engine) failMissingTargetRetry( +func (e *Engine) failMissingTarget( webhookDB *gorm.DB, webhookID string, d *database.Delivery, @@ -1249,13 +1281,37 @@ func (e *Engine) failMissingTargetRetry( 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) e.log.Warn( - "failing orphaned retrying delivery: "+ - "its target no longer exists", + "failing recovered delivery: its target no longer exists", "webhook_id", webhookID, "delivery_id", d.ID, + "status", d.Status, "target_id", d.TargetID, "target_type", targetType, ) @@ -1289,15 +1345,14 @@ func (e *Engine) missingTargetReason( if err != nil { return "", fmt.Sprintf( "target %s no longer exists; the delivery "+ - "cannot be retried and has been failed "+ - "terminally", + "has been failed terminally", targetID, ) } return target.Type, fmt.Sprintf( "target %q (type %s) was deleted; the delivery "+ - "cannot be retried and has been failed terminally", + "has been failed terminally", target.Name, target.Type, ) } @@ -1643,7 +1698,7 @@ func (e *Engine) hydrateEvent( return event, nil } -func (e *Engine) loadRetryDelivery( +func (e *Engine) loadDelivery( webhookDB *gorm.DB, deliveryID string, ) (*database.Delivery, error) { var d database.Delivery @@ -1996,13 +2051,30 @@ func (e *Engine) sendRecoveredDeliveries( target, ok := targetMap[deliveries[i].TargetID] if !ok { - e.log.Error( - "target not found for delivery", - "delivery_id", deliveries[i].ID, - "target_id", deliveries[i].TargetID, - ) + // A missing entry does not mean the target is gone: the + // map is also empty when its query failed. Only a lookup + // that finds no row ends the delivery; any other error + // 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( diff --git a/internal/delivery/engine_lifecycle_test.go b/internal/delivery/engine_lifecycle_test.go index 6e7ef2e..9fadf9d 100644 --- a/internal/delivery/engine_lifecycle_test.go +++ b/internal/delivery/engine_lifecycle_test.go @@ -2,6 +2,8 @@ package delivery_test import ( "context" + "fmt" + "path/filepath" "testing" "time" @@ -269,3 +271,88 @@ func TestEngine_StopHookHonoursStopTimeout(t *testing.T) { 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", + ) +} diff --git a/internal/delivery/engine_test.go b/internal/delivery/engine_test.go index 6c2e0a2..13d1625 100644 --- a/internal/delivery/engine_test.go +++ b/internal/delivery/engine_test.go @@ -17,6 +17,7 @@ import ( "time" "github.com/google/uuid" + "github.com/prometheus/client_golang/prometheus" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "gorm.io/driver/sqlite" @@ -24,6 +25,7 @@ import ( _ "modernc.org/sqlite" "sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/delivery" + "sneak.berlin/go/webhooker/internal/metrics" ) // 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) { t.Parallel() @@ -1070,6 +1166,10 @@ func TestIsForwardableHeader(t *testing.T) { assert.False(t, delivery.ExportIsForwardableHeader("Content-Length"), ) + + assert.False(t, + delivery.ExportIsForwardableHeader("Content-Type"), + ) } 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( t *testing.T, ) { diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index b71de96..9dd2531 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -101,6 +101,19 @@ func (e *Engine) ExportDeliverHTTP( 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. func (e *Engine) ExportDeliverDatabase( webhookDB *gorm.DB, d *database.Delivery, @@ -179,6 +192,14 @@ func (e *Engine) ExportGetCircuitBreaker( 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. func (e *Engine) ExportParseHTTPConfig( configJSON string, @@ -321,6 +342,31 @@ func (e *Engine) ExportRecoverRetryingDeliveries( 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. func (e *Engine) ExportDeliveryCh() chan Task { return e.deliveryCh diff --git a/internal/delivery/inflight_test.go b/internal/delivery/inflight_test.go index 1c86003..08a0f69 100644 --- a/internal/delivery/inflight_test.go +++ b/internal/delivery/inflight_test.go @@ -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 // of the pending reconcile. A second attempt that reached the receiver // and whose status write then failed sits at retrying holding a diff --git a/internal/delivery/metrics_test.go b/internal/delivery/metrics_test.go index 48492dd..c9d95a2 100644 --- a/internal/delivery/metrics_test.go +++ b/internal/delivery/metrics_test.go @@ -412,9 +412,10 @@ func TestDeliveryMetrics_BreakerBlockedIsNotAnAttempt( s.Engine.ExportProcessRetryTask(context.TODO(), &blocked) - // The breaker refused it: rescheduled, so the retry counter - // moved, but nothing was attempted or timed. - assert.InDelta(t, retriesBefore+1, + // The breaker refused it: rescheduled without rewriting the + // retrying status it already had, so the retry counter did not + // move, and nothing was attempted or timed. + assert.InDelta(t, retriesBefore, mCounter(t, reg, mRetries, mTypeHTTP), 0) assert.InDelta(t, threshold, mCounter(t, reg, mAttempts, mTypeHTTP), 0) diff --git a/internal/delivery/redirect_test.go b/internal/delivery/redirect_test.go index 31e3d66..e5d5347 100644 --- a/internal/delivery/redirect_test.go +++ b/internal/delivery/redirect_test.go @@ -170,7 +170,8 @@ func TestDelivery_CrossOriginRedirectDropsOriginScopedHeaders( // Stripping must not fire within the configured origin, or every // destination that redirects its own path would lose its // 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( t *testing.T, ) { @@ -338,10 +339,11 @@ func TestRedirectPolicy_StopsAtHopCap(t *testing.T) { // The set the redirect policy strips is whatever the delivery path // actually put on the wire, so a header added to the forward set is // covered without a second edit. A header the event never carried -// is not in the set, and the delivery path's own two are deliberately -// excluded: Content-Type describes the body, which a 307 carries -// across hosts, and the inbound User-Agent every real sender supplies -// is overwritten before the request goes out. +// is not in the set, and neither is the inbound Content-Type, because +// it is not forwarded. Two more are deliberately excluded: a +// Content-Type configured on the target describes the body, which a +// 307 carries across hosts, and the inbound User-Agent every real +// sender supplies is overwritten before the request goes out. func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) { t.Parallel() @@ -370,6 +372,7 @@ func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) { &delivery.HTTPTargetConfig{ Headers: map[string]string{ probeHeaderName: probeHeaderValue, + "Content-Type": testContentType, }, }, ) @@ -377,7 +380,11 @@ func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) { assert.Equal(t, []string{probeHeaderName, inboundHeaderName}, names, "both header classes are reported, and only those: "+ - "Host is never forwarded, Content-Type and "+ - "User-Agent are the delivery path's own", + "Host and the inbound Content-Type are never "+ + "forwarded, User-Agent is the delivery path's own", + ) + assert.NotContains(t, names, "Content-Type", + "a Content-Type configured on the target must survive "+ + "a cross-origin 307/308 with the body it describes", ) } diff --git a/internal/delivery/ssrf.go b/internal/delivery/ssrf.go index a718e6a..0fa73b3 100644 --- a/internal/delivery/ssrf.go +++ b/internal/delivery/ssrf.go @@ -26,7 +26,7 @@ var ( "hostname resolved to no IP addresses", ) errBlockedIP = errors.New( - "blocked private/reserved IP range", + "blocked private, reserved or cloud metadata address", ) errBlockedMetadata = errors.New( "blocked link-local or cloud instance metadata " + @@ -37,9 +37,10 @@ var ( ) ) -// blockedNetworks contains all private/reserved IP ranges -// that should be blocked to prevent SSRF attacks. An operator -// can permit specific blocks out of this set with +// blockedNetworks is the default blocklist: the private and +// reserved IP ranges, plus the public cloud metadata addresses, +// that are blocked to prevent SSRF attacks. An operator can +// permit specific blocks out of this set with // ALLOWED_EGRESS_CIDRS; see Guard. // //nolint:gochecknoglobals // package-level network list is appropriate here @@ -122,6 +123,8 @@ func init() { "::1/128", "fc00::/7", "fe80::/10", + // Azure WireServer, a public address that serves VM credentials. + "168.63.129.16/32", }) // Every entry is named. The set must not grow or shrink @@ -216,8 +219,8 @@ func matchesAny(networks []*net.IPNet, ip net.IP) bool { } // isBlockedIP checks whether an IP address falls within -// any blocked private/reserved network range, before any -// operator allowlist is considered. +// the default blocklist, before any operator allowlist is +// considered. func isBlockedIP(ip net.IP) bool { return matchesAny(blockedNetworks, ip) } @@ -320,7 +323,7 @@ func (g *Guard) allows(ip net.IP) bool { // // 1. alwaysBlockedNetworks is refused before the allowlist is // consulted, so no configured CIDR reaches link-local or a -// cloud instance metadata endpoint. +// cloud metadata endpoint at a non-public address. // 2. The allowlist is consulted next, so a listed private // network becomes reachable. // 3. Everything else keeps the default blocklist's answer. diff --git a/internal/delivery/ssrf_allowlist_test.go b/internal/delivery/ssrf_allowlist_test.go index 74f0f35..314865c 100644 --- a/internal/delivery/ssrf_allowlist_test.go +++ b/internal/delivery/ssrf_allowlist_test.go @@ -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 // validator and the dialer are not two policies that happen to // agree: both are defined in terms of checkIP, so the exported diff --git a/internal/delivery/target_database.go b/internal/delivery/target_database.go index d319551..7173b71 100644 --- a/internal/delivery/target_database.go +++ b/internal/delivery/target_database.go @@ -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 // expiry, without requiring a write. It returns nil (nothing to // do) when the archive file does not exist, so a sweep never diff --git a/internal/delivery/target_database_evict_test.go b/internal/delivery/target_database_evict_test.go index 14e7945..095823d 100644 --- a/internal/delivery/target_database_evict_test.go +++ b/internal/delivery/target_database_evict_test.go @@ -1,6 +1,7 @@ package delivery_test import ( + "context" "errors" "fmt" "net/http" @@ -361,3 +362,46 @@ func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) { "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", + ) +} diff --git a/internal/delivery/target_headers_test.go b/internal/delivery/target_headers_test.go index 1b1c7ae..3237132 100644 --- a/internal/delivery/target_headers_test.go +++ b/internal/delivery/target_headers_test.go @@ -11,10 +11,11 @@ import ( "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. const ( headerAuthorization = "Authorization" + headerContentType = "Content-Type" bearerValue = "Bearer abc" archiveConfigNever = "{\"expiry\":\"never\"}" ) diff --git a/internal/delivery/target_http.go b/internal/delivery/target_http.go index 127c7c4..9c7b539 100644 --- a/internal/delivery/target_http.go +++ b/internal/delivery/target_http.go @@ -197,10 +197,14 @@ func (c *httpCore) circuitBreakerBlock( "cooldown_remaining", remaining, ) - c.eng.settleStatus( - webhookDB, d, d.Target.Type, - database.DeliveryStatusRetrying, - ) + // A delivery already at retrying is left as it is, so a task + // the breaker keeps turning away writes nothing each time. + if d.Status != database.DeliveryStatusRetrying { + c.eng.settleStatus( + webhookDB, d, d.Target.Type, + database.DeliveryStatusRetrying, + ) + } retryTask := *task sched.ScheduleRetry(retryTask, remaining) @@ -537,6 +541,11 @@ func isForwardableHeader(name string) bool { "Upgrade", "Proxy-Authorization", "Proxy-Connection", "Content-Length": return false + case "Content-Type": + // applyRequestHeaders sets Content-Type itself. The receiver + // already stored this inbound value as the event's + // ContentType, so forwarding it too would send it twice. + return false default: return true } @@ -549,6 +558,10 @@ func isForwardableHeader(name string) bool { // policy strips exactly that set on a hop that leaves the origin, // so the forward set is decided here and only here — a header added // to it is covered off-origin without a second edit elsewhere. +// +// Content-Type goes out once: a Content-Type configured on the target +// wins, otherwise the event's ContentType, otherwise none. The inbound +// Content-Type in the event's headers is never forwarded. func applyRequestHeaders( req *http.Request, event *database.Event, @@ -569,10 +582,10 @@ func applyRequestHeaders( req.Header.Set("User-Agent", "webhooker/1.0") - // Content-Type describes the body being sent rather than the - // sender, and the delivery path sets it from the event itself. - // A 307/308 preserves the body across hosts, so stripping it - // would send that body untyped. + // A Content-Type configured on the target describes the body + // being sent rather than the sender. A 307/308 preserves the + // body across hosts, so stripping it would send that body + // untyped. delete(originScoped, "Content-Type") // User-Agent is overwritten just above, so an inbound one never diff --git a/internal/delivery/terminal_state_test.go b/internal/delivery/terminal_state_test.go index 382a6ff..8eda81a 100644 --- a/internal/delivery/terminal_state_test.go +++ b/internal/delivery/terminal_state_test.go @@ -18,16 +18,17 @@ import ( // https://git.eeqj.de/sneak/webhooker/issues/107: a delivery failed // 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 -// 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 // a target whose type was written by a build that knew a type this one // does not. const tUnknownType = database.TargetType("pubsub") -// tSeedDeletedTarget creates a target, a retrying delivery against it -// with one recorded failed attempt, and then deletes the target the -// way the source page does. +// tSeedDeletedTarget creates a target, a delivery against it at the +// given status with one recorded failed attempt, and then deletes the +// target the way the source page does. // // 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 @@ -36,6 +37,7 @@ func tSeedDeletedTarget( t *testing.T, s iSetup, name, url string, + status database.DeliveryStatus, ) string { t.Helper() @@ -51,8 +53,7 @@ func tSeedDeletedTarget( ) d := iSeedDelivery( - t, s.WebhookDB, event.ID, targetID, - database.DeliveryStatusRetrying, + t, s.WebhookDB, event.ID, targetID, status, ) iSeedFailedResult(t, s.WebhookDB, d.ID) @@ -173,6 +174,7 @@ func TestRecoverSingleRetry_TargetDeleted(t *testing.T) { deliveryID := tSeedDeletedTarget( t, s, "gone-on-recovery", "http://example.com/hook", + database.DeliveryStatusRetrying, ) s.Engine.ExportRecoverWebhookDeliveries( @@ -210,6 +212,7 @@ func TestSweepSingleRetry_TargetDeleted(t *testing.T) { deliveryID := tSeedDeletedTarget( t, s, "gone-on-sweep", "http://example.com/hook", + database.DeliveryStatusRetrying, ) // 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 -// path to the same rule as the existing one: no target row, and so no +// TestFailMissingTarget_WritesNoTargetRow holds the new terminal path +// 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 // database. See https://git.eeqj.de/sneak/webhooker/issues/206. -func TestFailMissingTargetRetry_WritesNoTargetRow( +func TestFailMissingTarget_WritesNoTargetRow( t *testing.T, ) { t.Parallel() @@ -297,6 +300,7 @@ func TestFailMissingTargetRetry_WritesNoTargetRow( deliveryID := tSeedDeletedTarget( t, s, "credential-bearing", hookURL, + database.DeliveryStatusRetrying, ) s.Engine.ExportSweepWebhookRetries( @@ -529,3 +533,260 @@ func TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone( 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()) +} diff --git a/internal/handlers/ui_copy_test.go b/internal/handlers/ui_copy_test.go index da9b6d6..8598de9 100644 --- a/internal/handlers/ui_copy_test.go +++ b/internal/handlers/ui_copy_test.go @@ -300,3 +300,80 @@ func TestEntrypointCopyButtonIsProgressiveEnhancement(t *testing.T) { "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", + ) +} diff --git a/internal/middleware/csrf_test.go b/internal/middleware/csrf_test.go index 4273964..f1660ed 100644 --- a/internal/middleware/csrf_test.go +++ b/internal/middleware/csrf_test.go @@ -380,9 +380,8 @@ func csrfTookStrictPath( // TestCSRF_ForwardedProtoSpellingsTakeStrictPath runs the header // spellings a real proxy emits through the middleware. The environment -// is dev -- the DEFAULT when WEBHOOKER_ENVIRONMENT is unset -- to pin -// that the routing is a per-request transport decision and owes -// nothing to configuration. +// is set to dev -- the permissive setting -- to pin that the routing is +// a per-request transport decision and owes nothing to configuration. func TestCSRF_ForwardedProtoSpellingsTakeStrictPath(t *testing.T) { t.Parallel() diff --git a/internal/middleware/recoverer.go b/internal/middleware/recoverer.go index 1b47d6c..8dcb76e 100644 --- a/internal/middleware/recoverer.go +++ b/internal/middleware/recoverer.go @@ -133,6 +133,11 @@ func (w *recoverResponseWriter) Unwrap() http.ResponseWriter { // what the access log records and the metrics count, and outside the // sentryhttp handler, whose Repanic option depends on something // further out recovering what it re-raises. +// +// Unlike http.Error on its own, it deletes any Set-Cookie the handler +// set before panicking, because a request that failed must not hand +// the client a credential; every other header is left to http.Error. +// See https://git.eeqj.de/sneak/webhooker/issues/193. func (s *Middleware) Recoverer() func(http.Handler) http.Handler { return func(next http.Handler) http.Handler { return http.HandlerFunc(func( @@ -164,6 +169,8 @@ func (s *Middleware) Recoverer() func(http.Handler) http.Handler { return } + rw.Header().Del("Set-Cookie") + http.Error( rw, http.StatusText( diff --git a/internal/middleware/recoverer_test.go b/internal/middleware/recoverer_test.go index bbfd309..fac5cc4 100644 --- a/internal/middleware/recoverer_test.go +++ b/internal/middleware/recoverer_test.go @@ -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 // panics after sending its status. The bytes are already on the wire, -// so a second WriteHeader would change nothing the client sees and -// would draw net/http's "superfluous response.WriteHeader" report. +// cookie included, so a second WriteHeader would change nothing the +// client sees and would draw net/http's "superfluous +// response.WriteHeader" report. func TestRecovererKeepsAnAlreadyCommittedResponse(t *testing.T) { t.Parallel() probe := newRecovererProbe( t, false, func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Set-Cookie", "session=x") w.WriteHeader(committedStatus) _, _ = w.Write([]byte("partial")) @@ -331,6 +359,7 @@ func TestRecovererKeepsAnAlreadyCommittedResponse(t *testing.T) { assert.Equal(t, committedStatus, resp.StatusCode) assert.Equal(t, "partial", string(body)) + assert.Len(t, resp.Cookies(), 1) record := probe.panicRecord(t) assert.Equal(t, panicMarker, record["panic"]) diff --git a/internal/reqtls/reqtls.go b/internal/reqtls/reqtls.go index f1704a4..b9e28b5 100644 --- a/internal/reqtls/reqtls.go +++ b/internal/reqtls/reqtls.go @@ -5,7 +5,7 @@ // several packages, by hand, and the answers disagreed. The session // cookie's Secure attribute was decided at startup from the configured // 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. // Everything kept working, which is exactly why nobody noticed. // diff --git a/internal/session/session.go b/internal/session/session.go index 901568e..97b6664 100644 --- a/internal/session/session.go +++ b/internal/session/session.go @@ -146,8 +146,8 @@ func newStore(key []byte) *sessions.CookieStore { // // This is decided per-request, not once at startup. Deciding it at // startup from the configured environment is what this replaces, and -// it got the DEFAULT posture wrong: "dev" is the environment when -// WEBHOOKER_ENVIRONMENT is unset, so a deployment terminating TLS at a +// it got the DEFAULT posture wrong: "dev" was then the environment when +// WEBHOOKER_ENVIRONMENT was unset, so a deployment terminating TLS at a // proxy without also setting the environment emitted the // authentication cookie with no Secure attribute -- silently, and on // the same response as a CSRF cookie that did have one. diff --git a/internal/session/session_test.go b/internal/session/session_test.go index 962df24..a52ae96 100644 --- a/internal/session/session_test.go +++ b/internal/session/session_test.go @@ -990,8 +990,8 @@ func sessionCookieFrom( // TestSave_SecureFollowsRequestTransport is the regression test for // the defect this replaces: Secure was fixed at startup from the -// configured environment, and "dev" is the environment when -// WEBHOOKER_ENVIRONMENT is unset. A deployment behind a TLS proxy in +// configured environment, and "dev" was then the environment when +// WEBHOOKER_ENVIRONMENT was unset. A deployment behind a TLS proxy in // that DEFAULT posture shipped the authentication cookie with no // Secure attribute and said nothing about it. // diff --git a/templates/source_detail.html b/templates/source_detail.html index bde8cae..3b26967 100644 --- a/templates/source_detail.html +++ b/templates/source_detail.html @@ -120,9 +120,12 @@ -
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.
0 is fire-and-forget: one attempt, no circuit breaker.
+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.