Compare commits
10
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
51186b347b | ||
|
|
8ad2a86e4b | ||
|
|
a891b726e5 | ||
|
|
ab63b5f777 | ||
|
|
9cf9cdd8eb | ||
|
|
f0adeafde3 | ||
|
|
f755c03110 | ||
|
|
4a724130ca | ||
|
|
d4f4ddf51f | ||
|
|
51580a2bc6 |
@@ -157,6 +157,11 @@ 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
|
||||||
@@ -195,15 +200,16 @@ 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 that
|
- **It cannot open link-local, or a cloud metadata endpoint at a
|
||||||
discloses credentials or user data.** An address is on the list below
|
non-public address that discloses credentials or user data.** An
|
||||||
when both of these hold: the provider fixes it, so it cannot collide
|
address is on the list below when it is not a public address and both
|
||||||
with anything you run; and reaching it hands out credentials, user
|
of these hold: the provider fixes it, so it cannot collide with
|
||||||
data or bootstrap material. Those stay blocked no matter what you
|
anything you run; and reaching it hands out credentials, user data or
|
||||||
list, including when you list them outright or list a supernet such
|
bootstrap material. Those stay blocked no matter what you list,
|
||||||
as `0.0.0.0/0`, `::/0`, `fd00::/8` or `100.64.0.0/10`. Treat this as
|
including when you list them outright or list a supernet such as
|
||||||
best effort rather than a guarantee — it is a hand-maintained list
|
`0.0.0.0/0`, `::/0`, `fd00::/8` or `100.64.0.0/10`. Treat this as best
|
||||||
and the caveat below the table applies:
|
effort rather than a guarantee — it is a hand-maintained list and the
|
||||||
|
caveat below the table applies:
|
||||||
|
|
||||||
| Blocked unconditionally | What it is |
|
| Blocked unconditionally | What it is |
|
||||||
| ----------------------- | ---------- |
|
| ----------------------- | ---------- |
|
||||||
@@ -242,7 +248,8 @@ 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.
|
escape hatch at all; Azure's `168.63.129.16` is refused by the default
|
||||||
|
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.
|
||||||
@@ -968,15 +975,10 @@ 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 `webhooker.db` and every `events-*.db`, which
|
A clean shutdown closes every database, which checkpoints and removes
|
||||||
checkpoints and removes their sidecars; a killed or crashed instance
|
its sidecars; a killed or crashed instance leaves them, and they must be
|
||||||
leaves them, and they must be carried with the `.db`. **Archive
|
carried with the `.db`. An archive the service has not opened since a
|
||||||
databases are different**: their handle is not closed at shutdown, so
|
crash keeps that crash's sidecars, even across a later clean stop.
|
||||||
`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
|
||||||
@@ -1051,10 +1053,9 @@ 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.
|
20 KB `.db` with no sidecars about a minute after its last write. A
|
||||||
Shutdown is **not** on that list: the archive handle is not closed when
|
clean stop closes it too. So either move `archive-{uuid}.db` together
|
||||||
the service stops. So either move `archive-{uuid}.db` together with any
|
with any `-wal`/`-shm` beside it, or wait until there are none.
|
||||||
`-wal`/`-shm` beside it, or wait until there are none.
|
|
||||||
|
|
||||||
### Restore
|
### Restore
|
||||||
|
|
||||||
@@ -1073,12 +1074,10 @@ the service stops. So either move `archive-{uuid}.db` together with any
|
|||||||
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 has none for `webhooker.db` or the `events-*.db`,
|
stop-and-copy set normally has none, because a clean stop closes
|
||||||
because a clean stop closes those and checkpoints their sidecars
|
every database and checkpoints its sidecars away; the exception is an
|
||||||
away — but it will normally have them for `archive-*.db`, whose
|
archive not opened since a crash. A copy salvaged from a crashed
|
||||||
handle stays open across shutdown, and those carry the archive's
|
instance has them for everything, and needs all of them.
|
||||||
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
|
||||||
@@ -2722,7 +2721,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`) |
|
||||||
| 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` |
|
| `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` |
|
||||||
| `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
|
||||||
@@ -3088,7 +3087,8 @@ 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`
|
4. `delivery.Engine` — waits for its workers, then closes the archive
|
||||||
|
databases
|
||||||
5. `healthcheck`
|
5. `healthcheck`
|
||||||
6. `WebhookDBManager`
|
6. `WebhookDBManager`
|
||||||
7. the database close
|
7. the database close
|
||||||
@@ -3280,3 +3280,5 @@ MIT
|
|||||||
## Author
|
## Author
|
||||||
|
|
||||||
[@sneak](https://sneak.berlin)
|
[@sneak](https://sneak.berlin)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ go 1.26.1
|
|||||||
|
|
||||||
require (
|
require (
|
||||||
github.com/99designs/basicauth-go v0.0.0-20230316000542-bf6f9cbbf0f8
|
github.com/99designs/basicauth-go v0.0.0-20230316000542-bf6f9cbbf0f8
|
||||||
|
github.com/dustin/go-humanize v1.0.1
|
||||||
github.com/getsentry/sentry-go v0.25.0
|
github.com/getsentry/sentry-go v0.25.0
|
||||||
github.com/go-chi/chi v1.5.5
|
github.com/go-chi/chi v1.5.5
|
||||||
github.com/go-chi/cors v1.2.1
|
github.com/go-chi/cors v1.2.1
|
||||||
@@ -29,7 +30,6 @@ require (
|
|||||||
github.com/beorn7/perks v1.0.1 // indirect
|
github.com/beorn7/perks v1.0.1 // indirect
|
||||||
github.com/cespare/xxhash/v2 v2.2.0 // indirect
|
github.com/cespare/xxhash/v2 v2.2.0 // indirect
|
||||||
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect
|
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect
|
||||||
github.com/dustin/go-humanize v1.0.1 // indirect
|
|
||||||
github.com/gorilla/securecookie v1.1.2 // indirect
|
github.com/gorilla/securecookie v1.1.2 // indirect
|
||||||
github.com/jinzhu/inflection v1.0.0 // indirect
|
github.com/jinzhu/inflection v1.0.0 // indirect
|
||||||
github.com/jinzhu/now v1.1.5 // indirect
|
github.com/jinzhu/now v1.1.5 // indirect
|
||||||
|
|||||||
@@ -192,9 +192,10 @@ 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 address; it is not exhaustive of every
|
// at a provider-fixed, non-public address; it is not
|
||||||
// cloud's metadata address. See alwaysBlockedNetworks for the
|
// exhaustive of every cloud's metadata address. See
|
||||||
// authoritative list and the criterion it is built from.
|
// alwaysBlockedNetworks for the authoritative list and the
|
||||||
|
// criterion it is built from.
|
||||||
AllowedEgressCIDRs []netip.Prefix
|
AllowedEgressCIDRs []netip.Prefix
|
||||||
|
|
||||||
params *ConfigParams
|
params *ConfigParams
|
||||||
@@ -746,12 +747,14 @@ 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 private/reserved networks. Anyone "+
|
"otherwise-blocked networks. Anyone who can create a "+
|
||||||
"who can create a delivery target can now make this "+
|
"delivery target can now make this process issue "+
|
||||||
"process issue requests into them, and read back the "+
|
"requests into them, and read back the response. Only "+
|
||||||
"response. Link-local and the known cloud instance "+
|
"the addresses the README lists as blocked "+
|
||||||
"metadata endpoints outside it stay blocked "+
|
"unconditionally stay blocked regardless of what is "+
|
||||||
"regardless of what is listed here.",
|
"listed here; a public cloud metadata address such as "+
|
||||||
|
"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,12 +834,13 @@ 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. Asserted on the clause naming the
|
// What stays shut is the whole unconditional set, not
|
||||||
// wider set rather than on "Link-local" alone, so the
|
// link-local alone; a public metadata address is not in
|
||||||
// string cannot narrow back to link-local only while
|
// it, so a listed block covering it opens it.
|
||||||
// the always-blocked set covers ULA, CGNAT and two
|
assert.Contains(t, logged, "blocked unconditionally")
|
||||||
// public metadata addresses as well.
|
assert.Contains(t, logged, "168.63.129.16 is reachable")
|
||||||
assert.Contains(t, logged, "metadata endpoints outside it")
|
// The listed blocks need not be private or reserved.
|
||||||
|
assert.NotContains(t, logged, "private/reserved")
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+70
-23
@@ -362,6 +362,15 @@ 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")
|
||||||
|
|
||||||
@@ -376,6 +385,8 @@ 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
|
||||||
@@ -728,9 +739,7 @@ func (e *Engine) recoverSingleRetry(
|
|||||||
// webhook on one bad read would be a far larger fault than
|
// webhook on one bad read would be a far larger fault than
|
||||||
// the strand it is meant to clear.
|
// the strand it is meant to clear.
|
||||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||||
e.failMissingTargetRetry(
|
e.failMissingTarget(webhookDB, webhookID, d)
|
||||||
webhookDB, webhookID, d,
|
|
||||||
)
|
|
||||||
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -1133,9 +1142,7 @@ func (e *Engine) sweepSingleRetry(
|
|||||||
// Deleted is terminal, unreadable is not; see
|
// Deleted is terminal, unreadable is not; see
|
||||||
// recoverSingleRetry.
|
// recoverSingleRetry.
|
||||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||||
e.failMissingTargetRetry(
|
e.failMissingTarget(webhookDB, webhookID, d)
|
||||||
webhookDB, webhookID, d,
|
|
||||||
)
|
|
||||||
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -1249,19 +1256,19 @@ func (e *Engine) failUnretryableRetry(
|
|||||||
e.failDelivery(webhookDB, d, target.Type, reason)
|
e.failDelivery(webhookDB, d, target.Type, reason)
|
||||||
}
|
}
|
||||||
|
|
||||||
// failMissingTargetRetry terminally fails an orphaned retrying
|
// failMissingTarget terminally fails a recovered delivery, pending or
|
||||||
// delivery whose target row is gone. Both restart recovery and the
|
// retrying, whose target row is gone. Restart recovery and the periodic
|
||||||
// periodic sweep call it, so the transition exists once.
|
// sweep call it for both statuses, so the transition exists once.
|
||||||
//
|
//
|
||||||
// Until it existed both paths logged the failed lookup and returned,
|
// Until it existed those paths logged the failed lookup and moved on,
|
||||||
// which left the delivery retrying for the life of the database and
|
// which left the delivery where it was for the life of the database and
|
||||||
// the sweep repeating the same error every minute forever. Failing it
|
// the sweep repeating the same error every minute forever. Failing it
|
||||||
// with a recorded reason is the treatment the other orphaned-retry
|
// with a recorded reason is the treatment the other orphaned-retry
|
||||||
// cases already get, so all of them read alike in the event log.
|
// cases already get, so all of them read alike in the event log.
|
||||||
//
|
//
|
||||||
// Logged at warn rather than error: a deleted target is an operator
|
// Logged at warn rather than error: a deleted target is an operator
|
||||||
// action, not a system fault.
|
// action, not a system fault.
|
||||||
func (e *Engine) failMissingTargetRetry(
|
func (e *Engine) failMissingTarget(
|
||||||
webhookDB *gorm.DB,
|
webhookDB *gorm.DB,
|
||||||
webhookID string,
|
webhookID string,
|
||||||
d *database.Delivery,
|
d *database.Delivery,
|
||||||
@@ -1274,13 +1281,37 @@ func (e *Engine) failMissingTargetRetry(
|
|||||||
|
|
||||||
defer e.inflight.release(d.ID)
|
defer e.inflight.release(d.ID)
|
||||||
|
|
||||||
|
// The batch was read before ownership was taken, and a worker may
|
||||||
|
// have settled the delivery and let it go in between. Only a row
|
||||||
|
// still in the status the batch read is failed.
|
||||||
|
row, err := e.loadDelivery(webhookDB, d.ID)
|
||||||
|
if err != nil {
|
||||||
|
e.log.Error(
|
||||||
|
"failed to load delivery",
|
||||||
|
"delivery_id", d.ID,
|
||||||
|
"error", err,
|
||||||
|
)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
if row.Status != d.Status {
|
||||||
|
e.log.Debug(
|
||||||
|
"delivery already handled, not failed",
|
||||||
|
"delivery_id", d.ID,
|
||||||
|
"status", row.Status,
|
||||||
|
)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
targetType, reason := e.missingTargetReason(d.TargetID)
|
targetType, reason := e.missingTargetReason(d.TargetID)
|
||||||
|
|
||||||
e.log.Warn(
|
e.log.Warn(
|
||||||
"failing orphaned retrying delivery: "+
|
"failing recovered delivery: its target no longer exists",
|
||||||
"its target no longer exists",
|
|
||||||
"webhook_id", webhookID,
|
"webhook_id", webhookID,
|
||||||
"delivery_id", d.ID,
|
"delivery_id", d.ID,
|
||||||
|
"status", d.Status,
|
||||||
"target_id", d.TargetID,
|
"target_id", d.TargetID,
|
||||||
"target_type", targetType,
|
"target_type", targetType,
|
||||||
)
|
)
|
||||||
@@ -1314,15 +1345,14 @@ func (e *Engine) missingTargetReason(
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return "", fmt.Sprintf(
|
return "", fmt.Sprintf(
|
||||||
"target %s no longer exists; the delivery "+
|
"target %s no longer exists; the delivery "+
|
||||||
"cannot be retried and has been failed "+
|
"has been failed terminally",
|
||||||
"terminally",
|
|
||||||
targetID,
|
targetID,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
return target.Type, fmt.Sprintf(
|
return target.Type, fmt.Sprintf(
|
||||||
"target %q (type %s) was deleted; the delivery "+
|
"target %q (type %s) was deleted; the delivery "+
|
||||||
"cannot be retried and has been failed terminally",
|
"has been failed terminally",
|
||||||
target.Name, target.Type,
|
target.Name, target.Type,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
@@ -2021,13 +2051,30 @@ func (e *Engine) sendRecoveredDeliveries(
|
|||||||
|
|
||||||
target, ok := targetMap[deliveries[i].TargetID]
|
target, ok := targetMap[deliveries[i].TargetID]
|
||||||
if !ok {
|
if !ok {
|
||||||
e.log.Error(
|
// A missing entry does not mean the target is gone: the
|
||||||
"target not found for delivery",
|
// map is also empty when its query failed. Only a lookup
|
||||||
"delivery_id", deliveries[i].ID,
|
// that finds no row ends the delivery; any other error
|
||||||
"target_id", deliveries[i].TargetID,
|
// leaves it pending for the next sweep. See
|
||||||
)
|
// recoverSingleRetry.
|
||||||
|
var err error
|
||||||
|
|
||||||
continue
|
target, err = e.loadTarget(deliveries[i].TargetID)
|
||||||
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||||
|
e.failMissingTarget(webhookDB, webhookID, &deliveries[i])
|
||||||
|
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
e.log.Error(
|
||||||
|
"failed to load target for recovered delivery",
|
||||||
|
"delivery_id", deliveries[i].ID,
|
||||||
|
"target_id", deliveries[i].TargetID,
|
||||||
|
"error", err,
|
||||||
|
)
|
||||||
|
|
||||||
|
continue
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if !e.takeForRedispatch(
|
if !e.takeForRedispatch(
|
||||||
|
|||||||
@@ -2,6 +2,8 @@ package delivery_test
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"path/filepath"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -269,3 +271,88 @@ 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",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|||||||
@@ -342,6 +342,31 @@ func (e *Engine) ExportRecoverRetryingDeliveries(
|
|||||||
e.recoverRetryingDeliveries(webhookDB, webhookID)
|
e.recoverRetryingDeliveries(webhookDB, webhookID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ExportFailMissingTarget exposes failMissingTarget, so a test can hand
|
||||||
|
// it a delivery as a batch read it earlier.
|
||||||
|
func (e *Engine) ExportFailMissingTarget(
|
||||||
|
webhookDB *gorm.DB,
|
||||||
|
webhookID string,
|
||||||
|
d *database.Delivery,
|
||||||
|
) {
|
||||||
|
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
|
||||||
|
|||||||
@@ -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 IP range",
|
"blocked private, reserved or cloud metadata address",
|
||||||
)
|
)
|
||||||
errBlockedMetadata = errors.New(
|
errBlockedMetadata = errors.New(
|
||||||
"blocked link-local or cloud instance metadata " +
|
"blocked link-local or cloud instance metadata " +
|
||||||
@@ -37,9 +37,10 @@ var (
|
|||||||
)
|
)
|
||||||
)
|
)
|
||||||
|
|
||||||
// blockedNetworks contains all private/reserved IP ranges
|
// blockedNetworks is the default blocklist: the private and
|
||||||
// that should be blocked to prevent SSRF attacks. An operator
|
// reserved IP ranges, plus the public cloud metadata addresses,
|
||||||
// can permit specific blocks out of this set with
|
// that are blocked to prevent SSRF attacks. An operator can
|
||||||
|
// 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
|
||||||
@@ -122,6 +123,8 @@ 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
|
||||||
@@ -216,8 +219,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
|
||||||
// any blocked private/reserved network range, before any
|
// the default blocklist, before any operator allowlist is
|
||||||
// operator allowlist is considered.
|
// considered.
|
||||||
func isBlockedIP(ip net.IP) bool {
|
func isBlockedIP(ip net.IP) bool {
|
||||||
return matchesAny(blockedNetworks, ip)
|
return matchesAny(blockedNetworks, ip)
|
||||||
}
|
}
|
||||||
@@ -320,7 +323,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 instance metadata endpoint.
|
// cloud metadata endpoint at a non-public address.
|
||||||
// 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,6 +390,41 @@ 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,6 +277,24 @@ 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,6 +1,7 @@
|
|||||||
package delivery_test
|
package delivery_test
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"net/http"
|
"net/http"
|
||||||
@@ -361,3 +362,46 @@ 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",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|||||||
@@ -18,16 +18,17 @@ import (
|
|||||||
// https://git.eeqj.de/sneak/webhooker/issues/107: a delivery failed
|
// https://git.eeqj.de/sneak/webhooker/issues/107: a delivery failed
|
||||||
// with nothing in its event log to say why, and a retrying delivery
|
// with nothing in its event log to say why, and a retrying delivery
|
||||||
// whose target was deleted, which used to keep sending and then never
|
// whose target was deleted, which used to keep sending and then never
|
||||||
// terminalise.
|
// terminalise. Section 4 is the same deleted-target gap for a pending
|
||||||
|
// delivery: https://git.eeqj.de/sneak/webhooker/issues/293.
|
||||||
|
|
||||||
// tUnknownType is a target type no build implements. It stands in for
|
// tUnknownType is a target type no build implements. It stands in for
|
||||||
// a target whose type was written by a build that knew a type this one
|
// a target whose type was written by a build that knew a type this one
|
||||||
// does not.
|
// does not.
|
||||||
const tUnknownType = database.TargetType("pubsub")
|
const tUnknownType = database.TargetType("pubsub")
|
||||||
|
|
||||||
// tSeedDeletedTarget creates a target, a retrying delivery against it
|
// tSeedDeletedTarget creates a target, a delivery against it at the
|
||||||
// with one recorded failed attempt, and then deletes the target the
|
// given status with one recorded failed attempt, and then deletes the
|
||||||
// way the source page does.
|
// target the way the source page does.
|
||||||
//
|
//
|
||||||
// It asserts the delete is soft, because that is the whole reason the
|
// It asserts the delete is soft, because that is the whole reason the
|
||||||
// engine could not tell a deleted target from a target id that never
|
// engine could not tell a deleted target from a target id that never
|
||||||
@@ -36,6 +37,7 @@ func tSeedDeletedTarget(
|
|||||||
t *testing.T,
|
t *testing.T,
|
||||||
s iSetup,
|
s iSetup,
|
||||||
name, url string,
|
name, url string,
|
||||||
|
status database.DeliveryStatus,
|
||||||
) string {
|
) string {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
@@ -51,8 +53,7 @@ func tSeedDeletedTarget(
|
|||||||
)
|
)
|
||||||
|
|
||||||
d := iSeedDelivery(
|
d := iSeedDelivery(
|
||||||
t, s.WebhookDB, event.ID, targetID,
|
t, s.WebhookDB, event.ID, targetID, status,
|
||||||
database.DeliveryStatusRetrying,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
iSeedFailedResult(t, s.WebhookDB, d.ID)
|
iSeedFailedResult(t, s.WebhookDB, d.ID)
|
||||||
@@ -173,6 +174,7 @@ func TestRecoverSingleRetry_TargetDeleted(t *testing.T) {
|
|||||||
|
|
||||||
deliveryID := tSeedDeletedTarget(
|
deliveryID := tSeedDeletedTarget(
|
||||||
t, s, "gone-on-recovery", "http://example.com/hook",
|
t, s, "gone-on-recovery", "http://example.com/hook",
|
||||||
|
database.DeliveryStatusRetrying,
|
||||||
)
|
)
|
||||||
|
|
||||||
s.Engine.ExportRecoverWebhookDeliveries(
|
s.Engine.ExportRecoverWebhookDeliveries(
|
||||||
@@ -210,6 +212,7 @@ func TestSweepSingleRetry_TargetDeleted(t *testing.T) {
|
|||||||
|
|
||||||
deliveryID := tSeedDeletedTarget(
|
deliveryID := tSeedDeletedTarget(
|
||||||
t, s, "gone-on-sweep", "http://example.com/hook",
|
t, s, "gone-on-sweep", "http://example.com/hook",
|
||||||
|
database.DeliveryStatusRetrying,
|
||||||
)
|
)
|
||||||
|
|
||||||
// Twice, because the bug was an error the sweep repeated every
|
// Twice, because the bug was an error the sweep repeated every
|
||||||
@@ -278,11 +281,11 @@ func TestSweepSingleRetry_TargetNeverExisted(t *testing.T) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestFailMissingTargetRetry_WritesNoTargetRow holds the new terminal
|
// TestFailMissingTarget_WritesNoTargetRow holds the new terminal path
|
||||||
// path to the same rule as the existing one: no target row, and so no
|
// to the same rule as the existing one: no target row, and so no
|
||||||
// plaintext target config, may be written into the per-webhook event
|
// plaintext target config, may be written into the per-webhook event
|
||||||
// database. See https://git.eeqj.de/sneak/webhooker/issues/206.
|
// database. See https://git.eeqj.de/sneak/webhooker/issues/206.
|
||||||
func TestFailMissingTargetRetry_WritesNoTargetRow(
|
func TestFailMissingTarget_WritesNoTargetRow(
|
||||||
t *testing.T,
|
t *testing.T,
|
||||||
) {
|
) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
@@ -297,6 +300,7 @@ func TestFailMissingTargetRetry_WritesNoTargetRow(
|
|||||||
|
|
||||||
deliveryID := tSeedDeletedTarget(
|
deliveryID := tSeedDeletedTarget(
|
||||||
t, s, "credential-bearing", hookURL,
|
t, s, "credential-bearing", hookURL,
|
||||||
|
database.DeliveryStatusRetrying,
|
||||||
)
|
)
|
||||||
|
|
||||||
s.Engine.ExportSweepWebhookRetries(
|
s.Engine.ExportSweepWebhookRetries(
|
||||||
@@ -529,3 +533,260 @@ func TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone(
|
|||||||
|
|
||||||
assert.Zero(t, s.Engine.ExportInflightHeld())
|
assert.Zero(t, s.Engine.ExportInflightHeld())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// --- 4. A pending delivery whose target is gone ---
|
||||||
|
|
||||||
|
func TestRecoverPending_TargetDeleted(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
s := newISetup(t)
|
||||||
|
|
||||||
|
deliveryID := tSeedDeletedTarget(
|
||||||
|
t, s, "gone-while-pending", "http://example.com/hook",
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
|
||||||
|
s.Engine.ExportRecoverWebhookDeliveries(
|
||||||
|
context.Background(), s.WebhookID,
|
||||||
|
)
|
||||||
|
|
||||||
|
iAssertStatus(
|
||||||
|
t, s.WebhookDB, deliveryID,
|
||||||
|
database.DeliveryStatusFailed,
|
||||||
|
)
|
||||||
|
|
||||||
|
last := tLastResult(t, s, deliveryID, 2)
|
||||||
|
|
||||||
|
assert.False(t, last.Success)
|
||||||
|
assert.Equal(t, 2, last.AttemptNum)
|
||||||
|
assert.Contains(t, last.Error, "gone-while-pending")
|
||||||
|
assert.Contains(t, last.Error, "was deleted")
|
||||||
|
|
||||||
|
assert.Empty(t, fDrain(s.Engine),
|
||||||
|
"a delivery whose target is gone was sent",
|
||||||
|
)
|
||||||
|
assert.Zero(t, s.Engine.ExportInflightHeld(),
|
||||||
|
"the terminal path leaked its ownership reference",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRecoverPending_TargetDeleted_LeavesAnOwnedDeliveryAlone: the
|
||||||
|
// terminal write takes ownership like every other recovery write, so a
|
||||||
|
// delivery the engine still holds is not failed underneath its worker.
|
||||||
|
func TestRecoverPending_TargetDeleted_LeavesAnOwnedDeliveryAlone(
|
||||||
|
t *testing.T,
|
||||||
|
) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
s := newISetup(t)
|
||||||
|
|
||||||
|
deliveryID := tSeedDeletedTarget(
|
||||||
|
t, s, "gone-but-owned", "http://example.com/hook",
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
|
||||||
|
require.True(t, s.Engine.ExportRetainDelivery(deliveryID))
|
||||||
|
|
||||||
|
s.Engine.ExportRecoverWebhookDeliveries(
|
||||||
|
context.Background(), s.WebhookID,
|
||||||
|
)
|
||||||
|
|
||||||
|
iAssertStatus(
|
||||||
|
t, s.WebhookDB, deliveryID,
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.Len(t, iResults(t, s.WebhookDB, deliveryID), 1,
|
||||||
|
"a delivery the engine owns was failed underneath it",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestFailMissingTarget_LeavesASettledDeliveryAlone: the recovery paths
|
||||||
|
// read their batch before taking ownership, and a worker may send a
|
||||||
|
// delivery and let it go in between. The terminal write goes by the row
|
||||||
|
// as it is now, not as the batch read it.
|
||||||
|
func TestFailMissingTarget_LeavesASettledDeliveryAlone(
|
||||||
|
t *testing.T,
|
||||||
|
) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
s := newISetup(t)
|
||||||
|
|
||||||
|
deliveryID := tSeedDeletedTarget(
|
||||||
|
t, s, "gone-after-sending", "http://example.com/hook",
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
|
||||||
|
var batch database.Delivery
|
||||||
|
|
||||||
|
require.NoError(t, s.WebhookDB.First(
|
||||||
|
&batch, "id = ?", deliveryID,
|
||||||
|
).Error)
|
||||||
|
|
||||||
|
// A worker settles the delivery after the batch was read.
|
||||||
|
require.NoError(t, s.WebhookDB.Model(&database.Delivery{}).
|
||||||
|
Where("id = ?", deliveryID).
|
||||||
|
Update("status", database.DeliveryStatusDelivered).Error)
|
||||||
|
|
||||||
|
s.Engine.ExportFailMissingTarget(
|
||||||
|
s.WebhookDB, s.WebhookID, &batch,
|
||||||
|
)
|
||||||
|
|
||||||
|
iAssertStatus(
|
||||||
|
t, s.WebhookDB, deliveryID,
|
||||||
|
database.DeliveryStatusDelivered,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.Len(t, iResults(t, s.WebhookDB, deliveryID), 1,
|
||||||
|
"a delivery settled after the batch read was then failed",
|
||||||
|
)
|
||||||
|
assert.Zero(t, s.Engine.ExportInflightHeld())
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestSweepPending_TargetDeleted sweeps twice over a batch that also
|
||||||
|
// holds a healthy stranded delivery. The one whose target is gone is
|
||||||
|
// failed once and then left alone; the healthy one is queued by the
|
||||||
|
// first sweep and not again by the second.
|
||||||
|
func TestSweepPending_TargetDeleted(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
liveTargetID := uuid.New().String()
|
||||||
|
s := fSweepSetup(t, liveTargetID, "still-there")
|
||||||
|
|
||||||
|
deliveryID := tSeedDeletedTarget(
|
||||||
|
t, s, "gone-on-pending-sweep", "http://example.com/hook",
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
rAgePending(t, s.WebhookDB, deliveryID)
|
||||||
|
|
||||||
|
event := iSeedEvent(
|
||||||
|
t, s.WebhookDB, s.WebhookID, `{"target":"live"}`,
|
||||||
|
)
|
||||||
|
|
||||||
|
healthy := iSeedDelivery(
|
||||||
|
t, s.WebhookDB, event.ID, liveTargetID,
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
rAgePending(t, s.WebhookDB, healthy.ID)
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
|
||||||
|
|
||||||
|
tasks := fDrain(s.Engine)
|
||||||
|
require.Len(t, tasks, 1,
|
||||||
|
"the first sweep did not queue the healthy delivery",
|
||||||
|
)
|
||||||
|
assert.Equal(t, healthy.ID, tasks[0].DeliveryID)
|
||||||
|
|
||||||
|
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
|
||||||
|
|
||||||
|
assert.Empty(t, fDrain(s.Engine),
|
||||||
|
"the second sweep queued a delivery again",
|
||||||
|
)
|
||||||
|
|
||||||
|
iAssertStatus(
|
||||||
|
t, s.WebhookDB, deliveryID,
|
||||||
|
database.DeliveryStatusFailed,
|
||||||
|
)
|
||||||
|
|
||||||
|
last := tLastResult(t, s, deliveryID, 2)
|
||||||
|
|
||||||
|
assert.Contains(t, last.Error, "gone-on-pending-sweep")
|
||||||
|
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
|
||||||
|
// read of the main database is not a deleted target. Restart recovery
|
||||||
|
// holds every pending delivery of the webhook in one batch, so failing
|
||||||
|
// on this would fail all of them.
|
||||||
|
func TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone(
|
||||||
|
t *testing.T,
|
||||||
|
) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
s := newISetup(t)
|
||||||
|
|
||||||
|
targetID := uuid.New().String()
|
||||||
|
|
||||||
|
iCreateTarget(
|
||||||
|
t, s.MainDB, targetID, s.WebhookID, "healthy",
|
||||||
|
database.TargetTypeLog, "", 0,
|
||||||
|
)
|
||||||
|
|
||||||
|
event := iSeedEvent(
|
||||||
|
t, s.WebhookDB, s.WebhookID, `{"still":"pending"}`,
|
||||||
|
)
|
||||||
|
|
||||||
|
d := iSeedDelivery(
|
||||||
|
t, s.WebhookDB, event.ID, targetID,
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
|
||||||
|
sqlDB, err := s.MainDB.DB()
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.NoError(t, sqlDB.Close())
|
||||||
|
|
||||||
|
s.Engine.ExportRecoverPendingDeliveries(
|
||||||
|
context.Background(), s.WebhookDB, s.WebhookID,
|
||||||
|
)
|
||||||
|
|
||||||
|
iAssertStatus(
|
||||||
|
t, s.WebhookDB, d.ID,
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.Empty(t, iResults(t, s.WebhookDB, d.ID),
|
||||||
|
"an unreadable main database produced a terminal "+
|
||||||
|
"failure row",
|
||||||
|
)
|
||||||
|
assert.Empty(t, fDrain(s.Engine))
|
||||||
|
assert.Zero(t, s.Engine.ExportInflightHeld())
|
||||||
|
}
|
||||||
|
|||||||
@@ -1,6 +1,8 @@
|
|||||||
package handlers
|
package handlers
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"time"
|
||||||
|
|
||||||
"sneak.berlin/go/webhooker/internal/delivery"
|
"sneak.berlin/go/webhooker/internal/delivery"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -25,7 +27,7 @@ const maxRenderedResponseBytes = 4096
|
|||||||
// cut, so an oversized stored response never becomes a Go
|
// cut, so an oversized stored response never becomes a Go
|
||||||
// string at all.
|
// string at all.
|
||||||
const deliveryResultColumns = "delivery_id, attempt_num, success, " +
|
const deliveryResultColumns = "delivery_id, attempt_num, success, " +
|
||||||
"status_code, error, duration, " +
|
"status_code, error, duration, created_at, " +
|
||||||
"substr(cast(response_body as blob), 1, ?) AS response_body, " +
|
"substr(cast(response_body as blob), 1, ?) AS response_body, " +
|
||||||
"length(cast(response_body as blob)) AS response_bytes"
|
"length(cast(response_body as blob)) AS response_bytes"
|
||||||
|
|
||||||
@@ -106,6 +108,10 @@ type deliveryResultRow struct {
|
|||||||
Duration int64
|
Duration int64
|
||||||
ResponseBody []byte
|
ResponseBody []byte
|
||||||
ResponseBytes int64
|
ResponseBytes int64
|
||||||
|
|
||||||
|
// CreatedAt is when the attempt's result was recorded, which
|
||||||
|
// is when the attempt finished.
|
||||||
|
CreatedAt time.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
// view projects a loaded row for rendering, stripping the
|
// view projects a loaded row for rendering, stripping the
|
||||||
|
|||||||
@@ -28,7 +28,7 @@ const (
|
|||||||
// maxBodyShift is the bit shift for 1 MB body limit.
|
// maxBodyShift is the bit shift for 1 MB body limit.
|
||||||
maxBodyShift = 20
|
maxBodyShift = 20
|
||||||
// recentEventLimit is the number of recent events to show.
|
// recentEventLimit is the number of recent events to show.
|
||||||
recentEventLimit = 20
|
recentEventLimit = 50
|
||||||
// paginationPerPage is the number of items per page.
|
// paginationPerPage is the number of items per page.
|
||||||
paginationPerPage = 25
|
paginationPerPage = 25
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,248 @@
|
|||||||
|
package handlers
|
||||||
|
|
||||||
|
import (
|
||||||
|
"net/http"
|
||||||
|
"strconv"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/dustin/go-humanize"
|
||||||
|
"gorm.io/gorm"
|
||||||
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
|
)
|
||||||
|
|
||||||
|
// recentEventColumns is the recent events list's projection. It
|
||||||
|
// reads the body's size and never the body itself, for the reason
|
||||||
|
// maxRenderedBodyBytes gives; the cast to blob makes length count
|
||||||
|
// bytes rather than characters.
|
||||||
|
const recentEventColumns = "id, created_at, method, content_type, " +
|
||||||
|
"resubmitted_from_id, length(cast(body as blob)) AS body_bytes"
|
||||||
|
|
||||||
|
// RecentEventView is one row of the recent events list on a
|
||||||
|
// webhook's page.
|
||||||
|
type RecentEventView struct {
|
||||||
|
Method string
|
||||||
|
ContentType string
|
||||||
|
|
||||||
|
// ResubmittedFromID names the event this one was copied from,
|
||||||
|
// empty for an event that arrived on the receiver.
|
||||||
|
ResubmittedFromID string
|
||||||
|
|
||||||
|
// Received is how long ago the event arrived, and ReceivedUTC
|
||||||
|
// the full timestamp the page shows on hover.
|
||||||
|
Received string
|
||||||
|
ReceivedUTC string
|
||||||
|
|
||||||
|
// Size is the size of the stored body.
|
||||||
|
Size string
|
||||||
|
|
||||||
|
// ProcessingTime is how long the event's slowest delivery
|
||||||
|
// took; see processingTime.
|
||||||
|
ProcessingTime string
|
||||||
|
|
||||||
|
// Status is what the webhook's HTTP target answered, and
|
||||||
|
// StatusClass its colour; see targetStatus. Both are empty
|
||||||
|
// unless the webhook has exactly one HTTP target.
|
||||||
|
Status string
|
||||||
|
StatusClass string
|
||||||
|
}
|
||||||
|
|
||||||
|
// recentEventRow is one row of recentEventColumns.
|
||||||
|
type recentEventRow struct {
|
||||||
|
ID string
|
||||||
|
CreatedAt time.Time
|
||||||
|
Method string
|
||||||
|
ContentType string
|
||||||
|
ResubmittedFromID *string
|
||||||
|
BodyBytes uint64
|
||||||
|
}
|
||||||
|
|
||||||
|
// singleHTTPTargetID returns the ID of the webhook's HTTP target
|
||||||
|
// when it has exactly one, and "" when it has none or several.
|
||||||
|
func singleHTTPTargetID(targets []database.Target) string {
|
||||||
|
id := ""
|
||||||
|
count := 0
|
||||||
|
|
||||||
|
for i := range targets {
|
||||||
|
if targets[i].Type == database.TargetTypeHTTP {
|
||||||
|
id = targets[i].ID
|
||||||
|
count++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if count != 1 {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
|
return id
|
||||||
|
}
|
||||||
|
|
||||||
|
// loadRecentEvents loads the webhook's recentEventLimit newest
|
||||||
|
// events for its page, newest first. statusTargetID is the
|
||||||
|
// webhook's only HTTP target, or "" when the list shows no status.
|
||||||
|
func (h *Handlers) loadRecentEvents(
|
||||||
|
webhookDB *gorm.DB, webhookID, statusTargetID string,
|
||||||
|
) ([]RecentEventView, error) {
|
||||||
|
var rows []recentEventRow
|
||||||
|
|
||||||
|
err := webhookDB.Model(&database.Event{}).
|
||||||
|
Select(recentEventColumns).
|
||||||
|
Where("webhook_id = ?", webhookID).
|
||||||
|
Order("created_at DESC").
|
||||||
|
Limit(recentEventLimit).
|
||||||
|
Find(&rows).Error
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
eventIDs := make([]string, len(rows))
|
||||||
|
for i := range rows {
|
||||||
|
eventIDs[i] = rows[i].ID
|
||||||
|
}
|
||||||
|
|
||||||
|
// Oldest first, so an event's last delivery to a target is its
|
||||||
|
// newest: a replay adds a delivery rather than changing the
|
||||||
|
// earlier one.
|
||||||
|
var deliveries []database.Delivery
|
||||||
|
|
||||||
|
err = webhookDB.
|
||||||
|
Select("id, event_id, target_id, status, created_at").
|
||||||
|
Where("event_id IN ?", eventIDs).
|
||||||
|
Order("created_at ASC").
|
||||||
|
Find(&deliveries).Error
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
byEvent := make(map[string][]database.Delivery, len(rows))
|
||||||
|
deliveryIDs := make([]string, len(deliveries))
|
||||||
|
|
||||||
|
for i := range deliveries {
|
||||||
|
eventID := deliveries[i].EventID
|
||||||
|
byEvent[eventID] = append(byEvent[eventID], deliveries[i])
|
||||||
|
deliveryIDs[i] = deliveries[i].ID
|
||||||
|
}
|
||||||
|
|
||||||
|
attempts, err := h.loadDeliveryResults(webhookDB, deliveryIDs)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
views := make([]RecentEventView, len(rows))
|
||||||
|
for i := range rows {
|
||||||
|
views[i] = rows[i].view(
|
||||||
|
byEvent[rows[i].ID], attempts, statusTargetID,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
return views, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// view projects a loaded row for rendering. deliveries is the
|
||||||
|
// event's deliveries, oldest first, and attempts their recorded
|
||||||
|
// attempts keyed by delivery ID.
|
||||||
|
func (r *recentEventRow) view(
|
||||||
|
deliveries []database.Delivery,
|
||||||
|
attempts map[string][]deliveryResultRow,
|
||||||
|
statusTargetID string,
|
||||||
|
) RecentEventView {
|
||||||
|
v := RecentEventView{
|
||||||
|
Method: r.Method,
|
||||||
|
ContentType: r.ContentType,
|
||||||
|
Received: humanize.Time(r.CreatedAt),
|
||||||
|
ReceivedUTC: r.CreatedAt.UTC().Format(time.DateTime) + " UTC",
|
||||||
|
Size: humanize.Bytes(r.BodyBytes),
|
||||||
|
ProcessingTime: processingTime(deliveries, attempts),
|
||||||
|
}
|
||||||
|
|
||||||
|
if r.ResubmittedFromID != nil {
|
||||||
|
v.ResubmittedFromID = *r.ResubmittedFromID
|
||||||
|
}
|
||||||
|
|
||||||
|
if statusTargetID != "" {
|
||||||
|
v.Status, v.StatusClass = targetStatus(
|
||||||
|
deliveries, attempts, statusTargetID,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
return v
|
||||||
|
}
|
||||||
|
|
||||||
|
// processingTime is how long the event's slowest delivery took,
|
||||||
|
// from being queued to its last recorded attempt, time spent
|
||||||
|
// waiting between retries included. A delivery is queued when its
|
||||||
|
// event is received, or when an operator replays it, so a replay
|
||||||
|
// is timed from the replay rather than from the event's arrival.
|
||||||
|
// It is "in progress" while any delivery is pending or retrying,
|
||||||
|
// and empty for an event with no deliveries.
|
||||||
|
func processingTime(
|
||||||
|
deliveries []database.Delivery,
|
||||||
|
attempts map[string][]deliveryResultRow,
|
||||||
|
) string {
|
||||||
|
if len(deliveries) == 0 {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
|
var slowest time.Duration
|
||||||
|
|
||||||
|
for i := range deliveries {
|
||||||
|
if !deliveries[i].Status.Terminal() {
|
||||||
|
return "in progress"
|
||||||
|
}
|
||||||
|
|
||||||
|
tries := attempts[deliveries[i].ID]
|
||||||
|
if len(tries) == 0 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
last := tries[len(tries)-1].CreatedAt
|
||||||
|
slowest = max(slowest, last.Sub(deliveries[i].CreatedAt))
|
||||||
|
}
|
||||||
|
|
||||||
|
return slowest.Round(time.Millisecond).String()
|
||||||
|
}
|
||||||
|
|
||||||
|
// targetStatus is what the target answered for the event, and the
|
||||||
|
// colour to show it in: the HTTP status code of the last attempt of
|
||||||
|
// the event's newest delivery to the target. Without a code it is
|
||||||
|
// "no response" when that attempt failed before a response
|
||||||
|
// arrived, the delivery's status ("pending") before any attempt,
|
||||||
|
// and "not sent" when the event has no delivery to the target.
|
||||||
|
func targetStatus(
|
||||||
|
deliveries []database.Delivery,
|
||||||
|
attempts map[string][]deliveryResultRow,
|
||||||
|
targetID string,
|
||||||
|
) (string, string) {
|
||||||
|
newest := -1
|
||||||
|
|
||||||
|
for i := range deliveries {
|
||||||
|
if deliveries[i].TargetID == targetID {
|
||||||
|
newest = i
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if newest < 0 {
|
||||||
|
return "not sent", "text-gray-400"
|
||||||
|
}
|
||||||
|
|
||||||
|
tries := attempts[deliveries[newest].ID]
|
||||||
|
if len(tries) == 0 {
|
||||||
|
return string(deliveries[newest].Status), "text-gray-400"
|
||||||
|
}
|
||||||
|
|
||||||
|
code := tries[len(tries)-1].StatusCode
|
||||||
|
|
||||||
|
switch {
|
||||||
|
case code == 0:
|
||||||
|
return "no response", "text-red-600"
|
||||||
|
case code >= http.StatusInternalServerError:
|
||||||
|
return strconv.Itoa(code), "text-red-600"
|
||||||
|
case code >= http.StatusBadRequest:
|
||||||
|
return strconv.Itoa(code), "text-yellow-600"
|
||||||
|
case code >= http.StatusMultipleChoices:
|
||||||
|
return strconv.Itoa(code), "text-gray-500"
|
||||||
|
case code >= http.StatusOK:
|
||||||
|
return strconv.Itoa(code), "text-green-600"
|
||||||
|
default:
|
||||||
|
return strconv.Itoa(code), "text-gray-500"
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,295 @@
|
|||||||
|
package handlers_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"net/http"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
"gorm.io/gorm"
|
||||||
|
"gorm.io/gorm/clause"
|
||||||
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
|
"sneak.berlin/go/webhooker/internal/handlers"
|
||||||
|
"sneak.berlin/go/webhooker/internal/session"
|
||||||
|
)
|
||||||
|
|
||||||
|
// statusTitle marks the status column's cell in a recent events
|
||||||
|
// row; it is absent from the page when the column is not shown.
|
||||||
|
const statusTitle = `title="HTTP status from the HTTP target"`
|
||||||
|
|
||||||
|
// recentEventsFixture is one started app and a webhook whose
|
||||||
|
// recent events list a test fills.
|
||||||
|
type recentEventsFixture struct {
|
||||||
|
h *handlers.Handlers
|
||||||
|
sess *session.Session
|
||||||
|
db *database.Database
|
||||||
|
webhook *database.Webhook
|
||||||
|
webhookDB *gorm.DB
|
||||||
|
}
|
||||||
|
|
||||||
|
func newRecentEventsFixture(t *testing.T) *recentEventsFixture {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
f := &recentEventsFixture{}
|
||||||
|
|
||||||
|
var dbMgr *database.WebhookDBManager
|
||||||
|
|
||||||
|
app := newTestApp(t, &f.h, &f.sess, &f.db, &dbMgr)
|
||||||
|
app.RequireStart()
|
||||||
|
|
||||||
|
t.Cleanup(app.RequireStop)
|
||||||
|
|
||||||
|
f.webhook = seedWebhook(t, f.db)
|
||||||
|
|
||||||
|
webhookDB, err := dbMgr.GetDB(f.webhook.ID)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
f.webhookDB = webhookDB
|
||||||
|
|
||||||
|
return f
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *recentEventsFixture) render(t *testing.T) string {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
return renderSourceDetailPage(t, f.h, f.sess, f.webhook.ID)
|
||||||
|
}
|
||||||
|
|
||||||
|
// event records an event received at receivedAt.
|
||||||
|
func (f *recentEventsFixture) event(
|
||||||
|
t *testing.T, contentType, body string, receivedAt time.Time,
|
||||||
|
) *database.Event {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
event := &database.Event{
|
||||||
|
WebhookID: f.webhook.ID,
|
||||||
|
Method: http.MethodPost,
|
||||||
|
Body: body,
|
||||||
|
ContentType: contentType,
|
||||||
|
}
|
||||||
|
event.CreatedAt = receivedAt
|
||||||
|
|
||||||
|
require.NoError(t, f.webhookDB.Omit(
|
||||||
|
clause.Associations,
|
||||||
|
).Create(event).Error)
|
||||||
|
|
||||||
|
return event
|
||||||
|
}
|
||||||
|
|
||||||
|
// delivery records a delivery of the event to the target, queued
|
||||||
|
// when the event was received.
|
||||||
|
func (f *recentEventsFixture) delivery(
|
||||||
|
t *testing.T,
|
||||||
|
event *database.Event,
|
||||||
|
targetID string,
|
||||||
|
status database.DeliveryStatus,
|
||||||
|
) *database.Delivery {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
return f.deliveryQueuedAt(
|
||||||
|
t, event, targetID, status, event.CreatedAt,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// deliveryQueuedAt records a delivery of the event to the target,
|
||||||
|
// queued at queuedAt, as a replay is.
|
||||||
|
func (f *recentEventsFixture) deliveryQueuedAt(
|
||||||
|
t *testing.T,
|
||||||
|
event *database.Event,
|
||||||
|
targetID string,
|
||||||
|
status database.DeliveryStatus,
|
||||||
|
queuedAt time.Time,
|
||||||
|
) *database.Delivery {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
dlv := &database.Delivery{
|
||||||
|
EventID: event.ID,
|
||||||
|
TargetID: targetID,
|
||||||
|
Status: status,
|
||||||
|
}
|
||||||
|
dlv.CreatedAt = queuedAt
|
||||||
|
|
||||||
|
require.NoError(t, f.webhookDB.Omit(
|
||||||
|
clause.Associations,
|
||||||
|
).Create(dlv).Error)
|
||||||
|
|
||||||
|
return dlv
|
||||||
|
}
|
||||||
|
|
||||||
|
// attempt records one attempt of the delivery that finished took
|
||||||
|
// after the delivery was queued, with HTTP status code (0 for no
|
||||||
|
// response).
|
||||||
|
func (f *recentEventsFixture) attempt(
|
||||||
|
t *testing.T, dlv *database.Delivery, code int, took time.Duration,
|
||||||
|
) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
result := &database.DeliveryResult{
|
||||||
|
DeliveryID: dlv.ID,
|
||||||
|
AttemptNum: 1,
|
||||||
|
StatusCode: code,
|
||||||
|
}
|
||||||
|
result.CreatedAt = dlv.CreatedAt.Add(took)
|
||||||
|
|
||||||
|
require.NoError(t, f.webhookDB.Omit(
|
||||||
|
clause.Associations,
|
||||||
|
).Create(result).Error)
|
||||||
|
}
|
||||||
|
|
||||||
|
// statusCell is the status column's cell as the page renders it.
|
||||||
|
func statusCell(class, text string) string {
|
||||||
|
return `<span class="font-medium ` + class + `" ` + statusTitle +
|
||||||
|
`>` + text + `</span>`
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestHandleSourceDetail_ShowsFiftyNewestEvents proves the list
|
||||||
|
// holds the 50 newest events, newest first, and not one more.
|
||||||
|
func TestHandleSourceDetail_ShowsFiftyNewestEvents(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
f := newRecentEventsFixture(t)
|
||||||
|
base := time.Now().Add(-time.Hour)
|
||||||
|
|
||||||
|
for i := range 51 {
|
||||||
|
f.event(
|
||||||
|
t, fmt.Sprintf("application/x-recent-%02d", i), "{}",
|
||||||
|
base.Add(time.Duration(i)*time.Second),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
body := f.render(t)
|
||||||
|
|
||||||
|
assert.Equal(t, 50, strings.Count(body, `title="Body size"`))
|
||||||
|
assert.NotContains(t, body, "application/x-recent-00")
|
||||||
|
assert.Contains(t, body, "application/x-recent-01")
|
||||||
|
assert.Less(
|
||||||
|
t,
|
||||||
|
strings.Index(body, "application/x-recent-50"),
|
||||||
|
strings.Index(body, "application/x-recent-49"),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestHandleSourceDetail_RecentEventColumns proves a row shows its
|
||||||
|
// time relative with the UTC timestamp on hover, its body size,
|
||||||
|
// and its processing time once every delivery has finished.
|
||||||
|
func TestHandleSourceDetail_RecentEventColumns(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
f := newRecentEventsFixture(t)
|
||||||
|
logTarget := seedTarget(t, f.db, f.webhook.ID, database.TargetTypeLog)
|
||||||
|
|
||||||
|
receivedAt := time.Now().Add(-210 * time.Second).
|
||||||
|
UTC().Truncate(time.Second)
|
||||||
|
|
||||||
|
done := f.event(
|
||||||
|
t, contentTypeJSON, strings.Repeat("x", 2048), receivedAt,
|
||||||
|
)
|
||||||
|
f.attempt(
|
||||||
|
t,
|
||||||
|
f.delivery(t, done, logTarget.ID, database.DeliveryStatusDelivered),
|
||||||
|
0, 1500*time.Millisecond,
|
||||||
|
)
|
||||||
|
|
||||||
|
waiting := f.event(t, "text/plain", "{}", receivedAt)
|
||||||
|
f.delivery(t, waiting, logTarget.ID, database.DeliveryStatusPending)
|
||||||
|
|
||||||
|
body := f.render(t)
|
||||||
|
|
||||||
|
assert.Contains(
|
||||||
|
t, body,
|
||||||
|
`<span title="`+receivedAt.Format(time.DateTime)+
|
||||||
|
` UTC">3 minutes ago</span>`,
|
||||||
|
)
|
||||||
|
assert.Contains(t, body, `<span title="Body size">2.0 kB</span>`)
|
||||||
|
assert.Contains(t, body, ">1.5s</span>")
|
||||||
|
assert.Contains(t, body, ">in progress</span>")
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestHandleSourceDetail_StatusWithSingleHTTPTarget proves that a
|
||||||
|
// webhook with exactly one HTTP target shows, colour-coded, what
|
||||||
|
// that target answered for each event. The log target beside it
|
||||||
|
// does not count against "exactly one".
|
||||||
|
func TestHandleSourceDetail_StatusWithSingleHTTPTarget(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
f := newRecentEventsFixture(t)
|
||||||
|
target := seedTarget(t, f.db, f.webhook.ID, database.TargetTypeHTTP)
|
||||||
|
seedTarget(t, f.db, f.webhook.ID, database.TargetTypeLog)
|
||||||
|
|
||||||
|
now := time.Now()
|
||||||
|
|
||||||
|
for _, code := range []int{204, 302, 404, 503, 0} {
|
||||||
|
dlv := f.delivery(
|
||||||
|
t, f.event(t, contentTypeJSON, "{}", now), target.ID,
|
||||||
|
database.DeliveryStatusDelivered,
|
||||||
|
)
|
||||||
|
f.attempt(t, dlv, code, time.Second)
|
||||||
|
}
|
||||||
|
|
||||||
|
f.delivery(
|
||||||
|
t, f.event(t, contentTypeJSON, "{}", now), target.ID,
|
||||||
|
database.DeliveryStatusPending,
|
||||||
|
)
|
||||||
|
f.event(t, contentTypeJSON, "{}", now)
|
||||||
|
|
||||||
|
// A replay is a newer delivery, and its answer is the one shown.
|
||||||
|
replayed := f.event(t, contentTypeJSON, "{}", now)
|
||||||
|
f.attempt(t, f.delivery(
|
||||||
|
t, replayed, target.ID, database.DeliveryStatusFailed,
|
||||||
|
), 502, time.Second)
|
||||||
|
f.attempt(t, f.deliveryQueuedAt(
|
||||||
|
t, replayed, target.ID, database.DeliveryStatusDelivered,
|
||||||
|
now.Add(time.Minute),
|
||||||
|
), 200, time.Second)
|
||||||
|
|
||||||
|
body := f.render(t)
|
||||||
|
|
||||||
|
assert.Contains(t, body, statusCell("text-green-600", "204"))
|
||||||
|
assert.Contains(t, body, statusCell("text-gray-500", "302"))
|
||||||
|
assert.Contains(t, body, statusCell("text-yellow-600", "404"))
|
||||||
|
assert.Contains(t, body, statusCell("text-red-600", "503"))
|
||||||
|
assert.Contains(t, body, statusCell("text-red-600", "no response"))
|
||||||
|
assert.Contains(t, body, statusCell("text-gray-400", "pending"))
|
||||||
|
assert.Contains(t, body, statusCell("text-gray-400", "not sent"))
|
||||||
|
assert.Contains(t, body, statusCell("text-green-600", "200"))
|
||||||
|
assert.NotContains(t, body, ">502<")
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestHandleSourceDetail_NoStatusWithoutSingleHTTPTarget proves the
|
||||||
|
// status column is absent when the webhook has no HTTP target or
|
||||||
|
// more than one.
|
||||||
|
func TestHandleSourceDetail_NoStatusWithoutSingleHTTPTarget(
|
||||||
|
t *testing.T,
|
||||||
|
) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
cases := map[string][]database.TargetType{
|
||||||
|
"none": {database.TargetTypeLog},
|
||||||
|
"several": {database.TargetTypeHTTP, database.TargetTypeHTTP},
|
||||||
|
}
|
||||||
|
|
||||||
|
for name, types := range cases {
|
||||||
|
t.Run(name, func(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
f := newRecentEventsFixture(t)
|
||||||
|
event := f.event(t, contentTypeJSON, "{}", time.Now())
|
||||||
|
|
||||||
|
for _, tt := range types {
|
||||||
|
target := seedTarget(t, f.db, f.webhook.ID, tt)
|
||||||
|
f.attempt(t, f.delivery(
|
||||||
|
t, event, target.ID,
|
||||||
|
database.DeliveryStatusDelivered,
|
||||||
|
), 200, time.Second)
|
||||||
|
}
|
||||||
|
|
||||||
|
body := f.render(t)
|
||||||
|
|
||||||
|
assert.Contains(t, body, `title="Body size"`)
|
||||||
|
assert.NotContains(t, body, statusTitle)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -415,16 +415,21 @@ func (h *Handlers) renderSourceDetail(
|
|||||||
"webhook_id = ?", webhook.ID,
|
"webhook_id = ?", webhook.ID,
|
||||||
).Find(&targets)
|
).Find(&targets)
|
||||||
|
|
||||||
var events []database.Event
|
var events []RecentEventView
|
||||||
|
|
||||||
if h.dbMgr.DBExists(webhook.ID) {
|
if h.dbMgr.DBExists(webhook.ID) {
|
||||||
webhookDB, dbErr := h.dbMgr.GetDB(webhook.ID)
|
webhookDB, dbErr := h.dbMgr.GetDB(webhook.ID)
|
||||||
if dbErr == nil {
|
if dbErr == nil {
|
||||||
webhookDB.Where(
|
events, dbErr = h.loadRecentEvents(
|
||||||
"webhook_id = ?", webhook.ID,
|
webhookDB, webhook.ID, singleHTTPTargetID(targets),
|
||||||
).Order("created_at DESC").Limit(
|
)
|
||||||
recentEventLimit,
|
}
|
||||||
).Find(&events)
|
|
||||||
|
if dbErr != nil {
|
||||||
|
h.log.Error(
|
||||||
|
"failed to load recent events",
|
||||||
|
"webhook_id", webhook.ID, "error", dbErr,
|
||||||
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -133,6 +133,11 @@ 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(
|
||||||
@@ -164,6 +169,8 @@ 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,16 +304,44 @@ 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,
|
||||||
// so a second WriteHeader would change nothing the client sees and
|
// cookie included, so a second WriteHeader would change nothing the
|
||||||
// would draw net/http's "superfluous response.WriteHeader" report.
|
// client sees and would draw net/http's "superfluous
|
||||||
|
// 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"))
|
||||||
|
|
||||||
@@ -331,6 +359,7 @@ 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"])
|
||||||
|
|||||||
@@ -92,11 +92,25 @@ func (s *Server) setupGlobalMiddleware() {
|
|||||||
func (s *Server) setupRoutes() {
|
func (s *Server) setupRoutes() {
|
||||||
s.router.Get("/", s.h.HandleIndex())
|
s.router.Get("/", s.h.HandleIndex())
|
||||||
|
|
||||||
s.router.Mount(
|
// Static assets answer GET and HEAD only. chi's default 405
|
||||||
"/s",
|
// carries no Allow header, so this group supplies its own.
|
||||||
http.StripPrefix("/s", http.FileServer(http.FS(static.Static))),
|
staticFiles := http.StripPrefix(
|
||||||
|
"/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.
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -396,13 +396,15 @@ func (e *testEnv) storedHash(t *testing.T, username string) string {
|
|||||||
|
|
||||||
// --- /s static group ---
|
// --- /s static group ---
|
||||||
|
|
||||||
// TestStaticServesEveryMethod pins what the static mount actually
|
// TestStaticServesOnlyGetAndHead pins the methods the static group
|
||||||
// answers. chi's Mount registers the handler for all methods and
|
// answers: GET and HEAD are served the asset, and the other methods
|
||||||
// http.FileServer only special-cases HEAD (by suppressing the body),
|
// chi routes (POST, PUT, DELETE and the rest) are refused with 405
|
||||||
// so a POST or a DELETE to an asset is served the file rather than
|
// and an Allow header naming those two. A method chi does not route,
|
||||||
// refused. The README documents this; the test is what keeps the two
|
// such as PROPFIND, is refused with 405 by the top-level router
|
||||||
// from drifting.
|
// before it reaches the static group, so it gets no Allow header.
|
||||||
func TestStaticServesEveryMethod(t *testing.T) {
|
// The README documents this; the test is what keeps the two from
|
||||||
|
// drifting.
|
||||||
|
func TestStaticServesOnlyGetAndHead(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
env := newTestEnv(t)
|
env := newTestEnv(t)
|
||||||
@@ -417,6 +419,7 @@ func TestStaticServesEveryMethod(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()
|
||||||
@@ -428,18 +431,38 @@ func TestStaticServesEveryMethod(t *testing.T) {
|
|||||||
w := httptest.NewRecorder()
|
w := httptest.NewRecorder()
|
||||||
env.router.ServeHTTP(w, req)
|
env.router.ServeHTTP(w, req)
|
||||||
|
|
||||||
assert.Equal(t, http.StatusOK, w.Code,
|
switch method {
|
||||||
"static mount answers every method")
|
case http.MethodGet:
|
||||||
|
assert.Equal(t, http.StatusOK, w.Code)
|
||||||
if method == http.MethodHead {
|
assert.Equal(t, body, w.Body.Bytes(),
|
||||||
|
"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":
|
||||||
return
|
assert.Equal(
|
||||||
|
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")
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -187,12 +187,24 @@
|
|||||||
<div class="divide-y divide-gray-100">
|
<div class="divide-y divide-gray-100">
|
||||||
{{range .Events}}
|
{{range .Events}}
|
||||||
<div class="p-4">
|
<div class="p-4">
|
||||||
<div class="flex items-center justify-between">
|
<div class="flex flex-wrap items-center justify-between gap-3">
|
||||||
<div class="flex items-center gap-3">
|
<div class="flex flex-wrap items-center gap-3">
|
||||||
<span class="badge-info">{{.Method}}</span>
|
<span class="badge-info">{{.Method}}</span>
|
||||||
<span class="text-sm text-gray-500">{{.ContentType}}</span>
|
<span class="text-sm text-gray-500 break-all">{{.ContentType}}</span>
|
||||||
|
{{if .ResubmittedFromID}}
|
||||||
|
<span class="text-xs text-gray-500" title="This event is a copy of {{.ResubmittedFromID}}">resubmitted copy</span>
|
||||||
|
{{end}}
|
||||||
|
</div>
|
||||||
|
<div class="flex flex-wrap items-center gap-3 text-xs text-gray-400">
|
||||||
|
<span title="Body size">{{.Size}}</span>
|
||||||
|
{{if .ProcessingTime}}
|
||||||
|
<span title="Processing time: how long the slowest delivery took, from being queued to its last attempt">{{.ProcessingTime}}</span>
|
||||||
|
{{end}}
|
||||||
|
{{if .Status}}
|
||||||
|
<span class="font-medium {{.StatusClass}}" title="HTTP status from the HTTP target">{{.Status}}</span>
|
||||||
|
{{end}}
|
||||||
|
<span title="{{.ReceivedUTC}}">{{.Received}}</span>
|
||||||
</div>
|
</div>
|
||||||
<span class="text-xs text-gray-400">{{.CreatedAt.Format "2006-01-02 15:04:05 UTC"}}</span>
|
|
||||||
</div>
|
</div>
|
||||||
</div>
|
</div>
|
||||||
{{else}}
|
{{else}}
|
||||||
|
|||||||
Reference in New Issue
Block a user