8 Commits

Author SHA1 Message Date
def52ae092 State the UUID-is-the-credential rule as a rule (closes #301)
All checks were successful
check / check (push) Successful in 3m35s
The receiver has authenticated on the entrypoint UUID alone since
inbound signature verification was removed in #279. The README
described that as the current state; it did not say it is the
decision. Restate it as the rule, so a proposal to add HMAC, a shared
secret or a bearer token to the receiver is contradicted by the docs
rather than merely unimplemented.

The rule now appears in the intro, in its own section, and in the
Authentication and Security lists, and carries the two consequences an
operator has to act on: the URL is a capability to be kept out of logs
and tickets, and rotation means minting a new entrypoint rather than
changing a key.

Also corrects one stale comment: a redirect test said the inbound
signature was one "the receiver verifies", which in this repo's
vocabulary names webhooker's own receiver. The endpoint that verifies
it is the delivery target's.
2026-08-25 20:38:47 +00:00
d61d9dc1c1 Drop the stale open-work claim from the TODO status section
All checks were successful
check / check (push) Successful in 2m56s
The milestone is the authoritative count and the section already says
so; the lead-in asserted work remaining independently of it.
2026-08-24 15:52:44 +00:00
b0a011f6b4 Render unknown for a zero CreatedAt in Slack/Mattermost messages (closes #298)
Some checks failed
check / check (push) Superseded by a newer commit; never tested
2026-08-24 17:52:25 +02:00
5976a4a98f Carry the event's receipt time into every delivery (closes #257)
All checks were successful
check / check (push) Successful in 3m4s
2026-08-24 06:44:24 +02:00
b2c9acdaa6 Correct four documentation claims ahead of the 1.0.0 tag
All checks were successful
check / check (push) Successful in 7s
2026-08-24 06:25:08 +02:00
af3703d748 Close the two remaining delivery terminal-state gaps (closes #107)
All checks were successful
check / check (push) Successful in 3m16s
2026-08-24 05:12:02 +02:00
322d9a6d6b Create every SQLite file 0600 (closes #255)
All checks were successful
check / check (push) Successful in 3m17s
2026-08-24 04:33:40 +02:00
b9f7db6901 Fail loudly on an unparseable SENTRY_DSN and a malformed .env (closes #283)
All checks were successful
check / check (push) Successful in 2m54s
2026-08-24 04:25:29 +02:00
24 changed files with 2222 additions and 118 deletions

105
README.md
View File

@@ -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
@@ -71,8 +78,18 @@ make clean # Remove bin/
### Configuration
All configuration is via environment variables. For local development,
you can place variables in a `.env` file in the project root (loaded
automatically via `godotenv/autoload`).
you can place variables in a `.env` file in the process working
directory, read once at startup before anything else looks at the
environment.
The file is optional and having none is the normal case for a
deployment. A file that is there but cannot be parsed aborts startup
with a message naming it, because a single malformed line makes none
of the file apply: every variable in it silently reverts to its
default, which is exactly the failure [Invalid values abort
startup](#invalid-values-abort-startup) exists to prevent, for all of
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:
@@ -125,7 +142,7 @@ TTY detection, and security headers are always applied.
| `MAINTENANCE_MODE` | Report `maintenanceMode: true` in the healthcheck JSON. It does not change how any request is served — no maintenance page exists | `false` |
| `METRICS_USERNAME` | Basic auth username for `/metrics`. Must be set together with `METRICS_PASSWORD`; one without the other fails startup | `""` |
| `METRICS_PASSWORD` | Basic auth password for `/metrics`. Must be set together with `METRICS_USERNAME`; one without the other fails startup | `""` |
| `SENTRY_DSN` | Sentry error reporting DSN | `""` |
| `SENTRY_DSN` | Sentry error reporting DSN. Unset leaves error reporting off; a value the Sentry SDK cannot parse fails startup rather than serving with reporting silently off | `""` |
| `RETENTION_SWEEP_INTERVAL` | How often the retention reaper and archive sweeper run (Go duration, must be positive) | `1h` |
| `SESSION_IDLE_TIMEOUT` | Idle session timeout (Go duration) | `24h` |
| `RECEIVER_RATE_LIMIT` | Receiver requests/minute per IP per entrypoint (10x that per IP across the route) | `120` |
@@ -464,10 +481,20 @@ startup), every entry in `TRUSTED_PROXIES` and
`ALLOWED_EGRESS_CIDRS` must be a CIDR block or a bare IP address, and
`BIND_ADDRESS` must be an IP address literal — `localhost`,
`127.0.0.1:8080` and `10.0.0.0/8` are each rejected rather than
resolved, split, or narrowed to something they do not say.
resolved, split, or narrowed to something they do not say — and
`SENTRY_DSN` must parse as a Sentry DSN.
`SESSION_IDLE_TIMEOUT` is the exception: a
non-positive value there means idle expiry is disabled, not invalid.
`SENTRY_DSN` is checked with the Sentry SDK's own parser, the same call
the SDK makes on the DSN it is later handed, so what configuration
accepts is exactly what will initialise. A typo in it is the one
configuration mistake nothing downstream can ever notice — the variable
is still set, so every later signal reports error reporting as on while
no report is being sent — which is why it aborts rather than starting
with reporting off. Leaving it unset is not a mistake and not affected:
error reporting is simply off and startup is normal.
Boolean variables (`DEBUG`, `MAINTENANCE_MODE`) accept exactly the
spellings Go's `strconv.ParseBool` accepts — `1`, `t`, `T`, `TRUE`,
`true`, `True`, `0`, `f`, `F`, `FALSE`, `false`, `False` — and nothing
@@ -1129,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
@@ -1214,6 +1265,16 @@ webhooker solves this by acting as a durable intermediary:
backoff. Every delivery attempt is logged with status codes, response
bodies, and timing.
**That guarantee is at-least-once, not exactly-once.** When a send
reaches its target but the write recording that outcome fails, the
delivery is deliberately left in a recoverable state rather than
marked done — losing a delivery is the worse failure — so the
pending sweep picks it up about fifteen minutes later, or the next
restart does, and the target receives a payload it already got.
webhooker adds no delivery identifier of its own to an outbound
request, so **make your receiver idempotent** against whatever the
payload itself carries.
3. **Observability** — Full request/response logging for every webhook
received and every delivery attempted. Prometheus metrics expose
volume, latency, and error rates. The web UI provides real-time
@@ -2596,12 +2657,15 @@ abuse limit later; they are tracked as future work.
| `POST` | `/source/{id}/edit` | Edit webhook submission |
| `POST` | `/source/{id}/delete` | Delete webhook |
| `GET` | `/source/{id}/logs` | Webhook event logs |
| `GET` | `/source/{id}/logs/{eventID}/body` | Download an event's full stored body. The log page renders each body only up to its cap, so this is the only route that serves a whole one; it is offered wherever a body is shown truncated |
| `POST` | `/source/{id}/deliveries/{deliveryID}/replay` | Replay a finished delivery: creates a new delivery for the same event against the target's current configuration (30 per minute per bucket, then `429`) |
| `POST` | `/source/{id}/events/{eventID}/resubmit` | Resubmit a stored event: creates a new event copying it and fans that out to every currently active target (30 per minute per bucket, then `429`) |
| `POST` | `/source/{id}/entrypoints` | Add entrypoint to webhook |
| `POST` | `/source/{id}/entrypoints/{entrypointID}/delete` | Delete an entrypoint |
| `POST` | `/source/{id}/entrypoints/{entrypointID}/toggle` | Enable or disable an entrypoint |
| `POST` | `/source/{id}/targets` | Add target to webhook |
| `GET` | `/source/{id}/targets/{targetID}/edit` | Edit target form. The one page that renders a target's destination URL and header values in full, rather than masked |
| `POST` | `/source/{id}/targets/{targetID}/edit` | Edit target submission |
| `POST` | `/source/{id}/targets/{targetID}/delete` | Delete a target |
| `POST` | `/source/{id}/targets/{targetID}/toggle` | Enable or disable a target |
@@ -2640,6 +2704,8 @@ webhooker/
├── internal/
│ ├── banner/
│ │ └── banner.go # Ruled block for the one credential shown in the clear
│ ├── ciscript/
│ │ └── doc.go # Tests for the CI shell scripts in script/; no runtime code
│ ├── resetpw/
│ │ └── resetpw.go # `webhooker resetpw`: set an account's password, stopped deployments only
│ ├── config/
@@ -2709,13 +2775,17 @@ webhooker/
│ │ ├── ratelimit.go # Per-IP rate limiting middleware (go-chi/httprate)
│ │ ├── loginguard.go # Login failure counters and the Argon2id verification semaphore
│ │ └── testing.go # NewForTest: Middleware without the fx lifecycle
│ ├── reqtls/
│ │ └── reqtls.go # IsTLS: the one TLS predicate, r.TLS or X-Forwarded-Proto
│ ├── server/
│ │ ├── server.go # Server struct, fx lifecycle, signal handling
│ │ ├── http.go # HTTP server setup with timeouts
│ │ └── routes.go # All route definitions
── session/
├── session.go # Cookie-based session management
└── testing.go # NewForTest: Session without the fx lifecycle
── session/
├── session.go # Cookie-based session management
└── testing.go # NewForTest: Session without the fx lifecycle
│ └── versionscript/
│ └── doc.go # Tests for script/version and the build files that use it
├── static/
│ ├── static.go # //go:embed directive
│ ├── css/input.css # Tailwind input, source for tailwind.css (make css)
@@ -2828,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
@@ -2867,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

41
TODO.md
View File

@@ -18,18 +18,27 @@ Issue branches do NOT touch this file — the manager maintains it on
# Status
1.0.0 is open, with work remaining. The milestone
(https://git.eeqj.de/sneak/webhooker/milestone/9) is the authoritative
list, and the only place to read a count or a state of play from. This
file records where the project is, not what is in flight: a sentence
whose truth depends on a branch being unmerged is wrong the moment it
merges, and this file has been wrong that way before.
The milestone (https://git.eeqj.de/sneak/webhooker/milestone/9) is the
authoritative list, and the only place to read a count or a state of
play from. This file records where the project is, not what is in
flight: a sentence whose truth depends on a branch being unmerged is
wrong the moment it merges, and this file has been wrong that way
before.
The tag is held on a durability defect
(https://git.eeqj.de/sneak/webhooker/issues/256): a concurrent reader
of a per-webhook event database strands delivered webhooks at
`pending`, and the next restart re-delivers them. That issue gates
`v1.0.0`, and is where the fix's own state is tracked.
The durability defect that held the tag has landed
(https://git.eeqj.de/sneak/webhooker/issues/256, commit `8d64259`).
Every SQLite handle opens with WAL journaling and a busy timeout, a
bookkeeping write that fails leaves its delivery in a recoverable
state rather than a lying one, and recovery skips a delivery that
already has a successful result row. Final pre-tag verification
exercised it and confirmed it holds. Whatever the milestone still
shows open is what remains before `v1.0.0`.
Delivery is at-least-once by design, not by accident: a send whose
result row does not land is attempted again, so a receiver can see a
duplicate. That is deliberate — the alternative is a silent lost
delivery — and the README says so under Rationale. It is not a defect
to re-file.
One caveat on reading a green check: a docs-only commit deliberately
replays from the layer cache
@@ -39,11 +48,11 @@ commit invalidates the `COPY` layer and genuinely executes.
# Next Step
Land https://git.eeqj.de/sneak/webhooker/issues/256, then clear the
rest of the open 1.0.0 milestone and tag `v1.0.0`. Merging `next` into
`main` is a separate act from tagging and waits on neither of those:
`next` is kept mergeable at all times, which is the point of the
branch.
Clear the rest of the open 1.0.0 milestone
(https://git.eeqj.de/sneak/webhooker/milestone/9) and tag `v1.0.0`.
Merging `next` into `main` is a separate act from tagging and waits on
neither of those: `next` is kept mergeable at all times, which is the
point of the branch.
# Completed Steps

View File

@@ -0,0 +1,107 @@
package main
import (
"bytes"
"os"
"path/filepath"
"strings"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/config"
)
// dotEnvKey is a throwaway variable name these tests write and read,
// so they cannot disturb real configuration.
const dotEnvKey = "WEBHOOKER_TEST_DISPATCH_VALUE"
// writeDotEnvInWorkingDir puts contents in a .env file in a fresh
// temporary directory and moves the process there.
//
// The callers are deliberately not parallel and must stay that way:
// t.Chdir moves the whole process. Go releases parallel tests only
// after every sequential test in the package has finished, so nothing
// else runs while these do.
func writeDotEnvInWorkingDir(t *testing.T, contents string) {
t.Helper()
dir := t.TempDir()
require.NoError(t, os.WriteFile(
filepath.Join(dir, config.DotEnvPath),
[]byte(contents), 0o600,
))
t.Chdir(dir)
}
// TestDispatch_MalformedDotEnvRefuses pins the second half of the
// defect. godotenv applies nothing at all when a file will not parse,
// so one mistyped line used to revert every variable in it to its
// default and start the server anyway, with no log line naming the
// file. The refusal has to arrive before any subcommand runs, which
// is why `help` — the one subcommand that touches nothing — is still
// refused here.
//
//nolint:paralleltest // t.Chdir moves the whole process.
func TestDispatch_MalformedDotEnvRefuses(t *testing.T) {
writeDotEnvInWorkingDir(t, "PORT 19615\n")
var stdout, stderr bytes.Buffer
code := dispatch(
[]string{helpCommand}, strings.NewReader(""), &stdout, &stderr,
)
require.Equal(t, 1, code, "a broken .env must exit non-zero")
assert.Contains(
t, stderr.String(), config.DotEnvPath,
"the refusal must name the file",
)
assert.Empty(
t, stdout.String(),
"the subcommand must not have run",
)
}
// TestDispatch_LoadsDotEnvBeforeSubcommands pins the ordering the
// godotenv/autoload import used to provide for free. It ran in an
// init(), so .env was in the environment before anything read it —
// including config.DataDir, which both the DATA_DIR lock and resetpw
// call outside the fx graph. Loading any later would let a .env that
// sets DATA_DIR lock one directory while the config opened databases
// in another.
func TestDispatch_LoadsDotEnvBeforeSubcommands(t *testing.T) {
t.Setenv(dotEnvKey, "placeholder")
require.NoError(t, os.Unsetenv(dotEnvKey))
writeDotEnvInWorkingDir(t, dotEnvKey+"=from-dot-env\n")
var stdout, stderr bytes.Buffer
code := dispatch(
[]string{helpCommand}, strings.NewReader(""), &stdout, &stderr,
)
require.Equal(t, 0, code)
assert.Equal(
t, "from-dot-env", os.Getenv(dotEnvKey),
"the file must be applied before the subcommand runs",
)
}
// TestDispatch_MissingDotEnvIsFine pins the case most deployments are
// in: no .env at all, which must stay a normal start.
//
//nolint:paralleltest // t.Chdir moves the whole process.
func TestDispatch_MissingDotEnvIsFine(t *testing.T) {
t.Chdir(t.TempDir())
var stdout, stderr bytes.Buffer
code := dispatch(
[]string{helpCommand}, strings.NewReader(""), &stdout, &stderr,
)
require.Equal(t, 0, code)
assert.Empty(t, stderr.String())
}

View File

@@ -54,6 +54,11 @@ const stopTimeout = 5 * time.Second
// caller can tell "called wrong" from "declined".
const exitUsage = 2
// helpCommand is the subcommand that prints usage. The flag spellings
// beside it in the switch are aliases; this is the name the usage text
// documents and the one tests invoke.
const helpCommand = "help"
// Build-time variables set via -ldflags.
//
//nolint:gochecknoglobals // Build-time variables injected by the linker.
@@ -75,11 +80,27 @@ func main() {
// every existing deployment invoke; that path is unchanged, including
// where the DATA_DIR lock is taken relative to building the fx graph
// and how fx propagates a non-zero exit itself.
//
// The optional .env file is read here, before any subcommand and so
// before anything reads the environment — config.DataDir, which both
// the DATA_DIR lock and resetpw call outside the fx graph, above all.
// It used to be read from an init() in internal/config, which put it
// earlier still but threw the error away: a single malformed line
// applied none of the file and said nothing about it. A file that is
// not there stays fine, since .env is optional and most deployments
// do not have one.
func dispatch(
args []string,
stdin io.Reader,
stdout, stderr io.Writer,
) int {
err := config.LoadDotEnv()
if err != nil {
_, _ = fmt.Fprintf(stderr, "%s: %v\n", appname, err)
return 1
}
if len(args) == 0 {
return run(stderr)
}
@@ -87,7 +108,7 @@ func dispatch(
switch args[0] {
case resetpw.Name:
return resetpw.Run(args[1:], stdin, stdout, stderr)
case "help", "-h", "-help", "--help":
case helpCommand, "-h", "-help", "--help":
usage(stdout)
return 0

View File

@@ -121,7 +121,7 @@ func TestDispatch_Help(t *testing.T) {
var stdout, stderr bytes.Buffer
code := dispatch(
[]string{"help"}, strings.NewReader(""), &stdout, &stderr,
[]string{helpCommand}, strings.NewReader(""), &stdout, &stderr,
)
require.Equal(t, 0, code)

View File

@@ -4,6 +4,7 @@ package config
import (
"errors"
"fmt"
"io/fs"
"log/slog"
"net/netip"
"os"
@@ -11,13 +12,11 @@ import (
"strings"
"time"
"github.com/getsentry/sentry-go"
"github.com/joho/godotenv"
"go.uber.org/fx"
"sneak.berlin/go/webhooker/internal/globals"
"sneak.berlin/go/webhooker/internal/logger"
// Populates the environment from a ./.env file automatically for
// development configuration. Kept in one place only (here).
_ "github.com/joho/godotenv/autoload"
)
const (
@@ -84,6 +83,12 @@ const (
// IPv6 prefix spends on the ::ffff:0:0/96 wrapper, so a /104
// covers the same addresses as an IPv4 /8.
mappedV4Offset = 96
// DotEnvPath is the optional file of KEY=value lines read into the
// environment at startup, relative to the process working
// directory. Exported so that documentation and tests name the
// same path the loader opens.
DotEnvPath = ".env"
)
// ErrInvalidEnvironment is returned when WEBHOOKER_ENVIRONMENT
@@ -107,6 +112,15 @@ var ErrInvalidCIDR = errors.New("invalid CIDR")
// something that is not an IP address literal.
var ErrInvalidBindAddress = errors.New("invalid bind address")
// ErrInvalidSentryDSN is returned when SENTRY_DSN is set to something
// the Sentry SDK cannot parse as a DSN.
var ErrInvalidSentryDSN = errors.New("invalid Sentry DSN")
// ErrDotEnvUnreadable is returned when the optional .env file exists
// but cannot be read or parsed. A file that is not there is not an
// error; a file that is there and broken is.
var ErrDotEnvUnreadable = errors.New("unreadable .env file")
// ErrIncompleteMetricsAuth is returned when exactly one of
// METRICS_USERNAME and METRICS_PASSWORD carries a value. Neither
// fallback is acceptable: serving /metrics on the username alone
@@ -212,12 +226,62 @@ func (c *Config) MetricsAuthEnabled() bool {
return c.MetricsUsername != "" && c.MetricsPassword != ""
}
// SentryEnabled reports whether error reporting is shipped to Sentry.
// It is the only answer to that question in the codebase: the SDK
// initialisation, the sentryhttp middleware registration and the
// startup log's sentryEnabled field all read this one method, so the
// log cannot report reporting as on while nothing is sending.
//
// A non-empty DSN is enough because loadFromEnv already parsed it with
// the SDK's own parser and refused to build a Config around one the
// SDK would reject, and because initialising the SDK with a DSN that
// parsed and failed anyway aborts the process rather than leaving this
// true and the client absent.
func (c *Config) SentryEnabled() bool {
return c.SentryDSN != ""
}
// envString returns the value of the named environment variable,
// or an empty string if not set.
func envString(key string) string {
return os.Getenv(key)
}
// LoadDotEnv reads DotEnvPath into the environment when that file is
// present, and reports a file that is present but broken.
//
// It has to run before anything reads the environment, so that every
// reader agrees on what the environment holds — the DATA_DIR lock
// taken before the fx graph exists as much as loadFromEnv itself. A
// variable already set in the real environment wins: godotenv never
// overwrites one.
//
// A missing file is not an error. It is a development convenience and
// most deployments set the environment directly.
//
// Any other failure is. godotenv parses the whole file before setting
// anything, so a single malformed line applies none of it: every
// variable in the file silently reverts to its default, which defeats
// the fail-loud guarantee for all of them at once.
func LoadDotEnv() error {
return loadDotEnvFile(DotEnvPath)
}
// loadDotEnvFile is LoadDotEnv over a named file, so tests can point
// at a temporary one instead of the process working directory.
func loadDotEnvFile(path string) error {
err := godotenv.Load(path)
if err == nil || errors.Is(err, fs.ErrNotExist) {
return nil
}
return fmt.Errorf(
"%w: %s: %w; nothing in it was applied, so fix the file or "+
"remove it",
ErrDotEnvUnreadable, path, err,
)
}
// DataDir resolves DATA_DIR, applying DefaultDataDir when it is unset
// or empty. It is exported so that entry points which must act on the
// data directory before the fx graph exists — taking the exclusive
@@ -462,6 +526,41 @@ func envBindAddress(key, defaultValue string) (string, error) {
return addr.String(), nil
}
// envSentryDSN returns the value of the named environment variable
// checked as a Sentry DSN. An unset (or empty, or whitespace-only)
// value yields "", which means error reporting stays off — the common
// case, and a normal start.
//
// A set value is parsed with sentry.NewDsn, which is the call
// sentry.Init makes on the DSN it is handed, so what passes here is
// exactly what the SDK will accept later and the two cannot disagree.
// Reproducing the check by hand instead would cost this package its
// dependency on the SDK — already a module dependency, already linked
// into the binary — in exchange for a second definition of "valid DSN"
// free to drift from the one that decides.
//
// A set value that does not parse is a hard error naming the key, so
// startup fails loudly. Losing error reporting is the failure this
// variable exists to prevent, and a typo in a DSN is silent forever:
// nothing later in the process can notice that reports are going
// nowhere. The bad value is quoted because it is a URL to a public
// endpoint carrying a public key, not a secret.
func envSentryDSN(key string) (string, error) {
v := strings.TrimSpace(os.Getenv(key))
if v == "" {
return "", nil
}
_, err := sentry.NewDsn(v)
if err != nil {
return "", fmt.Errorf(
"%w: %s: %q: %w", ErrInvalidSentryDSN, key, v, err,
)
}
return v, nil
}
// resolveMetricsAuth reads the /metrics basic-auth credentials and
// rejects a half-set pair, naming both variables either way. The
// error carries neither value: the password is a secret.
@@ -594,6 +693,11 @@ func loadFromEnv() (*Config, error) {
return nil, err
}
sentryDSN, err := envSentryDSN("SENTRY_DSN")
if err != nil {
return nil, err
}
return &Config{
DataDir: DataDir(),
Debug: debug,
@@ -603,7 +707,7 @@ func loadFromEnv() (*Config, error) {
MetricsPassword: metricsPassword,
Port: port,
BindAddress: bindAddress,
SentryDSN: envString("SENTRY_DSN"),
SentryDSN: sentryDSN,
RetentionSweepInterval: retentionSweepInterval,
SessionIdleTimeout: sessionIdleTimeout,
ReceiverRateLimit: receiverRateLimit,
@@ -738,7 +842,7 @@ func New(lc fx.Lifecycle, params ConfigParams) (*Config, error) {
"receiverRateLimit", s.ReceiverRateLimit,
"trustedProxies", len(s.TrustedProxies),
"allowedEgressCIDRs", len(s.AllowedEgressCIDRs),
"hasSentryDSN", s.SentryDSN != "",
"sentryEnabled", s.SentryEnabled(),
"hasMetricsAuth", s.MetricsAuthEnabled(),
)

View File

@@ -0,0 +1,158 @@
package config_test
import (
"os"
"path/filepath"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/config"
)
// dotEnvKey is a throwaway variable name the .env tests write and
// read, so they cannot disturb real configuration.
const dotEnvKey = "WEBHOOKER_TEST_DOTENV_VALUE"
// malformedDotEnv is a file godotenv cannot parse. The first line is
// the realistic typo — a space where the `=` belongs — and the rest
// make sure nothing downstream treats the file as salvageable line by
// line.
const malformedDotEnv = "PORT 19615\n" +
"this is not = valid ! syntax\n" +
"\"unclosed\n"
// unsetDotEnvKey makes dotEnvKey genuinely absent for the duration of
// the test and restores it afterwards. t.Setenv registers the restore;
// the Unsetenv that follows is what the test actually needs, because a
// variable set to the empty string is still present in os.Environ and
// godotenv would refuse to overwrite it.
func unsetDotEnvKey(t *testing.T) {
t.Helper()
t.Setenv(dotEnvKey, "placeholder")
require.NoError(t, os.Unsetenv(dotEnvKey))
}
// writeDotEnv writes contents to a .env file in a fresh temporary
// directory and returns its path.
func writeDotEnv(t *testing.T, contents string) string {
t.Helper()
path := filepath.Join(t.TempDir(), config.DotEnvPath)
require.NoError(t, os.WriteFile(path, []byte(contents), 0o600))
return path
}
// TestLoadDotEnv_MissingFileIsFine pins the case most deployments are
// in. The file is optional: it is a development convenience, and a
// deployment that configures the environment directly must start
// normally rather than be refused for a file it was never meant to
// have.
//
//nolint:paralleltest // unsetDotEnvKey uses t.Setenv.
func TestLoadDotEnv_MissingFileIsFine(t *testing.T) {
unsetDotEnvKey(t)
absent := filepath.Join(t.TempDir(), config.DotEnvPath)
require.NoError(t, config.LoadDotEnvFileForTest(absent))
_, present := os.LookupEnv(dotEnvKey)
assert.False(t, present, "nothing may be set from an absent file")
}
// TestLoadDotEnv_AppliesValues pins that a well-formed file still
// reaches the environment, which is the whole reason the file is read
// at all.
//
//nolint:paralleltest // unsetDotEnvKey uses t.Setenv.
func TestLoadDotEnv_AppliesValues(t *testing.T) {
unsetDotEnvKey(t)
path := writeDotEnv(t, "# a comment\n"+dotEnvKey+"=from-dot-env\n")
require.NoError(t, config.LoadDotEnvFileForTest(path))
assert.Equal(t, "from-dot-env", os.Getenv(dotEnvKey))
}
// TestLoadDotEnv_RealEnvironmentWins pins that the file cannot
// override a variable the process was actually started with. A
// deployment that sets DATA_DIR in its unit file must not have it
// silently replaced by a stale .env left in the working directory.
func TestLoadDotEnv_RealEnvironmentWins(t *testing.T) {
t.Setenv(dotEnvKey, "from-environment")
path := writeDotEnv(t, dotEnvKey+"=from-dot-env\n")
require.NoError(t, config.LoadDotEnvFileForTest(path))
assert.Equal(t, "from-environment", os.Getenv(dotEnvKey))
}
// TestLoadDotEnv_MalformedFileAborts is the defect this fixes. One bad
// line makes godotenv apply none of the file, so every variable in it
// reverts to its default; the process used to start that way with no
// log line naming the file at all.
//
//nolint:paralleltest // unsetDotEnvKey uses t.Setenv.
func TestLoadDotEnv_MalformedFileAborts(t *testing.T) {
unsetDotEnvKey(t)
path := writeDotEnv(
t, malformedDotEnv+dotEnvKey+"=from-dot-env\n",
)
err := config.LoadDotEnvFileForTest(path)
require.Error(t, err)
require.ErrorIs(t, err, config.ErrDotEnvUnreadable)
assert.Contains(
t, err.Error(), config.DotEnvPath,
"the failure must name the file it could not read",
)
_, present := os.LookupEnv(dotEnvKey)
assert.False(
t, present,
"a rejected file must apply nothing, not part of itself",
)
}
// TestLoadDotEnv_UnreadableFileAborts pins that only absence is
// tolerated. A .env that exists but cannot be read is a file the
// operator meant to be applied, so it fails like a malformed one
// rather than being treated as though it were not there.
func TestLoadDotEnv_UnreadableFileAborts(t *testing.T) {
t.Parallel()
// A directory in the file's place: open succeeds and the read
// fails, which no umask or root-ness can turn back into success
// the way a chmod could.
path := filepath.Join(t.TempDir(), config.DotEnvPath)
require.NoError(t, os.Mkdir(path, 0o750))
err := config.LoadDotEnvFileForTest(path)
require.Error(t, err)
require.ErrorIs(t, err, config.ErrDotEnvUnreadable)
}
// TestLoadDotEnv_ReadsTheWorkingDirectory pins the path LoadDotEnv
// itself opens, which the tests above bypass. It is relative to the
// process working directory, as it was under godotenv/autoload and as
// the README documents.
//
//nolint:paralleltest // t.Chdir moves the whole process.
func TestLoadDotEnv_ReadsTheWorkingDirectory(t *testing.T) {
unsetDotEnvKey(t)
dir := t.TempDir()
require.NoError(t, os.WriteFile(
filepath.Join(dir, config.DotEnvPath),
[]byte(dotEnvKey+"=from-working-directory\n"),
0o600,
))
t.Chdir(dir)
require.NoError(t, config.LoadDotEnv())
assert.Equal(t, "from-working-directory", os.Getenv(dotEnvKey))
}

View File

@@ -518,8 +518,20 @@ type badEnvValueCase struct {
}
// badEnvValueCases is the config.New table, kept out of the test body
// so the test itself stays readable.
// so the test itself stays readable. It is assembled from per-variable
// groups because one literal covering every variable outgrew the
// function-length budget.
func badEnvValueCases() []badEnvValueCase {
cases := listenerEnvValueCases()
cases = append(cases, flagEnvValueCases()...)
cases = append(cases, sentryEnvValueCases()...)
return cases
}
// listenerEnvValueCases covers the two variables that describe the
// HTTP listener.
func listenerEnvValueCases() []badEnvValueCase {
return []badEnvValueCase{
{
name: "valid PORT is used",
@@ -542,27 +554,6 @@ func badEnvValueCases() []badEnvValueCase {
value: "70000",
expectError: true,
},
{
name: "valid DEBUG is used",
key: envKeyDebug,
value: "true",
check: func(t *testing.T, cfg *config.Config) {
t.Helper()
assert.True(t, cfg.Debug)
},
},
{
name: "unparseable DEBUG aborts startup",
key: envKeyDebug,
value: "ture",
expectError: true,
},
{
name: "unparseable MAINTENANCE_MODE aborts startup",
key: envKeyMaintenanceMode,
value: "sometimes",
expectError: true,
},
{
name: "valid BIND_ADDRESS is used",
key: envKeyBindAddress,
@@ -595,6 +586,69 @@ func badEnvValueCases() []badEnvValueCase {
}
}
// flagEnvValueCases covers the boolean variables.
func flagEnvValueCases() []badEnvValueCase {
return []badEnvValueCase{
{
name: "valid DEBUG is used",
key: envKeyDebug,
value: "true",
check: func(t *testing.T, cfg *config.Config) {
t.Helper()
assert.True(t, cfg.Debug)
},
},
{
name: "unparseable DEBUG aborts startup",
key: envKeyDebug,
value: "ture",
expectError: true,
},
{
name: "unparseable MAINTENANCE_MODE aborts startup",
key: envKeyMaintenanceMode,
value: "sometimes",
expectError: true,
},
}
}
// sentryEnvValueCases covers SENTRY_DSN. The three rejected values are
// the ones measured on the defect: each initialised the SDK with an
// error and left the process serving with error reporting off.
func sentryEnvValueCases() []badEnvValueCase {
return []badEnvValueCase{
{
name: "valid SENTRY_DSN is used",
key: envKeySentryDSN,
value: validSentryDSN,
check: func(t *testing.T, cfg *config.Config) {
t.Helper()
assert.Equal(t, validSentryDSN, cfg.SentryDSN)
assert.True(t, cfg.SentryEnabled())
},
},
{
name: "unparseable SENTRY_DSN aborts startup",
key: envKeySentryDSN,
value: "not-a-dsn",
expectError: true,
},
{
name: "SENTRY_DSN that is not a URL aborts startup",
key: envKeySentryDSN,
value: "%%%",
expectError: true,
},
{
name: "keyless SENTRY_DSN aborts startup",
key: envKeySentryDSN,
value: "https://example.invalid/1",
expectError: true,
},
}
}
// TestNewUsesDefaultsWhenUnset proves the fail-loud behaviour did not
// break the legitimate unset case: absent variables still get their
// documented defaults.
@@ -603,7 +657,7 @@ func TestNewUsesDefaultsWhenUnset(t *testing.T) {
for _, key := range []string{
envKeyPort, envKeyDebug, envKeyMaintenanceMode,
envKeyBindAddress,
envKeyBindAddress, envKeySentryDSN,
} {
require.NoError(t, os.Unsetenv(key))
}
@@ -626,4 +680,9 @@ func TestNewUsesDefaultsWhenUnset(t *testing.T) {
t, config.DefaultBindAddressForTest, cfg.BindAddress,
)
assert.Equal(t, bindAddressDefault, cfg.BindAddress)
// An absent SENTRY_DSN is the common case and must stay a normal
// start with error reporting off, not a refusal.
assert.Empty(t, cfg.SentryDSN)
assert.False(t, cfg.SentryEnabled())
}

View File

@@ -51,6 +51,18 @@ func EnvPortForTest(key string, defaultValue int) (int, error) {
return envPort(key, defaultValue)
}
// EnvSentryDSNForTest exposes envSentryDSN.
func EnvSentryDSNForTest(key string) (string, error) {
return envSentryDSN(key)
}
// LoadDotEnvFileForTest exposes the loader LoadDotEnv runs, over a
// caller-named file rather than the process working directory, so
// each .env state can be covered without moving the test process.
func LoadDotEnvFileForTest(path string) error {
return loadDotEnvFile(path)
}
// EnvBindAddressForTest exposes envBindAddress.
func EnvBindAddressForTest(key, defaultValue string) (string, error) {
return envBindAddress(key, defaultValue)

View File

@@ -0,0 +1,141 @@
package config_test
import (
"os"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/config"
)
// envKeySentryDSN is the variable envSentryDSN reads in production.
const envKeySentryDSN = "SENTRY_DSN"
// validSentryDSN is a syntactically complete DSN. The host is under
// .invalid (RFC 2606), so nothing a test builds around it can reach a
// real Sentry installation.
const validSentryDSN = "https://abc123@sentry.invalid/42"
// envSentryDSNCase is one row of the envSentryDSN table.
type envSentryDSNCase struct {
name string
set bool
value string
expectError bool
expected string
}
// envSentryDSNCases is the envSentryDSN table. The three invalid
// values are the ones measured on the defect: each initialised the SDK
// with an error and left the process serving with reporting off.
func envSentryDSNCases() []envSentryDSNCase {
return []envSentryDSNCase{
{
name: "unset means reporting off",
expected: "",
},
{
name: "empty means reporting off",
set: true,
value: "",
expected: "",
},
{
name: "whitespace means reporting off",
set: true,
value: " ",
expected: "",
},
{
name: "a valid DSN is kept",
set: true,
value: validSentryDSN,
expected: validSentryDSN,
},
{
name: "surrounding whitespace is trimmed",
set: true,
value: " " + validSentryDSN + "\t",
expected: validSentryDSN,
},
{
name: "a value that is not a URL is rejected",
set: true,
value: "not-a-dsn",
expectError: true,
},
{
name: "an unparseable URL is rejected",
set: true,
value: "%%%",
expectError: true,
},
{
name: "a DSN without a public key is rejected",
set: true,
value: "https://example.invalid/1",
expectError: true,
},
{
name: "a DSN without a project id is rejected",
set: true,
value: "https://abc123@sentry.invalid/",
expectError: true,
},
{
name: "a non-HTTP scheme is rejected",
set: true,
value: "ftp://abc123@sentry.invalid/42",
expectError: true,
},
}
}
// TestEnvSentryDSN covers the helper directly. What it pins beyond the
// value is the failure shape: a set-but-unparseable DSN names the
// variable and the value, exactly as the other fail-loud helpers do,
// so an operator reads the fix off the message.
func TestEnvSentryDSN(t *testing.T) {
for _, tt := range envSentryDSNCases() {
t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests.
if tt.set {
t.Setenv(envKeySentryDSN, tt.value)
} else {
require.NoError(t, os.Unsetenv(envKeySentryDSN))
}
got, err := config.EnvSentryDSNForTest(envKeySentryDSN)
if tt.expectError {
require.Error(t, err)
require.ErrorIs(t, err, config.ErrInvalidSentryDSN)
assert.Contains(t, err.Error(), envKeySentryDSN)
assert.Contains(t, err.Error(), tt.value)
assert.Empty(t, got)
return
}
require.NoError(t, err)
assert.Equal(t, tt.expected, got)
})
}
}
// TestSentryEnabled_TracksTheDSN pins that the one method answering
// "is anything being reported" agrees with the DSN in every state. The
// startup log, the SDK initialisation and the sentryhttp middleware
// all read it, so a log field cannot report reporting as on while
// nothing is sending.
func TestSentryEnabled_TracksTheDSN(t *testing.T) {
t.Parallel()
assert.False(t, (&config.Config{}).SentryEnabled())
assert.True(
t,
(&config.Config{SentryDSN: validSentryDSN}).SentryEnabled(),
)
}

View File

@@ -93,11 +93,19 @@ func TestMainDatabaseFilesAreOwnerOnly(t *testing.T) {
t, filepath.Join(dataDir, database.MainDBFileName),
)
// The data directory stays group-readable. Deployments may rely on
// the group bit; the file mode is the barrier, not the directory.
// The data directory grants nothing to `other`. Asserted as a
// property rather than as an exact 0750, because MkdirAll applies
// the ambient umask: the exact mode is the developer's umask as
// much as the application's request, and pinning it would make
// `make check` pass or fail on where it is run. The group bits are
// deliberately left unasserted — deployments may rely on them.
info, err := os.Stat(dataDir)
require.NoError(t, err)
assert.Equal(t, fs.FileMode(0o750), info.Mode().Perm())
assert.Zero(
t,
info.Mode().Perm()&0o007,
"the data directory must not be world-accessible",
)
}
// TestPerWebhookEventDatabaseFilesAreOwnerOnly covers the events-*.db

View File

@@ -4,6 +4,7 @@ package delivery
import (
"context"
"errors"
"fmt"
"log/slog"
"net/http"
@@ -439,7 +440,7 @@ func (e *Engine) processNewTask(
event := buildEventFromTask(task)
event, err = e.resolveEventBody(
event, err = e.hydrateEvent(
webhookDB, event, task,
)
if err != nil {
@@ -505,9 +506,13 @@ func (e *Engine) processRetryTask(
return
}
if e.abandonRetryForMissingTarget(webhookDB, d, task) {
return
}
event := buildEventFromTask(task)
event, err = e.resolveEventBody(
event, err = e.hydrateEvent(
webhookDB, event, task,
)
if err != nil {
@@ -529,6 +534,64 @@ func (e *Engine) processRetryTask(
e.processDelivery(ctx, webhookDB, d, task)
}
// abandonRetryForMissingTarget stops a retry chain whose target has
// been deleted, and reports whether it did.
//
// A scheduled retry lives in memory as a time.AfterFunc holding the
// target's configuration as it was when the chain began, and nothing
// else on this path reads the target row. Without this check a
// deletion stops nothing: the timer keeps firing and keeps sending to
// the destination the operator removed, for the whole remaining
// backoff chain. Terminalising in the recovery and sweep paths alone
// is not enough, because those only see the delivery once nothing
// holds it in memory — which is to say after a restart.
//
// The worker already owns this delivery, so the terminal write happens
// here directly, exactly as a target's own Deliver fails one. Claiming
// it again through the recovery gate would only fail against the
// reference the worker itself is holding.
//
// A lookup that fails for any other reason is not a deletion — it is
// the main database being unreadable — and the delivery goes ahead as
// it did before. A guard that terminally failed deliveries on a
// transient fault would be worse than the bug it fixes.
func (e *Engine) abandonRetryForMissingTarget(
webhookDB *gorm.DB,
d *database.Delivery,
task *Task,
) bool {
_, err := e.loadTarget(task.TargetID)
if err == nil {
return false
}
if !errors.Is(err, gorm.ErrRecordNotFound) {
e.log.Warn(
"could not confirm the target of a retrying "+
"delivery still exists; attempting anyway",
"delivery_id", task.DeliveryID,
"target_id", task.TargetID,
"error", err,
)
return false
}
targetType, reason := e.missingTargetReason(task.TargetID)
e.log.Warn(
"abandoning scheduled retry: target is gone",
"webhook_id", task.WebhookID,
"delivery_id", task.DeliveryID,
"target_id", task.TargetID,
"target_type", targetType,
)
e.failDelivery(webhookDB, d, targetType, reason)
return true
}
func (e *Engine) recoverInFlight(ctx context.Context) {
var webhookIDs []string
@@ -633,6 +696,20 @@ func (e *Engine) recoverSingleRetry(
) {
target, err := e.loadTarget(d.TargetID)
if err != nil {
// A target that is merely gone is an operator action with a
// terminal answer. Any other failure is the main database
// refusing to read, which is transient and must leave the
// delivery alone: failing every retrying delivery of every
// 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,
)
return
}
e.log.Error(
"failed to load target for retrying "+
"delivery recovery",
@@ -1028,6 +1105,16 @@ func (e *Engine) sweepSingleRetry(
) {
target, err := e.loadTarget(d.TargetID)
if err != nil {
// Deleted is terminal, unreadable is not; see
// recoverSingleRetry.
if errors.Is(err, gorm.ErrRecordNotFound) {
e.failMissingTargetRetry(
webhookDB, webhookID, d,
)
return
}
e.log.Error(
"retry sweep: failed to load target",
"delivery_id", d.ID,
@@ -1134,6 +1221,113 @@ func (e *Engine) failUnretryableRetry(
target.Type,
)
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.
//
// Until it existed both paths logged the failed lookup and returned,
// which left the delivery retrying 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(
webhookDB *gorm.DB,
webhookID string,
d *database.Delivery,
) {
// Terminal, and reached from the recovery paths, so it takes
// ownership like every other write they make.
if !e.inflight.retainIdle(d.ID) {
return
}
defer e.inflight.release(d.ID)
targetType, reason := e.missingTargetReason(d.TargetID)
e.log.Warn(
"failing orphaned retrying delivery: "+
"its target no longer exists",
"webhook_id", webhookID,
"delivery_id", d.ID,
"target_id", d.TargetID,
"target_type", targetType,
)
e.failDelivery(webhookDB, d, targetType, reason)
}
// missingTargetReason describes a target id that no longer resolves,
// and returns the type of the deleted row where there still is one.
//
// The lookup is Unscoped because deletes are soft: the row survives
// with deleted_at set, invisible to loadTarget's default scope.
// Reading it is what separates "you deleted this target" from "this id
// never named a row" — different things to whoever reads the event
// log, and only the first is something an operator did. The widened
// scope is deliberately confined to this terminal path: the engine's
// normal target loading must go on refusing a deleted target, or
// deleting one would stop nothing.
//
// The type comes back so the caller can label the delivery's status
// transition with it. Where the row is gone entirely there is no type
// to give, and updateDeliveryStatus leaves the counter alone rather
// than opening a series named by the empty string.
func (e *Engine) missingTargetReason(
targetID string,
) (database.TargetType, string) {
var target database.Target
err := e.database.DB().Unscoped().
First(&target, "id = ?", targetID).Error
if err != nil {
return "", fmt.Sprintf(
"target %s no longer exists; the delivery "+
"cannot be retried and 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",
target.Name, target.Type,
)
}
// failDelivery records why a delivery is over and then marks it
// failed. The caller must already own the delivery: every call site is
// either a worker holding the reference runTask took, or a recovery
// path that took one through retainIdle.
//
// The result row is written first and a failure to write it stops the
// transition, which is what keeps a delivery from ending failed with
// an empty event log — the state that leaves an operator with nothing
// but a server log line to work out what happened. A delivery whose
// reason could not be recorded stays in the non-terminal state it
// already holds, where the sweep will find it again; see
// bookkeepingFailed.
//
// The target type is a parameter rather than read off d because the
// orphaned-retry callers deliberately hold a delivery loaded without
// its Target relation: populating d.Target would make GORM's
// SaveBeforeAssociations upsert the whole target row — plaintext
// config, which for a slack target is the credential — into the
// per-webhook event database. See
// https://git.eeqj.de/sneak/webhooker/issues/206.
func (e *Engine) failDelivery(
webhookDB *gorm.DB,
d *database.Delivery,
targetType database.TargetType,
reason string,
) {
err := e.recordResult(
webhookDB,
d,
@@ -1150,14 +1344,8 @@ func (e *Engine) failUnretryableRetry(
return
}
// The type is passed rather than assigned onto d: the delivery
// is loaded here without its target relation, and populating
// d.Target would make GORM's SaveBeforeAssociations upsert the
// whole target row — plaintext config, which for a slack target
// is the credential — into the per-webhook event database. See
// https://git.eeqj.de/sneak/webhooker/issues/206.
e.settleStatus(
webhookDB, d, target.Type,
webhookDB, d, targetType,
database.DeliveryStatusFailed,
)
}
@@ -1178,9 +1366,19 @@ func (e *Engine) processDelivery(
"type", d.Target.Type,
)
e.settleStatus(
// The reason is recorded, not just logged. This branch used
// to fail the delivery with no DeliveryResult at all, which
// showed in the event log as "failed, no attempts recorded
// yet" and left one server log line as the only account of
// why anywhere.
e.failDelivery(
webhookDB, d, d.Target.Type,
database.DeliveryStatusFailed,
fmt.Sprintf(
"unknown target type %q: this build has no "+
"delivery implementation for it, so no "+
"attempt was made",
d.Target.Type,
),
)
return
@@ -1349,6 +1547,11 @@ func truncate(s string, maxLen int) string {
// --- Helper functions ---
// buildEventFromTask reconstructs the event a Task describes, as far
// as the Task itself goes. The fields it cannot fill — the body when
// it was too large to inline, and the receipt time, which no Task
// carries — come from the stored row in hydrateEvent, which every
// caller of this function runs next.
func buildEventFromTask(task *Task) database.Event {
event := database.Event{
EntrypointID: task.EntrypointID,
@@ -1376,29 +1579,67 @@ func buildTargetFromTask(task *Task) database.Target {
return target
}
func (e *Engine) resolveEventBody(
// hydrateEvent fills in the event fields a Task does not carry, by
// reading the stored event row.
//
// CreatedAt is the event's receipt time and lives only in that row.
// The Slack target renders it into every message it sends, so an
// unhydrated event puts the zero time in front of a human on every
// notification the product delivers. See
// https://git.eeqj.de/sneak/webhooker/issues/257.
//
// The body comes from the same row when the Task did not inline it,
// which is the case for a body at or above MaxInlineBodySize.
//
// A read failure is fatal to the delivery only when the body depended
// on it. When the Task inlined the body, the delivery has everything
// it needs to be sent and goes ahead with the timestamp unset: the row
// can be gone under a retention reap while a queued delivery still
// holds its body, and dropping a deliverable event to protect one
// metadata field would be a worse failure than the one it prevents.
func (e *Engine) hydrateEvent(
webhookDB *gorm.DB,
event database.Event,
task *Task,
) (database.Event, error) {
if task.Body != nil {
columns := []string{"created_at"}
if task.Body == nil {
columns = append(columns, "body")
}
var dbEvent database.Event
err := webhookDB.Select(columns).
First(&dbEvent, "id = ?", task.EventID).Error
if err != nil {
if task.Body == nil {
return event, fmt.Errorf(
"fetching event body: %w", err,
)
}
e.log.Warn(
"could not read the stored event; delivering "+
"the inlined body without its receipt time",
"event_id", task.EventID,
"delivery_id", task.DeliveryID,
"error", err,
)
event.Body = *task.Body
return event, nil
}
var dbEvent database.Event
event.CreatedAt = dbEvent.CreatedAt
err := webhookDB.Select("body").
First(&dbEvent, "id = ?", task.EventID).Error
if err != nil {
return event, fmt.Errorf(
"fetching event body: %w", err,
)
if task.Body != nil {
event.Body = *task.Body
} else {
event.Body = dbEvent.Body
}
event.Body = dbEvent.Body
return event, nil
}

View File

@@ -377,6 +377,17 @@ func TestProcessRetryTask_SuccessfulRetry(t *testing.T) {
bodyStr := event.Body
cfg := iHTTPConfig(ts.URL)
// The target row exists because the engine confirms a scheduled
// retry's target has not been deleted before it runs it. A retry
// task whose target id names no row at all is a state the service
// does not produce: the handler read that target to build the
// task. See https://git.eeqj.de/sneak/webhooker/issues/107.
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, "retry-target",
database.TargetTypeHTTP, cfg, 5,
)
task := iTask(
d, event, s.WebhookID, targetID,
"retry-target", cfg, 5, 2, &bodyStr,
@@ -456,6 +467,12 @@ func TestProcessRetryTask_LargeBody_FetchFromDB(
)
cfg := iHTTPConfig(ts.URL)
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, "retry-large",
database.TargetTypeHTTP, cfg, 5,
)
task := iTask(
d, event, s.WebhookID, targetID,
"retry-large", cfg, 5, 2, nil,
@@ -558,6 +575,12 @@ func TestWorkerLifecycle_ProcessesRetryChannel(
bodyStr := event.Body
cfg := iHTTPConfig(ts.URL)
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, "retry-chan-test",
database.TargetTypeHTTP, cfg, 5,
)
task := iTask(
d, event, s.WebhookID, targetID,
"retry-chan-test", cfg, 5, 2, &bodyStr,

View File

@@ -96,7 +96,15 @@ func TestEventDBHoldsNoTargetRows(t *testing.T) {
)
assertNoTargetRows(t, dbPath)
// A retry.
// A retry. Its target exists in the main database, because the
// engine confirms a scheduled retry's target has not been
// deleted before running it; see
// https://git.eeqj.de/sneak/webhooker/issues/107.
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, "leaky-target",
database.TargetTypeHTTP, cfg, 5,
)
rd := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,

View File

@@ -0,0 +1,442 @@
package delivery_test
import (
"context"
"encoding/json"
"io"
"net/http"
"net/http/httptest"
"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"
"sneak.berlin/go/webhooker/internal/delivery"
)
// tsEventCreatedAt is the receipt time seeded on the events these
// tests deliver. It is far enough from both the zero time and from
// now that neither can be mistaken for it.
func tsEventCreatedAt() time.Time {
return time.Date(
2026, time.March, 4, 5, 6, 7, 0, time.UTC,
)
}
// tsZeroStamp is what a Slack message renders when the event handed
// to FormatSlackMessage carries no CreatedAt.
const tsZeroStamp = "*Timestamp:* `0001-01-01T00:00:00Z`"
// tsEventBody is the body seeded on every event in this file. It is
// small enough that a Task can inline it.
const tsEventBody = `{"hello":"world"}`
// tsUndeliverableHook stands in for a Slack incoming webhook on the
// tests that never send: the config parser requires a URL, but no
// request is made.
const tsUndeliverableHook = "https://hooks.slack.com/services/T/B/x"
// tsSink is a stand-in Slack incoming webhook that records the raw
// body posted to it.
type tsSink struct {
*httptest.Server
bodies chan []byte
}
func newTSSink(t *testing.T) *tsSink {
t.Helper()
s := &tsSink{bodies: make(chan []byte, 8)}
s.Server = httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
select {
case s.bodies <- body:
default:
}
w.WriteHeader(http.StatusOK)
},
))
t.Cleanup(s.Close)
return s
}
// text returns the Slack message text from the single payload the
// sink received.
func (s *tsSink) text(t *testing.T) string {
t.Helper()
select {
case raw := <-s.bodies:
t.Logf("raw slack payload: %s", raw)
var payload struct {
Text string `json:"text"`
}
require.NoError(t, json.Unmarshal(raw, &payload))
return payload.Text
case <-time.After(5 * time.Second):
t.Fatal("slack sink received no payload")
return ""
}
}
func tsSlackConfig(t *testing.T, url string) string {
t.Helper()
data, err := json.Marshal(
delivery.SlackTargetConfig{WebhookURL: url},
)
require.NoError(t, err)
return string(data)
}
// tsSeedEvent writes an event whose CreatedAt is tsEventCreatedAt
// rather than the write time, so an assertion on the rendered
// timestamp cannot pass by accident against "roughly now".
func tsSeedEvent(
t *testing.T, db *gorm.DB, webhookID string,
) database.Event {
t.Helper()
event := database.Event{
WebhookID: webhookID,
EntrypointID: uuid.New().String(),
Method: http.MethodPost,
Headers: `{}`,
Body: tsEventBody,
ContentType: "application/json",
}
event.ID = uuid.New().String()
event.CreatedAt = tsEventCreatedAt()
event.UpdatedAt = tsEventCreatedAt()
require.NoError(t, db.Create(&event).Error)
var stored database.Event
require.NoError(t,
db.First(&stored, "id = ?", event.ID).Error,
)
require.Equal(t,
tsEventCreatedAt().UTC(), stored.CreatedAt.UTC(),
"seeded created_at did not round-trip",
)
return event
}
// tsSeedTarget writes the slack target row into the main database.
// The retry path confirms the target still exists before sending.
func tsSeedTarget(
t *testing.T, mainDB *gorm.DB, webhookID, config string,
) database.Target {
t.Helper()
target := database.Target{
WebhookID: webhookID,
Name: "slack-sink",
Type: database.TargetTypeSlack,
Config: config,
Active: true,
}
require.NoError(t, mainDB.Create(&target).Error)
return target
}
func tsTask(
d database.Delivery,
event database.Event,
webhookID string,
target database.Target,
attemptNum int,
body *string,
) delivery.Task {
return delivery.Task{
DeliveryID: d.ID,
EventID: event.ID,
WebhookID: webhookID,
EntrypointID: event.EntrypointID,
TargetID: target.ID,
TargetName: target.Name,
TargetType: database.TargetTypeSlack,
TargetConfig: target.Config,
MaxRetries: 0,
Method: event.Method,
Headers: event.Headers,
ContentType: event.ContentType,
Body: body,
AttemptNum: attemptNum,
}
}
func tsAssertRealTimestamp(t *testing.T, text string) {
t.Helper()
assert.NotContains(t, text, tsZeroStamp,
"slack message carries the zero timestamp",
)
assert.Contains(t, text,
"*Timestamp:* `"+
tsEventCreatedAt().UTC().Format(time.RFC3339)+"`",
"slack message does not carry the event's receipt time",
)
}
// tsCase is one end-to-end delivery of a seeded event to a slack
// sink, over whichever engine path `process` names.
type tsCase struct {
// status is the delivery row's status before the engine runs.
// The retry path refuses a delivery that is not retrying.
status database.DeliveryStatus
// inlineBody mirrors a Task built for a body under
// MaxInlineBodySize. When false the engine reads the body back
// from the stored row.
inlineBody bool
attemptNum int
process func(
ctx context.Context, e *delivery.Engine, task *delivery.Task,
)
}
// run delivers one event through the named path and returns the
// Slack message text the sink received.
func (c tsCase) run(t *testing.T) (iSetup, database.Delivery, string) {
t.Helper()
s := newISetup(t)
sink := newTSSink(t)
cfg := tsSlackConfig(t, sink.URL)
target := tsSeedTarget(t, s.MainDB, s.WebhookID, cfg)
event := tsSeedEvent(t, s.WebhookDB, s.WebhookID)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, target.ID, c.status,
)
var body *string
if c.inlineBody {
bodyStr := event.Body
body = &bodyStr
}
task := tsTask(
d, event, s.WebhookID, target, c.attemptNum, body,
)
c.process(context.TODO(), s.Engine, &task)
return s, d, sink.text(t)
}
// TestSlackFirstAttemptCarriesEventTimestamp covers the path an
// event takes on its first delivery: the task comes from the
// receiver and the engine reconstructs the event from it.
func TestSlackFirstAttemptCarriesEventTimestamp(t *testing.T) {
t.Parallel()
s, d, text := tsCase{
status: database.DeliveryStatusPending,
inlineBody: true,
attemptNum: 1,
process: func(
ctx context.Context,
e *delivery.Engine,
task *delivery.Task,
) {
e.ExportProcessNewTask(ctx, task)
},
}.run(t)
tsAssertRealTimestamp(t, text)
iAssertStatus(t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered,
)
}
// TestSlackFirstAttemptLargeBodyCarriesEventTimestamp covers the
// first-attempt path for an event whose body exceeded
// MaxInlineBodySize, so the task carries no body and the engine
// reads it back from the stored row.
func TestSlackFirstAttemptLargeBodyCarriesEventTimestamp(
t *testing.T,
) {
t.Parallel()
_, _, text := tsCase{
status: database.DeliveryStatusPending,
inlineBody: false,
attemptNum: 1,
process: func(
ctx context.Context,
e *delivery.Engine,
task *delivery.Task,
) {
e.ExportProcessNewTask(ctx, task)
},
}.run(t)
tsAssertRealTimestamp(t, text)
}
// TestSlackRetryCarriesEventTimestamp covers the retry path, which
// reconstructs the event from the same task the first attempt used.
func TestSlackRetryCarriesEventTimestamp(t *testing.T) {
t.Parallel()
s, d, text := tsCase{
status: database.DeliveryStatusRetrying,
inlineBody: true,
attemptNum: 2,
process: func(
ctx context.Context,
e *delivery.Engine,
task *delivery.Task,
) {
e.ExportProcessRetryTask(ctx, task)
},
}.run(t)
tsAssertRealTimestamp(t, text)
iAssertStatus(t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered,
)
}
// TestFormatSlackMessageOverTaskReconstructedEvent asserts on the
// formatted message directly, over the event the delivery paths
// reconstruct from a Task. It is the unit-level guard under the
// end-to-end tests: revert the CreatedAt population in hydrateEvent
// and this fails on the zero timestamp.
func TestFormatSlackMessageOverTaskReconstructedEvent(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
cfg := tsSlackConfig(t, tsUndeliverableHook)
target := tsSeedTarget(t, s.MainDB, s.WebhookID, cfg)
event := tsSeedEvent(t, s.WebhookDB, s.WebhookID)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, target.ID,
database.DeliveryStatusPending,
)
bodyStr := event.Body
task := tsTask(d, event, s.WebhookID, target, 1, &bodyStr)
rebuilt, err := s.Engine.ExportEventForTask(
s.WebhookDB, &task,
)
require.NoError(t, err)
assert.False(t, rebuilt.CreatedAt.IsZero(),
"reconstructed event carries the zero time",
)
assert.Equal(t,
tsEventCreatedAt().UTC(), rebuilt.CreatedAt.UTC(),
)
tsAssertRealTimestamp(
t, delivery.FormatSlackMessage(&rebuilt),
)
}
// TestFormatSlackMessageZeroTimestamp asserts the rendering choice
// directly, without going through the engine: a zero CreatedAt (the
// shape a reaped-row fallback produces) renders as "unknown" rather
// than the year-1 zero time, while a real CreatedAt still renders as
// RFC3339.
func TestFormatSlackMessageZeroTimestamp(t *testing.T) {
t.Parallel()
zeroEvent := database.Event{
Method: http.MethodPost,
ContentType: testContentType,
Body: tsEventBody,
}
zeroText := delivery.FormatSlackMessage(&zeroEvent)
assert.NotContains(t, zeroText, "0001-01-01",
"slack message carries the zero-time year",
)
assert.Contains(t, zeroText, "*Timestamp:* `unknown`",
"slack message does not mark an unset receipt time as unknown",
)
nonZeroEvent := zeroEvent
nonZeroEvent.CreatedAt = tsEventCreatedAt()
nonZeroText := delivery.FormatSlackMessage(&nonZeroEvent)
assert.Contains(t, nonZeroText,
"*Timestamp:* `"+
tsEventCreatedAt().UTC().Format(time.RFC3339)+"`",
"slack message does not render a real receipt time as RFC3339",
)
}
// TestEventReconstructionSurvivesAReapedRow pins the fallback: an
// event row reaped by retention while its delivery still holds the
// body inline is still delivered, with the receipt time unset,
// rather than dropped.
func TestEventReconstructionSurvivesAReapedRow(t *testing.T) {
t.Parallel()
s := newISetup(t)
cfg := tsSlackConfig(t, tsUndeliverableHook)
target := tsSeedTarget(t, s.MainDB, s.WebhookID, cfg)
event := tsSeedEvent(t, s.WebhookDB, s.WebhookID)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, target.ID,
database.DeliveryStatusPending,
)
bodyStr := event.Body
task := tsTask(d, event, s.WebhookID, target, 1, &bodyStr)
require.NoError(t, s.WebhookDB.Unscoped().Delete(
&database.Event{}, "id = ?", event.ID,
).Error)
rebuilt, err := s.Engine.ExportEventForTask(
s.WebhookDB, &task,
)
require.NoError(t, err)
assert.Equal(t, bodyStr, rebuilt.Body)
assert.True(t, rebuilt.CreatedAt.IsZero())
// A task with no inlined body has nothing left to deliver, so
// the same reaped row is an error there.
noBody := task
noBody.Body = nil
_, err = s.Engine.ExportEventForTask(s.WebhookDB, &noBody)
require.Error(t, err)
}

View File

@@ -151,6 +151,16 @@ func (e *Engine) ExportProcessRetryTask(
e.processRetryTask(ctx, task)
}
// ExportEventForTask exposes the event reconstruction the delivery
// paths run: buildEventFromTask followed by hydrateEvent.
func (e *Engine) ExportEventForTask(
webhookDB *gorm.DB, task *Task,
) (database.Event, error) {
return e.hydrateEvent(
webhookDB, buildEventFromTask(task), task,
)
}
// ExportProcessDelivery exposes processDelivery.
func (e *Engine) ExportProcessDelivery(
ctx context.Context,

View File

@@ -229,6 +229,13 @@ func mExhaustRetries(t *testing.T, s iSetup) {
body := event.Body
cfg := iHTTPConfig(ts.URL)
// The retry below is only run if its target still exists; see
// https://git.eeqj.de/sneak/webhooker/issues/107.
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, "metrics-fail",
database.TargetTypeHTTP, cfg, 2,
)
first := iTask(
d, event, s.WebhookID, targetID,
"metrics-fail", cfg, 2, 1, &body,
@@ -289,6 +296,13 @@ func TestDeliveryMetrics_CircuitBreakerGauge(t *testing.T) {
// rather than the budget is what stops the delivery.
maxRetries := delivery.ExportDefaultFailureThreshold + 5
// The retries below are only run if their target still exists;
// see https://git.eeqj.de/sneak/webhooker/issues/107.
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, "metrics-trip",
database.TargetTypeHTTP, cfg, maxRetries,
)
first := iTask(
d, event, s.WebhookID, targetID,
"metrics-trip", cfg, maxRetries, 1, &body,
@@ -353,6 +367,11 @@ func TestDeliveryMetrics_BreakerBlockedIsNotAnAttempt(
cfg := iHTTPConfig(ts.URL)
maxRetries := delivery.ExportDefaultFailureThreshold + 5
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, "metrics-blocked",
database.TargetTypeHTTP, cfg, maxRetries,
)
first := iTask(
d, event, s.WebhookID, targetID,
"metrics-blocked", cfg, maxRetries, 1, &body,

View File

@@ -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,
) {

View File

@@ -231,10 +231,15 @@ func FormatSlackMessage(
event.ContentType,
)
timestamp := "unknown"
if !event.CreatedAt.IsZero() {
timestamp = event.CreatedAt.UTC().Format(time.RFC3339)
}
fmt.Fprintf(
&b,
"*Timestamp:* `%s`\n",
event.CreatedAt.UTC().Format(time.RFC3339),
timestamp,
)
fmt.Fprintf(

View File

@@ -0,0 +1,531 @@
package delivery_test
import (
"context"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
)
// The two terminal-state gaps of
// 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.
// 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.
//
// 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
// named a row: the surviving row is invisible to a scoped read.
func tSeedDeletedTarget(
t *testing.T,
s iSetup,
name, url string,
) string {
t.Helper()
targetID := uuid.New().String()
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, name,
database.TargetTypeHTTP, iHTTPConfig(url), 5,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"target":"deleted"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
iSeedFailedResult(t, s.WebhookDB, d.ID)
require.NoError(t, s.MainDB.Delete(
&database.Target{}, "id = ?", targetID,
).Error)
var scoped, unscoped int64
require.NoError(t, s.MainDB.
Model(&database.Target{}).
Where("id = ?", targetID).
Count(&scoped).Error)
require.NoError(t, s.MainDB.Unscoped().
Model(&database.Target{}).
Where("id = ?", targetID).
Count(&unscoped).Error)
require.Zero(t, scoped,
"the deleted target is still visible to a scoped read",
)
require.Equal(t, int64(1), unscoped,
"the delete was hard, so this test proves nothing about "+
"the soft-delete case it exists for",
)
return d.ID
}
// tLastResult returns a delivery's final recorded attempt, asserting
// the expected number of them.
func tLastResult(
t *testing.T,
s iSetup,
deliveryID string,
want int,
) database.DeliveryResult {
t.Helper()
results := iResults(t, s.WebhookDB, deliveryID)
require.Len(t, results, want)
return results[want-1]
}
// --- 1. A failure with nothing recorded ---
func TestProcessDelivery_UnknownTargetType_RecordsWhy(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"unknown":"type"}`,
)
seeded := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
target := database.Target{
Name: "mystery",
Type: tUnknownType,
Config: iHTTPConfig("http://example.com/hook"),
}
target.ID = targetID
d := database.Delivery{
EventID: event.ID,
TargetID: targetID,
Status: database.DeliveryStatusPending,
Event: event,
Target: target,
}
d.ID = seeded.ID
body := event.Body
task := iTask(
seeded, event, s.WebhookID, targetID, "mystery",
target.Config, 0, 1, &body,
)
task.TargetType = tUnknownType
s.Engine.ExportProcessDelivery(
context.Background(), s.WebhookDB, &d, &task,
)
iAssertStatus(
t, s.WebhookDB, d.ID, database.DeliveryStatusFailed,
)
last := tLastResult(t, s, d.ID, 1)
assert.False(t, last.Success)
assert.Equal(t, 1, last.AttemptNum)
assert.Contains(t, last.Error, string(tUnknownType),
"the recorded reason does not name the offending type",
)
}
// --- 2. A retrying delivery whose target is gone ---
func TestRecoverSingleRetry_TargetDeleted(t *testing.T) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "deleted-target-recovery",
)
deliveryID := tSeedDeletedTarget(
t, s, "gone-on-recovery", "http://example.com/hook",
)
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-on-recovery")
assert.Contains(t, last.Error, "was deleted")
assert.Empty(t, s.Engine.ExportRetryCh(),
"a delivery whose target is gone was rescheduled",
)
assert.Zero(t, s.Engine.ExportInflightHeld(),
"the terminal path leaked its ownership reference",
)
}
func TestSweepSingleRetry_TargetDeleted(t *testing.T) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "deleted-target-sweep",
)
deliveryID := tSeedDeletedTarget(
t, s, "gone-on-sweep", "http://example.com/hook",
)
// Twice, because the bug was an error the sweep repeated every
// minute for the life of the database: the second sweep must
// find nothing left to do.
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
iAssertStatus(
t, s.WebhookDB, deliveryID,
database.DeliveryStatusFailed,
)
last := tLastResult(t, s, deliveryID, 2)
assert.Contains(t, last.Error, "gone-on-sweep")
assert.Contains(t, last.Error, "was deleted")
assert.Empty(t, s.Engine.ExportRetryCh())
assert.Zero(t, s.Engine.ExportInflightHeld())
}
// TestSweepSingleRetry_TargetNeverExisted covers the other half of the
// soft-delete distinction: an id with no row at all, deleted or
// otherwise, must not be reported as something the operator deleted.
func TestSweepSingleRetry_TargetNeverExisted(t *testing.T) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "target-never-existed",
)
targetID := uuid.New().String()
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"target":"absent"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
iSeedFailedResult(t, s.WebhookDB, d.ID)
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
iAssertStatus(
t, s.WebhookDB, d.ID, database.DeliveryStatusFailed,
)
last := tLastResult(t, s, d.ID, 2)
assert.Contains(t, last.Error, targetID)
assert.Contains(t, last.Error, "no longer exists")
assert.NotContains(t, last.Error, "was deleted",
"an id that never named a row was reported as a deletion",
)
}
// TestFailMissingTargetRetry_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(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "no-target-row-deleted",
)
hookURL := "https://hooks.slack.com/services/T00/B00/x"
deliveryID := tSeedDeletedTarget(
t, s, "credential-bearing", hookURL,
)
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
iAssertStatus(
t, s.WebhookDB, deliveryID,
database.DeliveryStatusFailed,
)
var configs []string
require.NoError(t, s.WebhookDB.
Table("targets").
Pluck("config", &configs).Error)
assert.Empty(t, configs,
"the deleted-target terminal path wrote a target row "+
"into the per-webhook event database",
)
}
// --- 3. The scheduled retry chain ---
// tRetryChainSetup wires a counting sink and a retrying delivery
// against a live target pointing at it, and returns the task a
// scheduled retry would carry — config and all, snapshotted as
// ScheduleRetry snapshots it.
func tRetryChainSetup(
t *testing.T,
s iSetup,
name string,
hits *atomic.Int64,
) (delivery.Task, string) {
t.Helper()
ts := httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
hits.Add(1)
w.WriteHeader(http.StatusOK)
},
))
t.Cleanup(ts.Close)
iCreateWebhook(t, s.MainDB, s.WebhookID, name)
targetID := uuid.New().String()
cfg := iHTTPConfig(ts.URL)
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, name,
database.TargetTypeHTTP, cfg, 5,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"chain":"retry"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
iSeedFailedResult(t, s.WebhookDB, d.ID)
body := event.Body
return iTask(
d, event, s.WebhookID, targetID, name, cfg, 5, 2, &body,
), targetID
}
// TestProcessRetryTask_TargetDeleted_MakesNoAttempt is the half the
// deployability audit found worse than filed: terminalising on
// recovery and sweep alone leaves the already-scheduled timer chain
// running, and it holds the target's configuration from before the
// deletion, so it goes on sending to a destination that was removed.
func TestProcessRetryTask_TargetDeleted_MakesNoAttempt(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
var hits atomic.Int64
task, targetID := tRetryChainSetup(
t, s, "gone-mid-chain", &hits,
)
require.NoError(t, s.MainDB.Delete(
&database.Target{}, "id = ?", targetID,
).Error)
s.Engine.ExportProcessRetryTask(
context.Background(), &task,
)
assert.Zero(t, hits.Load(),
"a scheduled retry fired at a target the operator "+
"had already deleted",
)
iAssertStatus(
t, s.WebhookDB, task.DeliveryID,
database.DeliveryStatusFailed,
)
last := tLastResult(t, s, task.DeliveryID, 2)
assert.False(t, last.Success)
assert.Contains(t, last.Error, "was deleted")
assert.Zero(t, s.Engine.ExportInflightHeld())
}
// TestProcessRetryTask_TargetPresent_StillDelivers is the guard's
// mutation check: a liveness check that refused every retry would pass
// the test above and break every retry there is.
func TestProcessRetryTask_TargetPresent_StillDelivers(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
var hits atomic.Int64
task, _ := tRetryChainSetup(t, s, "still-there", &hits)
s.Engine.ExportProcessRetryTask(
context.Background(), &task,
)
assert.Equal(t, int64(1), hits.Load())
iAssertStatus(
t, s.WebhookDB, task.DeliveryID,
database.DeliveryStatusDelivered,
)
}
// TestProcessRetryTask_TargetUnreadable_StillDelivers pins the other
// half of the guard: only a target that is confirmed gone stops a
// retry. A main database that cannot be read is a transient fault, and
// a guard that abandoned deliveries on one would be a worse bug than
// the one it fixes.
func TestProcessRetryTask_TargetUnreadable_StillDelivers(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
var hits atomic.Int64
task, _ := tRetryChainSetup(t, s, "unreadable-main", &hits)
sqlDB, err := s.MainDB.DB()
require.NoError(t, err)
require.NoError(t, sqlDB.Close())
s.Engine.ExportProcessRetryTask(
context.Background(), &task,
)
assert.Equal(t, int64(1), hits.Load(),
"a retry was abandoned because the main database "+
"could not be read, not because its target was gone",
)
iAssertStatus(
t, s.WebhookDB, task.DeliveryID,
database.DeliveryStatusDelivered,
)
}
// TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone is the
// same rule on the recovery path. A read failure that is not
// "record not found" must leave every retrying delivery of every
// webhook exactly as it was.
func TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "unreadable-on-recovery",
)
targetID := uuid.New().String()
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, "healthy",
database.TargetTypeHTTP,
iHTTPConfig("http://example.com/hook"), 5,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"still":"retrying"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
iSeedFailedResult(t, s.WebhookDB, d.ID)
sqlDB, err := s.MainDB.DB()
require.NoError(t, err)
require.NoError(t, sqlDB.Close())
s.Engine.ExportRecoverRetryingDeliveries(
s.WebhookDB, s.WebhookID,
)
iAssertStatus(
t, s.WebhookDB, d.ID,
database.DeliveryStatusRetrying,
)
assert.Len(t, iResults(t, s.WebhookDB, d.ID), 1,
"an unreadable main database produced a terminal "+
"failure row",
)
assert.Zero(t, s.Engine.ExportInflightHeld())
}

View File

@@ -70,7 +70,7 @@ func (s *Server) serveUntilShutdown() {
err := s.httpServer.ListenAndServe()
if err != nil && !errors.Is(err, http.ErrServerClosed) {
s.log.Error("listen error", "error", err)
s.shutdownOnListenFailure()
s.shutdownWithFailure()
}
}

View File

@@ -93,7 +93,7 @@ func requireListenFailureExit(t *testing.T, env *testEnv) {
select {
case sig := <-app.Wait():
require.Equal(
t, server.ListenFailureExitCode, sig.ExitCode,
t, server.StartupFailureExitCode, sig.ExitCode,
"listen failure must exit non-zero",
)
case <-time.After(listenFailureDeadline):

View File

@@ -0,0 +1,99 @@
package server_test
import (
"context"
"net"
"strconv"
"testing"
"time"
"github.com/stretchr/testify/require"
"go.uber.org/fx"
"sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/globals"
"sneak.berlin/go/webhooker/internal/server"
)
// TestSentryInitFailure_ShutsDownTheApp pins that error reporting
// which is configured and cannot be started ends the application
// instead of serving without it.
//
// The measured defect logged `sentry init failure` and kept running,
// so the deployment served traffic with reporting off while every
// other signal — SENTRY_DSN still set, the startup summary's own
// field — said it was on. Nothing later in the process can notice
// that reports are going nowhere, which is why this exits rather than
// degrades.
//
// The DSN is placed on a hand-built Config, which is the only way to
// reach this branch at all: loadFromEnv now parses SENTRY_DSN with
// sentry.NewDsn, the same call sentry.Init makes, so a DSN that
// survives configuration cannot fail initialisation in the SDK
// version this pins. The branch stays because that is a property of
// the SDK's current implementation rather than of its contract.
func TestSentryInitFailure_ShutsDownTheApp(t *testing.T) {
t.Parallel()
port := freePort(t)
env := newTestEnvWithConfig(t, &config.Config{
DataDir: t.TempDir(),
Environment: config.EnvironmentDev,
BindAddress: loopbackV4,
Port: port,
SentryDSN: "not-a-dsn",
})
app := fx.New(
fx.NopLogger,
fx.Supply(env.log, env.cfg, env.mw, env.hnd),
fx.Provide(globals.New, server.New),
fx.Invoke(func(*server.Server) {}),
)
startCtx, cancelStart := context.WithTimeout(
context.Background(), lifecycleTimeout,
)
defer cancelStart()
require.NoError(t, app.Start(startCtx))
select {
case sig := <-app.Wait():
require.Equal(
t, server.StartupFailureExitCode, sig.ExitCode,
"a sentry failure must exit non-zero",
)
case <-time.After(listenFailureDeadline):
t.Fatal("a sentry failure left the app running")
}
// The stop sequence still has to complete: the failure must reach
// shutdown through fx rather than around it.
stopCtx, cancelStop := context.WithTimeout(
context.Background(), lifecycleTimeout,
)
defer cancelStop()
require.NoError(t, app.Stop(stopCtx))
// And it must give up before it listens. A process that bound the
// port and then exited would have accepted requests it could not
// report on, which is the state under test in miniature.
requireBindable(t, port)
}
// requireBindable asserts that the port is free, which it is only if
// the server under test never claimed it.
func requireBindable(t *testing.T, port int) {
t.Helper()
var listenCfg net.ListenConfig
listener, err := listenCfg.Listen(
t.Context(), "tcp",
net.JoinHostPort(loopbackV4, strconv.Itoa(port)),
)
require.NoError(t, err, "the server bound a port it then gave up")
require.NoError(t, listener.Close())
}

View File

@@ -51,12 +51,13 @@ const (
minSentryFlush = 250 * time.Millisecond
)
// ListenFailureExitCode is the status the process exits with when the
// HTTP listener cannot be established, or dies for a reason other
// than a requested shutdown. It must stay non-zero: systemd
// `Restart=on-failure` and Docker's restart policies key off it, and a
// zero exit would read as a deliberate stop.
const ListenFailureExitCode = 1
// StartupFailureExitCode is the status the process exits with when
// the serving goroutine gives up: the HTTP listener cannot be
// established or dies for a reason other than a requested shutdown, or
// error reporting is configured and cannot be started. It must stay
// non-zero: systemd `Restart=on-failure` and Docker's restart policies
// key off it, and a zero exit would read as a deliberate stop.
const StartupFailureExitCode = 1
// SentryFlushBudget reports how long the Sentry flush may run when
// remaining is the time left on the fx stop context after the HTTP
@@ -135,11 +136,25 @@ func New(lc fx.Lifecycle, params ServerParams) (*Server, error) {
}
// Run configures Sentry and starts serving HTTP requests.
//
// A Sentry failure ends the application instead of listening. It runs
// before the listener rather than after it so that the process never
// binds a port it is about to give up.
func (s *Server) Run() {
s.configure()
// logging before sentry, because sentry logs
s.enableSentry()
err := s.enableSentry()
if err != nil {
s.log.Error(
"SENTRY_DSN is set but error reporting could not be "+
"started; refusing to serve with it off",
"error", err,
)
s.shutdownWithFailure()
return
}
s.serve()
}
@@ -150,11 +165,23 @@ func (s *Server) MaintenanceMode() bool {
return s.params.Config.MaintenanceMode
}
func (s *Server) enableSentry() {
// enableSentry initialises the Sentry SDK when error reporting is
// configured, and reports the failure when it is configured and cannot
// be initialised. A DSN that is not set is not a failure: reporting
// stays off and the server starts normally.
//
// There is no fallback to running with reporting off. An operator who
// set SENTRY_DSN asked for failures to be visible, and serving traffic
// with reporting quietly off is the one state nothing can ever tell
// them about — the DSN is still set, so every later signal says it is
// on. Config already refused a DSN the SDK cannot parse, which is what
// a typo produces, so reaching this branch means the SDK refused
// something that parsed: not a condition to guess at either.
func (s *Server) enableSentry() error {
s.sentryEnabled.Store(false)
if s.params.Config.SentryDSN == "" {
return
if !s.params.Config.SentryEnabled() {
return nil
}
err := sentry.Init(sentryClientOptions(
@@ -166,19 +193,19 @@ func (s *Server) enableSentry() {
),
))
if err != nil {
s.log.Error("sentry init failure", "error", err)
// Don't use fatal since we still want the service to run
return
return fmt.Errorf("initialising sentry: %w", err)
}
s.log.Info("sentry error reporting activated")
s.sentryEnabled.Store(true)
return nil
}
// serve installs the signal watcher, starts the listener and blocks
// until the server's context is cancelled. The process exit status is
// fx's to decide — from a signal, or from the code
// shutdownOnListenFailure hands the Shutdowner — so this reports
// shutdownWithFailure hands the Shutdowner — so this reports
// nothing back to its caller.
func (s *Server) serve() {
ctx, cancelFunc := context.WithCancel(context.Background())
@@ -208,20 +235,24 @@ func (s *Server) serve() {
// Do not call cleanShutdown() here to avoid double invocation.
}
// shutdownOnListenFailure ends the application after the HTTP
// listener failed. The fx OnStart hook returns as soon as the serving
// goroutine is spawned, so nothing downstream of it ever learns that
// the listen failed: fx reports RUNNING and the process sits alive
// with nothing bound, which is invisible to systemd and Docker
// restart policies. Asking the Shutdowner to stop the app with a
// non-zero code is what turns that into a visible failure.
// shutdownWithFailure ends the application non-zero from the serving
// goroutine. It is how anything on that goroutine fails fatally: the
// fx OnStart hook returns as soon as the goroutine is spawned, so
// nothing downstream of it ever learns that the goroutine gave up. fx
// reports RUNNING and the process sits alive having done neither what
// it was asked nor anything visible instead, which systemd and
// Docker restart policies cannot see. Asking the Shutdowner to stop
// the app with a non-zero code is what turns that into a visible
// failure, and it is the whole of "fatal" here — no panic, no
// os.Exit, and every stop hook still runs.
//
// The context cancel that follows only unwinds serve()'s own wait.
// The shutdown itself runs through fx's normal stop sequence, so the
// clean-shutdown drain in cleanShutdown is reached unchanged.
func (s *Server) shutdownOnListenFailure() {
// The context cancel that follows only unwinds serve()'s own wait,
// and is skipped before serve has installed one. The shutdown itself
// runs through fx's normal stop sequence, so the clean-shutdown drain
// in cleanShutdown is reached unchanged.
func (s *Server) shutdownWithFailure() {
err := s.params.Shutdowner.Shutdown(
fx.ExitCode(ListenFailureExitCode),
fx.ExitCode(StartupFailureExitCode),
)
if err != nil {
s.log.Error("shutdown request failed", "error", err)