Compare commits

8 Commits
Author SHA1 Message Date
clawbot f01e5c45ad Give each event its own page and show bodies the same everywhere (closes #369)
check / check (push) Successful in 3m20s
Each row of the recent events on the webhook page links to the
event's own page, /hook/{id}/events/{eventID}, and expands to show its
body; only the newest starts expanded. The event's page shows its
details, its whole body and every delivery with its attempts.

One renderer, newBodyView with templates/event_body.html, shows a body
in all three places: whole up to 32 KiB, cut there in the lists with a
link to the event's page, JSON pretty-printed, a body of more than 200
lines or 32 KiB in a scrolling box, and a body that is not text left
out beside its download link. A resubmitted copy links to its
original's page.

Model: opus-5-5
2026-10-02 18:33:27 +00:00
clawbot 9526e961b5 Add a Download button that exports a database target's archive as gzipped JSON (closes #374)
check / check (push) Successful in 3m29s
Each database target on the webhook page has a Download button that streams its archive as gzipped JSON, archive-WEBHOOKNAME-TARGETNAME-TIME.json.gz, with names made safe by delivery.ArchiveFileName's function. The export reads one consistent snapshot through one cursor in a read-only transaction, so archive writes carry on, and holds the rename lock only while it reads the stored names and opens the file. It extends its write deadline as it writes, so a large archive downloads for as long as the client reads; a failure after the response has started aborts the connection so the browser marks the download failed. The request limit is now the service's own middleware, which no longer writes a 504 over a response already started.

Model: opus-5-5
2026-10-02 20:32:09 +02:00
clawbot 4915d60d8e Route the delivery tests' gorm.Open through gormlog (closes #462)
check / check (push) Successful in 3m17s
Six test-only gorm.Open calls in internal/delivery passed a bare gorm.Config, leaving the unfiltered idiom in the tree to be copied into production code, where every gorm.Open goes through gormlog.New. They now pass gormlog.New over a logger that discards, so no gorm.Open in the tree uses a bare gorm.Config. The stale sentence saying the tree has one test-only (*gorm.DB).Scan caller is corrected in the README and in the ParamsFilter comment: only tests call Scan, and what a test binds is fixture data. Test and documentation change only.

Model: opus-5-5
2026-10-02 20:20:49 +02:00
clawbot 0945831442 Clamp the HTTP drain by the tail-hook reserve (closes #170)
check / check (push) Successful in 3m17s
The HTTP drain at shutdown waited up to ShutdownTimeout regardless of how much of the stop budget earlier hooks had used, so a slow archive sweeper or retention reaper could eat the reserve the hooks after the server need, and the database close was skipped. The drain now waits at most the shorter of ShutdownTimeout and what is left of the budget less TailHookReserve, as the Sentry flush already does. The reserve is documented as derived from the two timeouts. Tests cover earlier hooks having spent part of the budget, on a clock that host speed cannot move, and pin that a drain on the full budget gets all of ShutdownTimeout.

Model: opus-5-5
2026-10-02 20:11:43 +02:00
clawbot f82b730c31 Pin the HTTP target's unpinned error checks (closes #285)
check / check (push) Successful in 3m20s
Seven error checks in internal/delivery/target_http.go could be removed without any test noticing, among them withRetry's check on writing the delivery result, the branch that leaves a sent delivery retrying and recoverable when its bookkeeping write fails. Each now fails a test when removed. The "send succeeded" case starts from a tripped circuit breaker, so a probe whose send succeeds but whose result write fails must still close the breaker. The checks in remainingBackoff and backoffElapsed stay unpinned: removing them gives the same answer, and they state a rule a reader needs. Test change only.

Model: opus-5-5
2026-10-02 19:37:36 +02:00
clawbot 1a1fee0874 Name the reaper's hard delete in the event body comments (closes #455)
check / check (push) Successful in 3m20s
The comments on eventBodyQuery and on TestHandleEventBodyDownload_ReapedEvent404s credited the soft-delete predicate for refusing a reaped event. The retention reaper deletes event rows outright and nothing soft-deletes an event, so a reaped event is simply gone. Both comments now say so; the test's "soft deleted" case is described as pinning the query's deleted_at predicate for a row no code produces today. Comments only.

Model: opus-5-5
2026-10-02 19:36:30 +02:00
clawbot 290925f184 Fix two resubmit comments and test the resubmit route's middleware (closes #252)
check / check (push) Successful in 3m18s
Two comments named the wrong mechanism: loadResubmitSource credited soft-delete for refusing a reaped event, though the retention reaper deletes event rows outright, and createAndFanOut claimed to be the only path that creates deliveries, though per-delivery replay creates one without an event. Both now say what the code does. The resubmit route's middleware had no tests through the router; new tests drive the production router to pin the refusal without a valid CSRF token, the rate limit, signed-out requests never spending it, and another webhook's event refused by the event lookup while the user's own event is accepted. Each fails with its check removed.

Model: opus-5-5
2026-10-02 19:20:40 +02:00
clawbot 73353bc8e5 Harden the (*gorm.DB).Scan guard test (closes #232)
check / check (push) Successful in 3m17s
The test that refuses production calls to (*gorm.DB).Scan, the one GORM path that bypasses the logger's value suppression, overstated what it checks and could pass while skipping a whole package. Its comments now say it matches receiver method names, not types, and name the evasion this leaves; GORM's Rows is dropped from the accepted names. Method values are stated as out of scope with the reason. The file-count floor is replaced by a check that every package the walk parses, static and templates included, was reached. The planted snippets are valid Go and cover each receiver form the guard claims to handle. Test change only.

Model: opus-5-5
2026-10-02 19:19:45 +02:00
36 changed files with 2332 additions and 151 deletions
+53 -22
View File
@@ -2033,6 +2033,29 @@ Because each `database` target has its own archive file, a target's
webhook with different expiries keep two archives, each pruned on its webhook with different expiries keep two archives, each pruned on its
own schedule. own schedule.
Each `database` target on the webhook page has a **Download** button,
which returns its archive as one gzipped JSON file,
`archive-{webhook_name}-{target_name}-{YYYYMMDDTHHMMSSZ}.json.gz`, the
names made safe as above and the time in UTC. The file holds one
object: `webhook` and `target`, each an `id` and a `name`;
`exported_at`; and `archived_events`, one object per archived row with
every column, keyed by column name. A body that is not valid UTF-8 is
written in base64, with `"body_encoding": "base64"` beside it. An
archive that does not exist yet, or was moved away, downloads with an
empty `archived_events`; the download never creates the file.
The download streams: each row is read and written out compressed
before the next is read, so neither the archive nor the JSON is held in
memory. It reads on a connection of its own, inside one read-only
transaction, so the file holds the archive as it stood when the
download started, and archive writes go on meanwhile, since under WAL a
reader never blocks a writer. While it runs, the `-wal` cannot be
checkpointed past what it reads, so a long download lets the `-wal`
grow. It finds the file by the stored names under the lock that webhook
edits, target edits and target creation hold, and lets go once the file
is open: a rename during the download moves the file without affecting
it.
Deleting a webhook releases its archives: the delivery engine's cached Deleting a webhook releases its archives: the delivery engine's cached
archive writers are dropped and their file handles closed, so nothing archive writers are dropped and their file handles closed, so nothing
lingers after the webhook is gone. The archive **files themselves are lingers after the webhook is gone. The archive **files themselves are
@@ -2611,11 +2634,10 @@ on all three arms of `Trace`, including the routine one an operator
reaches at `DEBUG`, which is the only level at which a successful reaches at `DEBUG`, which is the only level at which a successful
`INSERT` is written at all. One GORM path does not consult the filter — `INSERT` is written at all. One GORM path does not consult the filter —
`(*gorm.DB).Scan`, which records the statement through GORM's own trace `(*gorm.DB).Scan`, which records the statement through GORM's own trace
recorder. No production code path calls it; its one caller is recorder. No production code path calls it; only tests do, and what a
`internal/database/database_test.go:91`, whose `SELECT 1` binds test binds is fixture data. `internal/gormlog/scan_guard_test.go` fails
nothing, and `internal/gormlog/scan_guard_test.go` fails if a non-test if a non-test file calls it. `Pluck`, `Row` and `Raw` all run through
file calls it. `Pluck`, `Row` and `Raw` all run through the normal the normal callback processor and are filtered.
callback processor and are filtered.
See `#### What DEBUG=true exposes` under Configuration. See `#### What DEBUG=true exposes` under Configuration.
What that ceiling does **not** cover, stated here so the figure is not What that ceiling does **not** cover, stated here so the figure is not
@@ -2901,6 +2923,7 @@ returns to the page that was asked for.
| `POST` | `/hook/{id}/targets` | Add target to webhook | | `POST` | `/hook/{id}/targets` | Add target to webhook |
| `GET` | `/hook/{id}/targets/{targetID}/edit` | Edit target form. The one page that renders a target's destination URL and header values in full, rather than masked | | `GET` | `/hook/{id}/targets/{targetID}/edit` | Edit target form. The one page that renders a target's destination URL and header values in full, rather than masked |
| `POST` | `/hook/{id}/targets/{targetID}/edit` | Edit target submission | | `POST` | `/hook/{id}/targets/{targetID}/edit` | Edit target submission |
| `GET` | `/hook/{id}/targets/{targetID}/download` | Download a `database` target's archive as one gzipped JSON file. See [Database Architecture](#database-architecture) |
| `POST` | `/hook/{id}/targets/{targetID}/delete` | Delete a target | | `POST` | `/hook/{id}/targets/{targetID}/delete` | Delete a target |
| `POST` | `/hook/{id}/targets/{targetID}/toggle` | Enable or disable a target | | `POST` | `/hook/{id}/targets/{targetID}/toggle` | Enable or disable a target |
@@ -2982,6 +3005,7 @@ webhooker/
│ │ ├── target_slack.go # Slack/Mattermost incoming-webhook target │ │ ├── target_slack.go # Slack/Mattermost incoming-webhook target
│ │ ├── target_database.go # Database archive target │ │ ├── target_database.go # Database archive target
│ │ ├── target_database_archive.go # Archive file lifecycle and pruning │ │ ├── target_database_archive.go # Archive file lifecycle and pruning
│ │ ├── target_database_export.go # Archive download as gzipped JSON
│ │ ├── target_log.go # Log target (stdout) │ │ ├── target_log.go # Log target (stdout)
│ │ ├── target_config_view.go # Masked target config for templates │ │ ├── target_config_view.go # Masked target config for templates
│ │ ├── archive_sweeper.go # Periodic pruning of idle archives │ │ ├── archive_sweeper.go # Periodic pruning of idle archives
@@ -3248,9 +3272,9 @@ each hook. The order, read off the fx stop-hook log:
1. `ArchiveSweeper` 1. `ArchiveSweeper`
2. `RetentionReaper` 2. `RetentionReaper`
3. `server` — the HTTP drain, bounded separately by 3. `server` — the HTTP drain, bounded by `server.ShutdownTimeout`
`server.ShutdownTimeout` (**3 seconds**), then a Sentry flush if (**3 seconds**) and by what the hooks before it left, then a Sentry
`SENTRY_DSN` is set flush if `SENTRY_DSN` is set
4. `delivery.Engine` — waits for its workers, then closes the archive 4. `delivery.Engine` — waits for its workers, then closes the archive
databases databases
5. `healthcheck` 5. `healthcheck`
@@ -3270,23 +3294,30 @@ exhaust the sequence budget at the instant it finished, and every
later hook — the delivery engine, the healthcheck, the webhook DB later hook — the delivery engine, the healthcheck, the webhook DB
manager and the database close — would be skipped in exactly the manager and the database close — would be skipped in exactly the
case where the drain mattered. 3 seconds leaves 2 seconds case where the drain mattered. 3 seconds leaves 2 seconds
(`server.TailHookReserve`) for the tail, which is far more than the (`server.TailHookReserve`) for the tail. The reserve is that
microseconds it needs. remainder, not a figure sized to the tail, which takes about a
millisecond.
That reserve belongs to the tail hooks, not to the server hook, and That reserve belongs to the tail hooks, not to the server hook, and
the Sentry flush is what could take it: it runs after the drain the server hook could take it in two ways. The hooks before it may
**inside the same hook**, and `sentry.Flush` takes a bare duration already have spent part of the budget, so a full 3-second drain
and honours no context, so an unreachable Sentry endpoint would add would come out of the reserve; the drain is therefore also bounded
its own timeout on top of a full-length drain and consume the whole by whatever is left on the stop context minus the reserve. And the
sequence budget by itself. It is therefore clamped to whatever is Sentry flush runs after the drain **inside the same hook**, and
left on the stop context minus the reserve, and skipped when that `sentry.Flush` takes a bare duration and honours no context, so an
leaves too little to be worth attempting — so a full-length drain unreachable Sentry endpoint would add its own timeout on top of a
means Sentry events are dropped rather than the database close being full-length drain and consume the whole sequence budget by itself.
skipped. It is clamped the same way, and skipped when that leaves too little
to be worth attempting — so a full-length drain means Sentry events
are dropped rather than the database close being skipped.
This does not make the database close unconditional: a wedged This does not make the database close unconditional. A slow
`ArchiveSweeper` or `RetentionReaper` still runs first and can `ArchiveSweeper` or `RetentionReaper` is enough to cut the shutdown
consume the whole budget on its own. short, not only one that consumes the whole budget: what they spend
comes out of the drain first, so after 2 seconds of theirs a request
still in flight gets 1 second to finish, and after 3 it gets none.
Past 3 seconds they spend the reserve itself, and one that takes the
whole budget skips every hook after it, the database close included.
The value is chosen to sit inside the container stop grace period. The value is chosen to sit inside the container stop grace period.
Docker's default `docker stop` grace is 10 seconds and the Dockerfile Docker's default `docker stop` grace is 10 seconds and the Dockerfile
+9 -7
View File
@@ -38,17 +38,19 @@ import (
// hook that used the whole budget would exhaust it at that instant, // hook that used the whole budget would exhaust it at that instant,
// and fx would skip every hook after the server — the delivery // and fx would skip every hook after the server — the delivery
// engine, the healthcheck, the webhook DB manager and the database // engine, the healthcheck, the webhook DB manager and the database
// close. That hook is the 3s HTTP drain plus the Sentry flush that // close. That hook is the HTTP drain plus the Sentry flush that
// follows it in the same hook, so the flush is clamped to the stop // follows it in the same hook, and each is clamped to the stop
// context's remaining time less server.TailHookReserve rather than // context's remaining time less server.TailHookReserve rather than
// running for its own fixed 2s; the reserve is what the tail hooks // running for its own fixed 3s and 2s; the reserve is what the tail
// live on, and they are microsecond-scale in normal operation. // hooks live on, and they are microsecond-scale in normal operation.
// TestStopTimeout_LeavesHeadroomForTailHooks pins the arithmetic // TestStopTimeout_LeavesHeadroomForTailHooks pins the arithmetic
// across every drain length. // across every drain length and every amount of budget the hooks
// before the server may already have spent.
// //
// This does not make the database close unconditional: the // This does not make the database close unconditional: the
// ArchiveSweeper and RetentionReaper hooks run before the server // ArchiveSweeper and RetentionReaper hooks run before the server.
// and can still consume the whole budget on their own. // What they spend comes out of the drain first, but past 3s it comes
// out of the reserve, and they can consume the whole budget.
const stopTimeout = 5 * time.Second const stopTimeout = 5 * time.Second
// exitUsage is the status for a command line this binary cannot make // exitUsage is the status for a command line this binary cannot make
+27 -9
View File
@@ -252,22 +252,40 @@ const tailHeadroom = 2 * time.Second
// can produce, since a shorter drain leaves the flush more room and // can produce, since a shorter drain leaves the flush more room and
// the worst case is not necessarily at either extreme. // the worst case is not necessarily at either extreme.
// //
// Shrinking either budget, or unbounding the flush again, must fail // Nor does the hook start on a full budget: the ArchiveSweeper and
// here rather than silently recreating a hook that swallows the // RetentionReaper hooks run before it, and whatever they spent is
// whole sequence. // gone. The outer sweep walks every amount they can spend. Once they
// have eaten into the headroom themselves, the hook must spend
// nothing of what is left. A drain that starts on the full budget
// must still get all of ShutdownTimeout, so a smaller stopTimeout
// cannot silently shorten every drain.
//
// Shrinking either budget, or unbounding the drain or the flush
// again, must fail here rather than silently recreating a hook that
// swallows the whole sequence.
func TestStopTimeout_LeavesHeadroomForTailHooks(t *testing.T) { func TestStopTimeout_LeavesHeadroomForTailHooks(t *testing.T) {
t.Parallel() t.Parallel()
require.Less(t, server.ShutdownTimeout, stopTimeout) require.Less(t, server.ShutdownTimeout, stopTimeout)
require.Equal(
t, server.ShutdownTimeout, server.DrainBudget(stopTimeout),
"a drain that starts on the full stop budget is cut short",
)
const step = 10 * time.Millisecond const step = 10 * time.Millisecond
for drain := time.Duration(0); drain <= server.ShutdownTimeout; drain += step { for spent := time.Duration(0); spent <= stopTimeout; spent += step {
hook := drain + server.SentryFlushBudget(stopTimeout-drain) remaining := stopTimeout - spent
longest := max(server.DrainBudget(remaining), 0)
require.LessOrEqual( for drain := time.Duration(0); drain <= longest; drain += step {
t, hook+tailHeadroom, stopTimeout, hook := drain + server.SentryFlushBudget(remaining-drain)
"a %s drain leaves the tail hooks short", drain,
) require.GreaterOrEqual(
t, remaining-hook, min(remaining, tailHeadroom),
"a %s drain after %s of earlier hooks leaves "+
"the tail hooks short", drain, spent,
)
}
} }
} }
+2 -2
View File
@@ -32,8 +32,8 @@ type Event struct {
ContentType string `json:"contentType"` ContentType string `json:"contentType"`
// BodyBytes is the size of Body in bytes, recorded when the event // BodyBytes is the size of Body in bytes, recorded when the event
// is stored so the recent events list can show it without reading // is stored, so that the recent events list, which reads only the
// the body. // start of each body, knows the whole body's size.
BodyBytes int64 `gorm:"not null" json:"bodyBytes"` BodyBytes int64 `gorm:"not null" json:"bodyBytes"`
// ResubmittedFromID names the event this one was copied from by // ResubmittedFromID names the event this one was copied from by
+7 -3
View File
@@ -22,6 +22,7 @@ import (
_ "modernc.org/sqlite" // Pure Go SQLite driver. _ "modernc.org/sqlite" // Pure Go SQLite driver.
"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/gormlog"
) )
const ( const (
@@ -70,7 +71,8 @@ func setupArchiveTest(t *testing.T) *archiveEnv {
t.Cleanup(func() { _ = sqlDB.Close() }) t.Cleanup(func() { _ = sqlDB.Close() })
gdb, err := gorm.Open( gdb, err := gorm.Open(
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, sqlite.Dialector{Conn: sqlDB},
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
) )
require.NoError(t, err) require.NoError(t, err)
@@ -168,7 +170,8 @@ func (env *archiveEnv) seedArchiveRows(
require.NoError(t, err) require.NoError(t, err)
gdb, err := gorm.Open( gdb, err := gorm.Open(
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, sqlite.Dialector{Conn: sqlDB},
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
) )
require.NoError(t, err) require.NoError(t, err)
@@ -227,7 +230,8 @@ func countArchivedRows(path string) (int64, error) {
defer func() { _ = sqlDB.Close() }() defer func() { _ = sqlDB.Close() }()
gdb, err := gorm.Open( gdb, err := gorm.Open(
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, sqlite.Dialector{Conn: sqlDB},
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
) )
if err != nil { if err != nil {
return 0, err return 0, err
+29 -1
View File
@@ -23,6 +23,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/gormlog"
) )
// iSetup holds common integration test dependencies. // iSetup holds common integration test dependencies.
@@ -80,7 +81,8 @@ func iMainDB(t *testing.T) *gorm.DB {
t.Cleanup(func() { _ = sqlDB.Close() }) t.Cleanup(func() { _ = sqlDB.Close() })
db, err := gorm.Open( db, err := gorm.Open(
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, sqlite.Dialector{Conn: sqlDB},
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
) )
require.NoError(t, err) require.NoError(t, err)
@@ -1425,6 +1427,32 @@ func TestDeliverHTTP_InvalidConfig(t *testing.T) {
) )
} }
// TestDeliverHTTP_InvalidConfigUnrecordedStaysPending: a delivery is
// failed for an invalid config only once the reason is recorded.
// Unrecorded, it stays pending, where the sweep finds it again.
func TestDeliverHTTP_InvalidConfigUnrecordedStaysPending(t *testing.T) {
t.Parallel()
db := testWebhookDB(t)
e := testEngine(t, 1)
event, del := iSeedEventAndDelivery(
t, db, `{"config":"invalid"}`, "",
)
task, d := iHTTPTaskAndDelivery(
event, del, "bad-config", `not-json`, 0, 1,
)
require.NoError(t, db.Exec("drop table delivery_results").Error)
e.ExportDeliverHTTP(context.TODO(), db, d, task)
iAssertStatus(t, db, del.ID,
database.DeliveryStatusPending,
)
}
// --- Notify batching --- // --- Notify batching ---
func TestNotify_MultipleTasks(t *testing.T) { func TestNotify_MultipleTasks(t *testing.T) {
+74 -1
View File
@@ -5,6 +5,7 @@ import (
"context" "context"
"encoding/json" "encoding/json"
"fmt" "fmt"
"io"
"log/slog" "log/slog"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
@@ -25,6 +26,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/gormlog"
"sneak.berlin/go/webhooker/internal/metrics" "sneak.berlin/go/webhooker/internal/metrics"
) )
@@ -49,7 +51,8 @@ func testWebhookDB(t *testing.T) *gorm.DB {
t.Cleanup(func() { _ = sqlDB.Close() }) t.Cleanup(func() { _ = sqlDB.Close() })
db, err := gorm.Open( db, err := gorm.Open(
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, sqlite.Dialector{Conn: sqlDB},
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
) )
require.NoError(t, err) require.NoError(t, err)
@@ -1056,6 +1059,21 @@ func TestParseHTTPConfig_MissingURL(t *testing.T) {
) )
} }
func TestParseHTTPConfig_Undecodable(t *testing.T) {
t.Parallel()
e := testEngine(t, 1)
_, err := e.ExportParseHTTPConfig(
`{"url":"https://example.com/hook","timeout":"soon"}`,
)
assert.Error(t, err,
"config that does not decode should return error, "+
"even when the part that did names a URL",
)
}
func TestScheduleRetry_SendsToRetryChannel( func TestScheduleRetry_SendsToRetryChannel(
t *testing.T, t *testing.T,
) { ) {
@@ -1241,6 +1259,33 @@ func TestDoHTTPRequest_ForwardsHeaders(t *testing.T) {
) )
} }
// A response that ends before the length it announced is an error, not
// a short body.
func TestDoHTTPRequest_CutShortResponseIsAnError(t *testing.T) {
t.Parallel()
ts := httptest.NewServer(
http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Length", "100")
_, _ = w.Write([]byte("cut short"))
},
),
)
defer ts.Close()
e := testEngine(t, 1)
_, body, _, err := e.ExportDoHTTPRequest(
context.TODO(),
&delivery.HTTPTargetConfig{URL: ts.URL},
&database.Event{},
)
require.ErrorIs(t, err, io.ErrUnexpectedEOF)
assert.Empty(t, body)
}
// The event's stored inbound headers carry the same Content-Type the // The event's stored inbound headers carry the same Content-Type the
// receiver saved as the event's ContentType, so a delivery could send // 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 // it twice. It must go out exactly once, with a Content-Type configured
@@ -1317,6 +1362,34 @@ func TestApplyRequestHeaders_SendsOneContentType(t *testing.T) {
} }
} }
// Stored inbound headers that do not decode forward nothing, not the
// part of them that happened to decode.
func TestApplyRequestHeaders_UndecodableInboundForwardsNothing(
t *testing.T,
) {
t.Parallel()
req, err := http.NewRequestWithContext(
context.Background(),
http.MethodPost,
"https://target.example.com/hook",
http.NoBody,
)
require.NoError(t, err)
names := delivery.ExportApplyRequestHeaders(
req,
&database.Event{
Headers: `{"X-Custom":["value1"],"X-Broken":"not a list"}`,
},
&delivery.HTTPTargetConfig{},
"webhooker/dev",
)
assert.Empty(t, names)
assert.Empty(t, req.Header.Get("X-Custom"))
}
func TestProcessDelivery_RoutesToCorrectHandler( func TestProcessDelivery_RoutesToCorrectHandler(
t *testing.T, t *testing.T,
) { ) {
@@ -376,3 +376,97 @@ func TestFailedResultWriteLeavesDeliveryRecoverable(
database.DeliveryStatusPending, database.DeliveryStatusPending,
) )
} }
// TestFailedResultWriteWithRetriesLeavesDeliveryRecoverable is the same
// rule for a target with retries: whatever the receiver answered, the
// delivery stays pending and no retry is scheduled. The circuit breaker
// still learns the answer, because it describes the target's health,
// not the database's.
func TestFailedResultWriteWithRetriesLeavesDeliveryRecoverable(
t *testing.T,
) {
t.Parallel()
// The "send succeeded" case starts with the breaker tripped open,
// so the delivery goes out as its probe and only a recorded
// success closes it again.
tests := []struct {
name string
answer int
tripped bool
wantBreaker delivery.CircuitState
}{
{"send succeeded", http.StatusOK, true, delivery.CircuitClosed},
{"send failed", http.StatusBadGateway, false, delivery.CircuitOpen},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
ts := httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(tc.answer)
},
))
defer ts.Close()
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"unwritable":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
require.NoError(
t,
s.WebhookDB.Exec("drop table delivery_results").Error,
)
// A single failure opens this breaker, and with no
// cooldown an open breaker lets the next delivery
// through as a probe.
cb := delivery.NewTestCircuitBreaker(1, 0)
if tc.tripped {
cb.RecordFailure()
}
s.Engine.ExportSetCircuitBreaker(targetID, cb)
full := &database.Delivery{
EventID: event.ID,
TargetID: targetID,
Status: database.DeliveryStatusPending,
Event: event,
Target: database.Target{
Name: "unwritable",
Type: database.TargetTypeHTTP,
Config: iHTTPConfig(ts.URL),
MaxRetries: 3,
},
}
full.ID = d.ID
sched := &recordingScheduler{}
s.Engine.ExportDeliverHTTPWithScheduler(
context.Background(), s.WebhookDB, full,
&delivery.Task{
DeliveryID: d.ID,
TargetID: targetID,
AttemptNum: 1,
},
sched,
)
iAssertStatus(t, s.WebhookDB, d.ID, database.DeliveryStatusPending)
assert.Empty(t, sched.delays, "no retry may be scheduled")
assert.Equal(t, tc.wantBreaker, cb.State())
})
}
}
+6 -10
View File
@@ -3,7 +3,6 @@ package delivery
import ( import (
"context" "context"
"fmt" "fmt"
"path/filepath"
"strings" "strings"
"sync" "sync"
"time" "time"
@@ -277,10 +276,9 @@ func (t *databaseTarget) releaseSweepWriter(
} }
// newWriter builds the writer for a database target's archive. The // newWriter builds the writer for a database target's archive. The
// file lives beside the webhook's event database in the data // file is the one ArchivePath gives for the webhook and the target as
// directory and is named for the webhook and the target as the main // the main database names them now; from then on only rename changes
// database has them now; from then on only rename changes the name // the name the writer uses. It does not touch the archive file.
// the writer uses. It does not touch the archive file.
func (t *databaseTarget) newWriter( func (t *databaseTarget) newWriter(
targetID string, targetID string,
) (*archiveWriter, error) { ) (*archiveWriter, error) {
@@ -299,12 +297,10 @@ func (t *databaseTarget) newWriter(
) )
} }
dir := filepath.Dir(t.eng.dbManager.DBPath(target.WebhookID)) w := newArchiveWriter(
name := ArchiveFileName( ArchivePath(t.eng.dbManager, &target.Webhook, &target),
target.Webhook.Name, target.Name, target.ID, t.eng.log,
) )
w := newArchiveWriter(filepath.Join(dir, name), t.eng.log)
w.webhookID = target.WebhookID w.webhookID = target.WebhookID
return w, nil return w, nil
+275
View File
@@ -0,0 +1,275 @@
package delivery
import (
"compress/gzip"
"context"
"database/sql"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"log/slog"
"path/filepath"
"time"
"unicode/utf8"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/gormlog"
)
// archiveTableQuery counts the archive's table: 0 when the archive
// writer has created the file but not yet the table in it.
const archiveTableQuery = "SELECT count(*) FROM sqlite_master " +
"WHERE type = 'table' AND name = 'archived_events'"
// ArchivePath returns where a database target's archive file is: in
// the data directory, beside the webhook's event database, under the
// name ArchiveFileName gives it.
func ArchivePath(
dbMgr *database.WebhookDBManager,
webhook *database.Webhook,
target *database.Target,
) string {
return filepath.Join(
filepath.Dir(dbMgr.DBPath(webhook.ID)),
ArchiveFileName(webhook.Name, target.Name, target.ID),
)
}
// ArchiveExportFileName returns the name a database target's archive
// downloads under:
// archive-WEBHOOKNAME-TARGETNAME-YYYYMMDDTHHMMSSZ.json.gz, the names
// made safe as in ArchiveFileName and the time in UTC.
func ArchiveExportFileName(
webhookName, targetName string, at time.Time,
) string {
return "archive-" + archiveNamePart(webhookName) + "-" +
archiveNamePart(targetName) + "-" +
at.UTC().Format("20060102T150405Z") + ".json.gz"
}
// ArchiveExport is a database target's archive opened for download.
// It reads the file on its own connection, inside one read-only
// transaction, so it writes out the archive as it stood when
// OpenArchiveExport returned.
//
// Archives are in WAL mode, where a reader works from a snapshot and
// never blocks a writer: archive writes go on while an export is open,
// and the export does not see them. SQLite cannot checkpoint the -wal
// past an open snapshot, so the -wal grows until the export is closed.
type ArchiveExport struct {
db *sql.DB
tx *gorm.DB
// empty is true when there is nothing to read: no file, or a file
// without the archive's table yet.
empty bool
}
// exportedName is how an export names its webhook and its target.
type exportedName struct {
ID string `json:"id"`
Name string `json:"name"`
}
// OpenArchiveExport opens the archive file at path for export and
// takes the snapshot the export reads. It never creates the file: with
// no file at path, the export has no rows.
//
// Once it has returned, the file is open, so a rename or a move of it
// does not affect the export, which reads the same file under its new
// name.
//
// The transaction lasts as long as ctx does, so ctx must last for the
// whole export.
func OpenArchiveExport(
ctx context.Context, path string, log *slog.Logger,
) (*ArchiveExport, error) {
if !fileExists(path) {
return &ArchiveExport{empty: true}, nil
}
db, err := database.OpenSQLite(path, archiveModeExisting)
if err != nil {
return nil, fmt.Errorf("opening archive %s: %w", path, err)
}
gdb, err := gorm.Open(
sqlite.Dialector{Conn: db}, &gorm.Config{
// Never leave this at GORM's default. See
// internal/gormlog.
Logger: gormlog.New(log),
},
)
if err != nil {
_ = db.Close()
return nil, fmt.Errorf("opening archive %s: %w", path, err)
}
// ReadOnly makes the driver begin a deferred transaction in place
// of the BEGIN IMMEDIATE the connection string asks for, so the
// export never takes the archive's write lock.
tx := gdb.WithContext(ctx).Begin(&sql.TxOptions{ReadOnly: true})
if tx.Error != nil {
_ = db.Close()
return nil, fmt.Errorf("reading archive %s: %w", path, tx.Error)
}
// The transaction's first read is what takes the snapshot.
var tables int
err = tx.Raw(archiveTableQuery).Row().Scan(&tables)
if err != nil {
_ = tx.Rollback()
_ = db.Close()
return nil, fmt.Errorf("reading archive %s: %w", path, err)
}
return &ArchiveExport{db: db, tx: tx, empty: tables == 0}, nil
}
// WriteGzipJSON writes the export to w as one gzipped JSON object:
// webhook and target, each an id and a name; exported_at; and
// archived_events, one object per archived row, keyed by column name.
// A body that is not valid UTF-8 cannot be a JSON string, so it is
// written in base64, with "body_encoding": "base64" beside it.
//
// Each row is written out before the next is read, so neither the
// archive nor its JSON is ever held in memory whole. After an error
// the gzip stream is left unfinished, so what was written does not
// decompress as a whole file.
func (x *ArchiveExport) WriteGzipJSON(
ctx context.Context,
w io.Writer,
webhook *database.Webhook,
target *database.Target,
exportedAt time.Time,
) error {
head, err := json.Marshal(map[string]any{
"webhook": exportedName{ID: webhook.ID, Name: webhook.Name},
"target": exportedName{ID: target.ID, Name: target.Name},
"exported_at": exportedAt.UTC(),
})
if err != nil {
return fmt.Errorf("encoding archive export: %w", err)
}
zw := gzip.NewWriter(w)
err = x.writeJSON(ctx, zw, head)
if err != nil {
return fmt.Errorf("writing archive export: %w", err)
}
return zw.Close()
}
// Close ends the export's transaction and closes its connection.
func (x *ArchiveExport) Close() error {
if x.db == nil {
return nil
}
_ = x.tx.Rollback()
return x.db.Close()
}
// writeJSON writes head with archived_events added as its last key,
// the rows going into it one at a time.
func (x *ArchiveExport) writeJSON(
ctx context.Context, w io.Writer, head []byte,
) error {
// head goes out without its closing brace, so that
// archived_events can follow it.
_, err := w.Write(head[:len(head)-1])
if err != nil {
return err
}
_, err = io.WriteString(w, `,"archived_events":[`)
if err != nil {
return err
}
err = x.writeRows(ctx, w)
if err != nil {
return err
}
_, err = io.WriteString(w, "\n]}\n")
return err
}
// writeRows writes each archived row to w, oldest first, one per line,
// separated by commas.
func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error {
if x.empty {
return nil
}
rows, err := x.tx.WithContext(ctx).
Model(&archivedEvent{}).Order("id").Rows()
if err != nil {
return err
}
defer func() { _ = rows.Close() }()
for sep := "\n"; rows.Next(); sep = ",\n" {
var ev archivedEvent
err = x.tx.ScanRows(rows, &ev)
if err != nil {
return err
}
_, err = io.WriteString(w, sep)
if err != nil {
return err
}
err = writeRow(w, &ev)
if err != nil {
return err
}
}
return rows.Err()
}
// writeRow writes an archived row to w as a JSON object keyed by
// column name, its body in base64 when it is not valid UTF-8.
func writeRow(w io.Writer, ev *archivedEvent) error {
row := map[string]any{
"id": ev.ID,
"event_id": ev.EventID,
"webhook_id": ev.WebhookID,
"entrypoint_id": ev.EntrypointID,
"method": ev.Method,
"headers": ev.Headers,
"body": ev.Body,
"content_type": ev.ContentType,
"archived_at": ev.ArchivedAt.UTC(),
}
if !utf8.ValidString(ev.Body) {
row["body"] = base64.StdEncoding.EncodeToString([]byte(ev.Body))
row["body_encoding"] = "base64"
}
line, err := json.Marshal(row)
if err != nil {
return err
}
_, err = w.Write(line)
return err
}
@@ -0,0 +1,412 @@
package delivery_test
import (
"bufio"
"bytes"
"compress/gzip"
"crypto/rand"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"os"
"path/filepath"
"runtime"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
)
// The webhook and the target the export tests' archives belong to.
const (
exportWebhookID = "wh-export"
exportWebhookName = "Orders (EU)"
exportTargetID = "tgt-export"
exportTargetName = "Long-term archive"
)
const (
// binaryBody is a body that is not valid UTF-8.
binaryBody = "\xff\xfe\x00\x01binary\x80"
// openedEventID is the event the snapshot tests archive before
// they open the export.
openedEventID = "opened"
)
// writeExportTo writes export to w as the archive of the export tests'
// webhook and target, exported at 2026-10-02T12:03:04Z.
func writeExportTo(
t *testing.T, export *delivery.ArchiveExport, w io.Writer,
) error {
t.Helper()
return export.WriteGzipJSON(
t.Context(), w,
&database.Webhook{
BaseModel: database.BaseModel{ID: exportWebhookID},
Name: exportWebhookName,
},
&database.Target{
BaseModel: database.BaseModel{ID: exportTargetID},
Name: exportTargetName,
},
time.Date(2026, 10, 2, 12, 3, 4, 0, time.UTC),
)
}
// exportArchive runs a whole export of the archive at path and returns
// its JSON, decompressed and parsed.
func exportArchive(t *testing.T, path string) map[string]any {
t.Helper()
export, err := delivery.OpenArchiveExport(
t.Context(), path, archiveTestLogger(),
)
require.NoError(t, err)
defer func() { require.NoError(t, export.Close()) }()
return writeExport(t, export)
}
// writeExport writes an opened export and returns its JSON,
// decompressed and parsed. Reading to the end makes the gzip reader
// check that the stream was finished.
func writeExport(
t *testing.T, export *delivery.ArchiveExport,
) map[string]any {
t.Helper()
var buf bytes.Buffer
require.NoError(t, writeExportTo(t, export, &buf))
zr, err := gzip.NewReader(&buf)
require.NoError(t, err)
raw, err := io.ReadAll(zr)
require.NoError(t, err)
var got map[string]any
require.NoError(t, json.Unmarshal(raw, &got))
return got
}
// exportedEvents returns an export's archived_events.
func exportedEvents(t *testing.T, got map[string]any) []map[string]any {
t.Helper()
list, ok := got["archived_events"].([]any)
require.True(t, ok, "archived_events must be an array: %v", got)
events := make([]map[string]any, len(list))
for i, v := range list {
events[i], ok = v.(map[string]any)
require.True(t, ok, "an archived event must be an object: %v", v)
}
return events
}
// exportedEventIDs returns the event_id of each of an export's
// archived_events.
func exportedEventIDs(t *testing.T, got map[string]any) []string {
t.Helper()
events := exportedEvents(t, got)
ids := make([]string, 0, len(events))
for _, ev := range events {
ids = append(ids, fmt.Sprint(ev["event_id"]))
}
return ids
}
// TestArchiveExport_MatchesStoredRows proves an export holds the
// webhook, the target, the time, and every column of every stored
// row: a body that is valid UTF-8 as a string, and one that is not in
// base64, marked as such.
func TestArchiveExport_MatchesStoredRows(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "archive.db")
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
bodies := []string{`{"order":1}`, "plain text", "", binaryBody}
for i, body := range bodies {
require.NoError(t, w.Write(delivery.ExportArchivedEvent{
EventID: fmt.Sprintf("ev-%d", i),
WebhookID: exportWebhookID,
EntrypointID: "ep-1",
Method: "POST",
Headers: `{"X-Test":["yes"]}`,
Body: body,
ContentType: testContentType,
}, 0))
}
var stored []delivery.ExportArchivedEvent
require.NoError(t, openArchiveDBForRead(t, path).
Order("id").Find(&stored).Error)
got := exportArchive(t, path)
assert.Equal(t,
map[string]any{"id": exportWebhookID, "name": exportWebhookName},
got["webhook"],
)
assert.Equal(t,
map[string]any{"id": exportTargetID, "name": exportTargetName},
got["target"],
)
assert.Equal(t, "2026-10-02T12:03:04Z", got["exported_at"])
events := exportedEvents(t, got)
require.Len(t, events, len(bodies))
for i, row := range stored {
assertExportedRow(t, row, events[i])
}
}
// assertExportedRow checks that ev, from an export, holds every column
// of the stored row.
func assertExportedRow(
t *testing.T, row delivery.ExportArchivedEvent, ev map[string]any,
) {
t.Helper()
archivedAt, err := time.Parse(
time.RFC3339Nano, fmt.Sprint(ev["archived_at"]),
)
require.NoError(t, err)
assert.True(t, archivedAt.Equal(row.ArchivedAt))
assert.EqualValues(t, row.ID, ev["id"])
assert.Equal(t, row.EventID, ev["event_id"])
assert.Equal(t, row.WebhookID, ev["webhook_id"])
assert.Equal(t, row.EntrypointID, ev["entrypoint_id"])
assert.Equal(t, row.Method, ev["method"])
assert.Equal(t, row.Headers, ev["headers"])
assert.Equal(t, row.ContentType, ev["content_type"])
if row.Body != binaryBody {
assert.Equal(t, row.Body, ev["body"])
assert.Len(t, ev, 9, "the nine columns and nothing else: %v", ev)
return
}
body, err := base64.StdEncoding.DecodeString(fmt.Sprint(ev["body"]))
require.NoError(t, err)
assert.Equal(t, binaryBody, string(body))
assert.Equal(t, "base64", ev["body_encoding"])
assert.Len(t, ev, 10, "the nine columns and body_encoding: %v", ev)
}
// TestArchiveExport_Empty proves an archive with nothing in it exports
// as an empty archived_events: no file, which the export must not
// create; a file the archive writer has not yet put its table in; and
// a table with no rows.
func TestArchiveExport_Empty(t *testing.T) {
t.Parallel()
dir := t.TempDir()
missing := filepath.Join(dir, "missing.db")
noTable := filepath.Join(dir, "no-table.db")
noRows := filepath.Join(dir, "no-rows.db")
require.NoError(t, os.WriteFile(noTable, nil, 0o600))
require.NoError(t,
delivery.NewExportArchiveWriter(noRows, archiveTestLogger(), 0).
Open(0),
)
for _, path := range []string{missing, noTable, noRows} {
assert.Empty(t, exportedEvents(t, exportArchive(t, path)), path)
}
for _, suffix := range archiveFileSuffixes() {
assert.NoFileExists(t, missing+suffix)
}
}
// TestArchiveExport_ReadsOneSnapshot proves an export writes the
// archive as it was when it was opened, and holds up no archive
// write: a row written while the export is open is stored, and is not
// in the export. A write held up for the whole busy timeout would
// fail.
func TestArchiveExport_ReadsOneSnapshot(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "archive.db")
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: openedEventID}, 0))
export, err := delivery.OpenArchiveExport(
t.Context(), path, archiveTestLogger(),
)
require.NoError(t, err)
defer func() { require.NoError(t, export.Close()) }()
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "during"}, 0))
assert.Equal(t,
[]string{openedEventID}, exportedEventIDs(t, writeExport(t, export)),
)
var stored int64
require.NoError(t, openArchiveDBForRead(t, path).
Model(&delivery.ExportArchivedEvent{}).Count(&stored).Error)
assert.Equal(t, int64(2), stored)
}
// TestArchiveExport_SurvivesRename proves that renaming the archive
// while an export of it is open, as renaming its webhook or target
// does, leaves the export reading the same file.
func TestArchiveExport_SurvivesRename(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "archive-old.db")
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: openedEventID}, 0))
export, err := delivery.OpenArchiveExport(
t.Context(), path, archiveTestLogger(),
)
require.NoError(t, err)
defer func() { require.NoError(t, export.Close()) }()
require.NoError(t, w.Rename("archive-new.db"))
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "after"}, 0))
require.NoFileExists(t, path)
assert.Equal(t,
[]string{openedEventID}, exportedEventIDs(t, writeExport(t, export)),
)
}
// heapPeak is an io.Writer that discards what it is given and records
// the largest heap it saw at a write. It collects garbage before each
// reading, so the heap it reads is what is still held.
type heapPeak struct {
max uint64
}
func (p *heapPeak) Write(b []byte) (int, error) {
var m runtime.MemStats
runtime.GC()
runtime.ReadMemStats(&m)
p.max = max(p.max, m.HeapAlloc)
return len(b), nil
}
// exportHeapGrowth exports an archive of rows random bodies, each
// bodySize bytes of base64, and returns how far the heap rose above
// where it stood when the export began, at its highest.
func exportHeapGrowth(t *testing.T, rows, bodySize int) uint64 {
t.Helper()
path := filepath.Join(t.TempDir(), "archive.db")
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
// Base64 makes four characters of every three bytes.
random := make([]byte, bodySize/4*3)
for range rows {
_, _ = rand.Read(random)
require.NoError(t, w.Write(delivery.ExportArchivedEvent{
Body: base64.StdEncoding.EncodeToString(random),
}, 0))
}
export, err := delivery.OpenArchiveExport(
t.Context(), path, archiveTestLogger(),
)
require.NoError(t, err)
defer func() { require.NoError(t, export.Close()) }()
runtime.GC()
var start runtime.MemStats
runtime.ReadMemStats(&start)
// Through a buffer, the heap is read once per 8 KiB of output
// rather than at each of gzip's small writes, which takes far
// longer.
peak := &heapPeak{max: start.HeapAlloc}
buffered := bufio.NewWriterSize(peak, 8<<10)
require.NoError(t, writeExportTo(t, export, buffered))
require.NoError(t, buffered.Flush())
return peak.max - start.HeapAlloc
}
// TestArchiveExport_Streams proves an export holds neither the archive
// nor its output in memory whole: exporting 384 KiB more of archive
// raises the heap's peak by less than half of that. The export's own
// memory, mostly gzip's compressor, is the same for both archives, so
// it cancels out. The bodies are random bytes in base64, which gzip
// shrinks by only a quarter, so an export that read every row before
// writing, or built the JSON or the gzipped file before writing it,
// would raise the peak by at least three quarters of the difference.
//
// The smaller archive has two rows so that its export, too, writes
// out more than the 8 KiB buffer in exportHeapGrowth before it ends:
// the heap must be read while the export's own memory is held.
//
//nolint:paralleltest // It measures the heap, which tests share.
func TestArchiveExport_Streams(t *testing.T) {
const (
bodySize = 16 << 10
smallRows = 2
largeRows = smallRows + 24
limit = (largeRows - smallRows) * bodySize / 2
)
small := exportHeapGrowth(t, smallRows, bodySize)
large := exportHeapGrowth(t, largeRows, bodySize)
assert.Less(t, large, small+limit,
"the heap rose by %d for %d rows and by %d for %d rows",
small, smallRows, large, largeRows,
)
}
// TestArchiveExportFileName proves the download is named for the
// webhook and the target, with the names made safe as for the archive
// file, and the export time in UTC.
func TestArchiveExportFileName(t *testing.T) {
t.Parallel()
cest := time.FixedZone("CEST", int((2 * time.Hour).Seconds()))
assert.Equal(t,
"archive-orders-eu-long-term-archive-20261002T120304Z.json.gz",
delivery.ArchiveExportFileName(
exportWebhookName, exportTargetName,
time.Date(2026, 10, 2, 14, 3, 4, 0, cest),
),
)
}
+3 -1
View File
@@ -17,6 +17,7 @@ import (
_ "modernc.org/sqlite" // Pure Go SQLite driver. _ "modernc.org/sqlite" // Pure Go SQLite driver.
"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/gormlog"
) )
func archiveTestLogger() *slog.Logger { func archiveTestLogger() *slog.Logger {
@@ -42,7 +43,8 @@ func openArchiveDBForRead(
t.Cleanup(func() { _ = sqlDB.Close() }) t.Cleanup(func() { _ = sqlDB.Close() })
gdb, err := gorm.Open( gdb, err := gorm.Open(
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, sqlite.Dialector{Conn: sqlDB},
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
) )
require.NoError(t, err) require.NoError(t, err)
+21
View File
@@ -179,6 +179,27 @@ func TestDoHTTPRequest_TransportErrorMasksURL(t *testing.T) {
) )
} }
// TestDoHTTPRequest_UnparsableURLIsMasked is the same for an HTTP
// target URL that no request can be built from.
func TestDoHTTPRequest_UnparsableURLIsMasked(t *testing.T) {
t.Parallel()
e := testEngine(t, 1)
statusCode, _, _, reqErr := e.ExportDoHTTPRequest(
context.TODO(),
&delivery.HTTPTargetConfig{
URL: "https://hooks.example.com" + maskSecretPath + "\n",
},
&database.Event{},
)
require.Error(t, reqErr)
assert.Zero(t, statusCode)
assertNoCredential(t, reqErr.Error())
assert.Contains(t, reqErr.Error(), "invalid control character")
}
// TestValidateTargetURL_UnparsableURLIsMasked proves the SSRF // TestValidateTargetURL_UnparsableURLIsMasked proves the SSRF
// validator's error does not carry the submitted URL, which // validator's error does not carry the submitted URL, which
// the handler both logs and shows. // the handler both logs and shows.
+3 -3
View File
@@ -111,9 +111,9 @@ func (l *Logger) LogMode(gormlogger.LogLevel) gormlogger.Interface {
// //
// One GORM path does not consult this: (*gorm.DB).Scan records the // One GORM path does not consult this: (*gorm.DB).Scan records the
// statement through gorm's own traceRecorder, which does not implement // statement through gorm's own traceRecorder, which does not implement
// this interface. No production code path calls it; its one caller is // this interface. No production code path calls it; only tests do, and
// internal/database/database_test.go:91, whose SELECT 1 binds nothing. // what a test binds is fixture data. scan_guard_test.go fails if a
// scan_guard_test.go fails if a non-test file calls it. // non-test file calls it.
// (*gorm.DB).Pluck, Row and Raw all run through the normal callback // (*gorm.DB).Pluck, Row and Raw all run through the normal callback
// processor and are filtered. // processor and are filtered.
func (l *Logger) ParamsFilter( func (l *Logger) ParamsFilter(
+99 -40
View File
@@ -14,18 +14,16 @@ import (
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
) )
// minNonTestFiles guards the walk below against passing because it // isRowProducer reports whether name is GORM's Row or database/sql's
// found nothing to look at. The tree held 60 non-test .go files when // QueryRow or QueryRowContext, which return a *sql.Row whose Scan is
// this was written. // database/sql's and not (*gorm.DB).Scan. GORM's Rows is not listed:
const minNonTestFiles = 40 // it also returns an error, so Scan is never called on its result
// directly. It matches the method name only and resolves no types, so
// isRowProducer reports whether name is a method that returns a // a repo-local method with one of these names that returns *gorm.DB
// database/sql row handle. GORM's Row and Rows return *sql.Row and // gets past it: Scan on that method's result is not reported.
// *sql.Rows, so Scan on the result of one of them is database/sql's
// Scan and never (*gorm.DB).Scan.
func isRowProducer(name string) bool { func isRowProducer(name string) bool {
switch name { switch name {
case "Row", "Rows", "QueryRow", "QueryRowContext": case "Row", "QueryRow", "QueryRowContext":
return true return true
default: default:
return false return false
@@ -50,9 +48,14 @@ func receiverIsRowHandle(x ast.Expr) bool {
} }
// unguardedScans returns the position of every Scan call in file whose // unguardedScans returns the position of every Scan call in file whose
// receiver is not a row handle. It fails closed: a receiver it cannot // receiver is not a call to a row producer. It fails closed: any other
// resolve syntactically — a local variable, a struct field — is // receiver — a local variable, a struct field, a call to any other
// reported rather than assumed safe. // method — is reported rather than assumed safe.
//
// It sees only calls written x.Scan(...). A method value, f := db.Scan
// followed by f(&v), is out of scope: Scan is never the called
// expression there, and nobody writes a query that way by accident,
// which is the mistake this check exists to catch.
func unguardedScans( func unguardedScans(
fset *token.FileSet, file *ast.File, fset *token.FileSet, file *ast.File,
) []token.Position { ) []token.Position {
@@ -111,15 +114,15 @@ func skipDir(name string) bool {
} }
} }
// walkNonTestGo parses every non-test .go file under root and returns // walkNonTestGo parses every non-test .go file under root. It returns
// how many it parsed along with every unguarded Scan it found. // the directories, relative to root, it parsed a file in, along with
func walkNonTestGo(t *testing.T, root string) (int, []string) { // every unguarded Scan it found.
func walkNonTestGo(t *testing.T, root string) (map[string]bool, []string) {
t.Helper() t.Helper()
var ( walked := map[string]bool{}
parsed int
hits []string var hits []string
)
fset := token.NewFileSet() fset := token.NewFileSet()
@@ -147,7 +150,12 @@ func walkNonTestGo(t *testing.T, root string) (int, []string) {
return err return err
} }
parsed++ dir, err := filepath.Rel(root, filepath.Dir(path))
if err != nil {
return err
}
walked[dir] = true
for _, pos := range unguardedScans(fset, file) { for _, pos := range unguardedScans(fset, file) {
hits = append(hits, relPosition(root, pos)) hits = append(hits, relPosition(root, pos))
@@ -157,7 +165,7 @@ func walkNonTestGo(t *testing.T, root string) (int, []string) {
}, },
)) ))
return parsed, hits return walked, hits
} }
// isNonTestGo reports whether a file name is Go source this check // isNonTestGo reports whether a file name is Go source this check
@@ -189,19 +197,39 @@ func relPosition(root string, pos token.Position) string {
// logged with its values interpolated. The package comment states the // logged with its values interpolated. The package comment states the
// limit; this fails when someone adds a call site anyway. // limit; this fails when someone adds a call site anyway.
// //
// The current tree has one caller, internal/database/database_test.go, // Test files are not governed: what a test binds is fixture data.
// which this check does not govern: it is test-only and its SELECT 1
// binds nothing.
func TestGormScanIsNeverCalledOutsideTests(t *testing.T) { func TestGormScanIsNeverCalledOutsideTests(t *testing.T) {
t.Parallel() t.Parallel()
parsed, offenders := walkNonTestGo(t, moduleRoot(t)) root := moduleRoot(t)
walked, offenders := walkNonTestGo(t, root)
// The module's packages are static, templates, and every directory
// directly under cmd and internal. Each holds non-test code, so one
// the walk parsed nothing in was skipped, and a Scan there would
// pass unseen.
packages := []string{"static", "templates"}
for _, parent := range []string{"cmd", "internal"} {
entries, err := os.ReadDir(filepath.Join(root, parent))
require.NoError(t, err)
for _, entry := range entries {
if !entry.IsDir() {
continue
}
packages = append(packages, filepath.Join(parent, entry.Name()))
}
}
for _, dir := range packages {
require.True(
t, walked[dir],
"the walk parsed no non-test .go file in %s", dir,
)
}
require.GreaterOrEqual(
t, parsed, minNonTestFiles,
"parsed %d non-test .go files, so this check found "+
"nothing to look at", parsed,
)
require.Empty( require.Empty(
t, offenders, t, offenders,
"Scan called on a receiver this check cannot show is a "+ "Scan called on a receiver this check cannot show is a "+
@@ -222,18 +250,51 @@ type scanGuardCase struct {
want int want int
} }
// scanGuardCases covers each receiver form unguardedScans names, plus
// each row producer isRowProducer lets through. Each body is valid Go
// inside plantedFile.
func scanGuardCases() []scanGuardCase { func scanGuardCases() []scanGuardCase {
return []scanGuardCase{ return []scanGuardCase{
{"gorm chain", `db.DB().Raw("SELECT 1").Scan(&v)`, 1}, {"local variable", "q := gdb.Raw(\"SELECT 1\")\n\tq.Scan(&v)", 1},
{"gorm receiver", `gdb.Scan(&v)`, 1}, {"struct field", `s.db.Scan(&v)`, 1},
{"gorm via variable", "q := gdb.Raw(\"x\")\nq.Scan(&v)", 1}, {"gorm chain", `gdb.Raw("SELECT 1").Scan(&v)`, 1},
{"gorm model chain", `gdb.Model(&x).Scan(&v)`, 1}, {
{"sql row", `gdb.Raw("SELECT 1").Row().Scan(&v)`, 0}, "sql rows in a variable",
{"sql rows", `gdb.Raw("SELECT 1").Rows().Scan(&v)`, 0}, "rows, _ := gdb.Raw(\"SELECT 1\").Rows()\n\trows.Scan(&v)",
1,
},
{"gorm Row", `gdb.Raw("SELECT 1").Row().Scan(&v)`, 0},
{"sql QueryRow", `sqlDB.QueryRow("SELECT 1").Scan(&v)`, 0},
{
"sql QueryRowContext",
`sqlDB.QueryRowContext(ctx, "SELECT 1").Scan(&v)`,
0,
},
{"unrelated call", `gdb.Find(&v)`, 0}, {"unrelated call", `gdb.Find(&v)`, 0},
} }
} }
// plantedFile wraps one case body in a function that declares every
// name the bodies use, so each body is the Go it stands for. The result
// is parsed, never compiled.
const plantedFile = `package p
import (
"context"
"database/sql"
"gorm.io/gorm"
)
type store struct{ db *gorm.DB }
func f(ctx context.Context, gdb *gorm.DB, sqlDB *sql.DB, s store) {
var v int
%s
}
`
// TestScanGuard_ReportsPlantedCalls proves the check fires. Without it // TestScanGuard_ReportsPlantedCalls proves the check fires. Without it
// a detector that matched nothing would satisfy the walk above no // a detector that matched nothing would satisfy the walk above no
// matter what the tree contained. // matter what the tree contained.
@@ -245,9 +306,7 @@ func TestScanGuard_ReportsPlantedCalls(t *testing.T) {
t.Parallel() t.Parallel()
fset := token.NewFileSet() fset := token.NewFileSet()
src := fmt.Sprintf( src := fmt.Sprintf(plantedFile, tc.body)
"package p\n\nfunc f() {\n\t%s\n}\n", tc.body,
)
file, err := parser.ParseFile( file, err := parser.ParseFile(
fset, tc.name+".go", src, 0, fset, tc.name+".go", src, 0,
+6 -4
View File
@@ -15,10 +15,12 @@ import (
// eventBodyQuery reads one event's stored body as bytes. The cast // eventBodyQuery reads one event's stored body as bytes. The cast
// to blob is what makes the driver hand back the stored bytes // to blob is what makes the driver hand back the stored bytes
// rather than a string conversion, so Content-Length taken from // rather than a string conversion, so Content-Length taken from
// the result matches what goes on the wire. The soft-delete // the result matches what goes on the wire. The retention reaper
// predicate is spelled out because Raw bypasses GORM's default // deletes event rows outright, so a reaped event is simply gone
// scope, and it is what stops a reaped event still being // and the query finds no row. The deleted_at predicate repeats
// downloadable. // the soft-delete scope GORM adds to its own queries, which Raw
// bypasses; nothing soft-deletes an event, so today it excludes
// nothing.
const eventBodyQuery = "SELECT cast(body as blob) " + const eventBodyQuery = "SELECT cast(body as blob) " +
"FROM events WHERE id = ? AND webhook_id = ? AND deleted_at IS NULL" "FROM events WHERE id = ? AND webhook_id = ? AND deleted_at IS NULL"
+5 -4
View File
@@ -405,10 +405,11 @@ func TestHandleEventBodyDownload_UnknownEvent404s(t *testing.T) {
// route. The body is read in one query before any header is // route. The body is read in one query before any header is
// written, so a reaped event cannot produce a partial download: // written, so a reaped event cannot produce a partial download:
// it is a clean 404 with no Content-Length and no // it is a clean 404 with no Content-Length and no
// Content-Disposition. Both removals the codebase performs are // Content-Disposition. The reaper deletes event rows outright,
// covered — the reaper hard-deletes, and a soft-deleted row is // which is the "hard deleted" case. The "soft deleted" case
// excluded by the query's own deleted_at predicate rather than // covers a row no code produces today: it only pins the query's
// by GORM's default scope, which Raw bypasses. // own deleted_at predicate, the soft-delete condition Raw would
// otherwise skip.
func TestHandleEventBodyDownload_ReapedEvent404s(t *testing.T) { func TestHandleEventBodyDownload_ReapedEvent404s(t *testing.T) {
t.Parallel() t.Parallel()
+26 -23
View File
@@ -5,6 +5,7 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"io" "io"
"unicode"
"unicode/utf8" "unicode/utf8"
) )
@@ -22,11 +23,12 @@ const maxRenderedBodyBytes = 32 << 10
// scrolls, so that it does not make the page huge. // scrolls, so that it does not make the page huge.
const maxInlineBodyLines = 200 const maxInlineBodyLines = 200
// maxIndentGrowth is how many times its size a JSON body may grow // maxIndentDepth is how deeply a JSON body's objects and arrays may
// when it is pretty-printed; past that it is shown as received. // nest for it to be pretty-printed; a deeper one is shown as received.
// Indenting grows with nesting depth as well as with size: 20 KB of // Each level indents every line inside it two more spaces, so 10 KB of
// nested brackets indents to some 200 MB. // nested brackets would indent to some 50 MB; within this depth a body
const maxIndentGrowth = 4 // grows at most 35 times.
const maxIndentDepth = 16
// jsonIndent is the indent of a pretty-printed JSON body. // jsonIndent is the indent of a pretty-printed JSON body.
const jsonIndent = " " const jsonIndent = " "
@@ -70,9 +72,15 @@ func newBodyView(eventURL string, body []byte, size int64) BodyView {
v.ShownBytes = len(body) v.ShownBytes = len(body)
} }
// html/template shows invalid UTF-8 and NUL bytes as replacement // html/template shows invalid UTF-8 as replacement characters,
// characters, so a body holding either is not text. // and a browser shows a control character other than tab, line
if !utf8.Valid(body) || bytes.IndexByte(body, 0) >= 0 { // feed and carriage return as a box or not at all, so a body
// holding either is not text.
isControl := func(r rune) bool {
return unicode.IsControl(r) && r != '\t' && r != '\n' && r != '\r'
}
if !utf8.Valid(body) || bytes.IndexFunc(body, isControl) >= 0 {
v.Binary = true v.Binary = true
return v return v
@@ -83,7 +91,8 @@ func newBodyView(eventURL string, body []byte, size int64) BodyView {
body = indentJSON(body) body = indentJSON(body)
} }
lines := bytes.Count(body, []byte("\n")) + 1 // A final newline ends the last line rather than starting another.
lines := bytes.Count(bytes.TrimSuffix(body, []byte("\n")), []byte("\n")) + 1
v.Text = string(body) v.Text = string(body)
v.Scroll = lines > maxInlineBodyLines || size > maxRenderedBodyBytes v.Scroll = lines > maxInlineBodyLines || size > maxRenderedBodyBytes
@@ -92,8 +101,7 @@ func newBodyView(eventURL string, body []byte, size int64) BodyView {
} }
// indentJSON returns body pretty-printed when it is a JSON document, // indentJSON returns body pretty-printed when it is a JSON document,
// and unchanged when it is not or would grow more than // and unchanged when it is not or nests deeper than maxIndentDepth.
// maxIndentGrowth times.
func indentJSON(body []byte) []byte { func indentJSON(body []byte) []byte {
if !json.Valid(body) || !indentFits(body) { if !json.Valid(body) || !indentFits(body) {
return body return body
@@ -109,21 +117,17 @@ func indentJSON(body []byte) []byte {
return out.Bytes() return out.Bytes()
} }
// indentFits reports whether pretty-printing the JSON document body // indentFits reports whether the objects and arrays of the JSON
// keeps it within maxIndentGrowth times its size. It adds up an upper // document body nest at most maxIndentDepth deep.
// bound instead of indenting: each token starts at most one line,
// indented once per enclosing object or array, and a key gains the
// space after its colon.
func indentFits(body []byte) bool { func indentFits(body []byte) bool {
limit := maxIndentGrowth * len(body) depth := 0
size, depth := len(body), 0
dec := json.NewDecoder(bytes.NewReader(body)) dec := json.NewDecoder(bytes.NewReader(body))
// A number too large for a float64 is still valid JSON. // A number too large for a float64 is still valid JSON.
dec.UseNumber() dec.UseNumber()
for size <= limit { for {
tok, err := dec.Token() tok, err := dec.Token()
if errors.Is(err, io.EOF) { if errors.Is(err, io.EOF) {
return true return true
@@ -136,12 +140,11 @@ func indentFits(body []byte) bool {
switch tok { switch tok {
case json.Delim('{'), json.Delim('['): case json.Delim('{'), json.Delim('['):
depth++ depth++
if depth > maxIndentDepth {
return false
}
case json.Delim('}'), json.Delim(']'): case json.Delim('}'), json.Delim(']'):
depth-- depth--
} }
size += len("\n") + depth*len(jsonIndent) + len(" ")
} }
return false
} }
+52 -7
View File
@@ -40,6 +40,32 @@ func TestNewBodyView_FormatsValidJSON(t *testing.T) {
assert.False(t, v.Scroll) assert.False(t, v.Scroll)
} }
// TestNewBodyView_FormatsNestedJSON proves a document with a few
// levels of nesting is pretty-printed: only deep nesting is shown
// as received.
func TestNewBodyView_FormatsNestedJSON(t *testing.T) {
t.Parallel()
v := bodyView(`{"data":[[1,2,3],[4,5,6]]}`)
assert.Equal(t, []string{
`{`,
` "data": [`,
` [`,
` 1,`,
` 2,`,
` 3`,
` ],`,
` [`,
` 4,`,
` 5,`,
` 6`,
` ]`,
` ]`,
`}`,
}, strings.Split(v.Text, "\n"))
}
// TestNewBodyView_InvalidJSONAsReceived proves a body that is not // TestNewBodyView_InvalidJSONAsReceived proves a body that is not
// a JSON document is shown exactly as it arrived. // a JSON document is shown exactly as it arrived.
func TestNewBodyView_InvalidJSONAsReceived(t *testing.T) { func TestNewBodyView_InvalidJSONAsReceived(t *testing.T) {
@@ -54,12 +80,19 @@ func TestNewBodyView_InvalidJSONAsReceived(t *testing.T) {
} }
} }
// TestNewBodyView_DeepJSONAsReceived proves a JSON body that // TestNewBodyView_DeepJSONAsReceived proves a JSON body nested
// indenting would grow out of all proportion is shown as it // more than 16 levels deep is shown as it arrived. 10 KB of nested
// arrived. 10 KB of nested arrays would indent to some 50 MB. // arrays would indent to some 50 MB.
func TestNewBodyView_DeepJSONAsReceived(t *testing.T) { func TestNewBodyView_DeepJSONAsReceived(t *testing.T) {
t.Parallel() t.Parallel()
nested := func(depth int) string {
return strings.Repeat("[", depth) + "1" + strings.Repeat("]", depth)
}
assert.NotEqual(t, nested(16), bodyView(nested(16)).Text)
assert.Equal(t, nested(17), bodyView(nested(17)).Text)
body := strings.Repeat("[", 5000) + strings.Repeat("]", 5000) body := strings.Repeat("[", 5000) + strings.Repeat("]", 5000)
assert.Equal(t, body, bodyView(body).Text) assert.Equal(t, body, bodyView(body).Text)
@@ -74,6 +107,10 @@ func TestNewBodyView_ScrollsPast200Lines(t *testing.T) {
assert.False(t, bodyView(lines(200)).Scroll) assert.False(t, bodyView(lines(200)).Scroll)
assert.True(t, bodyView(lines(201)).Scroll) assert.True(t, bodyView(lines(201)).Scroll)
// A final newline ends the last line rather than starting another.
assert.False(t, bodyView(lines(200)+"\n").Scroll)
assert.True(t, bodyView(lines(201)+"\n").Scroll)
// One line as received, 201 once formatted: the brackets and // One line as received, 201 once formatted: the brackets and
// 199 elements. // 199 elements.
numbers := "[" + strings.TrimSuffix(strings.Repeat("1,", 199), ",") + "]" numbers := "[" + strings.TrimSuffix(strings.Repeat("1,", 199), ",") + "]"
@@ -94,17 +131,25 @@ func TestNewBodyView_LargeBodyScrolls(t *testing.T) {
} }
// TestNewBodyView_BinaryNotShown proves a body that is not text // TestNewBodyView_BinaryNotShown proves a body that is not text
// is never shown, since html/template would turn it into // is never shown: one that is not valid UTF-8, or that holds a
// replacement characters. // control character other than tab, line feed and carriage return.
func TestNewBodyView_BinaryNotShown(t *testing.T) { func TestNewBodyView_BinaryNotShown(t *testing.T) {
t.Parallel() t.Parallel()
for _, body := range []string{"\xff\xfe\xfd", "a\x00b"} { for _, body := range []string{
"\xff\xfe\xfd",
"a\x00b",
// A small protobuf message: valid UTF-8, but control bytes.
"\x08\x01\x12\x03abc",
"\x1b[31mred\x1b[0m",
"a\x7fb",
} {
v := bodyView(body) v := bodyView(body)
assert.True(t, v.Binary) assert.True(t, v.Binary, "%q", body)
assert.Empty(t, v.Text) assert.Empty(t, v.Text)
} }
assert.False(t, bodyView("snow "+snowman).Binary) assert.False(t, bodyView("snow "+snowman).Binary)
assert.False(t, bodyView("a\tb\r\nc\n").Binary)
} }
+3 -2
View File
@@ -145,8 +145,9 @@ func (h *Handlers) resubmitEvent(
// per-webhook database files — a sibling webhook's event is not in the // per-webhook database files — a sibling webhook's event is not in the
// database being queried at all — and is there so the scoping survives // database being queried at all — and is there so the scoping survives
// any future change that puts more than one webhook's events in one // any future change that puts more than one webhook's events in one
// file. Going through Model applies GORM's soft-delete scope, which is // file. A reaped event is not found because the retention reaper
// what stops a reaped event being resubmitted. // deletes its row outright rather than marking it deleted; see
// deleteEvents in internal/database/retention.go.
func loadResubmitSource( func loadResubmitSource(
webhookDB *gorm.DB, webhookDB *gorm.DB,
webhookID, eventID string, webhookID, eventID string,
+2 -1
View File
@@ -97,7 +97,8 @@ type Handlers struct {
// names through the archive rename, the save and any move back. // names through the archive rename, the save and any move back.
// Interleaved, one could rename an archive between another's // Interleaved, one could rename an archive between another's
// rename and save, leaving the file named for one edit and the // rename and save, leaving the file named for one edit and the
// stored names from the other. // stored names from the other. An archive download holds it while
// it reads the stored names and opens the file they give.
renameMu sync.Mutex renameMu sync.Mutex
// dummyVerifications counts the equivalent-cost verifications // dummyVerifications counts the equivalent-cost verifications
+5 -1
View File
@@ -191,6 +191,7 @@ func storedRetentionDays(
type sourceTestEnv struct { type sourceTestEnv struct {
handlers *handlers.Handlers handlers *handlers.Handlers
db *database.Database db *database.Database
dbMgr *database.WebhookDBManager
archives *recordingArchives archives *recordingArchives
cookies []*http.Cookie cookies []*http.Cookie
} }
@@ -204,9 +205,11 @@ func setupSourceTest(t *testing.T) *sourceTestEnv {
var db *database.Database var db *database.Database
var dbMgr *database.WebhookDBManager
var archives *recordingArchives var archives *recordingArchives
app := newTestApp(t, &h, &sess, &db, &archives) app := newTestApp(t, &h, &sess, &db, &dbMgr, &archives)
app.RequireStart() app.RequireStart()
t.Cleanup(app.RequireStop) t.Cleanup(app.RequireStop)
@@ -214,6 +217,7 @@ func setupSourceTest(t *testing.T) *sourceTestEnv {
return &sourceTestEnv{ return &sourceTestEnv{
handlers: h, handlers: h,
db: db, db: db,
dbMgr: dbMgr,
archives: archives, archives: archives,
cookies: authenticatedCookies( cookies: authenticatedCookies(
t, sess, sourceTestUserID, "sourceuser", t, sess, sourceTestUserID, "sourceuser",
+121
View File
@@ -0,0 +1,121 @@
package handlers
import (
"context"
"errors"
"net/http"
"time"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
)
// downloadWriteTimeout is how long one write of a download may wait
// for a client that has stopped reading.
const downloadWriteTimeout = 60 * time.Second
// HandleTargetDownload serves a database target's archive as one
// gzipped JSON file, named for the webhook, the target and the time;
// see delivery.ArchiveExport.WriteGzipJSON for what it holds. Other
// target types have no archive and are a 404.
//
// A download runs for as long as the client keeps reading: it reads
// under a context the request limit does not cancel, and gives each
// write its own deadline in place of the server's write timeout. It
// stops when a write fails.
func (h *Handlers) HandleTargetDownload() http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
ctx := context.WithoutCancel(r.Context())
webhook, target, export, ok := h.openTargetArchive(ctx, w, r)
if !ok {
return
}
defer func() { _ = export.Close() }()
now := time.Now()
w.Header().Set("Content-Type", "application/gzip")
w.Header().Set(
"Content-Disposition",
`attachment; filename="`+delivery.ArchiveExportFileName(
webhook.Name, target.Name, now,
)+`"`,
)
err := export.WriteGzipJSON(
ctx,
downloadWriter{w: w, rc: http.NewResponseController(w)},
&webhook, target, now,
)
if err != nil {
h.log.Error(
"failed to export archive",
"target_id", target.ID,
"error", err,
)
// The 200 has gone out. Aborting the connection is what
// tells the client the file is incomplete.
panic(http.ErrAbortHandler)
}
}
}
// downloadWriter writes a download to the client, giving each write
// downloadWriteTimeout to finish.
type downloadWriter struct {
w http.ResponseWriter
rc *http.ResponseController
}
func (d downloadWriter) Write(b []byte) (int, error) {
// A writer that has no write deadline, such as a test's recorder,
// answers http.ErrNotSupported and needs none extended.
err := d.rc.SetWriteDeadline(time.Now().Add(downloadWriteTimeout))
if err != nil && !errors.Is(err, http.ErrNotSupported) {
return 0, err
}
return d.w.Write(b)
}
// openTargetArchive opens the archive of the request's database target
// for export, with its reads under ctx. It reports false once it has
// written the response.
//
// It holds renameMu, which every archive rename runs under, while it
// reads the stored names and opens the file, so the file it opens is
// the one those names give. It lets go before the export is streamed:
// once the file is open, a rename does not affect the export.
func (h *Handlers) openTargetArchive(
ctx context.Context,
w http.ResponseWriter,
r *http.Request,
) (database.Webhook, *database.Target, *delivery.ArchiveExport, bool) {
h.renameMu.Lock()
defer h.renameMu.Unlock()
webhook, target, ok := h.ownedTarget(w, r)
if !ok {
return database.Webhook{}, nil, nil, false
}
if target.Type != database.TargetTypeDatabase {
h.renderError(w, r, http.StatusNotFound)
return database.Webhook{}, nil, nil, false
}
export, err := delivery.OpenArchiveExport(
ctx, delivery.ArchivePath(h.dbMgr, &webhook, target), h.log,
)
if err != nil {
h.serverError(w, r, "failed to open archive for export", err)
return database.Webhook{}, nil, nil, false
}
return webhook, target, export, true
}
+379
View File
@@ -0,0 +1,379 @@
package handlers_test
import (
"bytes"
"compress/gzip"
"context"
"crypto/rand"
"encoding/json"
"errors"
"io"
"log/slog"
"net"
"net/http"
"net/http/httptest"
"net/url"
"sync"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
"sneak.berlin/go/webhooker/internal/middleware"
)
// errClientGone is the write failure of a client that has gone away.
var errClientGone = errors.New("client gone")
// downloadPath is the archive download route of a target.
func downloadPath(webhookID, targetID string) string {
return "/hook/" + webhookID + "/targets/" + targetID + "/download"
}
// renameTarget submits the edit form renaming a target to Renamed.
func renameTarget(
env *sourceTestEnv, webhookID, targetID string,
) *httptest.ResponseRecorder {
form := url.Values{}
form.Set("name", "Renamed")
return submitTargetEdit(env, webhookID, targetID, form)
}
// TestHandleTargetDownload proves a database target's archive
// downloads as a gzipped JSON attachment named for the webhook, the
// target and the time, here with no archive file yet, so with no
// rows; and that a target of another type has no download.
func TestHandleTargetDownload(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
logTarget := seedTarget(t, env.db, wh.ID, database.TargetTypeLog)
w := serveTarget(
env, http.MethodGet, downloadPath(wh.ID, archive.ID), nil,
)
require.Equal(t, http.StatusOK, w.Code, w.Body.String())
assert.Equal(t, "application/gzip", w.Header().Get("Content-Type"))
assert.Regexp(t,
`^attachment; filename="archive-seeded-t-database-`+
`\d{8}T\d{6}Z\.json\.gz"$`,
w.Header().Get("Content-Disposition"),
)
zr, err := gzip.NewReader(w.Body)
require.NoError(t, err)
var got map[string]json.RawMessage
require.NoError(t, json.NewDecoder(zr).Decode(&got))
assert.JSONEq(t,
`{"id":"`+archive.ID+`","name":"t-database"}`,
string(got["target"]),
)
assert.JSONEq(t, `[]`, string(got["archived_events"]))
w = serveTarget(
env, http.MethodGet, downloadPath(wh.ID, logTarget.ID), nil,
)
assert.Equal(t, http.StatusNotFound, w.Code)
}
// TestHandleTargetDownload_WaitsForRename proves a download reads the
// target's names and opens its archive under the lock a rename holds:
// started while an edit is renaming the archive, it waits, and is
// named for the target's new name.
func TestHandleTargetDownload_WaitsForRename(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
renaming, release := env.archives.BlockNextRename()
edited := make(chan *httptest.ResponseRecorder, 1)
go func() {
edited <- renameTarget(env, wh.ID, archive.ID)
}()
<-renaming
downloaded := make(chan *httptest.ResponseRecorder, 1)
go func() {
downloaded <- serveTarget(
env, http.MethodGet, downloadPath(wh.ID, archive.ID), nil,
)
}()
select {
case <-downloaded:
release()
t.Fatal("the download did not wait for the rename")
case <-time.After(100 * time.Millisecond):
}
release()
require.Equal(t, http.StatusSeeOther, (<-edited).Code)
w := <-downloaded
require.Equal(t, http.StatusOK, w.Code)
assert.Contains(t,
w.Header().Get("Content-Disposition"), "archive-seeded-renamed-",
)
}
// stalledWriter is a response writer whose first write waits until
// resume is closed, closing writing when it starts to wait.
type stalledWriter struct {
*httptest.ResponseRecorder
once sync.Once
writing chan struct{}
resume chan struct{}
}
func (s *stalledWriter) Write(b []byte) (int, error) {
s.once.Do(func() {
close(s.writing)
<-s.resume
})
return s.ResponseRecorder.Write(b)
}
// TestHandleTargetDownload_StreamsWithoutTheLock proves a download
// lets go of the rename lock once its archive is open: while the
// download is stalled writing, an edit can still rename the target.
func TestHandleTargetDownload_StreamsWithoutTheLock(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
req := httptest.NewRequestWithContext(
t.Context(), http.MethodGet, downloadPath(wh.ID, archive.ID), nil,
)
for _, c := range env.cookies {
req.AddCookie(c)
}
sw := &stalledWriter{
ResponseRecorder: httptest.NewRecorder(),
writing: make(chan struct{}),
resume: make(chan struct{}),
}
downloaded := make(chan struct{})
go func() {
targetRouter(env).ServeHTTP(sw, req)
close(downloaded)
}()
<-sw.writing
edited := make(chan *httptest.ResponseRecorder, 1)
go func() {
edited <- renameTarget(env, wh.ID, archive.ID)
}()
select {
case w := <-edited:
assert.Equal(t, http.StatusSeeOther, w.Code)
case <-time.After(10 * time.Second):
t.Error("the rename waited for the download")
}
close(sw.resume)
<-downloaded
assert.Equal(t, http.StatusOK, sw.Code)
}
// seedArchive writes rows to the archive file at path, each with a
// body of bodySize random bytes, which do not compress. Its table has
// only the columns the test fills; an export writes the others empty.
func seedArchive(t *testing.T, path string, rows, bodySize int) {
t.Helper()
db, err := database.OpenSQLite(path, database.SQLiteModeCreate)
require.NoError(t, err)
defer func() { require.NoError(t, db.Close()) }()
_, err = db.ExecContext(t.Context(),
"CREATE TABLE archived_events (id INTEGER PRIMARY KEY, body TEXT)",
)
require.NoError(t, err)
body := make([]byte, bodySize)
for range rows {
_, _ = rand.Read(body)
_, err = db.ExecContext(t.Context(),
"INSERT INTO archived_events (body) VALUES (?)", string(body),
)
require.NoError(t, err)
}
}
// limitedServer serves the target routes as the server does, behind the
// access log, whose lines it returns, and the request limit, here
// limit, which is also its write timeout. Each connection's send buffer
// is a few KiB, so a larger response is still being written while its
// client is not reading.
func limitedServer(
t *testing.T, env *sourceTestEnv, limit time.Duration,
) (*httptest.Server, *bytes.Buffer) {
t.Helper()
const sendBuffer = 4 << 10
logBuf := new(bytes.Buffer)
mw := middleware.NewForTest(
slog.New(slog.NewJSONHandler(logBuf, nil)),
&config.Config{Environment: config.EnvironmentDev},
nil,
)
srv := httptest.NewUnstartedServer(
mw.Logging()(mw.Timeout(limit)(targetRouter(env))),
)
srv.Config.WriteTimeout = limit
srv.Config.ConnContext = func(
ctx context.Context, c net.Conn,
) context.Context {
tcp, ok := c.(*net.TCPConn)
if assert.True(t, ok) {
assert.NoError(t, tcp.SetWriteBuffer(sendBuffer))
}
return ctx
}
srv.Start()
t.Cleanup(srv.Close)
return srv, logBuf
}
// TestHandleTargetDownload_OutlastsTheRequestLimit proves a download
// runs for as long as the client keeps reading, and is logged as the
// 200 it was. Behind a request limit and a server write timeout of a
// tenth of a second, the client stops reading once the response has
// started, waits three times as long, and still gets the whole file.
// The archive is larger than the connection holds, so the download is
// still being written while the client waits.
func TestHandleTargetDownload_OutlastsTheRequestLimit(t *testing.T) {
t.Parallel()
const (
limit = 100 * time.Millisecond
rows = 8
bodySize = 64 << 10
)
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
seedArchive(
t, delivery.ArchivePath(env.dbMgr, &wh, archive), rows, bodySize,
)
srv, accessLog := limitedServer(t, env, limit)
req, err := http.NewRequestWithContext(
t.Context(), http.MethodGet,
srv.URL+downloadPath(wh.ID, archive.ID), nil,
)
require.NoError(t, err)
for _, c := range env.cookies {
req.AddCookie(c)
}
resp, err := srv.Client().Do(req)
require.NoError(t, err)
defer func() { _ = resp.Body.Close() }()
require.Equal(t, http.StatusOK, resp.StatusCode)
time.Sleep(3 * limit)
zr, err := gzip.NewReader(resp.Body)
require.NoError(t, err)
var (
got map[string]json.RawMessage
events []json.RawMessage
)
require.NoError(t, json.NewDecoder(zr).Decode(&got))
require.NoError(t, json.Unmarshal(got["archived_events"], &events))
assert.Len(t, events, rows)
// Reading to the end makes the gzip reader check that the file was
// finished.
_, err = io.ReadAll(zr)
require.NoError(t, err)
// Close waits for the handler, so the access log line is written.
srv.Close()
var access map[string]any
require.NoError(t, json.Unmarshal(accessLog.Bytes(), &access))
assert.EqualValues(t, http.StatusOK, access["status"])
assert.GreaterOrEqual(t,
access["latency_ms"], float64(limit.Milliseconds()),
"the download must outlast the request limit",
)
}
// brokenWriter is a response writer whose writes fail once the
// response has started, as they do when the client goes away.
type brokenWriter struct {
*httptest.ResponseRecorder
}
func (b brokenWriter) Write(p []byte) (int, error) {
if b.Body.Len() > 0 {
return 0, errClientGone
}
return b.ResponseRecorder.Write(p)
}
// TestHandleTargetDownload_AbortsWhenItFails proves a download that
// fails after its response has started aborts the connection, so the
// client sees a failed download rather than a file that looks
// complete and does not decompress.
func TestHandleTargetDownload_AbortsWhenItFails(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
req := httptest.NewRequestWithContext(
t.Context(), http.MethodGet, downloadPath(wh.ID, archive.ID), nil,
)
for _, c := range env.cookies {
req.AddCookie(c)
}
w := brokenWriter{ResponseRecorder: httptest.NewRecorder()}
assert.PanicsWithValue(t, http.ErrAbortHandler, func() {
targetRouter(env).ServeHTTP(w, req)
})
assert.Equal(t, http.StatusOK, w.Code)
}
+6 -2
View File
@@ -37,14 +37,18 @@ const (
editAuthHeader = "Authorization: Bearer " + editBearerSecret editAuthHeader = "Authorization: Bearer " + editBearerSecret
) )
// targetRouter mounts the target create and edit routes on a chi // targetRouter mounts the target create, edit and download routes on
// router so the handlers see the URL parameters they read. // a chi router so the handlers see the URL parameters they read.
func targetRouter(env *sourceTestEnv) *chi.Mux { func targetRouter(env *sourceTestEnv) *chi.Mux {
router := chi.NewRouter() router := chi.NewRouter()
router.Post( router.Post(
"/hook/{sourceID}/targets", "/hook/{sourceID}/targets",
env.handlers.HandleTargetCreate(), env.handlers.HandleTargetCreate(),
) )
router.Get(
"/hook/{sourceID}/targets/{targetID}/download",
env.handlers.HandleTargetDownload(),
)
router.Get( router.Get(
"/hook/{sourceID}/targets/{targetID}/edit", "/hook/{sourceID}/targets/{targetID}/edit",
env.handlers.HandleTargetEdit(), env.handlers.HandleTargetEdit(),
+6 -4
View File
@@ -272,10 +272,12 @@ func requestEventSource(
// createAndFanOut writes the event and one pending delivery per target, // createAndFanOut writes the event and one pending delivery per target,
// and adds them to the webhook's running totals, in a single // and adds them to the webhook's running totals, in a single
// transaction, then hands the tasks to the delivery engine. It is the // transaction, then hands the tasks to the delivery engine. Every
// only path by which an event and its deliveries are created, so a // event is created here, received or resubmitted, so a resubmitted
// resubmitted event is retried, SSRF-guarded and circuit-broken // event is retried, SSRF-guarded and circuit-broken exactly as a
// exactly as a received one is. // received one is. Per-delivery replay is the one other path that
// creates a delivery: it adds one to an existing event without
// coming through here.
// //
// The tasks are returned as well as queued, so a caller can report how // The tasks are returned as well as queued, so a caller can report how
// many targets the event went to. // many targets the event went to.
+71
View File
@@ -0,0 +1,71 @@
package middleware
import (
"context"
"errors"
"net/http"
"time"
)
// Timeout returns middleware that gives each request limit to finish:
// it cancels the request's context once limit has passed, and answers
// 504 when the handler then returns without having started its
// response.
//
// It replaces chi's middleware.Timeout, which writes that 504 even
// after the handler has sent its own status. A download that outlasts
// the limit has already sent its 200 and the whole file, so the late
// 504 changes nothing for the client: the access log and the metrics
// would record it in place of the 200, and net/http would complain of
// a superfluous WriteHeader.
func (s *Middleware) Timeout(
limit time.Duration,
) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(
w http.ResponseWriter,
r *http.Request,
) {
ctx, cancel := context.WithTimeout(r.Context(), limit)
defer cancel()
tw := &timeoutResponseWriter{ResponseWriter: w}
next.ServeHTTP(tw, r.WithContext(ctx))
if !tw.started &&
errors.Is(ctx.Err(), context.DeadlineExceeded) {
w.WriteHeader(http.StatusGatewayTimeout)
}
})
}
}
// timeoutResponseWriter records whether the handler has started its
// response.
type timeoutResponseWriter struct {
http.ResponseWriter
started bool
}
func (w *timeoutResponseWriter) WriteHeader(code int) {
w.started = true
w.ResponseWriter.WriteHeader(code)
}
func (w *timeoutResponseWriter) Write(b []byte) (int, error) {
// A Write without a WriteHeader starts the response too: net/http
// sends 200 in front of it.
w.started = true
//nolint:wrapcheck // Pass the writer's own error through unchanged.
return w.ResponseWriter.Write(b)
}
// Unwrap lets http.ResponseController reach the writer underneath, so
// a handler can still set a write deadline through this wrapper.
func (w *timeoutResponseWriter) Unwrap() http.ResponseWriter {
return w.ResponseWriter
}
+56
View File
@@ -0,0 +1,56 @@
package middleware_test
import (
"net/http"
"net/http/httptest"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// TestTimeout proves the request limit answers 504 to a handler that
// outlasts it without starting its response, and leaves a response the
// handler has started with the status it sent. Both are what the
// access log records.
func TestTimeout(t *testing.T) {
t.Parallel()
const limit = 10 * time.Millisecond
for _, tc := range []struct {
name string
sent int // the status the handler sends, or 0 for none
want int
}{
{name: "not started", sent: 0, want: http.StatusGatewayTimeout},
{name: "started", sent: http.StatusOK, want: http.StatusOK},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
m, buf := capturingMiddleware(t)
handler := m.Logging()(m.Timeout(limit)(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
if tc.sent != 0 {
w.WriteHeader(tc.sent)
}
<-r.Context().Done()
},
)))
w := httptest.NewRecorder()
handler.ServeHTTP(w, httptest.NewRequestWithContext(
t.Context(), http.MethodGet, "/", nil,
))
assert.Equal(t, tc.want, w.Code)
entries := accessLogEntries(t, buf)
require.Len(t, entries, 1)
assert.EqualValues(t, tc.want, entries[0]["status"])
})
}
}
+69
View File
@@ -0,0 +1,69 @@
package server_test
import (
"compress/gzip"
"encoding/json"
"net/http"
"regexp"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"gorm.io/gorm/clause"
"sneak.berlin/go/webhooker/internal/database"
)
// TestHook_DownloadArchive follows the Download link the webhook page
// shows for a database target, and only for it, and gets the archive
// as a gzipped JSON file. Signed out, the link leads to the login page.
func TestHook_DownloadArchive(t *testing.T) {
t.Parallel()
env := newTestEnv(t)
userID, _ := env.seedUser(t, "archivist", "somepassword")
cookies := env.authCookies(t, userID, "archivist")
wh := env.seedWebhook(t, userID)
env.seedTarget(t, wh.ID)
archive := &database.Target{
WebhookID: wh.ID,
Name: "kept",
Type: database.TargetTypeDatabase,
Active: true,
}
require.NoError(t,
env.db.DB().Omit(clause.Associations).Create(archive).Error,
)
page := env.get("/hook/"+wh.ID, cookies)
require.Equal(t, http.StatusOK, page.Code)
links := regexp.MustCompile(
`href="(/hook/[^/"]+/targets/[^/"]+/download)"`,
).FindAllStringSubmatch(page.Body.String(), -1)
require.Len(t, links, 1, "only the database target has a Download")
link := links[0][1]
assert.Equal(t,
"/hook/"+wh.ID+"/targets/"+archive.ID+"/download", link,
)
w := env.get(link, cookies)
require.Equal(t, http.StatusOK, w.Code)
assert.Equal(t, "application/gzip", w.Header().Get("Content-Type"))
zr, err := gzip.NewReader(w.Body)
require.NoError(t, err)
var got map[string]json.RawMessage
require.NoError(t, json.NewDecoder(zr).Decode(&got))
assert.JSONEq(t,
`{"id":"`+wh.ID+`","name":"routed"}`, string(got["webhook"]),
)
w = env.get(link, nil)
assert.Equal(t, http.StatusSeeOther, w.Code)
assert.Contains(t, w.Header().Get("Location"), "/pages/login")
}
+10
View File
@@ -1,6 +1,8 @@
package server package server
import ( import (
"context"
"log/slog"
"net/http" "net/http"
"testing" "testing"
@@ -37,6 +39,14 @@ func SentryClientOptionsForTest(
return sentryClientOptions(dsn, release) return sentryClientOptions(dsn, release)
} }
// CleanShutdownForTest runs the server's stop hook, cleanShutdown,
// against hs: a server the test started itself, so it can hold a
// request open across the drain. Sentry is off.
func CleanShutdownForTest(ctx context.Context, hs *http.Server) {
s := &Server{log: slog.New(slog.DiscardHandler), httpServer: hs}
s.cleanShutdown(ctx)
}
// newServerForTest builds a Server through New, as the application // newServerForTest builds a Server through New, as the application
// does, on a lifecycle that is never started: the hooks New adds to // does, on a lifecycle that is never started: the hooks New adds to
// it never run, so nothing listens. // it never run, so nothing listens.
+215
View File
@@ -0,0 +1,215 @@
package server_test
import (
"net/http"
"net/url"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// maxResubmits bounds the requests the tests below send to the
// resubmit route. The route's rate limit belongs to the middleware;
// this only has to sit well above it, so that a route without the
// limiter fails its test instead of looping.
const maxResubmits = 100
// resubmitPath is the resubmit route for one stored event.
func resubmitPath(webhookID, eventID string) string {
return "/hook/" + webhookID + "/events/" + eventID + "/resubmit"
}
// csrfForm is a resubmit form carrying the given CSRF token.
func csrfForm(token string) url.Values {
form := url.Values{}
form.Set("csrf_token", token)
return form
}
// TestEventResubmit_SignedOutRequestsNeverReachTheRateLimit pins
// RequireAuth on the resubmit route. The handler also turns away a
// request without a session, with the same redirect, so a refusal
// alone would pass without RequireAuth. What RequireAuth adds is that
// it refuses such a request before the route's rate limit, so a
// signed-out client cannot spend the budget a signed-in user
// resubmits from. Each request carries a CSRF token valid for its own
// cookie, so CSRF lets it through to RequireAuth.
func TestEventResubmit_SignedOutRequestsNeverReachTheRateLimit(
t *testing.T,
) {
t.Parallel()
env := newTestEnv(t)
userID, _ := env.seedUser(t, "resubmitter", "somepassword")
wh := env.seedWebhook(t, userID)
evt := env.seedEvent(t, wh.ID, `{"resubmit":"me"}`)
path := resubmitPath(wh.ID, evt.ID)
logsPath := "/hook/" + wh.ID + "/events"
token, signedOut := env.csrfFrom(t, "/pages/login", nil)
for i := range maxResubmits {
w := env.post(path, csrfForm(token), signedOut)
require.Equal(t, http.StatusSeeOther, w.Code, "request %d", i)
require.Equal(
t, "/pages/login", w.Header().Get("Location"),
"request %d", i,
)
}
require.Equal(
t, int64(1), env.countEvents(t, wh.ID),
"a signed-out request must store nothing",
)
token, cookies := env.csrfFrom(
t, logsPath, env.authCookies(t, userID, "resubmitter"),
)
env.requireNotice(
t, env.post(path, csrfForm(token), cookies),
logsPath, "resubmit-no-targets",
"this source has no active targets", cookies,
)
}
// TestEventResubmit_RefusedWithoutAValidCSRFToken pins CSRF on the
// resubmit route: a signed-in user's POST is refused with 403, and
// stores nothing, unless it carries the token issued to that user's
// own browser.
func TestEventResubmit_RefusedWithoutAValidCSRFToken(t *testing.T) {
t.Parallel()
env := newTestEnv(t)
userID, _ := env.seedUser(t, "resubmitter", "somepassword")
wh := env.seedWebhook(t, userID)
evt := env.seedEvent(t, wh.ID, `{"resubmit":"me"}`)
path := resubmitPath(wh.ID, evt.ID)
logsPath := "/hook/" + wh.ID + "/events"
token, cookies := env.csrfFrom(
t, logsPath, env.authCookies(t, userID, "resubmitter"),
)
otherBrowsers, _ := env.csrfFrom(t, "/pages/login", nil)
for name, form := range map[string]url.Values{
"no token": {},
"a malformed token": csrfForm("not-a-token"),
"another browser's token": csrfForm(otherBrowsers),
} {
assert.Equal(
t, http.StatusForbidden,
env.post(path, form, cookies).Code, name,
)
}
assert.Equal(
t, int64(1), env.countEvents(t, wh.ID),
"a refused request must store nothing",
)
// The same request with the user's own token goes through, so the
// refusals above were the token's doing.
env.requireNotice(
t, env.post(path, csrfForm(token), cookies),
logsPath, "resubmit-no-targets",
"this source has no active targets", cookies,
)
}
// TestEventResubmit_AnotherWebhooksEvent404s pins, on the route as
// registered, that a signed-in user gets 404, and nothing is stored,
// for an event of a webhook another user owns, which the handler's
// ownership check refuses, and for another webhook's event posted
// under a webhook the user does own, which the event lookup refuses.
// The user's own event, posted the same way, is accepted, so the
// second 404 comes from the lookup and not from a route that never
// passed the event ID to the handler.
func TestEventResubmit_AnotherWebhooksEvent404s(t *testing.T) {
t.Parallel()
env := newTestEnv(t)
ownerID, _ := env.seedUser(t, "owner", "somepassword")
owners := env.seedWebhook(t, ownerID)
ownersEvent := env.seedEvent(t, owners.ID, `{"owner":"only"}`)
intruderID, _ := env.seedUser(t, "intruder", "somepassword")
intruders := env.seedWebhook(t, intruderID)
intrudersEvent := env.seedEvent(
t, intruders.ID, `{"intruder":"own"}`,
)
intrudersLogs := "/hook/" + intruders.ID + "/events"
token, cookies := env.csrfFrom(
t, intrudersLogs, env.authCookies(t, intruderID, "intruder"),
)
for name, path := range map[string]string{
"another user's webhook": resubmitPath(
owners.ID, ownersEvent.ID,
),
"another webhook's event": resubmitPath(
intruders.ID, ownersEvent.ID,
),
} {
w := env.post(path, csrfForm(token), cookies)
assert.Equal(t, http.StatusNotFound, w.Code, name)
}
assert.Equal(t, int64(1), env.countEvents(t, owners.ID))
assert.Equal(t, int64(1), env.countEvents(t, intruders.ID))
env.requireNotice(
t,
env.post(
resubmitPath(intruders.ID, intrudersEvent.ID),
csrfForm(token), cookies,
),
intrudersLogs, "resubmit-no-targets",
"this source has no active targets", cookies,
)
}
// TestEventResubmit_RateLimited pins the rate limit on the resubmit
// route: a signed-in user's resubmits are accepted until the budget
// is spent, and then refused with 429.
func TestEventResubmit_RateLimited(t *testing.T) {
t.Parallel()
env := newTestEnv(t)
userID, _ := env.seedUser(t, "resubmitter", "somepassword")
wh := env.seedWebhook(t, userID)
evt := env.seedEvent(t, wh.ID, `{"resubmit":"me"}`)
path := resubmitPath(wh.ID, evt.ID)
token, cookies := env.csrfFrom(
t, "/hook/"+wh.ID+"/events",
env.authCookies(t, userID, "resubmitter"),
)
limited := false
for range maxResubmits {
code := env.post(path, csrfForm(token), cookies).Code
if code == http.StatusTooManyRequests {
limited = true
break
}
require.Equal(
t, http.StatusSeeOther, code,
"a resubmit within the budget must be accepted",
)
}
assert.True(
t, limited, "repeated resubmits must eventually be refused",
)
}
+5 -1
View File
@@ -73,7 +73,7 @@ func (s *Server) setupGlobalMiddleware() {
} }
s.router.Use(s.mw.CORS()) s.router.Use(s.mw.CORS())
s.router.Use(middleware.Timeout(requestTimeout)) s.router.Use(s.mw.Timeout(requestTimeout))
// Panic recovery, deliberately here rather than first. It has to // Panic recovery, deliberately here rather than first. It has to
// run inside every middleware that observes the response, so the // run inside every middleware that observes the response, so the
@@ -313,6 +313,10 @@ func (s *Server) setupSourceRoutes() {
"/targets/{targetID}/edit", "/targets/{targetID}/edit",
s.h.HandleTargetEditSubmit(), s.h.HandleTargetEditSubmit(),
) )
r.Get(
"/targets/{targetID}/download",
s.h.HandleTargetDownload(),
)
r.Post( r.Post(
"/targets/{targetID}/delete", "/targets/{targetID}/delete",
s.h.HandleTargetDelete(), s.h.HandleTargetDelete(),
+17
View File
@@ -445,6 +445,23 @@ func (e *testEnv) countDeliveries(
return count return count
} }
// countEvents reports how many events a webhook's database holds.
func (e *testEnv) countEvents(t *testing.T, webhookID string) int64 {
t.Helper()
webhookDB, err := e.dbMgr.GetDB(webhookID)
require.NoError(t, err)
var count int64
require.NoError(
t,
webhookDB.Model(&database.Event{}).Count(&count).Error,
)
return count
}
// storedHash reads the current password hash for a username. // storedHash reads the current password hash for a username.
func (e *testEnv) storedHash(t *testing.T, username string) string { func (e *testEnv) storedHash(t *testing.T, username string) string {
t.Helper() t.Helper()
+26 -3
View File
@@ -39,6 +39,12 @@ const (
// refuses to spend, leaving it for the hooks that run after the // refuses to spend, leaving it for the hooks that run after the
// server: the delivery engine, the healthcheck, the webhook DB // server: the delivery engine, the healthcheck, the webhook DB
// manager and the database close. // manager and the database close.
//
// Its value is not tuned to those hooks, which take about a
// millisecond between them. It is what the 5s fx stop timeout in
// cmd/webhooker leaves after a full ShutdownTimeout drain, so a
// drain that starts on a full budget still gets all of
// ShutdownTimeout.
TailHookReserve = 2 * time.Second TailHookReserve = 2 * time.Second
// sentryFlushTimeout is the longest wait for Sentry to flush // sentryFlushTimeout is the longest wait for Sentry to flush
@@ -59,6 +65,16 @@ const (
// key off it, and a zero exit would read as a deliberate stop. // key off it, and a zero exit would read as a deliberate stop.
const StartupFailureExitCode = 1 const StartupFailureExitCode = 1
// DrainBudget reports how long the HTTP drain may wait for in-flight
// requests when remaining is the time left on the fx stop context as
// the server's stop hook starts. The hooks before the server can
// already have spent part of the budget, so the drain takes its time
// out of what they left, never out of TailHookReserve. Zero or less
// means no wait at all.
func DrainBudget(remaining time.Duration) time.Duration {
return min(ShutdownTimeout, remaining-TailHookReserve)
}
// SentryFlushBudget reports how long the Sentry flush may run when // SentryFlushBudget reports how long the Sentry flush may run when
// remaining is the time left on the fx stop context after the HTTP // remaining is the time left on the fx stop context after the HTTP
// drain. sentry.Flush takes a bare duration and honours no context, // drain. sentry.Flush takes a bare duration and honours no context,
@@ -261,10 +277,17 @@ func (s *Server) cleanupForExit() {
s.log.Info("cleaning up") s.log.Info("cleaning up")
} }
// cleanShutdown drains the HTTP server and flushes Sentry inside what
// is left of the fx stop budget. A context carrying no deadline — a
// caller outside the fx lifecycle — gets the full ShutdownTimeout.
func (s *Server) cleanShutdown(ctx context.Context) { func (s *Server) cleanShutdown(ctx context.Context) {
ctxShutdown, shutdownCancel := context.WithTimeout( drain := ShutdownTimeout
ctx, ShutdownTimeout,
) if deadline, ok := ctx.Deadline(); ok {
drain = DrainBudget(time.Until(deadline))
}
ctxShutdown, shutdownCancel := context.WithTimeout(ctx, drain)
defer shutdownCancel() defer shutdownCancel()
err := s.httpServer.Shutdown(ctxShutdown) err := s.httpServer.Shutdown(ctxShutdown)
+135
View File
@@ -1,13 +1,148 @@
package server_test package server_test
import ( import (
"context"
"net"
"net/http"
"testing" "testing"
"testing/synctest"
"time" "time"
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/server" "sneak.berlin/go/webhooker/internal/server"
) )
// TestDrainBudget covers the clamp that keeps the HTTP drain from
// spending the tail hooks' share of the fx stop budget when the hooks
// before the server have already used part of it.
func TestDrainBudget(t *testing.T) {
t.Parallel()
tests := []struct {
name string
remaining time.Duration
want time.Duration
}{
{
name: "only the reserve is left",
remaining: server.TailHookReserve,
want: 0,
},
{
name: "earlier hooks spent part of the budget",
remaining: server.TailHookReserve + time.Second,
want: time.Second,
},
{
name: "capped at the nominal timeout",
remaining: time.Hour,
want: server.ShutdownTimeout,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
require.Equal(t, tt.want, server.DrainBudget(tt.remaining))
})
}
}
// TestCleanShutdown_LeavesTailHookReserve stops the server with a
// request still in flight, after the hooks before it have spent all
// of the stop budget but TailHookReserve. The drain must give up at
// once rather than wait for the request: what is left belongs to the
// hooks after the server, the database close among them. A drain
// bounded only by ShutdownTimeout waits until the stop context
// expires, and fx then skips those hooks.
//
// The test runs in a synctest bubble, whose clock moves only while
// every goroutine in it is blocked, so a drain that gives up at once
// leaves the stop context unexpired however slow the host is. The
// request travels over net.Pipe because a goroutine waiting on a
// real socket would stop that clock from moving at all.
func TestCleanShutdown_LeavesTailHookReserve(t *testing.T) {
t.Parallel()
synctest.Test(t, func(t *testing.T) {
entered := make(chan struct{})
release := make(chan struct{})
hs := &http.Server{
Handler: http.HandlerFunc(
func(http.ResponseWriter, *http.Request) {
close(entered)
<-release
},
),
ReadHeaderTimeout: time.Second,
}
srvConn, cliConn := net.Pipe()
listener := pipeListener{
conns: make(chan net.Conn, 1),
closed: make(chan struct{}),
}
listener.conns <- srvConn
go func() { _ = hs.Serve(listener) }()
// Cleanups run last first: the handler returns, then closing
// the client end ends the server's write of the response.
t.Cleanup(func() { _ = cliConn.Close() })
t.Cleanup(func() { close(release) })
_, err := cliConn.Write(
[]byte("GET / HTTP/1.1\r\nHost: webhooker.test\r\n\r\n"),
)
require.NoError(t, err)
<-entered
stopCtx, cancel := context.WithTimeout(
t.Context(), server.TailHookReserve,
)
defer cancel()
server.CleanShutdownForTest(stopCtx, hs)
require.NoError(
t, stopCtx.Err(), "the drain spent the tail hooks' reserve",
)
})
}
// pipeListener is the net.Listener http.Server.Serve needs to serve
// the server end of a net.Pipe: Accept returns that one connection,
// then waits until Close, as a real listener with no more clients
// does.
type pipeListener struct {
conns chan net.Conn
closed chan struct{}
}
func (l pipeListener) Accept() (net.Conn, error) {
select {
case conn := <-l.conns:
return conn, nil
case <-l.closed:
return nil, net.ErrClosed
}
}
func (l pipeListener) Close() error {
close(l.closed)
return nil
}
// Addr is never called by http.Server.Serve.
func (pipeListener) Addr() net.Addr {
return nil
}
// TestSentryFlushBudget covers the clamp that keeps the Sentry flush // TestSentryFlushBudget covers the clamp that keeps the Sentry flush
// from spending the tail hooks' share of the fx stop budget. // from spending the tail hooks' share of the fx stop budget.
// sentry.Flush ignores the stop context, so without the clamp a // sentry.Flush ignores the stop context, so without the clamp a
+3
View File
@@ -157,6 +157,9 @@
{{else}} {{else}}
<span class="badge-error">Inactive</span> <span class="badge-error">Inactive</span>
{{end}} {{end}}
{{if eq .Type "database"}}
<a href="/hook/{{$.Webhook.ID}}/targets/{{.ID}}/download" class="btn-small" title="Download the archive as gzipped JSON">Download</a>
{{end}}
<a href="/hook/{{$.Webhook.ID}}/targets/{{.ID}}/edit" class="btn-small" title="Edit">Edit</a> <a href="/hook/{{$.Webhook.ID}}/targets/{{.ID}}/edit" class="btn-small" title="Edit">Edit</a>
<form method="POST" action="/hook/{{$.Webhook.ID}}/targets/{{.ID}}/toggle" class="inline"> <form method="POST" action="/hook/{{$.Webhook.ID}}/targets/{{.ID}}/toggle" class="inline">
<input type="hidden" name="csrf_token" value="{{$.CSRFToken}}"> <input type="hidden" name="csrf_token" value="{{$.CSRFToken}}">