Author SHA1 Message Date
clawbot 6be031594f Add a Download button that exports a database target's archive as gzipped JSON (closes #374)
check / check (push) Waiting to run
Each database target on the webhook page links to
/hook/ID/targets/TARGETID/download, which streams the target's archive
as archive-WEBHOOK-TARGET-YYYYMMDDTHHMMSSZ.json.gz: the webhook, the
target, exported_at, and archived_events, one object per row keyed by
column, a body that is not valid UTF-8 in base64 with body_encoding.
The export reads through one cursor inside a read-only transaction:
one snapshot, no write lock. A download runs for as long as the client
keeps reading, and one that fails after it has started aborts the
connection. The request limit answers 504 only when the handler has
not started its response, so a long download is logged as its 200.

Model: opus-5-5
2026-10-02 16:04:11 +00:00
clawbot bf3df0312b Clear the environment in the cmd/webhooker tests that build a Config (closes #452)
check / check (push) Waiting to run
The two cmd/webhooker tests that build the app graph, TestNewApp_StopTimeout and TestNewApp_SendsFxEventsToTheLogger, built a Config without clearing the environment, so a variable exported in the developer's shell changed their result: a METRICS_USERNAME without METRICS_PASSWORD failed the second. Both now call config.ClearEnvForTest before setting their own variables, as the config tests and the first-boot test already do. No other test outside internal/config builds a Config through config.New. Test change only.

Model: opus-5-5
2026-10-02 18:03:31 +02:00
clawbot 503c57efd9 Assert the body a retry delivers, not only its status (closes #294)
check / check (push) Waiting to run
TestProcessRetryTask_LargeBody_FetchFromDB and TestProcessRetryTask_SuccessfulRetry checked only that the delivery ended delivered, so deleting the event-body fetch on the retry path, the behaviour the first is named for, left both green while a retry could deliver an empty or truncated body. Both now compare the body the target received with the stored event body byte for byte, and both fail when that fetch is deleted. The other retry-path tests are not about the body and are unchanged. Test change only.

Model: opus-5-5
2026-10-02 18:03:20 +02:00
clawbot d084f4f912 Pin the body cap's order in every page route group (closes #93)
check / check (push) Waiting to run
Follow-ups from an August review of the body cap, each checked against the current tree. One route test now requires an oversized POST, with no session and no CSRF token, to be refused with 413 before CSRF runs, in every page route group with a POST route and in /settings/, so moving a group's body cap after CSRF fails it. The three test router helpers build the server through New with a lifecycle that is never started, so no field is set by hand. The middleware test comment names runMaxBodySize, and the MaxBodySize doc comment says methods other than POST, PUT and PATCH pass uncapped on purpose. The README item was already settled.

Model: opus-5-5
2026-10-02 17:53:21 +02:00
clawbot e67fffb05d Check the log charge against every code point (closes #172)
check / check (push) Waiting to run
The access log's 2,560-byte line ceiling holds only if logfield.EncodedBytes charges each code point at least what the log handlers write for it, and the test checked that on a sample. TestEncodedBytes_ChargesAtLeastWhatTheHandlersEmit now covers every code point, surrogates aside, for both handlers. Below U+1000, where the handlers' escaping varies, each code point is measured alone, both in a bare value and in a quoted one, so undercharging any of them, DEL included, fails and names it. From U+1000 up it compares batch sums, which the doc comment says can hide one JSON-only overcharge. Reverting the astral charge to 6 fails the test. It adds under 2 seconds under -race.

Model: opus-5-5
2026-10-02 17:53:10 +02:00
clawbot 4958a6f2e4 Isolate config tests from the shell; any out-of-range PORT is ErrInvalidPort (closes #94)
check / check (push) Waiting to run
The config tests unset variables without restoring them and read whatever the developer's shell exported, so results could differ from one machine to the next. config.ClearEnvForTest, in internal/config/testing.go, unsets every variable in the process environment and, when the test ends, leaves it exactly as it found it; every config test and the first-boot test call it before setting their own. TestEnvPositiveInt and TestEnvPort share one table runner, and the port errors name the bad value. A PORT of zero, below zero, or too large to parse now matches ErrInvalidPort, as one above 65535 already did. The README and the Settings page say an unparseable or non-positive RETENTION_SWEEP_INTERVAL stops startup.

Model: opus-5-5
2026-10-02 17:44:30 +02:00
clawbot 385fbc1a6a Send fx's own events through the service's logger (closes #183)
check / check (push) Waiting to run
fx printed its dependency graph and lifecycle hooks through its own console logger on standard error, so an operator shipping the JSON log to a collector got a second shape on a second stream for every start. The production app now passes fx.WithLogger with a small FxLogger in internal/logger that writes fx's events through the service's logger: graph events at debug, lifecycle at info, failures at error. Its constructor takes the configuration, so DEBUG=true applies before fx replays the events it held back. go.uber.org/fx moves from v1.20.1 to v1.24.0. Tests keep fx.NopLogger. The README says a failure before the logger exists, and the Go runtime's own output, still go to standard error as plain text.

Model: opus-5-5
2026-10-02 17:22:05 +02:00
40 changed files with 2118 additions and 404 deletions
+47 -16
View File
@@ -139,7 +139,7 @@ TTY detection, and security headers are always applied.
| `METRICS_USERNAME` | Basic auth username for `/metrics`. Must be set together with `METRICS_PASSWORD`; one without the other fails startup | `""` | | `METRICS_USERNAME` | Basic auth username for `/metrics`. Must be set together with `METRICS_PASSWORD`; one without the other fails startup | `""` |
| `METRICS_PASSWORD` | Basic auth password for `/metrics`. Must be set together with `METRICS_USERNAME`; one without the other fails startup | `""` | | `METRICS_PASSWORD` | Basic auth password for `/metrics`. Must be set together with `METRICS_USERNAME`; one without the other fails startup | `""` |
| `SENTRY_DSN` | Sentry error reporting DSN. Unset leaves error reporting off; a value the Sentry SDK cannot parse fails startup rather than serving with reporting silently off | `""` | | `SENTRY_DSN` | Sentry error reporting DSN. Unset leaves error reporting off; a value the Sentry SDK cannot parse fails startup rather than serving with reporting silently off | `""` |
| `RETENTION_SWEEP_INTERVAL` | How often the retention reaper and archive sweeper run (Go duration, must be positive) | `1h` | | `RETENTION_SWEEP_INTERVAL` | How often the retention reaper and archive sweeper run (Go duration, must be positive). A value that does not parse, or is zero or negative, fails startup | `1h` |
| `SESSION_IDLE_TIMEOUT` | Idle session timeout (Go duration) | `24h` | | `SESSION_IDLE_TIMEOUT` | Idle session timeout (Go duration) | `24h` |
| `RECEIVER_RATE_LIMIT` | Receiver requests/minute per IP per entrypoint (10x that per IP across the route) | `120` | | `RECEIVER_RATE_LIMIT` | Receiver requests/minute per IP per entrypoint (10x that per IP across the route) | `120` |
| `TRUSTED_PROXIES` | CIDRs whose forwarded headers are trusted. A set value replaces the default. If any client can reach webhooker, or the proxy in front of it, from an RFC 1918 source address, set it to the proxy's address alone. See [Trusted proxies](#trusted-proxies) | `10.0.0.0/8,172.16.0.0/12,192.168.0.0/16` (RFC 1918) | | `TRUSTED_PROXIES` | CIDRs whose forwarded headers are trusted. A set value replaces the default. If any client can reach webhooker, or the proxy in front of it, from an RFC 1918 source address, set it to the proxy's address alone. See [Trusted proxies](#trusted-proxies) | `10.0.0.0/8,172.16.0.0/12,192.168.0.0/16` (RFC 1918) |
@@ -557,8 +557,8 @@ If it is lost, run `webhooker resetpw admin` on a stopped deployment.
``` ```
It is a banner rather than a log line because that is the only time it It is a banner rather than a log line because that is the only time it
is ever shown: as one `INFO` record it sat among the roughly 45 fx is ever shown: as one `INFO` record it would sit among the records fx
`PROVIDE`/`RUN`/`HOOK` lines a boot writes, and under `docker run -d` writes as each start hook runs, and under `docker run -d`
it is one line in a log subject to rotation. The database stores only it is one line in a log subject to rotation. The database stores only
its Argon2id hash. There is no second account and no forgot-password its Argon2id hash. There is no second account and no forgot-password
flow, so the banner and the reset command below are the only two ways flow, so the banner and the reset command below are the only two ways
@@ -628,7 +628,8 @@ Changing a password you still know needs none of this — use
`DEBUG=true` lowers the log level to `DEBUG`, which turns on every `DEBUG=true` lowers the log level to `DEBUG`, which turns on every
statement GORM runs, the two by-design lookup misses on the statement GORM runs, the two by-design lookup misses on the
unauthenticated routes, and the rate limiter's own rejections. It is unauthenticated routes, the rate limiter's own rejections, and fx's
records of building the dependency graph at startup. It is
meant to be safe to turn on while diagnosing a live service and safe to meant to be safe to turn on while diagnosing a live service and safe to
paste the output of into a bug report. paste the output of into a bug report.
@@ -2029,6 +2030,29 @@ Because each `database` target has its own archive file, a target's
webhook with different expiries keep two archives, each pruned on its webhook with different expiries keep two archives, each pruned on its
own schedule. own schedule.
Each `database` target on the webhook page has a **Download** button,
which returns its archive as one gzipped JSON file,
`archive-{webhook_name}-{target_name}-{YYYYMMDDTHHMMSSZ}.json.gz`, the
names made safe as above and the time in UTC. The file holds one
object: `webhook` and `target`, each an `id` and a `name`;
`exported_at`; and `archived_events`, one object per archived row with
every column, keyed by column name. A body that is not valid UTF-8 is
written in base64, with `"body_encoding": "base64"` beside it. An
archive that does not exist yet, or was moved away, downloads with an
empty `archived_events`; the download never creates the file.
The download streams: each row is read and written out compressed
before the next is read, so neither the archive nor the JSON is held in
memory. It reads on a connection of its own, inside one read-only
transaction, so the file holds the archive as it stood when the
download started, and archive writes go on meanwhile, since under WAL a
reader never blocks a writer. While it runs, the `-wal` cannot be
checkpointed past what it reads, so a long download lets the `-wal`
grow. It finds the file by the stored names under the lock that webhook
edits, target edits and target creation hold, and lets go once the file
is open: a rename during the download moves the file without affecting
it.
Deleting a webhook releases its archives: the delivery engine's cached Deleting a webhook releases its archives: the delivery engine's cached
archive writers are dropped and their file handles closed, so nothing archive writers are dropped and their file handles closed, so nothing
lingers after the webhook is gone. The archive **files themselves are lingers after the webhook is gone. The archive **files themselves are
@@ -2637,16 +2661,20 @@ read as more than it is:
that type on a specific webhook, and each line it writes is bounded that type on a specific webhook, and each line it writes is bounded
per event by the 1 MB receiver body cap. Adding one is a decision to per event by the 1 MB receiver body cap. Adding one is a decision to
spend log volume on that webhook's payloads. spend log volume on that webhook's payloads.
- **Two writers that do not go through `internal/logger` at all**, both - **The Go runtime**, which does not go through `internal/logger`. The
on standard error. `fx` prints the dependency graph and the lifecycle runtime writes an unrecovered panic or a fatal error itself, as plain
hooks through its default console logger at startup and shutdown — text on standard error, and that output cannot be redirected. A panic
nothing calls `fx.WithLogger`, and `fx.New` builds that logger over in a background worker rather than in a request handler is the case
`os.Stderr`. The Go runtime writes a panic or a fatal error itself; a that reaches it, since nothing recovers those. It carries no
panic in a background worker rather than in a request handler is the client-chosen value at a client-chosen length: the service's own
case that reaches it, since nothing recovers those. Neither carries a `panic` calls are invariant guards over constants and over
client-chosen value at a client-chosen length: the five `panic` calls `crypto/rand`, apart from the one that hands `http.ErrAbortHandler`
in this service are invariant guards over constants and over back to `net/http`, described below.
`crypto/rand`. - **A failure before fx's logger is built**, such as an invalid
configuration value. fx's logger takes the configuration, so when
that fails fx's own console logger still prints the failure as plain
text on standard error. Its values come from the operator's
environment, not from a client.
- **`net/http`'s own faults**, which are _not_ a separate writer. - **`net/http`'s own faults**, which are _not_ a separate writer.
`internal/server/http.go` builds its server with a nil `ErrorLog`, so `internal/server/http.go` builds its server with a nil `ErrorLog`, so
`net/http` falls back to the `log` package's default logger — and `net/http` falls back to the `log` package's default logger — and
@@ -2892,6 +2920,7 @@ returns to the page that was asked for.
| `POST` | `/hook/{id}/targets` | Add target to webhook | | `POST` | `/hook/{id}/targets` | Add target to webhook |
| `GET` | `/hook/{id}/targets/{targetID}/edit` | Edit target form. The one page that renders a target's destination URL and header values in full, rather than masked | | `GET` | `/hook/{id}/targets/{targetID}/edit` | Edit target form. The one page that renders a target's destination URL and header values in full, rather than masked |
| `POST` | `/hook/{id}/targets/{targetID}/edit` | Edit target submission | | `POST` | `/hook/{id}/targets/{targetID}/edit` | Edit target submission |
| `GET` | `/hook/{id}/targets/{targetID}/download` | Download a `database` target's archive as one gzipped JSON file. See [Database Architecture](#database-architecture) |
| `POST` | `/hook/{id}/targets/{targetID}/delete` | Delete a target | | `POST` | `/hook/{id}/targets/{targetID}/delete` | Delete a target |
| `POST` | `/hook/{id}/targets/{targetID}/toggle` | Enable or disable a target | | `POST` | `/hook/{id}/targets/{targetID}/toggle` | Enable or disable a target |
@@ -2937,7 +2966,8 @@ webhooker/
│ ├── resetpw/ │ ├── resetpw/
│ │ └── resetpw.go # `webhooker resetpw`: set an account's password, stopped deployments only │ │ └── resetpw.go # `webhooker resetpw`: set an account's password, stopped deployments only
│ ├── config/ │ ├── config/
│ │ └── config.go # Configuration loading from environment variables │ │ ├── config.go # Configuration loading from environment variables
│ │ └── testing.go # ClearEnvForTest: an empty environment for one test
│ ├── database/ │ ├── database/
│ │ ├── base_model.go # BaseModel with UUID primary keys │ │ ├── base_model.go # BaseModel with UUID primary keys
│ │ ├── database.go # GORM connection, migrations, admin seed │ │ ├── database.go # GORM connection, migrations, admin seed
@@ -2972,6 +3002,7 @@ webhooker/
│ │ ├── target_slack.go # Slack/Mattermost incoming-webhook target │ │ ├── target_slack.go # Slack/Mattermost incoming-webhook target
│ │ ├── target_database.go # Database archive target │ │ ├── target_database.go # Database archive target
│ │ ├── target_database_archive.go # Archive file lifecycle and pruning │ │ ├── target_database_archive.go # Archive file lifecycle and pruning
│ │ ├── target_database_export.go # Archive download as gzipped JSON
│ │ ├── target_log.go # Log target (stdout) │ │ ├── target_log.go # Log target (stdout)
│ │ ├── target_config_view.go # Masked target config for templates │ │ ├── target_config_view.go # Masked target config for templates
│ │ ├── archive_sweeper.go # Periodic pruning of idle archives │ │ ├── archive_sweeper.go # Periodic pruning of idle archives
@@ -2996,7 +3027,7 @@ webhooker/
│ ├── lifecycle/ │ ├── lifecycle/
│ │ └── lifecycle.go # Shared stop-hook waiter, bounded by the stop context │ │ └── lifecycle.go # Shared stop-hook waiter, bounded by the stop context
│ ├── logger/ │ ├── logger/
│ │ └── logger.go # slog setup with TTY detection │ │ └── logger.go # slog setup with TTY detection; fx's event logger
│ ├── metrics/ │ ├── metrics/
│ │ └── metrics.go # Delivery Prometheus collectors, labelled by target type │ │ └── metrics.go # Delivery Prometheus collectors, labelled by target type
│ ├── middleware/ │ ├── middleware/
+14
View File
@@ -8,6 +8,7 @@ import (
"time" "time"
"go.uber.org/fx" "go.uber.org/fx"
"go.uber.org/fx/fxevent"
"sneak.berlin/go/webhooker/internal/config" "sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/datadir" "sneak.berlin/go/webhooker/internal/datadir"
@@ -168,6 +169,19 @@ func run(stderr io.Writer) int {
func newApp() *fx.App { func newApp() *fx.App {
return fx.New( return fx.New(
fx.StopTimeout(stopTimeout), fx.StopTimeout(stopTimeout),
// fx's own events go through the service's logger, not fx's
// console logger on standard error. The exception is a failure
// before this logger is built, such as an invalid configuration
// value, which fx's console logger still prints there. fx holds
// its events back until this logger is built and then replays
// them, so it takes the configuration, which sets the level
// DEBUG=true asks for: without it the replay would run at INFO
// and drop every record of how the graph was built.
fx.WithLogger(
func(l *logger.Logger, _ *config.Config) fxevent.Logger {
return logger.NewFxLogger(l.Get())
},
),
fx.Provide( fx.Provide(
globals.New, globals.New,
logger.New, logger.New,
+102
View File
@@ -2,12 +2,19 @@ package main
import ( import (
"bytes" "bytes"
"encoding/json"
"io"
"log/slog"
"net"
"os"
"strconv"
"strings" "strings"
"testing" "testing"
"time" "time"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/datadir" "sneak.berlin/go/webhooker/internal/datadir"
"sneak.berlin/go/webhooker/internal/resetpw" "sneak.berlin/go/webhooker/internal/resetpw"
"sneak.berlin/go/webhooker/internal/server" "sneak.berlin/go/webhooker/internal/server"
@@ -30,6 +37,7 @@ const dockerStopGrace = 10 * time.Second
// fx.New applies options before it executes invokes, so the timeout // fx.New applies options before it executes invokes, so the timeout
// is set whether or not the graph itself can be constructed here. // is set whether or not the graph itself can be constructed here.
func TestNewApp_StopTimeout(t *testing.T) { func TestNewApp_StopTimeout(t *testing.T) {
config.ClearEnvForTest(t)
t.Setenv("DATA_DIR", t.TempDir()) t.Setenv("DATA_DIR", t.TempDir())
got := newApp().StopTimeout() got := newApp().StopTimeout()
@@ -38,6 +46,100 @@ func TestNewApp_StopTimeout(t *testing.T) {
require.Less(t, got, dockerStopGrace) require.Less(t, got, dockerStopGrace)
} }
// freePort returns a loopback TCP port that was free a moment ago, by
// taking one and releasing it.
func freePort(t *testing.T) int {
t.Helper()
var listenCfg net.ListenConfig
l, err := listenCfg.Listen(t.Context(), "tcp", "127.0.0.1:0")
require.NoError(t, err)
addr, ok := l.Addr().(*net.TCPAddr)
require.True(t, ok, "listener is not TCP")
require.NoError(t, l.Close())
return addr.Port
}
// TestNewApp_SendsFxEventsToTheLogger starts and stops the app main
// runs, with DEBUG=true, and reads back what reached the service's
// logger. fx's own events must arrive there as structured records:
// the start at INFO, and at DEBUG the records of how the graph was
// built.
//
// fx holds its events back until its logger is built and then replays
// them all at once, so the earliest of them arriving shows the replay
// ran at DEBUG: that globals.New was provided, which fx records before
// anything is built, and the run of logger.New, which happens before
// the configuration sets the level.
func TestNewApp_SendsFxEventsToTheLogger(t *testing.T) {
config.ClearEnvForTest(t)
t.Setenv("DATA_DIR", t.TempDir())
t.Setenv("PORT", strconv.Itoa(freePort(t)))
t.Setenv("DEBUG", "true")
// internal/logger writes to whatever os.Stdout is when it builds
// its handler. A file is not a terminal, so that handler is the
// JSON one the service uses in production.
out, err := os.CreateTemp(t.TempDir(), "stdout")
require.NoError(t, err)
stdout := os.Stdout
os.Stdout = out
t.Cleanup(func() {
os.Stdout = stdout
_ = out.Close()
})
app := newApp()
require.NoError(t, app.Start(t.Context()))
require.NoError(t, app.Stop(t.Context()))
_, err = out.Seek(0, io.SeekStart)
require.NoError(t, err)
written, err := io.ReadAll(out)
require.NoError(t, err)
type record struct {
Level string `json:"level"`
Msg string `json:"msg"`
Name string `json:"name"`
Constructor string `json:"constructor"`
}
var records []record
for line := range strings.Lines(string(written)) {
var r record
// The first-boot banner is plain text, not a record.
if json.Unmarshal([]byte(line), &r) == nil {
records = append(records, r)
}
}
const pkg = "sneak.berlin/go/webhooker/internal/"
info := slog.LevelInfo.String()
debug := slog.LevelDebug.String()
assert.Contains(t, records, record{Level: info, Msg: "started"})
assert.Contains(t, records, record{
Level: debug, Msg: "provided", Constructor: pkg + "globals.New()",
})
assert.Contains(t, records, record{
Level: debug, Msg: "run", Name: pkg + "logger.New()",
})
assert.Contains(t, records, record{Level: debug, Msg: "invoking"})
assert.Contains(t, records, record{
Level: debug, Msg: "initialized custom fxevent.Logger",
})
}
// TestRunRefusesLockedDataDir pins what an operator's second start // TestRunRefusesLockedDataDir pins what an operator's second start
// does. The entry point must refuse before it builds the fx graph — // does. The entry point must refuse before it builds the fx graph —
// nothing may open a database in a DATA_DIR another process holds — // nothing may open a database in a DATA_DIR another process holds —
+4 -6
View File
@@ -20,7 +20,7 @@ require (
github.com/prometheus/client_model v0.5.0 github.com/prometheus/client_model v0.5.0
github.com/slok/go-http-metrics v0.11.0 github.com/slok/go-http-metrics v0.11.0
github.com/stretchr/testify v1.11.1 github.com/stretchr/testify v1.11.1
go.uber.org/fx v1.20.1 go.uber.org/fx v1.24.0
golang.org/x/crypto v0.38.0 golang.org/x/crypto v0.38.0
gopkg.in/yaml.v3 v3.0.1 gopkg.in/yaml.v3 v3.0.1
gorm.io/driver/sqlite v1.5.4 gorm.io/driver/sqlite v1.5.4
@@ -42,7 +42,6 @@ require (
github.com/jinzhu/now v1.1.5 // indirect github.com/jinzhu/now v1.1.5 // indirect
github.com/kballard/go-shellquote v0.0.0-20180428030007-95032a82bc51 // indirect github.com/kballard/go-shellquote v0.0.0-20180428030007-95032a82bc51 // indirect
github.com/klauspost/cpuid/v2 v2.2.10 // indirect github.com/klauspost/cpuid/v2 v2.2.10 // indirect
github.com/kr/text v0.2.0 // indirect
github.com/mattn/go-isatty v0.0.20 // indirect github.com/mattn/go-isatty v0.0.20 // indirect
github.com/mattn/go-sqlite3 v1.14.17 // indirect github.com/mattn/go-sqlite3 v1.14.17 // indirect
github.com/matttproud/golang_protobuf_extensions/v2 v2.0.0 // indirect github.com/matttproud/golang_protobuf_extensions/v2 v2.0.0 // indirect
@@ -51,10 +50,9 @@ require (
github.com/prometheus/procfs v0.12.0 // indirect github.com/prometheus/procfs v0.12.0 // indirect
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
github.com/zeebo/xxh3 v1.0.2 // indirect github.com/zeebo/xxh3 v1.0.2 // indirect
go.uber.org/atomic v1.9.0 // indirect go.uber.org/dig v1.19.0 // indirect
go.uber.org/dig v1.17.0 // indirect go.uber.org/multierr v1.10.0 // indirect
go.uber.org/multierr v1.9.0 // indirect go.uber.org/zap v1.26.0 // indirect
go.uber.org/zap v1.23.0 // indirect
golang.org/x/mod v0.17.0 // indirect golang.org/x/mod v0.17.0 // indirect
golang.org/x/sync v0.14.0 // indirect golang.org/x/sync v0.14.0 // indirect
golang.org/x/sys v0.47.0 // indirect golang.org/x/sys v0.47.0 // indirect
+10 -20
View File
@@ -1,7 +1,5 @@
github.com/99designs/basicauth-go v0.0.0-20230316000542-bf6f9cbbf0f8 h1:nMpu1t4amK3vJWBibQ5X/Nv0aXL+b69TQf2uK5PH7Go= github.com/99designs/basicauth-go v0.0.0-20230316000542-bf6f9cbbf0f8 h1:nMpu1t4amK3vJWBibQ5X/Nv0aXL+b69TQf2uK5PH7Go=
github.com/99designs/basicauth-go v0.0.0-20230316000542-bf6f9cbbf0f8/go.mod h1:3cARGAK9CfW3HoxCy1a0G4TKrdiKke8ftOMEOHyySYs= github.com/99designs/basicauth-go v0.0.0-20230316000542-bf6f9cbbf0f8/go.mod h1:3cARGAK9CfW3HoxCy1a0G4TKrdiKke8ftOMEOHyySYs=
github.com/benbjohnson/clock v1.3.0 h1:ip6w0uFQkncKQ979AypyG0ER7mqUSBdKLOgAle/AT8A=
github.com/benbjohnson/clock v1.3.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA=
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw=
github.com/cespare/xxhash/v2 v2.2.0 h1:DC2CZ1Ep5Y4k3ZQ899DldepgrayRUGE6BBZ/cd9Cj44= github.com/cespare/xxhash/v2 v2.2.0 h1:DC2CZ1Ep5Y4k3ZQ899DldepgrayRUGE6BBZ/cd9Cj44=
@@ -12,9 +10,6 @@ github.com/chromedp/chromedp v0.16.0 h1:rOO4deOm4CbZgBCa8mD9g2rDyIoNs0BkgvNrlbp5
github.com/chromedp/chromedp v0.16.0/go.mod h1:rbuGKFT1vMcFcFqKfPIO1GpX/N+2s8onm2qMxZLbU5U= github.com/chromedp/chromedp v0.16.0/go.mod h1:rbuGKFT1vMcFcFqKfPIO1GpX/N+2s8onm2qMxZLbU5U=
github.com/chromedp/sysutil v1.1.0 h1:PUFNv5EcprjqXZD9nJb9b/c9ibAbxiYo4exNWZyipwM= github.com/chromedp/sysutil v1.1.0 h1:PUFNv5EcprjqXZD9nJb9b/c9ibAbxiYo4exNWZyipwM=
github.com/chromedp/sysutil v1.1.0/go.mod h1:WiThHUdltqCNKGc4gaU50XgYjwjYIhKWoHGPTUfWTJ8= github.com/chromedp/sysutil v1.1.0/go.mod h1:WiThHUdltqCNKGc4gaU50XgYjwjYIhKWoHGPTUfWTJ8=
github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM= github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY=
@@ -83,7 +78,6 @@ github.com/pingcap/errors v0.11.4 h1:lFuQV/oaUMGcD2tqt+01ROSmJs75VG1ToEOkZIZ4nE4
github.com/pingcap/errors v0.11.4/go.mod h1:Oi8TUi2kEtXXLMJk9l1cGmz20kV3TaQ0usTwv5KuLY8= github.com/pingcap/errors v0.11.4/go.mod h1:Oi8TUi2kEtXXLMJk9l1cGmz20kV3TaQ0usTwv5KuLY8=
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U=
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/prometheus/client_golang v1.18.0 h1:HzFfmkOzH5Q8L8G+kSJKUx5dtG87sewO+FoDDqP5Tbk= github.com/prometheus/client_golang v1.18.0 h1:HzFfmkOzH5Q8L8G+kSJKUx5dtG87sewO+FoDDqP5Tbk=
@@ -100,28 +94,24 @@ github.com/rogpeppe/go-internal v1.10.0 h1:TMyTOH3F/DB16zRVcYyreMH6GnZZrwQVAoYjR
github.com/rogpeppe/go-internal v1.10.0/go.mod h1:UQnix2H7Ngw/k4C5ijL5+65zddjncjaFoBhdsK/akog= github.com/rogpeppe/go-internal v1.10.0/go.mod h1:UQnix2H7Ngw/k4C5ijL5+65zddjncjaFoBhdsK/akog=
github.com/slok/go-http-metrics v0.11.0 h1:ABJUpekCZSkQT1wQrFvS4kGbhea/w6ndFJaWJeh3zL0= github.com/slok/go-http-metrics v0.11.0 h1:ABJUpekCZSkQT1wQrFvS4kGbhea/w6ndFJaWJeh3zL0=
github.com/slok/go-http-metrics v0.11.0/go.mod h1:ZGKeYG1ET6TEJpQx18BqAJAvxw9jBAZXCHU7bWQqqAc= github.com/slok/go-http-metrics v0.11.0/go.mod h1:ZGKeYG1ET6TEJpQx18BqAJAvxw9jBAZXCHU7bWQqqAc=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/objx v0.5.2 h1:xuMeJ0Sdp5ZMRXx/aWO6RZxdr3beISkG5/G/aIRr3pY= github.com/stretchr/objx v0.5.2 h1:xuMeJ0Sdp5ZMRXx/aWO6RZxdr3beISkG5/G/aIRr3pY=
github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA= github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA=
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ= github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ=
github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0= github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0=
github.com/zeebo/xxh3 v1.0.2 h1:xZmwmqxHZA8AI603jOQ0tMqmBr9lPeFwGg6d+xy9DC0= github.com/zeebo/xxh3 v1.0.2 h1:xZmwmqxHZA8AI603jOQ0tMqmBr9lPeFwGg6d+xy9DC0=
github.com/zeebo/xxh3 v1.0.2/go.mod h1:5NWz9Sef7zIDm2JHfFlcQvNekmcEl9ekUZQQKCYaDcA= github.com/zeebo/xxh3 v1.0.2/go.mod h1:5NWz9Sef7zIDm2JHfFlcQvNekmcEl9ekUZQQKCYaDcA=
go.uber.org/atomic v1.9.0 h1:ECmE8Bn/WFTYwEW/bpKD3M8VtR/zQVbavAoalC1PYyE= go.uber.org/dig v1.19.0 h1:BACLhebsYdpQ7IROQ1AGPjrXcP5dF80U3gKoFzbaq/4=
go.uber.org/atomic v1.9.0/go.mod h1:fEN4uk6kAWBTFdckzkM89CLk9XfWZrxpCo0nPH17wJc= go.uber.org/dig v1.19.0/go.mod h1:Us0rSJiThwCv2GteUN0Q7OKvU7n5J4dxZ9JKUXozFdE=
go.uber.org/dig v1.17.0 h1:5Chju+tUvcC+N7N6EV08BJz41UZuO3BmHcN4A287ZLI= go.uber.org/fx v1.24.0 h1:wE8mruvpg2kiiL1Vqd0CC+tr0/24XIB10Iwp2lLWzkg=
go.uber.org/dig v1.17.0/go.mod h1:rTxpf7l5I0eBTlE6/9RL+lDybC7WFwY2QH55ZSjy1mU= go.uber.org/fx v1.24.0/go.mod h1:AmDeGyS+ZARGKM4tlH4FY2Jr63VjbEDJHtqXTGP5hbo=
go.uber.org/fx v1.20.1 h1:zVwVQGS8zYvhh9Xxcu4w1M6ESyeMzebzj2NbSayZ4Mk= go.uber.org/goleak v1.2.0 h1:xqgm/S+aQvhWFTtR0XK3Jvg7z8kGV8P4X14IzwN3Eqk=
go.uber.org/fx v1.20.1/go.mod h1:iSYNbHf2y55acNCwCXKx7LbWb5WG1Bnue5RDXz1OREg= go.uber.org/goleak v1.2.0/go.mod h1:XJYK+MuIchqpmGmUSAzotztawfKvYLUIgg7guXrwVUo=
go.uber.org/goleak v1.1.11 h1:wy28qYRKZgnJTxGxvye5/wgWr1EKjmUDGYox5mGlRlI= go.uber.org/multierr v1.10.0 h1:S0h4aNzvfcFsC3dRF1jLoaov7oRaKqRGC/pUEJ2yvPQ=
go.uber.org/goleak v1.1.11/go.mod h1:cwTWslyiVhfpKIDGSZEM2HlOvcqm+tG4zioyIeLoqMQ= go.uber.org/multierr v1.10.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y=
go.uber.org/multierr v1.9.0 h1:7fIwc/ZtS0q++VgcfqFDxSBZVv/Xo49/SYnDFupUwlI= go.uber.org/zap v1.26.0 h1:sI7k6L95XOKS281NhVKOFCUNIvv9e0w4BF8N3u+tCRo=
go.uber.org/multierr v1.9.0/go.mod h1:X2jQV1h+kxSjClGpnseKVIxpmcjrj7MNnI0bnlfKTVQ= go.uber.org/zap v1.26.0/go.mod h1:dtElttAiwGvoJ/vj4IwHBS/gXsEu/pZ50mUIRWuG0so=
go.uber.org/zap v1.23.0 h1:OjGQ5KQDEUawVHxNwQgPpiypGHOxo2mNZsOqTak4fFY=
go.uber.org/zap v1.23.0/go.mod h1:D+nX8jyLsMHMYrln8A0rJjFt/T/9/bGgIhAqxv5URuY=
golang.org/x/crypto v0.38.0 h1:jt+WWG8IZlBnVbomuhg2Mdq0+BBQaHbtqHEFEigjUV8= golang.org/x/crypto v0.38.0 h1:jt+WWG8IZlBnVbomuhg2Mdq0+BBQaHbtqHEFEigjUV8=
golang.org/x/crypto v0.38.0/go.mod h1:MvrbAqul58NNYPKnOra203SB9vpuZW0e+RRZV+Ggqjw= golang.org/x/crypto v0.38.0/go.mod h1:MvrbAqul58NNYPKnOra203SB9vpuZW0e+RRZV+Ggqjw=
golang.org/x/mod v0.17.0 h1:zY54UmvipHiNd+pm+m0x9KhZ9hl1/7QNMyxXbc6ICqA= golang.org/x/mod v0.17.0 h1:zY54UmvipHiNd+pm+m0x9KhZ9hl1/7QNMyxXbc6ICqA=
+19 -10
View File
@@ -80,8 +80,7 @@ const (
// process over a Docker network or a private LAN connects from. // process over a Docker network or a private LAN connects from.
defaultTrustedProxies = "10.0.0.0/8,172.16.0.0/12,192.168.0.0/16" defaultTrustedProxies = "10.0.0.0/8,172.16.0.0/12,192.168.0.0/16"
// maxPort is the highest valid TCP port number. The lower // maxPort is the highest valid TCP port number.
// bound (at least 1) is enforced by envPositiveInt.
maxPort = 65535 maxPort = 65535
// mappedV4Offset is the number of leading bits an IPv4-mapped // mappedV4Offset is the number of leading bits an IPv4-mapped
@@ -105,7 +104,7 @@ var ErrInvalidEnvironment = errors.New("invalid environment")
var ErrNonPositiveValue = errors.New("value must be positive") var ErrNonPositiveValue = errors.New("value must be positive")
// ErrInvalidPort is returned when an environment variable holding a // ErrInvalidPort is returned when an environment variable holding a
// TCP port number is set above the valid port range. // TCP port number is set to a number outside 1 to 65535.
var ErrInvalidPort = errors.New("invalid port") var ErrInvalidPort = errors.New("invalid port")
// ErrInvalidCIDR is returned when an environment variable holding a // ErrInvalidCIDR is returned when an environment variable holding a
@@ -363,17 +362,27 @@ func envPositiveInt(
// envPort returns the value of the named environment variable parsed // envPort returns the value of the named environment variable parsed
// as a TCP port number. Returns defaultValue if not set. A set value // as a TCP port number. Returns defaultValue if not set. A set value
// that is unparseable, below 1, or above maxPort is a hard error // that is unparseable, below 1, or above maxPort is a hard error
// naming the key and the bad value. // naming the key and the bad value; every out-of-range value wraps
// ErrInvalidPort, including one too large or too small for an int.
func envPort(key string, defaultValue int) (int, error) { func envPort(key string, defaultValue int) (int, error) {
port, err := envPositiveInt(key, defaultValue) v := os.Getenv(key)
if err != nil { if v == "" {
return 0, err return defaultValue, nil
} }
if port > maxPort { // strconv.ErrRange means a number too large or too small for an
// int, which is outside the port range as well.
port, err := strconv.Atoi(v)
if err != nil && !errors.Is(err, strconv.ErrRange) {
return 0, fmt.Errorf( return 0, fmt.Errorf(
"%w: %s must be at most %d, got %d", "invalid integer for %s: %q: %w", key, v, err,
ErrInvalidPort, key, maxPort, port, )
}
if err != nil || port < 1 || port > maxPort {
return 0, fmt.Errorf(
"%w: %s must be from 1 to %d, got %q",
ErrInvalidPort, key, maxPort, v,
) )
} }
+16 -45
View File
@@ -3,7 +3,6 @@ package config_test
import ( import (
"bytes" "bytes"
"log/slog" "log/slog"
"os"
"testing" "testing"
"time" "time"
@@ -71,14 +70,12 @@ func TestEnvironmentConfig(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
if tt.envValue != "" { if tt.envValue != "" {
t.Setenv( t.Setenv(
"WEBHOOKER_ENVIRONMENT", tt.envValue, "WEBHOOKER_ENVIRONMENT", tt.envValue,
) )
} else {
require.NoError(t, os.Unsetenv(
"WEBHOOKER_ENVIRONMENT",
))
} }
for k, v := range tt.envVars { for k, v := range tt.envVars {
@@ -199,14 +196,11 @@ func TestRetentionSweepInterval(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
t.Setenv("WEBHOOKER_ENVIRONMENT", "dev") t.Setenv("WEBHOOKER_ENVIRONMENT", "dev")
if tt.set { if tt.set {
t.Setenv("RETENTION_SWEEP_INTERVAL", tt.value) t.Setenv("RETENTION_SWEEP_INTERVAL", tt.value)
} else {
require.NoError(t, os.Unsetenv(
"RETENTION_SWEEP_INTERVAL",
))
} }
if tt.expectError { if tt.expectError {
@@ -341,14 +335,11 @@ func TestSessionIdleTimeout(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
t.Setenv("WEBHOOKER_ENVIRONMENT", "dev") t.Setenv("WEBHOOKER_ENVIRONMENT", "dev")
if tt.set { if tt.set {
t.Setenv("SESSION_IDLE_TIMEOUT", tt.value) t.Setenv("SESSION_IDLE_TIMEOUT", tt.value)
} else {
require.NoError(t, os.Unsetenv(
"SESSION_IDLE_TIMEOUT",
))
} }
if tt.expectError { if tt.expectError {
@@ -397,16 +388,12 @@ func TestDefaultDataDir(t *testing.T) {
t.Run("env="+name, func(t *testing.T) { t.Run("env="+name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
if env != "" { if env != "" {
t.Setenv("WEBHOOKER_ENVIRONMENT", env) t.Setenv("WEBHOOKER_ENVIRONMENT", env)
} else {
require.NoError(t, os.Unsetenv(
"WEBHOOKER_ENVIRONMENT",
))
} }
require.NoError(t, os.Unsetenv("DATA_DIR"))
var cfg *config.Config var cfg *config.Config
app := fxtest.New( app := fxtest.New(
@@ -446,9 +433,9 @@ func TestDataDirHelper(t *testing.T) {
t.Run(name, func(t *testing.T) { t.Run(name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
if set == "" { config.ClearEnvForTest(t)
require.NoError(t, os.Unsetenv("DATA_DIR"))
} else { if set != "" {
t.Setenv("DATA_DIR", set) t.Setenv("DATA_DIR", set)
} }
@@ -511,14 +498,11 @@ func TestReceiverRateLimit(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
t.Setenv("WEBHOOKER_ENVIRONMENT", "dev") t.Setenv("WEBHOOKER_ENVIRONMENT", "dev")
if tt.set { if tt.set {
t.Setenv("RECEIVER_RATE_LIMIT", tt.value) t.Setenv("RECEIVER_RATE_LIMIT", tt.value)
} else {
require.NoError(t, os.Unsetenv(
"RECEIVER_RATE_LIMIT",
))
} }
if tt.expectError { if tt.expectError {
@@ -630,12 +614,11 @@ func TestTrustedProxies(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
t.Setenv("WEBHOOKER_ENVIRONMENT", "dev") t.Setenv("WEBHOOKER_ENVIRONMENT", "dev")
if tt.set { if tt.set {
t.Setenv("TRUSTED_PROXIES", tt.value) t.Setenv("TRUSTED_PROXIES", tt.value)
} else {
require.NoError(t, os.Unsetenv("TRUSTED_PROXIES"))
} }
if tt.expectError { if tt.expectError {
@@ -742,14 +725,11 @@ func TestAllowedEgressCIDRs(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
t.Setenv("WEBHOOKER_ENVIRONMENT", "dev") t.Setenv("WEBHOOKER_ENVIRONMENT", "dev")
if tt.set { if tt.set {
t.Setenv("ALLOWED_EGRESS_CIDRS", tt.value) t.Setenv("ALLOWED_EGRESS_CIDRS", tt.value)
} else {
require.NoError(
t, os.Unsetenv("ALLOWED_EGRESS_CIDRS"),
)
} }
if tt.expectError { if tt.expectError {
@@ -817,13 +797,10 @@ func TestEgressAllowlistWarning(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
t.Setenv("WEBHOOKER_ENVIRONMENT", config.EnvironmentDev) t.Setenv("WEBHOOKER_ENVIRONMENT", config.EnvironmentDev)
if tt.allowed == "" { if tt.allowed != "" {
require.NoError(
t, os.Unsetenv("ALLOWED_EGRESS_CIDRS"),
)
} else {
t.Setenv("ALLOWED_EGRESS_CIDRS", tt.allowed) t.Setenv("ALLOWED_EGRESS_CIDRS", tt.allowed)
} }
@@ -956,20 +933,14 @@ func TestMetricsAuthConfig(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
if tt.username.set { if tt.username.set {
t.Setenv("METRICS_USERNAME", tt.username.value) t.Setenv("METRICS_USERNAME", tt.username.value)
} else {
require.NoError(
t, os.Unsetenv("METRICS_USERNAME"),
)
} }
if tt.password.set { if tt.password.set {
t.Setenv("METRICS_PASSWORD", tt.password.value) t.Setenv("METRICS_PASSWORD", tt.password.value)
} else {
require.NoError(
t, os.Unsetenv("METRICS_PASSWORD"),
)
} }
if tt.expectError { if tt.expectError {
+7 -18
View File
@@ -22,17 +22,6 @@ const malformedDotEnv = "PORT 19615\n" +
"this is not = valid ! syntax\n" + "this is not = valid ! syntax\n" +
"\"unclosed\n" "\"unclosed\n"
// unsetDotEnvKey makes dotEnvKey genuinely absent for the duration of
// the test and restores it afterwards. t.Setenv registers the restore;
// the Unsetenv that follows is what the test actually needs, because a
// variable set to the empty string is still present in os.Environ and
// godotenv would refuse to overwrite it.
func unsetDotEnvKey(t *testing.T) {
t.Helper()
t.Setenv(dotEnvKey, "placeholder")
require.NoError(t, os.Unsetenv(dotEnvKey))
}
// writeDotEnv writes contents to a .env file in a fresh temporary // writeDotEnv writes contents to a .env file in a fresh temporary
// directory and returns its path. // directory and returns its path.
func writeDotEnv(t *testing.T, contents string) string { func writeDotEnv(t *testing.T, contents string) string {
@@ -50,9 +39,9 @@ func writeDotEnv(t *testing.T, contents string) string {
// normally rather than be refused for a file it was never meant to // normally rather than be refused for a file it was never meant to
// have. // have.
// //
//nolint:paralleltest // unsetDotEnvKey uses t.Setenv. //nolint:paralleltest // ClearEnvForTest uses t.Setenv.
func TestLoadDotEnv_MissingFileIsFine(t *testing.T) { func TestLoadDotEnv_MissingFileIsFine(t *testing.T) {
unsetDotEnvKey(t) config.ClearEnvForTest(t)
absent := filepath.Join(t.TempDir(), config.DotEnvPath) absent := filepath.Join(t.TempDir(), config.DotEnvPath)
require.NoError(t, config.LoadDotEnvFileForTest(absent)) require.NoError(t, config.LoadDotEnvFileForTest(absent))
@@ -65,9 +54,9 @@ func TestLoadDotEnv_MissingFileIsFine(t *testing.T) {
// reaches the environment, which is the whole reason the file is read // reaches the environment, which is the whole reason the file is read
// at all. // at all.
// //
//nolint:paralleltest // unsetDotEnvKey uses t.Setenv. //nolint:paralleltest // ClearEnvForTest uses t.Setenv.
func TestLoadDotEnv_AppliesValues(t *testing.T) { func TestLoadDotEnv_AppliesValues(t *testing.T) {
unsetDotEnvKey(t) config.ClearEnvForTest(t)
path := writeDotEnv(t, "# a comment\n"+dotEnvKey+"=from-dot-env\n") path := writeDotEnv(t, "# a comment\n"+dotEnvKey+"=from-dot-env\n")
@@ -93,9 +82,9 @@ func TestLoadDotEnv_RealEnvironmentWins(t *testing.T) {
// reverts to its default; the process used to start that way with no // reverts to its default; the process used to start that way with no
// log line naming the file at all. // log line naming the file at all.
// //
//nolint:paralleltest // unsetDotEnvKey uses t.Setenv. //nolint:paralleltest // ClearEnvForTest uses t.Setenv.
func TestLoadDotEnv_MalformedFileAborts(t *testing.T) { func TestLoadDotEnv_MalformedFileAborts(t *testing.T) {
unsetDotEnvKey(t) config.ClearEnvForTest(t)
path := writeDotEnv( path := writeDotEnv(
t, malformedDotEnv+dotEnvKey+"=from-dot-env\n", t, malformedDotEnv+dotEnvKey+"=from-dot-env\n",
@@ -143,7 +132,7 @@ func TestLoadDotEnv_UnreadableFileAborts(t *testing.T) {
// //
//nolint:paralleltest // t.Chdir moves the whole process. //nolint:paralleltest // t.Chdir moves the whole process.
func TestLoadDotEnv_ReadsTheWorkingDirectory(t *testing.T) { func TestLoadDotEnv_ReadsTheWorkingDirectory(t *testing.T) {
unsetDotEnvKey(t) config.ClearEnvForTest(t)
dir := t.TempDir() dir := t.TempDir()
require.NoError(t, os.WriteFile( require.NoError(t, os.WriteFile(
+78 -91
View File
@@ -1,7 +1,6 @@
package config_test package config_test
import ( import (
"os"
"testing" "testing"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
@@ -121,10 +120,10 @@ func TestEnvBool(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
if tt.set { if tt.set {
t.Setenv(testEnvKey, tt.value) t.Setenv(testEnvKey, tt.value)
} else {
require.NoError(t, os.Unsetenv(testEnvKey))
} }
got, err := config.EnvBoolForTest( got, err := config.EnvBoolForTest(
@@ -145,17 +144,62 @@ func TestEnvBool(t *testing.T) {
} }
} }
// envIntCase is one row of the envPositiveInt and envPort tables.
type envIntCase struct {
name string
set bool
value string
expectError bool
errIs error
expected int
}
// runEnvIntCases runs each row through parse, which is
// envPositiveInt or envPort, with testEnvKey set to the row's value
// or left unset.
func runEnvIntCases(
t *testing.T,
parse func(key string, defaultValue int) (int, error),
defaultValue int,
tests []envIntCase,
) {
t.Helper()
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests.
config.ClearEnvForTest(t)
if tt.set {
t.Setenv(testEnvKey, tt.value)
}
got, err := parse(testEnvKey, defaultValue)
if tt.expectError {
require.Error(t, err)
assert.Contains(t, err.Error(), testEnvKey)
assert.Contains(t, err.Error(), tt.value)
if tt.errIs != nil {
require.ErrorIs(t, err, tt.errIs)
}
return
}
require.NoError(t, err)
assert.Equal(t, tt.expected, got)
})
}
}
//nolint:paralleltest // runEnvIntCases uses t.Setenv.
func TestEnvPositiveInt(t *testing.T) { func TestEnvPositiveInt(t *testing.T) {
const defaultValue = 7 const defaultValue = 7
tests := []struct { runEnvIntCases(t, config.EnvPositiveIntForTest, defaultValue, []envIntCase{
name string
set bool
value string
expectError bool
errIs error
expected int
}{
{ {
name: "unset returns the default integer", name: "unset returns the default integer",
expected: defaultValue, expected: defaultValue,
@@ -192,51 +236,14 @@ func TestEnvPositiveInt(t *testing.T) {
expectError: true, expectError: true,
errIs: config.ErrNonPositiveValue, errIs: config.ErrNonPositiveValue,
}, },
} })
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests.
if tt.set {
t.Setenv(testEnvKey, tt.value)
} else {
require.NoError(t, os.Unsetenv(testEnvKey))
}
got, err := config.EnvPositiveIntForTest(
testEnvKey, defaultValue,
)
if tt.expectError {
require.Error(t, err)
assert.Contains(t, err.Error(), testEnvKey)
assert.Contains(t, err.Error(), tt.value)
if tt.errIs != nil {
require.ErrorIs(t, err, tt.errIs)
}
return
}
require.NoError(t, err)
assert.Equal(t, tt.expected, got)
})
}
} }
//nolint:paralleltest // runEnvIntCases uses t.Setenv.
func TestEnvPort(t *testing.T) { func TestEnvPort(t *testing.T) {
const defaultValue = 8080 const defaultValue = 8080
tests := []struct { runEnvIntCases(t, config.EnvPortForTest, defaultValue, []envIntCase{
name string
set bool
value string
expectError bool
errIs error
expected int
}{
{ {
name: "unset returns the default port", name: "unset returns the default port",
expected: defaultValue, expected: defaultValue,
@@ -264,7 +271,14 @@ func TestEnvPort(t *testing.T) {
set: true, set: true,
value: "0", value: "0",
expectError: true, expectError: true,
errIs: config.ErrNonPositiveValue, errIs: config.ErrInvalidPort,
},
{
name: "negative is rejected",
set: true,
value: "-1",
expectError: true,
errIs: config.ErrInvalidPort,
}, },
{ {
name: "above the port range is rejected", name: "above the port range is rejected",
@@ -273,37 +287,14 @@ func TestEnvPort(t *testing.T) {
expectError: true, expectError: true,
errIs: config.ErrInvalidPort, errIs: config.ErrInvalidPort,
}, },
} {
name: "too large for an int is rejected",
for _, tt := range tests { set: true,
t.Run(tt.name, func(t *testing.T) { value: "99999999999999999999",
// Cannot use t.Parallel() here because t.Setenv expectError: true,
// is incompatible with parallel subtests. errIs: config.ErrInvalidPort,
if tt.set { },
t.Setenv(testEnvKey, tt.value) })
} else {
require.NoError(t, os.Unsetenv(testEnvKey))
}
got, err := config.EnvPortForTest(
testEnvKey, defaultValue,
)
if tt.expectError {
require.Error(t, err)
assert.Contains(t, err.Error(), testEnvKey)
if tt.errIs != nil {
require.ErrorIs(t, err, tt.errIs)
}
return
}
require.NoError(t, err)
assert.Equal(t, tt.expected, got)
})
}
} }
// TestEnvBindAddress covers BIND_ADDRESS parsing. // TestEnvBindAddress covers BIND_ADDRESS parsing.
@@ -319,10 +310,10 @@ func TestEnvBindAddress(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
if tt.set { if tt.set {
t.Setenv(testEnvKey, tt.value) t.Setenv(testEnvKey, tt.value)
} else {
require.NoError(t, os.Unsetenv(testEnvKey))
} }
got, err := config.EnvBindAddressForTest( got, err := config.EnvBindAddressForTest(
@@ -485,6 +476,7 @@ func TestNewRejectsBadEnvValues(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
t.Setenv("WEBHOOKER_ENVIRONMENT", "dev") t.Setenv("WEBHOOKER_ENVIRONMENT", "dev")
t.Setenv(tt.key, tt.value) t.Setenv(tt.key, tt.value)
@@ -646,14 +638,9 @@ func sentryEnvValueCases() []badEnvValueCase {
// break the legitimate unset case: absent variables still get their // break the legitimate unset case: absent variables still get their
// documented defaults. // documented defaults.
func TestNewUsesDefaultsWhenUnset(t *testing.T) { func TestNewUsesDefaultsWhenUnset(t *testing.T) {
config.ClearEnvForTest(t)
t.Setenv("WEBHOOKER_ENVIRONMENT", "dev") t.Setenv("WEBHOOKER_ENVIRONMENT", "dev")
for _, key := range []string{
envKeyPort, envKeyDebug, envKeyBindAddress, envKeySentryDSN,
} {
require.NoError(t, os.Unsetenv(key))
}
cfg, err := buildConfig(t) cfg, err := buildConfig(t)
require.NoError(t, err) require.NoError(t, err)
require.NotNil(t, cfg) require.NotNil(t, cfg)
+2 -3
View File
@@ -1,7 +1,6 @@
package config_test package config_test
import ( import (
"os"
"testing" "testing"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
@@ -101,10 +100,10 @@ func TestEnvSentryDSN(t *testing.T) {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv // Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests. // is incompatible with parallel subtests.
config.ClearEnvForTest(t)
if tt.set { if tt.set {
t.Setenv(envKeySentryDSN, tt.value) t.Setenv(envKeySentryDSN, tt.value)
} else {
require.NoError(t, os.Unsetenv(envKeySentryDSN))
} }
got, err := config.EnvSentryDSNForTest(envKeySentryDSN) got, err := config.EnvSentryDSNForTest(envKeySentryDSN)
+50
View File
@@ -0,0 +1,50 @@
package config
import (
"os"
"strings"
"testing"
)
// ClearEnvForTest unsets every variable in the process environment
// for the rest of the test, so a test sees only the variables it sets
// itself, not whatever the developer's shell exports. When the test
// ends it leaves the environment exactly as it found it: each variable
// it unset is put back, and any variable added since is removed.
func ClearEnvForTest(t *testing.T) {
t.Helper()
present := make(map[string]bool)
for _, entry := range os.Environ() {
key, _, _ := strings.Cut(entry, "=")
present[key] = true
// t.Setenv registers the restore; the Unsetenv after it is
// what makes the key absent, since a key set to the empty
// string is still present, and godotenv will not overwrite a
// present key.
t.Setenv(key, "")
err := os.Unsetenv(key)
if err != nil {
t.Fatalf("unsetting %s: %v", key, err)
}
}
// A variable the test adds other than through t.Setenv, as loading
// a .env file does, has no restore of its own.
t.Cleanup(func() {
for _, entry := range os.Environ() {
key, _, _ := strings.Cut(entry, "=")
if present[key] {
continue
}
err := os.Unsetenv(key)
if err != nil {
t.Errorf("unsetting %s: %v", key, err)
}
}
})
}
+36
View File
@@ -0,0 +1,36 @@
package config_test
import (
"os"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/config"
)
// TestClearEnvForTest_RemovesAddedVariables pins that a variable set
// after the clear other than through t.Setenv, as a test's .env file
// sets one, is gone once the test ends, so it cannot reach the tests
// that run after it.
//
//nolint:paralleltest // ClearEnvForTest uses t.Setenv.
func TestClearEnvForTest_RemovesAddedVariables(t *testing.T) {
// The outer clear keeps a value of the key exported in the shell
// from making it a variable the inner clear has to put back.
config.ClearEnvForTest(t)
t.Run("loads a .env file after the clear", func(t *testing.T) {
config.ClearEnvForTest(t)
path := writeDotEnv(t, dotEnvKey+"=from-dot-env\n")
require.NoError(t, config.LoadDotEnvFileForTest(path))
require.Equal(t, "from-dot-env", os.Getenv(dotEnvKey))
})
_, present := os.LookupEnv(dotEnvKey)
assert.False(
t, present,
"a variable set after the clear must not outlive the test",
)
}
+2 -15
View File
@@ -45,13 +45,8 @@ type ArchiveSweeper struct {
eng *Engine eng *Engine
log *slog.Logger log *slog.Logger
interval time.Duration interval time.Duration
cancel context.CancelFunc
// cancel needs no lock: fx calls the stop hook only after the wg sync.WaitGroup
// start hook has returned, so stop never reads it while start
// is still setting it.
cancel context.CancelFunc
wg sync.WaitGroup
} }
// NewArchiveSweeper creates the archive sweeper and registers // NewArchiveSweeper creates the archive sweeper and registers
@@ -168,18 +163,10 @@ func (s *ArchiveSweeper) sweep(ctx context.Context) {
var targets []database.Target var targets []database.Target
err := s.db.DB(). err := s.db.DB().
WithContext(ctx).
Model(&database.Target{}). Model(&database.Target{}).
Where("type = ?", database.TargetTypeDatabase). Where("type = ?", database.TargetTypeDatabase).
Find(&targets).Error Find(&targets).Error
if err != nil { if err != nil {
// The app stopping as a sweep starts cancels the listing.
// Stopping is not a failure, so it must not produce an
// error line.
if ctx.Err() != nil {
return
}
s.log.Error( s.log.Error(
"archive sweep: failed to list database targets", "archive sweep: failed to list database targets",
"error", err, "error", err,
-60
View File
@@ -1,11 +1,9 @@
package delivery_test package delivery_test
import ( import (
"bytes"
"context" "context"
"database/sql" "database/sql"
"fmt" "fmt"
"log/slog"
"net/http" "net/http"
"os" "os"
"path/filepath" "path/filepath"
@@ -683,64 +681,6 @@ func TestArchiveSweep_ClosesHandleOfRegisteredWriter(
) )
} }
// TestArchiveSweep_ClosesHandleBeforeReopening proves the sweep
// closes the handle it finds open before it reopens the file.
// TestArchiveSweep_LeavesArchiveClosed cannot see this: without the
// close, the reopen replaces the handle without closing it, the
// sweep then closes only the new one, and one connection leaks per
// archive per sweep.
func TestArchiveSweep_ClosesHandleBeforeReopening(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "archive.db")
w := delivery.NewExportArchiveWriter(
path, archiveTestLogger(), 0,
)
require.NoError(t, w.Open(time.Hour))
before, err := w.DB().DB()
require.NoError(t, err)
require.NoError(t, w.SweepExpired(time.Hour))
assert.Error(
t, before.PingContext(t.Context()),
"the handle open before the sweep must be closed by it",
)
}
// TestArchiveSweep_CancelledSweepLogsNoError proves a sweep whose
// context is already cancelled, as when the app stops just as a
// sweep starts, returns without an error line: stopping is not a
// failure.
func TestArchiveSweep_CancelledSweepLogsNoError(t *testing.T) {
t.Parallel()
env := setupArchiveTest(t)
var errorLines bytes.Buffer
sweeper := delivery.NewTestArchiveSweeper(
env.mainDB, env.eng,
slog.New(slog.NewTextHandler(
&errorLines,
&slog.HandlerOptions{Level: slog.LevelError},
)),
)
ctx, cancel := context.WithCancel(context.Background())
cancel()
sweeper.ExportSweep(ctx)
assert.Empty(
t, errorLines.String(),
"a cancelled sweep must not log at error level",
)
}
// TestArchiveSweep_NeverExpiryUntouched proves the sweep is a // TestArchiveSweep_NeverExpiryUntouched proves the sweep is a
// no-op for the default retention policy, so archives with no // no-op for the default retention policy, so archives with no
// expiry (or the literal "never") behave exactly as before. // expiry (or the literal "never") behave exactly as before.
+16 -2
View File
@@ -355,9 +355,14 @@ func TestProcessRetryTask_SuccessfulRetry(t *testing.T) {
s := newISetup(t) s := newISetup(t)
var receivedBody string
ts := httptest.NewServer( ts := httptest.NewServer(
http.HandlerFunc( http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) { func(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
receivedBody = string(body)
w.WriteHeader(http.StatusOK) w.WriteHeader(http.StatusOK)
}, },
), ),
@@ -397,6 +402,8 @@ func TestProcessRetryTask_SuccessfulRetry(t *testing.T) {
context.TODO(), &task, context.TODO(), &task,
) )
assert.Equal(t, event.Body, receivedBody)
iAssertStatus(t, s.WebhookDB, d.ID, iAssertStatus(t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered, database.DeliveryStatusDelivered,
) )
@@ -443,9 +450,14 @@ func TestProcessRetryTask_LargeBody_FetchFromDB(
s := newISetup(t) s := newISetup(t)
var receivedBody string
ts := httptest.NewServer( ts := httptest.NewServer(
http.HandlerFunc( http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) { func(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
receivedBody = string(body)
w.WriteHeader(http.StatusOK) w.WriteHeader(http.StatusOK)
}, },
), ),
@@ -482,6 +494,8 @@ func TestProcessRetryTask_LargeBody_FetchFromDB(
context.TODO(), &task, context.TODO(), &task,
) )
assert.Equal(t, largeBody, receivedBody)
iAssertStatus(t, s.WebhookDB, d.ID, iAssertStatus(t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered, database.DeliveryStatusDelivered,
) )
+6 -10
View File
@@ -3,7 +3,6 @@ package delivery
import ( import (
"context" "context"
"fmt" "fmt"
"path/filepath"
"strings" "strings"
"sync" "sync"
"time" "time"
@@ -277,10 +276,9 @@ func (t *databaseTarget) releaseSweepWriter(
} }
// newWriter builds the writer for a database target's archive. The // newWriter builds the writer for a database target's archive. The
// file lives beside the webhook's event database in the data // file is the one ArchivePath gives for the webhook and the target as
// directory and is named for the webhook and the target as the main // the main database names them now; from then on only rename changes
// database has them now; from then on only rename changes the name // the name the writer uses. It does not touch the archive file.
// the writer uses. It does not touch the archive file.
func (t *databaseTarget) newWriter( func (t *databaseTarget) newWriter(
targetID string, targetID string,
) (*archiveWriter, error) { ) (*archiveWriter, error) {
@@ -299,12 +297,10 @@ func (t *databaseTarget) newWriter(
) )
} }
dir := filepath.Dir(t.eng.dbManager.DBPath(target.WebhookID)) w := newArchiveWriter(
name := ArchiveFileName( ArchivePath(t.eng.dbManager, &target.Webhook, &target),
target.Webhook.Name, target.Name, target.ID, t.eng.log,
) )
w := newArchiveWriter(filepath.Join(dir, name), t.eng.log)
w.webhookID = target.WebhookID w.webhookID = target.WebhookID
return w, nil return w, nil
+275
View File
@@ -0,0 +1,275 @@
package delivery
import (
"compress/gzip"
"context"
"database/sql"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"log/slog"
"path/filepath"
"time"
"unicode/utf8"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/gormlog"
)
// archiveTableQuery counts the archive's table: 0 when the archive
// writer has created the file but not yet the table in it.
const archiveTableQuery = "SELECT count(*) FROM sqlite_master " +
"WHERE type = 'table' AND name = 'archived_events'"
// ArchivePath returns where a database target's archive file is: in
// the data directory, beside the webhook's event database, under the
// name ArchiveFileName gives it.
func ArchivePath(
dbMgr *database.WebhookDBManager,
webhook *database.Webhook,
target *database.Target,
) string {
return filepath.Join(
filepath.Dir(dbMgr.DBPath(webhook.ID)),
ArchiveFileName(webhook.Name, target.Name, target.ID),
)
}
// ArchiveExportFileName returns the name a database target's archive
// downloads under:
// archive-WEBHOOKNAME-TARGETNAME-YYYYMMDDTHHMMSSZ.json.gz, the names
// made safe as in ArchiveFileName and the time in UTC.
func ArchiveExportFileName(
webhookName, targetName string, at time.Time,
) string {
return "archive-" + archiveNamePart(webhookName) + "-" +
archiveNamePart(targetName) + "-" +
at.UTC().Format("20060102T150405Z") + ".json.gz"
}
// ArchiveExport is a database target's archive opened for download.
// It reads the file on its own connection, inside one read-only
// transaction, so it writes out the archive as it stood when
// OpenArchiveExport returned.
//
// Archives are in WAL mode, where a reader works from a snapshot and
// never blocks a writer: archive writes go on while an export is open,
// and the export does not see them. SQLite cannot checkpoint the -wal
// past an open snapshot, so the -wal grows until the export is closed.
type ArchiveExport struct {
db *sql.DB
tx *gorm.DB
// empty is true when there is nothing to read: no file, or a file
// without the archive's table yet.
empty bool
}
// exportedName is how an export names its webhook and its target.
type exportedName struct {
ID string `json:"id"`
Name string `json:"name"`
}
// OpenArchiveExport opens the archive file at path for export and
// takes the snapshot the export reads. It never creates the file: with
// no file at path, the export has no rows.
//
// Once it has returned, the file is open, so a rename or a move of it
// does not affect the export, which reads the same file under its new
// name.
//
// The transaction lasts as long as ctx does, so ctx must last for the
// whole export.
func OpenArchiveExport(
ctx context.Context, path string, log *slog.Logger,
) (*ArchiveExport, error) {
if !fileExists(path) {
return &ArchiveExport{empty: true}, nil
}
db, err := database.OpenSQLite(path, archiveModeExisting)
if err != nil {
return nil, fmt.Errorf("opening archive %s: %w", path, err)
}
gdb, err := gorm.Open(
sqlite.Dialector{Conn: db}, &gorm.Config{
// Never leave this at GORM's default. See
// internal/gormlog.
Logger: gormlog.New(log),
},
)
if err != nil {
_ = db.Close()
return nil, fmt.Errorf("opening archive %s: %w", path, err)
}
// ReadOnly makes the driver begin a deferred transaction in place
// of the BEGIN IMMEDIATE the connection string asks for, so the
// export never takes the archive's write lock.
tx := gdb.WithContext(ctx).Begin(&sql.TxOptions{ReadOnly: true})
if tx.Error != nil {
_ = db.Close()
return nil, fmt.Errorf("reading archive %s: %w", path, tx.Error)
}
// The transaction's first read is what takes the snapshot.
var tables int
err = tx.Raw(archiveTableQuery).Row().Scan(&tables)
if err != nil {
_ = tx.Rollback()
_ = db.Close()
return nil, fmt.Errorf("reading archive %s: %w", path, err)
}
return &ArchiveExport{db: db, tx: tx, empty: tables == 0}, nil
}
// WriteGzipJSON writes the export to w as one gzipped JSON object:
// webhook and target, each an id and a name; exported_at; and
// archived_events, one object per archived row, keyed by column name.
// A body that is not valid UTF-8 cannot be a JSON string, so it is
// written in base64, with "body_encoding": "base64" beside it.
//
// Each row is written out before the next is read, so neither the
// archive nor its JSON is ever held in memory whole. After an error
// the gzip stream is left unfinished, so what was written does not
// decompress as a whole file.
func (x *ArchiveExport) WriteGzipJSON(
ctx context.Context,
w io.Writer,
webhook *database.Webhook,
target *database.Target,
exportedAt time.Time,
) error {
head, err := json.Marshal(map[string]any{
"webhook": exportedName{ID: webhook.ID, Name: webhook.Name},
"target": exportedName{ID: target.ID, Name: target.Name},
"exported_at": exportedAt.UTC(),
})
if err != nil {
return fmt.Errorf("encoding archive export: %w", err)
}
zw := gzip.NewWriter(w)
err = x.writeJSON(ctx, zw, head)
if err != nil {
return fmt.Errorf("writing archive export: %w", err)
}
return zw.Close()
}
// Close ends the export's transaction and closes its connection.
func (x *ArchiveExport) Close() error {
if x.db == nil {
return nil
}
_ = x.tx.Rollback()
return x.db.Close()
}
// writeJSON writes head with archived_events added as its last key,
// the rows going into it one at a time.
func (x *ArchiveExport) writeJSON(
ctx context.Context, w io.Writer, head []byte,
) error {
// head goes out without its closing brace, so that
// archived_events can follow it.
_, err := w.Write(head[:len(head)-1])
if err != nil {
return err
}
_, err = io.WriteString(w, `,"archived_events":[`)
if err != nil {
return err
}
err = x.writeRows(ctx, w)
if err != nil {
return err
}
_, err = io.WriteString(w, "\n]}\n")
return err
}
// writeRows writes each archived row to w, oldest first, one per line,
// separated by commas.
func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error {
if x.empty {
return nil
}
rows, err := x.tx.WithContext(ctx).
Model(&archivedEvent{}).Order("id").Rows()
if err != nil {
return err
}
defer func() { _ = rows.Close() }()
for sep := "\n"; rows.Next(); sep = ",\n" {
var ev archivedEvent
err = x.tx.ScanRows(rows, &ev)
if err != nil {
return err
}
_, err = io.WriteString(w, sep)
if err != nil {
return err
}
err = writeRow(w, &ev)
if err != nil {
return err
}
}
return rows.Err()
}
// writeRow writes an archived row to w as a JSON object keyed by
// column name, its body in base64 when it is not valid UTF-8.
func writeRow(w io.Writer, ev *archivedEvent) error {
row := map[string]any{
"id": ev.ID,
"event_id": ev.EventID,
"webhook_id": ev.WebhookID,
"entrypoint_id": ev.EntrypointID,
"method": ev.Method,
"headers": ev.Headers,
"body": ev.Body,
"content_type": ev.ContentType,
"archived_at": ev.ArchivedAt.UTC(),
}
if !utf8.ValidString(ev.Body) {
row["body"] = base64.StdEncoding.EncodeToString([]byte(ev.Body))
row["body_encoding"] = "base64"
}
line, err := json.Marshal(row)
if err != nil {
return err
}
_, err = w.Write(line)
return err
}
@@ -0,0 +1,393 @@
package delivery_test
import (
"bufio"
"bytes"
"compress/gzip"
"crypto/rand"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"os"
"path/filepath"
"runtime"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
)
// The webhook and the target the export tests' archives belong to.
const (
exportWebhookID = "wh-export"
exportWebhookName = "Orders (EU)"
exportTargetID = "tgt-export"
exportTargetName = "Long-term archive"
)
const (
// binaryBody is a body that is not valid UTF-8.
binaryBody = "\xff\xfe\x00\x01binary\x80"
// openedEventID is the event the snapshot tests archive before
// they open the export.
openedEventID = "opened"
)
// writeExportTo writes export to w as the archive of the export tests'
// webhook and target, exported at 2026-10-02T12:03:04Z.
func writeExportTo(
t *testing.T, export *delivery.ArchiveExport, w io.Writer,
) error {
t.Helper()
return export.WriteGzipJSON(
t.Context(), w,
&database.Webhook{
BaseModel: database.BaseModel{ID: exportWebhookID},
Name: exportWebhookName,
},
&database.Target{
BaseModel: database.BaseModel{ID: exportTargetID},
Name: exportTargetName,
},
time.Date(2026, 10, 2, 12, 3, 4, 0, time.UTC),
)
}
// exportArchive runs a whole export of the archive at path and returns
// its JSON, decompressed and parsed.
func exportArchive(t *testing.T, path string) map[string]any {
t.Helper()
export, err := delivery.OpenArchiveExport(
t.Context(), path, archiveTestLogger(),
)
require.NoError(t, err)
defer func() { require.NoError(t, export.Close()) }()
return writeExport(t, export)
}
// writeExport writes an opened export and returns its JSON,
// decompressed and parsed. Reading to the end makes the gzip reader
// check that the stream was finished.
func writeExport(
t *testing.T, export *delivery.ArchiveExport,
) map[string]any {
t.Helper()
var buf bytes.Buffer
require.NoError(t, writeExportTo(t, export, &buf))
zr, err := gzip.NewReader(&buf)
require.NoError(t, err)
raw, err := io.ReadAll(zr)
require.NoError(t, err)
var got map[string]any
require.NoError(t, json.Unmarshal(raw, &got))
return got
}
// exportedEvents returns an export's archived_events.
func exportedEvents(t *testing.T, got map[string]any) []map[string]any {
t.Helper()
list, ok := got["archived_events"].([]any)
require.True(t, ok, "archived_events must be an array: %v", got)
events := make([]map[string]any, len(list))
for i, v := range list {
events[i], ok = v.(map[string]any)
require.True(t, ok, "an archived event must be an object: %v", v)
}
return events
}
// exportedEventIDs returns the event_id of each of an export's
// archived_events.
func exportedEventIDs(t *testing.T, got map[string]any) []string {
t.Helper()
events := exportedEvents(t, got)
ids := make([]string, 0, len(events))
for _, ev := range events {
ids = append(ids, fmt.Sprint(ev["event_id"]))
}
return ids
}
// TestArchiveExport_MatchesStoredRows proves an export holds the
// webhook, the target, the time, and every column of every stored
// row: a body that is valid UTF-8 as a string, and one that is not in
// base64, marked as such.
func TestArchiveExport_MatchesStoredRows(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "archive.db")
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
bodies := []string{`{"order":1}`, "plain text", "", binaryBody}
for i, body := range bodies {
require.NoError(t, w.Write(delivery.ExportArchivedEvent{
EventID: fmt.Sprintf("ev-%d", i),
WebhookID: exportWebhookID,
EntrypointID: "ep-1",
Method: "POST",
Headers: `{"X-Test":["yes"]}`,
Body: body,
ContentType: testContentType,
}, 0))
}
var stored []delivery.ExportArchivedEvent
require.NoError(t, openArchiveDBForRead(t, path).
Order("id").Find(&stored).Error)
got := exportArchive(t, path)
assert.Equal(t,
map[string]any{"id": exportWebhookID, "name": exportWebhookName},
got["webhook"],
)
assert.Equal(t,
map[string]any{"id": exportTargetID, "name": exportTargetName},
got["target"],
)
assert.Equal(t, "2026-10-02T12:03:04Z", got["exported_at"])
events := exportedEvents(t, got)
require.Len(t, events, len(bodies))
for i, row := range stored {
assertExportedRow(t, row, events[i])
}
}
// assertExportedRow checks that ev, from an export, holds every column
// of the stored row.
func assertExportedRow(
t *testing.T, row delivery.ExportArchivedEvent, ev map[string]any,
) {
t.Helper()
archivedAt, err := time.Parse(
time.RFC3339Nano, fmt.Sprint(ev["archived_at"]),
)
require.NoError(t, err)
assert.True(t, archivedAt.Equal(row.ArchivedAt))
assert.EqualValues(t, row.ID, ev["id"])
assert.Equal(t, row.EventID, ev["event_id"])
assert.Equal(t, row.WebhookID, ev["webhook_id"])
assert.Equal(t, row.EntrypointID, ev["entrypoint_id"])
assert.Equal(t, row.Method, ev["method"])
assert.Equal(t, row.Headers, ev["headers"])
assert.Equal(t, row.ContentType, ev["content_type"])
if row.Body != binaryBody {
assert.Equal(t, row.Body, ev["body"])
assert.Len(t, ev, 9, "the nine columns and nothing else: %v", ev)
return
}
body, err := base64.StdEncoding.DecodeString(fmt.Sprint(ev["body"]))
require.NoError(t, err)
assert.Equal(t, binaryBody, string(body))
assert.Equal(t, "base64", ev["body_encoding"])
assert.Len(t, ev, 10, "the nine columns and body_encoding: %v", ev)
}
// TestArchiveExport_Empty proves an archive with nothing in it exports
// as an empty archived_events: no file, which the export must not
// create; a file the archive writer has not yet put its table in; and
// a table with no rows.
func TestArchiveExport_Empty(t *testing.T) {
t.Parallel()
dir := t.TempDir()
missing := filepath.Join(dir, "missing.db")
noTable := filepath.Join(dir, "no-table.db")
noRows := filepath.Join(dir, "no-rows.db")
require.NoError(t, os.WriteFile(noTable, nil, 0o600))
require.NoError(t,
delivery.NewExportArchiveWriter(noRows, archiveTestLogger(), 0).
Open(0),
)
for _, path := range []string{missing, noTable, noRows} {
assert.Empty(t, exportedEvents(t, exportArchive(t, path)), path)
}
for _, suffix := range archiveFileSuffixes() {
assert.NoFileExists(t, missing+suffix)
}
}
// TestArchiveExport_ReadsOneSnapshot proves an export writes the
// archive as it was when it was opened, and holds up no archive
// write: a row written while the export is open is stored, and is not
// in the export. A write held up for the whole busy timeout would
// fail.
func TestArchiveExport_ReadsOneSnapshot(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "archive.db")
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: openedEventID}, 0))
export, err := delivery.OpenArchiveExport(
t.Context(), path, archiveTestLogger(),
)
require.NoError(t, err)
defer func() { require.NoError(t, export.Close()) }()
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "during"}, 0))
assert.Equal(t,
[]string{openedEventID}, exportedEventIDs(t, writeExport(t, export)),
)
var stored int64
require.NoError(t, openArchiveDBForRead(t, path).
Model(&delivery.ExportArchivedEvent{}).Count(&stored).Error)
assert.Equal(t, int64(2), stored)
}
// TestArchiveExport_SurvivesRename proves that renaming the archive
// while an export of it is open, as renaming its webhook or target
// does, leaves the export reading the same file.
func TestArchiveExport_SurvivesRename(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "archive-old.db")
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: openedEventID}, 0))
export, err := delivery.OpenArchiveExport(
t.Context(), path, archiveTestLogger(),
)
require.NoError(t, err)
defer func() { require.NoError(t, export.Close()) }()
require.NoError(t, w.Rename("archive-new.db"))
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "after"}, 0))
require.NoFileExists(t, path)
assert.Equal(t,
[]string{openedEventID}, exportedEventIDs(t, writeExport(t, export)),
)
}
// heapPeak is an io.Writer that discards what it is given and records
// the largest heap it saw at a write. It collects garbage before each
// reading, so the heap it reads is what is still held.
type heapPeak struct {
max uint64
}
func (p *heapPeak) Write(b []byte) (int, error) {
var m runtime.MemStats
runtime.GC()
runtime.ReadMemStats(&m)
p.max = max(p.max, m.HeapAlloc)
return len(b), nil
}
// TestArchiveExport_Streams proves an export holds neither the archive
// nor its output in memory whole: exporting a 3 MiB archive grows the
// heap by less than half of that. The bodies are random bytes in
// base64, which gzip shrinks by only a quarter, so an export that read
// every row before writing, or built the JSON or the gzipped file
// before writing it, would hold at least 2.25 MiB at a write. A
// streaming export holds about 1 MiB, most of it gzip's compressor,
// which is why the archive is not smaller.
//
//nolint:paralleltest // It measures the heap, which tests share.
func TestArchiveExport_Streams(t *testing.T) {
const (
rows = 96
bodySize = 32 << 10
limit = rows * bodySize / 2
)
path := filepath.Join(t.TempDir(), "archive.db")
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
// Base64 makes four characters of every three bytes.
random := make([]byte, bodySize/4*3)
for range rows {
_, _ = rand.Read(random)
require.NoError(t, w.Write(delivery.ExportArchivedEvent{
Body: base64.StdEncoding.EncodeToString(random),
}, 0))
}
export, err := delivery.OpenArchiveExport(
t.Context(), path, archiveTestLogger(),
)
require.NoError(t, err)
defer func() { require.NoError(t, export.Close()) }()
// Through a buffer, the heap is read once per 256 KiB of output
// rather than at each of gzip's small writes, which takes far
// longer.
peak := &heapPeak{}
buffered := bufio.NewWriterSize(peak, 256<<10)
runtime.GC()
var start runtime.MemStats
runtime.ReadMemStats(&start)
require.NoError(t, writeExportTo(t, export, buffered))
require.NoError(t, buffered.Flush())
assert.Less(t, peak.max, start.HeapAlloc+limit,
"heap at the start %d, at its peak %d", start.HeapAlloc, peak.max,
)
}
// TestArchiveExportFileName proves the download is named for the
// webhook and the target, with the names made safe as for the archive
// file, and the export time in UTC.
func TestArchiveExportFileName(t *testing.T) {
t.Parallel()
cest := time.FixedZone("CEST", int((2 * time.Hour).Seconds()))
assert.Equal(t,
"archive-orders-eu-long-term-archive-20261002T120304Z.json.gz",
delivery.ArchiveExportFileName(
exportWebhookName, exportTargetName,
time.Date(2026, 10, 2, 14, 3, 4, 0, cest),
),
)
}
+3 -2
View File
@@ -117,8 +117,8 @@ func readFirstBootSecrets(
} }
// bootAtDebug starts and stops the real application graph against // bootAtDebug starts and stops the real application graph against
// dataDir with DEBUG=true, and returns everything it wrote to standard // dataDir with DEBUG=true and nothing else set, and returns everything
// output. // it wrote to standard output.
// //
// config.New reads DEBUG from the environment exactly as the binary // config.New reads DEBUG from the environment exactly as the binary
// does, internal/logger builds the handler it builds in production, // does, internal/logger builds the handler it builds in production,
@@ -128,6 +128,7 @@ func readFirstBootSecrets(
func bootAtDebug(t *testing.T, dataDir string) string { func bootAtDebug(t *testing.T, dataDir string) string {
t.Helper() t.Helper()
config.ClearEnvForTest(t)
t.Setenv("DEBUG", "true") t.Setenv("DEBUG", "true")
t.Setenv("DATA_DIR", dataDir) t.Setenv("DATA_DIR", dataDir)
+2 -1
View File
@@ -97,7 +97,8 @@ type Handlers struct {
// names through the archive rename, the save and any move back. // names through the archive rename, the save and any move back.
// Interleaved, one could rename an archive between another's // Interleaved, one could rename an archive between another's
// rename and save, leaving the file named for one edit and the // rename and save, leaving the file named for one edit and the
// stored names from the other. // stored names from the other. An archive download holds it while
// it reads the stored names and opens the file they give.
renameMu sync.Mutex renameMu sync.Mutex
// dummyVerifications counts the equivalent-cost verifications // dummyVerifications counts the equivalent-cost verifications
+2 -1
View File
@@ -77,7 +77,8 @@ func settingRows(cfg *config.Config) []settingRow {
{ {
"RETENTION_SWEEP_INTERVAL", "RETENTION_SWEEP_INTERVAL",
"How often the retention reaper and archive sweeper run " + "How often the retention reaper and archive sweeper run " +
"(Go duration, must be positive)", "(Go duration, must be positive). A value that does " +
"not parse, or is zero or negative, fails startup",
cfg.RetentionSweepInterval.String(), cfg.RetentionSweepInterval.String(),
}, },
{ {
+5 -1
View File
@@ -191,6 +191,7 @@ func storedRetentionDays(
type sourceTestEnv struct { type sourceTestEnv struct {
handlers *handlers.Handlers handlers *handlers.Handlers
db *database.Database db *database.Database
dbMgr *database.WebhookDBManager
archives *recordingArchives archives *recordingArchives
cookies []*http.Cookie cookies []*http.Cookie
} }
@@ -204,9 +205,11 @@ func setupSourceTest(t *testing.T) *sourceTestEnv {
var db *database.Database var db *database.Database
var dbMgr *database.WebhookDBManager
var archives *recordingArchives var archives *recordingArchives
app := newTestApp(t, &h, &sess, &db, &archives) app := newTestApp(t, &h, &sess, &db, &dbMgr, &archives)
app.RequireStart() app.RequireStart()
t.Cleanup(app.RequireStop) t.Cleanup(app.RequireStop)
@@ -214,6 +217,7 @@ func setupSourceTest(t *testing.T) *sourceTestEnv {
return &sourceTestEnv{ return &sourceTestEnv{
handlers: h, handlers: h,
db: db, db: db,
dbMgr: dbMgr,
archives: archives, archives: archives,
cookies: authenticatedCookies( cookies: authenticatedCookies(
t, sess, sourceTestUserID, "sourceuser", t, sess, sourceTestUserID, "sourceuser",
+121
View File
@@ -0,0 +1,121 @@
package handlers
import (
"context"
"errors"
"net/http"
"time"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
)
// downloadWriteTimeout is how long one write of a download may wait
// for a client that has stopped reading.
const downloadWriteTimeout = 60 * time.Second
// HandleTargetDownload serves a database target's archive as one
// gzipped JSON file, named for the webhook, the target and the time;
// see delivery.ArchiveExport.WriteGzipJSON for what it holds. Other
// target types have no archive and are a 404.
//
// A download runs for as long as the client keeps reading: it reads
// under a context the request limit does not cancel, and gives each
// write its own deadline in place of the server's write timeout. It
// stops when a write fails.
func (h *Handlers) HandleTargetDownload() http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
ctx := context.WithoutCancel(r.Context())
webhook, target, export, ok := h.openTargetArchive(ctx, w, r)
if !ok {
return
}
defer func() { _ = export.Close() }()
now := time.Now()
w.Header().Set("Content-Type", "application/gzip")
w.Header().Set(
"Content-Disposition",
`attachment; filename="`+delivery.ArchiveExportFileName(
webhook.Name, target.Name, now,
)+`"`,
)
err := export.WriteGzipJSON(
ctx,
downloadWriter{w: w, rc: http.NewResponseController(w)},
&webhook, target, now,
)
if err != nil {
h.log.Error(
"failed to export archive",
"target_id", target.ID,
"error", err,
)
// The 200 has gone out. Aborting the connection is what
// tells the client the file is incomplete.
panic(http.ErrAbortHandler)
}
}
}
// downloadWriter writes a download to the client, giving each write
// downloadWriteTimeout to finish.
type downloadWriter struct {
w http.ResponseWriter
rc *http.ResponseController
}
func (d downloadWriter) Write(b []byte) (int, error) {
// A writer that has no write deadline, such as a test's recorder,
// answers http.ErrNotSupported and needs none extended.
err := d.rc.SetWriteDeadline(time.Now().Add(downloadWriteTimeout))
if err != nil && !errors.Is(err, http.ErrNotSupported) {
return 0, err
}
return d.w.Write(b)
}
// openTargetArchive opens the archive of the request's database target
// for export, with its reads under ctx. It reports false once it has
// written the response.
//
// It holds renameMu, which every archive rename runs under, while it
// reads the stored names and opens the file, so the file it opens is
// the one those names give. It lets go before the export is streamed:
// once the file is open, a rename does not affect the export.
func (h *Handlers) openTargetArchive(
ctx context.Context,
w http.ResponseWriter,
r *http.Request,
) (database.Webhook, *database.Target, *delivery.ArchiveExport, bool) {
h.renameMu.Lock()
defer h.renameMu.Unlock()
webhook, target, ok := h.ownedTarget(w, r)
if !ok {
return database.Webhook{}, nil, nil, false
}
if target.Type != database.TargetTypeDatabase {
h.renderError(w, r, http.StatusNotFound)
return database.Webhook{}, nil, nil, false
}
export, err := delivery.OpenArchiveExport(
ctx, delivery.ArchivePath(h.dbMgr, &webhook, target), h.log,
)
if err != nil {
h.serverError(w, r, "failed to open archive for export", err)
return database.Webhook{}, nil, nil, false
}
return webhook, target, export, true
}
+379
View File
@@ -0,0 +1,379 @@
package handlers_test
import (
"bytes"
"compress/gzip"
"context"
"crypto/rand"
"encoding/json"
"errors"
"io"
"log/slog"
"net"
"net/http"
"net/http/httptest"
"net/url"
"sync"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
"sneak.berlin/go/webhooker/internal/middleware"
)
// errClientGone is the write failure of a client that has gone away.
var errClientGone = errors.New("client gone")
// downloadPath is the archive download route of a target.
func downloadPath(webhookID, targetID string) string {
return "/hook/" + webhookID + "/targets/" + targetID + "/download"
}
// renameTarget submits the edit form renaming a target to Renamed.
func renameTarget(
env *sourceTestEnv, webhookID, targetID string,
) *httptest.ResponseRecorder {
form := url.Values{}
form.Set("name", "Renamed")
return submitTargetEdit(env, webhookID, targetID, form)
}
// TestHandleTargetDownload proves a database target's archive
// downloads as a gzipped JSON attachment named for the webhook, the
// target and the time, here with no archive file yet, so with no
// rows; and that a target of another type has no download.
func TestHandleTargetDownload(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
logTarget := seedTarget(t, env.db, wh.ID, database.TargetTypeLog)
w := serveTarget(
env, http.MethodGet, downloadPath(wh.ID, archive.ID), nil,
)
require.Equal(t, http.StatusOK, w.Code, w.Body.String())
assert.Equal(t, "application/gzip", w.Header().Get("Content-Type"))
assert.Regexp(t,
`^attachment; filename="archive-seeded-t-database-`+
`\d{8}T\d{6}Z\.json\.gz"$`,
w.Header().Get("Content-Disposition"),
)
zr, err := gzip.NewReader(w.Body)
require.NoError(t, err)
var got map[string]json.RawMessage
require.NoError(t, json.NewDecoder(zr).Decode(&got))
assert.JSONEq(t,
`{"id":"`+archive.ID+`","name":"t-database"}`,
string(got["target"]),
)
assert.JSONEq(t, `[]`, string(got["archived_events"]))
w = serveTarget(
env, http.MethodGet, downloadPath(wh.ID, logTarget.ID), nil,
)
assert.Equal(t, http.StatusNotFound, w.Code)
}
// TestHandleTargetDownload_WaitsForRename proves a download reads the
// target's names and opens its archive under the lock a rename holds:
// started while an edit is renaming the archive, it waits, and is
// named for the target's new name.
func TestHandleTargetDownload_WaitsForRename(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
renaming, release := env.archives.BlockNextRename()
edited := make(chan *httptest.ResponseRecorder, 1)
go func() {
edited <- renameTarget(env, wh.ID, archive.ID)
}()
<-renaming
downloaded := make(chan *httptest.ResponseRecorder, 1)
go func() {
downloaded <- serveTarget(
env, http.MethodGet, downloadPath(wh.ID, archive.ID), nil,
)
}()
select {
case <-downloaded:
release()
t.Fatal("the download did not wait for the rename")
case <-time.After(100 * time.Millisecond):
}
release()
require.Equal(t, http.StatusSeeOther, (<-edited).Code)
w := <-downloaded
require.Equal(t, http.StatusOK, w.Code)
assert.Contains(t,
w.Header().Get("Content-Disposition"), "archive-seeded-renamed-",
)
}
// stalledWriter is a response writer whose first write waits until
// resume is closed, closing writing when it starts to wait.
type stalledWriter struct {
*httptest.ResponseRecorder
once sync.Once
writing chan struct{}
resume chan struct{}
}
func (s *stalledWriter) Write(b []byte) (int, error) {
s.once.Do(func() {
close(s.writing)
<-s.resume
})
return s.ResponseRecorder.Write(b)
}
// TestHandleTargetDownload_StreamsWithoutTheLock proves a download
// lets go of the rename lock once its archive is open: while the
// download is stalled writing, an edit can still rename the target.
func TestHandleTargetDownload_StreamsWithoutTheLock(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
req := httptest.NewRequestWithContext(
t.Context(), http.MethodGet, downloadPath(wh.ID, archive.ID), nil,
)
for _, c := range env.cookies {
req.AddCookie(c)
}
sw := &stalledWriter{
ResponseRecorder: httptest.NewRecorder(),
writing: make(chan struct{}),
resume: make(chan struct{}),
}
downloaded := make(chan struct{})
go func() {
targetRouter(env).ServeHTTP(sw, req)
close(downloaded)
}()
<-sw.writing
edited := make(chan *httptest.ResponseRecorder, 1)
go func() {
edited <- renameTarget(env, wh.ID, archive.ID)
}()
select {
case w := <-edited:
assert.Equal(t, http.StatusSeeOther, w.Code)
case <-time.After(10 * time.Second):
t.Error("the rename waited for the download")
}
close(sw.resume)
<-downloaded
assert.Equal(t, http.StatusOK, sw.Code)
}
// seedArchive writes rows to the archive file at path, each with a
// body of bodySize random bytes, which do not compress. Its table has
// only the columns the test fills; an export writes the others empty.
func seedArchive(t *testing.T, path string, rows, bodySize int) {
t.Helper()
db, err := database.OpenSQLite(path, database.SQLiteModeCreate)
require.NoError(t, err)
defer func() { require.NoError(t, db.Close()) }()
_, err = db.ExecContext(t.Context(),
"CREATE TABLE archived_events (id INTEGER PRIMARY KEY, body TEXT)",
)
require.NoError(t, err)
body := make([]byte, bodySize)
for range rows {
_, _ = rand.Read(body)
_, err = db.ExecContext(t.Context(),
"INSERT INTO archived_events (body) VALUES (?)", string(body),
)
require.NoError(t, err)
}
}
// limitedServer serves the target routes as the server does, behind the
// access log, whose lines it returns, and the request limit, here
// limit, which is also its write timeout. Each connection's send buffer
// is a few KiB, so a larger response is still being written while its
// client is not reading.
func limitedServer(
t *testing.T, env *sourceTestEnv, limit time.Duration,
) (*httptest.Server, *bytes.Buffer) {
t.Helper()
const sendBuffer = 4 << 10
logBuf := new(bytes.Buffer)
mw := middleware.NewForTest(
slog.New(slog.NewJSONHandler(logBuf, nil)),
&config.Config{Environment: config.EnvironmentDev},
nil,
)
srv := httptest.NewUnstartedServer(
mw.Logging()(mw.Timeout(limit)(targetRouter(env))),
)
srv.Config.WriteTimeout = limit
srv.Config.ConnContext = func(
ctx context.Context, c net.Conn,
) context.Context {
tcp, ok := c.(*net.TCPConn)
if assert.True(t, ok) {
assert.NoError(t, tcp.SetWriteBuffer(sendBuffer))
}
return ctx
}
srv.Start()
t.Cleanup(srv.Close)
return srv, logBuf
}
// TestHandleTargetDownload_OutlastsTheRequestLimit proves a download
// runs for as long as the client keeps reading, and is logged as the
// 200 it was. Behind a request limit and a server write timeout of a
// tenth of a second, the client stops reading once the response has
// started, waits three times as long, and still gets the whole file.
// The archive is larger than the connection holds, so the download is
// still being written while the client waits.
func TestHandleTargetDownload_OutlastsTheRequestLimit(t *testing.T) {
t.Parallel()
const (
limit = 100 * time.Millisecond
rows = 8
bodySize = 64 << 10
)
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
seedArchive(
t, delivery.ArchivePath(env.dbMgr, &wh, archive), rows, bodySize,
)
srv, accessLog := limitedServer(t, env, limit)
req, err := http.NewRequestWithContext(
t.Context(), http.MethodGet,
srv.URL+downloadPath(wh.ID, archive.ID), nil,
)
require.NoError(t, err)
for _, c := range env.cookies {
req.AddCookie(c)
}
resp, err := srv.Client().Do(req)
require.NoError(t, err)
defer func() { _ = resp.Body.Close() }()
require.Equal(t, http.StatusOK, resp.StatusCode)
time.Sleep(3 * limit)
zr, err := gzip.NewReader(resp.Body)
require.NoError(t, err)
var (
got map[string]json.RawMessage
events []json.RawMessage
)
require.NoError(t, json.NewDecoder(zr).Decode(&got))
require.NoError(t, json.Unmarshal(got["archived_events"], &events))
assert.Len(t, events, rows)
// Reading to the end makes the gzip reader check that the file was
// finished.
_, err = io.ReadAll(zr)
require.NoError(t, err)
// Close waits for the handler, so the access log line is written.
srv.Close()
var access map[string]any
require.NoError(t, json.Unmarshal(accessLog.Bytes(), &access))
assert.EqualValues(t, http.StatusOK, access["status"])
assert.GreaterOrEqual(t,
access["latency_ms"], float64(limit.Milliseconds()),
"the download must outlast the request limit",
)
}
// brokenWriter is a response writer whose writes fail once the
// response has started, as they do when the client goes away.
type brokenWriter struct {
*httptest.ResponseRecorder
}
func (b brokenWriter) Write(p []byte) (int, error) {
if b.Body.Len() > 0 {
return 0, errClientGone
}
return b.ResponseRecorder.Write(p)
}
// TestHandleTargetDownload_AbortsWhenItFails proves a download that
// fails after its response has started aborts the connection, so the
// client sees a failed download rather than a file that looks
// complete and does not decompress.
func TestHandleTargetDownload_AbortsWhenItFails(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
req := httptest.NewRequestWithContext(
t.Context(), http.MethodGet, downloadPath(wh.ID, archive.ID), nil,
)
for _, c := range env.cookies {
req.AddCookie(c)
}
w := brokenWriter{ResponseRecorder: httptest.NewRecorder()}
assert.PanicsWithValue(t, http.ErrAbortHandler, func() {
targetRouter(env).ServeHTTP(w, req)
})
assert.Equal(t, http.StatusOK, w.Code)
}
+6 -2
View File
@@ -37,14 +37,18 @@ const (
editAuthHeader = "Authorization: Bearer " + editBearerSecret editAuthHeader = "Authorization: Bearer " + editBearerSecret
) )
// targetRouter mounts the target create and edit routes on a chi // targetRouter mounts the target create, edit and download routes on
// router so the handlers see the URL parameters they read. // a chi router so the handlers see the URL parameters they read.
func targetRouter(env *sourceTestEnv) *chi.Mux { func targetRouter(env *sourceTestEnv) *chi.Mux {
router := chi.NewRouter() router := chi.NewRouter()
router.Post( router.Post(
"/hook/{sourceID}/targets", "/hook/{sourceID}/targets",
env.handlers.HandleTargetCreate(), env.handlers.HandleTargetCreate(),
) )
router.Get(
"/hook/{sourceID}/targets/{targetID}/download",
env.handlers.HandleTargetDownload(),
)
router.Get( router.Get(
"/hook/{sourceID}/targets/{targetID}/edit", "/hook/{sourceID}/targets/{targetID}/edit",
env.handlers.HandleTargetEdit(), env.handlers.HandleTargetEdit(),
+120 -61
View File
@@ -5,6 +5,7 @@ import (
"io" "io"
"log/slog" "log/slog"
"strings" "strings"
"sync"
"testing" "testing"
"unicode/utf8" "unicode/utf8"
@@ -18,15 +19,11 @@ import (
// width. // width.
const budget = 64 const budget = 64
// sampleRunes is how many runes wide the values in the charge test // batchRunes is how many consecutive code points the charge test logs
// are. The handlers add a constant per field — a pair of quotes when // in one value from U+1000 up. Logging each of those on its own line
// the value needs quoting — so the per-rune charge is only visible // is too slow for the suite under the race detector; 4,096 at a time
// once it is amortised over a run of them. // is 271 batches, each logged on two lines, so 542 lines per handler.
const sampleRunes = 64 const batchRunes = 4096
// quotingSlack is that constant: the pair of quotes a handler adds to
// a value that needs them and omits from one that does not.
const quotingSlack = 2
// newHandlers are the two handlers internal/logger can install. Time // newHandlers are the two handlers internal/logger can install. Time
// is dropped so a line's width is a function of its value alone — // is dropped so a line's width is a function of its value alone —
@@ -66,46 +63,48 @@ func renderedWidth(
return buf.Len() return buf.Len()
} }
// chargeTestRunes is the set of code points the charge test measures: // emittedBytes is what a handler writes for the runes of s alone, in a
// every rune in the first two planes' worth of the BMP that the // value that starts with prefix: the width of a line carrying prefix
// handlers are most likely to treat specially, the separators that // and then s twice, less that of a line carrying prefix and s once.
// only slog's JSON handler escapes, and a stratified sample across // Both values start the same way and hold the same runes, so the text
// the rest of Unicode so the astral charge is exercised on more than // handler quotes both or neither, and the quotes cancel along with the
// one hand-picked rune. // prefix and everything else on the line.
func chargeTestRunes() []rune { func emittedBytes(
const ( newHandler func(io.Writer) slog.Handler,
denseCeiling = 0x800 prefix, s string,
stride = 1021 ) int {
surrogateLo = 0xD800 return renderedWidth(newHandler, prefix+s+s) -
surrogateHi = 0xDFFF renderedWidth(newHandler, prefix+s)
) }
var runes []rune // firstUndercharged returns the first rune in s that the handler
// writes in more bytes than EncodedBytes charges for it, and how many
// runes in s are undercharged that way. It measures one rune per line,
// in a value of that rune alone and again after a space, which makes
// the text handler quote the value. The charge test calls it on the
// code points below U+1000, and from there up only on a batch that has
// already failed, to name the code points rather than just their range.
func firstUndercharged(
newHandler func(io.Writer) slog.Handler,
s string,
) (rune, int) {
first, count := rune(-1), 0
keep := func(r rune) { for _, r := range s {
if r >= surrogateLo && r <= surrogateHi { charge := logfield.EncodedBytes(r)
return if emittedBytes(newHandler, "", string(r)) <= charge &&
emittedBytes(newHandler, " ", string(r)) <= charge {
continue
} }
runes = append(runes, r) if count == 0 {
first = r
}
count++
} }
for r := range rune(denseCeiling) { return first, count
keep(r)
}
for _, r := range []rune{
0x2028, 0x2029, 0x200B, 0x4E00, 0xE000, 0xFFFD,
0x1000C, 0x1F600, 0xE0001, 0x10FFFF,
} {
keep(r)
}
for r := rune(denseCeiling); r <= utf8.MaxRune; r += stride {
keep(r)
}
return runes
} }
// TestEncodedBytes_ChargesAtLeastWhatTheHandlersEmit is the property // TestEncodedBytes_ChargesAtLeastWhatTheHandlersEmit is the property
@@ -114,33 +113,93 @@ func chargeTestRunes() []rune {
// how a stated ceiling becomes false without any test noticing, so // how a stated ceiling becomes false without any test noticing, so
// the charge is measured against what the handlers actually write // the charge is measured against what the handlers actually write
// rather than against the escaping rules as read. // rather than against the escaping rules as read.
//
// Every code point below U+1000 is checked on its own, for both
// handlers. That range holds the quote, the backslash and the control
// characters the handlers escape, next to code points each handler
// writes in fewer bytes than their charge, which in a sum would cover
// a neighbour charged too little. Each is measured in a value of it
// alone and again in one the text handler quotes, because that handler
// writes U+007F as one raw byte in a value it leaves bare but as \x7f,
// four bytes, in one it quotes.
//
// From U+1000 up the text handler writes every code point in exactly
// its charge, so the rest of Unicode is checked batchRunes at a time:
// each batch's summed charge must cover what the handler writes for
// the whole batch. The sums there can miss the JSON handler alone
// writing one code point in more bytes than its charge, when it writes
// others in the same batch in fewer.
func TestEncodedBytes_ChargesAtLeastWhatTheHandlersEmit(t *testing.T) { func TestEncodedBytes_ChargesAtLeastWhatTheHandlersEmit(t *testing.T) {
t.Parallel() t.Parallel()
var below strings.Builder
for r := range rune(0x1000) {
below.WriteRune(r)
}
var batches []string
for lo := rune(0x1000); lo <= utf8.MaxRune; lo += batchRunes {
var batch strings.Builder
for r := lo; r < lo+batchRunes; r++ {
// Surrogate halves are not runes a string can carry.
if utf8.ValidRune(r) {
batch.WriteRune(r)
}
}
batches = append(batches, batch.String())
}
// What EncodedBytes charges for each batch. Under -race -cover this
// takes longer than logging the batches, so it is worked out once,
// by whichever handler finishes logging first, while the other is
// still logging.
charged := sync.OnceValue(func() []int {
costs := make([]int, len(batches))
for i, batch := range batches {
for _, r := range batch {
costs[i] += logfield.EncodedBytes(r)
}
}
return costs
})
for name, newHandler := range newHandlers() { for name, newHandler := range newHandlers() {
t.Run(name, func(t *testing.T) { t.Run(name, func(t *testing.T) {
t.Parallel() t.Parallel()
// 'a' is a printable ASCII rune, charged exactly one if first, count := firstUndercharged(newHandler, below.String()); count > 0 {
// byte, so it is the zero point the other runes are t.Errorf(
// measured against. "%d code points below U+1000 cost more than "+
base := renderedWidth( "EncodedBytes charges, the first U+%04X",
newHandler, strings.Repeat("a", sampleRunes), count, first,
)
for _, r := range chargeTestRunes() {
got := renderedWidth(
newHandler,
strings.Repeat(string(r), sampleRunes),
) )
charged := sampleRunes * }
(logfield.EncodedBytes(r) - 1)
require.LessOrEqual( emitted := make([]int, len(batches))
t, got-base, charged+quotingSlack, for i, batch := range batches {
"U+%04X costs more on the line than "+ emitted[i] = emittedBytes(newHandler, "", batch)
"EncodedBytes charges for it", }
r,
for i, cost := range charged() {
if emitted[i] <= cost {
continue
}
lo := rune(0x1000 + i*batchRunes)
first, count := firstUndercharged(
newHandler, batches[i],
)
t.Errorf(
"U+%04X to U+%04X emit %d bytes but are "+
"charged %d; %d of them cost more than "+
"EncodedBytes charges, the first U+%04X",
lo, lo+batchRunes-1, emitted[i], cost,
count, first,
) )
} }
}) })
+38
View File
@@ -9,6 +9,7 @@ import (
"time" "time"
"go.uber.org/fx" "go.uber.org/fx"
"go.uber.org/fx/fxevent"
"sneak.berlin/go/webhooker/internal/globals" "sneak.berlin/go/webhooker/internal/globals"
) )
@@ -106,3 +107,40 @@ func (l *Logger) Identify() {
func (l *Logger) Writer() io.Writer { func (l *Logger) Writer() io.Writer {
return os.Stdout return os.Stdout
} }
// FxLogger writes fx's own events through a slog logger: how the
// dependency graph was built at DEBUG, since it repeats on every
// start; the start and stop hooks, the start itself and the signal
// that stops the service at INFO; every failure at ERROR.
//
// The formatting is fx's own fxevent.SlogLogger. That logger takes a
// single level for every event that is not a failure, so FxLogger
// holds one at each level and picks between them.
type FxLogger struct {
graph *fxevent.SlogLogger
lifecycle *fxevent.SlogLogger
}
// NewFxLogger returns an FxLogger that writes through log.
func NewFxLogger(log *slog.Logger) *FxLogger {
graph := &fxevent.SlogLogger{Logger: log}
graph.UseLogLevel(slog.LevelDebug)
lifecycle := &fxevent.SlogLogger{Logger: log}
lifecycle.UseLogLevel(slog.LevelInfo)
return &FxLogger{graph: graph, lifecycle: lifecycle}
}
// LogEvent implements fxevent.Logger.
func (f *FxLogger) LogEvent(event fxevent.Event) {
switch event.(type) {
case *fxevent.Supplied, *fxevent.Provided, *fxevent.Replaced,
*fxevent.Decorated, *fxevent.BeforeRun, *fxevent.Run,
*fxevent.Invoking, *fxevent.Invoked,
*fxevent.LoggerInitialized:
f.graph.LogEvent(event)
default:
f.lifecycle.LogEvent(event)
}
}
+54
View File
@@ -1,13 +1,23 @@
package logger_test package logger_test
import ( import (
"bytes"
"encoding/json"
"errors"
"log/slog"
"testing" "testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.uber.org/fx"
"go.uber.org/fx/fxevent"
"go.uber.org/fx/fxtest" "go.uber.org/fx/fxtest"
"sneak.berlin/go/webhooker/internal/globals" "sneak.berlin/go/webhooker/internal/globals"
"sneak.berlin/go/webhooker/internal/logger" "sneak.berlin/go/webhooker/internal/logger"
) )
var errStopHook = errors.New("stop hook failed on purpose")
func testGlobals() *globals.Globals { func testGlobals() *globals.Globals {
return &globals.Globals{ return &globals.Globals{
Appname: "test-app", Appname: "test-app",
@@ -57,3 +67,47 @@ func TestEnableDebugLogging(t *testing.T) {
// Test debug logging // Test debug logging
l.Get().Debug("debug message", "test", true) l.Get().Debug("debug message", "test", true)
} }
// TestFxLogger_Levels starts and stops an fx app that reports its own
// events through NewFxLogger, as cmd/webhooker does, and reads back
// what reached the handler: the graph at DEBUG, the start at INFO and
// a failed stop hook at ERROR, each as a structured record.
func TestFxLogger_Levels(t *testing.T) {
t.Parallel()
var out bytes.Buffer
log := slog.New(slog.NewJSONHandler(
&out, &slog.HandlerOptions{Level: slog.LevelDebug},
))
app := fx.New(
fx.WithLogger(func() fxevent.Logger {
return logger.NewFxLogger(log)
}),
fx.Invoke(func(lc fx.Lifecycle) {
lc.Append(fx.StopHook(func() error { return errStopHook }))
}),
)
require.NoError(t, app.Start(t.Context()))
require.ErrorIs(t, app.Stop(t.Context()), errStopHook)
levels := map[string]string{}
decoder := json.NewDecoder(&out)
for decoder.More() {
var record struct {
Level string `json:"level"`
Msg string `json:"msg"`
}
require.NoError(t, decoder.Decode(&record))
levels[record.Msg] = record.Level
}
assert.Equal(t, "DEBUG", levels["provided"])
assert.Equal(t, "INFO", levels["started"])
assert.Equal(t, "ERROR", levels["OnStop hook failed"])
}
+4 -1
View File
@@ -600,7 +600,10 @@ func bodyLimitedMethod(method string) bool {
} }
// MaxBodySize returns middleware that limits the size of // MaxBodySize returns middleware that limits the size of
// POST/PUT/PATCH request bodies to maxBytes. It must be registered // POST/PUT/PATCH request bodies to maxBytes. A request with any other
// method passes through uncapped, deliberately: no route behind it
// reads a body on GET, HEAD or DELETE. A handler that starts to needs
// its method added to bodyLimitedMethod first. It must be registered
// before any middleware that parses the body — notably CSRF, which // before any middleware that parses the body — notably CSRF, which
// calls r.PostFormValue — so that form parsing happens under this // calls r.PostFormValue — so that form parsing happens under this
// cap rather than net/http's 10 MB default. // cap rather than net/http's 10 MB default.
+6 -4
View File
@@ -730,10 +730,8 @@ func TestNoCache_SetsHeaders(t *testing.T) {
const testBodyLimit int64 = 64 const testBodyLimit int64 = 64
// maxBodySizeHandler wraps a sentinel handler in MaxBodySize with // maxBodySizeResult is what runMaxBodySize's sentinel handler saw,
// testBodyLimit. The sentinel records whether it ran and how much of // together with the response.
// the body it managed to read, so tests can distinguish "never
// reached" from "reached but truncated".
type maxBodySizeResult struct { type maxBodySizeResult struct {
called bool called bool
read int read int
@@ -741,6 +739,10 @@ type maxBodySizeResult struct {
response *httptest.ResponseRecorder response *httptest.ResponseRecorder
} }
// runMaxBodySize wraps a sentinel handler in MaxBodySize with
// testBodyLimit and serves req through it. The sentinel records
// whether it ran and how much of the body it managed to read, so
// tests can distinguish "never reached" from "reached but truncated".
func runMaxBodySize( func runMaxBodySize(
t *testing.T, t *testing.T,
req *http.Request, req *http.Request,
+71
View File
@@ -0,0 +1,71 @@
package middleware
import (
"context"
"errors"
"net/http"
"time"
)
// Timeout returns middleware that gives each request limit to finish:
// it cancels the request's context once limit has passed, and answers
// 504 when the handler then returns without having started its
// response.
//
// It replaces chi's middleware.Timeout, which writes that 504 even
// after the handler has sent its own status. A download that outlasts
// the limit has already sent its 200 and the whole file, so the late
// 504 changes nothing for the client: the access log and the metrics
// would record it in place of the 200, and net/http would complain of
// a superfluous WriteHeader.
func (s *Middleware) Timeout(
limit time.Duration,
) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(
w http.ResponseWriter,
r *http.Request,
) {
ctx, cancel := context.WithTimeout(r.Context(), limit)
defer cancel()
tw := &timeoutResponseWriter{ResponseWriter: w}
next.ServeHTTP(tw, r.WithContext(ctx))
if !tw.started &&
errors.Is(ctx.Err(), context.DeadlineExceeded) {
w.WriteHeader(http.StatusGatewayTimeout)
}
})
}
}
// timeoutResponseWriter records whether the handler has started its
// response.
type timeoutResponseWriter struct {
http.ResponseWriter
started bool
}
func (w *timeoutResponseWriter) WriteHeader(code int) {
w.started = true
w.ResponseWriter.WriteHeader(code)
}
func (w *timeoutResponseWriter) Write(b []byte) (int, error) {
// A Write without a WriteHeader starts the response too: net/http
// sends 200 in front of it.
w.started = true
//nolint:wrapcheck // Pass the writer's own error through unchanged.
return w.ResponseWriter.Write(b)
}
// Unwrap lets http.ResponseController reach the writer underneath, so
// a handler can still set a write deadline through this wrapper.
func (w *timeoutResponseWriter) Unwrap() http.ResponseWriter {
return w.ResponseWriter
}
+56
View File
@@ -0,0 +1,56 @@
package middleware_test
import (
"net/http"
"net/http/httptest"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// TestTimeout proves the request limit answers 504 to a handler that
// outlasts it without starting its response, and leaves a response the
// handler has started with the status it sent. Both are what the
// access log records.
func TestTimeout(t *testing.T) {
t.Parallel()
const limit = 10 * time.Millisecond
for _, tc := range []struct {
name string
sent int // the status the handler sends, or 0 for none
want int
}{
{name: "not started", sent: 0, want: http.StatusGatewayTimeout},
{name: "started", sent: http.StatusOK, want: http.StatusOK},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
m, buf := capturingMiddleware(t)
handler := m.Logging()(m.Timeout(limit)(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
if tc.sent != 0 {
w.WriteHeader(tc.sent)
}
<-r.Context().Done()
},
)))
w := httptest.NewRecorder()
handler.ServeHTTP(w, httptest.NewRequestWithContext(
t.Context(), http.MethodGet, "/", nil,
))
assert.Equal(t, tc.want, w.Code)
entries := accessLogEntries(t, buf)
require.Len(t, entries, 1)
assert.EqualValues(t, tc.want, entries[0]["status"])
})
}
}
+69
View File
@@ -0,0 +1,69 @@
package server_test
import (
"compress/gzip"
"encoding/json"
"net/http"
"regexp"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"gorm.io/gorm/clause"
"sneak.berlin/go/webhooker/internal/database"
)
// TestHook_DownloadArchive follows the Download link the webhook page
// shows for a database target, and only for it, and gets the archive
// as a gzipped JSON file. Signed out, the link leads to the login page.
func TestHook_DownloadArchive(t *testing.T) {
t.Parallel()
env := newTestEnv(t)
userID, _ := env.seedUser(t, "archivist", "somepassword")
cookies := env.authCookies(t, userID, "archivist")
wh := env.seedWebhook(t, userID)
env.seedTarget(t, wh.ID)
archive := &database.Target{
WebhookID: wh.ID,
Name: "kept",
Type: database.TargetTypeDatabase,
Active: true,
}
require.NoError(t,
env.db.DB().Omit(clause.Associations).Create(archive).Error,
)
page := env.get("/hook/"+wh.ID, cookies)
require.Equal(t, http.StatusOK, page.Code)
links := regexp.MustCompile(
`href="(/hook/[^/"]+/targets/[^/"]+/download)"`,
).FindAllStringSubmatch(page.Body.String(), -1)
require.Len(t, links, 1, "only the database target has a Download")
link := links[0][1]
assert.Equal(t,
"/hook/"+wh.ID+"/targets/"+archive.ID+"/download", link,
)
w := env.get(link, cookies)
require.Equal(t, http.StatusOK, w.Code)
assert.Equal(t, "application/gzip", w.Header().Get("Content-Type"))
zr, err := gzip.NewReader(w.Body)
require.NoError(t, err)
var got map[string]json.RawMessage
require.NoError(t, json.NewDecoder(zr).Decode(&got))
assert.JSONEq(t,
`{"id":"`+wh.ID+`","name":"routed"}`, string(got["webhook"]),
)
w = env.get(link, nil)
assert.Equal(t, http.StatusSeeOther, w.Code)
assert.Contains(t, w.Header().Get("Location"), "/pages/login")
}
+2 -2
View File
@@ -191,7 +191,7 @@ func TestErrorPage_PanicOnAdminPage(t *testing.T) {
w := serve( w := serve(
server.NewRouterWithPageProbeForTest( server.NewRouterWithPageProbeForTest(
env.log.Get(), env.cfg, env.mw, env.hnd, t, env.log, env.cfg, env.mw, env.hnd,
true, panicProbeHandler, true, panicProbeHandler,
), ),
server.PageProbePattern, server.PageProbePattern,
@@ -200,7 +200,7 @@ func TestErrorPage_PanicOnAdminPage(t *testing.T) {
w = serve( w = serve(
server.NewRouterWithProbeForTest( server.NewRouterWithProbeForTest(
env.log.Get(), env.cfg, env.mw, env.hnd, t, env.log, env.cfg, env.mw, env.hnd,
true, panicProbeHandler, true, panicProbeHandler,
), ),
server.ProbePattern, server.ProbePattern,
+47 -26
View File
@@ -1,13 +1,16 @@
package server package server
import ( import (
"log/slog"
"net/http" "net/http"
"testing"
"github.com/getsentry/sentry-go" "github.com/getsentry/sentry-go"
"github.com/go-chi/chi" "github.com/go-chi/chi"
"github.com/stretchr/testify/require"
"go.uber.org/fx/fxtest"
"sneak.berlin/go/webhooker/internal/config" "sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/handlers" "sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/middleware" "sneak.berlin/go/webhooker/internal/middleware"
) )
@@ -34,23 +37,45 @@ func SentryClientOptionsForTest(
return sentryClientOptions(dsn, release) return sentryClientOptions(dsn, release)
} }
// newServerForTest builds a Server through New, as the application
// does, on a lifecycle that is never started: the hooks New adds to
// it never run, so nothing listens.
func newServerForTest(
t *testing.T,
log *logger.Logger,
cfg *config.Config,
mw *middleware.Middleware,
h *handlers.Handlers,
) *Server {
t.Helper()
s, err := New(fxtest.NewLifecycle(t), ServerParams{
Logger: log,
Config: cfg,
Middleware: mw,
Handlers: h,
})
require.NoError(t, err)
return s
}
// NewRouterForTest builds the real route tree via SetupRoutes with // NewRouterForTest builds the real route tree via SetupRoutes with
// the supplied middleware and handlers, bypassing the fx lifecycle // the supplied middleware and handlers, on a Server from New whose
// and the HTTP listener. Tests use it so that route-group middleware // lifecycle is never started, so no HTTP listener runs. Tests use it
// registration order is exercised exactly as it ships, rather than // so that route-group middleware registration order is exercised
// against a hand-rebuilt chain that could drift from routes.go. // exactly as it ships, rather than against a hand-rebuilt chain that
// could drift from routes.go.
func NewRouterForTest( func NewRouterForTest(
log *slog.Logger, t *testing.T,
log *logger.Logger,
cfg *config.Config, cfg *config.Config,
mw *middleware.Middleware, mw *middleware.Middleware,
h *handlers.Handlers, h *handlers.Handlers,
) http.Handler { ) http.Handler {
s := &Server{ t.Helper()
log: log,
mw: mw, s := newServerForTest(t, log, cfg, mw, h)
h: h,
params: ServerParams{Config: cfg},
}
s.SetupRoutes() s.SetupRoutes()
return s.router return s.router
@@ -83,19 +108,17 @@ const ProbePattern = "/probe"
// option and the recoverer registered outside it is the thing a test // option and the recoverer registered outside it is the thing a test
// has to be able to pin. // has to be able to pin.
func NewRouterWithProbeForTest( func NewRouterWithProbeForTest(
log *slog.Logger, t *testing.T,
log *logger.Logger,
cfg *config.Config, cfg *config.Config,
mw *middleware.Middleware, mw *middleware.Middleware,
h *handlers.Handlers, h *handlers.Handlers,
sentryEnabled bool, sentryEnabled bool,
probe http.HandlerFunc, probe http.HandlerFunc,
) http.Handler { ) http.Handler {
s := &Server{ t.Helper()
log: log,
mw: mw, s := newServerForTest(t, log, cfg, mw, h)
h: h,
params: ServerParams{Config: cfg},
}
s.sentryEnabled.Store(sentryEnabled) s.sentryEnabled.Store(sentryEnabled)
s.SetupRoutes() s.SetupRoutes()
s.router.Handle(ProbePattern, probe) s.router.Handle(ProbePattern, probe)
@@ -113,19 +136,17 @@ const PageProbePattern = "/pages/probe"
// it, so the probe runs behind that group's own middleware exactly as // it, so the probe runs behind that group's own middleware exactly as
// the group's real routes do. // the group's real routes do.
func NewRouterWithPageProbeForTest( func NewRouterWithPageProbeForTest(
log *slog.Logger, t *testing.T,
log *logger.Logger,
cfg *config.Config, cfg *config.Config,
mw *middleware.Middleware, mw *middleware.Middleware,
h *handlers.Handlers, h *handlers.Handlers,
sentryEnabled bool, sentryEnabled bool,
probe http.HandlerFunc, probe http.HandlerFunc,
) http.Handler { ) http.Handler {
s := &Server{ t.Helper()
log: log,
mw: mw, s := newServerForTest(t, log, cfg, mw, h)
h: h,
params: ServerParams{Config: cfg},
}
s.sentryEnabled.Store(sentryEnabled) s.sentryEnabled.Store(sentryEnabled)
s.SetupRoutes() s.SetupRoutes()
+2 -2
View File
@@ -199,7 +199,7 @@ func TestPanicProbeChild(t *testing.T) {
env := newTestEnv(t) env := newTestEnv(t)
router := server.NewRouterWithProbeForTest( router := server.NewRouterWithProbeForTest(
env.log.Get(), env.cfg, env.mw, env.hnd, t, env.log, env.cfg, env.mw, env.hnd,
false, panicProbeHandler, false, panicProbeHandler,
) )
@@ -253,7 +253,7 @@ func TestSentryStillSeesAPanic(t *testing.T) {
require.NoError(t, err) require.NoError(t, err)
router := server.NewRouterWithProbeForTest( router := server.NewRouterWithProbeForTest(
env.log.Get(), env.cfg, env.mw, env.hnd, t, env.log, env.cfg, env.mw, env.hnd,
true, panicProbeHandler, true, panicProbeHandler,
) )
+2 -2
View File
@@ -67,11 +67,11 @@ func TestResponseControllerThroughProductionRouter(t *testing.T) {
routers := map[string]http.Handler{ routers := map[string]http.Handler{
server.ProbePattern: server.NewRouterWithProbeForTest( server.ProbePattern: server.NewRouterWithProbeForTest(
env.log.Get(), env.cfg, env.mw, env.hnd, t, env.log, env.cfg, env.mw, env.hnd,
tc.sentryEnabled, probe, tc.sentryEnabled, probe,
), ),
server.PageProbePattern: server.NewRouterWithPageProbeForTest( server.PageProbePattern: server.NewRouterWithPageProbeForTest(
env.log.Get(), env.cfg, env.mw, env.hnd, t, env.log, env.cfg, env.mw, env.hnd,
tc.sentryEnabled, probe, tc.sentryEnabled, probe,
), ),
} }
+5 -1
View File
@@ -73,7 +73,7 @@ func (s *Server) setupGlobalMiddleware() {
} }
s.router.Use(s.mw.CORS()) s.router.Use(s.mw.CORS())
s.router.Use(middleware.Timeout(requestTimeout)) s.router.Use(s.mw.Timeout(requestTimeout))
// Panic recovery, deliberately here rather than first. It has to // Panic recovery, deliberately here rather than first. It has to
// run inside every middleware that observes the response, so the // run inside every middleware that observes the response, so the
@@ -312,6 +312,10 @@ func (s *Server) setupSourceRoutes() {
"/targets/{targetID}/edit", "/targets/{targetID}/edit",
s.h.HandleTargetEditSubmit(), s.h.HandleTargetEditSubmit(),
) )
r.Get(
"/targets/{targetID}/download",
s.h.HandleTargetDownload(),
)
r.Post( r.Post(
"/targets/{targetID}/delete", "/targets/{targetID}/delete",
s.h.HandleTargetDelete(), s.h.HandleTargetDelete(),
+44 -2
View File
@@ -136,7 +136,7 @@ func newTestEnvWithConfig(
t.Cleanup(app.RequireStop) t.Cleanup(app.RequireStop)
return &testEnv{ return &testEnv{
router: server.NewRouterForTest(log.Get(), cfg, mw, hnd), router: server.NewRouterForTest(t, log, cfg, mw, hnd),
sess: sess, sess: sess,
db: db, db: db,
dbMgr: dbMgr, dbMgr: dbMgr,
@@ -531,6 +531,48 @@ func TestStaticServesOnlyGetAndHead(t *testing.T) {
} }
} }
// --- every page route group ---
// TestPageRouteGroups_OversizeBody_RejectedBeforeCSRF pins the body
// cap ahead of CSRF and RequireAuth in every page route group. The
// requests carry no session and no CSRF token, so if either ran first
// the answer would be a 403 or a redirect to the login page rather
// than 413, and CSRF would issue its cookie (see
// TestPagesLogin_UnderLimit_NoToken_CSRFRejects). /settings has no
// POST route, but its group's middleware runs before the method is
// matched, so a POST there still reaches CSRF's form parsing if the
// cap moves after it. The user and webhook in the paths need not
// exist: nothing after the cap runs.
func TestPageRouteGroups_OversizeBody_RejectedBeforeCSRF(
t *testing.T,
) {
t.Parallel()
env := newTestEnv(t)
form := url.Values{}
form.Set("name", oversizeValue())
for _, path := range []string{
"/pages/login",
"/user/nobody/password",
"/settings/",
"/hooks/new",
"/hook/nonexistent/edit",
} {
w := env.post(path, form, nil)
assert.Equal(
t, http.StatusRequestEntityTooLarge, w.Code, path,
)
assert.False(
t, csrfCookieSet(w),
"CSRF middleware must not run for an oversized body to %s",
path,
)
}
}
// --- /pages group --- // --- /pages group ---
// TestPagesLogin_OversizeBody_RejectedBeforeCSRF proves the cap runs // TestPagesLogin_OversizeBody_RejectedBeforeCSRF proves the cap runs
@@ -1610,7 +1652,7 @@ func TestTwoMetricsRoutersInOneProcess(t *testing.T) {
) )
third := &testEnv{ third := &testEnv{
router: server.NewRouterForTest( router: server.NewRouterForTest(
first.log.Get(), first.cfg, first.mw, first.hnd, t, first.log, first.cfg, first.mw, first.hnd,
), ),
} }
+3
View File
@@ -157,6 +157,9 @@
{{else}} {{else}}
<span class="badge-error">Inactive</span> <span class="badge-error">Inactive</span>
{{end}} {{end}}
{{if eq .Type "database"}}
<a href="/hook/{{$.Webhook.ID}}/targets/{{.ID}}/download" class="text-xs text-gray-500 hover:text-primary-600" title="Download the archive as gzipped JSON">Download</a>
{{end}}
<a href="/hook/{{$.Webhook.ID}}/targets/{{.ID}}/edit" class="text-xs text-gray-500 hover:text-primary-600" title="Edit">Edit</a> <a href="/hook/{{$.Webhook.ID}}/targets/{{.ID}}/edit" class="text-xs text-gray-500 hover:text-primary-600" title="Edit">Edit</a>
<form method="POST" action="/hook/{{$.Webhook.ID}}/targets/{{.ID}}/toggle" class="inline"> <form method="POST" action="/hook/{{$.Webhook.ID}}/targets/{{.ID}}/toggle" class="inline">
<input type="hidden" name="csrf_token" value="{{$.CSRFToken}}"> <input type="hidden" name="csrf_token" value="{{$.CSRFToken}}">