1 Commits
Author SHA1 Message Date
sneak af804be45c Report or refuse each unusable file webhooker reads (closes #290)
check / check (push) Waiting to run
Audit of the files webhooker reads configuration or required state
from. A missing or zero-length database is reported with the "created
a new, empty database" warning and its path: webhooker.db at start,
and a per-webhook database at the latest at the next start, since
restart recovery now opens every webhook's database. The main
database's open errors name webhooker.db
(#459). webhooker resetpw
refuses a zero-length webhooker.db as it refuses a missing one. A
directory in place of a database file or its -wal or -shm is refused,
naming it; beside a -shm directory SQLite opened the database
read-only without a word. The README says how each case is treated.

Model: opus-5-5
2026-10-02 18:12:33 +00:00
30 changed files with 383 additions and 1517 deletions
+49 -44
View File
@@ -79,7 +79,8 @@ directory, read once at startup before anything else looks at the
environment. environment.
The file is optional and having none is the normal case for a The file is optional and having none is the normal case for a
deployment. A file that is there but cannot be parsed aborts startup deployment. An empty file is the same as none: it has nothing in it to
apply. A file that is there but cannot be parsed aborts startup
with a message naming it, because a single malformed line makes none with a message naming it, because a single malformed line makes none
of the file apply: every variable in it silently reverts to its of the file apply: every variable in it silently reverts to its
default, which is exactly the failure [Invalid values abort default, which is exactly the failure [Invalid values abort
@@ -564,11 +565,13 @@ its Argon2id hash. There is no second account and no forgot-password
flow, so the banner and the reset command below are the only two ways flow, so the banner and the reset command below are the only two ways
in. in.
A start that finds no `webhooker.db` in `DATA_DIR` also logs A start that finds no `webhooker.db` in `DATA_DIR`, or a zero-length
one (which SQLite opens as an empty database), also logs
`created a new, empty database` at `WARN`, with the file's path, `created a new, empty database` at `WARN`, with the file's path,
shortly before the banner. On a deployment that has run before, that shortly before the banner. On a deployment that has run before, that
line means `DATA_DIR` was empty, most often because its volume is not line means `webhooker.db` was lost: either `DATA_DIR` was empty, most
mounted. often because its volume is not mounted, or the file was zero-length,
as a truncated copy leaves it.
#### Recovering a lost admin password #### Recovering a lost admin password
@@ -612,7 +615,8 @@ What it will not do:
the old password, so a reset underneath it would report a change the the old password, so a reset underneath it would report a change the
service does not honour. service does not honour.
- **Create anything.** A `DATA_DIR` that does not exist, or that holds - **Create anything.** A `DATA_DIR` that does not exist, or that holds
no `webhooker.db`, is an error rather than a new empty deployment — no `webhooker.db` or a zero-length one, is an error naming the path
rather than a new empty deployment —
a mistyped path must not be built out and then reported as a success. a mistyped path must not be built out and then reported as a success.
- **Create an account.** A username that does not exist is an error. - **Create an account.** A username that does not exist is an error.
`resetpw` changes an existing account's password and nothing else. `resetpw` changes an existing account's password and nothing else.
@@ -993,6 +997,16 @@ its sidecars; a killed or crashed instance leaves them, and they must be
carried with the `.db`. An archive the service has not opened since a carried with the `.db`. An archive the service has not opened since a
crash keeps that crash's sidecars, even across a later clean stop. crash keeps that crash's sidecars, even across a later clean stop.
A missing sidecar is therefore normal, and SQLite makes new ones, so a
`-wal` lost from a copy cannot be reported: the transactions it held
are simply gone. SQLite reads a `-wal` up to its first damaged frame,
as after a crash, and rebuilds a damaged `-shm`. A sidecar with the
wrong mode is set back to `0600` when its database is opened. A
directory in place of either is refused then, with an error naming it:
for `webhooker.db` the server and `webhooker resetpw` stop, and an event
or archive database fails as a damaged one does (see
[Database Architecture](#database-architecture)).
Configuration is **not** in `DATA_DIR` — it comes from the environment Configuration is **not** in `DATA_DIR` — it comes from the environment
and from a `.env` file read out of the process working directory. Back and from a `.env` file read out of the process working directory. Back
that up with your deployment config, separately. that up with your deployment config, separately.
@@ -1075,13 +1089,17 @@ with any `-wal`/`-shm` beside it, or wait until there are none.
1. Stop the service. 1. Stop the service.
2. Restore the **whole set together**: `webhooker.db` *and* every 2. Restore the **whole set together**: `webhooker.db` *and* every
`events-*.db` *and* every `archive-*.db`. A partial restore fails `events-*.db` *and* every `archive-*.db`. A partial restore is
quietly rather than loudly. Every database is opened `mode=rwc`, so a reported, not refused. Every database is opened `mode=rwc`, so a
missing `events-{uuid}.db` is **created empty** on first access missing `events-{uuid}.db` is **created empty**: the webhook comes
instead of erroring — the webhook comes back with its configuration back with its configuration intact and its entire event history
intact and its entire event history silently gone. Event databases gone. The first start after the restore logs
restored without `webhooker.db` are simply orphaned; nothing `created a new, empty database` at `WARN` for each such file, with
references their UUIDs. its path, as it does for a missing `webhooker.db`. A missing
`archive-*.db` is recreated at its target's next delivery without a
warning, since moving one away is a supported workflow. Event
databases restored without `webhooker.db` are simply orphaned;
nothing references their UUIDs.
3. Carry any `*.db-wal` and `*.db-shm` files that are in the backup. 3. Carry any `*.db-wal` and `*.db-shm` files that are in the backup.
They are part of the database, and dropping a `-wal` silently They are part of the database, and dropping a `-wal` silently
@@ -1939,10 +1957,19 @@ encryption key is generated and stored, and an `admin` user is created.
the deliveries per target, kept through retention the deliveries per target, kept through retention
Per-webhook databases are created automatically when a webhook is Per-webhook databases are created automatically when a webhook is
created (and lazily on first access for webhooks that predate this created. They are managed by the `WebhookDBManager` component, which
feature). They are managed by the `WebhookDBManager` component, which
handles connection pooling, lazy opening, migrations, and cleanup. handles connection pooling, lazy opening, migrations, and cleanup.
A per-webhook database that is missing or zero-length later means its
webhook's events and pending deliveries are gone. The next time it is
opened, an empty one is created in its place, so the webhook keeps
receiving, and `created a new, empty database` is logged at `WARN` with
the file's path. Every webhook's database is opened when the service
starts, so this appears at the latest at the first start after the
file was lost. A file there that SQLite cannot open fails that
webhook alone, with an `ERROR` naming the webhook on every access and a
500 to its senders, so one damaged file does not stop the others.
This separation provides: This separation provides:
- **Isolation** — a high-volume webhook won't cause lock contention or - **Isolation** — a high-volume webhook won't cause lock contention or
@@ -2007,7 +2034,9 @@ After each write the archive handle is closed
and reopened, debounced to at most once per second, so an operator can and reopened, debounced to at most once per second, so an operator can
move the archive file away for offline archiving without stopping the move the archive file away for offline archiving without stopping the
service; a moved or removed archive file is recreated automatically on service; a moved or removed archive file is recreated automatically on
the next write. An optional `expiry` in the target's config JSON (e.g. the next write. A zero-length archive file is written to as a new
archive: SQLite opens it as an empty database, so it holds nothing to
lose. An optional `expiry` in the target's config JSON (e.g.
`{"expiry":"720h"}`) is validated when the target is created — the `{"expiry":"720h"}`) is validated when the target is created — the
default (unset or the literal `never`) keeps rows forever — and rows default (unset or the literal `never`) keeps rows forever — and rows
older than the expiry are pruned each time the archive is (re)opened. An older than the expiry are pruned each time the archive is (re)opened. An
@@ -2031,29 +2060,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
@@ -2632,10 +2638,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
@@ -2920,7 +2927,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 |
@@ -3002,7 +3008,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
@@ -4,6 +4,7 @@ import (
"bytes" "bytes"
"context" "context"
"log/slog" "log/slog"
"os"
"path/filepath" "path/filepath"
"strings" "strings"
"testing" "testing"
@@ -119,3 +120,26 @@ func TestNewDatabase_IsLoggedWithItsPath(t *testing.T) {
t, second, created, "an existing database is not new", t, second, created, "an existing database is not new",
) )
} }
// TestZeroLengthDatabase_IsLoggedAsNew covers what
// https://git.eeqj.de/sneak/webhooker/issues/290 found: SQLite opens a
// zero-length file as an empty database, so a start on one is a first
// start, and it must say so exactly as a start with no file does.
func TestZeroLengthDatabase_IsLoggedAsNew(t *testing.T) {
t.Parallel()
dir := t.TempDir()
path := filepath.Join(dir, database.MainDBFileName)
require.NoError(t, os.WriteFile(path, nil, database.SQLiteFilePerm))
var out bytes.Buffer
db, err := database.Open(dir, slog.New(slog.NewTextHandler(&out, nil)))
require.NoError(t, err)
require.NoError(t, db.Close())
assert.Contains(
t, out.String(),
`level=WARN msg="created a new, empty database" path=`+path,
)
}
+12 -7
View File
@@ -8,7 +8,6 @@ import (
"errors" "errors"
"fmt" "fmt"
"io" "io"
"io/fs"
"log/slog" "log/slog"
"os" "os"
"path/filepath" "path/filepath"
@@ -203,8 +202,7 @@ func (d *Database) connectTo(dataDir string) error {
// Checked before opening, which creates the file. A DATA_DIR that // Checked before opening, which creates the file. A DATA_DIR that
// is unexpectedly empty -- its volume not mounted, say -- looks // is unexpectedly empty -- its volume not mounted, say -- looks
// exactly like a first start, so a new database is a warning. // exactly like a first start, so a new database is a warning.
_, statErr := os.Stat(dbPath) created := missingOrEmpty(dbPath)
created := errors.Is(statErr, fs.ErrNotExist)
// Opened through OpenSQLite so this handle carries the same WAL // Opened through OpenSQLite so this handle carries the same WAL
// journaling, busy timeout, immediate-transaction locking, and pool // journaling, busy timeout, immediate-transaction locking, and pool
@@ -213,13 +211,15 @@ func (d *Database) connectTo(dataDir string) error {
if err != nil { if err != nil {
d.log.Error( d.log.Error(
"failed to open database", "failed to open database",
"path", dbPath,
"error", err, "error", err,
) )
return err return err
} }
// Then use it with GORM // Then use it with GORM. Its errors are SQLite's alone and name no
// file, so the path is added to them here.
db, err := gorm.Open(sqlite.Dialector{ db, err := gorm.Open(sqlite.Dialector{
Conn: sqlDB, Conn: sqlDB,
}, &gorm.Config{ }, &gorm.Config{
@@ -229,10 +229,11 @@ func (d *Database) connectTo(dataDir string) error {
if err != nil { if err != nil {
d.log.Error( d.log.Error(
"failed to connect to database", "failed to connect to database",
"path", dbPath,
"error", err, "error", err,
) )
return err return fmt.Errorf("connecting to %s: %w", dbPath, err)
} }
d.db = db d.db = db
@@ -243,8 +244,12 @@ func (d *Database) connectTo(dataDir string) error {
d.log.Info("connected to database", "path", dbPath) d.log.Info("connected to database", "path", dbPath)
} }
// Run migrations err = d.migrate()
return d.migrate() if err != nil {
return fmt.Errorf("migrating %s: %w", dbPath, err)
}
return nil
} }
func (d *Database) migrate() error { func (d *Database) migrate() error {
+25
View File
@@ -1,9 +1,15 @@
package database_test package database_test
import ( import (
"bytes"
"context" "context"
"log/slog"
"os"
"path/filepath"
"testing" "testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.uber.org/fx/fxtest" "go.uber.org/fx/fxtest"
"sneak.berlin/go/webhooker/internal/config" "sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
@@ -100,3 +106,22 @@ func TestDatabaseConnection(t *testing.T) {
) )
} }
} }
// TestOpen_UnreadableDatabaseIsNamed pins
// https://git.eeqj.de/sneak/webhooker/issues/459: when SQLite cannot
// read webhooker.db, the error that stops the server and `webhooker
// resetpw` names the file, not only SQLite's own message.
func TestOpen_UnreadableDatabaseIsNamed(t *testing.T) {
t.Parallel()
dir := t.TempDir()
path := filepath.Join(dir, database.MainDBFileName)
require.NoError(t, os.WriteFile(
path, bytes.Repeat([]byte("junk"), 1024), database.SQLiteFilePerm,
))
_, err := database.Open(dir, slog.New(slog.DiscardHandler))
require.Error(t, err)
assert.Contains(t, err.Error(), path)
assert.Contains(t, err.Error(), "file is not a database")
}
+2 -2
View File
@@ -184,8 +184,8 @@ func (r *RetentionReaper) sweep(ctx context.Context) {
wh := webhooks[i] wh := webhooks[i]
// Nothing to reap if the per-webhook database has never // A missing database has nothing to reap. Restart recovery
// been created. // reports a lost one (see WebhookDBManager.GetDB).
if !r.dbManager.DBExists(wh.ID) { if !r.dbManager.DBExists(wh.ID) {
continue continue
} }
+23
View File
@@ -182,6 +182,29 @@ func TestOpenSQLiteTightensFilesLeftWorldReadable(t *testing.T) {
requireDatabaseSetOwnerOnly(t, path) requireDatabaseSetOwnerOnly(t, path)
} }
// TestOpenSQLiteRefusesADirectorySidecar covers a directory in place
// of -wal or -shm. Beside a -shm directory SQLite opens the database
// read-only without a word, and every write then fails naming no file,
// so the open must stop instead, naming the directory.
func TestOpenSQLiteRefusesADirectorySidecar(t *testing.T) {
t.Parallel()
for _, suffix := range []string{"-wal", "-shm"} {
t.Run(suffix, func(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), database.MainDBFileName)
require.NoError(t, os.Mkdir(path+suffix, 0o700))
_, err := database.OpenSQLite(
path, database.SQLiteModeCreate,
)
require.Error(t, err)
assert.Contains(t, err.Error(), path+suffix)
})
}
}
// TestOpenSQLiteExistingModeDoesNotCreateTheFile guards the mechanism // TestOpenSQLiteExistingModeDoesNotCreateTheFile guards the mechanism
// the fix uses: OpenSQLite now creates the database file itself, and // the fix uses: OpenSQLite now creates the database file itself, and
// must not do so for a caller that asked for an existing database. An // must not do so for a caller that asked for an existing database. An
+26 -2
View File
@@ -7,6 +7,7 @@ import (
"io/fs" "io/fs"
"net/url" "net/url"
"os" "os"
"syscall"
"time" "time"
_ "modernc.org/sqlite" // Pure Go SQLite driver _ "modernc.org/sqlite" // Pure Go SQLite driver
@@ -93,7 +94,8 @@ const (
const SQLiteFilePerm fs.FileMode = 0o600 const SQLiteFilePerm fs.FileMode = 0o600
// reserveSQLiteFile puts path at SQLiteFilePerm before the driver ever // reserveSQLiteFile puts path at SQLiteFilePerm before the driver ever
// touches it, and tightens any sidecar already on disk. // touches it, and tightens any sidecar already on disk. A directory in
// place of any of them is an error naming it.
// //
// The mode has to be settled here rather than by a chmod after opening, // The mode has to be settled here rather than by a chmod after opening,
// because SQLite picks it: robust_open substitutes // because SQLite picks it: robust_open substitutes
@@ -143,7 +145,15 @@ func reserveSQLiteFile(path string, create bool) error {
for _, p := range append( for _, p := range append(
[]string{path}, sqliteSidecarPaths(path)..., []string{path}, sqliteSidecarPaths(path)...,
) { ) {
err := os.Chmod(p, SQLiteFilePerm) // Chmod accepts a directory, and SQLite opens a database whose
// -shm is one read-only, without a word: every write then
// fails naming no file.
info, err := os.Stat(p)
if err == nil && info.IsDir() {
return fmt.Errorf("securing %s: %w", p, syscall.EISDIR)
}
err = os.Chmod(p, SQLiteFilePerm)
if err != nil && !errors.Is(err, fs.ErrNotExist) { if err != nil && !errors.Is(err, fs.ErrNotExist) {
return fmt.Errorf("securing %s: %w", p, err) return fmt.Errorf("securing %s: %w", p, err)
} }
@@ -152,6 +162,20 @@ func reserveSQLiteFile(path string, create bool) error {
return nil return nil
} }
// missingOrEmpty reports whether opening path in SQLiteModeCreate
// would start a new, empty database: the file is not there, or it is
// zero-length, which SQLite opens as an empty database. A file left at
// zero length by an interrupted first start or a truncated copy holds
// as little as a missing one, and must be reported the same way.
func missingOrEmpty(path string) bool {
info, err := os.Stat(path)
if errors.Is(err, fs.ErrNotExist) {
return true
}
return err == nil && info.Size() == 0
}
// sqliteSidecarPaths returns the files SQLite maintains beside a // sqliteSidecarPaths returns the files SQLite maintains beside a
// database under WAL. They carry the same rows as the database itself, // database under WAL. They carry the same rows as the database itself,
// so a fix that tightens only the main file has fixed nothing. // so a fix that tightens only the main file has fixed nothing.
+54 -28
View File
@@ -98,34 +98,18 @@ func NewWebhookDBManager(
return m, nil return m, nil
} }
// GetDB returns the database connection for a webhook, // GetDB returns the database connection for a webhook, opening it on
// creating the database file lazily if it doesn't exist. // first use.
//
// The file is made by CreateDB when the webhook is created. One that is
// missing or zero-length here means the webhook's events and pending
// deliveries are gone: an empty database is created in its place so
// the webhook keeps receiving, and that is logged as a warning naming
// the file, as a new main database is.
func (m *WebhookDBManager) GetDB( func (m *WebhookDBManager) GetDB(
webhookID string, webhookID string,
) (*gorm.DB, error) { ) (*gorm.DB, error) {
// Fast path: already open return m.getDB(webhookID, false)
if val, ok := m.dbs.Load(webhookID); ok {
return asGormDB(val, webhookID)
}
// Slow path: open the database under the lock, looking in the
// cache again first. A caller that raced another one here then
// waits for its handle instead of opening a second one.
m.mu.Lock()
defer m.mu.Unlock()
if val, ok := m.dbs.Load(webhookID); ok {
return asGormDB(val, webhookID)
}
db, err := m.openDB(webhookID)
if err != nil {
return nil, err
}
m.dbs.Store(webhookID, db)
return db, nil
} }
// asGormDB returns a value read from the cache as the database // asGormDB returns a value read from the cache as the database
@@ -143,12 +127,12 @@ func asGormDB(val any, webhookID string) (*gorm.DB, error) {
return db, nil return db, nil
} }
// CreateDB explicitly creates a new per-webhook database file // CreateDB creates a new webhook's database file and runs
// and runs migrations. // migrations.
func (m *WebhookDBManager) CreateDB( func (m *WebhookDBManager) CreateDB(
webhookID string, webhookID string,
) error { ) error {
_, err := m.GetDB(webhookID) _, err := m.getDB(webhookID, true)
return err return err
} }
@@ -266,6 +250,48 @@ func (m *WebhookDBManager) DBPath(
return m.dbPath(webhookID) return m.dbPath(webhookID)
} }
// getDB is GetDB, and CreateDB when isNew is true: the webhook has just
// been created, so a missing file is expected rather than lost.
func (m *WebhookDBManager) getDB(
webhookID string, isNew bool,
) (*gorm.DB, error) {
// Fast path: already open
if val, ok := m.dbs.Load(webhookID); ok {
return asGormDB(val, webhookID)
}
// Slow path: open the database under the lock, looking in the
// cache again first. A caller that raced another one here then
// waits for its handle instead of opening a second one.
m.mu.Lock()
defer m.mu.Unlock()
if val, ok := m.dbs.Load(webhookID); ok {
return asGormDB(val, webhookID)
}
// Checked before opening, which creates the file. See GetDB.
path := m.dbPath(webhookID)
replaced := !isNew && missingOrEmpty(path)
db, err := m.openDB(webhookID)
if err != nil {
return nil, err
}
if replaced {
m.log.Warn(
"created a new, empty database",
"webhook_id", webhookID,
"path", path,
)
}
m.dbs.Store(webhookID, db)
return db, nil
}
func (m *WebhookDBManager) dbPath( func (m *WebhookDBManager) dbPath(
webhookID string, webhookID string,
) string { ) string {
@@ -289,6 +289,75 @@ func TestWebhookDBManager_LazyCreation(t *testing.T) {
assert.True(t, mgr.DBExists(webhookID)) assert.True(t, mgr.DBExists(webhookID))
} }
// A webhook's database is made by CreateDB along with the webhook. One
// that GetDB finds missing or zero-length has lost the webhook's events
// and pending deliveries, so the empty database made in its place is
// logged as a warning naming the file
// (https://git.eeqj.de/sneak/webhooker/issues/290). CreateDB, and
// reopening a database that is there, log no such warning.
func TestWebhookDBManager_LostDatabaseIsLogged(t *testing.T) {
t.Parallel()
const created = `level=WARN msg="created a new, empty database"`
open := func(
t *testing.T, prepare func(*database.WebhookDBManager, string),
) (string, string) {
t.Helper()
var logs bytes.Buffer
mgr := database.NewTestWebhookDBManagerWithLogger(
t.TempDir(),
slog.New(slog.NewTextHandler(&logs, nil)),
)
webhookID := uuid.New().String()
prepare(mgr, webhookID)
_, err := mgr.GetDB(webhookID)
require.NoError(t, err)
require.NoError(t, mgr.CloseAll())
return logs.String(),
" webhook_id=" + webhookID + " path=" + mgr.DBPath(webhookID)
}
t.Run("missing", func(t *testing.T) {
t.Parallel()
logs, fields := open(
t, func(*database.WebhookDBManager, string) {},
)
assert.Contains(t, logs, created+fields)
})
t.Run("zero-length", func(t *testing.T) {
t.Parallel()
logs, fields := open(
t, func(mgr *database.WebhookDBManager, webhookID string) {
require.NoError(t, os.WriteFile(
mgr.DBPath(webhookID), nil, database.SQLiteFilePerm,
))
},
)
assert.Contains(t, logs, created+fields)
})
t.Run("created with the webhook, then reopened", func(t *testing.T) {
t.Parallel()
logs, _ := open(
t, func(mgr *database.WebhookDBManager, webhookID string) {
require.NoError(t, mgr.CreateDB(webhookID))
require.NoError(t, mgr.CloseAll())
},
)
assert.NotContains(t, logs, created)
})
}
func TestWebhookDBManager_DeliveryWorkflow(t *testing.T) { func TestWebhookDBManager_DeliveryWorkflow(t *testing.T) {
t.Parallel() t.Parallel()
+3 -7
View File
@@ -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
+3 -4
View File
@@ -699,10 +699,9 @@ func (e *Engine) recoverInFlight(ctx context.Context) {
default: default:
} }
if !e.dbManager.DBExists(webhookID) { // Opened even when its file is missing, so that GetDB reports
continue // a lost database at start, not when the webhook next receives
} // an event, which for a quiet webhook may be never.
e.recoverWebhookDeliveries(ctx, webhookID) e.recoverWebhookDeliveries(ctx, webhookID)
} }
} }
+36 -3
View File
@@ -1,6 +1,7 @@
package delivery_test package delivery_test
import ( import (
"bytes"
"context" "context"
"encoding/json" "encoding/json"
"fmt" "fmt"
@@ -23,7 +24,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 +81,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)
@@ -1136,6 +1135,40 @@ func TestRecoverInFlight_WithPendingDeliveries(
} }
} }
// TestRecoverInFlight_ReportsAMissingWebhookDatabase covers a webhook
// whose database file is gone, after a partial restore say. Restart
// recovery opens every webhook's database, so the empty one made in its
// place is reported at start, naming the file
// (https://git.eeqj.de/sneak/webhooker/issues/290).
func TestRecoverInFlight_ReportsAMissingWebhookDatabase(t *testing.T) {
t.Parallel()
mainDB := iMainDB(t)
webhookID := uuid.New().String()
iCreateWebhook(t, mainDB, webhookID, "lost-database")
var logs bytes.Buffer
dbMgr := database.NewTestWebhookDBManagerWithLogger(
t.TempDir(), slog.New(slog.NewTextHandler(&logs, nil)),
)
t.Cleanup(func() { _ = dbMgr.CloseAll() })
engine := delivery.NewTestEngineWithDB(
database.NewTestDatabase(mainDB), dbMgr,
slog.New(slog.DiscardHandler),
&http.Client{Timeout: 5 * time.Second}, 1,
)
engine.ExportRecoverInFlight(context.Background())
assert.Contains(
t, logs.String(),
`level=WARN msg="created a new, empty database" webhook_id=`+
webhookID+" path="+dbMgr.DBPath(webhookID),
)
}
// --- HTTP Config with custom headers --- // --- HTTP Config with custom headers ---
func TestDeliverHTTP_CustomTargetHeaders(t *testing.T) { func TestDeliverHTTP_CustomTargetHeaders(t *testing.T) {
+1 -3
View File
@@ -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)
+10 -6
View File
@@ -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
-275
View File
@@ -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),
),
)
}
+1 -3
View File
@@ -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)
+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; 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(
+1 -2
View File
@@ -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
+1 -5
View File
@@ -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",
-121
View File
@@ -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
}
-379
View File
@@ -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)
}
+2 -6
View File
@@ -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(),
-71
View File
@@ -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
}
-56
View File
@@ -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"])
})
}
}
+10 -1
View File
@@ -308,7 +308,7 @@ func checkDataDir(dir string) error {
dbPath := filepath.Join(dir, database.MainDBFileName) dbPath := filepath.Join(dir, database.MainDBFileName)
_, err = os.Stat(dbPath) dbInfo, err := os.Stat(dbPath)
switch { switch {
case errors.Is(err, fs.ErrNotExist): case errors.Is(err, fs.ErrNotExist):
@@ -319,6 +319,15 @@ func checkDataDir(dir string) error {
) )
case err != nil: case err != nil:
return fmt.Errorf("checking %s: %w", dbPath, err) return fmt.Errorf("checking %s: %w", dbPath, err)
case dbInfo.Size() == 0:
// SQLite opens a zero-length file as an empty database, so
// it holds no deployment either, and opening it would write
// an empty schema into it.
return fmt.Errorf(
"%w: %s is zero-length. The admin account is created by "+
"the first server start",
ErrNoDatabase, dbPath,
)
} }
return nil return nil
+27
View File
@@ -377,6 +377,33 @@ func TestMissingDatabaseCreatesNothing(t *testing.T) {
) )
} }
// TestZeroLengthDatabaseCreatesNothing covers a webhooker.db left at
// zero length, as a truncated copy leaves it. SQLite would open it as
// an empty database, so it is refused like a missing one and left as
// it is.
func TestZeroLengthDatabaseCreatesNothing(t *testing.T) {
dir := t.TempDir()
t.Setenv("DATA_DIR", dir)
dbPath := filepath.Join(dir, database.MainDBFileName)
require.NoError(
t, os.WriteFile(dbPath, nil, database.SQLiteFilePerm),
)
code, _, stderr := run(t, newPassword+"\n", operatorUser)
require.Equal(t, exitFailure, code)
assert.Contains(t, stderr, dbPath)
entries, err := os.ReadDir(dir)
require.NoError(t, err)
assert.Len(t, entries, 1, "nothing may be created beside it")
info, err := os.Stat(dbPath)
require.NoError(t, err)
assert.Zero(t, info.Size(), "nothing may be written into it")
}
// TestUnknownUserFails states the decision: resetpw changes an // TestUnknownUserFails states the decision: resetpw changes an
// existing account's password and never creates an account. A typo in // existing account's password and never creates an account. A typo in
// the username must say so rather than quietly adding a second user. // the username must say so rather than quietly adding a second user.
-69
View File
@@ -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 -5
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(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(),
-3
View File
@@ -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}}">