Author SHA1 Message Date
sneak 306c2cb9b1 Send Content-Type once on a delivery (closes #246)
check / check (push) Successful in 4m3s
A delivery set Content-Type from the event's ContentType and then
added the inbound Content-Type from the event's stored headers, so
the target could receive two values. The inbound Content-Type is no
longer forwarded from the stored headers; the receiver already saved
it as the event's ContentType.

Which value wins is now stated at applyRequestHeaders: a Content-Type
configured on the target, otherwise the event's ContentType,
otherwise none.

Model: opus-5-5
2026-09-29 03:56:16 +00:00
14 changed files with 92 additions and 672 deletions
+30 -38
View File
@@ -157,19 +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 is all the default blocklist covers: private and reserved space,
plus public addresses that serve cloud credentials. A cloud provider's
other services on public addresses are not refused — IBM Cloud's
`161.26.0.0/16` and `166.8.0.0/14`, for example, which carry its DNS
resolvers, time servers and package mirrors. They serve no credentials,
reaching them can be a legitimate delivery, and every cloud has some, so
a partial list would promise coverage it does not give.
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
@@ -208,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 |
| ----------------------- | ---------- | | ----------------------- | ---------- |
@@ -256,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.
@@ -983,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
@@ -1061,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
@@ -1082,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
@@ -3095,8 +3088,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
+9 -12
View File
@@ -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), ","),
) )
+6 -7
View File
@@ -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")
}) })
} }
} }
+23 -70
View File
@@ -362,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")
@@ -385,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
@@ -739,7 +728,9 @@ 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.failMissingTarget(webhookDB, webhookID, d) e.failMissingTargetRetry(
webhookDB, webhookID, d,
)
return return
} }
@@ -1142,7 +1133,9 @@ 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.failMissingTarget(webhookDB, webhookID, d) e.failMissingTargetRetry(
webhookDB, webhookID, d,
)
return return
} }
@@ -1256,19 +1249,19 @@ func (e *Engine) failUnretryableRetry(
e.failDelivery(webhookDB, d, target.Type, reason) e.failDelivery(webhookDB, d, target.Type, reason)
} }
// failMissingTarget terminally fails a recovered delivery, pending or // failMissingTargetRetry terminally fails an orphaned retrying
// retrying, whose target row is gone. Restart recovery and the periodic // delivery whose target row is gone. Both restart recovery and the
// sweep call it for both statuses, so the transition exists once. // periodic sweep call it, so the transition exists once.
// //
// Until it existed those paths logged the failed lookup and moved on, // Until it existed both paths logged the failed lookup and returned,
// which left the delivery where it was for the life of the database and // which left the delivery retrying 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) failMissingTarget( func (e *Engine) failMissingTargetRetry(
webhookDB *gorm.DB, webhookDB *gorm.DB,
webhookID string, webhookID string,
d *database.Delivery, d *database.Delivery,
@@ -1281,37 +1274,13 @@ func (e *Engine) failMissingTarget(
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 recovered delivery: its target no longer exists", "failing orphaned retrying delivery: "+
"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,
) )
@@ -1345,14 +1314,15 @@ 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 "+
"has been failed terminally", "cannot be retried and has been failed "+
"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 "+
"has been failed terminally", "cannot be retried and has been failed terminally",
target.Name, target.Type, target.Name, target.Type,
) )
} }
@@ -2051,30 +2021,13 @@ func (e *Engine) sendRecoveredDeliveries(
target, ok := targetMap[deliveries[i].TargetID] target, ok := targetMap[deliveries[i].TargetID]
if !ok { if !ok {
// A missing entry does not mean the target is gone: the e.log.Error(
// map is also empty when its query failed. Only a lookup "target not found for delivery",
// that finds no row ends the delivery; any other error "delivery_id", deliveries[i].ID,
// leaves it pending for the next sweep. See "target_id", deliveries[i].TargetID,
// recoverSingleRetry. )
var err error
target, err = e.loadTarget(deliveries[i].TargetID) continue
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,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",
)
}
-25
View File
@@ -342,31 +342,6 @@ 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
+6 -12
View File
@@ -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",
) )
} }
+7 -16
View File
@@ -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,18 +37,11 @@ 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.
// //
// A public address belongs here only if it serves cloud
// credentials; a provider's other services on public addresses,
// such as its DNS resolvers or package mirrors, stay out, since
// reaching them can be legitimate and no list of them could be
// complete.
//
//nolint:gochecknoglobals // package-level network list is appropriate here //nolint:gochecknoglobals // package-level network list is appropriate here
var blockedNetworks []*net.IPNet var blockedNetworks []*net.IPNet
@@ -129,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
@@ -225,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)
} }
@@ -329,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.
-35
View File
@@ -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
-18
View File
@@ -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",
)
}
+9 -270
View File
@@ -18,17 +18,16 @@ 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. Section 4 is the same deleted-target gap for a pending // terminalise.
// 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 delivery against it at the // tSeedDeletedTarget creates a target, a retrying delivery against it
// given status with one recorded failed attempt, and then deletes the // with one recorded failed attempt, and then deletes the target the
// target the way the source page does. // 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
@@ -37,7 +36,6 @@ 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()
@@ -53,7 +51,8 @@ func tSeedDeletedTarget(
) )
d := iSeedDelivery( d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID, status, t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
) )
iSeedFailedResult(t, s.WebhookDB, d.ID) iSeedFailedResult(t, s.WebhookDB, d.ID)
@@ -174,7 +173,6 @@ 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(
@@ -212,7 +210,6 @@ 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
@@ -281,11 +278,11 @@ func TestSweepSingleRetry_TargetNeverExisted(t *testing.T) {
) )
} }
// TestFailMissingTarget_WritesNoTargetRow holds the new terminal path // TestFailMissingTargetRetry_WritesNoTargetRow holds the new terminal
// to the same rule as the existing one: no target row, and so no // path 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 TestFailMissingTarget_WritesNoTargetRow( func TestFailMissingTargetRetry_WritesNoTargetRow(
t *testing.T, t *testing.T,
) { ) {
t.Parallel() t.Parallel()
@@ -300,7 +297,6 @@ func TestFailMissingTarget_WritesNoTargetRow(
deliveryID := tSeedDeletedTarget( deliveryID := tSeedDeletedTarget(
t, s, "credential-bearing", hookURL, t, s, "credential-bearing", hookURL,
database.DeliveryStatusRetrying,
) )
s.Engine.ExportSweepWebhookRetries( s.Engine.ExportSweepWebhookRetries(
@@ -533,260 +529,3 @@ 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())
}
-7
View File
@@ -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(
+2 -31
View File
@@ -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"])