From 888eaf526bbd8a6e150bdcc48e3d31cb58502e0f Mon Sep 17 00:00:00 2001 From: clawbot Date: Sun, 30 Aug 2026 04:05:38 +0200 Subject: [PATCH 01/15] State the UUID-is-the-credential rule as a rule (closes #301) (#302) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Closes https://git.eeqj.de/sneak/webhooker/issues/301. Docs-only apart from one test comment; no behaviour change. The receiver has authenticated on the entrypoint UUID alone since inbound signature verification was removed in https://git.eeqj.de/sneak/webhooker/pulls/279. The README described that as the current state. It did not say it is the decision, which leaves a future contributor free to propose HMAC as an improvement rather than as a reversal. What changed: - `## The entrypoint URL is the authentication secret` now states the rule: the v4 UUID at `/webhook/{uuid}` is the credential and the only one; no shared secret, HMAC signature, bearer token or second factor will be added, including as defence in depth. It names the removal that settled it, and it says explicitly that signature headers a sender sends anyway are stored and forwarded but never checked — the previous text left that ambiguous. - The same section carries the two consequences an operator has to act on: the URL is a capability, so keep it out of logs, tickets and screenshots; and rotation means minting a new entrypoint, not changing a key. - It also handles the case the rule will next be argued from: a sender that only supports signed payloads to a well-known URL is a constraint on that integration, to be raised on its own terms, not grounds to reintroduce shared secrets. - The rule is reachable without scrolling 1,100 lines: a pointer in the intro, a new first bullet under Authentication (which previously listed the web UI, the API and `/metrics` and said nothing about the receiver at all), and a sharpened bullet under Security. Stale language found and corrected: one, in `internal/delivery/redirect_test.go`. Its comment justified same-origin header retention partly by "the inbound signature the receiver verifies" — in this repo's vocabulary "the receiver" is `/webhook/{uuid}`, which verifies nothing. The endpoint that verifies it is the delivery target's, and the comment now says so. Two places that read like stale signing language were checked and left alone as accurate: `internal/delivery/redirect.go` and `internal/server/sentry.go` describe signature headers senders put on the receiver route, which do arrive and are forwarded — neither claims webhooker checks them. `REPO_POLICIES.md` was deliberately not touched. It is the cross-project policy document synced from `sneak/prompts` and carries `last_modified` front matter for that purpose, so a webhooker-specific carve-out does not belong in it. Worth knowing: its hardening section ends "if a standard security hardening measure exists for HTTP services and is not listed here, it is still expected. When in doubt, harden" — that is the sentence a future HMAC proposal will cite, and only the README now answers it. `TODO.md` is untouched per its own Workflow section (issue branches do not touch it). Co-authored-by: sneak Reviewed-on: https://git.eeqj.de/sneak/webhooker/pulls/302 Co-authored-by: clawbot Co-committed-by: clawbot --- README.md | 52 +++++++++++++++++++++++++----- internal/delivery/redirect_test.go | 3 +- 2 files changed, 46 insertions(+), 9 deletions(-) diff --git a/README.md b/README.md index c66e5ca..319989c 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 @@ -1149,14 +1156,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 @@ -2867,6 +2898,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 +2941,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 diff --git a/internal/delivery/redirect_test.go b/internal/delivery/redirect_test.go index 31e3d66..1534697 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, ) { -- 2.54.0 From 39afa69bfc3439b0a1d911eccfc83131f7ecadec Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Mon, 21 Sep 2026 10:01:51 +0200 Subject: [PATCH 02/15] Consolidate the data directory mode into one owner (closes #288) Two packages each declared the 0o750 mode for DATA_DIR and both created the directory. internal/datadir now exports DirPerm as the single definition, and internal/database uses it in both places it creates the directory. The value is unchanged, so existing deployments see no permission change. datadir owns it because guarding and creating DATA_DIR is that package's whole purpose and it imports nothing that would form a cycle. Model: opus-4-8 (implementation and review); fable-5-1 (merge) --- internal/database/database.go | 6 ++++-- internal/database/webhook_db_manager.go | 6 ++++-- internal/datadir/lock.go | 10 ++++++---- 3 files changed, 14 insertions(+), 8 deletions(-) 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/webhook_db_manager.go b/internal/database/webhook_db_manager.go index 81ca427..7da8169 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" ) @@ -53,8 +54,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", 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, -- 2.54.0 From b05182137050351e5359746ebcedca9e5b19f467 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Mon, 21 Sep 2026 18:33:09 +0200 Subject: [PATCH 03/15] Say max_retries is the total attempt count, not a retry count (closes #316) The help text under the field on both target forms and the max_retries rows in the README now say the number is the total number of delivery attempts: 0 is a single attempt with no retries and no circuit breaker, and N is N attempts in total. The delivery code already worked this way; only the wording was wrong, so an operator wanting one try plus two retries would have entered 2 instead of 3. A UI copy test renders both forms and pins the wording. Delivery behaviour is unchanged. Model: opus-4-8 (implementation); fable-5-1 (merge) --- README.md | 14 +++--- internal/handlers/ui_copy_test.go | 77 +++++++++++++++++++++++++++++++ templates/source_detail.html | 9 ++-- templates/target_edit.html | 2 +- 4 files changed, 91 insertions(+), 11 deletions(-) diff --git a/README.md b/README.md index 319989c..e1289af 100644 --- a/README.md +++ b/README.md @@ -1507,7 +1507,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. @@ -1515,12 +1515,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` 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/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.

diff --git a/templates/target_edit.html b/templates/target_edit.html index 9194721..a2ff56c 100644 --- a/templates/target_edit.html +++ b/templates/target_edit.html @@ -69,7 +69,7 @@
-

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.

{{end}} -- 2.54.0 From 7ed158844361261704426aa62c303a6154ccf1b1 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Mon, 28 Sep 2026 12:30:32 +0200 Subject: [PATCH 04/15] Document running webhooker under upaas (closes #323) Adds a "Running under upaas" section to the README: add no port mapping, since upaas publishes mapped ports on every host interface, and put the app on the reverse proxy's Docker network instead; one data volume at /var/lib/webhooker, created owned by UID 1000 before the first deploy; WEBHOOKER_ENVIRONMENT and TRUSTED_PROXIES; the health check upaas reads 60 seconds after a deploy; and where the first-run admin password appears and how to reset it. upaas bind-mounts a host directory it never creates, and one made by root stops the container at its data directory lock. The documented creation step removes that; the image is unchanged. Model: opus-5-5 --- README.md | 60 +++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 60 insertions(+) diff --git a/README.md b/README.md index e1289af..85275c4 100644 --- a/README.md +++ b/README.md @@ -731,6 +731,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 -- 2.54.0 From 237f131367db343e5b8cce0826d54a6f788dae95 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Mon, 28 Sep 2026 12:47:31 +0200 Subject: [PATCH 05/15] Default WEBHOOKER_ENVIRONMENT to prod (closes #307) An unset WEBHOOKER_ENVIRONMENT now means prod, not dev. The only thing dev still changes is CORS, which then answers every origin with Access-Control-Allow-Origin: *, so an operator who forgets the variable is no longer silently permissive; dev must be set explicitly. Cookie Secure and CSRF strictness follow each request's transport and are unaffected. The README, comments and tests no longer describe dev as the default: the deployment checklist asks only that the environment is not dev, the Docker and nginx examples drop the now-redundant setting, and the TRUSTED_PROXIES warning gives its real reason for firing in every environment. Model: opus-4-8 (implementation); opus-5-5 (rework) --- README.md | 39 ++++++++++++++++---------------- internal/config/config.go | 14 ++++++------ internal/config/config_test.go | 17 +++++--------- internal/middleware/csrf_test.go | 5 ++-- internal/reqtls/reqtls.go | 2 +- internal/session/session.go | 4 ++-- internal/session/session_test.go | 4 ++-- 7 files changed, 39 insertions(+), 46 deletions(-) diff --git a/README.md b/README.md index 85275c4..eb5ee27 100644 --- a/README.md +++ b/README.md @@ -44,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 @@ -92,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` | | -------- | ----------------------- | ---------------- | @@ -134,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` | @@ -399,10 +400,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. @@ -638,7 +638,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 ``` @@ -812,17 +811,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 @@ -914,7 +914,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 ``` diff --git a/internal/config/config.go b/internal/config/config.go index 6a7a628..69c5210 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -585,12 +585,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 && @@ -772,10 +774,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..0add1b2 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", @@ -848,10 +848,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 +870,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/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/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. // -- 2.54.0 From 6ebac4fa713f8a5e03c9cabf0cafc0100b6e9809 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Mon, 28 Sep 2026 14:13:22 +0200 Subject: [PATCH 06/15] Index the event-tier columns the sweeps, event log and retention scan (closes #314) The per-webhook event databases had no secondary indexes, so startup recovery, the retry and pending sweeps, the queue-depth sampler, the event log and retention each read whole tables. Indexes declared in the GORM model tags now serve them, and AutoMigrate adds them to new and existing databases alike. Each index also covers deleted_at: GORM adds deleted_at IS NULL to these queries, and SQLite, with no table statistics, otherwise prefers the existing deleted_at index. A test checks SQLite's plan for each statement as GORM builds it. Rule suppressed: lll on the three event-tier model structs, whose struct tags cannot wrap. The resubmitted_from_id scan is left to https://git.eeqj.de/sneak/webhooker/issues/325. Model: opus-4-8 (implementation); opus-5-5 (rework) --- README.md | 22 +++ 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 ++ 5 files changed, 230 insertions(+), 4 deletions(-) create mode 100644 internal/database/event_tier_indexes_test.go diff --git a/README.md b/README.md index eb5ee27..7db1199 100644 --- a/README.md +++ b/README.md @@ -1759,6 +1759,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`. 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 Date: Tue, 29 Sep 2026 04:12:14 +0200 Subject: [PATCH 07/15] Do not resend a delivery that restart recovery already sent (closes #299) Restart recovery could find a just-written delivery pending, send it and release it before the receiver's Notify queued the same delivery. Notify's claim then succeeded on the released id, and the worker sent it again because the new-task path never read the delivery's row. Before sending a new task the worker now reads the delivery's status by primary key and skips the task unless the row still says pending, as the retry path already does for retrying. Nothing else can change the row while the worker owns the delivery. A row left pending by a failed bookkeeping write is still sent again. loadRetryDelivery is renamed loadDelivery now that both paths use it. Model: opus-5-5 --- internal/delivery/engine.go | 29 ++++++++++++- internal/delivery/inflight_test.go | 67 ++++++++++++++++++++++++++++++ 2 files changed, 94 insertions(+), 2 deletions(-) diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index e023918..1143cbf 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -438,6 +438,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 +507,7 @@ func (e *Engine) processRetryTask( return } - d, err := e.loadRetryDelivery( + d, err := e.loadDelivery( webhookDB, task.DeliveryID, ) if err != nil { @@ -1643,7 +1668,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 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 -- 2.54.0 From 3cdab97930baad741bd5d5c17786e18e545436f5 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 04:30:14 +0200 Subject: [PATCH 08/15] Say make check needs make bootstrap on a fresh clone (closes #282) The third-party browser assets are not committed, so on a fresh clone make check fails in the tests until make bootstrap (or make assets) has fetched them. The Entrypoints section now says so up front, and why the check does not fetch them itself: it must not change files in the repo. Model: opus-5-5 --- README.md | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 7db1199..5b671a5 100644 --- a/README.md +++ b/README.md @@ -1256,8 +1256,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 -- 2.54.0 From 978eb01b292d380a0f1021f01cb483d93f1358aa Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 04:37:30 +0200 Subject: [PATCH 09/15] Open each event database once when callers race (closes #291) GetDB opened the database on a cache miss and then tried to cache it, so callers racing on a webhook's first use could each open the file, and the losers closed their copies. On a new file the parallel opens also create its tables at the same time, and one caller can fail with "table already exists". A mutex now covers the open: GetDB looks in the cache again under it, then opens and caches. DeleteDB and CloseAll take the same mutex, so neither runs while an open is under way. Reading an already cached database takes no lock. The new test starts many callers on one webhook at once and checks that exactly one open happened. Model: opus-5-5 --- internal/database/webhook_db_manager.go | 65 ++++++++++++-------- internal/database/webhook_db_manager_test.go | 52 ++++++++++++++++ 2 files changed, 90 insertions(+), 27 deletions(-) diff --git a/internal/database/webhook_db_manager.go b/internal/database/webhook_db_manager.go index 7da8169..a628b56 100644 --- a/internal/database/webhook_db_manager.go +++ b/internal/database/webhook_db_manager.go @@ -41,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 @@ -86,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 @@ -153,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 { @@ -186,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() -- 2.54.0 From e0b211f960f32970c959e21db837d06127ab43e9 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 05:48:20 +0200 Subject: [PATCH 10/15] Make deliveries refused while half-open wait a cooldown (closes #306) While the breaker was half-open, Allow refused every delivery but the probe and CooldownRemaining returned zero, so each queued task for the target went straight back onto the retry channel and rewrote its status on every pass until the probe finished. CooldownRemaining now returns the whole cooldown while half-open, so a refused delivery waits that long. A refused delivery already at retrying is not written again, so the retry counter now moves only when a refusal moves a delivery into retrying. Model: opus-5-5 --- README.md | 8 +- internal/delivery/circuit_breaker.go | 12 ++- internal/delivery/circuit_breaker_test.go | 8 +- internal/delivery/engine_test.go | 96 +++++++++++++++++++++++ internal/delivery/export_test.go | 21 +++++ internal/delivery/metrics_test.go | 7 +- internal/delivery/target_http.go | 12 ++- 7 files changed, 149 insertions(+), 15 deletions(-) diff --git a/README.md b/README.md index 5b671a5..b3afd54 100644 --- a/README.md +++ b/README.md @@ -2051,7 +2051,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:** @@ -2089,7 +2089,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 @@ -2103,7 +2105,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` | 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_test.go b/internal/delivery/engine_test.go index 6c2e0a2..213b13a 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() diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index b71de96..37c2695 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, 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/target_http.go b/internal/delivery/target_http.go index 127c7c4..fc264f9 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) -- 2.54.0 From 51580a2bc66b96b4f367da7db24a82e264abf5fb Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 06:48:19 +0200 Subject: [PATCH 11/15] Fail a pending delivery whose target was deleted (closes #293) Restart recovery and the pending sweep skipped a pending delivery whose target was missing from the batch's target map, every minute, for the life of the database. A miss now asks loadTarget: no row fails the delivery terminally with a recorded reason; any other error leaves it pending, since the map is also empty when its query failed; a target found there is used. The failure goes through the ownership-gated function the retrying paths already used, now failMissingTarget. Once it owns the delivery it re-reads the row and fails it only if the status is unchanged, so a delivery sent and settled in between is left alone. Model: opus-5-5 --- internal/delivery/engine.go | 82 +++++-- internal/delivery/export_test.go | 25 ++ internal/delivery/terminal_state_test.go | 279 ++++++++++++++++++++++- 3 files changed, 354 insertions(+), 32 deletions(-) diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index 1143cbf..c96dda7 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -728,9 +728,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 } @@ -1133,9 +1131,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 } @@ -1249,19 +1245,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, @@ -1274,13 +1270,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, ) @@ -1314,15 +1334,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, ) } @@ -2021,13 +2040,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/export_test.go b/internal/delivery/export_test.go index 37c2695..9dd2531 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -342,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/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()) +} -- 2.54.0 From d4f4ddf51f743159d25693bb783c86925602b4b4 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 07:11:55 +0200 Subject: [PATCH 12/15] Send Content-Type once on a delivery (closes #246) A delivery set Content-Type from the event's ContentType and then added the inbound Content-Type from the event's stored headers, so a target could receive two values. The inbound Content-Type is no longer forwarded from the stored headers; the receiver already saves it as the event's ContentType. Which value is sent is now stated at applyRequestHeaders: a Content-Type configured on the target, otherwise the event's ContentType, otherwise none. A configured one still survives a cross-origin 307/308 with its body. Model: opus-5-5 --- internal/delivery/engine_test.go | 79 ++++++++++++++++++++++++ internal/delivery/redirect_test.go | 18 ++++-- internal/delivery/target_headers_test.go | 3 +- internal/delivery/target_http.go | 17 +++-- 4 files changed, 106 insertions(+), 11 deletions(-) diff --git a/internal/delivery/engine_test.go b/internal/delivery/engine_test.go index 213b13a..13d1625 100644 --- a/internal/delivery/engine_test.go +++ b/internal/delivery/engine_test.go @@ -1166,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) { @@ -1247,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/redirect_test.go b/internal/delivery/redirect_test.go index 1534697..e5d5347 100644 --- a/internal/delivery/redirect_test.go +++ b/internal/delivery/redirect_test.go @@ -339,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() @@ -371,6 +372,7 @@ func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) { &delivery.HTTPTargetConfig{ Headers: map[string]string{ probeHeaderName: probeHeaderValue, + "Content-Type": testContentType, }, }, ) @@ -378,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/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 fc264f9..9c7b539 100644 --- a/internal/delivery/target_http.go +++ b/internal/delivery/target_http.go @@ -541,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 } @@ -553,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, @@ -573,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 -- 2.54.0 From 4a724130ca549ede9cda352f7918a61f429b3938 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 08:30:22 +0200 Subject: [PATCH 13/15] Close archive writers when the delivery engine stops (closes #280) The engine cached archive writers and never closed them at shutdown, so after a clean stop an archive's rows could sit in its -wal while the .db held no table. The engine's stop hook now evicts every cached writer once its workers have returned, the same way deleting a webhook does, so a clean stop leaves each archive as one file and a late write is refused. If the workers do not return within the stop budget, the writers are left open as a kill would leave them: closing would wait on a write in progress, and a still-running worker would open new ones. The README no longer says archives keep their sidecars across a clean stop. Model: opus-5-5 --- README.md | 33 +++---- internal/delivery/engine.go | 11 +++ internal/delivery/engine_lifecycle_test.go | 87 +++++++++++++++++++ internal/delivery/target_database.go | 18 ++++ .../delivery/target_database_evict_test.go | 44 ++++++++++ 5 files changed, 173 insertions(+), 20 deletions(-) diff --git a/README.md b/README.md index b3afd54..48a1dd0 100644 --- a/README.md +++ b/README.md @@ -968,15 +968,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 @@ -1051,10 +1046,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 @@ -1073,12 +1067,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 @@ -3088,7 +3080,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/delivery/engine.go b/internal/delivery/engine.go index c96dda7..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 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/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", + ) +} -- 2.54.0 From f755c03110706731eeb058b01defcb0162784ab0 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 10:22:07 +0200 Subject: [PATCH 14/15] Default-block Azure WireServer's public address (closes #245) Add 168.63.129.16 (Azure WireServer) to blockedNetworks, the default blocklist, not alwaysBlockedNetworks: it is public unicast, so an operator who lists it in ALLOWED_EGRESS_CIDRS can reach it again. The refusal message, the allowlist startup warning, the README and the comments no longer call every blocked address private/reserved, and no longer claim the allowlist cannot open any metadata endpoint. Sources: - https://learn.microsoft.com/en-us/azure/virtual-network/what-is-ip-address-168-63-129-16 - https://learn.microsoft.com/en-us/azure/virtual-machines/metadata-security-protocol/overview Deviation: 147.75.207.243 (Equinix Metal) is not added; Equinix documents only a hostname, and the service was sunset on 2026-06-30. Model: opus-5-5 --- README.md | 27 +++++++++++------- internal/config/config.go | 21 ++++++++------ internal/config/config_test.go | 13 +++++---- internal/delivery/ssrf.go | 17 +++++++----- internal/delivery/ssrf_allowlist_test.go | 35 ++++++++++++++++++++++++ 5 files changed, 81 insertions(+), 32 deletions(-) diff --git a/README.md b/README.md index 48a1dd0..5729606 100644 --- a/README.md +++ b/README.md @@ -157,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 @@ -195,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 | | ----------------------- | ---------- | @@ -242,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. diff --git a/internal/config/config.go b/internal/config/config.go index 69c5210..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 @@ -746,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), ","), ) diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 0add1b2..a7cc76e 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -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") }) } } 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 -- 2.54.0 From f0adeafde3ed98d42caad46ac70fe1a88d538dae Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 10:30:26 +0200 Subject: [PATCH 15/15] Drop Set-Cookie from the recovered 500 (closes #193) When a handler sets a cookie and then panics before sending anything, the recover middleware now deletes Set-Cookie before writing its 500, so a request that failed never hands the client a credential. Every other header, Location included, is left as http.Error leaves it, matching chi's Recoverer. A response that was already sent is untouched. Tests cover the uncommitted case (no cookie, Location kept) and assert the cookie still reaches the client when the response was committed before the panic. Model: opus-5-5 --- internal/middleware/recoverer.go | 7 ++++++ internal/middleware/recoverer_test.go | 33 +++++++++++++++++++++++++-- 2 files changed, 38 insertions(+), 2 deletions(-) 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"]) -- 2.54.0