Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e1e9ba85ea | ||
|
|
08894ce16e |
@@ -157,11 +157,6 @@ private and reserved ranges — RFC 1918, loopback, CGNAT, link-local and
|
|||||||
the rest — are refused, which stops a target from being used to make
|
the rest — are refused, which stops a target from being used to make
|
||||||
webhooker probe the network it sits in.
|
webhooker probe the network it sits in.
|
||||||
|
|
||||||
Besides the private and reserved ranges, the default blocklist refuses
|
|
||||||
public cloud metadata addresses: currently only `168.63.129.16`, Azure's
|
|
||||||
WireServer, which serves an Azure VM its credentials. Because it is a
|
|
||||||
public address, listing it in `ALLOWED_EGRESS_CIDRS` reopens it.
|
|
||||||
|
|
||||||
That default is also inconvenient for the thing webhooker is mostly
|
That default is also inconvenient for the thing webhooker is mostly
|
||||||
for: taking a public webhook and forwarding it to something on your own
|
for: taking a public webhook and forwarding it to something on your own
|
||||||
network. A container on the same Docker network, a box on `10.x`, a
|
network. A container on the same Docker network, a box on `10.x`, a
|
||||||
@@ -200,16 +195,15 @@ Two things this setting cannot do:
|
|||||||
the list is always an allowlist; an empty list (the default) means
|
the list is always an allowlist; an empty list (the default) means
|
||||||
every private and reserved range stays refused. Note that
|
every private and reserved range stays refused. Note that
|
||||||
`0.0.0.0/0` gets you most of the way there anyway, per above.
|
`0.0.0.0/0` gets you most of the way there anyway, per above.
|
||||||
- **It cannot open link-local, or a cloud metadata endpoint at a
|
- **It cannot open link-local, or a cloud metadata endpoint that
|
||||||
non-public address that discloses credentials or user data.** An
|
discloses credentials or user data.** An address is on the list below
|
||||||
address is on the list below when it is not a public address and both
|
when both of these hold: the provider fixes it, so it cannot collide
|
||||||
of these hold: the provider fixes it, so it cannot collide with
|
with anything you run; and reaching it hands out credentials, user
|
||||||
anything you run; and reaching it hands out credentials, user data or
|
data or bootstrap material. Those stay blocked no matter what you
|
||||||
bootstrap material. Those stay blocked no matter what you list,
|
list, including when you list them outright or list a supernet such
|
||||||
including when you list them outright or list a supernet such as
|
as `0.0.0.0/0`, `::/0`, `fd00::/8` or `100.64.0.0/10`. Treat this as
|
||||||
`0.0.0.0/0`, `::/0`, `fd00::/8` or `100.64.0.0/10`. Treat this as best
|
best effort rather than a guarantee — it is a hand-maintained list
|
||||||
effort rather than a guarantee — it is a hand-maintained list and the
|
and the caveat below the table applies:
|
||||||
caveat below the table applies:
|
|
||||||
|
|
||||||
| Blocked unconditionally | What it is |
|
| Blocked unconditionally | What it is |
|
||||||
| ----------------------- | ---------- |
|
| ----------------------- | ---------- |
|
||||||
@@ -248,8 +242,7 @@ Two things this setting cannot do:
|
|||||||
encodings, which the default blocklist does not match. A publicly
|
encodings, which the default blocklist does not match. A publicly
|
||||||
routable metadata address is not listed here, because nothing on this
|
routable metadata address is not listed here, because nothing on this
|
||||||
list can be reopened and blocking one that way would leave you no
|
list can be reopened and blocking one that way would leave you no
|
||||||
escape hatch at all; Azure's `168.63.129.16` is refused by the default
|
escape hatch at all.
|
||||||
blocklist instead, as described above.
|
|
||||||
|
|
||||||
This list is not exhaustive of every cloud's metadata address — if
|
This list is not exhaustive of every cloud's metadata address — if
|
||||||
yours is not here, do not allowlist the block that contains it.
|
yours is not here, do not allowlist the block that contains it.
|
||||||
@@ -975,10 +968,15 @@ scratch file**: it holds committed transactions that are not yet in the
|
|||||||
have no readable schema at all. `-shm` is regenerable, but there is no
|
have no readable schema at all. `-shm` is regenerable, but there is no
|
||||||
reason to separate the two — copy the directory and you have them.
|
reason to separate the two — copy the directory and you have them.
|
||||||
|
|
||||||
A clean shutdown closes every database, which checkpoints and removes
|
A clean shutdown closes `webhooker.db` and every `events-*.db`, which
|
||||||
its sidecars; a killed or crashed instance leaves them, and they must be
|
checkpoints and removes their sidecars; a killed or crashed instance
|
||||||
carried with the `.db`. An archive the service has not opened since a
|
leaves them, and they must be carried with the `.db`. **Archive
|
||||||
crash keeps that crash's sidecars, even across a later clean stop.
|
databases are different**: their handle is not closed at shutdown, so
|
||||||
|
`archive-*.db-wal` and `-shm` normally survive a clean stop and the
|
||||||
|
`-wal` can hold every row the archive has. Measured on a stopped
|
||||||
|
instance: `archive-….db` 4096 bytes with no table, its `-wal` 157 KB
|
||||||
|
holding all 8 archived events. Copying `DATA_DIR` in full is what makes
|
||||||
|
this a non-issue; copying `.db` files out of it by name is not.
|
||||||
|
|
||||||
Configuration is **not** in `DATA_DIR` — it comes from the environment
|
Configuration is **not** in `DATA_DIR` — it comes from the environment
|
||||||
and from a `.env` file read out of the process working directory. Back
|
and from a `.env` file read out of the process working directory. Back
|
||||||
@@ -1053,9 +1051,10 @@ The file becomes self-contained again when the handle closes, which
|
|||||||
happens on the next write past the debounce window, when the connection
|
happens on the next write past the debounce window, when the connection
|
||||||
pool retires the idle connection (about a minute after the last write),
|
pool retires the idle connection (about a minute after the last write),
|
||||||
or at the idle archive sweep — measured, the same file was a complete
|
or at the idle archive sweep — measured, the same file was a complete
|
||||||
20 KB `.db` with no sidecars about a minute after its last write. A
|
20 KB `.db` with no sidecars about a minute after its last write.
|
||||||
clean stop closes it too. So either move `archive-{uuid}.db` together
|
Shutdown is **not** on that list: the archive handle is not closed when
|
||||||
with any `-wal`/`-shm` beside it, or wait until there are none.
|
the service stops. So either move `archive-{uuid}.db` together with any
|
||||||
|
`-wal`/`-shm` beside it, or wait until there are none.
|
||||||
|
|
||||||
### Restore
|
### Restore
|
||||||
|
|
||||||
@@ -1074,10 +1073,12 @@ with any `-wal`/`-shm` beside it, or wait until there are none.
|
|||||||
They are part of the database, and dropping a `-wal` silently
|
They are part of the database, and dropping a `-wal` silently
|
||||||
discards every transaction it still holds. An `.backup` set will not
|
discards every transaction it still holds. An `.backup` set will not
|
||||||
contain any: it writes a single consolidated file per database. A
|
contain any: it writes a single consolidated file per database. A
|
||||||
stop-and-copy set normally has none, because a clean stop closes
|
stop-and-copy set has none for `webhooker.db` or the `events-*.db`,
|
||||||
every database and checkpoints its sidecars away; the exception is an
|
because a clean stop closes those and checkpoints their sidecars
|
||||||
archive not opened since a crash. A copy salvaged from a crashed
|
away — but it will normally have them for `archive-*.db`, whose
|
||||||
instance has them for everything, and needs all of them.
|
handle stays open across shutdown, and those carry the archive's
|
||||||
|
rows. A copy salvaged from a crashed instance has them for
|
||||||
|
everything, and needs all of them.
|
||||||
|
|
||||||
4. **Fix ownership.** The container runs as the non-root `webhooker`
|
4. **Fix ownership.** The container runs as the non-root `webhooker`
|
||||||
user, UID 1000 / GID 1000. Restored files must be owned by (or
|
user, UID 1000 / GID 1000. Restored files must be owned by (or
|
||||||
@@ -2050,7 +2051,7 @@ rescans the database anyway).
|
|||||||
| ----------- | -------- |
|
| ----------- | -------- |
|
||||||
| **Closed** | Normal operation. Deliveries flow through. Consecutive failures are counted. |
|
| **Closed** | Normal operation. Deliveries flow through. Consecutive failures are counted. |
|
||||||
| **Open** | Target appears down. Deliveries are skipped and rescheduled for after the cooldown. |
|
| **Open** | Target appears down. Deliveries are skipped and rescheduled for after the cooldown. |
|
||||||
| **Half-Open** | Cooldown expired. One probe delivery is allowed to test if the target has recovered. Other deliveries are rescheduled for one whole cooldown later. |
|
| **Half-Open** | Cooldown expired. One probe delivery is allowed to test if the target has recovered. |
|
||||||
|
|
||||||
**Transitions:**
|
**Transitions:**
|
||||||
|
|
||||||
@@ -2088,9 +2089,7 @@ operations), and log targets (stdout) do not use circuit breakers.
|
|||||||
When a circuit is open and a new delivery arrives, the engine marks the
|
When a circuit is open and a new delivery arrives, the engine marks the
|
||||||
delivery as `retrying` and schedules a retry timer for after the
|
delivery as `retrying` and schedules a retry timer for after the
|
||||||
remaining cooldown period. This ensures no deliveries are lost — they're
|
remaining cooldown period. This ensures no deliveries are lost — they're
|
||||||
just delayed until the target is healthy again. A delivery already in
|
just delayed until the target is healthy again.
|
||||||
`retrying` keeps that status without another database write each time
|
|
||||||
the breaker turns it away.
|
|
||||||
|
|
||||||
### Metrics
|
### Metrics
|
||||||
|
|
||||||
@@ -2104,7 +2103,7 @@ 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 actually dispatched to a target. A delivery a circuit breaker refused is not one: it is counted as a retry instead, but only when the refusal moves it into `retrying` |
|
| `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` |
|
||||||
@@ -2721,7 +2720,7 @@ abuse limit later; they are tracked as future work.
|
|||||||
| ------ | --------------------------- | ----------- |
|
| ------ | --------------------------- | ----------- |
|
||||||
| `GET` | `/` | Root redirect, 303 (authenticated → `/sources`, unauthenticated → `/pages/login`) |
|
| `GET` | `/` | Root redirect, 303 (authenticated → `/sources`, unauthenticated → `/pages/login`) |
|
||||||
| `GET` | `/.well-known/healthcheck` | Health check (JSON: `status`, `now`, `uptimeSeconds`, `uptimeHuman`, `version`, `appname`, `maintenanceMode`) |
|
| `GET` | `/.well-known/healthcheck` | Health check (JSON: `status`, `now`, `uptimeSeconds`, `uptimeHuman`, `version`, `appname`, `maintenanceMode`) |
|
||||||
| `GET`, `HEAD` | `/s/*` | Static file serving (embedded CSS, JS). `GET` and `HEAD` only — `POST`, `PUT`, `PATCH`, `DELETE`, `OPTIONS`, `TRACE` and `CONNECT` are answered `405 Method Not Allowed` with `Allow: GET, HEAD`. Any other method (such as `PROPFIND`) is refused by chi before it reaches this route, and gets `405` without an `Allow` header. Pinned by `TestStaticServesOnlyGetAndHead` |
|
| any | `/s/*` | Static file serving (embedded CSS, JS). Mounted for every method, not just `GET`/`HEAD`: chi's `Mount` registers all methods and `http.FileServer` special-cases only `HEAD` (by omitting the body), so a `POST` or `DELETE` to an asset is answered `200` with the file. Pinned by `TestStaticServesEveryMethod` |
|
||||||
| `POST` | `/webhook/{uuid}` | Webhook receiver endpoint. `POST` only — every other method is answered `405 Method Not Allowed` with `Allow: POST`. Rate limited (see [Rate Limiting](#rate-limiting)) |
|
| `POST` | `/webhook/{uuid}` | Webhook receiver endpoint. `POST` only — every other method is answered `405 Method Not Allowed` with `Allow: POST`. Rate limited (see [Rate Limiting](#rate-limiting)) |
|
||||||
|
|
||||||
#### Authentication Endpoints
|
#### Authentication Endpoints
|
||||||
@@ -3087,8 +3086,7 @@ each hook. The order, read off the fx stop-hook log:
|
|||||||
3. `server` — the HTTP drain, bounded separately by
|
3. `server` — the HTTP drain, bounded separately by
|
||||||
`server.ShutdownTimeout` (**3 seconds**), then a Sentry flush if
|
`server.ShutdownTimeout` (**3 seconds**), then a Sentry flush if
|
||||||
`SENTRY_DSN` is set
|
`SENTRY_DSN` is set
|
||||||
4. `delivery.Engine` — waits for its workers, then closes the archive
|
4. `delivery.Engine`
|
||||||
databases
|
|
||||||
5. `healthcheck`
|
5. `healthcheck`
|
||||||
6. `WebhookDBManager`
|
6. `WebhookDBManager`
|
||||||
7. the database close
|
7. the database close
|
||||||
|
|||||||
@@ -16,7 +16,6 @@ import (
|
|||||||
"sneak.berlin/go/webhooker/internal/handlers"
|
"sneak.berlin/go/webhooker/internal/handlers"
|
||||||
"sneak.berlin/go/webhooker/internal/healthcheck"
|
"sneak.berlin/go/webhooker/internal/healthcheck"
|
||||||
"sneak.berlin/go/webhooker/internal/logger"
|
"sneak.berlin/go/webhooker/internal/logger"
|
||||||
"sneak.berlin/go/webhooker/internal/metrics"
|
|
||||||
"sneak.berlin/go/webhooker/internal/middleware"
|
"sneak.berlin/go/webhooker/internal/middleware"
|
||||||
"sneak.berlin/go/webhooker/internal/resetpw"
|
"sneak.berlin/go/webhooker/internal/resetpw"
|
||||||
"sneak.berlin/go/webhooker/internal/server"
|
"sneak.berlin/go/webhooker/internal/server"
|
||||||
@@ -178,10 +177,6 @@ func newApp() *fx.App {
|
|||||||
healthcheck.New,
|
healthcheck.New,
|
||||||
session.New,
|
session.New,
|
||||||
handlers.New,
|
handlers.New,
|
||||||
// The registry /metrics serves, and the delivery
|
|
||||||
// collectors registered on it.
|
|
||||||
metrics.NewRegistry,
|
|
||||||
metrics.New,
|
|
||||||
middleware.New,
|
middleware.New,
|
||||||
// The one SSRF guard both target-creation validation
|
// The one SSRF guard both target-creation validation
|
||||||
// and the delivery dialer consult, so they cannot
|
// and the delivery dialer consult, so they cannot
|
||||||
|
|||||||
@@ -192,10 +192,9 @@ type Config struct {
|
|||||||
// alwaysBlockedNetworks stays blocked no matter what is listed
|
// alwaysBlockedNetworks stays blocked no matter what is listed
|
||||||
// here. That set is link-local plus the cloud metadata
|
// here. That set is link-local plus the cloud metadata
|
||||||
// endpoints outside it that disclose credentials or user data
|
// endpoints outside it that disclose credentials or user data
|
||||||
// at a provider-fixed, non-public address; it is not
|
// at a provider-fixed address; it is not exhaustive of every
|
||||||
// exhaustive of every cloud's metadata address. See
|
// cloud's metadata address. See alwaysBlockedNetworks for the
|
||||||
// alwaysBlockedNetworks for the authoritative list and the
|
// authoritative list and the criterion it is built from.
|
||||||
// criterion it is built from.
|
|
||||||
AllowedEgressCIDRs []netip.Prefix
|
AllowedEgressCIDRs []netip.Prefix
|
||||||
|
|
||||||
params *ConfigParams
|
params *ConfigParams
|
||||||
@@ -747,14 +746,12 @@ func (c *Config) warnEgressAllowlist(log *slog.Logger) {
|
|||||||
|
|
||||||
log.Warn(
|
log.Warn(
|
||||||
"ALLOWED_EGRESS_CIDRS lets delivery targets reach these "+
|
"ALLOWED_EGRESS_CIDRS lets delivery targets reach these "+
|
||||||
"otherwise-blocked networks. Anyone who can create a "+
|
"otherwise-blocked private/reserved networks. Anyone "+
|
||||||
"delivery target can now make this process issue "+
|
"who can create a delivery target can now make this "+
|
||||||
"requests into them, and read back the response. Only "+
|
"process issue requests into them, and read back the "+
|
||||||
"the addresses the README lists as blocked "+
|
"response. Link-local and the known cloud instance "+
|
||||||
"unconditionally stay blocked regardless of what is "+
|
"metadata endpoints outside it stay blocked "+
|
||||||
"listed here; a public cloud metadata address such as "+
|
"regardless of what is listed here.",
|
||||||
"168.63.129.16 is reachable once it, or a block "+
|
|
||||||
"covering it, is listed.",
|
|
||||||
"allowedEgressCIDRs",
|
"allowedEgressCIDRs",
|
||||||
strings.Join(PrefixStrings(c.AllowedEgressCIDRs), ","),
|
strings.Join(PrefixStrings(c.AllowedEgressCIDRs), ","),
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -834,13 +834,12 @@ func TestEgressAllowlistWarning(t *testing.T) {
|
|||||||
// to be able to read back which networks are open.
|
// to be able to read back which networks are open.
|
||||||
assert.Contains(t, logged, "10.0.0.0/8")
|
assert.Contains(t, logged, "10.0.0.0/8")
|
||||||
assert.Contains(t, logged, "127.0.0.0/8")
|
assert.Contains(t, logged, "127.0.0.0/8")
|
||||||
// What stays shut is the whole unconditional set, not
|
// What stays shut. Asserted on the clause naming the
|
||||||
// link-local alone; a public metadata address is not in
|
// wider set rather than on "Link-local" alone, so the
|
||||||
// it, so a listed block covering it opens it.
|
// string cannot narrow back to link-local only while
|
||||||
assert.Contains(t, logged, "blocked unconditionally")
|
// the always-blocked set covers ULA, CGNAT and two
|
||||||
assert.Contains(t, logged, "168.63.129.16 is reachable")
|
// public metadata addresses as well.
|
||||||
// The listed blocks need not be private or reserved.
|
assert.Contains(t, logged, "metadata endpoints outside it")
|
||||||
assert.NotContains(t, logged, "private/reserved")
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -76,20 +76,12 @@ func (cb *CircuitBreaker) Allow() bool {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// CooldownRemaining returns how long a delivery that Allow refused
|
// CooldownRemaining returns how much time is left before
|
||||||
// should wait before it is tried again. Closed, it returns zero.
|
// an open circuit transitions to half-open.
|
||||||
// Open, it returns what is left of the cooldown, or zero once that
|
|
||||||
// has passed. Half-open, it returns the whole cooldown: the one
|
|
||||||
// probe delivery is still in flight, and if it fails the circuit
|
|
||||||
// reopens for that long.
|
|
||||||
func (cb *CircuitBreaker) CooldownRemaining() time.Duration {
|
func (cb *CircuitBreaker) CooldownRemaining() time.Duration {
|
||||||
cb.mu.Lock()
|
cb.mu.Lock()
|
||||||
defer cb.mu.Unlock()
|
defer cb.mu.Unlock()
|
||||||
|
|
||||||
if cb.state == CircuitHalfOpen {
|
|
||||||
return cb.cooldown
|
|
||||||
}
|
|
||||||
|
|
||||||
if cb.state != CircuitOpen {
|
if cb.state != CircuitOpen {
|
||||||
return 0
|
return 0
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -267,7 +267,7 @@ func TestCircuitBreaker_CooldownRemaining_ClosedReturnsZero(
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsCooldown(
|
func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsZero(
|
||||||
t *testing.T,
|
t *testing.T,
|
||||||
) {
|
) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
@@ -282,11 +282,9 @@ func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsCooldown(
|
|||||||
|
|
||||||
require.True(t, cb.Allow())
|
require.True(t, cb.Allow())
|
||||||
|
|
||||||
// The cooldown newShortCooldownCB gives the breaker.
|
assert.Equal(t, time.Duration(0),
|
||||||
assert.Equal(t, 50*time.Millisecond,
|
|
||||||
cb.CooldownRemaining(),
|
cb.CooldownRemaining(),
|
||||||
"a delivery refused while half-open should wait "+
|
"half-open circuit should have zero cooldown remaining",
|
||||||
"a whole cooldown",
|
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -148,7 +148,6 @@ type EngineParams struct {
|
|||||||
DBManager *database.WebhookDBManager
|
DBManager *database.WebhookDBManager
|
||||||
Logger *logger.Logger
|
Logger *logger.Logger
|
||||||
SSRFGuard *Guard
|
SSRFGuard *Guard
|
||||||
Metrics *metrics.Set
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Engine processes queued deliveries in the background
|
// Engine processes queued deliveries in the background
|
||||||
@@ -168,10 +167,10 @@ type Engine struct {
|
|||||||
retryCh chan Task
|
retryCh chan Task
|
||||||
workers int
|
workers int
|
||||||
|
|
||||||
// mtr is the delivery metric set. Production wires the one
|
// mtr is the delivery metric set. Production wires the
|
||||||
// registered on the registry /metrics serves; a test can
|
// process-wide one; a test can substitute a set registered on
|
||||||
// substitute a set registered on a registry it holds, so it can
|
// a private registry so its assertions are not disturbed by
|
||||||
// gather what its own deliveries recorded.
|
// deliveries other tests are making at the same time.
|
||||||
mtr *metrics.Set
|
mtr *metrics.Set
|
||||||
|
|
||||||
// targets maps each target type to its implementation.
|
// targets maps each target type to its implementation.
|
||||||
@@ -205,7 +204,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,
|
||||||
mtr: params.Metrics,
|
mtr: metrics.Default(),
|
||||||
}
|
}
|
||||||
|
|
||||||
e.initTargets(&http.Client{
|
e.initTargets(&http.Client{
|
||||||
@@ -363,15 +362,6 @@ func (e *Engine) start() {
|
|||||||
// stop cancels the worker pool's context and waits for the pool
|
// stop cancels the worker pool's context and waits for the pool
|
||||||
// to drain, bounded by the stop hook's context: a wedged worker
|
// to drain, bounded by the stop hook's context: a wedged worker
|
||||||
// must not hang the process past fx's stop timeout.
|
// must not hang the process past fx's stop timeout.
|
||||||
//
|
|
||||||
// Once the pool has drained it closes the archive writers, so a
|
|
||||||
// clean stop leaves no archive -wal behind. Nothing else holds a
|
|
||||||
// writer for long by then: the archive sweeper stops before the
|
|
||||||
// engine, and deleting a webhook only closes one. If the pool did
|
|
||||||
// not drain in time, the writers are left open, as a kill would
|
|
||||||
// leave them. Closing them would wait for any write in progress,
|
|
||||||
// and a worker still running would then open new writers that
|
|
||||||
// nothing closes, so it gains nothing over a kill.
|
|
||||||
func (e *Engine) stop(ctx context.Context) error {
|
func (e *Engine) stop(ctx context.Context) error {
|
||||||
e.log.Info("delivery engine stopping")
|
e.log.Info("delivery engine stopping")
|
||||||
|
|
||||||
@@ -386,8 +376,6 @@ func (e *Engine) stop(ctx context.Context) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
e.dbTarget.evictAll()
|
|
||||||
|
|
||||||
e.log.Info("delivery engine stopped")
|
e.log.Info("delivery engine stopped")
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -2,8 +2,6 @@ package delivery_test
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
|
||||||
"path/filepath"
|
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -271,88 +269,3 @@ func TestEngine_StopHookHonoursStopTimeout(t *testing.T) {
|
|||||||
|
|
||||||
requireStopHookExpires(t, lc.hooks[0], "delivery engine")
|
requireStopHookExpires(t, lc.hooks[0], "delivery engine")
|
||||||
}
|
}
|
||||||
|
|
||||||
// deliverToArchive runs one delivery to a database target through
|
|
||||||
// the running engine and returns the webhook's archive file path.
|
|
||||||
// The archive writer holds the file open afterwards.
|
|
||||||
func deliverToArchive(t *testing.T, s iSetup) string {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
deliveryID, task := seedLogTask(t, s)
|
|
||||||
task.TargetType = database.TargetTypeDatabase
|
|
||||||
|
|
||||||
s.Engine.Notify([]delivery.Task{task})
|
|
||||||
|
|
||||||
iWaitForDelivered(t, s.WebhookDB, deliveryID)
|
|
||||||
|
|
||||||
return filepath.Join(
|
|
||||||
filepath.Dir(s.DBMgr.DBPath(s.WebhookID)),
|
|
||||||
fmt.Sprintf("archive-%s.db", s.WebhookID),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestEngine_StopHookClosesArchives is the regression test for an
|
|
||||||
// archive split across two files by a clean stop. The engine never
|
|
||||||
// closed its archive writers, so after a stop the archived rows
|
|
||||||
// could sit in archive-{id}.db-wal while archive-{id}.db held no
|
|
||||||
// table at all, and copying the .db on its own gave an empty
|
|
||||||
// database.
|
|
||||||
func TestEngine_StopHookClosesArchives(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
lc := startEngineViaHook(t, s.Engine)
|
|
||||||
|
|
||||||
path := deliverToArchive(t, s)
|
|
||||||
require.FileExists(
|
|
||||||
t, path+"-wal",
|
|
||||||
"an open archive should have a -wal for the stop to remove",
|
|
||||||
)
|
|
||||||
|
|
||||||
require.NoError(t, lc.hooks[0].OnStop(context.Background()))
|
|
||||||
|
|
||||||
wals, err := filepath.Glob(
|
|
||||||
filepath.Join(filepath.Dir(path), "archive-*.db-wal"),
|
|
||||||
)
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.Empty(
|
|
||||||
t, wals, "a clean stop must leave no archive -wal behind",
|
|
||||||
)
|
|
||||||
|
|
||||||
// With no -wal beside it, the row can only be in the .db.
|
|
||||||
count, err := countArchivedRows(path)
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.Equal(t, int64(1), count)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestEngine_StopHookTimeoutLeavesArchivesOpen covers a stop whose
|
|
||||||
// budget runs out while a worker is still running. The archive
|
|
||||||
// writers are left open, as a kill would leave them: closing them
|
|
||||||
// would wait for any write in progress, and that worker would then
|
|
||||||
// open new writers that nothing closes.
|
|
||||||
func TestEngine_StopHookTimeoutLeavesArchivesOpen(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
lc := startEngineViaHook(t, s.Engine)
|
|
||||||
|
|
||||||
deliverToArchive(t, s)
|
|
||||||
|
|
||||||
release := make(chan struct{})
|
|
||||||
|
|
||||||
t.Cleanup(func() {
|
|
||||||
close(release)
|
|
||||||
s.Engine.EvictWebhook(s.WebhookID)
|
|
||||||
})
|
|
||||||
|
|
||||||
s.Engine.ExportWedgeWorker(release)
|
|
||||||
|
|
||||||
requireStopHookExpires(t, lc.hooks[0], "delivery engine")
|
|
||||||
|
|
||||||
require.True(
|
|
||||||
t, s.Engine.ExportArchiveHandleOpen(s.WebhookID),
|
|
||||||
"a stop that timed out must not close archive writers",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -17,7 +17,6 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/google/uuid"
|
"github.com/google/uuid"
|
||||||
"github.com/prometheus/client_golang/prometheus"
|
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
"gorm.io/driver/sqlite"
|
"gorm.io/driver/sqlite"
|
||||||
@@ -25,7 +24,6 @@ import (
|
|||||||
_ "modernc.org/sqlite"
|
_ "modernc.org/sqlite"
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
"sneak.berlin/go/webhooker/internal/delivery"
|
"sneak.berlin/go/webhooker/internal/delivery"
|
||||||
"sneak.berlin/go/webhooker/internal/metrics"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// testContentType is the event content type used in tests.
|
// testContentType is the event content type used in tests.
|
||||||
@@ -896,100 +894,6 @@ func TestDeliverHTTP_CircuitBreakerBlocks(t *testing.T) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// recordingScheduler keeps the delay of every retry it is asked to
|
|
||||||
// schedule, and schedules nothing.
|
|
||||||
type recordingScheduler struct {
|
|
||||||
delays []time.Duration
|
|
||||||
}
|
|
||||||
|
|
||||||
func (s *recordingScheduler) ScheduleRetry(
|
|
||||||
_ delivery.Task, delay time.Duration,
|
|
||||||
) {
|
|
||||||
s.delays = append(s.delays, delay)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestDeliverHTTP_HalfOpenBreakerDelaysQueuedTasks proves that while a
|
|
||||||
// half-open breaker's one probe delivery is in flight, every other task
|
|
||||||
// for the target is put back with a whole cooldown as its delay rather
|
|
||||||
// than none, and that its status is written the first time the breaker
|
|
||||||
// turns it away and not on each pass after that.
|
|
||||||
func TestDeliverHTTP_HalfOpenBreakerDelaysQueuedTasks(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
db := testWebhookDB(t)
|
|
||||||
e := testEngine(t, 1)
|
|
||||||
|
|
||||||
// Every write of retrying moves the retry counter, so on a registry
|
|
||||||
// this test owns the counter is the number of those writes.
|
|
||||||
reg := prometheus.NewRegistry()
|
|
||||||
e.ExportSetMetrics(metrics.New(reg))
|
|
||||||
|
|
||||||
targetID := uuid.New().String()
|
|
||||||
cb := newShortCooldownCB(t)
|
|
||||||
e.ExportSetCircuitBreaker(targetID, cb)
|
|
||||||
|
|
||||||
for range delivery.ExportDefaultFailureThreshold {
|
|
||||||
cb.RecordFailure()
|
|
||||||
}
|
|
||||||
|
|
||||||
time.Sleep(60 * time.Millisecond)
|
|
||||||
|
|
||||||
require.True(t, cb.Allow(), "the probe delivery should go through")
|
|
||||||
require.Equal(t, delivery.CircuitHalfOpen, cb.State())
|
|
||||||
|
|
||||||
cfg := newHTTPTargetConfig(
|
|
||||||
"http://will-not-be-called.invalid",
|
|
||||||
)
|
|
||||||
sched := &recordingScheduler{}
|
|
||||||
|
|
||||||
const queued, passes = 3, 4
|
|
||||||
|
|
||||||
for range queued {
|
|
||||||
event := seedEvent(t, db, `{"cb":"half-open"}`)
|
|
||||||
dlv := seedDelivery(
|
|
||||||
t, db, event.ID, targetID,
|
|
||||||
database.DeliveryStatusPending,
|
|
||||||
)
|
|
||||||
|
|
||||||
for range passes {
|
|
||||||
// Each pass starts from the stored row, as a retry does.
|
|
||||||
var row database.Delivery
|
|
||||||
|
|
||||||
require.NoError(t, db.First(
|
|
||||||
&row, "id = ?", dlv.ID,
|
|
||||||
).Error)
|
|
||||||
|
|
||||||
fix := buildHTTPFixture(
|
|
||||||
row, event, targetID,
|
|
||||||
"test-cb-half-open", cfg, 5, 1,
|
|
||||||
)
|
|
||||||
|
|
||||||
e.ExportDeliverHTTPWithScheduler(
|
|
||||||
context.TODO(), db, fix.Delivery, fix.Task, sched,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
assertDeliveryStatus(t, db, dlv.ID,
|
|
||||||
database.DeliveryStatusRetrying,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
require.Len(t, sched.delays, queued*passes)
|
|
||||||
|
|
||||||
for _, delay := range sched.delays {
|
|
||||||
// The cooldown newShortCooldownCB gives the breaker.
|
|
||||||
assert.Equal(t, 50*time.Millisecond, delay,
|
|
||||||
"a task turned away while half-open should wait "+
|
|
||||||
"a whole cooldown",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
assert.InDelta(t, float64(queued),
|
|
||||||
mCounter(t, reg, mRetries, mTypeHTTP), 0,
|
|
||||||
"status should be written once per task, not once per pass",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestGetCircuitBreaker_CreatesOnDemand(t *testing.T) {
|
func TestGetCircuitBreaker_CreatesOnDemand(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
@@ -1166,10 +1070,6 @@ func TestIsForwardableHeader(t *testing.T) {
|
|||||||
assert.False(t,
|
assert.False(t,
|
||||||
delivery.ExportIsForwardableHeader("Content-Length"),
|
delivery.ExportIsForwardableHeader("Content-Length"),
|
||||||
)
|
)
|
||||||
|
|
||||||
assert.False(t,
|
|
||||||
delivery.ExportIsForwardableHeader("Content-Type"),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestTruncate(t *testing.T) {
|
func TestTruncate(t *testing.T) {
|
||||||
@@ -1251,81 +1151,6 @@ func TestDoHTTPRequest_ForwardsHeaders(t *testing.T) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// The event's stored inbound headers carry the same Content-Type the
|
|
||||||
// receiver saved as the event's ContentType, so a delivery could send
|
|
||||||
// it twice. It must go out exactly once, with a Content-Type configured
|
|
||||||
// on the target winning, then the event's ContentType.
|
|
||||||
func TestApplyRequestHeaders_SendsOneContentType(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
cases := map[string]struct {
|
|
||||||
inbound string
|
|
||||||
event string
|
|
||||||
configured string
|
|
||||||
want []string
|
|
||||||
}{
|
|
||||||
"inbound and event agree": {
|
|
||||||
inbound: testContentType,
|
|
||||||
event: testContentType,
|
|
||||||
want: []string{testContentType},
|
|
||||||
},
|
|
||||||
"inbound and event disagree": {
|
|
||||||
inbound: "text/plain",
|
|
||||||
event: testContentType,
|
|
||||||
want: []string{testContentType},
|
|
||||||
},
|
|
||||||
"event has none": {
|
|
||||||
inbound: testContentType,
|
|
||||||
want: nil,
|
|
||||||
},
|
|
||||||
"target configures its own": {
|
|
||||||
inbound: testContentType,
|
|
||||||
event: testContentType,
|
|
||||||
configured: "application/xml",
|
|
||||||
want: []string{"application/xml"},
|
|
||||||
},
|
|
||||||
}
|
|
||||||
|
|
||||||
for name, tc := range cases {
|
|
||||||
t.Run(name, func(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
inbound, err := json.Marshal(map[string][]string{
|
|
||||||
headerContentType: {tc.inbound},
|
|
||||||
})
|
|
||||||
require.NoError(t, err)
|
|
||||||
|
|
||||||
cfg := &delivery.HTTPTargetConfig{}
|
|
||||||
if tc.configured != "" {
|
|
||||||
cfg.Headers = map[string]string{
|
|
||||||
headerContentType: tc.configured,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
req, err := http.NewRequestWithContext(
|
|
||||||
context.Background(),
|
|
||||||
http.MethodPost,
|
|
||||||
"https://target.example.com/hook",
|
|
||||||
http.NoBody,
|
|
||||||
)
|
|
||||||
require.NoError(t, err)
|
|
||||||
|
|
||||||
delivery.ExportApplyRequestHeaders(
|
|
||||||
req,
|
|
||||||
&database.Event{
|
|
||||||
Headers: string(inbound),
|
|
||||||
ContentType: tc.event,
|
|
||||||
},
|
|
||||||
cfg,
|
|
||||||
)
|
|
||||||
|
|
||||||
assert.Equal(t,
|
|
||||||
tc.want, req.Header.Values(headerContentType),
|
|
||||||
)
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestProcessDelivery_RoutesToCorrectHandler(
|
func TestProcessDelivery_RoutesToCorrectHandler(
|
||||||
t *testing.T,
|
t *testing.T,
|
||||||
) {
|
) {
|
||||||
|
|||||||
@@ -9,7 +9,6 @@ import (
|
|||||||
"net/url"
|
"net/url"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/prometheus/client_golang/prometheus"
|
|
||||||
"go.uber.org/fx"
|
"go.uber.org/fx"
|
||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
@@ -102,19 +101,6 @@ func (e *Engine) ExportDeliverHTTP(
|
|||||||
e.httpTarget.Deliver(ctx, webhookDB, d, task, e)
|
e.httpTarget.Deliver(ctx, webhookDB, d, task, e)
|
||||||
}
|
}
|
||||||
|
|
||||||
// ExportDeliverHTTPWithScheduler delivers via the http target, handing
|
|
||||||
// any retry to sched instead of the engine, so a test can see the
|
|
||||||
// delay each retry is given.
|
|
||||||
func (e *Engine) ExportDeliverHTTPWithScheduler(
|
|
||||||
ctx context.Context,
|
|
||||||
webhookDB *gorm.DB,
|
|
||||||
d *database.Delivery,
|
|
||||||
task *Task,
|
|
||||||
sched Scheduler,
|
|
||||||
) {
|
|
||||||
e.httpTarget.Deliver(ctx, webhookDB, d, task, sched)
|
|
||||||
}
|
|
||||||
|
|
||||||
// ExportDeliverDatabase delivers via the database target.
|
// ExportDeliverDatabase delivers via the database target.
|
||||||
func (e *Engine) ExportDeliverDatabase(
|
func (e *Engine) ExportDeliverDatabase(
|
||||||
webhookDB *gorm.DB, d *database.Delivery,
|
webhookDB *gorm.DB, d *database.Delivery,
|
||||||
@@ -193,14 +179,6 @@ func (e *Engine) ExportGetCircuitBreaker(
|
|||||||
return e.httpTarget.getCircuitBreaker(targetID)
|
return e.httpTarget.getCircuitBreaker(targetID)
|
||||||
}
|
}
|
||||||
|
|
||||||
// ExportSetCircuitBreaker makes cb the http target's circuit breaker
|
|
||||||
// for targetID, so a test can use one with a short cooldown.
|
|
||||||
func (e *Engine) ExportSetCircuitBreaker(
|
|
||||||
targetID string, cb *CircuitBreaker,
|
|
||||||
) {
|
|
||||||
e.httpTarget.circuitBreakers.Store(targetID, cb)
|
|
||||||
}
|
|
||||||
|
|
||||||
// ExportParseHTTPConfig exposes parseHTTPConfig.
|
// ExportParseHTTPConfig exposes parseHTTPConfig.
|
||||||
func (e *Engine) ExportParseHTTPConfig(
|
func (e *Engine) ExportParseHTTPConfig(
|
||||||
configJSON string,
|
configJSON string,
|
||||||
@@ -353,21 +331,6 @@ func (e *Engine) ExportFailMissingTarget(
|
|||||||
e.failMissingTarget(webhookDB, webhookID, d)
|
e.failMissingTarget(webhookDB, webhookID, d)
|
||||||
}
|
}
|
||||||
|
|
||||||
// ExportSendRecoveredDeliveries exposes sendRecoveredDeliveries, so a
|
|
||||||
// test can hand it a target map that lacks a delivery's target.
|
|
||||||
func (e *Engine) ExportSendRecoveredDeliveries(
|
|
||||||
ctx context.Context,
|
|
||||||
webhookDB *gorm.DB,
|
|
||||||
deliveries []database.Delivery,
|
|
||||||
webhookID string,
|
|
||||||
targetMap map[string]database.Target,
|
|
||||||
settled map[string]struct{},
|
|
||||||
) {
|
|
||||||
e.sendRecoveredDeliveries(
|
|
||||||
ctx, webhookDB, deliveries, webhookID, targetMap, settled,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// ExportDeliveryCh returns the delivery channel.
|
// ExportDeliveryCh returns the delivery channel.
|
||||||
func (e *Engine) ExportDeliveryCh() chan Task {
|
func (e *Engine) ExportDeliveryCh() chan Task {
|
||||||
return e.deliveryCh
|
return e.deliveryCh
|
||||||
@@ -390,7 +353,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,
|
||||||
mtr: metrics.New(prometheus.NewRegistry()),
|
mtr: metrics.Default(),
|
||||||
}
|
}
|
||||||
e.initTargets(client)
|
e.initTargets(client)
|
||||||
|
|
||||||
@@ -405,7 +368,7 @@ func NewTestEngineSmallRetry(
|
|||||||
e := &Engine{
|
e := &Engine{
|
||||||
log: log,
|
log: log,
|
||||||
retryCh: make(chan Task, 1),
|
retryCh: make(chan Task, 1),
|
||||||
mtr: metrics.New(prometheus.NewRegistry()),
|
mtr: metrics.Default(),
|
||||||
}
|
}
|
||||||
e.initTargets(nil)
|
e.initTargets(nil)
|
||||||
|
|
||||||
@@ -428,7 +391,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,
|
||||||
mtr: metrics.New(prometheus.NewRegistry()),
|
mtr: metrics.Default(),
|
||||||
}
|
}
|
||||||
e.initTargets(client)
|
e.initTargets(client)
|
||||||
|
|
||||||
@@ -436,7 +399,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 registry it holds.
|
// assert on collectors registered on a private registry instead of
|
||||||
|
// the process-wide ones every other test is also moving.
|
||||||
func (e *Engine) ExportSetMetrics(mtr *metrics.Set) {
|
func (e *Engine) ExportSetMetrics(mtr *metrics.Set) {
|
||||||
e.mtr = mtr
|
e.mtr = mtr
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -35,8 +35,9 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
// mIsolate gives the setup's engine a metric set registered on a
|
// mIsolate gives the setup's engine a metric set registered on a
|
||||||
// registry this test holds, so its exact assertions can gather from
|
// private registry. The process-wide collectors are moved by every
|
||||||
// it.
|
// other delivery test running in parallel, so exact assertions are
|
||||||
|
// only possible against a registry this test owns.
|
||||||
func mIsolate(
|
func mIsolate(
|
||||||
t *testing.T, s iSetup,
|
t *testing.T, s iSetup,
|
||||||
) *prometheus.Registry {
|
) *prometheus.Registry {
|
||||||
@@ -411,10 +412,9 @@ func TestDeliveryMetrics_BreakerBlockedIsNotAnAttempt(
|
|||||||
|
|
||||||
s.Engine.ExportProcessRetryTask(context.TODO(), &blocked)
|
s.Engine.ExportProcessRetryTask(context.TODO(), &blocked)
|
||||||
|
|
||||||
// The breaker refused it: rescheduled without rewriting the
|
// The breaker refused it: rescheduled, so the retry counter
|
||||||
// retrying status it already had, so the retry counter did not
|
// moved, but nothing was attempted or timed.
|
||||||
// move, and nothing was attempted or timed.
|
assert.InDelta(t, retriesBefore+1,
|
||||||
assert.InDelta(t, retriesBefore,
|
|
||||||
mCounter(t, reg, mRetries, mTypeHTTP), 0)
|
mCounter(t, reg, mRetries, mTypeHTTP), 0)
|
||||||
assert.InDelta(t, threshold,
|
assert.InDelta(t, threshold,
|
||||||
mCounter(t, reg, mAttempts, mTypeHTTP), 0)
|
mCounter(t, reg, mAttempts, mTypeHTTP), 0)
|
||||||
|
|||||||
@@ -339,11 +339,10 @@ func TestRedirectPolicy_StopsAtHopCap(t *testing.T) {
|
|||||||
// The set the redirect policy strips is whatever the delivery path
|
// The set the redirect policy strips is whatever the delivery path
|
||||||
// actually put on the wire, so a header added to the forward set is
|
// actually put on the wire, so a header added to the forward set is
|
||||||
// covered without a second edit. A header the event never carried
|
// covered without a second edit. A header the event never carried
|
||||||
// is not in the set, and neither is the inbound Content-Type, because
|
// is not in the set, and the delivery path's own two are deliberately
|
||||||
// it is not forwarded. Two more are deliberately excluded: a
|
// excluded: Content-Type describes the body, which a 307 carries
|
||||||
// Content-Type configured on the target describes the body, which a
|
// across hosts, and the inbound User-Agent every real sender supplies
|
||||||
// 307 carries across hosts, and the inbound User-Agent every real
|
// is overwritten before the request goes out.
|
||||||
// sender supplies is overwritten before the request goes out.
|
|
||||||
func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
|
func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
@@ -372,7 +371,6 @@ func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
|
|||||||
&delivery.HTTPTargetConfig{
|
&delivery.HTTPTargetConfig{
|
||||||
Headers: map[string]string{
|
Headers: map[string]string{
|
||||||
probeHeaderName: probeHeaderValue,
|
probeHeaderName: probeHeaderValue,
|
||||||
"Content-Type": testContentType,
|
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
@@ -380,11 +378,7 @@ func TestApplyRequestHeaders_ReportsOriginScopedNames(t *testing.T) {
|
|||||||
assert.Equal(t,
|
assert.Equal(t,
|
||||||
[]string{probeHeaderName, inboundHeaderName}, names,
|
[]string{probeHeaderName, inboundHeaderName}, names,
|
||||||
"both header classes are reported, and only those: "+
|
"both header classes are reported, and only those: "+
|
||||||
"Host and the inbound Content-Type are never "+
|
"Host is never forwarded, Content-Type and "+
|
||||||
"forwarded, User-Agent is the delivery path's own",
|
"User-Agent are the delivery path's own",
|
||||||
)
|
|
||||||
assert.NotContains(t, names, "Content-Type",
|
|
||||||
"a Content-Type configured on the target must survive "+
|
|
||||||
"a cross-origin 307/308 with the body it describes",
|
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -26,7 +26,7 @@ var (
|
|||||||
"hostname resolved to no IP addresses",
|
"hostname resolved to no IP addresses",
|
||||||
)
|
)
|
||||||
errBlockedIP = errors.New(
|
errBlockedIP = errors.New(
|
||||||
"blocked private, reserved or cloud metadata address",
|
"blocked private/reserved IP range",
|
||||||
)
|
)
|
||||||
errBlockedMetadata = errors.New(
|
errBlockedMetadata = errors.New(
|
||||||
"blocked link-local or cloud instance metadata " +
|
"blocked link-local or cloud instance metadata " +
|
||||||
@@ -37,10 +37,9 @@ var (
|
|||||||
)
|
)
|
||||||
)
|
)
|
||||||
|
|
||||||
// blockedNetworks is the default blocklist: the private and
|
// blockedNetworks contains all private/reserved IP ranges
|
||||||
// reserved IP ranges, plus the public cloud metadata addresses,
|
// that should be blocked to prevent SSRF attacks. An operator
|
||||||
// that are blocked to prevent SSRF attacks. An operator can
|
// can permit specific blocks out of this set with
|
||||||
// permit specific blocks out of this set with
|
|
||||||
// ALLOWED_EGRESS_CIDRS; see Guard.
|
// ALLOWED_EGRESS_CIDRS; see Guard.
|
||||||
//
|
//
|
||||||
//nolint:gochecknoglobals // package-level network list is appropriate here
|
//nolint:gochecknoglobals // package-level network list is appropriate here
|
||||||
@@ -123,8 +122,6 @@ func init() {
|
|||||||
"::1/128",
|
"::1/128",
|
||||||
"fc00::/7",
|
"fc00::/7",
|
||||||
"fe80::/10",
|
"fe80::/10",
|
||||||
// Azure WireServer, a public address that serves VM credentials.
|
|
||||||
"168.63.129.16/32",
|
|
||||||
})
|
})
|
||||||
|
|
||||||
// Every entry is named. The set must not grow or shrink
|
// Every entry is named. The set must not grow or shrink
|
||||||
@@ -219,8 +216,8 @@ func matchesAny(networks []*net.IPNet, ip net.IP) bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// isBlockedIP checks whether an IP address falls within
|
// isBlockedIP checks whether an IP address falls within
|
||||||
// the default blocklist, before any operator allowlist is
|
// any blocked private/reserved network range, before any
|
||||||
// considered.
|
// operator allowlist is considered.
|
||||||
func isBlockedIP(ip net.IP) bool {
|
func isBlockedIP(ip net.IP) bool {
|
||||||
return matchesAny(blockedNetworks, ip)
|
return matchesAny(blockedNetworks, ip)
|
||||||
}
|
}
|
||||||
@@ -323,7 +320,7 @@ func (g *Guard) allows(ip net.IP) bool {
|
|||||||
//
|
//
|
||||||
// 1. alwaysBlockedNetworks is refused before the allowlist is
|
// 1. alwaysBlockedNetworks is refused before the allowlist is
|
||||||
// consulted, so no configured CIDR reaches link-local or a
|
// consulted, so no configured CIDR reaches link-local or a
|
||||||
// cloud metadata endpoint at a non-public address.
|
// cloud instance metadata endpoint.
|
||||||
// 2. The allowlist is consulted next, so a listed private
|
// 2. The allowlist is consulted next, so a listed private
|
||||||
// network becomes reachable.
|
// network becomes reachable.
|
||||||
// 3. Everything else keeps the default blocklist's answer.
|
// 3. Everything else keeps the default blocklist's answer.
|
||||||
|
|||||||
@@ -390,41 +390,6 @@ func TestGuardAllowlist_PublicUnaffected(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestGuardAllowlist_AzureWireServerReopenable covers Azure's
|
|
||||||
// WireServer, a public address that serves VM credentials. The
|
|
||||||
// default guard refuses it, but because it is public it sits in
|
|
||||||
// the default blocklist rather than the unconditional set, so an
|
|
||||||
// operator who lists it can reach it.
|
|
||||||
func TestGuardAllowlist_AzureWireServerReopenable(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
const wireServerIP = "168.63.129.16"
|
|
||||||
|
|
||||||
target := "http://" + wireServerIP + "/?comp=versions"
|
|
||||||
|
|
||||||
defaultGuard := delivery.NewTestGuard()
|
|
||||||
|
|
||||||
err := defaultGuard.ValidateTargetURL(context.Background(), target)
|
|
||||||
require.Error(t, err,
|
|
||||||
"WireServer must be refused with no allowlist set",
|
|
||||||
)
|
|
||||||
assert.NotContains(t, err.Error(), metadataRefusalClause,
|
|
||||||
"WireServer must be refused by the default blocklist, "+
|
|
||||||
"which an allowlist can override",
|
|
||||||
)
|
|
||||||
|
|
||||||
assertDialRefused(t, defaultGuard, target)
|
|
||||||
|
|
||||||
listed := delivery.NewTestGuard(
|
|
||||||
netip.MustParsePrefix(wireServerIP + "/32"),
|
|
||||||
)
|
|
||||||
|
|
||||||
assert.NoError(t,
|
|
||||||
listed.ValidateTargetURL(context.Background(), target),
|
|
||||||
"an operator who lists WireServer must be able to reach it",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestGuardCheckIP_BothPathsShareOneDecision asserts that the
|
// TestGuardCheckIP_BothPathsShareOneDecision asserts that the
|
||||||
// validator and the dialer are not two policies that happen to
|
// validator and the dialer are not two policies that happen to
|
||||||
// agree: both are defined in terms of checkIP, so the exported
|
// agree: both are defined in terms of checkIP, so the exported
|
||||||
|
|||||||
@@ -277,24 +277,6 @@ func (t *databaseTarget) evict(webhookID string) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// evictAll evicts every cached archive writer, exactly as evict
|
|
||||||
// does for one webhook. The engine calls it at shutdown, once its
|
|
||||||
// workers have returned. Closing the last handle on an archive
|
|
||||||
// moves the contents of its -wal into the .db and removes the
|
|
||||||
// -wal, so a clean stop leaves each archive as a single file.
|
|
||||||
func (t *databaseTarget) evictAll() {
|
|
||||||
t.mu.Lock()
|
|
||||||
|
|
||||||
writers := t.writers
|
|
||||||
t.writers = nil
|
|
||||||
|
|
||||||
t.mu.Unlock()
|
|
||||||
|
|
||||||
for _, w := range writers {
|
|
||||||
w.evict()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// sweepWebhook prunes one webhook's archive of rows older than
|
// sweepWebhook prunes one webhook's archive of rows older than
|
||||||
// expiry, without requiring a write. It returns nil (nothing to
|
// expiry, without requiring a write. It returns nil (nothing to
|
||||||
// do) when the archive file does not exist, so a sweep never
|
// do) when the archive file does not exist, so a sweep never
|
||||||
|
|||||||
@@ -1,7 +1,6 @@
|
|||||||
package delivery_test
|
package delivery_test
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"net/http"
|
"net/http"
|
||||||
@@ -362,46 +361,3 @@ func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) {
|
|||||||
"a later delivery should recreate the writer",
|
"a later delivery should recreate the writer",
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestEngineStop_WriteAfterStopIsRefused proves the engine's stop
|
|
||||||
// closes each archive writer the way deleting its webhook does: a
|
|
||||||
// write that reaches a writer after the stop is refused, reopens
|
|
||||||
// nothing and adds no row.
|
|
||||||
func TestEngineStop_WriteAfterStopIsRefused(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
eng, _ := evictTestEngine(t)
|
|
||||||
|
|
||||||
webhookDB := testWebhookDB(t)
|
|
||||||
event := seedEvent(t, webhookDB, `{"archived":true}`)
|
|
||||||
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
|
|
||||||
|
|
||||||
eng.ExportDeliverDatabase(webhookDB, d)
|
|
||||||
|
|
||||||
w := eng.ExportArchiveWriterFor(event.WebhookID)
|
|
||||||
require.NotNil(t, w)
|
|
||||||
require.True(t, w.HandleOpen())
|
|
||||||
|
|
||||||
require.NoError(t, eng.ExportStop(context.Background()))
|
|
||||||
|
|
||||||
err := w.Write(evictTestRow("ev-after-stop"), 0)
|
|
||||||
|
|
||||||
require.ErrorIs(
|
|
||||||
t, err, delivery.ErrExportArchiveWriterEvicted,
|
|
||||||
"a write after the stop must be refused",
|
|
||||||
)
|
|
||||||
assert.False(
|
|
||||||
t, w.HandleOpen(),
|
|
||||||
"a refused write must not reopen the archive",
|
|
||||||
)
|
|
||||||
assert.False(
|
|
||||||
t, eng.ExportHasArchiveWriter(event.WebhookID),
|
|
||||||
"the stop should empty the registry",
|
|
||||||
)
|
|
||||||
|
|
||||||
count, err := countArchivedRows(w.Path())
|
|
||||||
require.NoError(t, err)
|
|
||||||
assert.Equal(
|
|
||||||
t, int64(1), count, "the refused row must not be written",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -11,11 +11,10 @@ import (
|
|||||||
"sneak.berlin/go/webhooker/internal/delivery"
|
"sneak.berlin/go/webhooker/internal/delivery"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Literals these tests repeat, named so that the header names and the
|
// Literals these tests repeat, named so that the header name and the
|
||||||
// keep-forever archive config each have one definition.
|
// keep-forever archive config each have one definition.
|
||||||
const (
|
const (
|
||||||
headerAuthorization = "Authorization"
|
headerAuthorization = "Authorization"
|
||||||
headerContentType = "Content-Type"
|
|
||||||
bearerValue = "Bearer abc"
|
bearerValue = "Bearer abc"
|
||||||
archiveConfigNever = "{\"expiry\":\"never\"}"
|
archiveConfigNever = "{\"expiry\":\"never\"}"
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -197,14 +197,10 @@ func (c *httpCore) circuitBreakerBlock(
|
|||||||
"cooldown_remaining", remaining,
|
"cooldown_remaining", remaining,
|
||||||
)
|
)
|
||||||
|
|
||||||
// A delivery already at retrying is left as it is, so a task
|
c.eng.settleStatus(
|
||||||
// the breaker keeps turning away writes nothing each time.
|
webhookDB, d, d.Target.Type,
|
||||||
if d.Status != database.DeliveryStatusRetrying {
|
database.DeliveryStatusRetrying,
|
||||||
c.eng.settleStatus(
|
)
|
||||||
webhookDB, d, d.Target.Type,
|
|
||||||
database.DeliveryStatusRetrying,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
retryTask := *task
|
retryTask := *task
|
||||||
sched.ScheduleRetry(retryTask, remaining)
|
sched.ScheduleRetry(retryTask, remaining)
|
||||||
@@ -541,11 +537,6 @@ func isForwardableHeader(name string) bool {
|
|||||||
"Upgrade", "Proxy-Authorization",
|
"Upgrade", "Proxy-Authorization",
|
||||||
"Proxy-Connection", "Content-Length":
|
"Proxy-Connection", "Content-Length":
|
||||||
return false
|
return false
|
||||||
case "Content-Type":
|
|
||||||
// applyRequestHeaders sets Content-Type itself. The receiver
|
|
||||||
// already stored this inbound value as the event's
|
|
||||||
// ContentType, so forwarding it too would send it twice.
|
|
||||||
return false
|
|
||||||
default:
|
default:
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
@@ -558,10 +549,6 @@ func isForwardableHeader(name string) bool {
|
|||||||
// policy strips exactly that set on a hop that leaves the origin,
|
// policy strips exactly that set on a hop that leaves the origin,
|
||||||
// so the forward set is decided here and only here — a header added
|
// so the forward set is decided here and only here — a header added
|
||||||
// to it is covered off-origin without a second edit elsewhere.
|
// to it is covered off-origin without a second edit elsewhere.
|
||||||
//
|
|
||||||
// Content-Type goes out once: a Content-Type configured on the target
|
|
||||||
// wins, otherwise the event's ContentType, otherwise none. The inbound
|
|
||||||
// Content-Type in the event's headers is never forwarded.
|
|
||||||
func applyRequestHeaders(
|
func applyRequestHeaders(
|
||||||
req *http.Request,
|
req *http.Request,
|
||||||
event *database.Event,
|
event *database.Event,
|
||||||
@@ -582,10 +569,10 @@ func applyRequestHeaders(
|
|||||||
|
|
||||||
req.Header.Set("User-Agent", "webhooker/1.0")
|
req.Header.Set("User-Agent", "webhooker/1.0")
|
||||||
|
|
||||||
// A Content-Type configured on the target describes the body
|
// Content-Type describes the body being sent rather than the
|
||||||
// being sent rather than the sender. A 307/308 preserves the
|
// sender, and the delivery path sets it from the event itself.
|
||||||
// body across hosts, so stripping it would send that body
|
// A 307/308 preserves the body across hosts, so stripping it
|
||||||
// untyped.
|
// would send that body untyped.
|
||||||
delete(originScoped, "Content-Type")
|
delete(originScoped, "Content-Type")
|
||||||
|
|
||||||
// User-Agent is overwritten just above, so an inbound one never
|
// User-Agent is overwritten just above, so an inbound one never
|
||||||
|
|||||||
@@ -696,53 +696,6 @@ func TestSweepPending_TargetDeleted(t *testing.T) {
|
|||||||
assert.Contains(t, last.Error, "was deleted")
|
assert.Contains(t, last.Error, "was deleted")
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestSendRecoveredDeliveries_TargetMissingFromMap: the batch's target
|
|
||||||
// map is empty when its query failed, so every delivery in the batch is
|
|
||||||
// looked up on its own. A healthy one is sent to the target that lookup
|
|
||||||
// finds.
|
|
||||||
func TestSendRecoveredDeliveries_TargetMissingFromMap(
|
|
||||||
t *testing.T,
|
|
||||||
) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
targetID := uuid.New().String()
|
|
||||||
|
|
||||||
iCreateTarget(
|
|
||||||
t, s.MainDB, targetID, s.WebhookID, "found-on-lookup",
|
|
||||||
database.TargetTypeLog, "", 0,
|
|
||||||
)
|
|
||||||
|
|
||||||
event := iSeedEvent(
|
|
||||||
t, s.WebhookDB, s.WebhookID, `{"map":"empty"}`,
|
|
||||||
)
|
|
||||||
|
|
||||||
d := iSeedDelivery(
|
|
||||||
t, s.WebhookDB, event.ID, targetID,
|
|
||||||
database.DeliveryStatusPending,
|
|
||||||
)
|
|
||||||
|
|
||||||
s.Engine.ExportSendRecoveredDeliveries(
|
|
||||||
context.Background(), s.WebhookDB,
|
|
||||||
[]database.Delivery{d}, s.WebhookID,
|
|
||||||
map[string]database.Target{}, nil,
|
|
||||||
)
|
|
||||||
|
|
||||||
tasks := fDrain(s.Engine)
|
|
||||||
require.Len(t, tasks, 1,
|
|
||||||
"the healthy delivery was not queued exactly once",
|
|
||||||
)
|
|
||||||
assert.Equal(t, d.ID, tasks[0].DeliveryID)
|
|
||||||
assert.Equal(t, targetID, tasks[0].TargetID)
|
|
||||||
assert.Equal(t, database.TargetTypeLog, tasks[0].TargetType)
|
|
||||||
|
|
||||||
iAssertStatus(
|
|
||||||
t, s.WebhookDB, d.ID,
|
|
||||||
database.DeliveryStatusPending,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone: a failed
|
// TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone: a failed
|
||||||
// read of the main database is not a deleted target. Restart recovery
|
// read of the main database is not a deleted target. Restart recovery
|
||||||
// holds every pending delivery of the webhook in one batch, so failing
|
// holds every pending delivery of the webhook in one batch, so failing
|
||||||
|
|||||||
@@ -12,7 +12,6 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
|
|
||||||
"github.com/prometheus/client_golang/prometheus"
|
|
||||||
"go.uber.org/fx"
|
"go.uber.org/fx"
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
"sneak.berlin/go/webhooker/internal/delivery"
|
"sneak.berlin/go/webhooker/internal/delivery"
|
||||||
@@ -62,8 +61,6 @@ type HandlersParams struct {
|
|||||||
Notifier delivery.Notifier
|
Notifier delivery.Notifier
|
||||||
Evictor delivery.WebhookEvictor
|
Evictor delivery.WebhookEvictor
|
||||||
SSRFGuard *delivery.Guard
|
SSRFGuard *delivery.Guard
|
||||||
Metrics *metrics.Set
|
|
||||||
Registry *prometheus.Registry
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Handlers provides HTTP handler methods for all application
|
// Handlers provides HTTP handler methods for all application
|
||||||
@@ -125,7 +122,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.mtr = params.Metrics
|
s.mtr = metrics.Default()
|
||||||
s.ssrf = params.SSRFGuard
|
s.ssrf = params.SSRFGuard
|
||||||
|
|
||||||
// Parse all page templates once at startup
|
// Parse all page templates once at startup
|
||||||
|
|||||||
@@ -20,7 +20,6 @@ import (
|
|||||||
"sneak.berlin/go/webhooker/internal/handlers"
|
"sneak.berlin/go/webhooker/internal/handlers"
|
||||||
"sneak.berlin/go/webhooker/internal/healthcheck"
|
"sneak.berlin/go/webhooker/internal/healthcheck"
|
||||||
"sneak.berlin/go/webhooker/internal/logger"
|
"sneak.berlin/go/webhooker/internal/logger"
|
||||||
"sneak.berlin/go/webhooker/internal/metrics"
|
|
||||||
"sneak.berlin/go/webhooker/internal/middleware"
|
"sneak.berlin/go/webhooker/internal/middleware"
|
||||||
"sneak.berlin/go/webhooker/internal/session"
|
"sneak.berlin/go/webhooker/internal/session"
|
||||||
)
|
)
|
||||||
@@ -110,8 +109,6 @@ func newTestApp(
|
|||||||
func(r *recordingEvictor) delivery.WebhookEvictor {
|
func(r *recordingEvictor) delivery.WebhookEvictor {
|
||||||
return r
|
return r
|
||||||
},
|
},
|
||||||
metrics.NewRegistry,
|
|
||||||
metrics.New,
|
|
||||||
middleware.New,
|
middleware.New,
|
||||||
delivery.NewGuard,
|
delivery.NewGuard,
|
||||||
handlers.New,
|
handlers.New,
|
||||||
|
|||||||
@@ -1,20 +0,0 @@
|
|||||||
package handlers
|
|
||||||
|
|
||||||
import (
|
|
||||||
"net/http"
|
|
||||||
|
|
||||||
"github.com/prometheus/client_golang/prometheus/promhttp"
|
|
||||||
)
|
|
||||||
|
|
||||||
// HandleMetrics returns the Prometheus scrape handler for the
|
|
||||||
// registry every collector in this process registers on. It is what
|
|
||||||
// promhttp.Handler builds for the global default registry, including
|
|
||||||
// the promhttp_metric_handler_* series that count scrapes, pointed at
|
|
||||||
// that registry instead.
|
|
||||||
func (s *Handlers) HandleMetrics() http.HandlerFunc {
|
|
||||||
reg := s.params.Registry
|
|
||||||
|
|
||||||
return promhttp.InstrumentMetricHandler(
|
|
||||||
reg, promhttp.HandlerFor(reg, promhttp.HandlerOpts{}),
|
|
||||||
).ServeHTTP
|
|
||||||
}
|
|
||||||
+20
-26
@@ -3,17 +3,17 @@
|
|||||||
// deliveries are attempted, how they end, how long they take, how
|
// deliveries are attempted, how they end, how long they take, how
|
||||||
// deep the queues are, and how many circuit breakers are open.
|
// deep the queues are, and how many circuit breakers are open.
|
||||||
//
|
//
|
||||||
// It also builds the registry the authenticated /metrics route
|
// The inbound HTTP metrics come from the go-http-metrics recorder in
|
||||||
// serves. These collectors, the inbound HTTP metrics recorded in
|
// internal/middleware and land on prometheus.DefaultRegisterer. These
|
||||||
// internal/middleware, and the Go runtime and process collectors all
|
// collectors register there too, so both surfaces are gathered by the
|
||||||
// register on that one registry, never on Prometheus's global default.
|
// one promhttp handler mounted on the authenticated /metrics route.
|
||||||
package metrics
|
package metrics
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/prometheus/client_golang/prometheus"
|
"github.com/prometheus/client_golang/prometheus"
|
||||||
"github.com/prometheus/client_golang/prometheus/collectors"
|
|
||||||
"github.com/prometheus/client_golang/prometheus/promauto"
|
"github.com/prometheus/client_golang/prometheus/promauto"
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
)
|
)
|
||||||
@@ -57,31 +57,25 @@ var knownTargetTypes = []database.TargetType{
|
|||||||
database.TargetTypeSlack,
|
database.TargetTypeSlack,
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewRegistry returns the registry /metrics serves, carrying the Go
|
// defaultSet is the process-wide metric set, registered on the same
|
||||||
// runtime and process collectors that Prometheus's global default
|
// registry the HTTP middleware and the /metrics handler already use.
|
||||||
// registry carries, so the go_* and process_* series stay in the
|
// It is built on first use rather than in an init so that a test
|
||||||
// scrape.
|
// binary that never touches metrics never registers them.
|
||||||
//
|
//
|
||||||
// A registry of its own, rather than the global default, is what lets
|
//nolint:gochecknoglobals // one process-wide registration, by design
|
||||||
// two dependency graphs in one process — two tests, say — each
|
var defaultSet = sync.OnceValue(func() *Set {
|
||||||
// register their collectors without the second registration
|
return New(prometheus.DefaultRegisterer)
|
||||||
// panicking.
|
})
|
||||||
func NewRegistry() *prometheus.Registry {
|
|
||||||
reg := prometheus.NewRegistry()
|
|
||||||
reg.MustRegister(
|
|
||||||
collectors.NewGoCollector(),
|
|
||||||
collectors.NewProcessCollector(
|
|
||||||
collectors.ProcessCollectorOpts{},
|
|
||||||
),
|
|
||||||
)
|
|
||||||
|
|
||||||
return reg
|
// Default returns the process-wide metric set.
|
||||||
|
func Default() *Set {
|
||||||
|
return defaultSet()
|
||||||
}
|
}
|
||||||
|
|
||||||
// Set is one registered group of webhooker's delivery collectors.
|
// Set is one registered group of webhooker's delivery collectors.
|
||||||
// Production builds one on the registry /metrics serves; tests build
|
// Production uses the single Default set; tests build their own
|
||||||
// their own against a private registry so assertions are not
|
// against a private registry so assertions are not disturbed by
|
||||||
// disturbed by deliveries other tests are making concurrently.
|
// deliveries other tests are making concurrently.
|
||||||
type Set struct {
|
type Set struct {
|
||||||
eventsReceived prometheus.Counter
|
eventsReceived prometheus.Counter
|
||||||
deliveryAttempts *prometheus.CounterVec
|
deliveryAttempts *prometheus.CounterVec
|
||||||
@@ -99,7 +93,7 @@ type Set struct {
|
|||||||
// New registers a full set of delivery collectors on reg and returns
|
// New registers a full set of delivery collectors on reg and returns
|
||||||
// it. It panics if reg already holds them, which is the intended
|
// it. It panics if reg already holds them, which is the intended
|
||||||
// behaviour for a duplicate registration.
|
// behaviour for a duplicate registration.
|
||||||
func New(reg *prometheus.Registry) *Set {
|
func New(reg prometheus.Registerer) *Set {
|
||||||
factory := promauto.With(reg)
|
factory := promauto.With(reg)
|
||||||
|
|
||||||
s := &Set{
|
s := &Set{
|
||||||
|
|||||||
@@ -10,7 +10,8 @@ import (
|
|||||||
|
|
||||||
// MetricsMiddlewareForTest builds the metrics recording middleware
|
// MetricsMiddlewareForTest builds the metrics recording middleware
|
||||||
// against a caller-supplied recorder, so a test can gather from its
|
// against a caller-supplied recorder, so a test can gather from its
|
||||||
// own Prometheus registry without building a whole Middleware.
|
// own Prometheus registry rather than the process-wide default one
|
||||||
|
// that Middleware.Metrics uses.
|
||||||
func MetricsMiddlewareForTest(
|
func MetricsMiddlewareForTest(
|
||||||
rec httpmetrics.Recorder,
|
rec httpmetrics.Recorder,
|
||||||
) func(http.Handler) http.Handler {
|
) func(http.Handler) http.Handler {
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import (
|
|||||||
|
|
||||||
"github.com/go-chi/chi"
|
"github.com/go-chi/chi"
|
||||||
httpmetrics "github.com/slok/go-http-metrics/metrics"
|
httpmetrics "github.com/slok/go-http-metrics/metrics"
|
||||||
|
prommetrics "github.com/slok/go-http-metrics/metrics/prometheus"
|
||||||
ghmm "github.com/slok/go-http-metrics/middleware"
|
ghmm "github.com/slok/go-http-metrics/middleware"
|
||||||
"github.com/slok/go-http-metrics/middleware/std"
|
"github.com/slok/go-http-metrics/middleware/std"
|
||||||
)
|
)
|
||||||
@@ -151,14 +152,16 @@ func (r boundedLabelRecorder) AddInflightRequests(
|
|||||||
var _ httpmetrics.Recorder = boundedLabelRecorder{}
|
var _ httpmetrics.Recorder = boundedLabelRecorder{}
|
||||||
|
|
||||||
// Metrics returns middleware that records Prometheus HTTP metrics on
|
// Metrics returns middleware that records Prometheus HTTP metrics on
|
||||||
// the registry the /metrics route serves. Every call shares the one
|
// the default registry, which is the one the /metrics route gathers.
|
||||||
// recorder New built, so any number of routers can install it.
|
|
||||||
func (s *Middleware) Metrics() func(http.Handler) http.Handler {
|
func (s *Middleware) Metrics() func(http.Handler) http.Handler {
|
||||||
return metricsMiddleware(s.metricsRecorder)
|
return metricsMiddleware(
|
||||||
|
prommetrics.NewRecorder(prommetrics.Config{}),
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// metricsMiddleware builds the recording middleware against a given
|
// metricsMiddleware builds the recording middleware against a given
|
||||||
// recorder, so tests can gather from a registry of their own.
|
// recorder, so tests can gather from a registry of their own instead
|
||||||
|
// of the process-wide default.
|
||||||
func metricsMiddleware(
|
func metricsMiddleware(
|
||||||
rec httpmetrics.Recorder,
|
rec httpmetrics.Recorder,
|
||||||
) func(http.Handler) http.Handler {
|
) func(http.Handler) http.Handler {
|
||||||
|
|||||||
@@ -57,8 +57,9 @@ const (
|
|||||||
// Server.setupWebhookRoutes inside it. That ordering is the whole
|
// Server.setupWebhookRoutes inside it. That ordering is the whole
|
||||||
// defect, so a test that flattens it would prove nothing.
|
// defect, so a test that flattens it would prove nothing.
|
||||||
//
|
//
|
||||||
// The recorder writes to a registry of the test's own, so each test
|
// The recorder writes to a registry of the test's own rather than the
|
||||||
// observes only its own traffic.
|
// process-wide default one, so each test observes only its own
|
||||||
|
// traffic.
|
||||||
func metricsTestRouter(
|
func metricsTestRouter(
|
||||||
t *testing.T,
|
t *testing.T,
|
||||||
receiverLimit int,
|
receiverLimit int,
|
||||||
|
|||||||
@@ -13,9 +13,6 @@ import (
|
|||||||
"github.com/go-chi/chi"
|
"github.com/go-chi/chi"
|
||||||
"github.com/go-chi/chi/middleware"
|
"github.com/go-chi/chi/middleware"
|
||||||
"github.com/go-chi/cors"
|
"github.com/go-chi/cors"
|
||||||
"github.com/prometheus/client_golang/prometheus"
|
|
||||||
httpmetrics "github.com/slok/go-http-metrics/metrics"
|
|
||||||
prommetrics "github.com/slok/go-http-metrics/metrics/prometheus"
|
|
||||||
"go.uber.org/fx"
|
"go.uber.org/fx"
|
||||||
"sneak.berlin/go/webhooker/internal/config"
|
"sneak.berlin/go/webhooker/internal/config"
|
||||||
"sneak.berlin/go/webhooker/internal/globals"
|
"sneak.berlin/go/webhooker/internal/globals"
|
||||||
@@ -151,11 +148,10 @@ const (
|
|||||||
type MiddlewareParams struct {
|
type MiddlewareParams struct {
|
||||||
fx.In
|
fx.In
|
||||||
|
|
||||||
Logger *logger.Logger
|
Logger *logger.Logger
|
||||||
Globals *globals.Globals
|
Globals *globals.Globals
|
||||||
Config *config.Config
|
Config *config.Config
|
||||||
Session *session.Session
|
Session *session.Session
|
||||||
Registry *prometheus.Registry
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Middleware provides HTTP middleware for logging, CORS, auth, and
|
// Middleware provides HTTP middleware for logging, CORS, auth, and
|
||||||
@@ -165,12 +161,6 @@ type Middleware struct {
|
|||||||
params *MiddlewareParams
|
params *MiddlewareParams
|
||||||
session *session.Session
|
session *session.Session
|
||||||
|
|
||||||
// metricsRecorder records the inbound HTTP metrics on the
|
|
||||||
// registry /metrics serves. It is built once, in New, because
|
|
||||||
// building it registers its collectors, and a second
|
|
||||||
// registration on the same registry panics; see Metrics.
|
|
||||||
metricsRecorder httpmetrics.Recorder
|
|
||||||
|
|
||||||
// loginGuard counts failed credential verifications and bounds
|
// loginGuard counts failed credential verifications and bounds
|
||||||
// concurrent password hashing. It is built on first use so that
|
// concurrent password hashing. It is built on first use so that
|
||||||
// every construction path gets one; see guard().
|
// every construction path gets one; see guard().
|
||||||
@@ -189,9 +179,6 @@ func New(
|
|||||||
s.params = ¶ms
|
s.params = ¶ms
|
||||||
s.log = params.Logger.Get()
|
s.log = params.Logger.Get()
|
||||||
s.session = params.Session
|
s.session = params.Session
|
||||||
s.metricsRecorder = prommetrics.NewRecorder(
|
|
||||||
prommetrics.Config{Registry: params.Registry},
|
|
||||||
)
|
|
||||||
|
|
||||||
return s, nil
|
return s, nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -133,11 +133,6 @@ func (w *recoverResponseWriter) Unwrap() http.ResponseWriter {
|
|||||||
// what the access log records and the metrics count, and outside the
|
// what the access log records and the metrics count, and outside the
|
||||||
// sentryhttp handler, whose Repanic option depends on something
|
// sentryhttp handler, whose Repanic option depends on something
|
||||||
// further out recovering what it re-raises.
|
// further out recovering what it re-raises.
|
||||||
//
|
|
||||||
// Unlike http.Error on its own, it deletes any Set-Cookie the handler
|
|
||||||
// set before panicking, because a request that failed must not hand
|
|
||||||
// the client a credential; every other header is left to http.Error.
|
|
||||||
// See https://git.eeqj.de/sneak/webhooker/issues/193.
|
|
||||||
func (s *Middleware) Recoverer() func(http.Handler) http.Handler {
|
func (s *Middleware) Recoverer() func(http.Handler) http.Handler {
|
||||||
return func(next http.Handler) http.Handler {
|
return func(next http.Handler) http.Handler {
|
||||||
return http.HandlerFunc(func(
|
return http.HandlerFunc(func(
|
||||||
@@ -169,8 +164,6 @@ func (s *Middleware) Recoverer() func(http.Handler) http.Handler {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
rw.Header().Del("Set-Cookie")
|
|
||||||
|
|
||||||
http.Error(
|
http.Error(
|
||||||
rw,
|
rw,
|
||||||
http.StatusText(
|
http.StatusText(
|
||||||
|
|||||||
@@ -304,44 +304,16 @@ func TestRecovererRepanicsErrAbortHandler(t *testing.T) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestRecovererDropsSetCookieFromTheRecovered500 covers a handler that
|
|
||||||
// sets a cookie and a redirect target and then panics before sending
|
|
||||||
// anything. A request that failed must not hand the client a
|
|
||||||
// credential, so the 500 carries no cookie; Location is left alone.
|
|
||||||
func TestRecovererDropsSetCookieFromTheRecovered500(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
probe := newRecovererProbe(
|
|
||||||
t, false,
|
|
||||||
func(w http.ResponseWriter, _ *http.Request) {
|
|
||||||
w.Header().Set("Set-Cookie", "session=x")
|
|
||||||
w.Header().Set("Location", "/after")
|
|
||||||
|
|
||||||
panic(panicMarker)
|
|
||||||
},
|
|
||||||
)
|
|
||||||
|
|
||||||
resp, err := probe.get(t)
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NoError(t, resp.Body.Close())
|
|
||||||
|
|
||||||
assert.Equal(t, http.StatusInternalServerError, resp.StatusCode)
|
|
||||||
assert.Empty(t, resp.Cookies())
|
|
||||||
assert.Equal(t, "/after", resp.Header.Get("Location"))
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestRecovererKeepsAnAlreadyCommittedResponse covers a handler that
|
// TestRecovererKeepsAnAlreadyCommittedResponse covers a handler that
|
||||||
// panics after sending its status. The bytes are already on the wire,
|
// panics after sending its status. The bytes are already on the wire,
|
||||||
// cookie included, so a second WriteHeader would change nothing the
|
// so a second WriteHeader would change nothing the client sees and
|
||||||
// client sees and would draw net/http's "superfluous
|
// would draw net/http's "superfluous response.WriteHeader" report.
|
||||||
// response.WriteHeader" report.
|
|
||||||
func TestRecovererKeepsAnAlreadyCommittedResponse(t *testing.T) {
|
func TestRecovererKeepsAnAlreadyCommittedResponse(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
probe := newRecovererProbe(
|
probe := newRecovererProbe(
|
||||||
t, false,
|
t, false,
|
||||||
func(w http.ResponseWriter, _ *http.Request) {
|
func(w http.ResponseWriter, _ *http.Request) {
|
||||||
w.Header().Set("Set-Cookie", "session=x")
|
|
||||||
w.WriteHeader(committedStatus)
|
w.WriteHeader(committedStatus)
|
||||||
_, _ = w.Write([]byte("partial"))
|
_, _ = w.Write([]byte("partial"))
|
||||||
|
|
||||||
@@ -359,7 +331,6 @@ func TestRecovererKeepsAnAlreadyCommittedResponse(t *testing.T) {
|
|||||||
|
|
||||||
assert.Equal(t, committedStatus, resp.StatusCode)
|
assert.Equal(t, committedStatus, resp.StatusCode)
|
||||||
assert.Equal(t, "partial", string(body))
|
assert.Equal(t, "partial", string(body))
|
||||||
assert.Len(t, resp.Cookies(), 1)
|
|
||||||
|
|
||||||
record := probe.panicRecord(t)
|
record := probe.panicRecord(t)
|
||||||
assert.Equal(t, panicMarker, record["panic"])
|
assert.Equal(t, panicMarker, record["panic"])
|
||||||
|
|||||||
@@ -24,7 +24,6 @@ import (
|
|||||||
"sneak.berlin/go/webhooker/internal/handlers"
|
"sneak.berlin/go/webhooker/internal/handlers"
|
||||||
"sneak.berlin/go/webhooker/internal/healthcheck"
|
"sneak.berlin/go/webhooker/internal/healthcheck"
|
||||||
"sneak.berlin/go/webhooker/internal/logger"
|
"sneak.berlin/go/webhooker/internal/logger"
|
||||||
"sneak.berlin/go/webhooker/internal/metrics"
|
|
||||||
"sneak.berlin/go/webhooker/internal/middleware"
|
"sneak.berlin/go/webhooker/internal/middleware"
|
||||||
"sneak.berlin/go/webhooker/internal/resetpw"
|
"sneak.berlin/go/webhooker/internal/resetpw"
|
||||||
"sneak.berlin/go/webhooker/internal/session"
|
"sneak.berlin/go/webhooker/internal/session"
|
||||||
@@ -164,8 +163,6 @@ func newServerApp(
|
|||||||
session.New,
|
session.New,
|
||||||
func() delivery.Notifier { return &noopNotifier{} },
|
func() delivery.Notifier { return &noopNotifier{} },
|
||||||
func() delivery.WebhookEvictor { return &noopEvictor{} },
|
func() delivery.WebhookEvictor { return &noopEvictor{} },
|
||||||
metrics.NewRegistry,
|
|
||||||
metrics.New,
|
|
||||||
middleware.New,
|
middleware.New,
|
||||||
delivery.NewGuard,
|
delivery.NewGuard,
|
||||||
handlers.New,
|
handlers.New,
|
||||||
|
|||||||
+10
-18
@@ -7,6 +7,7 @@ import (
|
|||||||
sentryhttp "github.com/getsentry/sentry-go/http"
|
sentryhttp "github.com/getsentry/sentry-go/http"
|
||||||
"github.com/go-chi/chi"
|
"github.com/go-chi/chi"
|
||||||
"github.com/go-chi/chi/middleware"
|
"github.com/go-chi/chi/middleware"
|
||||||
|
"github.com/prometheus/client_golang/prometheus/promhttp"
|
||||||
"sneak.berlin/go/webhooker/static"
|
"sneak.berlin/go/webhooker/static"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -91,25 +92,11 @@ func (s *Server) setupGlobalMiddleware() {
|
|||||||
func (s *Server) setupRoutes() {
|
func (s *Server) setupRoutes() {
|
||||||
s.router.Get("/", s.h.HandleIndex())
|
s.router.Get("/", s.h.HandleIndex())
|
||||||
|
|
||||||
// Static assets answer GET and HEAD only. chi's default 405
|
s.router.Mount(
|
||||||
// carries no Allow header, so this group supplies its own.
|
"/s",
|
||||||
staticFiles := http.StripPrefix(
|
http.StripPrefix("/s", http.FileServer(http.FS(static.Static))),
|
||||||
"/s", http.FileServer(http.FS(static.Static)),
|
|
||||||
)
|
)
|
||||||
|
|
||||||
s.router.Route("/s", func(r chi.Router) {
|
|
||||||
r.MethodNotAllowed(func(w http.ResponseWriter, _ *http.Request) {
|
|
||||||
w.Header().Set("Allow", "GET, HEAD")
|
|
||||||
http.Error(
|
|
||||||
w,
|
|
||||||
"Method Not Allowed",
|
|
||||||
http.StatusMethodNotAllowed,
|
|
||||||
)
|
|
||||||
})
|
|
||||||
r.Method(http.MethodGet, "/*", staticFiles)
|
|
||||||
r.Method(http.MethodHead, "/*", staticFiles)
|
|
||||||
})
|
|
||||||
|
|
||||||
s.router.Route("/api/v1", func(_ chi.Router) {
|
s.router.Route("/api/v1", func(_ chi.Router) {
|
||||||
// API routes will be added here.
|
// API routes will be added here.
|
||||||
})
|
})
|
||||||
@@ -129,7 +116,12 @@ func (s *Server) setupRoutes() {
|
|||||||
if s.params.Config.MetricsAuthEnabled() {
|
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("/metrics", s.h.HandleMetrics())
|
r.Get(
|
||||||
|
"/metrics",
|
||||||
|
http.HandlerFunc(
|
||||||
|
promhttp.Handler().ServeHTTP,
|
||||||
|
),
|
||||||
|
)
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -23,7 +23,6 @@ import (
|
|||||||
"sneak.berlin/go/webhooker/internal/handlers"
|
"sneak.berlin/go/webhooker/internal/handlers"
|
||||||
"sneak.berlin/go/webhooker/internal/healthcheck"
|
"sneak.berlin/go/webhooker/internal/healthcheck"
|
||||||
"sneak.berlin/go/webhooker/internal/logger"
|
"sneak.berlin/go/webhooker/internal/logger"
|
||||||
"sneak.berlin/go/webhooker/internal/metrics"
|
|
||||||
"sneak.berlin/go/webhooker/internal/middleware"
|
"sneak.berlin/go/webhooker/internal/middleware"
|
||||||
"sneak.berlin/go/webhooker/internal/server"
|
"sneak.berlin/go/webhooker/internal/server"
|
||||||
"sneak.berlin/go/webhooker/internal/session"
|
"sneak.berlin/go/webhooker/internal/session"
|
||||||
@@ -113,8 +112,6 @@ func newTestEnvWithConfig(
|
|||||||
session.New,
|
session.New,
|
||||||
func() delivery.Notifier { return &noopNotifier{} },
|
func() delivery.Notifier { return &noopNotifier{} },
|
||||||
func() delivery.WebhookEvictor { return &noopEvictor{} },
|
func() delivery.WebhookEvictor { return &noopEvictor{} },
|
||||||
metrics.NewRegistry,
|
|
||||||
metrics.New,
|
|
||||||
middleware.New,
|
middleware.New,
|
||||||
delivery.NewGuard,
|
delivery.NewGuard,
|
||||||
handlers.New,
|
handlers.New,
|
||||||
@@ -399,15 +396,13 @@ func (e *testEnv) storedHash(t *testing.T, username string) string {
|
|||||||
|
|
||||||
// --- /s static group ---
|
// --- /s static group ---
|
||||||
|
|
||||||
// TestStaticServesOnlyGetAndHead pins the methods the static group
|
// TestStaticServesEveryMethod pins what the static mount actually
|
||||||
// answers: GET and HEAD are served the asset, and the other methods
|
// answers. chi's Mount registers the handler for all methods and
|
||||||
// chi routes (POST, PUT, DELETE and the rest) are refused with 405
|
// http.FileServer only special-cases HEAD (by suppressing the body),
|
||||||
// and an Allow header naming those two. A method chi does not route,
|
// so a POST or a DELETE to an asset is served the file rather than
|
||||||
// such as PROPFIND, is refused with 405 by the top-level router
|
// refused. The README documents this; the test is what keeps the two
|
||||||
// before it reaches the static group, so it gets no Allow header.
|
// from drifting.
|
||||||
// The README documents this; the test is what keeps the two from
|
func TestStaticServesEveryMethod(t *testing.T) {
|
||||||
// drifting.
|
|
||||||
func TestStaticServesOnlyGetAndHead(t *testing.T) {
|
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
env := newTestEnv(t)
|
env := newTestEnv(t)
|
||||||
@@ -422,7 +417,6 @@ func TestStaticServesOnlyGetAndHead(t *testing.T) {
|
|||||||
http.MethodPost,
|
http.MethodPost,
|
||||||
http.MethodPut,
|
http.MethodPut,
|
||||||
http.MethodDelete,
|
http.MethodDelete,
|
||||||
"PROPFIND",
|
|
||||||
} {
|
} {
|
||||||
t.Run(method, func(t *testing.T) {
|
t.Run(method, func(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
@@ -434,38 +428,18 @@ func TestStaticServesOnlyGetAndHead(t *testing.T) {
|
|||||||
w := httptest.NewRecorder()
|
w := httptest.NewRecorder()
|
||||||
env.router.ServeHTTP(w, req)
|
env.router.ServeHTTP(w, req)
|
||||||
|
|
||||||
switch method {
|
assert.Equal(t, http.StatusOK, w.Code,
|
||||||
case http.MethodGet:
|
"static mount answers every method")
|
||||||
assert.Equal(t, http.StatusOK, w.Code)
|
|
||||||
assert.Equal(t, body, w.Body.Bytes(),
|
if method == http.MethodHead {
|
||||||
"the asset itself is returned")
|
|
||||||
case http.MethodHead:
|
|
||||||
assert.Equal(t, http.StatusOK, w.Code)
|
|
||||||
assert.Empty(t, w.Body.Bytes(),
|
assert.Empty(t, w.Body.Bytes(),
|
||||||
"HEAD must not carry a body")
|
"HEAD must not carry a body")
|
||||||
case "PROPFIND":
|
|
||||||
assert.Equal(
|
return
|
||||||
t, http.StatusMethodNotAllowed, w.Code,
|
|
||||||
)
|
|
||||||
assert.Empty(t, w.Header().Get("Allow"),
|
|
||||||
"chi refuses a method it does not route "+
|
|
||||||
"before the static group runs")
|
|
||||||
assert.NotContains(
|
|
||||||
t, w.Body.String(), string(body),
|
|
||||||
"a refused method must not get the asset",
|
|
||||||
)
|
|
||||||
default:
|
|
||||||
assert.Equal(
|
|
||||||
t, http.StatusMethodNotAllowed, w.Code,
|
|
||||||
)
|
|
||||||
assert.Equal(
|
|
||||||
t, "GET, HEAD", w.Header().Get("Allow"),
|
|
||||||
)
|
|
||||||
assert.NotContains(
|
|
||||||
t, w.Body.String(), string(body),
|
|
||||||
"a refused method must not get the asset",
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
assert.Equal(t, body, w.Body.Bytes(),
|
||||||
|
"the asset itself is returned")
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -964,43 +938,3 @@ func TestMetricsRouteUnmountedOnHalfSetConfig(t *testing.T) {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestTwoMetricsRoutersInOneProcess pins
|
|
||||||
// https://git.eeqj.de/sneak/webhooker/issues/227: a second
|
|
||||||
// metrics-enabled router in one process used to panic, because the
|
|
||||||
// HTTP metrics registered on Prometheus's global default registry.
|
|
||||||
// Two routers are built over separate dependency graphs and a third
|
|
||||||
// over the first graph again, and each must still serve the HTTP,
|
|
||||||
// delivery and Go runtime series.
|
|
||||||
func TestTwoMetricsRoutersInOneProcess(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
first := newTestEnvWithConfig(
|
|
||||||
t, metricsConfig(t, metricsUser, metricsAuthValue),
|
|
||||||
)
|
|
||||||
second := newTestEnvWithConfig(
|
|
||||||
t, metricsConfig(t, metricsUser, metricsAuthValue),
|
|
||||||
)
|
|
||||||
third := &testEnv{
|
|
||||||
router: server.NewRouterForTest(
|
|
||||||
first.log.Get(), first.cfg, first.mw, first.hnd,
|
|
||||||
),
|
|
||||||
}
|
|
||||||
|
|
||||||
for _, env := range []*testEnv{first, second, third} {
|
|
||||||
env.get("/", nil)
|
|
||||||
|
|
||||||
scrape := env.metricsRequest(metricsUser, metricsAuthValue)
|
|
||||||
require.Equal(t, http.StatusOK, scrape.Code)
|
|
||||||
|
|
||||||
for _, series := range []string{
|
|
||||||
"http_request_duration_seconds",
|
|
||||||
"http_response_size_bytes",
|
|
||||||
"http_requests_inflight",
|
|
||||||
"webhooker_events_received_total",
|
|
||||||
"go_goroutines",
|
|
||||||
} {
|
|
||||||
assert.Contains(t, scrape.Body.String(), series)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
Reference in New Issue
Block a user