Author SHA1 Message Date
clawbot 7b09802967 State what the default blocklist covers (closes #244)
check / check (push) Successful in 4m23s
The default blocklist covers private and reserved space, plus public
addresses that serve cloud credentials. A provider's other services on
public addresses, such as IBM Cloud's 161.26.0.0/16 and 166.8.0.0/14,
are deliberately not on it: they serve no credentials, reaching them
can be legitimate, and every cloud has some, so a partial list would
promise coverage it does not give.

The README's egress section and the comment above blockedNetworks now
state this rule, so nobody infers wider coverage and a future candidate
can be accepted or refused against it. No list change.

Model: opus-5-5
2026-09-29 08:44:00 +00:00
clawbot f0adeafde3 Drop Set-Cookie from the recovered 500 (closes #193)
check / check (push) Successful in 3m35s
When a handler sets a cookie and then panics before sending
anything, the recover middleware now deletes Set-Cookie before
writing its 500, so a request that failed never hands the client
a credential. Every other header, Location included, is left as
http.Error leaves it, matching chi's Recoverer. A response that
was already sent is untouched.

Tests cover the uncommitted case (no cookie, Location kept) and
assert the cookie still reaches the client when the response was
committed before the panic.

Model: opus-5-5
2026-09-29 10:30:26 +02:00
clawbot f755c03110 Default-block Azure WireServer's public address (closes #245)
check / check (push) Successful in 4m34s
Add 168.63.129.16 (Azure WireServer) to blockedNetworks, the default
blocklist, not alwaysBlockedNetworks: it is public unicast, so an
operator who lists it in ALLOWED_EGRESS_CIDRS can reach it again. The
refusal message, the allowlist startup warning, the README and the
comments no longer call every blocked address private/reserved, and
no longer claim the allowlist cannot open any metadata endpoint.

Sources:
- https://learn.microsoft.com/en-us/azure/virtual-network/what-is-ip-address-168-63-129-16
- https://learn.microsoft.com/en-us/azure/virtual-machines/metadata-security-protocol/overview

Deviation: 147.75.207.243 (Equinix Metal) is not added; Equinix
documents only a hostname, and the service was sunset on 2026-06-30.

Model: opus-5-5
2026-09-29 10:22:07 +02:00
clawbot 4a724130ca Close archive writers when the delivery engine stops (closes #280)
check / check (push) Successful in 3m45s
The engine cached archive writers and never closed them at shutdown,
so after a clean stop an archive's rows could sit in its -wal while
the .db held no table. The engine's stop hook now evicts every cached
writer once its workers have returned, the same way deleting a webhook
does, so a clean stop leaves each archive as one file and a late write
is refused. If the workers do not return within the stop budget, the
writers are left open as a kill would leave them: closing would wait
on a write in progress, and a still-running worker would open new
ones.

The README no longer says archives keep their sidecars across a clean
stop.

Model: opus-5-5
2026-09-29 08:30:22 +02:00
clawbot d4f4ddf51f Send Content-Type once on a delivery (closes #246)
check / check (push) Successful in 3m20s
A delivery set Content-Type from the event's ContentType and then
added the inbound Content-Type from the event's stored headers, so a
target could receive two values. The inbound Content-Type is no
longer forwarded from the stored headers; the receiver already saves
it as the event's ContentType.

Which value is sent is now stated at applyRequestHeaders: a
Content-Type configured on the target, otherwise the event's
ContentType, otherwise none. A configured one still survives a
cross-origin 307/308 with its body.

Model: opus-5-5
2026-09-29 07:11:55 +02:00
clawbot 51580a2bc6 Fail a pending delivery whose target was deleted (closes #293)
check / check (push) Successful in 3m38s
Restart recovery and the pending sweep skipped a pending delivery
whose target was missing from the batch's target map, every minute,
for the life of the database. A miss now asks loadTarget: no row
fails the delivery terminally with a recorded reason; any other error
leaves it pending, since the map is also empty when its query failed;
a target found there is used.

The failure goes through the ownership-gated function the retrying
paths already used, now failMissingTarget. Once it owns the delivery
it re-reads the row and fails it only if the status is unchanged, so
a delivery sent and settled in between is left alone.

Model: opus-5-5
2026-09-29 06:48:19 +02:00
clawbot e0b211f960 Make deliveries refused while half-open wait a cooldown (closes #306)
check / check (push) Successful in 5m0s
While the breaker was half-open, Allow refused every delivery but the
probe and CooldownRemaining returned zero, so each queued task for the
target went straight back onto the retry channel and rewrote its status
on every pass until the probe finished.

CooldownRemaining now returns the whole cooldown while half-open, so a
refused delivery waits that long. A refused delivery already at
retrying is not written again, so the retry counter now moves only
when a refusal moves a delivery into retrying.

Model: opus-5-5
2026-09-29 05:48:20 +02:00
20 changed files with 623 additions and 80 deletions
+43 -33
View File
@@ -157,6 +157,19 @@ 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
@@ -195,15 +208,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 +256,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 +983,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 +1061,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 +1082,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
@@ -2051,7 +2058,7 @@ rescans the database anyway).
| ----------- | -------- | | ----------- | -------- |
| **Closed** | Normal operation. Deliveries flow through. Consecutive failures are counted. | | **Closed** | Normal operation. Deliveries flow through. Consecutive failures are counted. |
| **Open** | Target appears down. Deliveries are skipped and rescheduled for after the cooldown. | | **Open** | Target appears down. Deliveries are skipped and rescheduled for after the cooldown. |
| **Half-Open** | Cooldown expired. One probe delivery is allowed to test if the target has recovered. | | **Half-Open** | Cooldown expired. One probe delivery is allowed to test if the target has recovered. Other deliveries are rescheduled for one whole cooldown later. |
**Transitions:** **Transitions:**
@@ -2089,7 +2096,9 @@ operations), and log targets (stdout) do not use circuit breakers.
When a circuit is open and a new delivery arrives, the engine marks the When a circuit is open and a new delivery arrives, the engine marks the
delivery as `retrying` and schedules a retry timer for after the delivery as `retrying` and schedules a retry timer for after the
remaining cooldown period. This ensures no deliveries are lost — they're remaining cooldown period. This ensures no deliveries are lost — they're
just delayed until the target is healthy again. just delayed until the target is healthy again. A delivery already in
`retrying` keeps that status without another database write each time
the breaker turns it away.
### Metrics ### Metrics
@@ -2103,7 +2112,7 @@ arriving and being stored, they are just not getting anywhere.
| Metric | Type | Meaning | | Metric | Type | Meaning |
| ------ | ---- | ------- | | ------ | ---- | ------- |
| `webhooker_events_received_total` | counter | Events received and durably stored. Compare against the delivery counters on one dashboard | | `webhooker_events_received_total` | counter | Events received and durably stored. Compare against the delivery counters on one dashboard |
| `webhooker_delivery_attempts_total` | counter | Delivery attempts actually dispatched to a target. A delivery an open circuit breaker refused is not one: it is counted as a retry instead | | `webhooker_delivery_attempts_total` | counter | Delivery attempts actually dispatched to a target. A delivery a circuit breaker refused is not one: it is counted as a retry instead, but only when the refusal moves it into `retrying` |
| `webhooker_deliveries_succeeded_total` | counter | Deliveries that reached `delivered` | | `webhooker_deliveries_succeeded_total` | counter | Deliveries that reached `delivered` |
| `webhooker_deliveries_failed_total` | counter | Deliveries that failed terminally and will not be retried | | `webhooker_deliveries_failed_total` | counter | Deliveries that failed terminally and will not be retried |
| `webhooker_delivery_retries_total` | counter | Deliveries put back into `retrying` | | `webhooker_delivery_retries_total` | counter | Deliveries put back into `retrying` |
@@ -3086,7 +3095,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
+12 -9
View File
@@ -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), ","),
) )
+7 -6
View File
@@ -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")
}) })
} }
} }
+10 -2
View File
@@ -76,12 +76,20 @@ func (cb *CircuitBreaker) Allow() bool {
} }
} }
// CooldownRemaining returns how much time is left before // CooldownRemaining returns how long a delivery that Allow refused
// an open circuit transitions to half-open. // should wait before it is tried again. Closed, it returns zero.
// Open, it returns what is left of the cooldown, or zero once that
// has passed. Half-open, it returns the whole cooldown: the one
// probe delivery is still in flight, and if it fails the circuit
// reopens for that long.
func (cb *CircuitBreaker) CooldownRemaining() time.Duration { func (cb *CircuitBreaker) CooldownRemaining() time.Duration {
cb.mu.Lock() cb.mu.Lock()
defer cb.mu.Unlock() defer cb.mu.Unlock()
if cb.state == CircuitHalfOpen {
return cb.cooldown
}
if cb.state != CircuitOpen { if cb.state != CircuitOpen {
return 0 return 0
} }
+5 -3
View File
@@ -267,7 +267,7 @@ func TestCircuitBreaker_CooldownRemaining_ClosedReturnsZero(
) )
} }
func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsZero( func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsCooldown(
t *testing.T, t *testing.T,
) { ) {
t.Parallel() t.Parallel()
@@ -282,9 +282,11 @@ func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsZero(
require.True(t, cb.Allow()) require.True(t, cb.Allow())
assert.Equal(t, time.Duration(0), // The cooldown newShortCooldownCB gives the breaker.
assert.Equal(t, 50*time.Millisecond,
cb.CooldownRemaining(), cb.CooldownRemaining(),
"half-open circuit should have zero cooldown remaining", "a delivery refused while half-open should wait "+
"a whole cooldown",
) )
} }
+11
View File
@@ -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
@@ -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",
)
}
+175
View File
@@ -17,6 +17,7 @@ import (
"time" "time"
"github.com/google/uuid" "github.com/google/uuid"
"github.com/prometheus/client_golang/prometheus"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
"gorm.io/driver/sqlite" "gorm.io/driver/sqlite"
@@ -24,6 +25,7 @@ import (
_ "modernc.org/sqlite" _ "modernc.org/sqlite"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/delivery"
"sneak.berlin/go/webhooker/internal/metrics"
) )
// testContentType is the event content type used in tests. // testContentType is the event content type used in tests.
@@ -894,6 +896,100 @@ func TestDeliverHTTP_CircuitBreakerBlocks(t *testing.T) {
) )
} }
// recordingScheduler keeps the delay of every retry it is asked to
// schedule, and schedules nothing.
type recordingScheduler struct {
delays []time.Duration
}
func (s *recordingScheduler) ScheduleRetry(
_ delivery.Task, delay time.Duration,
) {
s.delays = append(s.delays, delay)
}
// TestDeliverHTTP_HalfOpenBreakerDelaysQueuedTasks proves that while a
// half-open breaker's one probe delivery is in flight, every other task
// for the target is put back with a whole cooldown as its delay rather
// than none, and that its status is written the first time the breaker
// turns it away and not on each pass after that.
func TestDeliverHTTP_HalfOpenBreakerDelaysQueuedTasks(t *testing.T) {
t.Parallel()
db := testWebhookDB(t)
e := testEngine(t, 1)
// Every write of retrying moves the retry counter, so on a registry
// this test owns the counter is the number of those writes.
reg := prometheus.NewRegistry()
e.ExportSetMetrics(metrics.New(reg))
targetID := uuid.New().String()
cb := newShortCooldownCB(t)
e.ExportSetCircuitBreaker(targetID, cb)
for range delivery.ExportDefaultFailureThreshold {
cb.RecordFailure()
}
time.Sleep(60 * time.Millisecond)
require.True(t, cb.Allow(), "the probe delivery should go through")
require.Equal(t, delivery.CircuitHalfOpen, cb.State())
cfg := newHTTPTargetConfig(
"http://will-not-be-called.invalid",
)
sched := &recordingScheduler{}
const queued, passes = 3, 4
for range queued {
event := seedEvent(t, db, `{"cb":"half-open"}`)
dlv := seedDelivery(
t, db, event.ID, targetID,
database.DeliveryStatusPending,
)
for range passes {
// Each pass starts from the stored row, as a retry does.
var row database.Delivery
require.NoError(t, db.First(
&row, "id = ?", dlv.ID,
).Error)
fix := buildHTTPFixture(
row, event, targetID,
"test-cb-half-open", cfg, 5, 1,
)
e.ExportDeliverHTTPWithScheduler(
context.TODO(), db, fix.Delivery, fix.Task, sched,
)
}
assertDeliveryStatus(t, db, dlv.ID,
database.DeliveryStatusRetrying,
)
}
require.Len(t, sched.delays, queued*passes)
for _, delay := range sched.delays {
// The cooldown newShortCooldownCB gives the breaker.
assert.Equal(t, 50*time.Millisecond, delay,
"a task turned away while half-open should wait "+
"a whole cooldown",
)
}
assert.InDelta(t, float64(queued),
mCounter(t, reg, mRetries, mTypeHTTP), 0,
"status should be written once per task, not once per pass",
)
}
func TestGetCircuitBreaker_CreatesOnDemand(t *testing.T) { func TestGetCircuitBreaker_CreatesOnDemand(t *testing.T) {
t.Parallel() t.Parallel()
@@ -1070,6 +1166,10 @@ func TestIsForwardableHeader(t *testing.T) {
assert.False(t, assert.False(t,
delivery.ExportIsForwardableHeader("Content-Length"), delivery.ExportIsForwardableHeader("Content-Length"),
) )
assert.False(t,
delivery.ExportIsForwardableHeader("Content-Type"),
)
} }
func TestTruncate(t *testing.T) { func TestTruncate(t *testing.T) {
@@ -1151,6 +1251,81 @@ func TestDoHTTPRequest_ForwardsHeaders(t *testing.T) {
) )
} }
// The event's stored inbound headers carry the same Content-Type the
// receiver saved as the event's ContentType, so a delivery could send
// it twice. It must go out exactly once, with a Content-Type configured
// on the target winning, then the event's ContentType.
func TestApplyRequestHeaders_SendsOneContentType(t *testing.T) {
t.Parallel()
cases := map[string]struct {
inbound string
event string
configured string
want []string
}{
"inbound and event agree": {
inbound: testContentType,
event: testContentType,
want: []string{testContentType},
},
"inbound and event disagree": {
inbound: "text/plain",
event: testContentType,
want: []string{testContentType},
},
"event has none": {
inbound: testContentType,
want: nil,
},
"target configures its own": {
inbound: testContentType,
event: testContentType,
configured: "application/xml",
want: []string{"application/xml"},
},
}
for name, tc := range cases {
t.Run(name, func(t *testing.T) {
t.Parallel()
inbound, err := json.Marshal(map[string][]string{
headerContentType: {tc.inbound},
})
require.NoError(t, err)
cfg := &delivery.HTTPTargetConfig{}
if tc.configured != "" {
cfg.Headers = map[string]string{
headerContentType: tc.configured,
}
}
req, err := http.NewRequestWithContext(
context.Background(),
http.MethodPost,
"https://target.example.com/hook",
http.NoBody,
)
require.NoError(t, err)
delivery.ExportApplyRequestHeaders(
req,
&database.Event{
Headers: string(inbound),
ContentType: tc.event,
},
cfg,
)
assert.Equal(t,
tc.want, req.Header.Values(headerContentType),
)
})
}
}
func TestProcessDelivery_RoutesToCorrectHandler( func TestProcessDelivery_RoutesToCorrectHandler(
t *testing.T, t *testing.T,
) { ) {
+36
View File
@@ -101,6 +101,19 @@ func (e *Engine) ExportDeliverHTTP(
e.httpTarget.Deliver(ctx, webhookDB, d, task, e) e.httpTarget.Deliver(ctx, webhookDB, d, task, e)
} }
// ExportDeliverHTTPWithScheduler delivers via the http target, handing
// any retry to sched instead of the engine, so a test can see the
// delay each retry is given.
func (e *Engine) ExportDeliverHTTPWithScheduler(
ctx context.Context,
webhookDB *gorm.DB,
d *database.Delivery,
task *Task,
sched Scheduler,
) {
e.httpTarget.Deliver(ctx, webhookDB, d, task, sched)
}
// ExportDeliverDatabase delivers via the database target. // ExportDeliverDatabase delivers via the database target.
func (e *Engine) ExportDeliverDatabase( func (e *Engine) ExportDeliverDatabase(
webhookDB *gorm.DB, d *database.Delivery, webhookDB *gorm.DB, d *database.Delivery,
@@ -179,6 +192,14 @@ func (e *Engine) ExportGetCircuitBreaker(
return e.httpTarget.getCircuitBreaker(targetID) return e.httpTarget.getCircuitBreaker(targetID)
} }
// ExportSetCircuitBreaker makes cb the http target's circuit breaker
// for targetID, so a test can use one with a short cooldown.
func (e *Engine) ExportSetCircuitBreaker(
targetID string, cb *CircuitBreaker,
) {
e.httpTarget.circuitBreakers.Store(targetID, cb)
}
// ExportParseHTTPConfig exposes parseHTTPConfig. // ExportParseHTTPConfig exposes parseHTTPConfig.
func (e *Engine) ExportParseHTTPConfig( func (e *Engine) ExportParseHTTPConfig(
configJSON string, configJSON string,
@@ -331,6 +352,21 @@ func (e *Engine) ExportFailMissingTarget(
e.failMissingTarget(webhookDB, webhookID, d) e.failMissingTarget(webhookDB, webhookID, d)
} }
// ExportSendRecoveredDeliveries exposes sendRecoveredDeliveries, so a
// test can hand it a target map that lacks a delivery's target.
func (e *Engine) ExportSendRecoveredDeliveries(
ctx context.Context,
webhookDB *gorm.DB,
deliveries []database.Delivery,
webhookID string,
targetMap map[string]database.Target,
settled map[string]struct{},
) {
e.sendRecoveredDeliveries(
ctx, webhookDB, deliveries, webhookID, targetMap, settled,
)
}
// ExportDeliveryCh returns the delivery channel. // ExportDeliveryCh returns the delivery channel.
func (e *Engine) ExportDeliveryCh() chan Task { func (e *Engine) ExportDeliveryCh() chan Task {
return e.deliveryCh return e.deliveryCh
+4 -3
View File
@@ -412,9 +412,10 @@ func TestDeliveryMetrics_BreakerBlockedIsNotAnAttempt(
s.Engine.ExportProcessRetryTask(context.TODO(), &blocked) s.Engine.ExportProcessRetryTask(context.TODO(), &blocked)
// The breaker refused it: rescheduled, so the retry counter // The breaker refused it: rescheduled without rewriting the
// moved, but nothing was attempted or timed. // retrying status it already had, so the retry counter did not
assert.InDelta(t, retriesBefore+1, // move, and nothing was attempted or timed.
assert.InDelta(t, retriesBefore,
mCounter(t, reg, mRetries, mTypeHTTP), 0) mCounter(t, reg, mRetries, mTypeHTTP), 0)
assert.InDelta(t, threshold, assert.InDelta(t, threshold,
mCounter(t, reg, mAttempts, mTypeHTTP), 0) mCounter(t, reg, mAttempts, mTypeHTTP), 0)
+12 -6
View File
@@ -339,10 +339,11 @@ 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 the delivery path's own two are deliberately // is not in the set, and neither is the inbound Content-Type, because
// excluded: Content-Type describes the body, which a 307 carries // it is not forwarded. Two more are deliberately excluded: a
// across hosts, and the inbound User-Agent every real sender supplies // Content-Type configured on the target describes the body, which a
// is overwritten before the request goes out. // 307 carries across hosts, and the inbound User-Agent every real
// 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()
@@ -371,6 +372,7 @@ 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,
}, },
}, },
) )
@@ -378,7 +380,11 @@ 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 is never forwarded, Content-Type and "+ "Host and the inbound Content-Type are never "+
"User-Agent are the delivery path's own", "forwarded, User-Agent is 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",
) )
} }
+16 -7
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 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,11 +37,18 @@ 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.
// //
// 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
@@ -122,6 +129,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 +225,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 +329,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.
+35
View File
@@ -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
+18
View File
@@ -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",
)
}
+2 -1
View File
@@ -11,10 +11,11 @@ import (
"sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/delivery"
) )
// Literals these tests repeat, named so that the header name and the // Literals these tests repeat, named so that the header names and the
// keep-forever archive config each have one definition. // keep-forever archive config each have one definition.
const ( const (
headerAuthorization = "Authorization" headerAuthorization = "Authorization"
headerContentType = "Content-Type"
bearerValue = "Bearer abc" bearerValue = "Bearer abc"
archiveConfigNever = "{\"expiry\":\"never\"}" archiveConfigNever = "{\"expiry\":\"never\"}"
) )
+21 -8
View File
@@ -197,10 +197,14 @@ func (c *httpCore) circuitBreakerBlock(
"cooldown_remaining", remaining, "cooldown_remaining", remaining,
) )
c.eng.settleStatus( // A delivery already at retrying is left as it is, so a task
webhookDB, d, d.Target.Type, // the breaker keeps turning away writes nothing each time.
database.DeliveryStatusRetrying, if d.Status != database.DeliveryStatusRetrying {
) c.eng.settleStatus(
webhookDB, d, d.Target.Type,
database.DeliveryStatusRetrying,
)
}
retryTask := *task retryTask := *task
sched.ScheduleRetry(retryTask, remaining) sched.ScheduleRetry(retryTask, remaining)
@@ -537,6 +541,11 @@ func isForwardableHeader(name string) bool {
"Upgrade", "Proxy-Authorization", "Upgrade", "Proxy-Authorization",
"Proxy-Connection", "Content-Length": "Proxy-Connection", "Content-Length":
return false return false
case "Content-Type":
// applyRequestHeaders sets Content-Type itself. The receiver
// already stored this inbound value as the event's
// ContentType, so forwarding it too would send it twice.
return false
default: default:
return true return true
} }
@@ -549,6 +558,10 @@ func isForwardableHeader(name string) bool {
// policy strips exactly that set on a hop that leaves the origin, // policy strips exactly that set on a hop that leaves the origin,
// so the forward set is decided here and only here — a header added // so the forward set is decided here and only here — a header added
// to it is covered off-origin without a second edit elsewhere. // to it is covered off-origin without a second edit elsewhere.
//
// Content-Type goes out once: a Content-Type configured on the target
// wins, otherwise the event's ContentType, otherwise none. The inbound
// Content-Type in the event's headers is never forwarded.
func applyRequestHeaders( func applyRequestHeaders(
req *http.Request, req *http.Request,
event *database.Event, event *database.Event,
@@ -569,10 +582,10 @@ func applyRequestHeaders(
req.Header.Set("User-Agent", "webhooker/1.0") req.Header.Set("User-Agent", "webhooker/1.0")
// Content-Type describes the body being sent rather than the // A Content-Type configured on the target describes the body
// sender, and the delivery path sets it from the event itself. // being sent rather than the sender. A 307/308 preserves the
// A 307/308 preserves the body across hosts, so stripping it // body across hosts, so stripping it would send that body
// would send that body untyped. // untyped.
delete(originScoped, "Content-Type") delete(originScoped, "Content-Type")
// User-Agent is overwritten just above, so an inbound one never // User-Agent is overwritten just above, so an inbound one never
+47
View File
@@ -696,6 +696,53 @@ func TestSweepPending_TargetDeleted(t *testing.T) {
assert.Contains(t, last.Error, "was deleted") assert.Contains(t, last.Error, "was deleted")
} }
// TestSendRecoveredDeliveries_TargetMissingFromMap: the batch's target
// map is empty when its query failed, so every delivery in the batch is
// looked up on its own. A healthy one is sent to the target that lookup
// finds.
func TestSendRecoveredDeliveries_TargetMissingFromMap(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
iCreateTarget(
t, s.MainDB, targetID, s.WebhookID, "found-on-lookup",
database.TargetTypeLog, "", 0,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"map":"empty"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
s.Engine.ExportSendRecoveredDeliveries(
context.Background(), s.WebhookDB,
[]database.Delivery{d}, s.WebhookID,
map[string]database.Target{}, nil,
)
tasks := fDrain(s.Engine)
require.Len(t, tasks, 1,
"the healthy delivery was not queued exactly once",
)
assert.Equal(t, d.ID, tasks[0].DeliveryID)
assert.Equal(t, targetID, tasks[0].TargetID)
assert.Equal(t, database.TargetTypeLog, tasks[0].TargetType)
iAssertStatus(
t, s.WebhookDB, d.ID,
database.DeliveryStatusPending,
)
}
// TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone: a failed // TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone: a failed
// read of the main database is not a deleted target. Restart recovery // read of the main database is not a deleted target. Restart recovery
// holds every pending delivery of the webhook in one batch, so failing // holds every pending delivery of the webhook in one batch, so failing
+7
View File
@@ -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(
+31 -2
View File
@@ -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"])