Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
83740b1de1 | ||
|
|
aeeeca5ea1 | ||
|
|
6ebac4fa71 | ||
|
|
237f131367 | ||
|
|
7ed1588443 | ||
|
|
b051821370 | ||
|
|
39afa69bfc | ||
|
|
888eaf526b |
@@ -7,6 +7,13 @@ services, durably stores them, and delivers them to configured targets
|
|||||||
with retry support, logging, and observability. Category: infrastructure
|
with retry support, logging, and observability. Category: infrastructure
|
||||||
/ web service. License: MIT.
|
/ web service. License: MIT.
|
||||||
|
|
||||||
|
Each entrypoint is a version 4 UUID served at `/webhook/{uuid}`, and
|
||||||
|
that UUID is the entrypoint's only credential. webhooker does not use
|
||||||
|
shared secrets, HMAC signatures or token headers on the receiver, and
|
||||||
|
will not add them — read
|
||||||
|
[The entrypoint URL is the authentication secret](#the-entrypoint-url-is-the-authentication-secret)
|
||||||
|
before deploying one.
|
||||||
|
|
||||||
## Getting Started
|
## Getting Started
|
||||||
|
|
||||||
### Prerequisites
|
### Prerequisites
|
||||||
@@ -37,9 +44,9 @@ make bootstrap
|
|||||||
# Run all checks (test, lint, format check)
|
# Run all checks (test, lint, format check)
|
||||||
make check
|
make check
|
||||||
|
|
||||||
# Run in development mode. DATA_DIR defaults to /var/lib/webhooker in
|
# Run the server from the clone. DATA_DIR defaults to
|
||||||
# every environment, so set it (in .env or the shell) to a writable
|
# /var/lib/webhooker in every environment, so set it (in .env or the
|
||||||
# directory when running from a clone.
|
# shell) to a writable directory.
|
||||||
DATA_DIR=./data make dev
|
DATA_DIR=./data make dev
|
||||||
|
|
||||||
# Build Docker image
|
# Build Docker image
|
||||||
@@ -85,7 +92,8 @@ them at once. A variable already present in the real environment wins
|
|||||||
over the file's value for the same name.
|
over the file's value for the same name.
|
||||||
|
|
||||||
The environment is selected by setting `WEBHOOKER_ENVIRONMENT` to `dev`
|
The environment is selected by setting `WEBHOOKER_ENVIRONMENT` to `dev`
|
||||||
or `prod` (default: `dev`). The setting controls exactly one behavior:
|
or `prod` (default: `prod`; `dev` must be set explicitly). The setting
|
||||||
|
controls exactly one behavior:
|
||||||
|
|
||||||
| Behavior | `dev` | `prod` |
|
| Behavior | `dev` | `prod` |
|
||||||
| -------- | ----------------------- | ---------------- |
|
| -------- | ----------------------- | ---------------- |
|
||||||
@@ -127,7 +135,7 @@ TTY detection, and security headers are always applied.
|
|||||||
|
|
||||||
| Variable | Description | Default |
|
| Variable | Description | Default |
|
||||||
| ----------------------- | ----------------------------------- | -------- |
|
| ----------------------- | ----------------------------------- | -------- |
|
||||||
| `WEBHOOKER_ENVIRONMENT` | `dev` or `prod` | `dev` |
|
| `WEBHOOKER_ENVIRONMENT` | `dev` or `prod` | `prod` |
|
||||||
| `PORT` | HTTP listen port | `8080` |
|
| `PORT` | HTTP listen port | `8080` |
|
||||||
| `BIND_ADDRESS` | IP address the HTTP listener binds. Loopback by default, so the cleartext listener is not published on every interface. The Docker image ships `0.0.0.0` instead. See [Bind address](#bind-address) | `127.0.0.1` (image: `0.0.0.0`) |
|
| `BIND_ADDRESS` | IP address the HTTP listener binds. Loopback by default, so the cleartext listener is not published on every interface. The Docker image ships `0.0.0.0` instead. See [Bind address](#bind-address) | `127.0.0.1` (image: `0.0.0.0`) |
|
||||||
| `DATA_DIR` | Directory for all SQLite databases | `/var/lib/webhooker` |
|
| `DATA_DIR` | Directory for all SQLite databases | `/var/lib/webhooker` |
|
||||||
@@ -392,10 +400,9 @@ the bucket is. See [Rate Limiting](#rate-limiting).
|
|||||||
|
|
||||||
The remedy is to set `TRUSTED_PROXIES` to your reverse proxy's
|
The remedy is to set `TRUSTED_PROXIES` to your reverse proxy's
|
||||||
address, which restores per-client buckets. webhooker logs a warning
|
address, which restores per-client buckets. webhooker logs a warning
|
||||||
at startup whenever `TRUSTED_PROXIES` is empty, in every environment —
|
at startup whenever `TRUSTED_PROXIES` is empty, in every environment,
|
||||||
not only when `WEBHOOKER_ENVIRONMENT=prod`, because that variable
|
because behind a proxy every client shares one bucket in `dev` and
|
||||||
defaults to `dev` and an operator who never set it is precisely the
|
`prod` alike. The warning is informational when nothing proxies to the
|
||||||
one at risk. The warning is informational when nothing proxies to the
|
|
||||||
process: with no proxy in front, the peer address is the client's own
|
process: with no proxy in front, the peer address is the client's own
|
||||||
and the buckets are already per-client. See
|
and the buckets are already per-client. See
|
||||||
[Rate Limiting](#rate-limiting) for what each limit shares.
|
[Rate Limiting](#rate-limiting) for what each limit shares.
|
||||||
@@ -631,7 +638,6 @@ decision:
|
|||||||
docker run -d \
|
docker run -d \
|
||||||
-p 127.0.0.1:8080:8080 \
|
-p 127.0.0.1:8080:8080 \
|
||||||
-v /path/to/data:/var/lib/webhooker \
|
-v /path/to/data:/var/lib/webhooker \
|
||||||
-e WEBHOOKER_ENVIRONMENT=prod \
|
|
||||||
-e BIND_ADDRESS=0.0.0.0 \
|
-e BIND_ADDRESS=0.0.0.0 \
|
||||||
webhooker:latest
|
webhooker:latest
|
||||||
```
|
```
|
||||||
@@ -724,6 +730,66 @@ listing the directory and learning your webhook UUIDs from the
|
|||||||
`events-{uuid}.db` filenames — not the barrier protecting the
|
`events-{uuid}.db` filenames — not the barrier protecting the
|
||||||
credentials.
|
credentials.
|
||||||
|
|
||||||
|
### Running under upaas
|
||||||
|
|
||||||
|
[upaas](https://git.eeqj.de/sneak/upaas) builds the image from this
|
||||||
|
repository's `Dockerfile` and runs it. The app needs:
|
||||||
|
|
||||||
|
- **Network and port:** add no port mapping in upaas. upaas publishes
|
||||||
|
every mapped port on all interfaces of the host
|
||||||
|
([upaas issue 113](https://git.eeqj.de/sneak/upaas/issues/113)),
|
||||||
|
which would put the plain-HTTP admin UI and receiver there. Instead,
|
||||||
|
set the app's Docker Network in upaas to your reverse proxy's Docker
|
||||||
|
network; the proxy then reaches the app at `upaas-` followed by the
|
||||||
|
app name, port `8080`. Leave `PORT` unset: the image's health check
|
||||||
|
probes `8080`.
|
||||||
|
- **Volume:** one host directory mounted at `/var/lib/webhooker`.
|
||||||
|
upaas bind-mounts the host path it is given and does not create it,
|
||||||
|
and the container does not start unless UID 1000 owns it (see
|
||||||
|
[Running with Docker](#running-with-docker)). Create it before the
|
||||||
|
first deploy:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
mkdir -p /path/to/data
|
||||||
|
chown 1000:1000 /path/to/data
|
||||||
|
chmod 750 /path/to/data
|
||||||
|
```
|
||||||
|
|
||||||
|
- **Environment variables:**
|
||||||
|
- `WEBHOOKER_ENVIRONMENT=prod`
|
||||||
|
- `TRUSTED_PROXIES`: your reverse proxy's address on that Docker
|
||||||
|
network. The `remoteIP` field of the `http request` log line for a
|
||||||
|
request that came through the proxy shows it; the health check's
|
||||||
|
own lines show `::1`. See [Trusted proxies](#trusted-proxies).
|
||||||
|
- Leave `BIND_ADDRESS` and `DATA_DIR` unset: the image sets
|
||||||
|
`BIND_ADDRESS` to `0.0.0.0`, and `DATA_DIR` defaults to
|
||||||
|
`/var/lib/webhooker`.
|
||||||
|
- Everything else is optional; see [Configuration](#configuration).
|
||||||
|
- **Health check:** the image's own, which requests
|
||||||
|
`/.well-known/healthcheck`. upaas reads the container's health 60
|
||||||
|
seconds after a deploy and marks the deploy failed unless it is
|
||||||
|
`healthy`.
|
||||||
|
- **First run:** the first start prints the `admin` password once, in
|
||||||
|
the banner described under [The admin account](#the-admin-account),
|
||||||
|
to the container's log. upaas names the container `upaas-` followed
|
||||||
|
by the app name, so for an app named `webhooker`:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
docker logs upaas-webhooker
|
||||||
|
```
|
||||||
|
|
||||||
|
If the password is lost, stop the container, set a new password with
|
||||||
|
the app's own image and volume, and start it again (see
|
||||||
|
[Recovering a lost admin password](#recovering-a-lost-admin-password)):
|
||||||
|
|
||||||
|
```bash
|
||||||
|
docker stop upaas-webhooker
|
||||||
|
docker run --rm --volumes-from upaas-webhooker \
|
||||||
|
"$(docker inspect -f '{{.Image}}' upaas-webhooker)" \
|
||||||
|
/app/webhooker resetpw -generate admin
|
||||||
|
docker start upaas-webhooker
|
||||||
|
```
|
||||||
|
|
||||||
## Deployment behind a reverse proxy
|
## Deployment behind a reverse proxy
|
||||||
|
|
||||||
webhooker terminates no TLS of its own. It serves plaintext HTTP and
|
webhooker terminates no TLS of its own. It serves plaintext HTTP and
|
||||||
@@ -745,17 +811,18 @@ reports.
|
|||||||
serves the admin login form and the unauthenticated receiver with
|
serves the admin login form and the unauthenticated receiver with
|
||||||
no TLS at all, and the proxy in front of it changes nothing about
|
no TLS at all, and the proxy in front of it changes nothing about
|
||||||
that.
|
that.
|
||||||
2. **Set `WEBHOOKER_ENVIRONMENT=prod`, and make sure the proxy sends
|
2. **Make sure the environment is not `dev` (leave
|
||||||
`X-Forwarded-Proto`.** These are two requirements, not one. The
|
`WEBHOOKER_ENVIRONMENT` unset or set it to `prod`), and make sure
|
||||||
environment setting decides CORS and nothing else: the default
|
the proxy sends `X-Forwarded-Proto`.** These are two requirements,
|
||||||
|
not one. The environment setting decides CORS and nothing else:
|
||||||
`dev` answers every origin with `Access-Control-Allow-Origin: *`
|
`dev` answers every origin with `Access-Control-Allow-Origin: *`
|
||||||
(without credentials), which a server-rendered production
|
(without credentials), which a server-rendered production
|
||||||
deployment has no use for. Cookie `Secure` and the strict
|
deployment has no use for, and `prod` — the default — disables it.
|
||||||
Origin/Referer mode are **not** tied to it — they are decided per
|
Cookie `Secure` and the strict Origin/Referer mode are **not** tied
|
||||||
request from the transport, which behind a proxy means the
|
to it — they are decided per request from the transport, which
|
||||||
`X-Forwarded-Proto` header. The block below sets it; without it
|
behind a proxy means the `X-Forwarded-Proto` header. The block below
|
||||||
every request is read as plaintext and cookies ship without
|
sets it; without it every request is read as plaintext and cookies
|
||||||
`Secure`. See [Configuration](#configuration).
|
ship without `Secure`. See [Configuration](#configuration).
|
||||||
3. **Set `TRUSTED_PROXIES` to the proxy's address.** Unset, every rate
|
3. **Set `TRUSTED_PROXIES` to the proxy's address.** Unset, every rate
|
||||||
limiter keys on the connecting peer, which behind a proxy is the
|
limiter keys on the connecting peer, which behind a proxy is the
|
||||||
proxy on every request: all clients collapse into one global bucket
|
proxy on every request: all clients collapse into one global bucket
|
||||||
@@ -847,7 +914,6 @@ sent — `$scheme` above does.
|
|||||||
With that block, webhooker's environment is:
|
With that block, webhooker's environment is:
|
||||||
|
|
||||||
```sh
|
```sh
|
||||||
WEBHOOKER_ENVIRONMENT=prod
|
|
||||||
BIND_ADDRESS=127.0.0.1 # the default; stated here to be explicit
|
BIND_ADDRESS=127.0.0.1 # the default; stated here to be explicit
|
||||||
TRUSTED_PROXIES=127.0.0.1
|
TRUSTED_PROXIES=127.0.0.1
|
||||||
```
|
```
|
||||||
@@ -1149,14 +1215,38 @@ backups at rest and restrict who can read them.
|
|||||||
|
|
||||||
## The entrypoint URL is the authentication secret
|
## The entrypoint URL is the authentication secret
|
||||||
|
|
||||||
The receiver verifies nothing about an inbound request. The UUID in an
|
**The entrypoint UUID is the credential, and it is the only one.**
|
||||||
entrypoint's URL is its credential: anyone who holds that URL can
|
webhooker mints a version 4 UUID per entrypoint and serves it at
|
||||||
submit events to it, and the receiver checks nothing else about the
|
`/webhook/{uuid}`. Possession of that URL is the authentication:
|
||||||
sender. Treat an entrypoint URL the way you would treat an API token.
|
anyone who holds it can submit events to the entrypoint, and the
|
||||||
|
receiver verifies nothing else about the sender.
|
||||||
|
|
||||||
There is no way to rotate the UUID in place. To retire one, delete the
|
There is no shared secret, no HMAC signature, no bearer token and no
|
||||||
entrypoint (or deactivate it, which answers `410`) and create a new
|
second factor on the receiver, and none will be added. This was
|
||||||
one, then point the sender at the new URL.
|
considered and rejected; the implementation that existed was removed
|
||||||
|
in [PR #279](https://git.eeqj.de/sneak/webhooker/pulls/279), closing
|
||||||
|
[issue #67](https://git.eeqj.de/sneak/webhooker/issues/67) and
|
||||||
|
[issue #241](https://git.eeqj.de/sneak/webhooker/issues/241). A
|
||||||
|
proposal to reintroduce any of them — including as "defence in depth"
|
||||||
|
alongside the UUID — is answered by this section. Inbound signature
|
||||||
|
headers a sender sends anyway (`X-Hub-Signature` and its
|
||||||
|
per-provider equivalents) are stored and forwarded as ordinary
|
||||||
|
headers; nothing checks them.
|
||||||
|
|
||||||
|
What that means for an operator:
|
||||||
|
|
||||||
|
- **The URL is a capability, so treat it as a secret.** Keep it out of
|
||||||
|
logs, ticket bodies, chat messages and screenshots. Anyone who reads
|
||||||
|
it anywhere can post events as that sender.
|
||||||
|
- **Rotating means minting a new entrypoint, not changing a key.**
|
||||||
|
There is no way to rotate the UUID in place. To retire one, delete
|
||||||
|
the entrypoint (or deactivate it, which answers `410`) and create a
|
||||||
|
new one, then point the sender at the new URL.
|
||||||
|
- **A sender that cannot be given a secret URL is a constraint on that
|
||||||
|
integration, not a reason to change this.** If a service only
|
||||||
|
supports signed payloads to a well-known URL, raise it as its own
|
||||||
|
problem — pick a different integration path, or accept that it
|
||||||
|
cannot be used. It is not grounds to reintroduce shared secrets.
|
||||||
|
|
||||||
## Entrypoints
|
## Entrypoints
|
||||||
|
|
||||||
@@ -1476,7 +1566,7 @@ events should be forwarded.
|
|||||||
| `type` | TargetType | One of: `http`, `slack`, `database`, `log` |
|
| `type` | TargetType | One of: `http`, `slack`, `database`, `log` |
|
||||||
| `active` | boolean | Whether deliveries are enabled (default: true) |
|
| `active` | boolean | Whether deliveries are enabled (default: true) |
|
||||||
| `config` | JSON text | Type-specific configuration |
|
| `config` | JSON text | Type-specific configuration |
|
||||||
| `max_retries` | integer | Maximum retry attempts for `http` and `slack` targets (0 = fire-and-forget, >0 = retries with backoff and a circuit breaker). Ignored by `database` and `log` targets |
|
| `max_retries` | integer | Total delivery attempts for `http` and `slack` targets, not retries on top of the first: 0 is a single fire-and-forget attempt with no retries and no circuit breaker, and a value of N makes N attempts in all, with exponential backoff and a per-target circuit breaker. Ignored by `database` and `log` targets |
|
||||||
| `max_queue_size` | integer | Stored and shown on the target's detail view, but not enforced anywhere yet: nothing in the delivery engine consults it. Queue depth is set by the two fixed 10,000-entry channels |
|
| `max_queue_size` | integer | Stored and shown on the target's detail view, but not enforced anywhere yet: nothing in the delivery engine consults it. Queue depth is set by the two fixed 10,000-entry channels |
|
||||||
|
|
||||||
**Relations:** Belongs to Webhook. Has many Deliveries.
|
**Relations:** Belongs to Webhook. Has many Deliveries.
|
||||||
@@ -1484,12 +1574,12 @@ events should be forwarded.
|
|||||||
**Target types:**
|
**Target types:**
|
||||||
|
|
||||||
- **`http`** — Forward the event as an HTTP POST to a configured URL.
|
- **`http`** — Forward the event as an HTTP POST to a configured URL.
|
||||||
Behavior depends on `max_retries`: when `max_retries` is 0 (the
|
`max_retries` is the total number of delivery attempts, not retries on
|
||||||
default), the target operates in fire-and-forget mode — a single
|
top of the first: when `max_retries` is 0 (the default), the target
|
||||||
attempt with no retries and no circuit breaker. When `max_retries` is
|
operates in fire-and-forget mode, a single attempt with no retries and
|
||||||
greater than 0, failed deliveries are retried with exponential backoff
|
no circuit breaker; a value of N makes up to N attempts in all,
|
||||||
up to `max_retries` attempts, protected by a per-target circuit
|
retrying failed deliveries with exponential backoff and protecting them
|
||||||
breaker.
|
with a per-target circuit breaker.
|
||||||
- **`slack`** — Post the event as a formatted message to a
|
- **`slack`** — Post the event as a formatted message to a
|
||||||
Slack-compatible incoming webhook URL (`webhookUrl` in `config`). It
|
Slack-compatible incoming webhook URL (`webhookUrl` in `config`). It
|
||||||
is built on the same HTTP core as `http` and honours `max_retries`
|
is built on the same HTTP core as `http` and honours `max_retries`
|
||||||
@@ -1669,6 +1759,28 @@ retries) is individually logged for full observability.
|
|||||||
|
|
||||||
**Relations:** Belongs to Delivery.
|
**Relations:** Belongs to Delivery.
|
||||||
|
|
||||||
|
#### Event-tier indexes
|
||||||
|
|
||||||
|
These indexes on the per-webhook event databases are declared in the model
|
||||||
|
tags, so `AutoMigrate` creates them on a fresh and on an existing database:
|
||||||
|
|
||||||
|
| Table | Columns | Serves |
|
||||||
|
| ------------------ | --------------------------- | ------ |
|
||||||
|
| `deliveries` | `status`, `deleted_at` | Startup recovery, the retry and pending sweeps every 60 seconds and the queue-depth sampler every 30 seconds, which select deliveries by status |
|
||||||
|
| `deliveries` | `event_id`, `deleted_at` | The event log, which loads each event's deliveries, and retention, which selects and deletes the deliveries of expired events |
|
||||||
|
| `delivery_results` | `delivery_id`, `deleted_at` | The event log, which loads the attempts of a page's deliveries, and retention, which deletes the attempts of expired events |
|
||||||
|
| `events` | `deleted_at`, `created_at` | Retention, which selects expired events by age |
|
||||||
|
| `events` | `created_at` | Retention's delete of the expired events themselves |
|
||||||
|
|
||||||
|
GORM's soft delete adds `deleted_at IS NULL` to these queries; retention's
|
||||||
|
deletes leave it out, but their lookups of expired rows keep it. SQLite keeps
|
||||||
|
no statistics on these tables, and without them it rates the `deleted_at`
|
||||||
|
index, which every live row matches, above an index on a column matched
|
||||||
|
against several values or compared with `<`. So every index but the last also
|
||||||
|
covers `deleted_at`. It comes second, so that retention's deletes can use the
|
||||||
|
index without it, except in `events`, where `created_at` is compared with `<`
|
||||||
|
and SQLite narrows by a `<` only on the last column it uses.
|
||||||
|
|
||||||
#### Common Fields
|
#### Common Fields
|
||||||
|
|
||||||
Every entity except `Setting` includes these fields from `BaseModel`.
|
Every entity except `Setting` includes these fields from `BaseModel`.
|
||||||
@@ -2867,6 +2979,10 @@ check, see [The login endpoint](#the-login-endpoint).
|
|||||||
|
|
||||||
### Authentication
|
### Authentication
|
||||||
|
|
||||||
|
- **Webhook receiver:** the entrypoint UUID in the URL, and nothing
|
||||||
|
else. No shared secret, no HMAC signature, no token header, and none
|
||||||
|
will be added — see
|
||||||
|
[The entrypoint URL is the authentication secret](#the-entrypoint-url-is-the-authentication-secret).
|
||||||
- **Web UI:** Cookie-based sessions using gorilla/sessions with
|
- **Web UI:** Cookie-based sessions using gorilla/sessions with
|
||||||
encrypted cookies. Sessions are configured with HttpOnly, SameSite
|
encrypted cookies. Sessions are configured with HttpOnly, SameSite
|
||||||
Lax, and Secure whenever the request is on TLS — the flag follows the
|
Lax, and Secure whenever the request is on TLS — the flag follows the
|
||||||
@@ -2906,7 +3022,8 @@ check, see [The login endpoint](#the-login-endpoint).
|
|||||||
mode
|
mode
|
||||||
- **The entrypoint URL is the receiver's only credential.** Nothing
|
- **The entrypoint URL is the receiver's only credential.** Nothing
|
||||||
about an inbound request is verified; possession of the UUID
|
about an inbound request is verified; possession of the UUID
|
||||||
authorises submission (see
|
authorises submission, and no shared secret or signature check will
|
||||||
|
be added alongside it (see
|
||||||
[The entrypoint URL is the authentication secret](#the-entrypoint-url-is-the-authentication-secret))
|
[The entrypoint URL is the authentication secret](#the-entrypoint-url-is-the-authentication-secret))
|
||||||
- **SSRF prevention** for HTTP delivery targets: private/reserved IP
|
- **SSRF prevention** for HTTP delivery targets: private/reserved IP
|
||||||
ranges (RFC 1918, loopback, link-local, cloud metadata) are blocked
|
ranges (RFC 1918, loopback, link-local, cloud metadata) are blocked
|
||||||
|
|||||||
@@ -585,12 +585,14 @@ func resolveMetricsAuth() (string, string, error) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// resolveEnvironment reads WEBHOOKER_ENVIRONMENT, defaulting to
|
// resolveEnvironment reads WEBHOOKER_ENVIRONMENT, defaulting to prod
|
||||||
// dev, and rejects unrecognised values.
|
// when it is unset so a deployment that forgets the variable is not
|
||||||
|
// silently permissive; dev must be set explicitly. It rejects
|
||||||
|
// unrecognised values.
|
||||||
func resolveEnvironment() (string, error) {
|
func resolveEnvironment() (string, error) {
|
||||||
environment := os.Getenv("WEBHOOKER_ENVIRONMENT")
|
environment := os.Getenv("WEBHOOKER_ENVIRONMENT")
|
||||||
if environment == "" {
|
if environment == "" {
|
||||||
environment = EnvironmentDev
|
environment = EnvironmentProd
|
||||||
}
|
}
|
||||||
|
|
||||||
if environment != EnvironmentDev &&
|
if environment != EnvironmentDev &&
|
||||||
@@ -772,10 +774,8 @@ func (c *Config) warnEgressAllowlist(log *slog.Logger) {
|
|||||||
// everyone else's wrong passwords, and the receiver's limits become
|
// everyone else's wrong passwords, and the receiver's limits become
|
||||||
// service-wide ceilings.
|
// service-wide ceilings.
|
||||||
//
|
//
|
||||||
// The warning is deliberately not gated on WEBHOOKER_ENVIRONMENT. That
|
// The warning is deliberately not gated on WEBHOOKER_ENVIRONMENT:
|
||||||
// variable defaults to dev, so gating on it would silence the warning
|
// behind a proxy every client shares one bucket in dev and prod alike.
|
||||||
// for exactly the operator who forgot to configure the deployment —
|
|
||||||
// the case it exists to catch.
|
|
||||||
//
|
//
|
||||||
// The default of trusting nobody is deliberate — trusting forwarded
|
// The default of trusting nobody is deliberate — trusting forwarded
|
||||||
// headers from arbitrary peers lets any client choose its own bucket —
|
// headers from arbitrary peers lets any client choose its own bucket —
|
||||||
|
|||||||
@@ -44,9 +44,9 @@ func TestEnvironmentConfig(t *testing.T) {
|
|||||||
isProd bool
|
isProd bool
|
||||||
}{
|
}{
|
||||||
{
|
{
|
||||||
name: "default is dev",
|
name: "default is prod",
|
||||||
isDev: true,
|
isDev: false,
|
||||||
isProd: false,
|
isProd: true,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
name: "explicit dev",
|
name: "explicit dev",
|
||||||
@@ -848,10 +848,9 @@ func TestEgressAllowlistWarning(t *testing.T) {
|
|||||||
// tells an operator a deployment behind a reverse proxy shares one
|
// tells an operator a deployment behind a reverse proxy shares one
|
||||||
// rate-limit bucket between every client, which turns the receiver
|
// rate-limit bucket between every client, which turns the receiver
|
||||||
// limits into service-wide ceilings and collapses login failure
|
// limits into service-wide ceilings and collapses login failure
|
||||||
// counting. It must fire whenever TRUSTED_PROXIES is empty,
|
// counting. It must fire whenever TRUSTED_PROXIES is empty, in any
|
||||||
// in any environment: WEBHOOKER_ENVIRONMENT defaults to dev, so gating
|
// environment, because behind a proxy every client shares one bucket
|
||||||
// on it would silence the warning for exactly the operator who never
|
// in dev and prod alike. It stays quiet once proxies are named.
|
||||||
// configured the deployment. It stays quiet once proxies are named.
|
|
||||||
func TestSharedRateLimitBucketWarning(t *testing.T) {
|
func TestSharedRateLimitBucketWarning(t *testing.T) {
|
||||||
tests := []struct {
|
tests := []struct {
|
||||||
name string
|
name string
|
||||||
@@ -871,10 +870,6 @@ func TestSharedRateLimitBucketWarning(t *testing.T) {
|
|||||||
expectWarning: false,
|
expectWarning: false,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
// The default environment. An internet-exposed
|
|
||||||
// deployment whose operator never set
|
|
||||||
// WEBHOOKER_ENVIRONMENT lands here and has exactly
|
|
||||||
// the exposure the warning announces.
|
|
||||||
name: "dev without trusted proxies warns",
|
name: "dev without trusted proxies warns",
|
||||||
environment: config.EnvironmentDev,
|
environment: config.EnvironmentDev,
|
||||||
expectWarning: true,
|
expectWarning: true,
|
||||||
|
|||||||
@@ -17,12 +17,12 @@ import (
|
|||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
"sneak.berlin/go/webhooker/internal/banner"
|
"sneak.berlin/go/webhooker/internal/banner"
|
||||||
"sneak.berlin/go/webhooker/internal/config"
|
"sneak.berlin/go/webhooker/internal/config"
|
||||||
|
"sneak.berlin/go/webhooker/internal/datadir"
|
||||||
"sneak.berlin/go/webhooker/internal/gormlog"
|
"sneak.berlin/go/webhooker/internal/gormlog"
|
||||||
"sneak.berlin/go/webhooker/internal/logger"
|
"sneak.berlin/go/webhooker/internal/logger"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
dataDirPerm = 0750
|
|
||||||
randomPasswordLen = 16
|
randomPasswordLen = 16
|
||||||
sessionKeyLen = 32
|
sessionKeyLen = 32
|
||||||
)
|
)
|
||||||
@@ -185,7 +185,9 @@ func (d *Database) connect() error {
|
|||||||
// caller's decision.
|
// caller's decision.
|
||||||
func (d *Database) connectTo(dataDir string) error {
|
func (d *Database) connectTo(dataDir string) error {
|
||||||
// Ensure the data directory exists before opening the database.
|
// Ensure the data directory exists before opening the database.
|
||||||
err := os.MkdirAll(dataDir, dataDirPerm)
|
// datadir.DirPerm is the single source of the directory mode; this
|
||||||
|
// package creates the directory too, since either may run first.
|
||||||
|
err := os.MkdirAll(dataDir, datadir.DirPerm)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf(
|
return fmt.Errorf(
|
||||||
"creating data directory %s: %w",
|
"creating data directory %s: %w",
|
||||||
|
|||||||
@@ -0,0 +1,166 @@
|
|||||||
|
package database_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/google/uuid"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
"gorm.io/gorm"
|
||||||
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
|
)
|
||||||
|
|
||||||
|
// TestWebhookDBManager_OpenAddsEventTierIndexes verifies that opening a
|
||||||
|
// per-webhook database that predates these indexes creates them. It
|
||||||
|
// stands in for an older database file by dropping the indexes
|
||||||
|
// AutoMigrate just created, then reopening the same file.
|
||||||
|
func TestWebhookDBManager_OpenAddsEventTierIndexes(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
indexes := []struct {
|
||||||
|
model any
|
||||||
|
name string
|
||||||
|
}{
|
||||||
|
{&database.Delivery{}, "idx_deliveries_status"},
|
||||||
|
{&database.Delivery{}, "idx_deliveries_event_id"},
|
||||||
|
{&database.DeliveryResult{}, "idx_delivery_results_delivery_id"},
|
||||||
|
{&database.Event{}, "idx_events_deleted_at_created_at"},
|
||||||
|
{&database.Event{}, "idx_events_created_at"},
|
||||||
|
}
|
||||||
|
|
||||||
|
mgr, lc := setupTestWebhookDBManager(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
require.NoError(t, lc.Start(ctx))
|
||||||
|
|
||||||
|
defer func() { require.NoError(t, lc.Stop(ctx)) }()
|
||||||
|
|
||||||
|
webhookID := uuid.New().String()
|
||||||
|
|
||||||
|
db, err := mgr.GetDB(webhookID)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// A fresh database has them.
|
||||||
|
for _, ix := range indexes {
|
||||||
|
require.True(t, db.Migrator().HasIndex(ix.model, ix.name))
|
||||||
|
}
|
||||||
|
|
||||||
|
// Stand in for a database file created before the indexes existed.
|
||||||
|
for _, ix := range indexes {
|
||||||
|
require.NoError(t, db.Migrator().DropIndex(ix.model, ix.name))
|
||||||
|
require.False(t, db.Migrator().HasIndex(ix.model, ix.name))
|
||||||
|
}
|
||||||
|
|
||||||
|
// Drop the cached connection so the next open reopens the file and
|
||||||
|
// runs AutoMigrate against it, as a restart would.
|
||||||
|
require.NoError(t, mgr.CloseAll())
|
||||||
|
|
||||||
|
db, err = mgr.GetDB(webhookID)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
for _, ix := range indexes {
|
||||||
|
assert.True(t, db.Migrator().HasIndex(ix.model, ix.name),
|
||||||
|
"opening the existing database should create %s", ix.name)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestEventTierQueriesUseTheirIndexes verifies that the statements the
|
||||||
|
// indexes are for use them. GORM builds each statement in a dry run as
|
||||||
|
// the code named above it does, soft-delete condition included, and
|
||||||
|
// SQLite, which keeps no statistics on these tables, must plan to seek
|
||||||
|
// on each index listed by the columns in parentheses.
|
||||||
|
func TestEventTierQueriesUseTheirIndexes(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
mgr, lc := setupTestWebhookDBManager(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
require.NoError(t, lc.Start(ctx))
|
||||||
|
|
||||||
|
defer func() { require.NoError(t, lc.Stop(ctx)) }()
|
||||||
|
|
||||||
|
db, err := mgr.GetDB(uuid.New().String())
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
dry := db.Session(&gorm.Session{DryRun: true})
|
||||||
|
ids := []string{
|
||||||
|
uuid.New().String(), uuid.New().String(), uuid.New().String(),
|
||||||
|
}
|
||||||
|
cutoff := time.Now()
|
||||||
|
|
||||||
|
var (
|
||||||
|
deliveries []database.Delivery
|
||||||
|
results []database.DeliveryResult
|
||||||
|
depths []struct{ Depth int }
|
||||||
|
)
|
||||||
|
|
||||||
|
byStatus := "idx_deliveries_status (status=? AND deleted_at=?)"
|
||||||
|
byEvent := "idx_deliveries_event_id (event_id=? AND deleted_at=?)"
|
||||||
|
byAge := "idx_events_deleted_at_created_at (deleted_at=? AND created_at<?)"
|
||||||
|
|
||||||
|
// The delivery engine: recovery and the retry sweep, the sweep for
|
||||||
|
// stranded pending deliveries, and the queue depth count.
|
||||||
|
assertPlanUses(t, db, dry.Where(
|
||||||
|
"status = ?", database.DeliveryStatusRetrying,
|
||||||
|
).Find(&deliveries), byStatus)
|
||||||
|
assertPlanUses(t, db, dry.Where(
|
||||||
|
"status = ? AND updated_at < ?",
|
||||||
|
database.DeliveryStatusPending, cutoff,
|
||||||
|
).Limit(500).Find(&deliveries), byStatus)
|
||||||
|
assertPlanUses(t, db, dry.Model(&database.Delivery{}).
|
||||||
|
Select("target_id", "status", "count(*) as depth").
|
||||||
|
Where("status IN ?", []database.DeliveryStatus{
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
database.DeliveryStatusRetrying,
|
||||||
|
}).Group("target_id, status").Find(&depths), byStatus)
|
||||||
|
|
||||||
|
// The event log: each event's deliveries, then their attempts
|
||||||
|
// (loadEventsWithDeliveries, loadDeliveryResults).
|
||||||
|
assertPlanUses(t, db, dry.Where("event_id = ?", ids[0]).
|
||||||
|
Find(&deliveries), byEvent)
|
||||||
|
assertPlanUses(t, db, dry.Where("delivery_id IN ?", ids).
|
||||||
|
Order("attempt_num ASC").Find(&results),
|
||||||
|
"idx_delivery_results_delivery_id (delivery_id=? AND deleted_at=?)")
|
||||||
|
|
||||||
|
// Retention's three deletes (reapExpired), whose subqueries are built
|
||||||
|
// afresh for each statement as it builds them.
|
||||||
|
expiredEventIDs := func() *gorm.DB {
|
||||||
|
return dry.Model(&database.Event{}).Select("id").
|
||||||
|
Where("created_at < ?", cutoff)
|
||||||
|
}
|
||||||
|
|
||||||
|
assertPlanUses(t, db, dry.Unscoped().Where(
|
||||||
|
"delivery_id IN (?)", dry.Model(&database.Delivery{}).
|
||||||
|
Select("id").Where("event_id IN (?)", expiredEventIDs()),
|
||||||
|
).Delete(&database.DeliveryResult{}),
|
||||||
|
"idx_delivery_results_delivery_id (delivery_id=?)", byEvent, byAge)
|
||||||
|
assertPlanUses(t, db, dry.Unscoped().Where(
|
||||||
|
"event_id IN (?)", expiredEventIDs(),
|
||||||
|
).Delete(&database.Delivery{}),
|
||||||
|
"idx_deliveries_event_id (event_id=?)", byAge)
|
||||||
|
assertPlanUses(t, db, dry.Unscoped().Where(
|
||||||
|
"created_at < ?", cutoff,
|
||||||
|
).Delete(&database.Event{}), "idx_events_created_at (created_at<?)")
|
||||||
|
}
|
||||||
|
|
||||||
|
// assertPlanUses asserts that SQLite's plan for a statement GORM built
|
||||||
|
// in a dry run, run with the same SQL and arguments GORM would send,
|
||||||
|
// names each of the given indexes.
|
||||||
|
func assertPlanUses(
|
||||||
|
t *testing.T, db, built *gorm.DB, indexes ...string,
|
||||||
|
) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
var plan []struct{ Detail string }
|
||||||
|
|
||||||
|
require.NoError(t, db.Raw(
|
||||||
|
"EXPLAIN QUERY PLAN "+built.Statement.SQL.String(),
|
||||||
|
built.Statement.Vars...,
|
||||||
|
).Scan(&plan).Error)
|
||||||
|
|
||||||
|
for _, index := range indexes {
|
||||||
|
assert.Contains(t, fmt.Sprint(plan), index,
|
||||||
|
built.Statement.SQL.String())
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,5 +1,7 @@
|
|||||||
package database
|
package database
|
||||||
|
|
||||||
|
import "gorm.io/gorm"
|
||||||
|
|
||||||
// DeliveryStatus represents the status of a delivery
|
// DeliveryStatus represents the status of a delivery
|
||||||
type DeliveryStatus string
|
type DeliveryStatus string
|
||||||
|
|
||||||
@@ -29,12 +31,19 @@ func (s DeliveryStatus) Terminal() bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Delivery represents a delivery attempt for an event to a target
|
// Delivery represents a delivery attempt for an event to a target
|
||||||
|
//
|
||||||
|
//nolint:lll // a struct tag cannot wrap
|
||||||
type Delivery struct {
|
type Delivery struct {
|
||||||
BaseModel
|
BaseModel
|
||||||
|
|
||||||
EventID string `gorm:"type:uuid;not null" json:"eventId"`
|
EventID string `gorm:"type:uuid;not null;index:idx_deliveries_event_id,priority:1" json:"eventId"`
|
||||||
TargetID string `gorm:"type:uuid;not null" json:"targetId"`
|
TargetID string `gorm:"type:uuid;not null" json:"targetId"`
|
||||||
Status DeliveryStatus `gorm:"not null;default:'pending'" json:"status"`
|
Status DeliveryStatus `gorm:"not null;default:'pending';index:idx_deliveries_status,priority:1" json:"status"`
|
||||||
|
|
||||||
|
// DeletedAt repeats the BaseModel field only to be the second column
|
||||||
|
// of the event_id and status indexes, for the reason DeliveryResult
|
||||||
|
// gives.
|
||||||
|
DeletedAt gorm.DeletedAt `gorm:"index:idx_deliveries_event_id,priority:2;index:idx_deliveries_status,priority:2" json:"deletedAt,omitzero"`
|
||||||
|
|
||||||
// Relations
|
// Relations
|
||||||
Event Event `json:"event,omitzero"`
|
Event Event `json:"event,omitzero"`
|
||||||
|
|||||||
@@ -1,10 +1,21 @@
|
|||||||
package database
|
package database
|
||||||
|
|
||||||
|
import "gorm.io/gorm"
|
||||||
|
|
||||||
// DeliveryResult represents the result of a delivery attempt
|
// DeliveryResult represents the result of a delivery attempt
|
||||||
|
//
|
||||||
|
//nolint:lll // a struct tag cannot wrap
|
||||||
type DeliveryResult struct {
|
type DeliveryResult struct {
|
||||||
BaseModel
|
BaseModel
|
||||||
|
|
||||||
DeliveryID string `gorm:"type:uuid;not null" json:"deliveryId"`
|
// DeliveryID and DeletedAt make up one index, in that order.
|
||||||
|
// DeletedAt repeats the BaseModel field only to join it: GORM adds
|
||||||
|
// "deleted_at IS NULL" to almost every query, and where a column is
|
||||||
|
// matched against several values SQLite otherwise reads through the
|
||||||
|
// deleted_at index, which every live row matches.
|
||||||
|
DeliveryID string `gorm:"type:uuid;not null;index:idx_delivery_results_delivery_id,priority:1" json:"deliveryId"`
|
||||||
|
DeletedAt gorm.DeletedAt `gorm:"index:idx_delivery_results_delivery_id,priority:2" json:"deletedAt,omitzero"`
|
||||||
|
|
||||||
AttemptNum int `gorm:"not null" json:"attemptNum"`
|
AttemptNum int `gorm:"not null" json:"attemptNum"`
|
||||||
Success bool `json:"success"`
|
Success bool `json:"success"`
|
||||||
StatusCode int `json:"statusCode,omitempty"`
|
StatusCode int `json:"statusCode,omitempty"`
|
||||||
|
|||||||
@@ -1,9 +1,27 @@
|
|||||||
package database
|
package database
|
||||||
|
|
||||||
|
import (
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"gorm.io/gorm"
|
||||||
|
)
|
||||||
|
|
||||||
// Event represents a captured webhook event
|
// Event represents a captured webhook event
|
||||||
|
//
|
||||||
|
//nolint:lll // a struct tag cannot wrap
|
||||||
type Event struct {
|
type Event struct {
|
||||||
BaseModel
|
BaseModel
|
||||||
|
|
||||||
|
// CreatedAt and DeletedAt repeat the BaseModel fields only to index
|
||||||
|
// them for retention, which finds events by age. Its lookups carry
|
||||||
|
// GORM's "deleted_at IS NULL" (see DeliveryResult) and compare
|
||||||
|
// created_at with <, so their index has deleted_at first: SQLite
|
||||||
|
// narrows by a < only on the last column it uses. Its final delete
|
||||||
|
// has no deleted_at condition and uses the index on created_at
|
||||||
|
// alone. The other tables keep the unindexed BaseModel created_at.
|
||||||
|
CreatedAt time.Time `gorm:"index;index:idx_events_deleted_at_created_at,priority:2" json:"createdAt"`
|
||||||
|
DeletedAt gorm.DeletedAt `gorm:"index:idx_events_deleted_at_created_at,priority:1" json:"deletedAt,omitzero"`
|
||||||
|
|
||||||
WebhookID string `gorm:"type:uuid;not null" json:"webhookId"`
|
WebhookID string `gorm:"type:uuid;not null" json:"webhookId"`
|
||||||
EntrypointID string `gorm:"type:uuid;not null" json:"entrypointId"`
|
EntrypointID string `gorm:"type:uuid;not null" json:"entrypointId"`
|
||||||
|
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ import (
|
|||||||
"gorm.io/driver/sqlite"
|
"gorm.io/driver/sqlite"
|
||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
"sneak.berlin/go/webhooker/internal/config"
|
"sneak.berlin/go/webhooker/internal/config"
|
||||||
|
"sneak.berlin/go/webhooker/internal/datadir"
|
||||||
"sneak.berlin/go/webhooker/internal/gormlog"
|
"sneak.berlin/go/webhooker/internal/gormlog"
|
||||||
"sneak.berlin/go/webhooker/internal/logger"
|
"sneak.berlin/go/webhooker/internal/logger"
|
||||||
)
|
)
|
||||||
@@ -53,8 +54,9 @@ func NewWebhookDBManager(
|
|||||||
log: params.Logger.Get(),
|
log: params.Logger.Get(),
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create data directory if it doesn't exist
|
// Create data directory if it doesn't exist. datadir.DirPerm is the
|
||||||
err := os.MkdirAll(m.dataDir, dataDirPerm)
|
// single source of the directory mode; either package may run first.
|
||||||
|
err := os.MkdirAll(m.dataDir, datadir.DirPerm)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf(
|
return nil, fmt.Errorf(
|
||||||
"creating data directory %s: %w",
|
"creating data directory %s: %w",
|
||||||
|
|||||||
@@ -29,9 +29,11 @@ import (
|
|||||||
// process that was killed with SIGKILL blocks nothing.
|
// process that was killed with SIGKILL blocks nothing.
|
||||||
const LockFileName = "webhooker.lock"
|
const LockFileName = "webhooker.lock"
|
||||||
|
|
||||||
// dirPerm is the mode Acquire creates DATA_DIR with. It matches what
|
// DirPerm is the mode DATA_DIR is created with. It is the single
|
||||||
// internal/database uses, since whichever runs first creates it.
|
// source of that mode: internal/database consumes it rather than
|
||||||
const dirPerm = 0o750
|
// keeping its own copy, so the two packages that both create the
|
||||||
|
// directory cannot drift into disagreeing about its permissions.
|
||||||
|
const DirPerm = 0o750
|
||||||
|
|
||||||
// ErrLocked reports that another live process holds the data
|
// ErrLocked reports that another live process holds the data
|
||||||
// directory. Callers that need to know whether a deployment is running
|
// directory. Callers that need to know whether a deployment is running
|
||||||
@@ -64,7 +66,7 @@ func Acquire(dir string) (*Lock, error) {
|
|||||||
return nil, ErrNoDir
|
return nil, ErrNoDir
|
||||||
}
|
}
|
||||||
|
|
||||||
err := os.MkdirAll(dir, dirPerm)
|
err := os.MkdirAll(dir, DirPerm)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf(
|
return nil, fmt.Errorf(
|
||||||
"creating data directory %s: %w", dir, err,
|
"creating data directory %s: %w", dir, err,
|
||||||
|
|||||||
@@ -438,6 +438,31 @@ func (e *Engine) processNewTask(
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Restart recovery can send and release this delivery before the
|
||||||
|
// receiver's Notify queues it. Ownership cannot refuse a delivery
|
||||||
|
// nobody holds, so the row decides whether it still needs sending.
|
||||||
|
row, err := e.loadDelivery(webhookDB, task.DeliveryID)
|
||||||
|
if err != nil {
|
||||||
|
e.log.Error(
|
||||||
|
"failed to load delivery",
|
||||||
|
"delivery_id", task.DeliveryID,
|
||||||
|
"error", err,
|
||||||
|
)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
if row.Status != database.DeliveryStatusPending {
|
||||||
|
e.log.Info(
|
||||||
|
"delivery already handled, not sent again",
|
||||||
|
"delivery_id", task.DeliveryID,
|
||||||
|
"event_id", task.EventID,
|
||||||
|
"status", row.Status,
|
||||||
|
)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
event := buildEventFromTask(task)
|
event := buildEventFromTask(task)
|
||||||
|
|
||||||
event, err = e.hydrateEvent(
|
event, err = e.hydrateEvent(
|
||||||
@@ -482,7 +507,7 @@ func (e *Engine) processRetryTask(
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
d, err := e.loadRetryDelivery(
|
d, err := e.loadDelivery(
|
||||||
webhookDB, task.DeliveryID,
|
webhookDB, task.DeliveryID,
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -1643,7 +1668,7 @@ func (e *Engine) hydrateEvent(
|
|||||||
return event, nil
|
return event, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (e *Engine) loadRetryDelivery(
|
func (e *Engine) loadDelivery(
|
||||||
webhookDB *gorm.DB, deliveryID string,
|
webhookDB *gorm.DB, deliveryID string,
|
||||||
) (*database.Delivery, error) {
|
) (*database.Delivery, error) {
|
||||||
var d database.Delivery
|
var d database.Delivery
|
||||||
|
|||||||
@@ -273,6 +273,73 @@ func TestOwnershipIsReleasedAfterDelivery(t *testing.T) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestNotifyAfterRecoveryDoesNotSendAgain is the startup race of
|
||||||
|
// https://git.eeqj.de/sneak/webhooker/issues/299. The receiver has
|
||||||
|
// written a delivery, restart recovery finds it pending, sends it and
|
||||||
|
// releases it, and only then does the receiver's Notify for it arrive.
|
||||||
|
// Nothing owns the delivery by then, so Notify takes it.
|
||||||
|
func TestNotifyAfterRecoveryDoesNotSendAgain(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
targetID := uuid.New().String()
|
||||||
|
s := fSweepSetup(t, targetID, "recovered")
|
||||||
|
|
||||||
|
event := iSeedEvent(
|
||||||
|
t, s.WebhookDB, s.WebhookID, `{"recovered":true}`,
|
||||||
|
)
|
||||||
|
|
||||||
|
d := iSeedDelivery(
|
||||||
|
t, s.WebhookDB, event.ID, targetID,
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
|
||||||
|
s.Engine.ExportStart()
|
||||||
|
|
||||||
|
defer func() {
|
||||||
|
require.NoError(
|
||||||
|
t, s.Engine.ExportStop(context.Background()),
|
||||||
|
)
|
||||||
|
}()
|
||||||
|
|
||||||
|
// Restart recovery sends the delivery and lets it go.
|
||||||
|
iWaitForDelivered(t, s.WebhookDB, d.ID)
|
||||||
|
require.Eventually(
|
||||||
|
t,
|
||||||
|
func() bool {
|
||||||
|
return s.Engine.ExportInflightHeld() == 0
|
||||||
|
},
|
||||||
|
5*time.Second, 20*time.Millisecond,
|
||||||
|
)
|
||||||
|
|
||||||
|
body := event.Body
|
||||||
|
|
||||||
|
s.Engine.Notify([]delivery.Task{{
|
||||||
|
DeliveryID: d.ID,
|
||||||
|
EventID: event.ID,
|
||||||
|
WebhookID: s.WebhookID,
|
||||||
|
TargetID: targetID,
|
||||||
|
TargetName: "recovered",
|
||||||
|
TargetType: database.TargetTypeLog,
|
||||||
|
Body: &body,
|
||||||
|
EntrypointID: event.EntrypointID,
|
||||||
|
}})
|
||||||
|
|
||||||
|
// Notify took the delivery, and a worker releases it once it has
|
||||||
|
// run the task.
|
||||||
|
require.Eventually(
|
||||||
|
t,
|
||||||
|
func() bool {
|
||||||
|
return s.Engine.ExportInflightHeld() == 0
|
||||||
|
},
|
||||||
|
5*time.Second, 20*time.Millisecond,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.Len(
|
||||||
|
t, iResults(t, s.WebhookDB, d.ID), 1,
|
||||||
|
"the delivery was sent a second time",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
// TestRetryingRecoverySkipsASuccessfulResult is the retrying-side twin
|
// TestRetryingRecoverySkipsASuccessfulResult is the retrying-side twin
|
||||||
// of the pending reconcile. A second attempt that reached the receiver
|
// of the pending reconcile. A second attempt that reached the receiver
|
||||||
// and whose status write then failed sits at retrying holding a
|
// and whose status write then failed sits at retrying holding a
|
||||||
|
|||||||
@@ -170,7 +170,8 @@ func TestDelivery_CrossOriginRedirectDropsOriginScopedHeaders(
|
|||||||
// Stripping must not fire within the configured origin, or every
|
// Stripping must not fire within the configured origin, or every
|
||||||
// destination that redirects its own path would lose its
|
// destination that redirects its own path would lose its
|
||||||
// credential and start answering 401 — and would lose the inbound
|
// credential and start answering 401 — and would lose the inbound
|
||||||
// signature the receiver verifies.
|
// signature header the target endpoint verifies. webhooker's own
|
||||||
|
// receiver verifies no signature; it only forwards the header.
|
||||||
func TestDelivery_SameOriginRedirectKeepsOriginScopedHeaders(
|
func TestDelivery_SameOriginRedirectKeepsOriginScopedHeaders(
|
||||||
t *testing.T,
|
t *testing.T,
|
||||||
) {
|
) {
|
||||||
|
|||||||
@@ -300,3 +300,80 @@ func TestEntrypointCopyButtonIsProgressiveEnhancement(t *testing.T) {
|
|||||||
"the page must render to completion, not abort partway",
|
"the page must render to completion, not abort partway",
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// maxRetriesHelp is the wording both target forms must carry. The
|
||||||
|
// delivery core makes max_retries attempts in total, not that many
|
||||||
|
// retries on top of a first try (a fresh delivery starts at attempt 1
|
||||||
|
// and target_http gives up once the attempt number reaches
|
||||||
|
// max_retries), and 0 is special-cased to a single fire-and-forget
|
||||||
|
// attempt with no circuit breaker.
|
||||||
|
const maxRetriesHelp = "This is the total number of delivery attempts, " +
|
||||||
|
"not retries on top of the first: a value of 3 makes three attempts " +
|
||||||
|
"in all. 0 means a single attempt with no retries and no circuit " +
|
||||||
|
"breaker."
|
||||||
|
|
||||||
|
// TestTargetFormMaxRetriesCopyMatchesBehaviour pins the max_retries
|
||||||
|
// help text on both the create form (the add-target form on the webhook
|
||||||
|
// detail page) and the edit form, so the copy cannot drift back to
|
||||||
|
// calling the number a retry count.
|
||||||
|
func TestTargetFormMaxRetriesCopyMatchesBehaviour(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
var h *handlers.Handlers
|
||||||
|
|
||||||
|
var sess *session.Session
|
||||||
|
|
||||||
|
app := newTestApp(t, &h, &sess)
|
||||||
|
app.RequireStart()
|
||||||
|
|
||||||
|
t.Cleanup(app.RequireStop)
|
||||||
|
|
||||||
|
webhook := &database.Webhook{Name: "wh", RetentionDays: 14}
|
||||||
|
webhook.ID = testWebhookID
|
||||||
|
|
||||||
|
entrypoint := database.Entrypoint{Path: "abc123"}
|
||||||
|
entrypoint.ID = "ep-1"
|
||||||
|
|
||||||
|
createBody := renderPage(
|
||||||
|
t, h, sess, "source_detail.html", map[string]any{
|
||||||
|
dataKeyWebhook: webhook,
|
||||||
|
"Entrypoints": handlers.NewEntrypointViews(
|
||||||
|
[]database.Entrypoint{entrypoint},
|
||||||
|
),
|
||||||
|
"Targets": delivery.NewTargetViews(nil),
|
||||||
|
"Events": []database.Event{},
|
||||||
|
"BaseURL": "https://hooks.example.com",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.Contains(
|
||||||
|
t, createBody, maxRetriesHelp,
|
||||||
|
"the add-target form must explain max_retries as total attempts",
|
||||||
|
)
|
||||||
|
|
||||||
|
// A slack target exercises the same max_retries field while needing
|
||||||
|
// only Config.URL from the edit template, so the test data stays
|
||||||
|
// minimal. The Target key mirrors the field names the template reads
|
||||||
|
// off the handler's view value.
|
||||||
|
editBody := renderPage(
|
||||||
|
t, h, sess, "target_edit.html", map[string]any{
|
||||||
|
dataKeyWebhook: webhook,
|
||||||
|
"Target": map[string]any{
|
||||||
|
"ID": "tg-1",
|
||||||
|
"Name": "t",
|
||||||
|
"Type": "slack",
|
||||||
|
"Active": true,
|
||||||
|
"MaxRetries": 3,
|
||||||
|
"Config": map[string]any{
|
||||||
|
"URL": "https://hooks.slack.com/services/x",
|
||||||
|
},
|
||||||
|
},
|
||||||
|
dataKeyError: "",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.Contains(
|
||||||
|
t, editBody, maxRetriesHelp,
|
||||||
|
"the target edit form must explain max_retries as total attempts",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|||||||
@@ -380,9 +380,8 @@ func csrfTookStrictPath(
|
|||||||
|
|
||||||
// TestCSRF_ForwardedProtoSpellingsTakeStrictPath runs the header
|
// TestCSRF_ForwardedProtoSpellingsTakeStrictPath runs the header
|
||||||
// spellings a real proxy emits through the middleware. The environment
|
// spellings a real proxy emits through the middleware. The environment
|
||||||
// is dev -- the DEFAULT when WEBHOOKER_ENVIRONMENT is unset -- to pin
|
// is set to dev -- the permissive setting -- to pin that the routing is
|
||||||
// that the routing is a per-request transport decision and owes
|
// a per-request transport decision and owes nothing to configuration.
|
||||||
// nothing to configuration.
|
|
||||||
func TestCSRF_ForwardedProtoSpellingsTakeStrictPath(t *testing.T) {
|
func TestCSRF_ForwardedProtoSpellingsTakeStrictPath(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
|
|||||||
@@ -5,7 +5,7 @@
|
|||||||
// several packages, by hand, and the answers disagreed. The session
|
// several packages, by hand, and the answers disagreed. The session
|
||||||
// cookie's Secure attribute was decided at startup from the configured
|
// cookie's Secure attribute was decided at startup from the configured
|
||||||
// environment while the CSRF cookie's was decided per-request, so a
|
// environment while the CSRF cookie's was decided per-request, so a
|
||||||
// deployment behind a TLS proxy in the default environment emitted one
|
// deployment behind a TLS proxy in the dev environment emitted one
|
||||||
// Secure cookie and one non-Secure cookie on the same response.
|
// Secure cookie and one non-Secure cookie on the same response.
|
||||||
// Everything kept working, which is exactly why nobody noticed.
|
// Everything kept working, which is exactly why nobody noticed.
|
||||||
//
|
//
|
||||||
|
|||||||
@@ -146,8 +146,8 @@ func newStore(key []byte) *sessions.CookieStore {
|
|||||||
//
|
//
|
||||||
// This is decided per-request, not once at startup. Deciding it at
|
// This is decided per-request, not once at startup. Deciding it at
|
||||||
// startup from the configured environment is what this replaces, and
|
// startup from the configured environment is what this replaces, and
|
||||||
// it got the DEFAULT posture wrong: "dev" is the environment when
|
// it got the DEFAULT posture wrong: "dev" was then the environment when
|
||||||
// WEBHOOKER_ENVIRONMENT is unset, so a deployment terminating TLS at a
|
// WEBHOOKER_ENVIRONMENT was unset, so a deployment terminating TLS at a
|
||||||
// proxy without also setting the environment emitted the
|
// proxy without also setting the environment emitted the
|
||||||
// authentication cookie with no Secure attribute -- silently, and on
|
// authentication cookie with no Secure attribute -- silently, and on
|
||||||
// the same response as a CSRF cookie that did have one.
|
// the same response as a CSRF cookie that did have one.
|
||||||
|
|||||||
@@ -990,8 +990,8 @@ func sessionCookieFrom(
|
|||||||
|
|
||||||
// TestSave_SecureFollowsRequestTransport is the regression test for
|
// TestSave_SecureFollowsRequestTransport is the regression test for
|
||||||
// the defect this replaces: Secure was fixed at startup from the
|
// the defect this replaces: Secure was fixed at startup from the
|
||||||
// configured environment, and "dev" is the environment when
|
// configured environment, and "dev" was then the environment when
|
||||||
// WEBHOOKER_ENVIRONMENT is unset. A deployment behind a TLS proxy in
|
// WEBHOOKER_ENVIRONMENT was unset. A deployment behind a TLS proxy in
|
||||||
// that DEFAULT posture shipped the authentication cookie with no
|
// that DEFAULT posture shipped the authentication cookie with no
|
||||||
// Secure attribute and said nothing about it.
|
// Secure attribute and said nothing about it.
|
||||||
//
|
//
|
||||||
|
|||||||
@@ -120,10 +120,13 @@
|
|||||||
<label class="text-sm text-gray-700">Timeout (seconds, blank = default):</label>
|
<label class="text-sm text-gray-700">Timeout (seconds, blank = default):</label>
|
||||||
<input type="number" name="timeout" min="0" max="300" :disabled="targetType !== 'http'" class="input text-sm w-24">
|
<input type="number" name="timeout" min="0" max="300" :disabled="targetType !== 'http'" class="input text-sm w-24">
|
||||||
</div>
|
</div>
|
||||||
<div x-show="targetType === 'http'" class="flex gap-2 items-center">
|
<div x-show="targetType === 'http'">
|
||||||
<label class="text-sm text-gray-700">Max retries (0 = fire-and-forget):</label>
|
<div class="flex gap-2 items-center">
|
||||||
|
<label class="text-sm text-gray-700">Max retries:</label>
|
||||||
<input type="number" name="max_retries" value="0" min="0" max="20" class="input text-sm w-24">
|
<input type="number" name="max_retries" value="0" min="0" max="20" class="input text-sm w-24">
|
||||||
</div>
|
</div>
|
||||||
|
<p class="text-xs text-gray-500 mt-1">This is the total number of delivery attempts, not retries on top of the first: a value of 3 makes three attempts in all. 0 means a single attempt with no retries and no circuit breaker.</p>
|
||||||
|
</div>
|
||||||
<div x-show="targetType === 'slack'">
|
<div x-show="targetType === 'slack'">
|
||||||
<input type="url" name="url" placeholder="https://hooks.slack.com/services/..." :disabled="targetType !== 'slack'" class="input text-sm">
|
<input type="url" name="url" placeholder="https://hooks.slack.com/services/..." :disabled="targetType !== 'slack'" class="input text-sm">
|
||||||
<p class="text-xs text-gray-500 mt-1">Slack or Mattermost incoming webhook URL. Payloads are pretty-printed in code blocks.</p>
|
<p class="text-xs text-gray-500 mt-1">Slack or Mattermost incoming webhook URL. Payloads are pretty-printed in code blocks.</p>
|
||||||
|
|||||||
@@ -69,7 +69,7 @@
|
|||||||
<div class="form-group">
|
<div class="form-group">
|
||||||
<label for="max_retries" class="label">Max retries</label>
|
<label for="max_retries" class="label">Max retries</label>
|
||||||
<input type="number" id="max_retries" name="max_retries" value="{{.Target.MaxRetries}}" min="0" max="20" class="input">
|
<input type="number" id="max_retries" name="max_retries" value="{{.Target.MaxRetries}}" min="0" max="20" class="input">
|
||||||
<p class="text-xs text-gray-500 mt-1">0 is fire-and-forget: one attempt, no circuit breaker.</p>
|
<p class="text-xs text-gray-500 mt-1">This is the total number of delivery attempts, not retries on top of the first: a value of 3 makes three attempts in all. 0 means a single attempt with no retries and no circuit breaker.</p>
|
||||||
</div>
|
</div>
|
||||||
{{end}}
|
{{end}}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user