Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9b7ba5150b |
@@ -2033,29 +2033,6 @@ 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
|
||||||
@@ -2634,10 +2611,11 @@ 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; only tests do, and what a
|
recorder. No production code path calls it; its one caller is
|
||||||
test binds is fixture data. `internal/gormlog/scan_guard_test.go` fails
|
`internal/database/database_test.go:91`, whose `SELECT 1` binds
|
||||||
if a non-test file calls it. `Pluck`, `Row` and `Raw` all run through
|
nothing, and `internal/gormlog/scan_guard_test.go` fails if a non-test
|
||||||
the normal callback processor and are filtered.
|
file calls it. `Pluck`, `Row` and `Raw` all run through the normal
|
||||||
|
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
|
||||||
@@ -2922,7 +2900,6 @@ 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 |
|
||||||
|
|
||||||
@@ -3004,7 +2981,6 @@ 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
|
||||||
@@ -3271,9 +3247,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 by `server.ShutdownTimeout`
|
3. `server` — the HTTP drain, bounded separately by
|
||||||
(**3 seconds**) and by what the hooks before it left, then a Sentry
|
`server.ShutdownTimeout` (**3 seconds**), then a Sentry flush if
|
||||||
flush if `SENTRY_DSN` is set
|
`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`
|
||||||
@@ -3293,30 +3269,23 @@ 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. The reserve is that
|
(`server.TailHookReserve`) for the tail, which is far more than the
|
||||||
remainder, not a figure sized to the tail, which takes about a
|
microseconds it needs.
|
||||||
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 server hook could take it in two ways. The hooks before it may
|
the Sentry flush is what could take it: it runs after the drain
|
||||||
already have spent part of the budget, so a full 3-second drain
|
**inside the same hook**, and `sentry.Flush` takes a bare duration
|
||||||
would come out of the reserve; the drain is therefore also bounded
|
and honours no context, so an unreachable Sentry endpoint would add
|
||||||
by whatever is left on the stop context minus the reserve. And the
|
its own timeout on top of a full-length drain and consume the whole
|
||||||
Sentry flush runs after the drain **inside the same hook**, and
|
sequence budget by itself. It is therefore clamped to whatever is
|
||||||
`sentry.Flush` takes a bare duration and honours no context, so an
|
left on the stop context minus the reserve, and skipped when that
|
||||||
unreachable Sentry endpoint would add its own timeout on top of a
|
leaves too little to be worth attempting — so a full-length drain
|
||||||
full-length drain and consume the whole sequence budget by itself.
|
means Sentry events are dropped rather than the database close being
|
||||||
It is clamped the same way, and skipped when that leaves too little
|
skipped.
|
||||||
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 slow
|
This does not make the database close unconditional: a wedged
|
||||||
`ArchiveSweeper` or `RetentionReaper` is enough to cut the shutdown
|
`ArchiveSweeper` or `RetentionReaper` still runs first and can
|
||||||
short, not only one that consumes the whole budget: what they spend
|
consume the whole budget on its own.
|
||||||
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
|
||||||
@@ -3381,9 +3350,8 @@ version is fixed independently of the compiler's:
|
|||||||
|
|
||||||
1. **Lint stage** (`golangci/golangci-lint:v2.12.2`, Debian-based) —
|
1. **Lint stage** (`golangci/golangci-lint:v2.12.2`, Debian-based) —
|
||||||
installs `make`, downloads dependencies, copies the source, and runs
|
installs `make`, downloads dependencies, copies the source, and runs
|
||||||
`make fmt-check`, then `script/assets` to extract Alpine.js from
|
`make fmt-check`, then `golangci-lint config verify` and
|
||||||
`3p/`, then `golangci-lint config verify` and `golangci-lint run`,
|
`golangci-lint run`, both with `--network=none`.
|
||||||
both with `--network=none`.
|
|
||||||
2. **Builder stage** (`golang:1.26.1-bookworm`) — depends on the lint
|
2. **Builder stage** (`golang:1.26.1-bookworm`) — depends on the lint
|
||||||
stage passing (it copies a file from it), runs `make test` and
|
stage passing (it copies a file from it), runs `make test` and
|
||||||
`make build` (both extract Alpine.js from `3p/` first), and finally
|
`make build` (both extract Alpine.js from `3p/` first), and finally
|
||||||
|
|||||||
@@ -38,19 +38,17 @@ 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 HTTP drain plus the Sentry flush that
|
// close. That hook is the 3s HTTP drain plus the Sentry flush that
|
||||||
// follows it in the same hook, and each is clamped to the stop
|
// follows it in the same hook, so the flush 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 3s and 2s; the reserve is what the tail
|
// running for its own fixed 2s; the reserve is what the tail hooks
|
||||||
// hooks live on, and they are microsecond-scale in normal operation.
|
// live on, and they are microsecond-scale in normal operation.
|
||||||
// TestStopTimeout_LeavesHeadroomForTailHooks pins the arithmetic
|
// TestStopTimeout_LeavesHeadroomForTailHooks pins the arithmetic
|
||||||
// across every drain length and every amount of budget the hooks
|
// across every drain length.
|
||||||
// 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
|
||||||
// What they spend comes out of the drain first, but past 3s it comes
|
// and can still consume the whole budget on their own.
|
||||||
// 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
|
||||||
|
|||||||
@@ -252,40 +252,22 @@ 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.
|
||||||
//
|
//
|
||||||
// Nor does the hook start on a full budget: the ArchiveSweeper and
|
// Shrinking either budget, or unbounding the flush again, must fail
|
||||||
// RetentionReaper hooks run before it, and whatever they spent is
|
// here rather than silently recreating a hook that swallows the
|
||||||
// gone. The outer sweep walks every amount they can spend. Once they
|
// whole sequence.
|
||||||
// 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 spent := time.Duration(0); spent <= stopTimeout; spent += step {
|
for drain := time.Duration(0); drain <= server.ShutdownTimeout; drain += step {
|
||||||
remaining := stopTimeout - spent
|
hook := drain + server.SentryFlushBudget(stopTimeout-drain)
|
||||||
longest := max(server.DrainBudget(remaining), 0)
|
|
||||||
|
|
||||||
for drain := time.Duration(0); drain <= longest; drain += step {
|
require.LessOrEqual(
|
||||||
hook := drain + server.SentryFlushBudget(remaining-drain)
|
t, hook+tailHeadroom, stopTimeout,
|
||||||
|
"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,
|
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -22,7 +22,6 @@ 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 (
|
||||||
@@ -71,8 +70,7 @@ 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},
|
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
|
||||||
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
|
|
||||||
)
|
)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
@@ -170,8 +168,7 @@ func (env *archiveEnv) seedArchiveRows(
|
|||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
gdb, err := gorm.Open(
|
gdb, err := gorm.Open(
|
||||||
sqlite.Dialector{Conn: sqlDB},
|
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
|
||||||
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
|
|
||||||
)
|
)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
@@ -230,8 +227,7 @@ 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},
|
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
|
||||||
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
|
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return 0, err
|
return 0, err
|
||||||
|
|||||||
@@ -23,7 +23,6 @@ import (
|
|||||||
_ "modernc.org/sqlite"
|
_ "modernc.org/sqlite"
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
"sneak.berlin/go/webhooker/internal/delivery"
|
"sneak.berlin/go/webhooker/internal/delivery"
|
||||||
"sneak.berlin/go/webhooker/internal/gormlog"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// iSetup holds common integration test dependencies.
|
// iSetup holds common integration test dependencies.
|
||||||
@@ -81,8 +80,7 @@ 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},
|
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
|
||||||
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
|
|
||||||
)
|
)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
|||||||
@@ -26,7 +26,6 @@ import (
|
|||||||
_ "modernc.org/sqlite"
|
_ "modernc.org/sqlite"
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
"sneak.berlin/go/webhooker/internal/delivery"
|
"sneak.berlin/go/webhooker/internal/delivery"
|
||||||
"sneak.berlin/go/webhooker/internal/gormlog"
|
|
||||||
"sneak.berlin/go/webhooker/internal/metrics"
|
"sneak.berlin/go/webhooker/internal/metrics"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -51,8 +50,7 @@ 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},
|
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
|
||||||
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
|
|
||||||
)
|
)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package delivery
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
@@ -276,9 +277,10 @@ 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 is the one ArchivePath gives for the webhook and the target as
|
// file lives beside the webhook's event database in the data
|
||||||
// the main database names them now; from then on only rename changes
|
// directory and is named for the webhook and the target as the main
|
||||||
// the name the writer uses. It does not touch the archive file.
|
// database has them now; from then on only rename changes the name
|
||||||
|
// 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) {
|
||||||
@@ -297,10 +299,12 @@ func (t *databaseTarget) newWriter(
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
w := newArchiveWriter(
|
dir := filepath.Dir(t.eng.dbManager.DBPath(target.WebhookID))
|
||||||
ArchivePath(t.eng.dbManager, &target.Webhook, &target),
|
name := ArchiveFileName(
|
||||||
t.eng.log,
|
target.Webhook.Name, target.Name, target.ID,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
w := newArchiveWriter(filepath.Join(dir, name), t.eng.log)
|
||||||
w.webhookID = target.WebhookID
|
w.webhookID = target.WebhookID
|
||||||
|
|
||||||
return w, nil
|
return w, nil
|
||||||
|
|||||||
@@ -1,275 +0,0 @@
|
|||||||
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
|
|
||||||
}
|
|
||||||
@@ -1,412 +0,0 @@
|
|||||||
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),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
@@ -17,7 +17,6 @@ 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 {
|
||||||
@@ -43,8 +42,7 @@ func openArchiveDBForRead(
|
|||||||
t.Cleanup(func() { _ = sqlDB.Close() })
|
t.Cleanup(func() { _ = sqlDB.Close() })
|
||||||
|
|
||||||
gdb, err := gorm.Open(
|
gdb, err := gorm.Open(
|
||||||
sqlite.Dialector{Conn: sqlDB},
|
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
|
||||||
&gorm.Config{Logger: gormlog.New(slog.New(slog.DiscardHandler))},
|
|
||||||
)
|
)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
|||||||
@@ -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; only tests do, and
|
// this interface. No production code path calls it; its one caller is
|
||||||
// what a test binds is fixture data. scan_guard_test.go fails if a
|
// internal/database/database_test.go:91, whose SELECT 1 binds nothing.
|
||||||
// non-test file calls it.
|
// scan_guard_test.go fails if a 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(
|
||||||
|
|||||||
@@ -97,8 +97,7 @@ 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. An archive download holds it while
|
// stored names from the other.
|
||||||
// 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
|
||||||
|
|||||||
@@ -191,7 +191,6 @@ 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
|
||||||
}
|
}
|
||||||
@@ -205,11 +204,9 @@ 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, &dbMgr, &archives)
|
app := newTestApp(t, &h, &sess, &db, &archives)
|
||||||
app.RequireStart()
|
app.RequireStart()
|
||||||
|
|
||||||
t.Cleanup(app.RequireStop)
|
t.Cleanup(app.RequireStop)
|
||||||
@@ -217,7 +214,6 @@ 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",
|
||||||
|
|||||||
@@ -1,121 +0,0 @@
|
|||||||
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
|
|
||||||
}
|
|
||||||
@@ -1,379 +0,0 @@
|
|||||||
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)
|
|
||||||
}
|
|
||||||
@@ -37,18 +37,14 @@ const (
|
|||||||
editAuthHeader = "Authorization: Bearer " + editBearerSecret
|
editAuthHeader = "Authorization: Bearer " + editBearerSecret
|
||||||
)
|
)
|
||||||
|
|
||||||
// targetRouter mounts the target create, edit and download routes on
|
// targetRouter mounts the target create and edit routes on a chi
|
||||||
// a chi router so the handlers see the URL parameters they read.
|
// 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(),
|
||||||
|
|||||||
@@ -1,71 +0,0 @@
|
|||||||
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
|
|
||||||
}
|
|
||||||
@@ -1,56 +0,0 @@
|
|||||||
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"])
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,69 +0,0 @@
|
|||||||
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")
|
|
||||||
}
|
|
||||||
@@ -1,8 +1,6 @@
|
|||||||
package server
|
package server
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
|
||||||
"log/slog"
|
|
||||||
"net/http"
|
"net/http"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
@@ -39,14 +37,6 @@ 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.
|
||||||
|
|||||||
@@ -73,7 +73,7 @@ func (s *Server) setupGlobalMiddleware() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
s.router.Use(s.mw.CORS())
|
s.router.Use(s.mw.CORS())
|
||||||
s.router.Use(s.mw.Timeout(requestTimeout))
|
s.router.Use(middleware.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
|
||||||
@@ -312,10 +312,6 @@ 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(),
|
||||||
|
|||||||
@@ -39,12 +39,6 @@ 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
|
||||||
@@ -65,16 +59,6 @@ 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,
|
||||||
@@ -277,17 +261,10 @@ 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) {
|
||||||
drain := ShutdownTimeout
|
ctxShutdown, shutdownCancel := context.WithTimeout(
|
||||||
|
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)
|
||||||
|
|||||||
@@ -1,148 +1,13 @@
|
|||||||
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
|
||||||
|
|||||||
@@ -157,9 +157,6 @@
|
|||||||
{{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}}">
|
||||||
|
|||||||
Reference in New Issue
Block a user