Compare commits
3 Commits
b8c8b75e04
...
dfd559417e
| Author | SHA1 | Date | |
|---|---|---|---|
| dfd559417e | |||
| a13e5b7ded | |||
| bb30b3ad64 |
38
README.md
38
README.md
@@ -107,14 +107,26 @@ TTY detection, and security headers are always applied.
|
|||||||
| `DATA_DIR` | Directory for all SQLite databases | `/var/lib/webhooker` |
|
| `DATA_DIR` | Directory for all SQLite databases | `/var/lib/webhooker` |
|
||||||
| `DEBUG` | Enable debug logging | `false` |
|
| `DEBUG` | Enable debug logging | `false` |
|
||||||
| `MAINTENANCE_MODE` | Report `maintenanceMode: true` in the healthcheck JSON. It does not change how any request is served — no maintenance page exists | `false` |
|
| `MAINTENANCE_MODE` | Report `maintenanceMode: true` in the healthcheck JSON. It does not change how any request is served — no maintenance page exists | `false` |
|
||||||
| `METRICS_USERNAME` | Basic auth username for `/metrics` | `""` |
|
| `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` | `""` |
|
| `METRICS_PASSWORD` | Basic auth password for `/metrics`. Must be set together with `METRICS_USERNAME`; one without the other fails startup | `""` |
|
||||||
| `SENTRY_DSN` | Sentry error reporting DSN | `""` |
|
| `SENTRY_DSN` | Sentry error reporting DSN | `""` |
|
||||||
| `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) | `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 (unset: all clients behind a proxy share one rate-limit bucket; a correct login password is never throttled either way) | `""` (none) |
|
| `TRUSTED_PROXIES` | CIDRs whose forwarded headers are trusted (unset: all clients behind a proxy share one rate-limit bucket; a correct login password is never throttled either way) | `""` (none) |
|
||||||
|
|
||||||
|
#### Metrics credentials
|
||||||
|
|
||||||
|
`METRICS_USERNAME` and `METRICS_PASSWORD` are set together or not at
|
||||||
|
all. With both set, `/metrics` is served behind basic auth. With
|
||||||
|
neither set, the route is not registered and returns 404. With one set
|
||||||
|
and the other empty or unset, the process refuses to start and exits
|
||||||
|
non-zero with an error naming both variables — mounting the endpoint
|
||||||
|
on the username alone would publish it behind a password that is the
|
||||||
|
empty string, and quietly withholding it would deny an endpoint that
|
||||||
|
was asked for. The `hasMetricsAuth` field in the startup log and the
|
||||||
|
existence of the route are the same value, so they cannot disagree.
|
||||||
|
|
||||||
#### Trusted proxies
|
#### Trusted proxies
|
||||||
|
|
||||||
`TRUSTED_PROXIES` is a comma-separated list of CIDR blocks (a bare
|
`TRUSTED_PROXIES` is a comma-separated list of CIDR blocks (a bare
|
||||||
@@ -1134,11 +1146,11 @@ arriving and being stored, they are just not getting anywhere.
|
|||||||
| Metric | Type | Meaning |
|
| Metric | Type | Meaning |
|
||||||
| ------ | ---- | ------- |
|
| ------ | ---- | ------- |
|
||||||
| `webhooker_events_received_total` | counter | Events received and durably stored. Compare against the delivery counters on one dashboard |
|
| `webhooker_events_received_total` | counter | Events received and durably stored. Compare against the delivery counters on one dashboard |
|
||||||
| `webhooker_delivery_attempts_total` | counter | Delivery attempts dispatched to a target |
|
| `webhooker_delivery_attempts_total` | counter | Delivery attempts actually dispatched to a target. A delivery an open circuit breaker refused is not one: it is counted as a retry instead |
|
||||||
| `webhooker_deliveries_succeeded_total` | counter | Deliveries that reached `delivered` |
|
| `webhooker_deliveries_succeeded_total` | counter | Deliveries that reached `delivered` |
|
||||||
| `webhooker_deliveries_failed_total` | counter | Deliveries that failed terminally and will not be retried |
|
| `webhooker_deliveries_failed_total` | counter | Deliveries that failed terminally and will not be retried |
|
||||||
| `webhooker_delivery_retries_total` | counter | Deliveries put back into `retrying` |
|
| `webhooker_delivery_retries_total` | counter | Deliveries put back into `retrying` |
|
||||||
| `webhooker_delivery_duration_seconds` | histogram | Wall time of a single delivery attempt |
|
| `webhooker_delivery_duration_seconds` | histogram | Wall time of a single dispatched delivery attempt, the same duration the attempt's `DeliveryResult` records |
|
||||||
| `webhooker_deliveries_pending` | gauge | Deliveries currently in `pending` |
|
| `webhooker_deliveries_pending` | gauge | Deliveries currently in `pending` |
|
||||||
| `webhooker_deliveries_retrying` | gauge | Deliveries currently in `retrying` |
|
| `webhooker_deliveries_retrying` | gauge | Deliveries currently in `retrying` |
|
||||||
| `webhooker_circuit_breakers_open` | gauge | Circuit breakers currently open |
|
| `webhooker_circuit_breakers_open` | gauge | Circuit breakers currently open |
|
||||||
@@ -1160,6 +1172,14 @@ delta would have to be seeded at startup from rows a previous process
|
|||||||
wrote, and would drift permanently on any transition that failed to
|
wrote, and would drift permanently on any transition that failed to
|
||||||
persist.
|
persist.
|
||||||
|
|
||||||
|
Those two gauges also publish an `unknown` series, from startup rather
|
||||||
|
than on first occurrence. Deliveries queued against a target that has
|
||||||
|
since been deleted are counted there: that backlog is the one nobody is
|
||||||
|
watching, so it is the one that must not silently vanish from the
|
||||||
|
gauge. The outcome counters move only after the status change has been
|
||||||
|
written, so a transition the database rejected is never reported as an
|
||||||
|
outcome that happened.
|
||||||
|
|
||||||
### Rate Limiting
|
### Rate Limiting
|
||||||
|
|
||||||
Global blanket rate limiting middleware (e.g., a per-IP throttle shared
|
Global blanket rate limiting middleware (e.g., a per-IP throttle shared
|
||||||
@@ -1720,7 +1740,7 @@ abuse limit later; they are tracked as future work.
|
|||||||
|
|
||||||
| Method | Path | Description |
|
| Method | Path | Description |
|
||||||
| ------ | ---------- | ----------- |
|
| ------ | ---------- | ----------- |
|
||||||
| `GET` | `/metrics` | Prometheus metrics, behind basic auth. The route is registered only when `METRICS_USERNAME` is set; otherwise it does not exist and returns 404 |
|
| `GET` | `/metrics` | Prometheus metrics, behind basic auth. The route is registered only when `METRICS_USERNAME` and `METRICS_PASSWORD` are both set; with neither set it does not exist and returns 404, and with only one set the process refuses to start |
|
||||||
|
|
||||||
#### API (Planned)
|
#### API (Planned)
|
||||||
|
|
||||||
@@ -1878,7 +1898,8 @@ Applied to all routes in this order:
|
|||||||
Permissions-Policy)
|
Permissions-Policy)
|
||||||
3. **Logging** — Structured request logging (method, URL, status,
|
3. **Logging** — Structured request logging (method, URL, status,
|
||||||
latency, remote IP, user agent, request ID)
|
latency, remote IP, user agent, request ID)
|
||||||
4. **Metrics** — Prometheus HTTP metrics (if `METRICS_USERNAME` is set)
|
4. **Metrics** — Prometheus HTTP metrics (if `METRICS_USERNAME` and
|
||||||
|
`METRICS_PASSWORD` are both set)
|
||||||
5. **CORS** — Cross-origin resource sharing headers
|
5. **CORS** — Cross-origin resource sharing headers
|
||||||
6. **Timeout** — 60-second request timeout
|
6. **Timeout** — 60-second request timeout
|
||||||
7. **Recoverer** — Panic recovery: one `ERROR` record through
|
7. **Recoverer** — Panic recovery: one `ERROR` record through
|
||||||
@@ -1910,8 +1931,9 @@ being read and without reaching CSRF, the route group's remaining
|
|||||||
middleware, or the handler. It is not rejected before *any* other
|
middleware, or the handler. It is not rejected before *any* other
|
||||||
middleware, though: the global entries listed above all run first, so
|
middleware, though: the global entries listed above all run first, so
|
||||||
such a request is still logged and given the security headers — and
|
such a request is still logged and given the security headers — and
|
||||||
counted in the metrics, on a deployment where `METRICS_USERNAME` is
|
counted in the metrics, on a deployment where the `/metrics`
|
||||||
set and the Metrics middleware is therefore registered at all. The
|
credentials are set and the Metrics middleware is therefore registered
|
||||||
|
at all. The
|
||||||
rejection itself is logged at `WARN` with the method, path and
|
rejection itself is logged at `WARN` with the method, path and
|
||||||
declared length. A chunked request, or
|
declared length. A chunked request, or
|
||||||
one that lies about its length, is hard-capped by
|
one that lies about its length, is hard-capped by
|
||||||
|
|||||||
@@ -71,6 +71,16 @@ var ErrInvalidPort = errors.New("invalid port")
|
|||||||
// nor a bare IP address.
|
// nor a bare IP address.
|
||||||
var ErrInvalidCIDR = errors.New("invalid CIDR")
|
var ErrInvalidCIDR = errors.New("invalid CIDR")
|
||||||
|
|
||||||
|
// ErrIncompleteMetricsAuth is returned when exactly one of
|
||||||
|
// METRICS_USERNAME and METRICS_PASSWORD carries a value. Neither
|
||||||
|
// fallback is acceptable: serving /metrics on the username alone
|
||||||
|
// publishes an endpoint whose password is the empty string, and
|
||||||
|
// silently leaving it unmounted withholds an endpoint the operator
|
||||||
|
// asked for. Half-set is a configuration error, so startup fails.
|
||||||
|
var ErrIncompleteMetricsAuth = errors.New(
|
||||||
|
"incomplete metrics credentials",
|
||||||
|
)
|
||||||
|
|
||||||
//nolint:revive // ConfigParams is a standard fx naming convention.
|
//nolint:revive // ConfigParams is a standard fx naming convention.
|
||||||
type ConfigParams struct {
|
type ConfigParams struct {
|
||||||
fx.In
|
fx.In
|
||||||
@@ -128,6 +138,21 @@ func (c *Config) IsProd() bool {
|
|||||||
return c.Environment == EnvironmentProd
|
return c.Environment == EnvironmentProd
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// MetricsAuthEnabled reports whether /metrics is served behind basic
|
||||||
|
// auth. It is the only answer to that question in the codebase: the
|
||||||
|
// route mount, the Prometheus recording middleware and the startup
|
||||||
|
// log's hasMetricsAuth field all read this one method, so the log
|
||||||
|
// cannot report auth as off while the route is mounted.
|
||||||
|
//
|
||||||
|
// It requires both credentials rather than the username alone.
|
||||||
|
// loadFromEnv already rejects a half-set pair, but a Config built in
|
||||||
|
// code bypasses that, and the failure mode this guards is an endpoint
|
||||||
|
// mounted with a credential map whose only password is the empty
|
||||||
|
// string.
|
||||||
|
func (c *Config) MetricsAuthEnabled() bool {
|
||||||
|
return c.MetricsUsername != "" && c.MetricsPassword != ""
|
||||||
|
}
|
||||||
|
|
||||||
// envString returns the value of the named environment variable,
|
// envString returns the value of the named environment variable,
|
||||||
// or an empty string if not set.
|
// or an empty string if not set.
|
||||||
func envString(key string) string {
|
func envString(key string) string {
|
||||||
@@ -329,6 +354,30 @@ func envPrefixList(key string) ([]netip.Prefix, error) {
|
|||||||
return prefixes, nil
|
return prefixes, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// resolveMetricsAuth reads the /metrics basic-auth credentials and
|
||||||
|
// rejects a half-set pair, naming both variables either way. The
|
||||||
|
// error carries neither value: the password is a secret.
|
||||||
|
func resolveMetricsAuth() (string, string, error) {
|
||||||
|
username := envString("METRICS_USERNAME")
|
||||||
|
password := envString("METRICS_PASSWORD")
|
||||||
|
|
||||||
|
if (username == "") == (password == "") {
|
||||||
|
return username, password, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
set, empty := "METRICS_USERNAME", "METRICS_PASSWORD"
|
||||||
|
if username == "" {
|
||||||
|
set, empty = empty, set
|
||||||
|
}
|
||||||
|
|
||||||
|
return "", "", fmt.Errorf(
|
||||||
|
"%w: %s is set but %s is empty; METRICS_USERNAME and "+
|
||||||
|
"METRICS_PASSWORD must both be set to serve /metrics, "+
|
||||||
|
"or both be empty to leave it unmounted",
|
||||||
|
ErrIncompleteMetricsAuth, set, empty,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
// resolveEnvironment reads WEBHOOKER_ENVIRONMENT, defaulting to
|
// resolveEnvironment reads WEBHOOKER_ENVIRONMENT, defaulting to
|
||||||
// dev, and rejects unrecognised values.
|
// dev, and rejects unrecognised values.
|
||||||
func resolveEnvironment() (string, error) {
|
func resolveEnvironment() (string, error) {
|
||||||
@@ -406,13 +455,18 @@ func loadFromEnv() (*Config, error) {
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
metricsUsername, metricsPassword, err := resolveMetricsAuth()
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
return &Config{
|
return &Config{
|
||||||
DataDir: envString("DATA_DIR"),
|
DataDir: envString("DATA_DIR"),
|
||||||
Debug: debug,
|
Debug: debug,
|
||||||
MaintenanceMode: maintenanceMode,
|
MaintenanceMode: maintenanceMode,
|
||||||
Environment: environment,
|
Environment: environment,
|
||||||
MetricsUsername: envString("METRICS_USERNAME"),
|
MetricsUsername: metricsUsername,
|
||||||
MetricsPassword: envString("METRICS_PASSWORD"),
|
MetricsPassword: metricsPassword,
|
||||||
Port: port,
|
Port: port,
|
||||||
SentryDSN: envString("SENTRY_DSN"),
|
SentryDSN: envString("SENTRY_DSN"),
|
||||||
RetentionSweepInterval: retentionSweepInterval,
|
RetentionSweepInterval: retentionSweepInterval,
|
||||||
@@ -512,8 +566,7 @@ func New(lc fx.Lifecycle, params ConfigParams) (*Config, error) {
|
|||||||
"receiverRateLimit", s.ReceiverRateLimit,
|
"receiverRateLimit", s.ReceiverRateLimit,
|
||||||
"trustedProxies", len(s.TrustedProxies),
|
"trustedProxies", len(s.TrustedProxies),
|
||||||
"hasSentryDSN", s.SentryDSN != "",
|
"hasSentryDSN", s.SentryDSN != "",
|
||||||
"hasMetricsAuth",
|
"hasMetricsAuth", s.MetricsAuthEnabled(),
|
||||||
s.MetricsUsername != "" && s.MetricsPassword != "",
|
|
||||||
)
|
)
|
||||||
|
|
||||||
s.warnSharedRateLimitBucket(log)
|
s.warnSharedRateLimitBucket(log)
|
||||||
|
|||||||
@@ -26,6 +26,12 @@ const (
|
|||||||
// cidrPrivateV4 is the sample trusted-proxy block the
|
// cidrPrivateV4 is the sample trusted-proxy block the
|
||||||
// TRUSTED_PROXIES cases are built from.
|
// TRUSTED_PROXIES cases are built from.
|
||||||
cidrPrivateV4 = "10.0.0.0/8"
|
cidrPrivateV4 = "10.0.0.0/8"
|
||||||
|
|
||||||
|
// metricsAuthValue is the sample METRICS_PASSWORD the metrics
|
||||||
|
// credential cases are built from. It is asserted absent from
|
||||||
|
// the startup error, so it must not be a substring of either
|
||||||
|
// variable name that error prints.
|
||||||
|
metricsAuthValue = "s3cret"
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestEnvironmentConfig(t *testing.T) {
|
func TestEnvironmentConfig(t *testing.T) {
|
||||||
@@ -726,3 +732,168 @@ func TestSharedRateLimitBucketWarning(t *testing.T) {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// metricsEnv describes what one subtest below puts in the
|
||||||
|
// environment for a single METRICS_ variable. A variable that is
|
||||||
|
// set to the empty string and one that is not set at all are
|
||||||
|
// distinct inputs here, because the reported bug arrived through
|
||||||
|
// the first of them.
|
||||||
|
type metricsEnv struct {
|
||||||
|
set bool
|
||||||
|
value string
|
||||||
|
}
|
||||||
|
|
||||||
|
// unset leaves the variable out of the environment entirely.
|
||||||
|
func unset() metricsEnv {
|
||||||
|
return metricsEnv{set: false, value: ""}
|
||||||
|
}
|
||||||
|
|
||||||
|
// setTo sets the variable, including to the empty string.
|
||||||
|
func setTo(value string) metricsEnv {
|
||||||
|
return metricsEnv{set: true, value: value}
|
||||||
|
}
|
||||||
|
|
||||||
|
// metricsAuthCase is one row of the table in TestMetricsAuthConfig,
|
||||||
|
// named so the table can live in its own function and keep the test
|
||||||
|
// itself short.
|
||||||
|
type metricsAuthCase struct {
|
||||||
|
name string
|
||||||
|
username metricsEnv
|
||||||
|
password metricsEnv
|
||||||
|
expectError bool
|
||||||
|
expectAuth bool
|
||||||
|
}
|
||||||
|
|
||||||
|
// metricsAuthCases enumerates every combination of the two
|
||||||
|
// credentials, counting "set to the empty string" and "not set at
|
||||||
|
// all" as separate inputs on each side.
|
||||||
|
func metricsAuthCases() []metricsAuthCase {
|
||||||
|
return []metricsAuthCase{
|
||||||
|
{
|
||||||
|
name: "both unset leaves metrics unmounted",
|
||||||
|
username: unset(),
|
||||||
|
password: unset(),
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "both empty leaves metrics unmounted",
|
||||||
|
username: setTo(""),
|
||||||
|
password: setTo(""),
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "both set enables metrics auth",
|
||||||
|
username: setTo("metrics"),
|
||||||
|
password: setTo(metricsAuthValue),
|
||||||
|
expectAuth: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "username with unset password fails",
|
||||||
|
username: setTo("metrics"),
|
||||||
|
password: unset(),
|
||||||
|
expectError: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "username with empty password fails",
|
||||||
|
username: setTo("metrics"),
|
||||||
|
password: setTo(""),
|
||||||
|
expectError: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "password with unset username fails",
|
||||||
|
username: unset(),
|
||||||
|
password: setTo(metricsAuthValue),
|
||||||
|
expectError: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "password with empty username fails",
|
||||||
|
username: setTo(""),
|
||||||
|
password: setTo(metricsAuthValue),
|
||||||
|
expectError: true,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestMetricsAuthConfig covers every combination of METRICS_USERNAME
|
||||||
|
// and METRICS_PASSWORD. Either both carry a value, in which case
|
||||||
|
// /metrics is served behind basic auth, or neither does, in which
|
||||||
|
// case the route is never mounted. One without the other is a
|
||||||
|
// startup error rather than a fallback: mounting on the username
|
||||||
|
// alone published /metrics behind a credential map that accepted an
|
||||||
|
// empty password, which is the defect this test exists to pin. See
|
||||||
|
// https://git.eeqj.de/sneak/webhooker/issues/205.
|
||||||
|
func TestMetricsAuthConfig(t *testing.T) {
|
||||||
|
for _, tt := range metricsAuthCases() {
|
||||||
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
|
// Cannot use t.Parallel() here because t.Setenv
|
||||||
|
// is incompatible with parallel subtests.
|
||||||
|
if tt.username.set {
|
||||||
|
t.Setenv("METRICS_USERNAME", tt.username.value)
|
||||||
|
} else {
|
||||||
|
require.NoError(
|
||||||
|
t, os.Unsetenv("METRICS_USERNAME"),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
if tt.password.set {
|
||||||
|
t.Setenv("METRICS_PASSWORD", tt.password.value)
|
||||||
|
} else {
|
||||||
|
require.NoError(
|
||||||
|
t, os.Unsetenv("METRICS_PASSWORD"),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
if tt.expectError {
|
||||||
|
assertMetricsAuthRejected(t)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
assertMetricsAuthAccepted(t, tt.expectAuth)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// assertMetricsAuthRejected requires that fx refused to build the
|
||||||
|
// graph, that the failure is ErrIncompleteMetricsAuth, and that the
|
||||||
|
// operator is told both variable names — the point of failing here
|
||||||
|
// rather than degrading is that the message says what to fix.
|
||||||
|
func assertMetricsAuthRejected(t *testing.T) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
var cfg *config.Config
|
||||||
|
|
||||||
|
app := fx.New(
|
||||||
|
fx.NopLogger,
|
||||||
|
fx.Provide(globals.New, logger.New, config.New),
|
||||||
|
fx.Populate(&cfg),
|
||||||
|
)
|
||||||
|
|
||||||
|
err := app.Err()
|
||||||
|
require.Error(t, err)
|
||||||
|
require.ErrorIs(t, err, config.ErrIncompleteMetricsAuth)
|
||||||
|
assert.Contains(t, err.Error(), "METRICS_USERNAME")
|
||||||
|
assert.Contains(t, err.Error(), "METRICS_PASSWORD")
|
||||||
|
// The password is a secret and must not reach a startup error.
|
||||||
|
assert.NotContains(t, err.Error(), metricsAuthValue)
|
||||||
|
}
|
||||||
|
|
||||||
|
// assertMetricsAuthAccepted requires that startup succeeded and that
|
||||||
|
// MetricsAuthEnabled — the single value the /metrics mount and the
|
||||||
|
// startup log both read — reports what the environment asked for.
|
||||||
|
func assertMetricsAuthAccepted(t *testing.T, expectAuth bool) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
var cfg *config.Config
|
||||||
|
|
||||||
|
app := fxtest.New(
|
||||||
|
t,
|
||||||
|
fx.Provide(globals.New, logger.New, config.New),
|
||||||
|
fx.Populate(&cfg),
|
||||||
|
)
|
||||||
|
require.NoError(t, app.Err())
|
||||||
|
|
||||||
|
app.RequireStart()
|
||||||
|
|
||||||
|
defer app.RequireStop()
|
||||||
|
|
||||||
|
assert.Equal(t, expectAuth, cfg.MetricsAuthEnabled())
|
||||||
|
}
|
||||||
|
|||||||
@@ -140,11 +140,11 @@ type Engine struct {
|
|||||||
retryCh chan Task
|
retryCh chan Task
|
||||||
workers int
|
workers int
|
||||||
|
|
||||||
// mx is the delivery metric set. Production wires the
|
// mtr is the delivery metric set. Production wires the
|
||||||
// process-wide one; a test can substitute a set registered on
|
// process-wide one; a test can substitute a set registered on
|
||||||
// a private registry so its assertions are not disturbed by
|
// a private registry so its assertions are not disturbed by
|
||||||
// deliveries other tests are making at the same time.
|
// deliveries other tests are making at the same time.
|
||||||
mx *metrics.Set
|
mtr *metrics.Set
|
||||||
|
|
||||||
// targets maps each target type to its implementation.
|
// targets maps each target type to its implementation.
|
||||||
targets map[database.TargetType]Target
|
targets map[database.TargetType]Target
|
||||||
@@ -171,7 +171,7 @@ func New(
|
|||||||
deliveryCh: make(chan Task, deliveryChannelSize),
|
deliveryCh: make(chan Task, deliveryChannelSize),
|
||||||
retryCh: make(chan Task, retryChannelSize),
|
retryCh: make(chan Task, retryChannelSize),
|
||||||
workers: defaultWorkers,
|
workers: defaultWorkers,
|
||||||
mx: metrics.Default(),
|
mtr: metrics.Default(),
|
||||||
}
|
}
|
||||||
|
|
||||||
e.initTargets(&http.Client{
|
e.initTargets(&http.Client{
|
||||||
@@ -838,11 +838,6 @@ func (e *Engine) failUnretryableRetry(
|
|||||||
target.Type,
|
target.Type,
|
||||||
)
|
)
|
||||||
|
|
||||||
// The delivery was loaded without its target relation, so
|
|
||||||
// attach it: updateDeliveryStatus reads the type to label the
|
|
||||||
// terminal failure it is about to count.
|
|
||||||
d.Target = *target
|
|
||||||
|
|
||||||
e.recordResult(
|
e.recordResult(
|
||||||
webhookDB,
|
webhookDB,
|
||||||
d,
|
d,
|
||||||
@@ -854,28 +849,26 @@ func (e *Engine) failUnretryableRetry(
|
|||||||
0,
|
0,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// The type is passed rather than assigned onto d: the delivery
|
||||||
|
// is loaded here without its target relation, and populating
|
||||||
|
// d.Target would make GORM's SaveBeforeAssociations upsert the
|
||||||
|
// whole target row — plaintext config, which for a slack target
|
||||||
|
// is the credential — into the per-webhook event database. See
|
||||||
|
// https://git.eeqj.de/sneak/webhooker/issues/206.
|
||||||
e.updateDeliveryStatus(
|
e.updateDeliveryStatus(
|
||||||
webhookDB, d, database.DeliveryStatusFailed,
|
webhookDB, d, target.Type,
|
||||||
|
database.DeliveryStatusFailed,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// processDelivery dispatches a delivery to the target that
|
// processDelivery dispatches a delivery to the target that
|
||||||
// owns its type. Unknown target types fail the delivery.
|
// owns its type. Unknown target types fail the delivery.
|
||||||
//
|
|
||||||
// It is also where the attempt counter and the duration histogram
|
|
||||||
// are recorded, because it is the one point every target type
|
|
||||||
// passes through on every attempt: a target added later is
|
|
||||||
// instrumented without touching it, and the duration measured is
|
|
||||||
// the whole cost of the attempt rather than whatever each target
|
|
||||||
// happens to time for itself.
|
|
||||||
func (e *Engine) processDelivery(
|
func (e *Engine) processDelivery(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
webhookDB *gorm.DB,
|
webhookDB *gorm.DB,
|
||||||
d *database.Delivery,
|
d *database.Delivery,
|
||||||
task *Task,
|
task *Task,
|
||||||
) {
|
) {
|
||||||
e.mx.DeliveryAttempted(d.Target.Type)
|
|
||||||
|
|
||||||
target, ok := e.targets[d.Target.Type]
|
target, ok := e.targets[d.Target.Type]
|
||||||
if !ok {
|
if !ok {
|
||||||
e.log.Error(
|
e.log.Error(
|
||||||
@@ -885,19 +878,32 @@ func (e *Engine) processDelivery(
|
|||||||
)
|
)
|
||||||
|
|
||||||
e.updateDeliveryStatus(
|
e.updateDeliveryStatus(
|
||||||
webhookDB, d, database.DeliveryStatusFailed,
|
webhookDB, d, d.Target.Type,
|
||||||
|
database.DeliveryStatusFailed,
|
||||||
)
|
)
|
||||||
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
start := time.Now()
|
|
||||||
|
|
||||||
target.Deliver(ctx, webhookDB, d, task, e)
|
target.Deliver(ctx, webhookDB, d, task, e)
|
||||||
|
}
|
||||||
|
|
||||||
e.mx.ObserveDeliveryDuration(
|
// observeAttempt counts one delivery attempt that was actually
|
||||||
d.Target.Type, time.Since(start),
|
// dispatched to a target, and records how long it took.
|
||||||
)
|
//
|
||||||
|
// It is called from the dispatch paths rather than from around
|
||||||
|
// Target.Deliver, because Deliver is also entered for deliveries
|
||||||
|
// that never reach the wire: a delivery an open circuit breaker
|
||||||
|
// refuses sends nothing, records no DeliveryResult, and is
|
||||||
|
// rescheduled. Counting those would climb the attempts counter with
|
||||||
|
// no traffic behind it and fill the duration histogram with
|
||||||
|
// microsecond samples, which would make the delivery-duration
|
||||||
|
// quantiles improve during exactly the outage they exist to reveal.
|
||||||
|
func (e *Engine) observeAttempt(
|
||||||
|
t database.TargetType, dur time.Duration,
|
||||||
|
) {
|
||||||
|
e.mtr.DeliveryAttempted(t)
|
||||||
|
e.mtr.ObserveDeliveryDuration(t, dur)
|
||||||
}
|
}
|
||||||
|
|
||||||
// recordResult persists a DeliveryResult row describing a
|
// recordResult persists a DeliveryResult row describing a
|
||||||
@@ -936,13 +942,21 @@ func (e *Engine) recordResult(
|
|||||||
// It is a cross-target helper the targets call, and therefore the
|
// It is a cross-target helper the targets call, and therefore the
|
||||||
// single point where a delivery's outcome — delivered, terminally
|
// single point where a delivery's outcome — delivered, terminally
|
||||||
// failed, or put back into retry — is counted.
|
// failed, or put back into retry — is counted.
|
||||||
|
//
|
||||||
|
// The target type is a parameter rather than read off d.Target
|
||||||
|
// because one caller — failUnretryableRetry — deliberately holds a
|
||||||
|
// delivery loaded without its target relation, and must keep it that
|
||||||
|
// way: a populated d.Target makes GORM upsert the target row, config
|
||||||
|
// and all, into the per-webhook database.
|
||||||
|
//
|
||||||
|
// The counter moves only after the row is written, so a transition
|
||||||
|
// the database rejected is not claimed as an outcome that happened.
|
||||||
func (e *Engine) updateDeliveryStatus(
|
func (e *Engine) updateDeliveryStatus(
|
||||||
webhookDB *gorm.DB,
|
webhookDB *gorm.DB,
|
||||||
d *database.Delivery,
|
d *database.Delivery,
|
||||||
|
targetType database.TargetType,
|
||||||
status database.DeliveryStatus,
|
status database.DeliveryStatus,
|
||||||
) {
|
) {
|
||||||
e.mx.DeliveryStatusChanged(d.Target.Type, status)
|
|
||||||
|
|
||||||
err := webhookDB.Model(d).
|
err := webhookDB.Model(d).
|
||||||
Update("status", status).Error
|
Update("status", status).Error
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -952,7 +966,11 @@ func (e *Engine) updateDeliveryStatus(
|
|||||||
"status", status,
|
"status", status,
|
||||||
"error", err,
|
"error", err,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
e.mtr.DeliveryStatusChanged(targetType, status)
|
||||||
}
|
}
|
||||||
|
|
||||||
func truncate(s string, maxLen int) string {
|
func truncate(s string, maxLen int) string {
|
||||||
|
|||||||
@@ -886,6 +886,82 @@ func TestSweepSingleRetry_TypeNoLongerRetries(
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestFailUnretryableRetry_WritesNoTargetRow proves the
|
||||||
|
// orphaned-retry terminal path leaves no target row — and so no
|
||||||
|
// plaintext target config — in the per-webhook event database.
|
||||||
|
//
|
||||||
|
// That path loads the delivery without its Target relation on
|
||||||
|
// purpose. Populating d.Target makes GORM's SaveBeforeAssociations
|
||||||
|
// upsert the whole target row on the status UPDATE, which for a slack
|
||||||
|
// target writes the incoming-webhook credential into events-*.db.
|
||||||
|
// See https://git.eeqj.de/sneak/webhooker/issues/206.
|
||||||
|
func TestFailUnretryableRetry_WritesNoTargetRow(
|
||||||
|
t *testing.T,
|
||||||
|
) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
s := newISetup(t)
|
||||||
|
|
||||||
|
iCreateWebhook(
|
||||||
|
t, s.MainDB, s.WebhookID, "no-target-row",
|
||||||
|
)
|
||||||
|
|
||||||
|
targetID := uuid.New().String()
|
||||||
|
|
||||||
|
// A Slack incoming-webhook URL: the target config IS the
|
||||||
|
// credential, which is what makes a leaked target row a
|
||||||
|
// disclosure rather than a curiosity.
|
||||||
|
hookURL := "https://hooks.slack.com/services/T00/B00/x"
|
||||||
|
|
||||||
|
iCreateTarget(t, s.MainDB, targetID,
|
||||||
|
s.WebhookID, "credential-bearing",
|
||||||
|
database.TargetTypeLog, iHTTPConfig(hookURL), 5,
|
||||||
|
)
|
||||||
|
|
||||||
|
event := iSeedEvent(
|
||||||
|
t, s.WebhookDB, s.WebhookID, `{"orphaned":"retry"}`,
|
||||||
|
)
|
||||||
|
|
||||||
|
d := iSeedDelivery(
|
||||||
|
t, s.WebhookDB, event.ID, targetID,
|
||||||
|
database.DeliveryStatusRetrying,
|
||||||
|
)
|
||||||
|
|
||||||
|
iSeedFailedResult(t, s.WebhookDB, d.ID)
|
||||||
|
|
||||||
|
s.Engine.ExportSweepWebhookRetries(
|
||||||
|
context.Background(), s.WebhookID,
|
||||||
|
)
|
||||||
|
|
||||||
|
iAssertStatus(t, s.WebhookDB, d.ID,
|
||||||
|
database.DeliveryStatusFailed,
|
||||||
|
)
|
||||||
|
|
||||||
|
// The table exists in the per-webhook database because GORM
|
||||||
|
// migrates the Delivery relation's model alongside it. It must
|
||||||
|
// stay empty.
|
||||||
|
var targetRows int64
|
||||||
|
|
||||||
|
require.NoError(t, s.WebhookDB.
|
||||||
|
Table("targets").
|
||||||
|
Count(&targetRows).Error)
|
||||||
|
|
||||||
|
assert.Zero(t, targetRows,
|
||||||
|
"orphaned-retry terminal failure wrote a target row "+
|
||||||
|
"into the per-webhook event database",
|
||||||
|
)
|
||||||
|
|
||||||
|
var configs []string
|
||||||
|
|
||||||
|
require.NoError(t, s.WebhookDB.
|
||||||
|
Table("targets").
|
||||||
|
Pluck("config", &configs).Error)
|
||||||
|
|
||||||
|
assert.NotContains(
|
||||||
|
t, strings.Join(configs, " "), hookURL,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
func TestRecoverSingleRetry_UnknownTargetType(
|
func TestRecoverSingleRetry_UnknownTargetType(
|
||||||
t *testing.T,
|
t *testing.T,
|
||||||
) {
|
) {
|
||||||
|
|||||||
@@ -254,7 +254,7 @@ func NewTestEngine(
|
|||||||
deliveryCh: make(chan Task, deliveryChannelSize),
|
deliveryCh: make(chan Task, deliveryChannelSize),
|
||||||
retryCh: make(chan Task, retryChannelSize),
|
retryCh: make(chan Task, retryChannelSize),
|
||||||
workers: workers,
|
workers: workers,
|
||||||
mx: metrics.Default(),
|
mtr: metrics.Default(),
|
||||||
}
|
}
|
||||||
e.initTargets(client)
|
e.initTargets(client)
|
||||||
|
|
||||||
@@ -269,7 +269,7 @@ func NewTestEngineSmallRetry(
|
|||||||
e := &Engine{
|
e := &Engine{
|
||||||
log: log,
|
log: log,
|
||||||
retryCh: make(chan Task, 1),
|
retryCh: make(chan Task, 1),
|
||||||
mx: metrics.Default(),
|
mtr: metrics.Default(),
|
||||||
}
|
}
|
||||||
e.initTargets(nil)
|
e.initTargets(nil)
|
||||||
|
|
||||||
@@ -292,7 +292,7 @@ func NewTestEngineWithDB(
|
|||||||
deliveryCh: make(chan Task, deliveryChannelSize),
|
deliveryCh: make(chan Task, deliveryChannelSize),
|
||||||
retryCh: make(chan Task, retryChannelSize),
|
retryCh: make(chan Task, retryChannelSize),
|
||||||
workers: workers,
|
workers: workers,
|
||||||
mx: metrics.Default(),
|
mtr: metrics.Default(),
|
||||||
}
|
}
|
||||||
e.initTargets(client)
|
e.initTargets(client)
|
||||||
|
|
||||||
@@ -302,8 +302,8 @@ func NewTestEngineWithDB(
|
|||||||
// ExportSetMetrics substitutes the engine's metric set, so a test can
|
// ExportSetMetrics substitutes the engine's metric set, so a test can
|
||||||
// assert on collectors registered on a private registry instead of
|
// assert on collectors registered on a private registry instead of
|
||||||
// the process-wide ones every other test is also moving.
|
// the process-wide ones every other test is also moving.
|
||||||
func (e *Engine) ExportSetMetrics(mx *metrics.Set) {
|
func (e *Engine) ExportSetMetrics(mtr *metrics.Set) {
|
||||||
e.mx = mx
|
e.mtr = mtr
|
||||||
}
|
}
|
||||||
|
|
||||||
// ExportSampleQueueDepths runs one queue depth sample synchronously.
|
// ExportSampleQueueDepths runs one queue depth sample synchronously.
|
||||||
|
|||||||
@@ -29,8 +29,9 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
mTypeHTTP = "http"
|
mTypeHTTP = "http"
|
||||||
mTypeLog = "log"
|
mTypeLog = "log"
|
||||||
|
mTypeUnknown = "unknown"
|
||||||
)
|
)
|
||||||
|
|
||||||
// mIsolate gives the setup's engine a metric set registered on a
|
// mIsolate gives the setup's engine a metric set registered on a
|
||||||
@@ -105,14 +106,15 @@ func mGauge(
|
|||||||
GetGauge().GetValue()
|
GetGauge().GetValue()
|
||||||
}
|
}
|
||||||
|
|
||||||
func mObservations(
|
// mHTTPDurations returns how many samples the delivery duration
|
||||||
t *testing.T,
|
// histogram holds for the http target type, which is the type every
|
||||||
reg *prometheus.Registry,
|
// test here times.
|
||||||
name, targetType string,
|
func mHTTPDurations(
|
||||||
|
t *testing.T, reg *prometheus.Registry,
|
||||||
) uint64 {
|
) uint64 {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
return mFind(t, reg, name, targetType).
|
return mFind(t, reg, mDuration, mTypeHTTP).
|
||||||
GetHistogram().GetSampleCount()
|
GetHistogram().GetSampleCount()
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -138,7 +140,7 @@ func TestDeliveryMetrics_SuccessAndRetryExhaustion(
|
|||||||
assert.InDelta(t, 0.0,
|
assert.InDelta(t, 0.0,
|
||||||
mCounter(t, reg, mRetries, mTypeHTTP), 0)
|
mCounter(t, reg, mRetries, mTypeHTTP), 0)
|
||||||
assert.Equal(t, uint64(1),
|
assert.Equal(t, uint64(1),
|
||||||
mObservations(t, reg, mDuration, mTypeHTTP))
|
mHTTPDurations(t, reg))
|
||||||
|
|
||||||
mExhaustRetries(t, s)
|
mExhaustRetries(t, s)
|
||||||
|
|
||||||
@@ -153,7 +155,7 @@ func TestDeliveryMetrics_SuccessAndRetryExhaustion(
|
|||||||
assert.InDelta(t, 1.0,
|
assert.InDelta(t, 1.0,
|
||||||
mCounter(t, reg, mFailed, mTypeHTTP), 0)
|
mCounter(t, reg, mFailed, mTypeHTTP), 0)
|
||||||
assert.Equal(t, uint64(3),
|
assert.Equal(t, uint64(3),
|
||||||
mObservations(t, reg, mDuration, mTypeHTTP))
|
mHTTPDurations(t, reg))
|
||||||
|
|
||||||
// Two consecutive failures are below the trip threshold.
|
// Two consecutive failures are below the trip threshold.
|
||||||
assert.InDelta(t, 0.0,
|
assert.InDelta(t, 0.0,
|
||||||
@@ -313,6 +315,128 @@ func TestDeliveryMetrics_CircuitBreakerGauge(t *testing.T) {
|
|||||||
mGauge(t, reg, mBreakers, mTypeHTTP), 0)
|
mGauge(t, reg, mBreakers, mTypeHTTP), 0)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestDeliveryMetrics_BreakerBlockedIsNotAnAttempt proves a delivery
|
||||||
|
// an open circuit breaker refuses is neither counted as an attempt
|
||||||
|
// nor observed in the duration histogram.
|
||||||
|
//
|
||||||
|
// It sends nothing and records no result row, so counting it would
|
||||||
|
// climb the attempts counter with no traffic behind it and pull the
|
||||||
|
// duration quantiles down with near-zero samples for as long as the
|
||||||
|
// breaker stayed open — the metric moving the wrong way during the
|
||||||
|
// outage it exists to reveal.
|
||||||
|
func TestDeliveryMetrics_BreakerBlockedIsNotAnAttempt(
|
||||||
|
t *testing.T,
|
||||||
|
) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
s := newISetup(t)
|
||||||
|
reg := mIsolate(t, s)
|
||||||
|
|
||||||
|
ts := httptest.NewServer(http.HandlerFunc(
|
||||||
|
func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
w.WriteHeader(http.StatusInternalServerError)
|
||||||
|
},
|
||||||
|
))
|
||||||
|
defer ts.Close()
|
||||||
|
|
||||||
|
event := iSeedEvent(
|
||||||
|
t, s.WebhookDB, s.WebhookID, `{"blocked":true}`,
|
||||||
|
)
|
||||||
|
targetID := uuid.New().String()
|
||||||
|
|
||||||
|
d := iSeedDelivery(
|
||||||
|
t, s.WebhookDB, event.ID, targetID,
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
|
||||||
|
body := event.Body
|
||||||
|
cfg := iHTTPConfig(ts.URL)
|
||||||
|
maxRetries := delivery.ExportDefaultFailureThreshold + 5
|
||||||
|
|
||||||
|
first := iTask(
|
||||||
|
d, event, s.WebhookID, targetID,
|
||||||
|
"metrics-blocked", cfg, maxRetries, 1, &body,
|
||||||
|
)
|
||||||
|
|
||||||
|
s.Engine.ExportProcessNewTask(context.TODO(), &first)
|
||||||
|
|
||||||
|
for attempt := 2; attempt <= delivery.
|
||||||
|
ExportDefaultFailureThreshold; attempt++ {
|
||||||
|
task := iTask(
|
||||||
|
d, event, s.WebhookID, targetID,
|
||||||
|
"metrics-blocked", cfg, maxRetries, attempt, &body,
|
||||||
|
)
|
||||||
|
|
||||||
|
s.Engine.ExportProcessRetryTask(context.TODO(), &task)
|
||||||
|
}
|
||||||
|
|
||||||
|
require.InDelta(t, 1.0,
|
||||||
|
mGauge(t, reg, mBreakers, mTypeHTTP), 0,
|
||||||
|
"breaker should be open before the blocked attempt")
|
||||||
|
|
||||||
|
threshold := float64(
|
||||||
|
delivery.ExportDefaultFailureThreshold,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.InDelta(t, threshold,
|
||||||
|
mCounter(t, reg, mAttempts, mTypeHTTP), 0)
|
||||||
|
assert.Equal(t, uint64(threshold),
|
||||||
|
mHTTPDurations(t, reg))
|
||||||
|
|
||||||
|
retriesBefore := mCounter(t, reg, mRetries, mTypeHTTP)
|
||||||
|
|
||||||
|
blocked := iTask(
|
||||||
|
d, event, s.WebhookID, targetID,
|
||||||
|
"metrics-blocked", cfg, maxRetries,
|
||||||
|
delivery.ExportDefaultFailureThreshold+1, &body,
|
||||||
|
)
|
||||||
|
|
||||||
|
s.Engine.ExportProcessRetryTask(context.TODO(), &blocked)
|
||||||
|
|
||||||
|
// The breaker refused it: rescheduled, so the retry counter
|
||||||
|
// moved, but nothing was attempted or timed.
|
||||||
|
assert.InDelta(t, retriesBefore+1,
|
||||||
|
mCounter(t, reg, mRetries, mTypeHTTP), 0)
|
||||||
|
assert.InDelta(t, threshold,
|
||||||
|
mCounter(t, reg, mAttempts, mTypeHTTP), 0)
|
||||||
|
assert.Equal(t, uint64(threshold),
|
||||||
|
mHTTPDurations(t, reg))
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestDeliveryMetrics_OrphanedRetryFailureLabelled proves the
|
||||||
|
// terminal failure of a delivery whose target no longer retries is
|
||||||
|
// counted against the target's real type, not against unknown. The
|
||||||
|
// type is threaded in as an argument because populating d.Target on
|
||||||
|
// that path would write the target row into the per-webhook database
|
||||||
|
// (https://git.eeqj.de/sneak/webhooker/issues/206).
|
||||||
|
func TestDeliveryMetrics_OrphanedRetryFailureLabelled(
|
||||||
|
t *testing.T,
|
||||||
|
) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
s := newISetup(t)
|
||||||
|
reg := mIsolate(t, s)
|
||||||
|
|
||||||
|
iCreateWebhook(
|
||||||
|
t, s.MainDB, s.WebhookID, "orphaned-label",
|
||||||
|
)
|
||||||
|
|
||||||
|
deliveryID := iSeedRetryingWithType(
|
||||||
|
t, s, database.TargetTypeLog,
|
||||||
|
)
|
||||||
|
|
||||||
|
s.Engine.ExportSweepWebhookRetries(
|
||||||
|
context.Background(), s.WebhookID,
|
||||||
|
)
|
||||||
|
|
||||||
|
iAssertStatus(t, s.WebhookDB, deliveryID,
|
||||||
|
database.DeliveryStatusFailed,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.InDelta(t, 1.0,
|
||||||
|
mCounter(t, reg, mFailed, mTypeLog), 0)
|
||||||
|
}
|
||||||
|
|
||||||
// TestDeliveryMetrics_QueueDepthGauges proves the sampler publishes
|
// TestDeliveryMetrics_QueueDepthGauges proves the sampler publishes
|
||||||
// the queued deliveries it finds in the per-webhook databases, and
|
// the queued deliveries it finds in the per-webhook databases, and
|
||||||
// that a drained queue reads zero rather than keeping its last
|
// that a drained queue reads zero rather than keeping its last
|
||||||
@@ -376,3 +500,46 @@ func TestDeliveryMetrics_QueueDepthGauges(t *testing.T) {
|
|||||||
assert.InDelta(t, 0.0,
|
assert.InDelta(t, 0.0,
|
||||||
mGauge(t, reg, mRetrying, mTypeHTTP), 0)
|
mGauge(t, reg, mRetrying, mTypeHTTP), 0)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestDeliveryMetrics_QueueDepthDeletedTarget proves a backlog queued
|
||||||
|
// against a target that has since been deleted stays visible, in the
|
||||||
|
// unknown series, instead of being dropped. That backlog is the one
|
||||||
|
// nobody is watching, so losing it would defeat the queue-depth
|
||||||
|
// alerting this metric exists for.
|
||||||
|
func TestDeliveryMetrics_QueueDepthDeletedTarget(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
s := newISetup(t)
|
||||||
|
reg := mIsolate(t, s)
|
||||||
|
|
||||||
|
iCreateWebhook(
|
||||||
|
t, s.MainDB, s.WebhookID, "deleted-target",
|
||||||
|
)
|
||||||
|
|
||||||
|
// No target row is created: this is a delivery whose target was
|
||||||
|
// deleted out from under it.
|
||||||
|
targetID := uuid.New().String()
|
||||||
|
|
||||||
|
event := iSeedEvent(
|
||||||
|
t, s.WebhookDB, s.WebhookID, `{"orphan":true}`,
|
||||||
|
)
|
||||||
|
|
||||||
|
iSeedDelivery(
|
||||||
|
t, s.WebhookDB, event.ID, targetID,
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
|
||||||
|
iSeedDelivery(
|
||||||
|
t, s.WebhookDB, event.ID, targetID,
|
||||||
|
database.DeliveryStatusRetrying,
|
||||||
|
)
|
||||||
|
|
||||||
|
s.Engine.ExportSampleQueueDepths(context.Background())
|
||||||
|
|
||||||
|
assert.InDelta(t, 1.0,
|
||||||
|
mGauge(t, reg, mPending, mTypeUnknown), 0)
|
||||||
|
assert.InDelta(t, 1.0,
|
||||||
|
mGauge(t, reg, mRetrying, mTypeUnknown), 0)
|
||||||
|
assert.InDelta(t, 0.0,
|
||||||
|
mGauge(t, reg, mPending, mTypeHTTP), 0)
|
||||||
|
}
|
||||||
|
|||||||
@@ -89,7 +89,7 @@ func (e *Engine) sampleQueueDepths(ctx context.Context) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
e.mx.SetQueueDepths(pending, retrying)
|
e.mtr.SetQueueDepths(pending, retrying)
|
||||||
}
|
}
|
||||||
|
|
||||||
// targetTypesByID maps every configured target id to its type. The
|
// targetTypesByID maps every configured target id to its type. The
|
||||||
@@ -121,9 +121,13 @@ func (e *Engine) targetTypesByID() (
|
|||||||
}
|
}
|
||||||
|
|
||||||
// sampleWebhookQueueDepths adds one webhook's queued deliveries into
|
// sampleWebhookQueueDepths adds one webhook's queued deliveries into
|
||||||
// the running totals. A delivery whose target has since been deleted
|
// the running totals.
|
||||||
// resolves to the empty type and lands in the unknown bucket rather
|
//
|
||||||
// than being dropped.
|
// A delivery whose target has since been deleted is not in the type
|
||||||
|
// map and so counts under the empty target type. Set.SetQueueDepths
|
||||||
|
// folds that into the unknown series rather than dropping it: a
|
||||||
|
// backlog stuck behind a deleted target is a backlog that still needs
|
||||||
|
// to be alertable.
|
||||||
func (e *Engine) sampleWebhookQueueDepths(
|
func (e *Engine) sampleWebhookQueueDepths(
|
||||||
webhookID string,
|
webhookID string,
|
||||||
types map[string]database.TargetType,
|
types map[string]database.TargetType,
|
||||||
|
|||||||
@@ -27,6 +27,12 @@ type Scheduler interface {
|
|||||||
// own circuit breaker, and reschedules via the injected
|
// own circuit breaker, and reschedules via the injected
|
||||||
// Scheduler. Fire-and-forget targets simply record a single
|
// Scheduler. Fire-and-forget targets simply record a single
|
||||||
// attempt.
|
// attempt.
|
||||||
|
//
|
||||||
|
// An implementation reports each attempt it actually dispatches to
|
||||||
|
// Engine.observeAttempt, alongside the DeliveryResult it records for
|
||||||
|
// it. Deliver is also entered for attempts that never happen — an
|
||||||
|
// open circuit breaker refuses one — so the count cannot be taken
|
||||||
|
// from around this call.
|
||||||
type Target interface {
|
type Target interface {
|
||||||
Deliver(
|
Deliver(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
@@ -74,6 +80,12 @@ type attemptResult struct {
|
|||||||
errMsg string
|
errMsg string
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// elapsed returns how long the attempt took. The field is stored in
|
||||||
|
// milliseconds because that is what DeliveryResult persists.
|
||||||
|
func (r attemptResult) elapsed() time.Duration {
|
||||||
|
return time.Duration(r.duration) * time.Millisecond
|
||||||
|
}
|
||||||
|
|
||||||
// initTargets builds the target registry, wiring each target
|
// initTargets builds the target registry, wiring each target
|
||||||
// to the engine's persistence helpers and giving the HTTP and
|
// to the engine's persistence helpers and giving the HTTP and
|
||||||
// Slack targets the shared SSRF-safe client. It is called by
|
// Slack targets the shared SSRF-safe client. It is called by
|
||||||
|
|||||||
@@ -42,7 +42,14 @@ func (t *databaseTarget) Deliver(
|
|||||||
_ *Task,
|
_ *Task,
|
||||||
_ Scheduler,
|
_ Scheduler,
|
||||||
) {
|
) {
|
||||||
|
start := time.Now()
|
||||||
|
|
||||||
err := t.archive(d)
|
err := t.archive(d)
|
||||||
|
|
||||||
|
elapsed := time.Since(start)
|
||||||
|
|
||||||
|
t.eng.observeAttempt(d.Target.Type, elapsed)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.eng.log.Error(
|
t.eng.log.Error(
|
||||||
"failed to archive event to database target",
|
"failed to archive event to database target",
|
||||||
@@ -53,22 +60,25 @@ func (t *databaseTarget) Deliver(
|
|||||||
|
|
||||||
t.eng.recordResult(
|
t.eng.recordResult(
|
||||||
webhookDB, d, 1, false, 0, "",
|
webhookDB, d, 1, false, 0, "",
|
||||||
err.Error(), 0,
|
err.Error(), elapsed.Milliseconds(),
|
||||||
)
|
)
|
||||||
|
|
||||||
t.eng.updateDeliveryStatus(
|
t.eng.updateDeliveryStatus(
|
||||||
webhookDB, d, database.DeliveryStatusFailed,
|
webhookDB, d, d.Target.Type,
|
||||||
|
database.DeliveryStatusFailed,
|
||||||
)
|
)
|
||||||
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
t.eng.recordResult(
|
t.eng.recordResult(
|
||||||
webhookDB, d, 1, true, 0, "", "", 0,
|
webhookDB, d, 1, true, 0, "", "",
|
||||||
|
elapsed.Milliseconds(),
|
||||||
)
|
)
|
||||||
|
|
||||||
t.eng.updateDeliveryStatus(
|
t.eng.updateDeliveryStatus(
|
||||||
webhookDB, d, database.DeliveryStatusDelivered,
|
webhookDB, d, d.Target.Type,
|
||||||
|
database.DeliveryStatusDelivered,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -74,6 +74,8 @@ func (c *httpCore) fireAndForget(
|
|||||||
d *database.Delivery,
|
d *database.Delivery,
|
||||||
res attemptResult,
|
res attemptResult,
|
||||||
) {
|
) {
|
||||||
|
c.eng.observeAttempt(d.Target.Type, res.elapsed())
|
||||||
|
|
||||||
c.eng.recordResult(
|
c.eng.recordResult(
|
||||||
webhookDB, d, 1, res.success,
|
webhookDB, d, 1, res.success,
|
||||||
res.statusCode, res.respBody, res.errMsg,
|
res.statusCode, res.respBody, res.errMsg,
|
||||||
@@ -82,7 +84,7 @@ func (c *httpCore) fireAndForget(
|
|||||||
|
|
||||||
if res.success {
|
if res.success {
|
||||||
c.eng.updateDeliveryStatus(
|
c.eng.updateDeliveryStatus(
|
||||||
webhookDB, d,
|
webhookDB, d, d.Target.Type,
|
||||||
database.DeliveryStatusDelivered,
|
database.DeliveryStatusDelivered,
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -90,7 +92,8 @@ func (c *httpCore) fireAndForget(
|
|||||||
}
|
}
|
||||||
|
|
||||||
c.eng.updateDeliveryStatus(
|
c.eng.updateDeliveryStatus(
|
||||||
webhookDB, d, database.DeliveryStatusFailed,
|
webhookDB, d, d.Target.Type,
|
||||||
|
database.DeliveryStatusFailed,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -116,6 +119,8 @@ func (c *httpCore) withRetry(
|
|||||||
|
|
||||||
res := attempt()
|
res := attempt()
|
||||||
|
|
||||||
|
c.eng.observeAttempt(d.Target.Type, res.elapsed())
|
||||||
|
|
||||||
c.eng.recordResult(
|
c.eng.recordResult(
|
||||||
webhookDB, d, attemptNum, res.success,
|
webhookDB, d, attemptNum, res.success,
|
||||||
res.statusCode, res.respBody, res.errMsg,
|
res.statusCode, res.respBody, res.errMsg,
|
||||||
@@ -126,7 +131,7 @@ func (c *httpCore) withRetry(
|
|||||||
cb.RecordSuccess()
|
cb.RecordSuccess()
|
||||||
|
|
||||||
c.eng.updateDeliveryStatus(
|
c.eng.updateDeliveryStatus(
|
||||||
webhookDB, d,
|
webhookDB, d, d.Target.Type,
|
||||||
database.DeliveryStatusDelivered,
|
database.DeliveryStatusDelivered,
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -164,7 +169,7 @@ func (c *httpCore) circuitBreakerBlock(
|
|||||||
)
|
)
|
||||||
|
|
||||||
c.eng.updateDeliveryStatus(
|
c.eng.updateDeliveryStatus(
|
||||||
webhookDB, d,
|
webhookDB, d, d.Target.Type,
|
||||||
database.DeliveryStatusRetrying,
|
database.DeliveryStatusRetrying,
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -184,7 +189,7 @@ func (c *httpCore) handleRetry(
|
|||||||
) {
|
) {
|
||||||
if attemptNum >= maxRetries {
|
if attemptNum >= maxRetries {
|
||||||
c.eng.updateDeliveryStatus(
|
c.eng.updateDeliveryStatus(
|
||||||
webhookDB, d,
|
webhookDB, d, d.Target.Type,
|
||||||
database.DeliveryStatusFailed,
|
database.DeliveryStatusFailed,
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -192,7 +197,8 @@ func (c *httpCore) handleRetry(
|
|||||||
}
|
}
|
||||||
|
|
||||||
c.eng.updateDeliveryStatus(
|
c.eng.updateDeliveryStatus(
|
||||||
webhookDB, d, database.DeliveryStatusRetrying,
|
webhookDB, d, d.Target.Type,
|
||||||
|
database.DeliveryStatusRetrying,
|
||||||
)
|
)
|
||||||
|
|
||||||
backoff := calcBackoff(attemptNum)
|
backoff := calcBackoff(attemptNum)
|
||||||
@@ -241,7 +247,7 @@ func (c *httpCore) publishCircuitState(
|
|||||||
return true
|
return true
|
||||||
})
|
})
|
||||||
|
|
||||||
c.eng.mx.SetCircuitBreakersOpen(targetType, open)
|
c.eng.mtr.SetCircuitBreakersOpen(targetType, open)
|
||||||
}
|
}
|
||||||
|
|
||||||
// remainingBackoff returns how long remains of the backoff
|
// remainingBackoff returns how long remains of the backoff
|
||||||
@@ -331,7 +337,8 @@ func (t *httpTarget) Deliver(
|
|||||||
)
|
)
|
||||||
|
|
||||||
t.eng.updateDeliveryStatus(
|
t.eng.updateDeliveryStatus(
|
||||||
webhookDB, d, database.DeliveryStatusFailed,
|
webhookDB, d, d.Target.Type,
|
||||||
|
database.DeliveryStatusFailed,
|
||||||
)
|
)
|
||||||
|
|
||||||
return
|
return
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package delivery
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"time"
|
||||||
|
|
||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
@@ -34,6 +35,8 @@ func (t *logTarget) Deliver(
|
|||||||
_ *Task,
|
_ *Task,
|
||||||
_ Scheduler,
|
_ Scheduler,
|
||||||
) {
|
) {
|
||||||
|
start := time.Now()
|
||||||
|
|
||||||
t.eng.log.Info(
|
t.eng.log.Info(
|
||||||
"webhook event delivered to log target",
|
"webhook event delivered to log target",
|
||||||
"delivery_id", d.ID,
|
"delivery_id", d.ID,
|
||||||
@@ -48,11 +51,17 @@ func (t *logTarget) Deliver(
|
|||||||
"body", d.Event.Body,
|
"body", d.Event.Body,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
elapsed := time.Since(start)
|
||||||
|
|
||||||
|
t.eng.observeAttempt(d.Target.Type, elapsed)
|
||||||
|
|
||||||
t.eng.recordResult(
|
t.eng.recordResult(
|
||||||
webhookDB, d, 1, true, 0, "", "", 0,
|
webhookDB, d, 1, true, 0, "", "",
|
||||||
|
elapsed.Milliseconds(),
|
||||||
)
|
)
|
||||||
|
|
||||||
t.eng.updateDeliveryStatus(
|
t.eng.updateDeliveryStatus(
|
||||||
webhookDB, d, database.DeliveryStatusDelivered,
|
webhookDB, d, d.Target.Type,
|
||||||
|
database.DeliveryStatusDelivered,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -101,7 +101,8 @@ func (t *slackTarget) failConfig(
|
|||||||
)
|
)
|
||||||
|
|
||||||
t.eng.updateDeliveryStatus(
|
t.eng.updateDeliveryStatus(
|
||||||
webhookDB, d, database.DeliveryStatusFailed,
|
webhookDB, d, d.Target.Type,
|
||||||
|
database.DeliveryStatusFailed,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -74,7 +74,7 @@ type Handlers struct {
|
|||||||
mw *middleware.Middleware
|
mw *middleware.Middleware
|
||||||
notifier delivery.Notifier
|
notifier delivery.Notifier
|
||||||
evictor delivery.WebhookEvictor
|
evictor delivery.WebhookEvictor
|
||||||
mx *metrics.Set
|
mtr *metrics.Set
|
||||||
templates map[string]*template.Template
|
templates map[string]*template.Template
|
||||||
|
|
||||||
// dummyVerifications counts the equivalent-cost verifications
|
// dummyVerifications counts the equivalent-cost verifications
|
||||||
@@ -116,7 +116,7 @@ func New(
|
|||||||
s.mw = params.Middleware
|
s.mw = params.Middleware
|
||||||
s.notifier = params.Notifier
|
s.notifier = params.Notifier
|
||||||
s.evictor = params.Evictor
|
s.evictor = params.Evictor
|
||||||
s.mx = metrics.Default()
|
s.mtr = metrics.Default()
|
||||||
|
|
||||||
// Parse all page templates once at startup
|
// Parse all page templates once at startup
|
||||||
s.templates = map[string]*template.Template{
|
s.templates = map[string]*template.Template{
|
||||||
|
|||||||
@@ -220,7 +220,7 @@ func (h *Handlers) createAndDeliverEvent(
|
|||||||
// Counted here, after the commit: an event is received once it
|
// Counted here, after the commit: an event is received once it
|
||||||
// is durably stored, which is what the delivery counters are
|
// is durably stored, which is what the delivery counters are
|
||||||
// compared against on a dashboard.
|
// compared against on a dashboard.
|
||||||
h.mx.EventReceived()
|
h.mtr.EventReceived()
|
||||||
|
|
||||||
h.finishWebhookResponse(w, event, entrypoint, tasks)
|
h.finishWebhookResponse(w, event, entrypoint, tasks)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -171,22 +171,56 @@ func (s *Set) DeliveryStatusChanged(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// SetQueueDepths publishes the pending and retrying queue depths from
|
// SetQueueDepths publishes the pending and retrying queue depths from
|
||||||
// one sample. Every known target type is written on every call, so a
|
// one sample. Every label in the queue domain is written on every
|
||||||
// type whose queue has drained reads zero instead of holding its last
|
// call, so a type whose queue has drained reads zero instead of
|
||||||
// value forever.
|
// holding its last value forever.
|
||||||
func (s *Set) SetQueueDepths(
|
func (s *Set) SetQueueDepths(
|
||||||
pending, retrying map[database.TargetType]int,
|
pending, retrying map[database.TargetType]int,
|
||||||
) {
|
) {
|
||||||
for _, t := range knownTargetTypes {
|
pendingByLabel := foldToLabels(pending)
|
||||||
label := string(t)
|
retryingByLabel := foldToLabels(retrying)
|
||||||
|
|
||||||
|
for _, label := range queueDepthLabels() {
|
||||||
s.deliveriesPending.WithLabelValues(label).
|
s.deliveriesPending.WithLabelValues(label).
|
||||||
Set(float64(pending[t]))
|
Set(float64(pendingByLabel[label]))
|
||||||
s.deliveriesRetrying.WithLabelValues(label).
|
s.deliveriesRetrying.WithLabelValues(label).
|
||||||
Set(float64(retrying[t]))
|
Set(float64(retryingByLabel[label]))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// queueDepthLabels is the label domain of the two queue-depth gauges:
|
||||||
|
// the known target types plus unknown.
|
||||||
|
//
|
||||||
|
// Unknown is a real bucket here, not a safety net. A delivery queued
|
||||||
|
// against a target that has since been deleted carries a target id no
|
||||||
|
// longer in the targets table, so the sample resolves it to the empty
|
||||||
|
// type; folding it into unknown is what keeps that backlog visible.
|
||||||
|
// Dropping it would hide the one queue nobody is watching.
|
||||||
|
func queueDepthLabels() []string {
|
||||||
|
labels := make([]string, 0, len(knownTargetTypes)+1)
|
||||||
|
|
||||||
|
for _, t := range knownTargetTypes {
|
||||||
|
labels = append(labels, string(t))
|
||||||
|
}
|
||||||
|
|
||||||
|
return append(labels, unknownTargetType)
|
||||||
|
}
|
||||||
|
|
||||||
|
// foldToLabels collapses a per-target-type count onto the bounded
|
||||||
|
// label domain, summing everything outside the known set into
|
||||||
|
// unknown.
|
||||||
|
func foldToLabels(
|
||||||
|
counts map[database.TargetType]int,
|
||||||
|
) map[string]int {
|
||||||
|
byLabel := make(map[string]int, len(counts))
|
||||||
|
|
||||||
|
for t, n := range counts {
|
||||||
|
byLabel[normalizeTargetType(t)] += n
|
||||||
|
}
|
||||||
|
|
||||||
|
return byLabel
|
||||||
|
}
|
||||||
|
|
||||||
// SetCircuitBreakersOpen publishes how many of a target type's
|
// SetCircuitBreakersOpen publishes how many of a target type's
|
||||||
// circuit breakers are currently open.
|
// circuit breakers are currently open.
|
||||||
func (s *Set) SetCircuitBreakersOpen(
|
func (s *Set) SetCircuitBreakersOpen(
|
||||||
@@ -274,6 +308,12 @@ func (s *Set) registerGauges(factory promauto.Factory) {
|
|||||||
// initSeries materialises every known-target-type series at zero, so
|
// initSeries materialises every known-target-type series at zero, so
|
||||||
// a dashboard and an alert rule see a target type that has not
|
// a dashboard and an alert rule see a target type that has not
|
||||||
// delivered yet rather than a missing series.
|
// delivered yet rather than a missing series.
|
||||||
|
//
|
||||||
|
// The queue-depth gauges additionally get their unknown series, which
|
||||||
|
// holds deliveries queued against a deleted target. That backlog can
|
||||||
|
// predate the process — it is read out of the databases, not counted
|
||||||
|
// from transitions — so its series has to exist from the first scrape
|
||||||
|
// rather than appearing only once a backlog has already built up.
|
||||||
func (s *Set) initSeries() {
|
func (s *Set) initSeries() {
|
||||||
for _, t := range knownTargetTypes {
|
for _, t := range knownTargetTypes {
|
||||||
label := string(t)
|
label := string(t)
|
||||||
@@ -286,6 +326,9 @@ func (s *Set) initSeries() {
|
|||||||
s.deliveriesRetrying.WithLabelValues(label)
|
s.deliveriesRetrying.WithLabelValues(label)
|
||||||
s.circuitBreakersOpen.WithLabelValues(label)
|
s.circuitBreakersOpen.WithLabelValues(label)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
s.deliveriesPending.WithLabelValues(unknownTargetType)
|
||||||
|
s.deliveriesRetrying.WithLabelValues(unknownTargetType)
|
||||||
}
|
}
|
||||||
|
|
||||||
// normalizeTargetType maps a target type onto the bounded label
|
// normalizeTargetType maps a target type onto the bounded label
|
||||||
|
|||||||
@@ -11,6 +11,12 @@ import (
|
|||||||
"sneak.berlin/go/webhooker/internal/metrics"
|
"sneak.berlin/go/webhooker/internal/metrics"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// knownLabels is the target_type label domain built from the target
|
||||||
|
// types the delivery engine implements.
|
||||||
|
func knownLabels() []string {
|
||||||
|
return []string{"http", "database", "log", "slack"}
|
||||||
|
}
|
||||||
|
|
||||||
// labelValues returns the target_type label values a metric family
|
// labelValues returns the target_type label values a metric family
|
||||||
// currently carries.
|
// currently carries.
|
||||||
func labelValues(
|
func labelValues(
|
||||||
@@ -105,9 +111,7 @@ func TestUnknownTargetTypeCollapses(t *testing.T) {
|
|||||||
)
|
)
|
||||||
|
|
||||||
assert.ElementsMatch(t,
|
assert.ElementsMatch(t,
|
||||||
[]string{
|
append(knownLabels(), "unknown"),
|
||||||
"http", "database", "log", "slack", "unknown",
|
|
||||||
},
|
|
||||||
values,
|
values,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
@@ -163,16 +167,72 @@ func TestKnownSeriesExistBeforeAnyDelivery(t *testing.T) {
|
|||||||
"webhooker_deliveries_succeeded_total",
|
"webhooker_deliveries_succeeded_total",
|
||||||
"webhooker_deliveries_failed_total",
|
"webhooker_deliveries_failed_total",
|
||||||
"webhooker_delivery_retries_total",
|
"webhooker_delivery_retries_total",
|
||||||
"webhooker_deliveries_pending",
|
|
||||||
"webhooker_deliveries_retrying",
|
|
||||||
"webhooker_circuit_breakers_open",
|
"webhooker_circuit_breakers_open",
|
||||||
} {
|
} {
|
||||||
assert.ElementsMatch(t,
|
assert.ElementsMatch(t,
|
||||||
[]string{"http", "database", "log", "slack"},
|
knownLabels(),
|
||||||
labelValues(t, reg, name),
|
labelValues(t, reg, name),
|
||||||
"metric %s", name,
|
"metric %s", name,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// The queue gauges additionally publish unknown from
|
||||||
|
// registration: a backlog queued against a deleted target lands
|
||||||
|
// there, and it can predate the process, so the series has to
|
||||||
|
// exist before the first sample rather than appearing only once
|
||||||
|
// something is already stuck.
|
||||||
|
for _, name := range []string{
|
||||||
|
"webhooker_deliveries_pending",
|
||||||
|
"webhooker_deliveries_retrying",
|
||||||
|
} {
|
||||||
|
assert.ElementsMatch(t,
|
||||||
|
append(knownLabels(), "unknown"),
|
||||||
|
labelValues(t, reg, name),
|
||||||
|
"metric %s", name,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestSetQueueDepthsFoldsUnknownTypes proves a queued delivery whose
|
||||||
|
// target type is not a known one — a target deleted out from under it
|
||||||
|
// resolves to the empty type — is summed into the unknown series
|
||||||
|
// instead of being dropped, and that the fold is a sum rather than a
|
||||||
|
// last-writer-wins.
|
||||||
|
func TestSetQueueDepthsFoldsUnknownTypes(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
reg := prometheus.NewRegistry()
|
||||||
|
set := metrics.New(reg)
|
||||||
|
|
||||||
|
set.SetQueueDepths(
|
||||||
|
map[database.TargetType]int{
|
||||||
|
database.TargetTypeHTTP: 1,
|
||||||
|
database.TargetType(""): 4,
|
||||||
|
database.TargetType("retired-type"): 3,
|
||||||
|
},
|
||||||
|
map[database.TargetType]int{
|
||||||
|
database.TargetType(""): 2,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.InDelta(t, 7.0, gaugeValue(
|
||||||
|
t, reg, "webhooker_deliveries_pending", "unknown",
|
||||||
|
), 0)
|
||||||
|
assert.InDelta(t, 2.0, gaugeValue(
|
||||||
|
t, reg, "webhooker_deliveries_retrying", "unknown",
|
||||||
|
), 0)
|
||||||
|
assert.InDelta(t, 1.0, gaugeValue(
|
||||||
|
t, reg, "webhooker_deliveries_pending", "http",
|
||||||
|
), 0)
|
||||||
|
|
||||||
|
set.SetQueueDepths(
|
||||||
|
map[database.TargetType]int{},
|
||||||
|
map[database.TargetType]int{},
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.InDelta(t, 0.0, gaugeValue(
|
||||||
|
t, reg, "webhooker_deliveries_pending", "unknown",
|
||||||
|
), 0)
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestDeliveryStatusChangedCounts maps each persisted status onto the
|
// TestDeliveryStatusChangedCounts maps each persisted status onto the
|
||||||
|
|||||||
@@ -43,10 +43,7 @@ func (s *Server) serveUntilShutdown() {
|
|||||||
err := s.httpServer.ListenAndServe()
|
err := s.httpServer.ListenAndServe()
|
||||||
if err != nil && !errors.Is(err, http.ErrServerClosed) {
|
if err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||||
s.log.Error("listen error", "error", err)
|
s.log.Error("listen error", "error", err)
|
||||||
|
s.shutdownOnListenFailure()
|
||||||
if s.cancelFunc != nil {
|
|
||||||
s.cancelFunc()
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
95
internal/server/listen_failure_test.go
Normal file
95
internal/server/listen_failure_test.go
Normal file
@@ -0,0 +1,95 @@
|
|||||||
|
package server_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"net"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
"go.uber.org/fx"
|
||||||
|
"sneak.berlin/go/webhooker/internal/globals"
|
||||||
|
"sneak.berlin/go/webhooker/internal/server"
|
||||||
|
)
|
||||||
|
|
||||||
|
// listenFailureDeadline is how long the app gets to give up after a
|
||||||
|
// listen it cannot satisfy. The defect this pins left the process
|
||||||
|
// reporting RUNNING for 183 seconds with nothing bound; a bind error
|
||||||
|
// is known instantly, so anything past a moment here is that defect
|
||||||
|
// back.
|
||||||
|
const listenFailureDeadline = 2 * time.Second
|
||||||
|
|
||||||
|
// lifecycleTimeout bounds the app's start and stop sequences so a
|
||||||
|
// wedged hook fails the test instead of hanging it.
|
||||||
|
const lifecycleTimeout = 15 * time.Second
|
||||||
|
|
||||||
|
// TestListenFailure_ShutsDownTheApp pins that a listener the server
|
||||||
|
// cannot bind terminates the application with a non-zero status.
|
||||||
|
//
|
||||||
|
// The fx OnStart hook returns as soon as the serving goroutine is
|
||||||
|
// spawned, so a bind failure is discovered after fx has already
|
||||||
|
// reported RUNNING. Nothing else in the graph observes it, and the
|
||||||
|
// process used to stay alive with no listener: down, but indis-
|
||||||
|
// tinguishable from healthy to systemd's Restart=on-failure and to
|
||||||
|
// Docker's restart policies, which is the state this test exists to
|
||||||
|
// keep from returning.
|
||||||
|
//
|
||||||
|
// The port is occupied by a listener this test holds open, on a
|
||||||
|
// kernel-chosen port, so the failure is the real EADDRINUSE the
|
||||||
|
// operator hits when a second instance starts. Loopback is enough to
|
||||||
|
// collide with the server's wildcard bind: a listening socket on a
|
||||||
|
// specific address blocks the wildcard from claiming the same port.
|
||||||
|
func TestListenFailure_ShutsDownTheApp(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
var listenCfg net.ListenConfig
|
||||||
|
|
||||||
|
occupied, err := listenCfg.Listen(
|
||||||
|
t.Context(), "tcp", "127.0.0.1:0",
|
||||||
|
)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
t.Cleanup(func() { _ = occupied.Close() })
|
||||||
|
|
||||||
|
addr, ok := occupied.Addr().(*net.TCPAddr)
|
||||||
|
require.True(t, ok, "listener is not TCP")
|
||||||
|
|
||||||
|
// The collaborators come from the wired graph rather than stubs,
|
||||||
|
// so the Server under test is the one that ships. Only the port
|
||||||
|
// is test-specific.
|
||||||
|
env := newTestEnv(t)
|
||||||
|
env.cfg.Port = addr.Port
|
||||||
|
|
||||||
|
app := fx.New(
|
||||||
|
fx.NopLogger,
|
||||||
|
fx.Supply(env.log, env.cfg, env.mw, env.hnd),
|
||||||
|
fx.Provide(globals.New, server.New),
|
||||||
|
fx.Invoke(func(*server.Server) {}),
|
||||||
|
)
|
||||||
|
|
||||||
|
startCtx, cancelStart := context.WithTimeout(
|
||||||
|
context.Background(), lifecycleTimeout,
|
||||||
|
)
|
||||||
|
defer cancelStart()
|
||||||
|
|
||||||
|
require.NoError(t, app.Start(startCtx))
|
||||||
|
|
||||||
|
select {
|
||||||
|
case sig := <-app.Wait():
|
||||||
|
require.Equal(
|
||||||
|
t, server.ListenFailureExitCode, sig.ExitCode,
|
||||||
|
"listen failure must exit non-zero",
|
||||||
|
)
|
||||||
|
case <-time.After(listenFailureDeadline):
|
||||||
|
t.Fatal("listen failure left the app running")
|
||||||
|
}
|
||||||
|
|
||||||
|
// The stop sequence still has to complete: the fix must reach
|
||||||
|
// shutdown through fx rather than around it.
|
||||||
|
stopCtx, cancelStop := context.WithTimeout(
|
||||||
|
context.Background(), lifecycleTimeout,
|
||||||
|
)
|
||||||
|
defer cancelStop()
|
||||||
|
|
||||||
|
require.NoError(t, app.Stop(stopCtx))
|
||||||
|
}
|
||||||
@@ -55,8 +55,11 @@ func (s *Server) setupGlobalMiddleware() {
|
|||||||
s.router.Use(s.mw.SecurityHeaders())
|
s.router.Use(s.mw.SecurityHeaders())
|
||||||
s.router.Use(s.mw.Logging())
|
s.router.Use(s.mw.Logging())
|
||||||
|
|
||||||
// Metrics middleware (only if credentials are configured)
|
// Metrics recording middleware, registered only when the
|
||||||
if s.params.Config.MetricsUsername != "" {
|
// endpoint that exposes what it records is served. The
|
||||||
|
// condition is the same MetricsAuthEnabled the /metrics mount
|
||||||
|
// in setupRoutes reads.
|
||||||
|
if s.params.Config.MetricsAuthEnabled() {
|
||||||
s.router.Use(s.mw.Metrics())
|
s.router.Use(s.mw.Metrics())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -103,8 +106,14 @@ func (s *Server) setupRoutes() {
|
|||||||
s.h.HandleHealthCheck(),
|
s.h.HandleHealthCheck(),
|
||||||
)
|
)
|
||||||
|
|
||||||
// set up authenticated /metrics route:
|
// Authenticated /metrics route. The condition is
|
||||||
if s.params.Config.MetricsUsername != "" {
|
// Config.MetricsAuthEnabled and never the username alone: a
|
||||||
|
// username with an empty password would otherwise mount the
|
||||||
|
// endpoint behind a credential map that accepts an empty
|
||||||
|
// password. Config rejects that combination at startup, and
|
||||||
|
// this reads the same value the startup log reports, so the
|
||||||
|
// two cannot disagree about whether the route exists.
|
||||||
|
if s.params.Config.MetricsAuthEnabled() {
|
||||||
s.router.Group(func(r chi.Router) {
|
s.router.Group(func(r chi.Router) {
|
||||||
r.Use(s.mw.MetricsAuth())
|
r.Use(s.mw.MetricsAuth())
|
||||||
r.Get(
|
r.Get(
|
||||||
|
|||||||
@@ -34,6 +34,13 @@ import (
|
|||||||
// the CSRF middleware executed.
|
// the CSRF middleware executed.
|
||||||
const csrfCookieName = "_gorilla_csrf"
|
const csrfCookieName = "_gorilla_csrf"
|
||||||
|
|
||||||
|
const (
|
||||||
|
// metricsUser and metricsAuthValue are the /metrics basic-auth
|
||||||
|
// credentials the metrics routing tests below configure.
|
||||||
|
metricsUser = "metrics"
|
||||||
|
metricsAuthValue = "s3cret"
|
||||||
|
)
|
||||||
|
|
||||||
type noopNotifier struct{}
|
type noopNotifier struct{}
|
||||||
|
|
||||||
func (n *noopNotifier) Notify([]delivery.Task) {}
|
func (n *noopNotifier) Notify([]delivery.Task) {}
|
||||||
@@ -69,9 +76,23 @@ type testEnv struct {
|
|||||||
func newTestEnv(t *testing.T) *testEnv {
|
func newTestEnv(t *testing.T) *testEnv {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
|
return newTestEnvWithConfig(t, &config.Config{
|
||||||
|
DataDir: t.TempDir(),
|
||||||
|
Environment: config.EnvironmentDev,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
// newTestEnvWithConfig is newTestEnv over a caller-supplied Config,
|
||||||
|
// for the routes whose existence the configuration decides. The same
|
||||||
|
// pointer reaches the router and every middleware, so a test cannot
|
||||||
|
// accidentally configure one and not the other.
|
||||||
|
func newTestEnvWithConfig(
|
||||||
|
t *testing.T, cfg *config.Config,
|
||||||
|
) *testEnv {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
var (
|
var (
|
||||||
log *logger.Logger
|
log *logger.Logger
|
||||||
cfg *config.Config
|
|
||||||
mw *middleware.Middleware
|
mw *middleware.Middleware
|
||||||
hnd *handlers.Handlers
|
hnd *handlers.Handlers
|
||||||
sess *session.Session
|
sess *session.Session
|
||||||
@@ -84,12 +105,7 @@ func newTestEnv(t *testing.T) *testEnv {
|
|||||||
fx.Provide(
|
fx.Provide(
|
||||||
globals.New,
|
globals.New,
|
||||||
logger.New,
|
logger.New,
|
||||||
func() *config.Config {
|
func() *config.Config { return cfg },
|
||||||
return &config.Config{
|
|
||||||
DataDir: t.TempDir(),
|
|
||||||
Environment: config.EnvironmentDev,
|
|
||||||
}
|
|
||||||
},
|
|
||||||
database.New,
|
database.New,
|
||||||
database.NewWebhookDBManager,
|
database.NewWebhookDBManager,
|
||||||
healthcheck.New,
|
healthcheck.New,
|
||||||
@@ -99,7 +115,7 @@ func newTestEnv(t *testing.T) *testEnv {
|
|||||||
middleware.New,
|
middleware.New,
|
||||||
handlers.New,
|
handlers.New,
|
||||||
),
|
),
|
||||||
fx.Populate(&log, &cfg, &mw, &hnd, &sess, &db, &dbMgr),
|
fx.Populate(&log, &mw, &hnd, &sess, &db, &dbMgr),
|
||||||
)
|
)
|
||||||
app.RequireStart()
|
app.RequireStart()
|
||||||
t.Cleanup(app.RequireStop)
|
t.Cleanup(app.RequireStop)
|
||||||
@@ -657,3 +673,119 @@ func TestSourceLogsBody_OtherUser404s(t *testing.T) {
|
|||||||
assert.Equal(t, http.StatusSeeOther, anon.Code)
|
assert.Equal(t, http.StatusSeeOther, anon.Code)
|
||||||
assert.Equal(t, "/pages/login", anon.Header().Get("Location"))
|
assert.Equal(t, "/pages/login", anon.Header().Get("Location"))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// metricsConfig is a Config differing from the routing default only
|
||||||
|
// in the two /metrics credentials.
|
||||||
|
func metricsConfig(
|
||||||
|
t *testing.T, username, password string,
|
||||||
|
) *config.Config {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
return &config.Config{
|
||||||
|
DataDir: t.TempDir(),
|
||||||
|
Environment: config.EnvironmentDev,
|
||||||
|
MetricsUsername: username,
|
||||||
|
MetricsPassword: password,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// metricsRequest asks the real router for /metrics with the given
|
||||||
|
// basic-auth credentials, or with no Authorization header when
|
||||||
|
// username is empty.
|
||||||
|
func (e *testEnv) metricsRequest(
|
||||||
|
username, password string,
|
||||||
|
) *httptest.ResponseRecorder {
|
||||||
|
req := httptest.NewRequestWithContext(
|
||||||
|
context.Background(), http.MethodGet, "/metrics", nil,
|
||||||
|
)
|
||||||
|
|
||||||
|
if username != "" {
|
||||||
|
req.SetBasicAuth(username, password)
|
||||||
|
}
|
||||||
|
|
||||||
|
w := httptest.NewRecorder()
|
||||||
|
e.router.ServeHTTP(w, req)
|
||||||
|
|
||||||
|
return w
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestMetricsRouteUnmountedWithoutCredentials pins that with neither
|
||||||
|
// credential configured the route does not exist, which is the
|
||||||
|
// documented behaviour and the only valid way for /metrics to be
|
||||||
|
// absent.
|
||||||
|
func TestMetricsRouteUnmountedWithoutCredentials(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
env := newTestEnvWithConfig(t, metricsConfig(t, "", ""))
|
||||||
|
|
||||||
|
assert.Equal(
|
||||||
|
t, http.StatusNotFound,
|
||||||
|
env.metricsRequest("", "").Code,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestMetricsRouteRequiresCredentials pins that with both credentials
|
||||||
|
// configured the route exists and every request that does not carry
|
||||||
|
// the configured pair is refused — including the empty password that
|
||||||
|
// a half-set configuration used to make sufficient.
|
||||||
|
func TestMetricsRouteRequiresCredentials(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
env := newTestEnvWithConfig(
|
||||||
|
t, metricsConfig(t, metricsUser, metricsAuthValue),
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.Equal(
|
||||||
|
t, http.StatusUnauthorized,
|
||||||
|
env.metricsRequest("", "").Code,
|
||||||
|
"no credentials must not reach the metrics handler",
|
||||||
|
)
|
||||||
|
assert.Equal(
|
||||||
|
t, http.StatusUnauthorized,
|
||||||
|
env.metricsRequest(metricsUser, "").Code,
|
||||||
|
"an empty password must not reach the metrics handler",
|
||||||
|
)
|
||||||
|
assert.Equal(
|
||||||
|
t, http.StatusUnauthorized,
|
||||||
|
env.metricsRequest(metricsUser, "wrong").Code,
|
||||||
|
)
|
||||||
|
|
||||||
|
ok := env.metricsRequest(metricsUser, metricsAuthValue)
|
||||||
|
assert.Equal(t, http.StatusOK, ok.Code)
|
||||||
|
assert.Contains(t, ok.Body.String(), "go_goroutines")
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestMetricsRouteUnmountedOnHalfSetConfig pins the defect from
|
||||||
|
// https://git.eeqj.de/sneak/webhooker/issues/205 at the routing
|
||||||
|
// layer. Config rejects a half-set pair at startup, so this Config
|
||||||
|
// cannot be reached from the environment; the assertion is that the
|
||||||
|
// route tree does not publish an endpoint accepting an empty
|
||||||
|
// password even when handed one anyway, because the mount and the
|
||||||
|
// startup log's hasMetricsAuth read the same value.
|
||||||
|
func TestMetricsRouteUnmountedOnHalfSetConfig(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
for _, tc := range []struct {
|
||||||
|
name string
|
||||||
|
username string
|
||||||
|
password string
|
||||||
|
}{
|
||||||
|
{name: "username only", username: metricsUser},
|
||||||
|
{name: "password only", password: metricsAuthValue},
|
||||||
|
} {
|
||||||
|
t.Run(tc.name, func(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
cfg := metricsConfig(t, tc.username, tc.password)
|
||||||
|
env := newTestEnvWithConfig(t, cfg)
|
||||||
|
|
||||||
|
assert.False(t, cfg.MetricsAuthEnabled())
|
||||||
|
assert.Equal(
|
||||||
|
t, http.StatusNotFound,
|
||||||
|
env.metricsRequest(
|
||||||
|
tc.username, tc.password,
|
||||||
|
).Code,
|
||||||
|
)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -50,6 +50,13 @@ const (
|
|||||||
minSentryFlush = 250 * time.Millisecond
|
minSentryFlush = 250 * time.Millisecond
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// ListenFailureExitCode is the status the process exits with when the
|
||||||
|
// HTTP listener cannot be established, or dies for a reason other
|
||||||
|
// than a requested shutdown. It must stay non-zero: systemd
|
||||||
|
// `Restart=on-failure` and Docker's restart policies key off it, and a
|
||||||
|
// zero exit would read as a deliberate stop.
|
||||||
|
const ListenFailureExitCode = 1
|
||||||
|
|
||||||
// SentryFlushBudget reports how long the Sentry flush may run when
|
// SentryFlushBudget reports how long the Sentry flush may run when
|
||||||
// remaining is the time left on the fx stop context after the HTTP
|
// remaining is the time left on the fx stop context after the HTTP
|
||||||
// drain. sentry.Flush takes a bare duration and honours no context,
|
// drain. sentry.Flush takes a bare duration and honours no context,
|
||||||
@@ -75,13 +82,13 @@ type ServerParams struct {
|
|||||||
Config *config.Config
|
Config *config.Config
|
||||||
Middleware *middleware.Middleware
|
Middleware *middleware.Middleware
|
||||||
Handlers *handlers.Handlers
|
Handlers *handlers.Handlers
|
||||||
|
Shutdowner fx.Shutdowner
|
||||||
}
|
}
|
||||||
|
|
||||||
// Server is the main HTTP server that wires up routes and manages
|
// Server is the main HTTP server that wires up routes and manages
|
||||||
// graceful shutdown.
|
// graceful shutdown.
|
||||||
type Server struct {
|
type Server struct {
|
||||||
startupTime time.Time
|
startupTime time.Time
|
||||||
exitCode int
|
|
||||||
sentryEnabled bool
|
sentryEnabled bool
|
||||||
log *slog.Logger
|
log *slog.Logger
|
||||||
cancelFunc context.CancelFunc
|
cancelFunc context.CancelFunc
|
||||||
@@ -159,7 +166,12 @@ func (s *Server) enableSentry() {
|
|||||||
s.sentryEnabled = true
|
s.sentryEnabled = true
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) serve() int {
|
// serve installs the signal watcher, starts the listener and blocks
|
||||||
|
// until the server's context is cancelled. The process exit status is
|
||||||
|
// fx's to decide — from a signal, or from the code
|
||||||
|
// shutdownOnListenFailure hands the Shutdowner — so this reports
|
||||||
|
// nothing back to its caller.
|
||||||
|
func (s *Server) serve() {
|
||||||
ctx, cancelFunc := context.WithCancel(context.Background())
|
ctx, cancelFunc := context.WithCancel(context.Background())
|
||||||
s.cancelFunc = cancelFunc
|
s.cancelFunc = cancelFunc
|
||||||
|
|
||||||
@@ -185,7 +197,30 @@ func (s *Server) serve() int {
|
|||||||
<-ctx.Done()
|
<-ctx.Done()
|
||||||
// Shutdown is handled by the fx OnStop hook (cleanShutdown).
|
// Shutdown is handled by the fx OnStop hook (cleanShutdown).
|
||||||
// Do not call cleanShutdown() here to avoid double invocation.
|
// Do not call cleanShutdown() here to avoid double invocation.
|
||||||
return s.exitCode
|
}
|
||||||
|
|
||||||
|
// shutdownOnListenFailure ends the application after the HTTP
|
||||||
|
// listener failed. The fx OnStart hook returns as soon as the serving
|
||||||
|
// goroutine is spawned, so nothing downstream of it ever learns that
|
||||||
|
// the listen failed: fx reports RUNNING and the process sits alive
|
||||||
|
// with nothing bound, which is invisible to systemd and Docker
|
||||||
|
// restart policies. Asking the Shutdowner to stop the app with a
|
||||||
|
// non-zero code is what turns that into a visible failure.
|
||||||
|
//
|
||||||
|
// The context cancel that follows only unwinds serve()'s own wait.
|
||||||
|
// The shutdown itself runs through fx's normal stop sequence, so the
|
||||||
|
// clean-shutdown drain in cleanShutdown is reached unchanged.
|
||||||
|
func (s *Server) shutdownOnListenFailure() {
|
||||||
|
err := s.params.Shutdowner.Shutdown(
|
||||||
|
fx.ExitCode(ListenFailureExitCode),
|
||||||
|
)
|
||||||
|
if err != nil {
|
||||||
|
s.log.Error("shutdown request failed", "error", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if s.cancelFunc != nil {
|
||||||
|
s.cancelFunc()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) cleanupForExit() {
|
func (s *Server) cleanupForExit() {
|
||||||
@@ -193,9 +228,6 @@ func (s *Server) cleanupForExit() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) cleanShutdown(ctx context.Context) {
|
func (s *Server) cleanShutdown(ctx context.Context) {
|
||||||
// initiate clean shutdown
|
|
||||||
s.exitCode = 0
|
|
||||||
|
|
||||||
ctxShutdown, shutdownCancel := context.WithTimeout(
|
ctxShutdown, shutdownCancel := context.WithTimeout(
|
||||||
ctx, ShutdownTimeout,
|
ctx, ShutdownTimeout,
|
||||||
)
|
)
|
||||||
|
|||||||
Reference in New Issue
Block a user