Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
af804be45c |
@@ -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,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|||||||
@@ -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 {
|
||||||
|
|||||||
@@ -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")
|
||||||
|
}
|
||||||
|
|||||||
@@ -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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
@@ -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()
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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) {
|
||||||
|
|||||||
@@ -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"])
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -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
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
@@ -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")
|
|
||||||
}
|
|
||||||
@@ -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(),
|
||||||
|
|||||||
@@ -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