3 Commits

Author SHA1 Message Date
clawbot
43f72e0fd8 Make SQLite durable under concurrent readers and stop re-delivering stranded webhooks (closes #256)
All checks were successful
check / check (push) Successful in 3m33s
An operator running `sqlite3 <db> .dump` against their own per-webhook
database wedged it: inbound webhooks rejected with HTTP 500, delivered
webhooks stranded at `pending`, and every one of them POSTed a second
time on the next restart while the event log recorded a single attempt.

Durability. Every SQLite file — main, per-webhook, and archive — now
opens through one path, `internal/database/sqlite_open.go`, in WAL
journal mode with a 10-second busy timeout, `BEGIN IMMEDIATE`
transactions, and a bounded connection pool. WAL is what stops a reader
blocking writers at all. `_txlock=immediate` is what stops a `COMMIT`
failing while its transaction stays open on a pooled connection, which
is how four `database is locked` errors became 593 `cannot start a
transaction within a transaction`. `cache=shared` is gone, because
under it an in-process conflict is SQLITE_LOCKED, which the busy
handler does not retry. The busy timeout is applied before
journal_mode: the driver runs DSN pragmas in order on every new
connection, and `PRAGMA journal_mode` takes a lock, so the reverse
order leaves the one pragma that can block uncovered by the handler
meant to cover it.

Eligibility. `internal/delivery/inflight.go` holds the set of
deliveries the engine owns — taken when a task is queued, when a
target schedules a retry, and by every recovery path before it
re-dispatches; dropped when the worker that ran the task returns.
Recovery and both sweep arms re-dispatch only what the set does not
hold. Nothing decides that from a row's age: a delivery waiting in a
10000-deep channel is arbitrarily old and perfectly healthy, and
reasoning from age re-sends it. `takeForRedispatch` is the single gate
every re-dispatch goes through — ownership first, then a conditional
update confirming the row is still in the status the batch read.

Bookkeeping. `recordResult` and `updateDeliveryStatus` return their
errors instead of logging and dropping them, and a caller whose
bookkeeping write failed writes nothing at all: the delivery keeps
whichever non-terminal status it already held, and the sweeps recover
it. Every recovery path — pending and retrying alike — first settles
any delivery that already holds a successful `DeliveryResult` rather
than sending it again. Recovery continues each delivery's own attempt
numbering instead of restarting at 1. The sweep gains a
`pending`-with-age-bound arm, so a stranded delivery no longer waits
for a restart.

Docs. WAL produces `-wal`/`-shm` sidecars, so the backup and restore
procedures in README.md are corrected against measurement: both
documented procedures were re-run against a live instance, a `-wal`
left by a crash carries data the `.db` alone does not, and an archive
file normally holds its rows in a `-wal` rather than in the `.db`.
2026-08-24 00:55:39 +00:00
5fda446c71 Name a deleted target on its historical deliveries (closes #211)
All checks were successful
check / check (push) Successful in 3m13s
2026-08-24 02:03:23 +02:00
763d8f8058 Bound the /metrics method label (closes #261)
Some checks failed
check / check (push) Superseded by a newer commit; never tested
2026-08-24 02:03:18 +02:00
29 changed files with 2792 additions and 189 deletions

139
README.md
View File

@@ -566,9 +566,25 @@ is both the simplest and the only complete rule:
`events-3f2a1c9e-....db`. The only other file is `webhooker.lock`, the `events-3f2a1c9e-....db`. The only other file is `webhooker.lock`, the
always-empty [single-instance lock](#single-instance-lock); it holds no always-empty [single-instance lock](#single-instance-lock); it holds no
state and is not part of the backup set — a copied one is stale and state and is not part of the backup set — a copied one is stale and
blocks nothing. No `-wal` or `-shm` files are produced (see below); a blocks nothing.
transient `{name}.db-journal` may exist beside a database while a write
is in flight and is not part of the backup set either. **`-wal` and `-shm` sidecars.** Every database runs in WAL journal mode,
so while the service is running each `{name}.db` has a `{name}.db-wal`
and a `{name}.db-shm` beside it. **`-wal` is part of the database, not a
scratch file**: it holds committed transactions that are not yet in the
`.db`, so a copy of the `.db` without its `-wal` is missing data and may
have no readable schema at all. `-shm` is regenerable, but there is no
reason to separate the two — copy the directory and you have them.
A clean shutdown closes `webhooker.db` and every `events-*.db`, which
checkpoints and removes their sidecars; a killed or crashed instance
leaves them, and they must be carried with the `.db`. **Archive
databases are different**: their handle is not closed at shutdown, so
`archive-*.db-wal` and `-shm` normally survive a clean stop and the
`-wal` can hold every row the archive has. Measured on a stopped
instance: `archive-….db` 4096 bytes with no table, its `-wal` 157 KB
holding all 8 archived events. Copying `DATA_DIR` in full is what makes
this a non-issue; copying `.db` files out of it by name is not.
Configuration is **not** in `DATA_DIR` — it comes from the environment Configuration is **not** in `DATA_DIR` — it comes from the environment
and from a `.env` file read out of the process working directory. Back and from a `.env` file read out of the process working directory. Back
@@ -576,17 +592,19 @@ that up with your deployment config, separately.
### A hot copy is not safe ### A hot copy is not safe
No `journal_mode` pragma is ever issued on any database webhooker opens, Every database webhooker opens runs in WAL journal mode. The main and
so all of them run on SQLite's default rollback journal. There is no event databases are also held open for the entire process lifetime —
WAL. The main and event databases are also held open for the entire `WebhookDBManager` caches event database handles and closes them only on
process lifetime — `WebhookDBManager` caches event database handles and webhook deletion or shutdown — so "it looked idle" is not a guarantee
closes them only on webhook deletion or shutdown — so "it looked idle" that nothing was mid-transaction.
is not a guarantee that nothing was mid-transaction.
That means `cp`, `rsync`, `tar` or a filesystem snapshot taken against a That means `cp`, `rsync`, `tar` or a filesystem snapshot taken against a
running instance can capture a database mid-transaction and yield a file running instance can capture a database and its `-wal` at two different
that is corrupt or missing state the journal would have rolled back. Use instants and yield a file that is corrupt or missing state. Copying a
one of the two procedures below instead. `.db` on its own is worse and fails loudly: recently written pages,
including the schema itself on a young database, live in the `-wal`, so
the copy reads back as an empty or table-less database. Use one of the
two procedures below instead.
**Stop, copy, start.** The simplest, needs no extra tooling, and the **Stop, copy, start.** The simplest, needs no extra tooling, and the
only one that gives a single point in time across every file: only one that gives a single point in time across every file:
@@ -605,8 +623,9 @@ for db in /path/to/data/*.db; do
done done
``` ```
`.backup` takes the proper locks and writes a consistent file. Two `.backup` reads through the WAL and writes a single consistent file with
caveats. First, the runtime image is `alpine:3.21` with only no sidecars of its own, so the destination is complete as it stands.
Two caveats. First, the runtime image is `alpine:3.21` with only
`ca-certificates` added — the `sqlite3` CLI is **not** in it, so run `ca-certificates` added — the `sqlite3` CLI is **not** in it, so run
this on the host against the volume path, or from a throwaway container this on the host against the volume path, or from a throwaway container
that mounts the volume. Second, each file is captured at its own that mounts the volume. Second, each file is captured at its own
@@ -614,15 +633,37 @@ instant, so a webhook created or an event delivered between two files
being copied lands in one and not the other. If you need the whole set being copied lands in one and not the other. If you need the whole set
coherent as of a single moment, stop the service. coherent as of a single moment, stop the service.
Note that `sqlite3 <db> .dump` is **not** one of these procedures: it is
an export, it holds a read transaction open for as long as it runs, and
it pins the WAL against checkpointing for that whole time. It is safe to
run — it does not block ingestion — but back up with `.backup` or a
stopped copy.
Archive databases are the one exception the service is built for: the Archive databases are the one exception the service is built for: the
archive writer closes its handle after each write (debounced to at most archive writer closes and reopens its handle around writes (debounced
one reopen per second), so an operator can move `archive-{uuid}.db` to at most one reopen per second), so an operator can move
away for offline retention while the service runs, and it is recreated `archive-{uuid}.db` away for offline retention while the service runs,
on the next write (see and it is recreated on the next write. See
[Database Architecture](#database-architecture)). That is a [Database Architecture](#database-architecture). That is a
move-the-file-away workflow, not a substitute for the backup procedures move-the-file-away workflow, not a substitute for the backup procedures
above. above.
**Move the sidecars with it.** Under WAL that workflow is no longer a
single file, and the common case is the dangerous one. The reopen
happens on the *next* write after the debounce window elapses, so after
the last write of a burst nothing checkpoints: measured, 20 s after ten
events the `archive-….db` was 4096 bytes — a header, no table — with
all ten rows sitting in a 189 KB `-wal`. Copying the `.db` alone at that
moment yields a file that opens with `no such table: archived_events`.
The file becomes self-contained again when the handle closes, which
happens on the next write past the debounce window, when the connection
pool retires the idle connection (about a minute after the last write),
or at the idle archive sweep — measured, the same file was a complete
20 KB `.db` with no sidecars about a minute after its last write.
Shutdown is **not** on that list: the archive handle is not closed when
the service stops. So either move `archive-{uuid}.db` together with any
`-wal`/`-shm` beside it, or wait until there are none.
### Restore ### Restore
1. Stop the service. 1. Stop the service.
@@ -636,14 +677,22 @@ above.
restored without `webhooker.db` are simply orphaned; nothing restored without `webhooker.db` are simply orphaned; nothing
references their UUIDs. references their UUIDs.
3. Do not carry `*.db-journal` files into the restore. Backups taken by 3. Carry any `*.db-wal` and `*.db-shm` files that are in the backup.
either procedure above are self-consistent and do not need one. They are part of the database, and dropping a `-wal` silently
discards every transaction it still holds. An `.backup` set will not
contain any: it writes a single consolidated file per database. A
stop-and-copy set has none for `webhooker.db` or the `events-*.db`,
because a clean stop closes those and checkpoints their sidecars
away — but it will normally have them for `archive-*.db`, whose
handle stays open across shutdown, and those carry the archive's
rows. A copy salvaged from a crashed instance has them for
everything, and needs all of them.
4. **Fix ownership.** The container runs as the non-root `webhooker` 4. **Fix ownership.** The container runs as the non-root `webhooker`
user, UID 1000 / GID 1000. Restored files must be owned by (or user, UID 1000 / GID 1000. Restored files must be owned by (or
writable by) that UID, and so must the directory itself — SQLite writable by) that UID, and so must the directory itself — SQLite
creates the rollback journal beside the database, so a writable file creates the `-wal` and `-shm` sidecars beside the database, so a
inside a directory it cannot write is not enough: writable file inside a directory it cannot write is not enough:
```bash ```bash
chown -R 1000:1000 /path/to/data chown -R 1000:1000 /path/to/data
@@ -1391,10 +1440,12 @@ This separation provides:
only, or disables cleanup entirely when set to `0` (retain forever). only, or disables cleanup entirely when set to `0` (retain forever).
- **Performance** — each webhook's database has its own page cache and - **Performance** — each webhook's database has its own page cache and
its own lock, so concurrent event ingestion across webhooks won't its own lock, so concurrent event ingestion across webhooks won't
contend. No write-ahead log is involved: both DSNs are contend. Every database — main, per-webhook, and archive — is opened
`file:{path}?cache=shared&mode=rwc` and no `journal_mode` pragma is through one code path (`internal/database/sqlite_open.go`) in WAL
ever issued, so every database runs on SQLite's default rollback journal mode, with a 10-second busy timeout, `BEGIN IMMEDIATE`
journal. transactions, and a bounded connection pool. Under WAL a reader never
blocks a writer, so an operator reading a database does not stall
event ingestion into it.
The **database target type** builds on this architecture to provide The **database target type** builds on this architecture to provide
long-term archiving, separate from the per-webhook event database (which long-term archiving, separate from the per-webhook event database (which
@@ -1677,6 +1728,40 @@ gauge. The outcome counters move only after the status change has been
written, so a transition the database rejected is never reported as an written, so a transition the database rejected is never reported as an
outcome that happened. outcome that happened.
#### Inbound HTTP metrics
The middleware records three more on the same registry:
| Metric | Type | Labels |
| ------ | ---- | ------ |
| `http_request_duration_seconds` | histogram | `service`, `handler`, `method`, `code` |
| `http_response_size_bytes` | histogram | `service`, `handler`, `method`, `code` |
| `http_requests_inflight` | gauge | `service`, `handler` |
Two of those labels are written once per request from bytes the client
chose, so both are bounded to something this service registers:
- `handler` is the chi route pattern — `/webhook/{uuid}`, never the
concrete path. A request matching no route carries `(unmatched)`,
and no entrypoint UUID ever reaches a label.
- `method` is the request method when the router can route it, and
`(unmatched)` otherwise. `net/http` accepts any RFC 9110 token as a
method, so the raw value bounds the label at nothing; the nine chi
matches routes for stay distinguishable, and a token that could only
ever have produced a 405 does not get a series of its own.
The other two are not request-controlled: `code` is the status one of
this service's own handlers wrote, and `service` is a fixed empty
string.
`http_requests_inflight` is deliberately aggregate — its `handler` is
always `(all)`, one series counting the requests in flight across the
whole service. The gauge is incremented before routing and decremented
after the handler returns, and the route pattern exists only between
those two moments, so labelling it by pattern would increment one
series and decrement another, leaving every pattern permanently off by
the number of requests it served.
### Rate Limiting ### Rate Limiting
Global blanket rate limiting middleware (e.g., a per-IP throttle shared Global blanket rate limiting middleware (e.g., a per-IP throttle shared

View File

@@ -4,7 +4,6 @@ package database
import ( import (
"context" "context"
"crypto/rand" "crypto/rand"
"database/sql"
"encoding/base64" "encoding/base64"
"errors" "errors"
"fmt" "fmt"
@@ -16,7 +15,6 @@ import (
"go.uber.org/fx" "go.uber.org/fx"
"gorm.io/driver/sqlite" "gorm.io/driver/sqlite"
"gorm.io/gorm" "gorm.io/gorm"
_ "modernc.org/sqlite" // Pure Go SQLite driver
"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/gormlog" "sneak.berlin/go/webhooker/internal/gormlog"
@@ -198,13 +196,11 @@ func (d *Database) connectTo(dataDir string) error {
// Construct the main application database path inside DATA_DIR. // Construct the main application database path inside DATA_DIR.
dbPath := filepath.Join(dataDir, MainDBFileName) dbPath := filepath.Join(dataDir, MainDBFileName)
dbURL := fmt.Sprintf(
"file:%s?cache=shared&mode=rwc",
dbPath,
)
// Open the database with the pure Go SQLite driver // Opened through OpenSQLite so this handle carries the same WAL
sqlDB, err := sql.Open("sqlite", dbURL) // journaling, busy timeout, immediate-transaction locking, and pool
// bounds as every other database file. See sqlite_open.go.
sqlDB, err := OpenSQLite(dbPath, SQLiteModeCreate)
if err != nil { if err != nil {
d.log.Error( d.log.Error(
"failed to open database", "failed to open database",

View File

@@ -0,0 +1,157 @@
package database
import (
"database/sql"
"fmt"
"net/url"
"time"
_ "modernc.org/sqlite" // Pure Go SQLite driver
)
// Every SQLite file this service opens — the main database, the
// per-webhook event databases, and the archive databases — is opened
// through OpenSQLite, so the durability settings below are properties
// of the service rather than of one call site.
//
// modernc.org/sqlite installs no busy handler and issues no pragmas of
// its own: it executes only the pragmas named in explicit `_pragma=`
// DSN parameters, and gorm.io/driver/sqlite adds none when it is
// handed an existing *sql.DB. Every setting therefore has to be
// spelled out here or it is simply not in effect.
// SQLite URI open modes.
const (
// SQLiteModeCreate creates the database file when it is missing.
SQLiteModeCreate = "rwc"
// SQLiteModeExisting requires the file to exist already.
SQLiteModeExisting = "rw"
)
const (
// SQLiteBusyTimeout is how long SQLite retries a lock conflict
// before returning SQLITE_BUSY.
//
// Under WAL a reader never blocks a writer, so the only conflict
// left is writer against writer: this process's delivery workers
// against each other, or against another process holding the write
// lock. Those clear in milliseconds. Ten seconds is far above that
// and still well inside the receiver's request budget, so an
// inbound webhook waits rather than being rejected with a 500.
SQLiteBusyTimeout = 10 * time.Second
// sqliteMaxOpenConns bounds the connection pool for one database
// file.
//
// The pool needs a bound at all because database/sql cannot detect
// a connection left mid-transaction: modernc.org/sqlite implements
// neither driver.Validator nor driver.SessionResetter, so a
// connection whose COMMIT failed is returned to the pool with its
// transaction still open and handed out again indefinitely. That is
// what turned four `database is locked` errors into 593
// `cannot start a transaction within a transaction` in
// https://git.eeqj.de/sneak/webhooker/issues/256.
//
// Four is above the one writer SQLite allows at a time, so reads
// still proceed while a write is in flight, and low enough that
// contention is resolved by the busy handler rather than by piling
// up connections against a lock only one of them can hold.
sqliteMaxOpenConns = 4
// sqliteMaxIdleConns keeps the pool warm without holding every
// connection open through an idle period.
sqliteMaxIdleConns = 2
// sqliteConnMaxLifetime and sqliteConnMaxIdleTime retire pooled
// connections on a schedule. With _txlock=immediate a failed
// COMMIT should no longer be reachable, but these bound the damage
// if one happens anyway: a poisoned connection is closed and
// replaced within the lifetime instead of wedging the file until
// the process restarts.
sqliteConnMaxLifetime = 5 * time.Minute
sqliteConnMaxIdleTime = time.Minute
)
// SQLiteDSN builds the connection string for one database file.
//
// mode is the SQLite URI open mode: "rwc" to create the file when it
// is missing, "rw" to require that it already exists.
//
// Three settings carry the fix for
// https://git.eeqj.de/sneak/webhooker/issues/256 and none of them is
// optional:
//
// - journal_mode=WAL, so a reader — an operator running
// `sqlite3 <db> .dump` over their own data — takes a snapshot
// instead of blocking every writer behind it.
//
// - busy_timeout, so a writer that does meet a lock waits for it.
// Without one SQLite gives up immediately; nothing above it
// retries.
//
// - _txlock=immediate, so every transaction takes the write lock at
// BEGIN. A deferred transaction acquires it lazily on its first
// write, and that upgrade returns SQLITE_BUSY *without* consulting
// the busy handler, because SQLite cannot block a transaction that
// may already hold a read snapshot. Such a COMMIT then fails while
// the transaction stays open on the connection. A busy timeout
// alone does not prevent this; BEGIN IMMEDIATE does, by putting
// the wait somewhere the handler applies.
//
// Note what is absent: `cache=shared`. Under a shared cache an
// in-process conflict is reported as SQLITE_LOCKED rather than
// SQLITE_BUSY, and the busy handler does not retry SQLITE_LOCKED — so
// leaving it in would have defeated the busy timeout for exactly the
// contention this service generates. Dropping it is part of the fix,
// not housekeeping.
//
// synchronous is deliberately left at SQLite's default of FULL: this
// is a webhook receiver whose one promise is that an event it answered
// 200 for is durable.
// The order of the _pragma parameters is load-bearing.
// modernc.org/sqlite executes them in the order they appear, on every
// new connection, before the connection is handed to the pool. Setting
// journal_mode first means that pragma itself runs with no busy
// handler installed: the pool opens connections lazily, so the moment
// a new one is created is a moment the database is under load, and
// PRAGMA journal_mode takes a lock. It would fail immediately with
// SQLITE_BUSY and fail the query that caused the connection to be
// opened. busy_timeout is therefore set first, so every pragma after
// it — and the whole life of the connection — is covered.
func SQLiteDSN(path, mode string) string {
q := url.Values{}
q.Set("mode", mode)
q.Set("_txlock", "immediate")
q.Add(
"_pragma",
fmt.Sprintf(
"busy_timeout(%d)",
SQLiteBusyTimeout.Milliseconds(),
),
)
q.Add("_pragma", "journal_mode(WAL)")
return "file:" + path + "?" + q.Encode()
}
// OpenSQLite opens the SQLite file at path with the service's
// durability settings and pool bounds applied. mode is the SQLite URI
// open mode ("rwc" or "rw").
//
// The handle is returned rather than a *gorm.DB because the callers
// wrap it in gorm themselves with their own logger.
func OpenSQLite(path, mode string) (*sql.DB, error) {
sqlDB, err := sql.Open("sqlite", SQLiteDSN(path, mode))
if err != nil {
return nil, fmt.Errorf(
"opening sqlite database %s: %w", path, err,
)
}
sqlDB.SetMaxOpenConns(sqliteMaxOpenConns)
sqlDB.SetMaxIdleConns(sqliteMaxIdleConns)
sqlDB.SetConnMaxLifetime(sqliteConnMaxLifetime)
sqlDB.SetConnMaxIdleTime(sqliteConnMaxIdleTime)
return sqlDB, nil
}

View File

@@ -0,0 +1,178 @@
package database_test
import (
"context"
"path/filepath"
"strings"
"testing"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
)
// livePragma reads a pragma off a live handle. Reading the DSN back
// would prove only that the string was built; these tests assert that
// SQLite actually applied it.
func livePragma(t *testing.T, db *gorm.DB, name string) string {
t.Helper()
var v string
row := db.Raw("pragma " + name).Row()
require.NoError(t, row.Scan(&v))
return v
}
func TestSQLiteDSNCarriesTheDurabilitySettings(t *testing.T) {
t.Parallel()
dsn := database.SQLiteDSN(
"/var/lib/webhooker/webhooker.db",
database.SQLiteModeCreate,
)
assert.Contains(t, dsn, "journal_mode%28WAL%29")
assert.Contains(t, dsn, "busy_timeout%2810000%29")
assert.Contains(t, dsn, "_txlock=immediate")
assert.Contains(t, dsn, "mode=rwc")
// busy_timeout must come first. The driver runs these in order on
// every new connection, and PRAGMA journal_mode takes a lock — a
// connection opened while the database is busy would fail on that
// pragma, with no busy handler yet installed to wait it out.
assert.Less(
t,
strings.Index(dsn, "busy_timeout"),
strings.Index(dsn, "journal_mode"),
"busy_timeout must be applied before journal_mode",
)
// cache=shared turns an in-process conflict into SQLITE_LOCKED,
// which the busy handler does not retry. It must never come back.
// See https://git.eeqj.de/sneak/webhooker/issues/256.
assert.NotContains(t, strings.ToLower(dsn), "cache=shared")
}
// TestPerWebhookDBAppliesPragmasOnALiveHandle is the check the issue
// asks for by name: the settings are confirmed by querying the running
// database, not by inspecting the connection string.
func TestPerWebhookDBAppliesPragmasOnALiveHandle(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)) }()
webhookID := uuid.New().String()
db, err := mgr.GetDB(webhookID)
require.NoError(t, err)
assert.Equal(
t, "wal",
strings.ToLower(livePragma(t, db, "journal_mode")),
)
assert.Equal(
t, "10000", livePragma(t, db, "busy_timeout"),
)
}
func TestMainDBAppliesPragmasOnALiveHandle(t *testing.T) {
t.Parallel()
ctx := context.Background()
dir := t.TempDir()
sqlDB, err := database.OpenSQLite(
filepath.Join(dir, database.MainDBFileName),
database.SQLiteModeCreate,
)
require.NoError(t, err)
defer func() { require.NoError(t, sqlDB.Close()) }()
var journal string
require.NoError(t, sqlDB.
QueryRowContext(ctx, "pragma journal_mode").
Scan(&journal))
assert.Equal(t, "wal", strings.ToLower(journal))
var busy string
require.NoError(t, sqlDB.
QueryRowContext(ctx, "pragma busy_timeout").
Scan(&busy))
assert.Equal(t, "10000", busy)
}
// TestConcurrentReaderDoesNotBlockWrites is the unit-scale form of the
// reproduction in
// https://git.eeqj.de/sneak/webhooker/issues/256: an operator's
// long-held read of their own data used to make every concurrent write
// fail. Under WAL the reader takes a snapshot and the writes proceed.
func TestConcurrentReaderDoesNotBlockWrites(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)) }()
webhookID := uuid.New().String()
db, err := mgr.GetDB(webhookID)
require.NoError(t, err)
// A second handle on the same file, holding a read transaction
// open across every write below — what `sqlite3 <db> .dump` is.
readerSQL, err := database.OpenSQLite(
mgr.DBPath(webhookID), database.SQLiteModeExisting,
)
require.NoError(t, err)
defer func() { require.NoError(t, readerSQL.Close()) }()
readerConn, err := readerSQL.Conn(ctx)
require.NoError(t, err)
defer func() { require.NoError(t, readerConn.Close()) }()
_, err = readerConn.ExecContext(ctx, "begin deferred")
require.NoError(t, err)
_, err = readerConn.ExecContext(
ctx, "select count(*) from events",
)
require.NoError(t, err)
for range 25 {
err = db.Transaction(func(tx *gorm.DB) error {
return tx.Create(&database.Event{
WebhookID: webhookID,
EntrypointID: uuid.New().String(),
Method: "POST",
Body: "{}",
}).Error
})
require.NoError(t, err)
}
_, err = readerConn.ExecContext(ctx, "commit")
require.NoError(t, err)
var count int64
require.NoError(
t,
db.Model(&database.Event{}).Count(&count).Error,
)
assert.Equal(t, int64(25), count)
}

View File

@@ -2,7 +2,6 @@ package database
import ( import (
"context" "context"
"database/sql"
"errors" "errors"
"fmt" "fmt"
"log/slog" "log/slog"
@@ -234,12 +233,11 @@ func (m *WebhookDBManager) openDB(
webhookID string, webhookID string,
) (*gorm.DB, error) { ) (*gorm.DB, error) {
path := m.dbPath(webhookID) path := m.dbPath(webhookID)
dbURL := fmt.Sprintf(
"file:%s?cache=shared&mode=rwc",
path,
)
sqlDB, err := sql.Open("sqlite", dbURL) // See sqlite_open.go: WAL, a busy timeout, immediate-transaction
// locking, and a bounded pool, all of which this file needs most —
// it is the one every delivery worker writes to concurrently.
sqlDB, err := OpenSQLite(path, SQLiteModeCreate)
if err != nil { if err != nil {
return nil, fmt.Errorf( return nil, fmt.Errorf(
"opening webhook database %s: %w", "opening webhook database %s: %w",

View File

@@ -41,6 +41,31 @@ const (
// sweep runs. // sweep runs.
retrySweepInterval = 60 * time.Second retrySweepInterval = 60 * time.Second
// pendingSweepMinAge is how long a delivery must have sat
// untouched at pending before the sweep will look at it.
//
// It is not what keeps the sweep off live work — inflightSet is,
// and it is exact. This bound sets the re-dispatch cadence for a
// delivery that really is stranded: without it, a delivery the
// database will not let the engine settle would be re-sent on
// every 60-second tick.
//
// It is nonetheless set clear of the longest legitimate attempt,
// so that the two guards do not both have to be right. That
// length is MaxTargetTimeoutSeconds (300s), the per-target
// timeout the target form accepts — not httpClientTimeout, which
// is merely the default. Fifteen minutes leaves a margin of
// three times the ceiling rather than the zero margin the two
// equal values would have given.
pendingSweepMinAge = 15 * time.Minute
// pendingSweepBatch bounds how many stranded pending deliveries
// one sweep of one webhook re-dispatches. The sweep runs every
// retrySweepInterval, so a larger backlog drains across
// successive sweeps instead of arriving as one burst against a
// database that was already struggling to accept writes.
pendingSweepBatch = 500
// MaxInlineBodySize is the maximum event body size that // MaxInlineBodySize is the maximum event body size that
// will be carried inline in a Task through the channel. // will be carried inline in a Task through the channel.
// Bodies at or above this size are left nil and fetched // Bodies at or above this size are left nil and fetched
@@ -157,6 +182,12 @@ type Engine struct {
// dbTarget is retained so the engine can reach the archive // dbTarget is retained so the engine can reach the archive
// writer registry for webhook eviction and the idle sweep. // writer registry for webhook eviction and the idle sweep.
dbTarget *databaseTarget dbTarget *databaseTarget
// inflight is the set of deliveries this engine currently owns.
// Recovery and the sweeps re-dispatch only what it does not
// hold. Held by value: its zero value works, so no constructor
// can leave it out. See inflight.go.
inflight inflightSet
} }
// New creates and registers the delivery engine with the // New creates and registers the delivery engine with the
@@ -189,12 +220,27 @@ func New(
// are ready. // are ready.
func (e *Engine) Notify(tasks []Task) { func (e *Engine) Notify(tasks []Task) {
for i := range tasks { for i := range tasks {
// Owned before it is queued, and until the worker that runs
// it returns. A task can sit in a 10000-deep channel for a
// long time on a healthy system, and nothing may re-send it
// while it waits. See inflight.go.
if !e.inflight.retainIdle(tasks[i].DeliveryID) {
e.log.Warn(
"delivery already in flight, not queued again",
"delivery_id", tasks[i].DeliveryID,
"event_id", tasks[i].EventID,
)
continue
}
select { select {
case e.deliveryCh <- tasks[i]: case e.deliveryCh <- tasks[i]:
default: default:
e.inflight.release(tasks[i].DeliveryID)
e.log.Warn( e.log.Warn(
"delivery channel full, "+ "delivery channel full, "+
"task will be recovered on restart", "task will be recovered by the sweep",
"delivery_id", tasks[i].DeliveryID, "delivery_id", tasks[i].DeliveryID,
"event_id", tasks[i].EventID, "event_id", tasks[i].EventID,
) )
@@ -229,10 +275,20 @@ func (e *Engine) ScheduleRetry(
"next_attempt", task.AttemptNum, "next_attempt", task.AttemptNum,
) )
// The reference is taken here rather than when the timer fires,
// so the delivery stays owned across the whole backoff window.
// Its caller is a target inside Deliver, so the engine already
// owns it; this second reference is what keeps that ownership
// alive after the worker returns and the row sits at retrying
// with nothing running. Without it the sweep finds the row
// orphaned and sends it again.
e.inflight.retain(task.DeliveryID)
time.AfterFunc(delay, func() { time.AfterFunc(delay, func() {
select { select {
case e.retryCh <- task: case e.retryCh <- task:
default: default:
e.inflight.release(task.DeliveryID)
e.log.Warn( e.log.Warn(
"retry channel full, delivery "+ "retry channel full, delivery "+
"will be recovered by periodic sweep", "will be recovered by periodic sweep",
@@ -332,13 +388,35 @@ func (e *Engine) worker(ctx context.Context) {
case <-ctx.Done(): case <-ctx.Done():
return return
case task := <-e.deliveryCh: case task := <-e.deliveryCh:
e.processNewTask(ctx, &task) e.runTask(ctx, &task, e.processNewTask)
case task := <-e.retryCh: case task := <-e.retryCh:
e.processRetryTask(ctx, &task) e.runTask(ctx, &task, e.processRetryTask)
} }
} }
} }
// runTask runs one task and then drops the reference the queueing
// side took on its delivery.
//
// The release is deferred rather than written after the call because
// every early return inside the processing paths must drop it too: a
// delivery whose database could not be opened is one the engine has
// stopped working on, and leaving it owned would hide it from the
// sweep forever.
//
// Ownership does not necessarily end here. A target that scheduled a
// retry took its own reference before this one is dropped, so the
// delivery stays owned through the backoff window.
func (e *Engine) runTask(
ctx context.Context,
task *Task,
run func(context.Context, *Task),
) {
defer e.inflight.release(task.DeliveryID)
run(ctx, task)
}
func (e *Engine) recoverPending(ctx context.Context) { func (e *Engine) recoverPending(ctx context.Context) {
defer e.wg.Done() defer e.wg.Done()
@@ -526,7 +604,16 @@ func (e *Engine) recoverRetryingDeliveries(
return return
} }
settled := e.reconcileDelivered(
webhookDB, webhookID, retrying,
e.loadTargetMap(retrying),
)
for i := range retrying { for i := range retrying {
if _, ok := settled[retrying[i].ID]; ok {
continue
}
e.recoverSingleRetry( e.recoverSingleRetry(
webhookDB, webhookID, &retrying[i], webhookDB, webhookID, &retrying[i],
) )
@@ -588,6 +675,10 @@ func (e *Engine) recoverSingleRetry(
d, webhookID, &event, &target, attemptNum+1, d, webhookID, &event, &target, attemptNum+1,
) )
if !e.rescheduleRecovered(webhookDB, task, remaining) {
return
}
e.log.Info( e.log.Info(
"recovering retrying delivery", "recovering retrying delivery",
"webhook_id", webhookID, "webhook_id", webhookID,
@@ -595,8 +686,6 @@ func (e *Engine) recoverSingleRetry(
"attempt", attemptNum, "attempt", attemptNum,
"remaining_backoff", remaining, "remaining_backoff", remaining,
) )
e.ScheduleRetry(task, remaining)
} }
func (e *Engine) recoverPendingDeliveries( func (e *Engine) recoverPendingDeliveries(
@@ -606,12 +695,14 @@ func (e *Engine) recoverPendingDeliveries(
) { ) {
var deliveries []database.Delivery var deliveries []database.Delivery
// No Preload: event bodies are read one at a time in
// sendRecoveredDeliveries, and only for the deliveries actually
// being sent.
result := webhookDB. result := webhookDB.
Where( Where(
"status = ?", "status = ?",
database.DeliveryStatusPending, database.DeliveryStatusPending,
). ).
Preload("Event").
Find(&deliveries) Find(&deliveries)
if result.Error != nil { if result.Error != nil {
@@ -634,11 +725,137 @@ func (e *Engine) recoverPendingDeliveries(
"count", len(deliveries), "count", len(deliveries),
) )
e.recoverPendingBatch(
ctx, webhookDB, webhookID, deliveries,
)
}
// recoverPendingBatch settles every delivery in the batch that was
// already delivered, and re-dispatches only the rest. Both the
// restart-time recovery and the periodic sweep go through it, so a
// pending delivery is treated the same however it was found.
func (e *Engine) recoverPendingBatch(
ctx context.Context,
webhookDB *gorm.DB,
webhookID string,
deliveries []database.Delivery,
) {
targetMap := e.loadTargetMap(deliveries) targetMap := e.loadTargetMap(deliveries)
e.sendRecoveredDeliveries( settled := e.reconcileDelivered(
ctx, deliveries, webhookID, targetMap, webhookDB, webhookID, deliveries, targetMap,
) )
e.sendRecoveredDeliveries(
ctx, webhookDB, deliveries, webhookID,
targetMap, settled,
)
}
// reconcileDelivered finds the deliveries in a recovered batch that
// already have a successful DeliveryResult, marks them delivered, and
// returns their ids so the caller does not send them a second time.
//
// This is the state the engine previously had no way to represent. A
// delivery is left in a non-terminal state by a failed bookkeeping
// write, and that covers two different histories: nothing was ever
// sent, or the send reached the receiver and only the status write
// failed. Re-sending was the sole option, so every stranded row
// produced a duplicate at the receiver and an event log that recorded
// one attempt for two POSTs. A successful result row distinguishes
// them: it is written before the status, so its presence means the
// wire I/O happened and was recorded, and all that is missing is the
// status.
//
// Every recovery path runs this, not only the pending one. A delivery
// abandoned at retrying can hold a successful result just as a pending
// one can — a second attempt that reached the receiver and whose status
// write then failed sits at retrying with success recorded — and
// re-sending it is the same duplicate.
//
// Deliveries whose result row itself never landed are not in the
// returned set and are re-sent, recorded as the further attempt they
// are. That is honest at-least-once delivery rather than a silent
// duplicate.
func (e *Engine) reconcileDelivered(
webhookDB *gorm.DB,
webhookID string,
deliveries []database.Delivery,
targetMap map[string]database.Target,
) map[string]struct{} {
settled := make(map[string]struct{})
if len(deliveries) == 0 {
return settled
}
ids := make([]string, 0, len(deliveries))
for i := range deliveries {
ids = append(ids, deliveries[i].ID)
}
var deliveredIDs []string
err := webhookDB.
Model(&database.DeliveryResult{}).
Where(
"delivery_id IN ? AND success = ?", ids, true,
).
Distinct().
Pluck("delivery_id", &deliveredIDs).Error
if err != nil {
// Every delivery stays out of the settled set, so the batch
// is re-sent exactly as it was before this check existed.
// That is the safe direction: a duplicate delivery beats
// declaring a delivery successful on a query that failed.
e.log.Error(
"failed to query successful delivery results; "+
"pending deliveries will be re-sent",
"webhook_id", webhookID,
"error", err,
)
return settled
}
for _, id := range deliveredIDs {
settled[id] = struct{}{}
}
if len(settled) == 0 {
return settled
}
e.log.Info(
"settling recovered deliveries that already succeeded",
"webhook_id", webhookID,
"count", len(settled),
)
for i := range deliveries {
if _, ok := settled[deliveries[i].ID]; !ok {
continue
}
// A delivery the engine is working on right now settles
// itself; writing over it from here would race that worker.
if !e.inflight.retainIdle(deliveries[i].ID) {
delete(settled, deliveries[i].ID)
continue
}
e.settleStatus(
webhookDB,
&deliveries[i],
targetMap[deliveries[i].TargetID].Type,
database.DeliveryStatusDelivered,
)
e.inflight.release(deliveries[i].ID)
}
return settled
} }
func (e *Engine) retrySweep(ctx context.Context) { func (e *Engine) retrySweep(ctx context.Context) {
@@ -722,6 +939,11 @@ func (e *Engine) sweepWebhookRetries(
return return
} }
settled := e.reconcileDelivered(
webhookDB, webhookID, retrying,
e.loadTargetMap(retrying),
)
for i := range retrying { for i := range retrying {
select { select {
case <-ctx.Done(): case <-ctx.Done():
@@ -729,10 +951,70 @@ func (e *Engine) sweepWebhookRetries(
default: default:
} }
if _, ok := settled[retrying[i].ID]; ok {
continue
}
e.sweepSingleRetry( e.sweepSingleRetry(
webhookDB, webhookID, &retrying[i], webhookDB, webhookID, &retrying[i],
) )
} }
e.sweepWebhookPending(ctx, webhookDB, webhookID)
}
// sweepWebhookPending recovers deliveries stranded at pending.
//
// A delivery is created pending and leaves that state only when its
// outcome is written, so a pending row the engine does not own is one
// whose bookkeeping write failed — the state that used to sit there
// until a restart, and then produce a duplicate at the receiver. The
// sweep gives it the same reconcile-then-dispatch treatment restart
// recovery gets, so it costs a minute rather than an operator
// noticing.
//
// What keeps the sweep off live work is ownership, checked per
// delivery in takeForRedispatch, not the age bound in this query.
// A delivery waiting in deliveryCh is pending and arbitrarily old —
// the channel holds 10000 tasks and 10 workers drain it — so
// reasoning from the row's age alone re-sends it. See inflight.go.
func (e *Engine) sweepWebhookPending(
ctx context.Context,
webhookDB *gorm.DB,
webhookID string,
) {
var pending []database.Delivery
err := webhookDB.
Where(
"status = ? AND updated_at < ?",
database.DeliveryStatusPending,
time.Now().Add(-pendingSweepMinAge),
).
Limit(pendingSweepBatch).
Find(&pending).Error
if err != nil {
e.log.Error(
"retry sweep: "+
"failed to query pending deliveries",
"webhook_id", webhookID,
"error", err,
)
return
}
if len(pending) == 0 {
return
}
e.log.Info(
"retry sweep: recovering stranded pending deliveries",
"webhook_id", webhookID,
"count", len(pending),
)
e.recoverPendingBatch(ctx, webhookDB, webhookID, pending)
} }
// sweepSingleRetry re-enqueues an orphaned retrying delivery // sweepSingleRetry re-enqueues an orphaned retrying delivery
@@ -789,8 +1071,13 @@ func (e *Engine) sweepSingleRetry(
d, webhookID, &event, &target, attemptNum+1, d, webhookID, &event, &target, attemptNum+1,
) )
select { if !e.redispatch(
case e.retryCh <- task: e.retryCh, webhookDB, task,
database.DeliveryStatusRetrying,
) {
return
}
e.log.Info( e.log.Info(
"retry sweep: "+ "retry sweep: "+
"recovered orphaned retrying delivery", "recovered orphaned retrying delivery",
@@ -798,8 +1085,6 @@ func (e *Engine) sweepSingleRetry(
"webhook_id", webhookID, "webhook_id", webhookID,
"attempt", attemptNum+1, "attempt", attemptNum+1,
) )
default:
}
} }
// failUnretryableRetry terminally fails an orphaned retrying // failUnretryableRetry terminally fails an orphaned retrying
@@ -822,6 +1107,16 @@ func (e *Engine) failUnretryableRetry(
d *database.Delivery, d *database.Delivery,
target *database.Target, target *database.Target,
) { ) {
// Terminal, and reached from the recovery paths, so it takes
// ownership like every other write they make: a delivery the
// engine is still attempting must not be failed underneath the
// worker running it.
if !e.inflight.retainIdle(d.ID) {
return
}
defer e.inflight.release(d.ID)
e.log.Warn( e.log.Warn(
"failing orphaned retrying delivery: target "+ "failing orphaned retrying delivery: target "+
"type no longer supports retries", "type no longer supports retries",
@@ -839,7 +1134,7 @@ func (e *Engine) failUnretryableRetry(
target.Type, target.Type,
) )
e.recordResult( err := e.recordResult(
webhookDB, webhookDB,
d, d,
e.countAttempts(webhookDB, d.ID)+1, e.countAttempts(webhookDB, d.ID)+1,
@@ -849,6 +1144,11 @@ func (e *Engine) failUnretryableRetry(
reason, reason,
0, 0,
) )
if err != nil {
e.bookkeepingFailed(d, err)
return
}
// The type is passed rather than assigned onto d: the delivery // The type is passed rather than assigned onto d: the delivery
// is loaded here without its target relation, and populating // is loaded here without its target relation, and populating
@@ -856,7 +1156,7 @@ func (e *Engine) failUnretryableRetry(
// whole target row — plaintext config, which for a slack target // whole target row — plaintext config, which for a slack target
// is the credential — into the per-webhook event database. See // is the credential — into the per-webhook event database. See
// https://git.eeqj.de/sneak/webhooker/issues/206. // https://git.eeqj.de/sneak/webhooker/issues/206.
e.updateDeliveryStatus( e.settleStatus(
webhookDB, d, target.Type, webhookDB, d, target.Type,
database.DeliveryStatusFailed, database.DeliveryStatusFailed,
) )
@@ -878,7 +1178,7 @@ func (e *Engine) processDelivery(
"type", d.Target.Type, "type", d.Target.Type,
) )
e.updateDeliveryStatus( e.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusFailed, database.DeliveryStatusFailed,
) )
@@ -910,6 +1210,14 @@ func (e *Engine) observeAttempt(
// recordResult persists a DeliveryResult row describing a // recordResult persists a DeliveryResult row describing a
// single attempt. It is a cross-target helper the targets // single attempt. It is a cross-target helper the targets
// call. // call.
//
// It returns its error rather than swallowing it. A DeliveryResult
// row is the only record that an attempt happened at all, so a
// caller that ignored a failed write would go on to mark the
// delivery delivered — leaving the event log claiming one attempt
// for a receiver that got two. Every caller must instead stop
// advancing the delivery's status and let it stay in the
// non-terminal state it already holds; see bookkeepingFailed.
func (e *Engine) recordResult( func (e *Engine) recordResult(
webhookDB *gorm.DB, webhookDB *gorm.DB,
d *database.Delivery, d *database.Delivery,
@@ -918,7 +1226,7 @@ func (e *Engine) recordResult(
statusCode int, statusCode int,
respBody, errMsg string, respBody, errMsg string,
durationMs int64, durationMs int64,
) { ) error {
result := &database.DeliveryResult{ result := &database.DeliveryResult{
DeliveryID: d.ID, DeliveryID: d.ID,
AttemptNum: attemptNum, AttemptNum: attemptNum,
@@ -931,12 +1239,44 @@ func (e *Engine) recordResult(
err := webhookDB.Create(result).Error err := webhookDB.Create(result).Error
if err != nil { if err != nil {
e.log.Error( return fmt.Errorf(
"failed to record delivery result", "recording delivery result for %s: %w", d.ID, err,
"delivery_id", d.ID,
"error", err,
) )
} }
return nil
}
// bookkeepingFailed reports that a delivery's own record of what
// happened could not be written, and deliberately writes nothing in
// response.
//
// Leaving the row alone is the whole point. A delivery is created
// pending and only ever leaves that state through
// updateDeliveryStatus, so a delivery whose bookkeeping write failed
// is still pending or retrying — the two non-terminal states, per
// DeliveryStatus.Terminal — and both are swept and recovered. Writing
// anything here would need the very database that just refused a
// write, and would be one more thing to fail; not writing cannot.
//
// The cost is honest at-least-once behaviour: a send that reached the
// receiver but whose result row did not land is attempted again, and
// recorded as the further attempt it is. What no longer happens is the
// silent duplicate — a second POST the event log denies ever
// occurred. Where the result row *did* land and only the status write
// failed, reconcileDelivered settles the row without re-sending.
func (e *Engine) bookkeepingFailed(
d *database.Delivery, err error,
) {
e.log.Error(
"delivery bookkeeping write failed; leaving delivery "+
"in a recoverable state",
"delivery_id", d.ID,
"event_id", d.EventID,
"target_id", d.TargetID,
"status", d.Status,
"error", err,
)
} }
// updateDeliveryStatus persists a new status for a delivery. // updateDeliveryStatus persists a new status for a delivery.
@@ -952,28 +1292,53 @@ func (e *Engine) recordResult(
// //
// The counter moves only after the row is written, so a transition // The counter moves only after the row is written, so a transition
// the database rejected is not claimed as an outcome that happened. // the database rejected is not claimed as an outcome that happened.
// For the same reason the error is returned rather than logged and
// dropped: a delivery whose status write failed has not reached that
// status, and its caller must not act as though it had.
func (e *Engine) updateDeliveryStatus( func (e *Engine) updateDeliveryStatus(
webhookDB *gorm.DB, webhookDB *gorm.DB,
d *database.Delivery, d *database.Delivery,
targetType database.TargetType, targetType database.TargetType,
status database.DeliveryStatus, status database.DeliveryStatus,
) { ) error {
err := webhookDB.Model(d). err := webhookDB.Model(d).
Update("status", status).Error Update("status", status).Error
if err != nil { if err != nil {
e.log.Error( return fmt.Errorf(
"failed to update delivery status", "updating delivery %s to status %s: %w",
"delivery_id", d.ID, d.ID, status, err,
"status", status,
"error", err,
) )
return
} }
// An empty type means the target row is gone — a delivery being
// settled long after its target was deleted. The row still has to
// be settled, but the counter is left alone rather than given a
// series labelled with the empty string.
if targetType != "" {
e.mtr.DeliveryStatusChanged(targetType, status) e.mtr.DeliveryStatusChanged(targetType, status)
} }
return nil
}
// settleStatus moves a delivery to its outcome status and reports a
// failed write through bookkeepingFailed, which leaves the row
// recoverable. It exists so the target call sites read as one
// statement rather than four lines of identical error handling.
func (e *Engine) settleStatus(
webhookDB *gorm.DB,
d *database.Delivery,
targetType database.TargetType,
status database.DeliveryStatus,
) {
err := e.updateDeliveryStatus(
webhookDB, d, targetType, status,
)
if err != nil {
e.bookkeepingFailed(d, err)
}
}
func truncate(s string, maxLen int) string { func truncate(s string, maxLen int) string {
if len(s) <= maxLen { if len(s) <= maxLen {
return s return s
@@ -1065,6 +1430,184 @@ func (e *Engine) countAttempts(
return int(resultCount) return int(resultCount)
} }
// takeForRedispatch decides whether a recovered delivery may be sent
// again, and takes it if so. It is the single gate every re-dispatch
// path goes through, and it asks two separate questions in order.
//
// First, does the engine already own this delivery? Ownership is
// exact and mutually exclusive, so a delivery queued, being attempted,
// or waiting out a retry backoff is refused here, and two dispatchers
// racing for the same delivery cannot both win. See inflight.go.
//
// Second, is the row still in the status that made it eligible? The
// batch was read some time ago and a worker may have settled a row
// since. The check is a conditional update rather than a read so the
// answer cannot go stale between asking and acting.
//
// Stamping updated_at is the same statement, and it is a cadence
// control rather than a claim: the pending sweep selects on that
// column, so a delivery handed out now is not selected again on the
// next tick a minute later but after pendingSweepMinAge. A delivery
// the database refuses to settle is therefore retried on that
// interval instead of every tick.
//
// A failed write is a refusal. It means the database is not accepting
// writes, which is the condition that stranded the delivery in the
// first place; an attempt that cannot be recorded is exactly the
// unlogged duplicate this is all here to prevent.
//
// The caller must release ownership if it then fails to queue the
// task.
func (e *Engine) takeForRedispatch(
webhookDB *gorm.DB,
deliveryID string,
eligible database.DeliveryStatus,
) bool {
if !e.inflight.retainIdle(deliveryID) {
return false
}
res := webhookDB.
Model(&database.Delivery{}).
Where(
"id = ? AND status = ?", deliveryID, eligible,
).
UpdateColumn("updated_at", time.Now())
if res.Error != nil {
e.log.Error(
"failed to mark delivery for re-dispatch; "+
"leaving it for a later sweep",
"delivery_id", deliveryID,
"error", res.Error,
)
e.inflight.release(deliveryID)
return false
}
if res.RowsAffected != 1 {
// Settled underneath us between the query and here.
e.inflight.release(deliveryID)
return false
}
return true
}
// queueRecovered puts an owned delivery's task on a worker channel,
// dropping the ownership the gate took if it does not fit.
func (e *Engine) queueRecovered(
ch chan<- Task, task Task,
) bool {
select {
case ch <- task:
return true
default:
e.inflight.release(task.DeliveryID)
e.log.Warn(
"worker channel full during recovery; "+
"delivery will be recovered by a later sweep",
"delivery_id", task.DeliveryID,
"webhook_id", task.WebhookID,
)
return false
}
}
// redispatch hands a recovered delivery to a worker channel through
// takeForRedispatch, and reports whether the task was queued.
func (e *Engine) redispatch(
ch chan<- Task,
webhookDB *gorm.DB,
task Task,
eligible database.DeliveryStatus,
) bool {
if !e.takeForRedispatch(
webhookDB, task.DeliveryID, eligible,
) {
return false
}
return e.queueRecovered(ch, task)
}
// rescheduleRecovered hands an orphaned retrying delivery back to the
// retry timer, through the same gate. It reports whether the delivery
// was rescheduled.
//
// The reference taken by the gate is dropped as soon as ScheduleRetry
// has taken its own, which it does before returning: what keeps the
// delivery owned through the backoff window is ScheduleRetry's
// reference, not this one.
func (e *Engine) rescheduleRecovered(
webhookDB *gorm.DB, task Task, delay time.Duration,
) bool {
if !e.takeForRedispatch(
webhookDB, task.DeliveryID,
database.DeliveryStatusRetrying,
) {
return false
}
defer e.inflight.release(task.DeliveryID)
e.ScheduleRetry(task, delay)
return true
}
// countAttemptsBatch counts the recorded attempts of every delivery
// in a batch with one grouped query, keyed by delivery id. Deliveries
// with no attempts are simply absent from the result, which reads back
// as the zero this caller wants.
//
// One query rather than one per delivery: this runs on the recovery
// path, which is a burst of writes against a database that has just
// been under enough contention to strand these rows in the first
// place. See https://git.eeqj.de/sneak/webhooker/issues/256.
func (e *Engine) countAttemptsBatch(
webhookDB *gorm.DB, deliveries []database.Delivery,
) map[string]int {
counts := make(map[string]int, len(deliveries))
if len(deliveries) == 0 {
return counts
}
ids := make([]string, 0, len(deliveries))
for i := range deliveries {
ids = append(ids, deliveries[i].ID)
}
// One delivery id per recorded attempt, tallied here rather than
// grouped in SQL: internal/gormlog forbids (*gorm.DB).Scan, which
// a GROUP BY into a struct would need, and an attempt row per
// delivery is bounded by the target's MaxRetries.
var attemptIDs []string
err := webhookDB.
Model(&database.DeliveryResult{}).
Where("delivery_id IN ?", ids).
Pluck("delivery_id", &attemptIDs).Error
if err != nil {
e.log.Error(
"failed to count delivery attempts for recovery",
"error", err,
)
return counts
}
for _, id := range attemptIDs {
counts[id]++
}
return counts
}
func (e *Engine) loadEvent( func (e *Engine) loadEvent(
webhookDB *gorm.DB, eventID string, webhookDB *gorm.DB, eventID string,
) (database.Event, error) { ) (database.Event, error) {
@@ -1132,6 +1675,10 @@ func buildRecoveryTask(
func (e *Engine) loadTargetMap( func (e *Engine) loadTargetMap(
deliveries []database.Delivery, deliveries []database.Delivery,
) map[string]database.Target { ) map[string]database.Target {
if len(deliveries) == 0 {
return nil
}
seen := make(map[string]bool) seen := make(map[string]bool)
targetIDs := make([]string, 0, len(deliveries)) targetIDs := make([]string, 0, len(deliveries))
@@ -1168,12 +1715,33 @@ func (e *Engine) loadTargetMap(
return targetMap return targetMap
} }
// sendRecoveredDeliveries re-dispatches pending deliveries, skipping
// the ids in settled — those already reached their receiver and have
// been marked delivered by reconcileDelivered.
//
// The skip and takeForRedispatch's status check answer different
// questions and neither replaces the other. This one is "did this
// delivery already succeed", which is what settles the row to
// delivered instead of sending it, and which is the only thing that
// keeps the retrying paths from terminally failing a delivery that
// reconcile just settled. The status check is "is the row still what
// the batch query said it was", which catches a worker settling it to
// anything at all in between.
func (e *Engine) sendRecoveredDeliveries( func (e *Engine) sendRecoveredDeliveries(
ctx context.Context, ctx context.Context,
webhookDB *gorm.DB,
deliveries []database.Delivery, deliveries []database.Delivery,
webhookID string, webhookID string,
targetMap map[string]database.Target, targetMap map[string]database.Target,
settled map[string]struct{},
) { ) {
// The attempt number continues each delivery's own history
// rather than restarting at 1. A recovered delivery may already
// have recorded attempts, and numbering the next one 1 again
// both collides in the event log and hands the retry path a
// backoff computed from the wrong attempt.
attempts := e.countAttemptsBatch(webhookDB, deliveries)
for i := range deliveries { for i := range deliveries {
select { select {
case <-ctx.Done(): case <-ctx.Done():
@@ -1181,6 +1749,10 @@ func (e *Engine) sendRecoveredDeliveries(
default: default:
} }
if _, ok := settled[deliveries[i].ID]; ok {
continue
}
target, ok := targetMap[deliveries[i].TargetID] target, ok := targetMap[deliveries[i].TargetID]
if !ok { if !ok {
e.log.Error( e.log.Error(
@@ -1192,22 +1764,40 @@ func (e *Engine) sendRecoveredDeliveries(
continue continue
} }
task := buildRecoveryTask( if !e.takeForRedispatch(
&deliveries[i], webhookID, webhookDB, deliveries[i].ID,
&deliveries[i].Event, &target, 1, database.DeliveryStatusPending,
) ) {
continue
}
select { // The body is read here, one delivery at a time and only for
case e.deliveryCh <- task: // deliveries that are actually being sent, rather than
default: // preloaded across the whole batch. A batch is up to
e.log.Warn( // pendingSweepBatch rows at up to the 1 MB ingest cap, and
"delivery channel full during "+ // most of a sweep's batch is refused by the gate above — so
"recovery, remaining deliveries "+ // preloading would hold hundreds of megabytes per webhook per
"will be recovered on next restart", // tick to build tasks it then discards.
event, err := e.loadEvent(
webhookDB, deliveries[i].EventID,
)
if err != nil {
e.log.Error(
"failed to load event for recovered delivery",
"delivery_id", deliveries[i].ID, "delivery_id", deliveries[i].ID,
"event_id", deliveries[i].EventID,
"error", err,
)
e.inflight.release(deliveries[i].ID)
continue
}
task := buildRecoveryTask(
&deliveries[i], webhookID, &event, &target,
attempts[deliveries[i].ID]+1,
) )
return e.queueRecovered(e.deliveryCh, task)
}
} }
} }

View File

@@ -2,7 +2,6 @@ package delivery_test
import ( import (
"context" "context"
"database/sql"
"encoding/json" "encoding/json"
"fmt" "fmt"
"io" "io"
@@ -70,11 +69,12 @@ func iMainDB(t *testing.T) *gorm.DB {
t.TempDir(), "main-test.db", t.TempDir(), "main-test.db",
) )
dsn := fmt.Sprintf( // Opened the way the service opens the main database, so these
"file:%s?cache=shared&mode=rwc", dbPath, // tests cannot pass against journal and locking settings
// production does not use.
sqlDB, err := database.OpenSQLite(
dbPath, database.SQLiteModeCreate,
) )
sqlDB, err := sql.Open("sqlite", dsn)
require.NoError(t, err) require.NoError(t, err)
t.Cleanup(func() { _ = sqlDB.Close() }) t.Cleanup(func() { _ = sqlDB.Close() })

View File

@@ -3,7 +3,6 @@ package delivery_test
import ( import (
"bytes" "bytes"
"context" "context"
"database/sql"
"encoding/json" "encoding/json"
"fmt" "fmt"
"log/slog" "log/slog"
@@ -37,11 +36,12 @@ func testWebhookDB(t *testing.T) *gorm.DB {
t.TempDir(), "events-test.db", t.TempDir(), "events-test.db",
) )
dsn := fmt.Sprintf( // Opened the way the service opens a per-webhook database, so
"file:%s?cache=shared&mode=rwc", dbPath, // these tests cannot pass against journal and locking settings
// production does not use.
sqlDB, err := database.OpenSQLite(
dbPath, database.SQLiteModeCreate,
) )
sqlDB, err := sql.Open("sqlite", dsn)
require.NoError(t, err) require.NoError(t, err)
t.Cleanup(func() { _ = sqlDB.Close() }) t.Cleanup(func() { _ = sqlDB.Close() })

View File

@@ -33,6 +33,11 @@ const (
// response is written against this number, so a test has to // response is written against this number, so a test has to
// be able to name it. // be able to name it.
ExportMaxBodyLog = maxBodyLog ExportMaxBodyLog = maxBodyLog
// ExportPendingSweepMinAge is how long a delivery must sit at
// pending before the sweep treats it as stranded. A test has to
// name it to age a row past the bound.
ExportPendingSweepMinAge = pendingSweepMinAge
) )
// ExportIsBlockedIP exposes isBlockedIP for testing. // ExportIsBlockedIP exposes isBlockedIP for testing.
@@ -286,6 +291,26 @@ func (e *Engine) ExportWedgeWorker(release <-chan struct{}) {
}) })
} }
// ExportInflightHeld reports how many deliveries the engine currently
// owns, so a test can prove ownership is released rather than leaked.
func (e *Engine) ExportInflightHeld() int {
return e.inflight.held()
}
// ExportRetainDelivery takes the first reference on a delivery, as the
// queueing side does. It lets a test put a delivery into the state a
// worker or a full channel would, without running the pool.
func (e *Engine) ExportRetainDelivery(deliveryID string) bool {
return e.inflight.retainIdle(deliveryID)
}
// ExportRecoverRetryingDeliveries exposes recoverRetryingDeliveries.
func (e *Engine) ExportRecoverRetryingDeliveries(
webhookDB *gorm.DB, webhookID string,
) {
e.recoverRetryingDeliveries(webhookDB, webhookID)
}
// ExportDeliveryCh returns the delivery channel. // ExportDeliveryCh returns the delivery channel.
func (e *Engine) ExportDeliveryCh() chan Task { func (e *Engine) ExportDeliveryCh() chan Task {
return e.deliveryCh return e.deliveryCh

View File

@@ -0,0 +1,110 @@
package delivery
import "sync"
// inflightSet records which deliveries the engine currently owns.
//
// A delivery is owned from the moment a task for it is handed to a
// channel or to a retry timer until the engine has no further plan for
// it in memory. Restart recovery and both arms of the periodic sweep
// re-dispatch only deliveries the set does not hold, which is what
// makes them exact rather than a guess about how long a row has sat at
// pending.
//
// This replaces reasoning from timestamps. A delivery's row says
// pending from creation until its outcome is written, which covers
// four different situations — never dispatched, waiting in a channel,
// being attempted right now, and genuinely stranded — and no column
// distinguishes them. Only the engine knows which, and it knows
// exactly. `deliveryChannelSize` is 10000 against 10 workers, so a
// perfectly healthy delivery can wait far longer than any age bound
// worth setting before its attempt even begins; an age bound alone
// re-sends it. See
// https://git.eeqj.de/sneak/webhooker/issues/256.
//
// In-memory state is sufficient because a data directory admits one
// process: internal/datadir takes an flock on it at startup and a
// second instance refuses to run. Deliveries owned by a process that
// died are not in any successor's set, and restart recovery is what
// picks those up.
//
// References are counted rather than held as a plain set because
// ownership outlives the worker that took it. A target that schedules
// a retry from inside Deliver adds a reference while the worker still
// holds one, so the delivery stays owned across the gap between the
// worker returning and the timer firing — the window in which a sweep
// would otherwise find the row at retrying and send it again.
//
// The zero value is ready to use, and the Engine holds one by value.
// That is deliberate: an engine built by a constructor that forgot to
// initialise this would not refuse to re-dispatch anything, and the
// symptom would be duplicate deliveries rather than a failure anybody
// notices.
type inflightSet struct {
mu sync.Mutex
ids map[string]int
}
// retain adds a reference to a delivery the caller already knows the
// engine owns, so that ownership survives the current holder letting
// go. It cannot fail.
func (s *inflightSet) retain(deliveryID string) {
s.mu.Lock()
defer s.mu.Unlock()
if s.ids == nil {
s.ids = make(map[string]int)
}
s.ids[deliveryID]++
}
// retainIdle takes the first reference on a delivery, and reports
// whether it got it. It fails when the engine already owns the
// delivery, which is what makes two claimants — restart recovery and
// the sweep run concurrently, or two sweep arms — mutually exclusive
// rather than merely atomic.
func (s *inflightSet) retainIdle(deliveryID string) bool {
s.mu.Lock()
defer s.mu.Unlock()
if s.ids[deliveryID] > 0 {
return false
}
if s.ids == nil {
s.ids = make(map[string]int)
}
s.ids[deliveryID] = 1
return true
}
// release drops one reference. The delivery becomes eligible for
// re-dispatch again once the last one goes.
func (s *inflightSet) release(deliveryID string) {
s.mu.Lock()
defer s.mu.Unlock()
n := s.ids[deliveryID] - 1
if n <= 0 {
delete(s.ids, deliveryID)
return
}
s.ids[deliveryID] = n
}
// held reports how many deliveries the engine currently owns. It
// exists so a test can assert that ownership is released rather than
// leaked: a reference that is never dropped hides its delivery from
// every sweep for the life of the process, which is the one way this
// mechanism can fail silently.
func (s *inflightSet) held() int {
s.mu.Lock()
defer s.mu.Unlock()
return len(s.ids)
}

View File

@@ -0,0 +1,428 @@
package delivery_test
import (
"context"
"path/filepath"
"strings"
"sync"
"testing"
"time"
"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"
)
// These tests pin the rule that decides whether a delivery may be
// handed back to a worker: the engine re-dispatches only what it does
// not already own. Age alone is not that rule — a healthy delivery
// waiting in a 10000-deep channel is old and must not be re-sent. See
// https://git.eeqj.de/sneak/webhooker/issues/256.
// fSweepSetup seeds the main database with the webhook row the sweep
// enumerates, and returns the setup.
func fSweepSetup(
t *testing.T, targetID, name string,
) iSetup {
t.Helper()
s := newISetup(t)
iCreateTarget(t, s.MainDB, targetID,
s.WebhookID, name,
database.TargetTypeLog, "", 0,
)
require.NoError(t, s.MainDB.Create(&database.Webhook{
BaseModel: database.BaseModel{ID: s.WebhookID},
UserID: uuid.New().String(),
Name: name,
}).Error)
return s
}
// fDrain collects every task the engine has queued.
//
// Every caller drives the dispatch paths synchronously and has already
// waited for them to return, so anything they queued is in the channel
// by now. The short grace covers nothing but scheduler jitter, and is
// kept small because one of these tests runs the drain forty times.
func fDrain(e *delivery.Engine) []delivery.Task {
var out []delivery.Task
for {
select {
case task := <-e.ExportDeliveryCh():
out = append(out, task)
case task := <-e.ExportRetryCh():
out = append(out, task)
case <-time.After(25 * time.Millisecond):
return out
}
}
}
// TestArchiveHandleIsWAL closes the last gap in the durability
// evidence: the main and per-webhook tiers each assert their journal
// mode on a live handle, and the archive tier gets its settings from
// the same code path but nothing checked the running file.
func TestArchiveHandleIsWAL(t *testing.T) {
t.Parallel()
w := delivery.NewExportArchiveWriter(
filepath.Join(t.TempDir(), "archive-wal.db"),
archiveTestLogger(), 0,
)
require.NoError(t, w.Open(0))
var mode string
row := w.DB().Raw("pragma journal_mode").Row()
require.NoError(t, row.Scan(&mode))
assert.Equal(t, "wal", strings.ToLower(mode))
var busy string
row = w.DB().Raw("pragma busy_timeout").Row()
require.NoError(t, row.Scan(&busy))
assert.Equal(t, "10000", busy)
}
// TestSweepLeavesAQueuedDeliveryAlone is the case the age bound cannot
// see. The delivery is queued and untouched, so its row is arbitrarily
// old and still perfectly healthy; only ownership distinguishes it
// from a stranded one.
func TestSweepLeavesAQueuedDeliveryAlone(t *testing.T) {
t.Parallel()
targetID := uuid.New().String()
s := fSweepSetup(t, targetID, "queued")
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"queued":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
rAgePending(t, s.WebhookDB, d.ID)
// Queued exactly as the receiver queues it, and never dequeued:
// no workers are running in this engine.
s.Engine.Notify([]delivery.Task{{
DeliveryID: d.ID,
EventID: event.ID,
WebhookID: s.WebhookID,
TargetID: targetID,
}})
require.Equal(t, 1, s.Engine.ExportInflightHeld())
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
tasks := fDrain(s.Engine)
assert.Len(
t, tasks, 1,
"the sweep must not queue a delivery that is "+
"already waiting for a worker",
)
}
// TestRecoveryAndSweepDoNotDoubleDispatch drives the two entry points
// the engine starts concurrently against one aged pending row. Before
// ownership they both dispatched it.
func TestRecoveryAndSweepDoNotDoubleDispatch(t *testing.T) {
t.Parallel()
targetID := uuid.New().String()
s := fSweepSetup(t, targetID, "racing")
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"racing":true}`,
)
ctx := context.Background()
for range 40 {
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
rAgePending(t, s.WebhookDB, d.ID)
var wg sync.WaitGroup
wg.Go(func() {
s.Engine.ExportRecoverPendingDeliveries(
ctx, s.WebhookDB, s.WebhookID,
)
})
wg.Go(func() {
s.Engine.ExportSweepWebhookRetries(
ctx, s.WebhookID,
)
})
wg.Wait()
tasks := fDrain(s.Engine)
require.Len(
t, tasks, 1,
"delivery %s dispatched %d times",
d.ID, len(tasks),
)
// No worker runs in this engine, so the reference the winner
// took is never released and earlier iterations' deliveries
// stay owned — which is itself the property under test, since
// both paths see them on every subsequent pass.
}
}
// TestConcurrentClaimsOfOneDeliveryYieldOneOwner exercises the
// exclusion directly, rather than arguing it from a SQL predicate.
func TestConcurrentClaimsOfOneDeliveryYieldOneOwner(
t *testing.T,
) {
t.Parallel()
eng := newISetup(t).Engine
deliveryID := uuid.New().String()
var (
wg sync.WaitGroup
mu sync.Mutex
won int
)
for range 64 {
wg.Go(func() {
if eng.ExportRetainDelivery(deliveryID) {
mu.Lock()
won++
mu.Unlock()
}
})
}
wg.Wait()
assert.Equal(t, 1, won)
assert.Equal(t, 1, eng.ExportInflightHeld())
}
// TestOwnershipIsReleasedAfterDelivery guards the other direction: a
// leaked reference hides a delivery from every sweep for the life of
// the process.
func TestOwnershipIsReleasedAfterDelivery(t *testing.T) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
iCreateTarget(t, s.MainDB, targetID,
s.WebhookID, "released",
database.TargetTypeLog, "", 0,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"released":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
s.Engine.ExportStart()
defer func() {
require.NoError(
t, s.Engine.ExportStop(context.Background()),
)
}()
body := `{"released":true}`
s.Engine.Notify([]delivery.Task{{
DeliveryID: d.ID,
EventID: event.ID,
WebhookID: s.WebhookID,
TargetID: targetID,
TargetName: "released",
TargetType: database.TargetTypeLog,
Body: &body,
EntrypointID: event.EntrypointID,
}})
iWaitForDelivered(t, s.WebhookDB, d.ID)
assert.Eventually(
t,
func() bool {
return s.Engine.ExportInflightHeld() == 0
},
2*time.Second, 20*time.Millisecond,
"the delivery stayed owned after it was delivered",
)
}
// TestRetryingRecoverySkipsASuccessfulResult is the retrying-side twin
// of the pending reconcile. A second attempt that reached the receiver
// and whose status write then failed sits at retrying holding a
// successful result, and re-sending it is the same duplicate.
func TestRetryingRecoverySkipsASuccessfulResult(t *testing.T) {
t.Parallel()
targetID := uuid.New().String()
s := fSweepSetup(t, targetID, "retry-settled")
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"retry":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
rSeedResult(t, s.WebhookDB, d.ID, 1, false)
rSeedResult(t, s.WebhookDB, d.ID, 2, true)
s.Engine.ExportRecoverRetryingDeliveries(
s.WebhookDB, s.WebhookID,
)
assert.Empty(
t, fDrain(s.Engine),
"a retrying delivery holding a successful result "+
"must not be sent again",
)
iAssertStatus(
t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered,
)
}
// TestRetryingSweepSkipsASuccessfulResult is the same rule on the
// periodic sweep's retrying arm.
func TestRetryingSweepSkipsASuccessfulResult(t *testing.T) {
t.Parallel()
targetID := uuid.New().String()
s := fSweepSetup(t, targetID, "retry-swept")
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"swept":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
rSeedResult(t, s.WebhookDB, d.ID, 1, false)
rSeedResult(t, s.WebhookDB, d.ID, 2, true)
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
assert.Empty(t, fDrain(s.Engine))
iAssertStatus(
t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered,
)
var attempts int64
require.NoError(t, s.WebhookDB.
Model(&database.DeliveryResult{}).
Where("delivery_id = ?", d.ID).
Count(&attempts).Error)
assert.Equal(
t, int64(2), attempts,
"settling must not invent an attempt",
)
}
// TestScheduledRetryIsNotSweptDuringBackoff closes the window between
// a target scheduling a retry and the timer firing. The row says
// retrying and nothing is running, which is exactly what an orphaned
// retry looks like from the database.
func TestScheduledRetryIsNotSweptDuringBackoff(t *testing.T) {
t.Parallel()
targetID := uuid.New().String()
s := fSweepSetup(t, targetID, "backoff")
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"backoff":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
s.Engine.ExportScheduleRetry(delivery.Task{
DeliveryID: d.ID,
EventID: event.ID,
WebhookID: s.WebhookID,
TargetID: targetID,
AttemptNum: 2,
}, time.Hour)
require.Equal(t, 1, s.Engine.ExportInflightHeld())
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
assert.Empty(
t, fDrain(s.Engine),
"the sweep must not duplicate a retry that is "+
"already scheduled",
)
}
// TestRedispatchStampsTheRow pins the cadence control: a stranded
// delivery that has just been handed out is not selected again by the
// next tick a minute later.
func TestRedispatchStampsTheRow(t *testing.T) {
t.Parallel()
targetID := uuid.New().String()
s := fSweepSetup(t, targetID, "stamped")
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"stamped":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
rAgePending(t, s.WebhookDB, d.ID)
ctx := context.Background()
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
require.Len(t, fDrain(s.Engine), 1)
var row database.Delivery
require.NoError(t, s.WebhookDB.
First(&row, "id = ?", d.ID).Error)
assert.WithinDuration(
t, time.Now(), row.UpdatedAt, time.Minute,
"a re-dispatched delivery must be stamped so the "+
"next tick does not select it again",
)
}

View File

@@ -3,8 +3,6 @@ package delivery_test
import ( import (
"bytes" "bytes"
"context" "context"
"database/sql"
"fmt"
"log/slog" "log/slog"
"net/http" "net/http"
"path/filepath" "path/filepath"
@@ -54,12 +52,10 @@ func (q *qdSyncBuf) String() string {
func qdMainDB(t *testing.T, log *slog.Logger) *gorm.DB { func qdMainDB(t *testing.T, log *slog.Logger) *gorm.DB {
t.Helper() t.Helper()
dsn := fmt.Sprintf( sqlDB, err := database.OpenSQLite(
"file:%s?cache=shared&mode=rwc",
filepath.Join(t.TempDir(), "main-gormlog.db"), filepath.Join(t.TempDir(), "main-gormlog.db"),
database.SQLiteModeCreate,
) )
sqlDB, err := sql.Open("sqlite", dsn)
require.NoError(t, err) require.NoError(t, err)
t.Cleanup(func() { _ = sqlDB.Close() }) t.Cleanup(func() { _ = sqlDB.Close() })

View File

@@ -0,0 +1,378 @@
package delivery_test
import (
"context"
"net/http"
"net/http/httptest"
"sync/atomic"
"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"
)
// These tests cover the delivery half of
// https://git.eeqj.de/sneak/webhooker/issues/256: a delivery that
// reached its receiver but whose bookkeeping write failed used to be
// left at pending and re-sent on the next restart, giving the receiver
// a second copy while the event log recorded one attempt.
// rSeedResult records a DeliveryResult against a delivery, standing in
// for the attempt row the send path writes before the status.
func rSeedResult(
t *testing.T,
db *gorm.DB,
deliveryID string,
attemptNum int,
success bool,
) {
t.Helper()
require.NoError(t, db.Create(&database.DeliveryResult{
DeliveryID: deliveryID,
AttemptNum: attemptNum,
Success: success,
}).Error)
}
// rAgePending backdates a delivery past the sweep's age bound, which is
// what separates a stranded delivery from one a worker still holds.
func rAgePending(
t *testing.T, db *gorm.DB, deliveryID string,
) {
t.Helper()
old := time.Now().Add(
-2 * delivery.ExportPendingSweepMinAge,
)
require.NoError(t, db.Model(&database.Delivery{}).
Where("id = ?", deliveryID).
UpdateColumn("updated_at", old).Error)
}
func TestRecoverySkipsPendingWithSuccessfulResult(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
iCreateTarget(t, s.MainDB, targetID,
s.WebhookID, "already-delivered",
database.TargetTypeLog, "", 0,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"delivered":true}`,
)
// The delivery whose send succeeded and whose result row landed:
// only the status write failed, so it sits at pending.
done := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
rSeedResult(t, s.WebhookDB, done.ID, 1, true)
// A delivery that was genuinely never attempted.
fresh := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
s.Engine.ExportRecoverPendingDeliveries(
context.Background(), s.WebhookDB, s.WebhookID,
)
select {
case task := <-s.Engine.ExportDeliveryCh():
assert.Equal(
t, fresh.ID, task.DeliveryID,
"only the unattempted delivery may be re-sent",
)
case <-time.After(2 * time.Second):
t.Fatal("expected the unattempted delivery")
}
select {
case task := <-s.Engine.ExportDeliveryCh():
t.Fatalf(
"re-sent an already delivered delivery: %s",
task.DeliveryID,
)
case <-time.After(200 * time.Millisecond):
}
// It is settled rather than merely skipped: leaving it pending
// would strand it again on the next sweep.
iAssertStatus(
t, s.WebhookDB, done.ID,
database.DeliveryStatusDelivered,
)
}
// TestRecoveryContinuesTheAttemptNumbering pins the audit trail: a
// recovered delivery that already recorded two attempts is re-sent as
// attempt three, not as attempt one again.
func TestRecoveryContinuesTheAttemptNumbering(t *testing.T) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
iCreateTarget(t, s.MainDB, targetID,
s.WebhookID, "numbering",
database.TargetTypeLog, "", 0,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"numbering":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
rSeedResult(t, s.WebhookDB, d.ID, 1, false)
rSeedResult(t, s.WebhookDB, d.ID, 2, false)
s.Engine.ExportRecoverPendingDeliveries(
context.Background(), s.WebhookDB, s.WebhookID,
)
select {
case task := <-s.Engine.ExportDeliveryCh():
assert.Equal(t, d.ID, task.DeliveryID)
assert.Equal(t, 3, task.AttemptNum)
case <-time.After(2 * time.Second):
t.Fatal("expected the delivery to be recovered")
}
}
// TestSweepRecoversStrandedPending is the half that removes the
// restart requirement: a delivery left at pending is picked up by the
// periodic sweep.
func TestSweepRecoversStrandedPending(t *testing.T) {
t.Parallel()
targetID := uuid.New().String()
s := fSweepSetup(t, targetID, "stranded")
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"stranded":true}`,
)
stranded := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
rAgePending(t, s.WebhookDB, stranded.ID)
// A delivery a worker may still be holding: young, and therefore
// none of the sweep's business.
inFlight := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
select {
case task := <-s.Engine.ExportDeliveryCh():
assert.Equal(t, stranded.ID, task.DeliveryID)
case <-time.After(2 * time.Second):
t.Fatal("expected the stranded delivery")
}
select {
case task := <-s.Engine.ExportDeliveryCh():
t.Fatalf(
"swept an in-flight delivery: %s",
task.DeliveryID,
)
case <-time.After(200 * time.Millisecond):
}
iAssertStatus(
t, s.WebhookDB, inFlight.ID,
database.DeliveryStatusPending,
)
}
// TestSweepClaimsAStrandedDeliveryOnlyOnce guards the repeat the sweep
// would otherwise be: the row stays pending for as long as the attempt
// runs, and a sweep a minute later must not send it a second time.
func TestSweepClaimsAStrandedDeliveryOnlyOnce(t *testing.T) {
t.Parallel()
targetID := uuid.New().String()
s := fSweepSetup(t, targetID, "claimed")
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"claimed":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
rAgePending(t, s.WebhookDB, d.ID)
ctx := context.Background()
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
select {
case task := <-s.Engine.ExportDeliveryCh():
assert.Equal(t, d.ID, task.DeliveryID)
case <-time.After(2 * time.Second):
t.Fatal("expected the stranded delivery")
}
// The delivery is still pending — nothing has run it yet — but
// the claim must keep the next sweep off it.
iAssertStatus(
t, s.WebhookDB, d.ID,
database.DeliveryStatusPending,
)
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
select {
case task := <-s.Engine.ExportDeliveryCh():
t.Fatalf(
"sent a claimed delivery again: %s",
task.DeliveryID,
)
case <-time.After(200 * time.Millisecond):
}
}
// TestSweepSettlesStrandedPendingWithoutResending is the sweep's own
// version of the reconcile: a stranded delivery holding a successful
// result is settled where it stands, and the receiver hears nothing.
func TestSweepSettlesStrandedPendingWithoutResending(
t *testing.T,
) {
t.Parallel()
targetID := uuid.New().String()
s := fSweepSetup(t, targetID, "settled")
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"settled":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
rSeedResult(t, s.WebhookDB, d.ID, 1, true)
rAgePending(t, s.WebhookDB, d.ID)
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
select {
case task := <-s.Engine.ExportDeliveryCh():
t.Fatalf(
"re-sent a delivery that already succeeded: %s",
task.DeliveryID,
)
case <-time.After(200 * time.Millisecond):
}
iAssertStatus(
t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered,
)
var attempts int64
require.NoError(t, s.WebhookDB.
Model(&database.DeliveryResult{}).
Where("delivery_id = ?", d.ID).
Count(&attempts).Error)
assert.Equal(
t, int64(1), attempts,
"settling must not invent an attempt",
)
}
// TestFailedResultWriteLeavesDeliveryRecoverable is the rule the
// targets now follow: a bookkeeping write that fails must not advance
// the status, because pending and retrying are the states the sweeps
// recover and delivered is a claim the database refused to record.
func TestFailedResultWriteLeavesDeliveryRecoverable(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
var hits atomic.Int64
ts := httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
hits.Add(1)
w.WriteHeader(http.StatusOK)
},
))
defer ts.Close()
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"unwritable":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
// Drop the table the attempt row goes in, so the send succeeds
// and only the bookkeeping write fails.
require.NoError(
t,
s.WebhookDB.Exec("drop table delivery_results").Error,
)
full := &database.Delivery{
EventID: event.ID,
TargetID: targetID,
Status: database.DeliveryStatusPending,
Event: event,
Target: database.Target{
Name: "unwritable",
Type: database.TargetTypeHTTP,
Config: iHTTPConfig(ts.URL),
},
}
full.ID = d.ID
s.Engine.ExportDeliverHTTP(
context.Background(), s.WebhookDB, full,
&delivery.Task{DeliveryID: d.ID, AttemptNum: 1},
)
assert.Equal(
t, int64(1), hits.Load(),
"the send itself must still happen",
)
iAssertStatus(
t, s.WebhookDB, d.ID,
database.DeliveryStatusPending,
)
}

View File

@@ -23,6 +23,13 @@ type ConfigField struct {
Value string Value string
} }
// deletedNameSuffix marks the name of a target that no longer
// exists. Deletes are soft and delivery history outlives the
// target, so the event log shows names of targets that are gone;
// an operator reading one needs to know it cannot be delivered
// to, replayed to, or configured.
const deletedNameSuffix = " (deleted)"
// TargetView is the display-safe projection of a target for // TargetView is the display-safe projection of a target for
// the UI. It deliberately has no raw configuration field, so // the UI. It deliberately has no raw configuration field, so
// no template — present or future — can render the stored // no template — present or future — can render the stored
@@ -30,14 +37,37 @@ type ConfigField struct {
type TargetView struct { type TargetView struct {
ID string ID string
Name string Name string
// Deleted reports that this target's row is soft deleted.
// Only views built for historical display carry it set:
// every other projection is of a live row.
Deleted bool
Type database.TargetType Type database.TargetType
Active bool Active bool
Config []ConfigField Config []ConfigField
} }
// DisplayName is the name to render, marked when the target has
// been deleted. Templates showing a name against historical data
// must use it rather than Name, which stays the stored name.
func (v TargetView) DisplayName() string {
if v.Deleted {
return v.Name + deletedNameSuffix
}
return v.Name
}
// NewTargetViews projects targets for rendering, replacing // NewTargetViews projects targets for rendering, replacing
// each stored configuration blob with named, display-safe // each stored configuration blob with named, display-safe
// fields. // fields.
//
// A soft-deleted row projects exactly as a live one does, minus
// the deleted marker on its name: masking is a property of the
// projection, not of the row's state, so a deleted target's
// credential is as unreachable from a template as a live
// target's.
func NewTargetViews( func NewTargetViews(
targets []database.Target, targets []database.Target,
) []TargetView { ) []TargetView {
@@ -49,6 +79,7 @@ func NewTargetViews(
views = append(views, TargetView{ views = append(views, TargetView{
ID: t.ID, ID: t.ID,
Name: t.Name, Name: t.Name,
Deleted: t.DeletedAt.Valid,
Type: t.Type, Type: t.Type,
Active: t.Active, Active: t.Active,
Config: targetConfigFields(t), Config: targetConfigFields(t),

View File

@@ -2,9 +2,11 @@ package delivery_test
import ( import (
"testing" "testing"
"time"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/delivery"
) )
@@ -17,6 +19,14 @@ const (
slackWebhookURL = "https://hooks.slack.com" + slackWebhookURL = "https://hooks.slack.com" +
slackSecretPath slackSecretPath
// slackMaskedURL is what a Slack webhook URL renders as
// once masked: scheme and host, path elided.
slackMaskedURL = "https://hooks.slack.com/..."
// slackTargetName is the target name the Slack projection
// tests use.
slackTargetName = "slack-target"
viewExampleOrigin = "https://example.com" viewExampleOrigin = "https://example.com"
viewExampleHook = viewExampleOrigin + "/hook" viewExampleHook = viewExampleOrigin + "/hook"
viewMaskedOrigin = viewExampleOrigin + "/..." viewMaskedOrigin = viewExampleOrigin + "/..."
@@ -33,7 +43,7 @@ func TestMaskedWebhookURL(t *testing.T) {
}{ }{
"slack webhook": { "slack webhook": {
url: slackWebhookURL, url: slackWebhookURL,
want: "https://hooks.slack.com/...", want: slackMaskedURL,
}, },
"query string dropped": { "query string dropped": {
url: viewExampleOrigin + "/a?token=secret", url: viewExampleOrigin + "/a?token=secret",
@@ -125,23 +135,61 @@ func viewFor(
return views[0] return views[0]
} }
func TestNewTargetViews_Slack(t *testing.T) { // TestNewTargetViews_DeletedTarget proves the projection marks
// a soft-deleted target's name and masks its configuration by
// the same rules a live target's is. Delivery history outlives
// the target it names, so this projection is what an operator
// reads about a target that no longer exists.
func TestNewTargetViews_DeletedTarget(t *testing.T) {
t.Parallel() t.Parallel()
view := viewFor(t, database.Target{ target := slackTarget()
Name: "slack-target", target.DeletedAt = gorm.DeletedAt{
Time: time.Now(),
Valid: true,
}
view := viewFor(t, target)
assert.True(t, view.Deleted)
assert.Equal(t, slackTargetName, view.Name)
assert.Equal(
t, slackTargetName+" (deleted)", view.DisplayName(),
)
assert.Equal(
t,
map[string]string{"Webhook URL": slackMaskedURL},
fieldMap(view.Config),
)
}
// slackTarget is the live Slack target the projection tests
// share.
func slackTarget() database.Target {
return database.Target{
Name: slackTargetName,
Type: database.TargetTypeSlack, Type: database.TargetTypeSlack,
Active: true, Active: true,
Config: `{"webhookUrl":"` + Config: `{"webhookUrl":"` +
slackWebhookURL + `"}`, slackWebhookURL + `"}`,
}) }
}
func TestNewTargetViews_Slack(t *testing.T) {
t.Parallel()
view := viewFor(t, slackTarget())
assert.Equal(t, slackTargetName, view.Name)
// A live target is never marked, so the marker cannot
// reach a name that still exists.
assert.False(t, view.Deleted)
assert.Equal(t, slackTargetName, view.DisplayName())
assert.Equal(t, "slack-target", view.Name)
assert.Equal( assert.Equal(
t, t,
map[string]string{ map[string]string{"Webhook URL": slackMaskedURL},
"Webhook URL": "https://hooks.slack.com/...",
},
fieldMap(view.Config), fieldMap(view.Config),
) )
} }
@@ -212,7 +260,7 @@ func TestNewTargetViews_HTTPMasksDestinationURL(t *testing.T) {
assert.Equal( assert.Equal(
t, t,
"https://hooks.slack.com/...", slackMaskedURL,
fields["Destination URL"], fields["Destination URL"],
) )

View File

@@ -58,12 +58,17 @@ func (t *databaseTarget) Deliver(
"error", err, "error", err,
) )
t.eng.recordResult( recErr := t.eng.recordResult(
webhookDB, d, 1, false, 0, "", webhookDB, d, 1, false, 0, "",
err.Error(), elapsed.Milliseconds(), err.Error(), elapsed.Milliseconds(),
) )
if recErr != nil {
t.eng.bookkeepingFailed(d, recErr)
t.eng.updateDeliveryStatus( return
}
t.eng.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusFailed, database.DeliveryStatusFailed,
) )
@@ -71,12 +76,17 @@ func (t *databaseTarget) Deliver(
return return
} }
t.eng.recordResult( recErr := t.eng.recordResult(
webhookDB, d, 1, true, 0, "", "", webhookDB, d, 1, true, 0, "", "",
elapsed.Milliseconds(), elapsed.Milliseconds(),
) )
if recErr != nil {
t.eng.bookkeepingFailed(d, recErr)
t.eng.updateDeliveryStatus( return
}
t.eng.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusDelivered, database.DeliveryStatusDelivered,
) )

View File

@@ -1,7 +1,6 @@
package delivery package delivery
import ( import (
"database/sql"
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
@@ -12,6 +11,7 @@ import (
"gorm.io/driver/sqlite" "gorm.io/driver/sqlite"
"gorm.io/gorm" "gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/gormlog" "sneak.berlin/go/webhooker/internal/gormlog"
) )
@@ -30,13 +30,13 @@ const (
// path: open the archive file, creating it if missing, so a // path: open the archive file, creating it if missing, so a
// first write (or a write after the operator moved the file // first write (or a write after the operator moved the file
// away) recreates it. // away) recreates it.
archiveModeCreate = "rwc" archiveModeCreate = database.SQLiteModeCreate
// archiveModeExisting is the SQLite URI mode used by the idle // archiveModeExisting is the SQLite URI mode used by the idle
// sweep: open read-write but never create. A sweep must never // sweep: open read-write but never create. A sweep must never
// conjure an empty archive file for a webhook that has a // conjure an empty archive file for a webhook that has a
// database target but has never received an event. // database target but has never received an event.
archiveModeExisting = "rw" archiveModeExisting = database.SQLiteModeExisting
) )
var ( var (
@@ -273,9 +273,11 @@ func (w *archiveWriter) open(expiry time.Duration) error {
func (w *archiveWriter) openMode( func (w *archiveWriter) openMode(
mode string, expiry time.Duration, mode string, expiry time.Duration,
) error { ) error {
dbURL := fmt.Sprintf("file:%s?mode=%s", w.path, mode) // Opened through database.OpenSQLite so an archive file carries
// the same WAL journaling, busy timeout, immediate-transaction
sqlDB, err := sql.Open("sqlite", dbURL) // locking, and pool bounds as every other database file. See
// internal/database/sqlite_open.go.
sqlDB, err := database.OpenSQLite(w.path, mode)
if err != nil { if err != nil {
return fmt.Errorf( return fmt.Errorf(
"opening archive database %s: %w", w.path, err, "opening archive database %s: %w", w.path, err,

View File

@@ -77,14 +77,19 @@ func (c *httpCore) fireAndForget(
) { ) {
c.eng.observeAttempt(d.Target.Type, res.elapsed()) c.eng.observeAttempt(d.Target.Type, res.elapsed())
c.eng.recordResult( err := c.eng.recordResult(
webhookDB, d, 1, res.success, webhookDB, d, 1, res.success,
res.statusCode, res.respBody, res.errMsg, res.statusCode, res.respBody, res.errMsg,
res.duration, res.duration,
) )
if err != nil {
c.eng.bookkeepingFailed(d, err)
return
}
if res.success { if res.success {
c.eng.updateDeliveryStatus( c.eng.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusDelivered, database.DeliveryStatusDelivered,
) )
@@ -92,7 +97,7 @@ func (c *httpCore) fireAndForget(
return return
} }
c.eng.updateDeliveryStatus( c.eng.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusFailed, database.DeliveryStatusFailed,
) )
@@ -122,16 +127,25 @@ func (c *httpCore) withRetry(
c.eng.observeAttempt(d.Target.Type, res.elapsed()) c.eng.observeAttempt(d.Target.Type, res.elapsed())
c.eng.recordResult( err := c.eng.recordResult(
webhookDB, d, attemptNum, res.success, webhookDB, d, attemptNum, res.success,
res.statusCode, res.respBody, res.errMsg, res.statusCode, res.respBody, res.errMsg,
res.duration, res.duration,
) )
if err != nil {
// The breaker still learns the outcome: it describes the
// target's health, which is unaffected by this database's.
c.recordCircuitOutcome(cb, res.success)
c.eng.bookkeepingFailed(d, err)
return
}
if res.success { if res.success {
cb.RecordSuccess() cb.RecordSuccess()
c.eng.updateDeliveryStatus( c.eng.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusDelivered, database.DeliveryStatusDelivered,
) )
@@ -146,6 +160,20 @@ func (c *httpCore) withRetry(
) )
} }
// recordCircuitOutcome feeds one attempt's outcome to the target's
// circuit breaker.
func (c *httpCore) recordCircuitOutcome(
cb *CircuitBreaker, success bool,
) {
if success {
cb.RecordSuccess()
return
}
cb.RecordFailure()
}
func (c *httpCore) circuitBreakerBlock( func (c *httpCore) circuitBreakerBlock(
webhookDB *gorm.DB, webhookDB *gorm.DB,
d *database.Delivery, d *database.Delivery,
@@ -169,7 +197,7 @@ func (c *httpCore) circuitBreakerBlock(
"cooldown_remaining", remaining, "cooldown_remaining", remaining,
) )
c.eng.updateDeliveryStatus( c.eng.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusRetrying, database.DeliveryStatusRetrying,
) )
@@ -189,7 +217,7 @@ func (c *httpCore) handleRetry(
attemptNum int, attemptNum int,
) { ) {
if attemptNum >= maxRetries { if attemptNum >= maxRetries {
c.eng.updateDeliveryStatus( c.eng.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusFailed, database.DeliveryStatusFailed,
) )
@@ -197,7 +225,7 @@ func (c *httpCore) handleRetry(
return return
} }
c.eng.updateDeliveryStatus( c.eng.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusRetrying, database.DeliveryStatusRetrying,
) )
@@ -332,12 +360,17 @@ func (t *httpTarget) Deliver(
"error", err, "error", err,
) )
t.eng.recordResult( recErr := t.eng.recordResult(
webhookDB, d, task.AttemptNum, webhookDB, d, task.AttemptNum,
false, 0, "", err.Error(), 0, false, 0, "", err.Error(), 0,
) )
if recErr != nil {
t.eng.bookkeepingFailed(d, recErr)
t.eng.updateDeliveryStatus( return
}
t.eng.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusFailed, database.DeliveryStatusFailed,
) )

View File

@@ -55,12 +55,17 @@ func (t *logTarget) Deliver(
t.eng.observeAttempt(d.Target.Type, elapsed) t.eng.observeAttempt(d.Target.Type, elapsed)
t.eng.recordResult( err := t.eng.recordResult(
webhookDB, d, 1, true, 0, "", "", webhookDB, d, 1, true, 0, "", "",
elapsed.Milliseconds(), elapsed.Milliseconds(),
) )
if err != nil {
t.eng.bookkeepingFailed(d, err)
t.eng.updateDeliveryStatus( return
}
t.eng.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusDelivered, database.DeliveryStatusDelivered,
) )

View File

@@ -95,12 +95,17 @@ func (t *slackTarget) failConfig(
d *database.Delivery, d *database.Delivery,
err error, err error,
) { ) {
t.eng.recordResult( recErr := t.eng.recordResult(
webhookDB, d, 1, webhookDB, d, 1,
false, 0, "", err.Error(), 0, false, 0, "", err.Error(), 0,
) )
if recErr != nil {
t.eng.bookkeepingFailed(d, recErr)
t.eng.updateDeliveryStatus( return
}
t.eng.settleStatus(
webhookDB, d, d.Target.Type, webhookDB, d, d.Target.Type,
database.DeliveryStatusFailed, database.DeliveryStatusFailed,
) )

View File

@@ -92,11 +92,10 @@ func (h *Handlers) HandleEventBodyDownload() http.HandlerFunc {
// once per range. // once per range.
// //
// One consequence is worth keeping in view: the read finishes // One consequence is worth keeping in view: the read finishes
// before the client is written to, so no read lock is held for // before the client is written to, so nothing is held open for
// the length of a slow download. These per-webhook databases // the length of a slow download. Under WAL a read no longer
// run in SQLite's default journal mode rather than WAL, so a // blocks the receiver, but it does pin the WAL against
// lock held that long would block the receiver from recording // checkpointing, and a download can last minutes.
// new events.
func (h *Handlers) serveEventBody( func (h *Handlers) serveEventBody(
w http.ResponseWriter, w http.ResponseWriter,
r *http.Request, r *http.Request,

View File

@@ -143,10 +143,10 @@ func (h *Handlers) resubmitEvent(
} }
// Read before the write transaction is opened. The body can be up // Read before the write transaction is opened. The body can be up
// to the 1 MB ingest cap, and holding a read of it inside the // to the 1 MB ingest cap, and every transaction on these files
// transaction would extend how long the per-webhook database is // takes the write lock at BEGIN (_txlock=immediate, see
// locked against the receiver, which runs these files in // internal/database/sqlite_open.go), so reading inside it would
// SQLite's default journal mode rather than WAL. // hold that lock against the receiver for the length of the read.
src, found, err := loadResubmitSource( src, found, err := loadResubmitSource(
webhookDB, webhook.ID, eventID.String(), webhookDB, webhook.ID, eventID.String(),
) )

View File

@@ -0,0 +1,144 @@
package handlers_test
import (
"net/http"
"net/http/httptest"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/session"
)
// deletedMarker is the suffix the event log appends to the name
// of a target that no longer exists.
const deletedMarker = " (deleted)"
// deleteTargetThroughHandler removes a target through the real
// deletion handler, so the test soft-deletes exactly the way the
// UI does rather than by writing the timestamp itself.
func deleteTargetThroughHandler(
t *testing.T,
h *handlers.Handlers,
sess *session.Session,
webhookID, targetID string,
) {
t.Helper()
req := postRequest(
"/source/"+webhookID+"/targets/"+targetID+"/delete",
authenticatedCookies(
t, sess, deleteTestUserID, deleteTestUsername,
),
map[string]string{
paramSourceID: webhookID,
paramTargetID: targetID,
},
)
w := httptest.NewRecorder()
h.HandleTargetDelete().ServeHTTP(w, req)
require.Equal(t, http.StatusSeeOther, w.Code)
}
// TestHandleSourceLogs_NamesDeletedTarget proves a delivery
// produced by a since-deleted target still names it on the event
// log, marked as deleted.
//
// Deletes are soft and deliveries carry no foreign key to the
// target row, so the history outlives the target. Against a
// scoped lookup the delivery resolves to a zero view and the page
// renders ": delivered" with nothing saying what it was delivered
// to.
func TestHandleSourceLogs_NamesDeletedTarget(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
dbMgr *database.WebhookDBManager
)
app := newTestApp(t, &h, &sess, &db, &dbMgr)
app.RequireStart()
t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db)
tgt := seedTarget(t, db, wh.ID, database.TargetTypeLog)
seedDeliveredEvent(t, dbMgr, wh.ID, tgt.ID)
// The control: the name is on the page while the target
// lives, and is not yet marked as deleted.
before := renderSourceLogsPage(t, h, sess, wh.ID)
assert.Contains(t, before, tgt.Name)
assert.NotContains(t, before, tgt.Name+deletedMarker)
deleteTargetThroughHandler(t, h, sess, wh.ID, tgt.ID)
after := renderSourceLogsPage(t, h, sess, wh.ID)
assert.Contains(
t, after, tgt.Name+deletedMarker,
"a delivery from a deleted target must keep its name, "+
"marked as no longer existing",
)
assert.Contains(
t, after, "delivered",
"the delivery history itself must survive the delete",
)
}
// TestHandleSourceLogs_MasksDeletedTargetConfig proves that
// naming a deleted target does not widen what the page shows of
// it: its stored configuration stays masked by exactly the rules
// a live target's is.
//
// The lookup behind the name reads soft-deleted rows, so it
// carries a full target row — credential blob included — into the
// place a zero value used to sit. The projection to TargetView is
// what keeps that blob away from the template, and it must hold
// for a deleted row too.
func TestHandleSourceLogs_MasksDeletedTargetConfig(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
dbMgr *database.WebhookDBManager
)
app := newTestApp(t, &h, &sess, &db, &dbMgr)
app.RequireStart()
t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db)
tgt := seedConfiguredTarget(
t, db, wh.ID,
database.TargetTypeSlack,
`{"webhookUrl":"`+slackWebhookURL+`"}`,
)
seedDeliveredEvent(t, dbMgr, wh.ID, tgt.ID)
deleteTargetThroughHandler(t, h, sess, wh.ID, tgt.ID)
body := renderSourceLogsPage(t, h, sess, wh.ID)
assert.NotContains(t, body, slackSecretPath)
assert.NotContains(t, body, "T00000000")
assert.NotContains(t, body, "B00000000")
assert.NotContains(
t, body, "XXXXXXXXXXXXXXXXXXXXXXXX",
)
assert.NotContains(t, body, "webhookUrl")
// The name is there; only the credential is not.
assert.Contains(t, body, tgt.Name+deletedMarker)
}

View File

@@ -860,11 +860,16 @@ func (h *Handlers) HandleSourceLogs() http.HandlerFunc {
// //
// The load is Unscoped because deleting a target only soft // The load is Unscoped because deleting a target only soft
// deletes the row while its deliveries survive in the // deletes the row while its deliveries survive in the
// per-webhook database: a scoped load leaves those deliveries // per-webhook database. Both halves of the map need those rows:
// with a zero redactor, which renders their response bodies // a scoped load leaves an old delivery with a zero redactor,
// unredacted. Only the redactor half of the map is built from // which renders its response bodies unredacted, and with a zero
// deleted rows. The view half, which is what the page lists, // view, which renders its target as a blank name.
// stays scoped. //
// This map is historical display only. It is built for the event
// log page and reaches nothing but DeliveryView.Target: the
// target list on the source detail page, the edit form and the
// replay path each resolve targets themselves, and a deleted row
// is refused there as before.
func (h *Handlers) loadTargetMap( func (h *Handlers) loadTargetMap(
webhookID string, webhookID string,
) (map[string]eventLogTarget, error) { ) (map[string]eventLogTarget, error) {
@@ -880,21 +885,18 @@ func (h *Handlers) loadTargetMap(
targetMap := make( targetMap := make(
map[string]eventLogTarget, len(targets), map[string]eventLogTarget, len(targets),
) )
live := make([]database.Target, 0, len(targets))
for i := range targets { for i := range targets {
targetMap[targets[i].ID] = eventLogTarget{ targetMap[targets[i].ID] = eventLogTarget{
Redactor: delivery.NewRedactor(&targets[i]), Redactor: delivery.NewRedactor(&targets[i]),
} }
if !targets[i].DeletedAt.Valid {
live = append(live, targets[i])
}
} }
// The views come from NewTargetViews rather than being // The views come from NewTargetViews rather than being
// rebuilt here, so the masking rules stay in one place. // rebuilt here, so the masking rules stay in one place and a
for _, v := range delivery.NewTargetViews(live) { // deleted target's configuration is masked by the same code
// that masks a live one's.
for _, v := range delivery.NewTargetViews(targets) {
entry := targetMap[v.ID] entry := targetMap[v.ID]
entry.View = v entry.View = v
targetMap[v.ID] = entry targetMap[v.ID] = entry

View File

@@ -26,6 +26,10 @@ const UnmatchedRouteConst = unmatchedRoute
// inflight gauge. // inflight gauge.
const InflightHandlerConst = inflightHandler const InflightHandlerConst = inflightHandler
// UnmatchedMethodConst exposes the sentinel that stands in for a
// method the router can never route.
const UnmatchedMethodConst = unmatchedMethod
// NewLoggingResponseWriterForTest wraps newLoggingResponseWriter // NewLoggingResponseWriterForTest wraps newLoggingResponseWriter
// for use in external test packages. // for use in external test packages.
func NewLoggingResponseWriterForTest( func NewLoggingResponseWriterForTest(

View File

@@ -26,6 +26,15 @@ import (
// counting the requests in flight across the whole service. // counting the requests in flight across the whole service.
const inflightHandler = "(all)" const inflightHandler = "(all)"
// unmatchedMethod is the `method` label for a request whose method
// the router can never route.
//
// It is deliberately the same sentinel as unmatchedRoute rather than
// a spelling of its own: both stand for a client-chosen token that
// matched nothing this service registers, and giving one idea two
// spellings would read in a scrape as two different unmatched states.
const unmatchedMethod = unmatchedRoute
// routePatternID is the `handler` label for a request: the chi route // routePatternID is the `handler` label for a request: the chi route
// pattern, never the concrete path. // pattern, never the concrete path.
// //
@@ -50,9 +59,45 @@ func routePatternID(ctx context.Context) string {
return unmatchedRoute return unmatchedRoute
} }
// routePatternRecorder wraps a go-http-metrics recorder and replaces // methodID is the `method` label for a request: the request method
// the handler id on every observation with the request's route // when the router can route it, and the unmatched sentinel otherwise.
// pattern. //
// net/http accepts any RFC 9110 token as a method and hands it
// through verbatim, so the raw method is client-chosen bytes and
// bounds the label at nothing — the same unauthenticated
// series-minting the handler label carried, reached through a second
// dimension. What bounds it is the set chi's router will match a
// route for: its methodMap, which is unexported, so it is restated
// here against the net/http constants it is built from. A token
// outside that set can only ever produce chi's 405, so folding every
// one of them onto a single series loses no information a scrape
// could have used, while the nine methods that can reach a handler
// stay distinguishable.
//
// chi.RegisterMethod would extend the router's set at runtime; this
// service never calls it, and a caller that started to would have to
// extend this switch with it.
func methodID(method string) string {
switch method {
case http.MethodConnect,
http.MethodDelete,
http.MethodGet,
http.MethodHead,
http.MethodOptions,
http.MethodPatch,
http.MethodPost,
http.MethodPut,
http.MethodTrace:
return method
default:
return unmatchedMethod
}
}
// boundedLabelRecorder wraps a go-http-metrics recorder and replaces
// the request-controlled labels on every observation with bounded
// ones: the handler id becomes the request's route pattern, and the
// method becomes one the router can route.
// //
// This is the seam that makes the pattern usable at all. The metrics // This is the seam that makes the pattern usable at all. The metrics
// middleware is global (see Server.setupGlobalMiddleware), so it is // middleware is global (see Server.setupGlobalMiddleware), so it is
@@ -71,29 +116,31 @@ func routePatternID(ctx context.Context) string {
// Those never reach a handler, but chi has already matched the route // Those never reach a handler, but chi has already matched the route
// by the time the limiter runs, so their 429s land on the pattern // by the time the limiter runs, so their 429s land on the pattern
// like any other response. // like any other response.
type routePatternRecorder struct { type boundedLabelRecorder struct {
inner httpmetrics.Recorder inner httpmetrics.Recorder
} }
func (r routePatternRecorder) ObserveHTTPRequestDuration( func (r boundedLabelRecorder) ObserveHTTPRequestDuration(
ctx context.Context, ctx context.Context,
props httpmetrics.HTTPReqProperties, props httpmetrics.HTTPReqProperties,
duration time.Duration, duration time.Duration,
) { ) {
props.ID = routePatternID(ctx) props.ID = routePatternID(ctx)
props.Method = methodID(props.Method)
r.inner.ObserveHTTPRequestDuration(ctx, props, duration) r.inner.ObserveHTTPRequestDuration(ctx, props, duration)
} }
func (r routePatternRecorder) ObserveHTTPResponseSize( func (r boundedLabelRecorder) ObserveHTTPResponseSize(
ctx context.Context, ctx context.Context,
props httpmetrics.HTTPReqProperties, props httpmetrics.HTTPReqProperties,
sizeBytes int64, sizeBytes int64,
) { ) {
props.ID = routePatternID(ctx) props.ID = routePatternID(ctx)
props.Method = methodID(props.Method)
r.inner.ObserveHTTPResponseSize(ctx, props, sizeBytes) r.inner.ObserveHTTPResponseSize(ctx, props, sizeBytes)
} }
func (r routePatternRecorder) AddInflightRequests( func (r boundedLabelRecorder) AddInflightRequests(
ctx context.Context, ctx context.Context,
props httpmetrics.HTTPProperties, props httpmetrics.HTTPProperties,
quantity int, quantity int,
@@ -102,7 +149,7 @@ func (r routePatternRecorder) AddInflightRequests(
r.inner.AddInflightRequests(ctx, props, quantity) r.inner.AddInflightRequests(ctx, props, quantity)
} }
var _ httpmetrics.Recorder = routePatternRecorder{} var _ httpmetrics.Recorder = boundedLabelRecorder{}
// Metrics returns middleware that records Prometheus HTTP metrics on // Metrics returns middleware that records Prometheus HTTP metrics on
// the default registry, which is the one the /metrics route gathers. // the default registry, which is the one the /metrics route gathers.
@@ -119,14 +166,14 @@ func metricsMiddleware(
rec httpmetrics.Recorder, rec httpmetrics.Recorder,
) func(http.Handler) http.Handler { ) func(http.Handler) http.Handler {
mdlw := ghmm.New(ghmm.Config{ mdlw := ghmm.New(ghmm.Config{
Recorder: routePatternRecorder{inner: rec}, Recorder: boundedLabelRecorder{inner: rec},
}) })
return func(next http.Handler) http.Handler { return func(next http.Handler) http.Handler {
// The handler id is unmatchedRoute rather than "" so that // The handler id is unmatchedRoute rather than "" so that
// the client-chosen URL path never enters the metrics // the client-chosen URL path never enters the metrics
// pipeline at all: an empty id is the library's signal to // pipeline at all: an empty id is the library's signal to
// substitute it. routePatternRecorder overwrites this value // substitute it. boundedLabelRecorder overwrites this value
// on every observation, so it is reachable only if that // on every observation, so it is reachable only if that
// decorator is removed — in which case the metrics collapse // decorator is removed — in which case the metrics collapse
// to one series instead of leaking again. // to one series instead of leaking again.

View File

@@ -0,0 +1,309 @@
package middleware_test
import (
"net/http"
"net/http/httptest"
"strings"
"testing"
"github.com/google/uuid"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promhttp"
dto "github.com/prometheus/client_model/go"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/middleware"
)
const (
// metricsProbeMethods is how many distinct invented method tokens
// each cardinality assertion drives. The measurement on the issue
// took 300 tokens from 106 exposition lines to 7,631 — about 25
// permanent lines per token, never reclaimed — so a probe of this
// size puts a regression thousands of lines over the bound rather
// than leaving it to a rounding argument.
metricsProbeMethods = 300
// probeMethodLen is how many characters each invented method
// token carries, matching the 12 the issue measured with.
probeMethodLen = 12
// methodLabel is the label these tests are about.
methodLabel = "method"
)
// realMethods is the positive control's domain: the methods chi's
// router can match a route for, every one of which a client
// legitimately sends and every one of which must keep a series of its
// own. Bounding the label by collapsing these into one bucket would
// destroy the metric it is meant to protect.
func realMethods() []string {
return []string{
http.MethodConnect, http.MethodDelete, http.MethodGet,
http.MethodHead, http.MethodOptions, http.MethodPatch,
http.MethodPost, http.MethodPut, http.MethodTrace,
}
}
// methodProbePath returns the one receiver path a method probe
// targets. Holding the path fixed leaves the method as the only
// dimension varying, so any series growth a probe produces is the
// method label's and nothing else's.
func methodProbePath() string {
return "/webhook/" + uuid.NewString()
}
// inventedMethods returns n distinct RFC 9110 method tokens that no
// router will ever match: uppercase hex from a fresh UUID, which is
// both the shape and the length an unauthenticated flood would send.
// net/http accepts any token as a method, so every one of these
// reaches the metrics pipeline exactly as a real method does.
func inventedMethods(n int) []string {
methods := make([]string, 0, n)
for range n {
token := strings.ToUpper(
strings.ReplaceAll(uuid.NewString(), "-", ""),
)
methods = append(methods, token[:probeMethodLen])
}
return methods
}
// driveMethods sends one request per supplied method to a single
// fixed path.
func driveMethods(
t *testing.T,
h http.Handler,
path string,
methods []string,
) map[int]int {
t.Helper()
probes := make([]probe, 0, len(methods))
for _, m := range methods {
probes = append(probes, probe{method: m, path: path})
}
return drive(t, h, probes)
}
// methodLabels returns the set of distinct `method` values across
// every gathered series that carries the label at all. The inflight
// gauge does not carry it, and so contributes nothing rather than an
// empty-string member.
func methodLabels(families []*dto.MetricFamily) map[string]struct{} {
seen := make(map[string]struct{})
for _, fam := range families {
for _, m := range fam.GetMetric() {
for _, pair := range m.GetLabel() {
if pair.GetName() == methodLabel {
seen[pair.GetValue()] = struct{}{}
}
}
}
}
return seen
}
// scrapeLines renders the registry through the same promhttp handler
// /metrics is mounted on and counts the sample lines it produced.
//
// This is the quantity the issue measured and the one a Prometheus
// server pays for on every scrape: one histogram label set is a
// single gathered series but around 25 lines of exposition, which is
// why 300 method tokens cost thousands of lines rather than hundreds.
func scrapeLines(t *testing.T, reg *prometheus.Registry) int {
t.Helper()
h := promhttp.HandlerFor(reg, promhttp.HandlerOpts{})
req := httptest.NewRequestWithContext(
t.Context(), http.MethodGet, "/metrics", nil,
)
w := httptest.NewRecorder()
h.ServeHTTP(w, req)
require.Equal(t, http.StatusOK, w.Code)
lines := 0
for line := range strings.SplitSeq(w.Body.String(), "\n") {
if line == "" || strings.HasPrefix(line, "#") {
continue
}
lines++
}
return lines
}
// TestMetrics_MethodSentinelIsTheRouteSentinel pins the convention
// rather than the mechanism. An unroutable method and an unmatched
// path are the same fact — a client-chosen token matching nothing
// this service registers — so they carry one spelling. Two spellings
// would read in a scrape as two different unmatched states.
func TestMetrics_MethodSentinelIsTheRouteSentinel(t *testing.T) {
t.Parallel()
assert.Equal(
t,
middleware.UnmatchedRouteConst,
middleware.UnmatchedMethodConst,
"the unmatched sentinel must have exactly one spelling",
)
}
// TestMetrics_InventedMethodsMintOneLabelSet is the direct assertion
// the issue asks for: N requests carrying N distinct invented method
// tokens must produce exactly ONE method label. Before the fix this
// produced N of them, on an unauthenticated route with no rate
// limiter.
func TestMetrics_InventedMethodsMintOneLabelSet(t *testing.T) {
t.Parallel()
h, reg := metricsTestRouter(t, generousReceiverLimit)
methods := inventedMethods(metricsProbeMethods)
codes := driveMethods(t, h, methodProbePath(), methods)
require.Equal(
t, metricsProbeMethods, codes[http.StatusMethodNotAllowed],
"every invented token should have been unroutable",
)
labels := methodLabels(gatherMetrics(t, reg))
// Asserted on the count rather than on the set, so that a
// regression reports one number instead of dumping every token it
// minted.
distinct := len(labels)
assert.Equal(
t, 1, distinct,
"invented methods must collapse onto one label",
)
assert.Contains(
t, keys(labels), middleware.UnmatchedMethodConst,
"that one label must be the unmatched sentinel",
)
// The scrape must not republish the tokens it was driven with
// either: a label that merely looks bounded while still echoing
// client bytes is the same defect wearing a different name.
echoed := 0
for _, m := range methods {
for label := range labels {
if strings.Contains(label, m) {
echoed++
}
}
}
assert.Equal(
t, 0, echoed,
"invented method tokens reached the metrics labels",
)
}
// TestMetrics_MethodSeriesCountIsFlatUnderAFlood reproduces the
// measurement on the issue in miniature: scrape, drive several
// hundred distinct method tokens, scrape again, and require the
// second scrape to be no larger than the first. The first batch
// establishes every label set the route can produce; a flood five
// times its size must land on exactly those.
func TestMetrics_MethodSeriesCountIsFlatUnderAFlood(t *testing.T) {
t.Parallel()
h, reg := metricsTestRouter(t, generousReceiverLimit)
path := methodProbePath()
driveMethods(t, h, path, inventedMethods(metricsProbeMethods))
seededSeries := seriesCount(gatherMetrics(t, reg))
seededLines := scrapeLines(t, reg)
driveMethods(t, h, path, inventedMethods(metricsProbeMethods*4))
floodedSeries := seriesCount(gatherMetrics(t, reg))
floodedLines := scrapeLines(t, reg)
t.Logf(
"after %d invented methods: %d series, %d lines; "+
"after %d more: %d series, %d lines",
metricsProbeMethods, seededSeries, seededLines,
metricsProbeMethods*4, floodedSeries, floodedLines,
)
assert.Equal(
t, seededSeries, floodedSeries,
"a flood of invented methods must not mint series",
)
assert.Equal(
t, seededLines, floodedLines,
"a flood of invented methods must not grow the scrape",
)
}
// TestMetrics_RealMethodsStayDistinct is the positive control. The
// bound is worth nothing if it is bought by flattening the metric:
// every method the router can route must still carry a series of its
// own, one sample each, under the route pattern it was sent to.
func TestMetrics_RealMethodsStayDistinct(t *testing.T) {
t.Parallel()
h, reg := metricsTestRouter(t, generousReceiverLimit)
methods := realMethods()
codes := driveMethods(t, h, methodProbePath(), methods)
require.Equal(
t, len(methods), codes[http.StatusNotFound],
"every real method should have reached the receiver",
)
families := gatherMetrics(t, reg)
want := make(map[string]struct{}, len(methods))
for _, m := range methods {
want[m] = struct{}{}
}
assert.Equal(
t, want, methodLabels(families),
"real methods must remain distinguishable",
)
// Appearing somewhere in the scrape is not enough: each method
// must own its duration series, holding the one sample it sent.
observed := 0
for _, fam := range families {
if !strings.HasSuffix(fam.GetName(), "request_duration_seconds") {
continue
}
for _, m := range fam.GetMetric() {
observed++
assert.Equal(
t, receiverRoutePattern,
labelValue(m, "handler"),
)
assert.Equal(
t, uint64(1),
m.GetHistogram().GetSampleCount(),
"method %q shares a series",
labelValue(m, methodLabel),
)
}
}
assert.Equal(
t, len(methods), observed,
"one duration series per routable method",
)
}

View File

@@ -97,20 +97,24 @@ func metricsTestRouter(
return r, reg return r, reg
} }
// drivePaths sends one POST per supplied path and returns how many // probe is one request a cardinality assertion sends. Both label
// responses carried each status code. // dimensions that have leaked are request-controlled — the path and
func drivePaths( // the method — so both vary here and one driver sends them.
t *testing.T, type probe struct {
h http.Handler, method string
paths []string, path string
) map[int]int { }
// drive sends every probe and returns how many responses carried each
// status code.
func drive(t *testing.T, h http.Handler, probes []probe) map[int]int {
t.Helper() t.Helper()
codes := make(map[int]int) codes := make(map[int]int)
for _, p := range paths { for _, p := range probes {
req := httptest.NewRequestWithContext( req := httptest.NewRequestWithContext(
t.Context(), http.MethodPost, p, nil, t.Context(), p.method, p.path, nil,
) )
w := httptest.NewRecorder() w := httptest.NewRecorder()
h.ServeHTTP(w, req) h.ServeHTTP(w, req)
@@ -120,6 +124,25 @@ func drivePaths(
return codes return codes
} }
// drivePaths sends one POST per supplied path.
func drivePaths(
t *testing.T,
h http.Handler,
paths []string,
) map[int]int {
t.Helper()
probes := make([]probe, 0, len(paths))
for _, p := range paths {
probes = append(
probes, probe{method: http.MethodPost, path: p},
)
}
return drive(t, h, probes)
}
// receiverPaths returns n distinct /webhook/ paths, each naming a // receiverPaths returns n distinct /webhook/ paths, each naming a
// fresh UUID exactly as an unauthenticated flood would. // fresh UUID exactly as an unauthenticated flood would.
func receiverPaths(n int) []string { func receiverPaths(n int) []string {

View File

@@ -39,7 +39,7 @@
<div class="flex items-center gap-4"> <div class="flex items-center gap-4">
{{range .Deliveries}} {{range .Deliveries}}
<span class="text-xs {{if eq .Status "delivered"}}text-green-600{{else if eq .Status "failed"}}text-red-600{{else if eq .Status "retrying"}}text-yellow-600{{else}}text-gray-400{{end}}"> <span class="text-xs {{if eq .Status "delivered"}}text-green-600{{else if eq .Status "failed"}}text-red-600{{else if eq .Status "retrying"}}text-yellow-600{{else}}text-gray-400{{end}}">
{{.Target.Name}}: {{.Status}} {{.Target.DisplayName}}: {{.Status}}
</span> </span>
{{end}} {{end}}
<span class="text-xs text-gray-400">{{.CreatedAt.Format "2006-01-02 15:04:05"}}</span> <span class="text-xs text-gray-400">{{.CreatedAt.Format "2006-01-02 15:04:05"}}</span>
@@ -74,7 +74,7 @@
<div class="py-2" x-data="{ attempts: false }"> <div class="py-2" x-data="{ attempts: false }">
<div class="flex items-center justify-between cursor-pointer" @click="attempts = !attempts"> <div class="flex items-center justify-between cursor-pointer" @click="attempts = !attempts">
<div class="flex items-center gap-3"> <div class="flex items-center gap-3">
<span class="text-sm text-gray-700">{{.Target.Name}}</span> <span class="text-sm text-gray-700">{{.Target.DisplayName}}</span>
<span class="text-xs {{if eq .Status "delivered"}}text-green-600{{else if eq .Status "failed"}}text-red-600{{else if eq .Status "retrying"}}text-yellow-600{{else}}text-gray-400{{end}}">{{.Status}}</span> <span class="text-xs {{if eq .Status "delivered"}}text-green-600{{else if eq .Status "failed"}}text-red-600{{else if eq .Status "retrying"}}text-yellow-600{{else}}text-gray-400{{end}}">{{.Status}}</span>
</div> </div>
<div class="flex items-center gap-3"> <div class="flex items-center gap-3">