Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e55134532c | ||
|
|
bcdd4791ec | ||
|
|
f1da5e73dd | ||
|
|
d2ecb83923 | ||
|
|
643077021d |
@@ -1368,12 +1368,13 @@ hides until the form closes, Cancel hides the form and drops what was typed, as
|
|||||||
does leaving the page and going back to it, and Save changes the description;
|
does leaving the page and going back to it, and Save changes the description;
|
||||||
of the recent events on the webhook page only the newest starts expanded, each
|
of the recent events on the webhook page only the newest starts expanded, each
|
||||||
expands and collapses, and Open leads to the event's own page; an event in the
|
expands and collapses, and Open leads to the event's own page; an event in the
|
||||||
event log expands and collapses, and so do a delivery's attempts inside it; and
|
event log expands and collapses when its row's caret or its ID is clicked, and
|
||||||
at phone width the menu button opens and closes the mobile menu. It also fails
|
from the keyboard, but not when its ID is selected with the mouse, and a
|
||||||
if the browser reports a console warning or error, an uncaught exception, or
|
delivery's attempts inside it expand and collapse; and at phone width the menu
|
||||||
anything the policy refused. `make check`
|
button opens and closes the mobile menu. It also fails if the browser reports a
|
||||||
and the image build lint it but do not run it, and `make test` leaves it out
|
console warning or error, an uncaught exception, or anything the policy refused.
|
||||||
(its file is built only with the `browser` build tag). Run it with
|
`make check` and the image build lint it but do not run it, and `make test`
|
||||||
|
leaves it out (its file is built only with the `browser` build tag). Run it with
|
||||||
`make test-browser` after changing `templates/` or `static/js/`: that builds
|
`make test-browser` after changing `templates/` or `static/js/`: that builds
|
||||||
`Dockerfile.browser`, which runs the test in a digest-pinned headless browser
|
`Dockerfile.browser`, which runs the test in a digest-pinned headless browser
|
||||||
image, so the host needs no browser.
|
image, so the host needs no browser.
|
||||||
@@ -1862,7 +1863,11 @@ deliver where the destination has since been fixed. A target that has
|
|||||||
been deleted or deactivated therefore refuses the replay with a
|
been deleted or deactivated therefore refuses the replay with a
|
||||||
message on the event log rather than delivering from stale
|
message on the event log rather than delivering from stale
|
||||||
configuration, and a replay is refused while an earlier one for the
|
configuration, and a replay is refused while an earlier one for the
|
||||||
same event and target is still pending or retrying.
|
same event and target is still pending or retrying. A delivery whose
|
||||||
|
target has been deleted shows no **Replay** action at all: recreating
|
||||||
|
the target makes a new one that the old delivery does not name, so
|
||||||
|
**Resubmit** is how that event reaches the webhook's currently active
|
||||||
|
targets.
|
||||||
|
|
||||||
**Resubmit.** Replay recovers one delivery; **resubmit** re-injects one
|
**Resubmit.** Replay recovers one delivery; **resubmit** re-injects one
|
||||||
EVENT. The event log offers a per-event **Resubmit** action that stores
|
EVENT. The event log offers a per-event **Resubmit** action that stores
|
||||||
@@ -1900,6 +1905,10 @@ retries) is individually logged for full observability.
|
|||||||
| `error` | string | Error message (on failure) |
|
| `error` | string | Error message (on failure) |
|
||||||
| `duration` | integer | Request duration in milliseconds |
|
| `duration` | integer | Request duration in milliseconds |
|
||||||
|
|
||||||
|
A `database` or `log` target sends no HTTP request, so in the event log and
|
||||||
|
on the event's page its attempts show no status: a successful one reads
|
||||||
|
"archived" or "written to the log".
|
||||||
|
|
||||||
**Relations:** Belongs to Delivery.
|
**Relations:** Belongs to Delivery.
|
||||||
|
|
||||||
#### EventTotals, TargetTotals and EntrypointTotals
|
#### EventTotals, TargetTotals and EntrypointTotals
|
||||||
@@ -2158,10 +2167,12 @@ the move-the-file-away workflow keeps working. Archives with no expiry,
|
|||||||
or the expiry `never`, are not touched by the sweep at all.
|
or the expiry `never`, are not touched by the sweep at all.
|
||||||
|
|
||||||
For a target with several files, the write path prunes only the file it
|
For a target with several files, the write path prunes only the file it
|
||||||
writes to, and the sweep prunes every one of them. A file named for a
|
writes to, and the sweep prunes every one of them, taking the target's
|
||||||
period that the sweep leaves empty is deleted, with any `-wal` and
|
lock for one file at a time, so a write to the target waits for at most
|
||||||
`-shm` beside it; the file without a period is kept even when empty, as
|
one file's prune. A file that is gone by the time the sweep reaches it
|
||||||
it always has been.
|
is skipped. A file named for a period that the sweep leaves empty is
|
||||||
|
deleted, with any `-wal` and `-shm` beside it; the file without a period
|
||||||
|
is kept even when empty, as it always has been.
|
||||||
|
|
||||||
Because each `database` target has its own archive file, a target's
|
Because each `database` target has its own archive file, a target's
|
||||||
`expiry` governs only its own archive. Two `database` targets on one
|
`expiry` governs only its own archive. Two `database` targets on one
|
||||||
@@ -2196,16 +2207,20 @@ file.
|
|||||||
|
|
||||||
The download streams: each row is read and written out compressed
|
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
|
before the next is read, so neither the archive nor the JSON is held in
|
||||||
memory. When it starts it opens every one of the target's files, each on
|
memory. When it starts it lists the target's files by the stored names,
|
||||||
a connection of its own inside one read-only transaction, so the
|
under the lock that webhook edits, target edits and target creation
|
||||||
download holds the archive as it stood then, and archive writes go on
|
hold, and lets go. It then opens one file at a time, only when its rows
|
||||||
meanwhile, since under WAL a reader never blocks a writer. Each file is
|
are about to be written out, and closes it before it opens the next, so
|
||||||
closed once its rows are written out; until then its `-wal` cannot be
|
it never has more than one of the target's files open. To open each, it
|
||||||
checkpointed past what the download reads, so a long download lets the
|
takes the lock again just long enough to find the file by its period
|
||||||
`-wal` grow. It finds the files by the stored names under the lock that
|
under the names stored then, so a rename during the download loses no
|
||||||
webhook edits, target edits and target creation hold, and lets go once
|
file; a file that is gone by then, emptied by the sweep or moved away,
|
||||||
the files are open: a rename during the download moves the files without
|
is skipped. Each file is read on a connection of its own inside one
|
||||||
affecting it.
|
read-only transaction, so it is written out as it stood when it was
|
||||||
|
opened, and archive writes go on meanwhile, since under WAL a reader
|
||||||
|
never blocks a writer. Until the open file is closed its `-wal` cannot
|
||||||
|
be checkpointed past what the download reads, so a long download lets
|
||||||
|
that `-wal` grow.
|
||||||
|
|
||||||
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
|
||||||
@@ -2402,6 +2417,20 @@ just delayed until the target is healthy again. A delivery already in
|
|||||||
`retrying` keeps that status without another database write each time
|
`retrying` keeps that status without another database write each time
|
||||||
the breaker turns it away.
|
the breaker turns it away.
|
||||||
|
|
||||||
|
While a target's breaker is open, the target's row on the webhook page
|
||||||
|
says its deliveries are paused until the cooldown ends, in UTC and as a
|
||||||
|
time from now. Each of its `retrying` deliveries shows as waiting in the
|
||||||
|
event log and on the event's page, with the earliest it can be tried
|
||||||
|
next: the later of the cooldown's end and the end of its own backoff
|
||||||
|
after its last attempt. It is only the earliest: when the cooldown ends,
|
||||||
|
one of the target's waiting deliveries is sent to test it while the
|
||||||
|
others wait at least one more cooldown, as the row also says. A time not
|
||||||
|
on the current UTC day is shown with its date. While the breaker is
|
||||||
|
half-open, the row says instead that deliveries are held while one
|
||||||
|
delivery tests whether the target has recovered, with no time, and the
|
||||||
|
target's deliveries show their plain status, since any of them may be
|
||||||
|
the one being sent.
|
||||||
|
|
||||||
### Metrics
|
### Metrics
|
||||||
|
|
||||||
`/metrics` serves one Prometheus registry behind basic auth (see
|
`/metrics` serves one Prometheus registry behind basic auth (see
|
||||||
|
|||||||
@@ -212,6 +212,10 @@ func newApp() *fx.App {
|
|||||||
// or renaming a webhook or target reaches its archive
|
// or renaming a webhook or target reaches its archive
|
||||||
// files.
|
// files.
|
||||||
func(e *delivery.Engine) delivery.Archives { return e },
|
func(e *delivery.Engine) delivery.Archives { return e },
|
||||||
|
// Wire *delivery.Engine as delivery.CircuitBreakers so
|
||||||
|
// the pages can show a target whose deliveries are
|
||||||
|
// paused.
|
||||||
|
func(e *delivery.Engine) delivery.CircuitBreakers { return e },
|
||||||
server.New,
|
server.New,
|
||||||
),
|
),
|
||||||
fx.Invoke(
|
fx.Invoke(
|
||||||
|
|||||||
@@ -102,6 +102,20 @@ func (cb *CircuitBreaker) CooldownRemaining() time.Duration {
|
|||||||
return remaining
|
return remaining
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// StateAndCooldown returns the circuit state and, while the circuit is
|
||||||
|
// open, what is left of the cooldown, or zero once that has passed.
|
||||||
|
// Both are read under one lock, so they always agree.
|
||||||
|
func (cb *CircuitBreaker) StateAndCooldown() (CircuitState, time.Duration) {
|
||||||
|
cb.mu.Lock()
|
||||||
|
defer cb.mu.Unlock()
|
||||||
|
|
||||||
|
if cb.state != CircuitOpen {
|
||||||
|
return cb.state, 0
|
||||||
|
}
|
||||||
|
|
||||||
|
return cb.state, max(cb.cooldown-time.Since(cb.lastFailure), 0)
|
||||||
|
}
|
||||||
|
|
||||||
// RecordSuccess records a successful delivery and resets
|
// RecordSuccess records a successful delivery and resets
|
||||||
// the circuit breaker to closed state.
|
// the circuit breaker to closed state.
|
||||||
func (cb *CircuitBreaker) RecordSuccess() {
|
func (cb *CircuitBreaker) RecordSuccess() {
|
||||||
|
|||||||
@@ -143,6 +143,15 @@ type Archives interface {
|
|||||||
Rename(targetID, webhookName, targetName string) error
|
Rename(targetID, webhookName, targetName string) error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// CircuitBreakers is how the handlers read a target's circuit
|
||||||
|
// breaker, so the webhook page and the event log can say that
|
||||||
|
// deliveries to the target are paused and until when. Like Archives,
|
||||||
|
// it keeps the handlers free of the engine's internals and is
|
||||||
|
// trivially faked in tests.
|
||||||
|
type CircuitBreakers interface {
|
||||||
|
StateAndCooldown(targetID string) (CircuitState, time.Duration)
|
||||||
|
}
|
||||||
|
|
||||||
// EngineParams are the fx dependencies for the delivery
|
// EngineParams are the fx dependencies for the delivery
|
||||||
// engine.
|
// engine.
|
||||||
type EngineParams struct {
|
type EngineParams struct {
|
||||||
@@ -186,9 +195,11 @@ type Engine struct {
|
|||||||
// targets maps each target type to its implementation.
|
// targets maps each target type to its implementation.
|
||||||
targets map[database.TargetType]Target
|
targets map[database.TargetType]Target
|
||||||
|
|
||||||
// httpTarget is retained so tests can reach the HTTP
|
// httpTarget and slackTarget are retained so StateAndCooldown
|
||||||
// target's shared client and circuit breakers.
|
// can read their circuit breakers, and so tests can reach the
|
||||||
httpTarget *httpTarget
|
// HTTP target's shared client.
|
||||||
|
httpTarget *httpTarget
|
||||||
|
slackTarget *slackTarget
|
||||||
|
|
||||||
// dbTarget is retained so the engine can reach the archive
|
// dbTarget is retained so the engine can reach the archive
|
||||||
// writer registry for eviction, renames and the idle sweep.
|
// writer registry for eviction, renames and the idle sweep.
|
||||||
@@ -301,6 +312,28 @@ func (e *Engine) Rename(
|
|||||||
return e.dbTarget.rename(targetID, webhookName, targetName)
|
return e.dbTarget.rename(targetID, webhookName, targetName)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// StateAndCooldown implements CircuitBreakers. It returns the state of
|
||||||
|
// the target's circuit breaker and, while the breaker is open, what is
|
||||||
|
// left of its cooldown; the cooldown is zero once that has passed and
|
||||||
|
// in any other state. A target with no breaker reads as closed with no
|
||||||
|
// cooldown, and reading never creates one.
|
||||||
|
func (e *Engine) StateAndCooldown(
|
||||||
|
targetID string,
|
||||||
|
) (CircuitState, time.Duration) {
|
||||||
|
for _, core := range []*httpCore{
|
||||||
|
e.httpTarget.httpCore, e.slackTarget.httpCore,
|
||||||
|
} {
|
||||||
|
val, ok := core.circuitBreakers.Load(targetID)
|
||||||
|
if ok {
|
||||||
|
cb, _ := val.(*CircuitBreaker)
|
||||||
|
|
||||||
|
return cb.StateAndCooldown()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return CircuitClosed, 0
|
||||||
|
}
|
||||||
|
|
||||||
// ScheduleRetry schedules a task to be re-enqueued onto the
|
// ScheduleRetry schedules a task to be re-enqueued onto the
|
||||||
// retry channel after delay. It implements the Scheduler
|
// retry channel after delay. It implements the Scheduler
|
||||||
// interface the targets use to own their durable retries.
|
// interface the targets use to own their durable retries.
|
||||||
|
|||||||
@@ -1018,6 +1018,62 @@ func TestGetCircuitBreaker_CreatesOnDemand(t *testing.T) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestStateAndCooldown_ReadsHTTPAndSlackBreakers proves the engine
|
||||||
|
// reads the state of an http or a slack target's circuit breaker, with
|
||||||
|
// what is left of its cooldown while it is open, and no cooldown while
|
||||||
|
// it is half-open, once it closes, or for a target with no breaker.
|
||||||
|
func TestStateAndCooldown_ReadsHTTPAndSlackBreakers(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
e := testEngine(t, 1)
|
||||||
|
|
||||||
|
httpID := uuid.New().String()
|
||||||
|
slackID := uuid.New().String()
|
||||||
|
|
||||||
|
state, cooldown := e.StateAndCooldown(httpID)
|
||||||
|
assert.Equal(t, delivery.CircuitClosed, state, "no breaker")
|
||||||
|
assert.Zero(t, cooldown, "no breaker")
|
||||||
|
|
||||||
|
httpCB := delivery.NewTestCircuitBreaker(1, time.Hour)
|
||||||
|
e.ExportSetCircuitBreaker(httpID, httpCB)
|
||||||
|
|
||||||
|
slackCB := delivery.NewTestCircuitBreaker(1, time.Hour)
|
||||||
|
e.ExportSetSlackCircuitBreaker(slackID, slackCB)
|
||||||
|
|
||||||
|
httpCB.RecordFailure()
|
||||||
|
slackCB.RecordFailure()
|
||||||
|
|
||||||
|
for _, id := range []string{httpID, slackID} {
|
||||||
|
state, cooldown := e.StateAndCooldown(id)
|
||||||
|
assert.Equal(t, delivery.CircuitOpen, state)
|
||||||
|
assert.Greater(t, cooldown, 59*time.Minute)
|
||||||
|
assert.LessOrEqual(t, cooldown, time.Hour)
|
||||||
|
}
|
||||||
|
|
||||||
|
httpCB.RecordSuccess()
|
||||||
|
slackCB.RecordSuccess()
|
||||||
|
|
||||||
|
for _, id := range []string{httpID, slackID} {
|
||||||
|
state, cooldown := e.StateAndCooldown(id)
|
||||||
|
assert.Equal(t, delivery.CircuitClosed, state, "closed")
|
||||||
|
assert.Zero(t, cooldown, "closed")
|
||||||
|
}
|
||||||
|
|
||||||
|
// A breaker with no cooldown goes half-open on the first Allow
|
||||||
|
// after it trips, letting that one delivery through to test the
|
||||||
|
// target.
|
||||||
|
halfOpenID := uuid.New().String()
|
||||||
|
halfOpenCB := delivery.NewTestCircuitBreaker(1, 0)
|
||||||
|
e.ExportSetCircuitBreaker(halfOpenID, halfOpenCB)
|
||||||
|
|
||||||
|
halfOpenCB.RecordFailure()
|
||||||
|
require.True(t, halfOpenCB.Allow())
|
||||||
|
|
||||||
|
state, cooldown = e.StateAndCooldown(halfOpenID)
|
||||||
|
assert.Equal(t, delivery.CircuitHalfOpen, state)
|
||||||
|
assert.Zero(t, cooldown, "half-open")
|
||||||
|
}
|
||||||
|
|
||||||
func TestParseHTTPConfig_Valid(t *testing.T) {
|
func TestParseHTTPConfig_Valid(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
|
|||||||
@@ -212,6 +212,14 @@ func (e *Engine) ExportSetCircuitBreaker(
|
|||||||
e.httpTarget.circuitBreakers.Store(targetID, cb)
|
e.httpTarget.circuitBreakers.Store(targetID, cb)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ExportSetSlackCircuitBreaker is ExportSetCircuitBreaker for the
|
||||||
|
// slack target.
|
||||||
|
func (e *Engine) ExportSetSlackCircuitBreaker(
|
||||||
|
targetID string, cb *CircuitBreaker,
|
||||||
|
) {
|
||||||
|
e.slackTarget.circuitBreakers.Store(targetID, cb)
|
||||||
|
}
|
||||||
|
|
||||||
// ExportParseHTTPConfig exposes parseHTTPConfig.
|
// ExportParseHTTPConfig exposes parseHTTPConfig.
|
||||||
func (e *Engine) ExportParseHTTPConfig(
|
func (e *Engine) ExportParseHTTPConfig(
|
||||||
configJSON string,
|
configJSON string,
|
||||||
|
|||||||
@@ -105,6 +105,7 @@ func (e *Engine) initTargets(client *http.Client) {
|
|||||||
dbT := &databaseTarget{eng: e}
|
dbT := &databaseTarget{eng: e}
|
||||||
|
|
||||||
e.httpTarget = httpT
|
e.httpTarget = httpT
|
||||||
|
e.slackTarget = slackT
|
||||||
e.dbTarget = dbT
|
e.dbTarget = dbT
|
||||||
|
|
||||||
e.targets = map[database.TargetType]Target{
|
e.targets = map[database.TargetType]Target{
|
||||||
|
|||||||
@@ -379,21 +379,54 @@ func (w *archiveWriter) close() {
|
|||||||
|
|
||||||
// sweepExpired prunes the target's archive files, which may have
|
// sweepExpired prunes the target's archive files, which may have
|
||||||
// gone idle, with no write to trigger the usual on-reopen prune. It
|
// gone idle, with no write to trigger the usual on-reopen prune. It
|
||||||
// takes the writer's own mutex for the whole operation, so a sweep
|
// lists the files under the writer's own mutex, then takes the mutex
|
||||||
// is ordered against concurrent writes rather than reaching around
|
// again for one file at a time, so a write waits for at most one
|
||||||
// them to the files.
|
// file's prune, and each prune is ordered against concurrent writes
|
||||||
|
// rather than reaching around them to the file.
|
||||||
//
|
//
|
||||||
// It never creates an archive file: it prunes only the files
|
// It never creates an archive file: it prunes only the files
|
||||||
// archiveFiles finds, and opens each with archiveModeExisting so
|
// archiveFiles lists, skips one that is gone by the time it is
|
||||||
// SQLite itself refuses to create one if the file disappears
|
// reached (moved away, or renamed since the listing), and opens each
|
||||||
// between the listing and the open. A file named for a period that
|
// with archiveModeExisting so SQLite itself refuses to create one if
|
||||||
// the prune leaves empty is deleted.
|
// the file disappears between the check and the open. A file named
|
||||||
|
// for a period that the prune leaves empty is deleted.
|
||||||
//
|
//
|
||||||
// The archive is left CLOSED afterwards. An idle archive holding
|
// The archive is left CLOSED afterwards. An idle archive holding
|
||||||
// no handle is what keeps the operator's move-the-file-away
|
// no handle is what keeps the operator's move-the-file-away
|
||||||
// workflow working; the next write reopens (and recreates) the
|
// workflow working; the next write reopens (and recreates) the
|
||||||
// file as it always has.
|
// file as it always has.
|
||||||
func (w *archiveWriter) sweepExpired(expiry time.Duration) error {
|
func (w *archiveWriter) sweepExpired(expiry time.Duration) error {
|
||||||
|
w.mu.Lock()
|
||||||
|
files, err := archiveFiles(w.path)
|
||||||
|
w.mu.Unlock()
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
var errs []error
|
||||||
|
|
||||||
|
for _, file := range files {
|
||||||
|
err = w.sweepFile(file, expiry)
|
||||||
|
if errors.Is(err, errArchiveWriterEvicted) {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
errs = append(errs, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return errors.Join(errs...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// sweepFile prunes one of the target's archive files for sweepExpired,
|
||||||
|
// holding w.mu while it does. It skips a file that is gone, and deletes
|
||||||
|
// the file, with its -wal and -shm, when it is named for a period and
|
||||||
|
// the prune leaves it empty.
|
||||||
|
func (w *archiveWriter) sweepFile(
|
||||||
|
file archiveFile, expiry time.Duration,
|
||||||
|
) error {
|
||||||
w.mu.Lock()
|
w.mu.Lock()
|
||||||
defer w.mu.Unlock()
|
defer w.mu.Unlock()
|
||||||
|
|
||||||
@@ -403,30 +436,10 @@ func (w *archiveWriter) sweepExpired(expiry time.Duration) error {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
files, err := archiveFiles(w.path)
|
if !fileExists(file.path) {
|
||||||
if err != nil {
|
return nil
|
||||||
return err
|
|
||||||
}
|
}
|
||||||
|
|
||||||
var errs []error
|
|
||||||
|
|
||||||
for _, file := range files {
|
|
||||||
err = w.sweepFile(file, expiry)
|
|
||||||
if err != nil {
|
|
||||||
errs = append(errs, err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return errors.Join(errs...)
|
|
||||||
}
|
|
||||||
|
|
||||||
// sweepFile prunes one of the target's archive files for
|
|
||||||
// sweepExpired, which holds w.mu, and deletes the file, with its -wal
|
|
||||||
// and -shm, when it is named for a period and the prune leaves it
|
|
||||||
// empty.
|
|
||||||
func (w *archiveWriter) sweepFile(
|
|
||||||
file archiveFile, expiry time.Duration,
|
|
||||||
) error {
|
|
||||||
// Drop any live handle first so the prune runs against a
|
// Drop any live handle first so the prune runs against a
|
||||||
// freshly opened file, matching the write path's semantics.
|
// freshly opened file, matching the write path's semantics.
|
||||||
w.close()
|
w.close()
|
||||||
|
|||||||
@@ -9,8 +9,11 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
|
"io/fs"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
"unicode/utf8"
|
"unicode/utf8"
|
||||||
|
|
||||||
@@ -51,18 +54,35 @@ func ArchiveExportFileName(
|
|||||||
at.UTC().Format("20060102T150405Z") + ".json.gz"
|
at.UTC().Format("20060102T150405Z") + ".json.gz"
|
||||||
}
|
}
|
||||||
|
|
||||||
// ArchiveExport is a database target's archive opened for download:
|
// ArchiveExport is a database target's archive listed for download. It
|
||||||
// every one of its files, each read on its own connection inside one
|
// opens one of the target's files at a time, only when its rows are
|
||||||
// read-only transaction, so it writes out the archive as it stood when
|
// about to be written out, and closes it before it opens the next, so
|
||||||
// OpenArchiveExport returned.
|
// an export holds at most one file open however many the target has.
|
||||||
//
|
//
|
||||||
// Archives are in WAL mode, where a reader works from a snapshot and
|
// Each file is read on its own connection inside one read-only
|
||||||
// never blocks a writer: archive writes go on while an export is open,
|
// transaction, so its rows are written out as the file stood when it
|
||||||
// and the export does not see them. SQLite cannot checkpoint a -wal
|
// was opened. Archives are in WAL mode, where a reader works from a
|
||||||
// past an open snapshot, so a file's -wal grows until the export has
|
// snapshot and never blocks a writer: archive writes go on while a file
|
||||||
// written out that file.
|
// is open, and the export does not see them. SQLite cannot checkpoint
|
||||||
|
// a -wal past an open snapshot, so the open file's -wal grows until the
|
||||||
|
// export has written that file out.
|
||||||
type ArchiveExport struct {
|
type ArchiveExport struct {
|
||||||
files []*exportFile
|
// periods are the periods of the target's files when the export
|
||||||
|
// was listed, "" for the file without one, in the order
|
||||||
|
// archiveFiles lists them.
|
||||||
|
periods []string
|
||||||
|
|
||||||
|
// lock is held while currentPath is called and a file is opened,
|
||||||
|
// so that a rename, which holds it too, cannot move the file in
|
||||||
|
// between.
|
||||||
|
lock sync.Locker
|
||||||
|
|
||||||
|
// currentPath returns the path ArchivePath gives the target under
|
||||||
|
// the names stored for it now, which a rename may have changed since
|
||||||
|
// the export was listed.
|
||||||
|
currentPath func() (string, error)
|
||||||
|
|
||||||
|
log *slog.Logger
|
||||||
}
|
}
|
||||||
|
|
||||||
// exportFile is one archive file opened for an export.
|
// exportFile is one archive file opened for an export.
|
||||||
@@ -83,43 +103,40 @@ type exportedName struct {
|
|||||||
Name string `json:"name"`
|
Name string `json:"name"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// OpenArchiveExport opens every one of a database target's archive
|
// NewArchiveExport lists a database target's archive files for export,
|
||||||
// files for export, given the path ArchivePath gives it (see
|
// given the path ArchivePath gives it (see archiveFiles). It opens none
|
||||||
// archiveFiles), and takes the snapshots the export reads. It never
|
// of them. Its caller holds lock, which every rename of the target's
|
||||||
// creates a file: with no files, the export has no rows.
|
// files runs under, from reading the names path is made of until it
|
||||||
|
// returns, so the files it lists are the ones those names give.
|
||||||
//
|
//
|
||||||
// Once it has returned, the files are open, so a rename or a move of
|
// WriteGzipJSON, called without lock held, finds each file again by its
|
||||||
// them does not affect the export, which reads the same files under
|
// period under the path currentPath gives, holding lock while it does
|
||||||
// their new names.
|
// and while it opens the file, so a rename during the export loses no
|
||||||
//
|
// file. A file that is gone by then, emptied by the sweep or moved
|
||||||
// The transactions last as long as ctx does, so ctx must last for the
|
// away, is skipped. The export never creates a file: with no files, it
|
||||||
// whole export.
|
// has no rows.
|
||||||
func OpenArchiveExport(
|
func NewArchiveExport(
|
||||||
ctx context.Context, path string, log *slog.Logger,
|
path string,
|
||||||
|
lock sync.Locker,
|
||||||
|
currentPath func() (string, error),
|
||||||
|
log *slog.Logger,
|
||||||
) (*ArchiveExport, error) {
|
) (*ArchiveExport, error) {
|
||||||
files, err := archiveFiles(path)
|
files, err := archiveFiles(path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
x := &ArchiveExport{}
|
x := &ArchiveExport{lock: lock, currentPath: currentPath, log: log}
|
||||||
|
|
||||||
for _, file := range files {
|
for _, file := range files {
|
||||||
f, err := openExportFile(ctx, file, log)
|
x.periods = append(x.periods, file.period)
|
||||||
if err != nil {
|
|
||||||
_ = x.Close()
|
|
||||||
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
x.files = append(x.files, f)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return x, nil
|
return x, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// openExportFile opens one archive file for an export and takes its
|
// openExportFile opens one archive file for an export and takes its
|
||||||
// snapshot.
|
// snapshot. The transaction lasts as long as ctx does.
|
||||||
func openExportFile(
|
func openExportFile(
|
||||||
ctx context.Context, file archiveFile, log *slog.Logger,
|
ctx context.Context, file archiveFile, log *slog.Logger,
|
||||||
) (*exportFile, error) {
|
) (*exportFile, error) {
|
||||||
@@ -179,9 +196,9 @@ func openExportFile(
|
|||||||
//
|
//
|
||||||
// Each row is written out before the next is read, so neither the
|
// Each row is written out before the next is read, so neither the
|
||||||
// archive nor its JSON is ever held in memory whole, and each file is
|
// archive nor its JSON is ever held in memory whole, and each file is
|
||||||
// closed once its rows are written. After an error the gzip stream is
|
// closed once its rows are written, before the next is opened. When it
|
||||||
// left unfinished, so what was written does not decompress as a whole
|
// returns, no file is open. After an error the gzip stream is left
|
||||||
// file.
|
// unfinished, so what was written does not decompress as a whole file.
|
||||||
func (x *ArchiveExport) WriteGzipJSON(
|
func (x *ArchiveExport) WriteGzipJSON(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
w io.Writer,
|
w io.Writer,
|
||||||
@@ -208,30 +225,35 @@ func (x *ArchiveExport) WriteGzipJSON(
|
|||||||
return zw.Close()
|
return zw.Close()
|
||||||
}
|
}
|
||||||
|
|
||||||
// Close ends the transactions of the files still open and closes
|
// openFile finds the target's archive file for period under the path
|
||||||
// their connections.
|
// currentPath gives now, and opens it for the export, holding x.lock
|
||||||
func (x *ArchiveExport) Close() error {
|
// for both. For a file that is gone, the error wraps fs.ErrNotExist.
|
||||||
errs := make([]error, 0, len(x.files))
|
func (x *ArchiveExport) openFile(
|
||||||
|
ctx context.Context, period string,
|
||||||
|
) (*exportFile, error) {
|
||||||
|
x.lock.Lock()
|
||||||
|
defer x.lock.Unlock()
|
||||||
|
|
||||||
for _, f := range x.files {
|
path, err := x.currentPath()
|
||||||
errs = append(errs, f.close())
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("finding archive file: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
return errors.Join(errs...)
|
file := archiveFile{path: archivePeriodPath(path, period), period: period}
|
||||||
|
|
||||||
|
_, err = os.Stat(file.path)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return openExportFile(ctx, file, x.log)
|
||||||
}
|
}
|
||||||
|
|
||||||
// close ends the file's transaction and closes its connection, once.
|
// close ends the file's transaction and closes its connection.
|
||||||
func (f *exportFile) close() error {
|
func (f *exportFile) close() error {
|
||||||
if f.db == nil {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
_ = f.tx.Rollback()
|
_ = f.tx.Rollback()
|
||||||
|
|
||||||
err := f.db.Close()
|
return f.db.Close()
|
||||||
f.db = nil
|
|
||||||
|
|
||||||
return err
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// writeJSON writes head with archived_events added as its last key,
|
// writeJSON writes head with archived_events added as its last key,
|
||||||
@@ -262,19 +284,24 @@ func (x *ArchiveExport) writeJSON(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// writeRows writes the archived rows of each file to w, one per line,
|
// writeRows writes the archived rows of each file to w, one per line,
|
||||||
// separated by commas, and closes each file once its rows are written.
|
// separated by commas, opening each file in turn and closing it once
|
||||||
|
// its rows are written.
|
||||||
func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error {
|
func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error {
|
||||||
sep := "\n"
|
sep := "\n"
|
||||||
|
|
||||||
for _, f := range x.files {
|
for _, period := range x.periods {
|
||||||
var err error
|
f, err := x.openFile(ctx, period)
|
||||||
|
if errors.Is(err, fs.ErrNotExist) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
sep, err = f.writeRows(ctx, w, sep)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
err = f.close()
|
sep, err = f.writeRows(ctx, w, sep)
|
||||||
|
|
||||||
|
err = errors.Join(err, f.close())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,6 +12,8 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"runtime"
|
"runtime"
|
||||||
|
"strings"
|
||||||
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -29,14 +31,8 @@ const (
|
|||||||
exportTargetName = "Long-term archive"
|
exportTargetName = "Long-term archive"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
// binaryBody is a body that is not valid UTF-8.
|
||||||
// binaryBody is a body that is not valid UTF-8.
|
const binaryBody = "\xff\xfe\x00\x01binary\x80"
|
||||||
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'
|
// writeExportTo writes export to w as the archive of the export tests'
|
||||||
// webhook and target, exported at 2026-10-02T12:03:04Z.
|
// webhook and target, exported at 2026-10-02T12:03:04Z.
|
||||||
@@ -59,19 +55,41 @@ func writeExportTo(
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// listExport lists the archive at path for export, as the archive of a
|
||||||
|
// target whose names do not change.
|
||||||
|
func listExport(t *testing.T, path string) *delivery.ArchiveExport {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
return newExport(t, path, &sync.Mutex{}, func() (string, error) {
|
||||||
|
return path, nil
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
// newExport lists the archive at path for export, to find each file
|
||||||
|
// again under the path currentPath gives, holding lock while it does.
|
||||||
|
// Nothing else takes lock while it lists, so it does not hold lock.
|
||||||
|
func newExport(
|
||||||
|
t *testing.T,
|
||||||
|
path string,
|
||||||
|
lock sync.Locker,
|
||||||
|
currentPath func() (string, error),
|
||||||
|
) *delivery.ArchiveExport {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
export, err := delivery.NewArchiveExport(
|
||||||
|
path, lock, currentPath, archiveTestLogger(),
|
||||||
|
)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
return export
|
||||||
|
}
|
||||||
|
|
||||||
// exportArchive runs a whole export of the archive at path and returns
|
// exportArchive runs a whole export of the archive at path and returns
|
||||||
// its JSON, decompressed and parsed.
|
// its JSON, decompressed and parsed.
|
||||||
func exportArchive(t *testing.T, path string) map[string]any {
|
func exportArchive(t *testing.T, path string) map[string]any {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
export, err := delivery.OpenArchiveExport(
|
return writeExport(t, listExport(t, path))
|
||||||
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,
|
// writeExport writes an opened export and returns its JSON,
|
||||||
@@ -241,63 +259,85 @@ func TestArchiveExport_Empty(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestArchiveExport_ReadsOneSnapshot proves an export writes the
|
// unlockHook is a sync.Locker that runs fn each time it is unlocked. An
|
||||||
// archive as it was when it was opened, and holds up no archive
|
// export unlocks its lock right after it opens a file.
|
||||||
// write: a row written while the export is open is stored, and is not
|
type unlockHook struct {
|
||||||
// in the export. A write held up for the whole busy timeout would
|
sync.Mutex
|
||||||
// fail.
|
|
||||||
|
fn func()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (u *unlockHook) Unlock() {
|
||||||
|
u.Mutex.Unlock()
|
||||||
|
u.fn()
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestArchiveExport_ReadsOneSnapshot proves an export writes a file
|
||||||
|
// out as it was when the export opened it, and holds up no archive
|
||||||
|
// write: a row written after the export was listed but before the file
|
||||||
|
// was opened is in the export, and one written while the file is open
|
||||||
|
// is stored, and is not. A write held up for the whole busy timeout
|
||||||
|
// would fail.
|
||||||
func TestArchiveExport_ReadsOneSnapshot(t *testing.T) {
|
func TestArchiveExport_ReadsOneSnapshot(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
path := filepath.Join(t.TempDir(), "archive.db")
|
path := filepath.Join(t.TempDir(), "archive.db")
|
||||||
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
|
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
|
||||||
|
|
||||||
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: openedEventID}, 0))
|
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "listed"}, 0))
|
||||||
|
|
||||||
export, err := delivery.OpenArchiveExport(
|
opened := &unlockHook{fn: func() {
|
||||||
t.Context(), path, archiveTestLogger(),
|
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "during"}, 0))
|
||||||
)
|
}}
|
||||||
require.NoError(t, err)
|
export := newExport(t, path, opened, func() (string, error) {
|
||||||
|
return path, nil
|
||||||
|
})
|
||||||
|
|
||||||
defer func() { require.NoError(t, export.Close()) }()
|
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "before-open"}, 0))
|
||||||
|
|
||||||
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "during"}, 0))
|
|
||||||
|
|
||||||
assert.Equal(t,
|
assert.Equal(t,
|
||||||
[]string{openedEventID}, exportedEventIDs(t, writeExport(t, export)),
|
[]string{"listed", "before-open"},
|
||||||
|
exportedEventIDs(t, writeExport(t, export)),
|
||||||
)
|
)
|
||||||
|
|
||||||
var stored int64
|
var stored int64
|
||||||
|
|
||||||
require.NoError(t, openArchiveDBForRead(t, path).
|
require.NoError(t, openArchiveDBForRead(t, path).
|
||||||
Model(&delivery.ExportArchivedEvent{}).Count(&stored).Error)
|
Model(&delivery.ExportArchivedEvent{}).Count(&stored).Error)
|
||||||
assert.Equal(t, int64(2), stored)
|
assert.Equal(t, int64(3), stored)
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestArchiveExport_SurvivesRename proves that renaming the archive
|
// TestArchiveExport_FindsFilesAfterRename proves that renaming the
|
||||||
// while an export of it is open, as renaming its webhook or target
|
// archive after an export has listed it, as renaming its webhook or
|
||||||
// does, leaves the export reading the same file.
|
// target does, loses no file: the export finds each file again by its
|
||||||
func TestArchiveExport_SurvivesRename(t *testing.T) {
|
// period under the new name. A file moved away by then is skipped.
|
||||||
|
func TestArchiveExport_FindsFilesAfterRename(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
path := filepath.Join(t.TempDir(), "archive-old.db")
|
dir := t.TempDir()
|
||||||
|
path := filepath.Join(dir, "archive-old.db")
|
||||||
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
|
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
|
||||||
|
|
||||||
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: openedEventID}, 0))
|
for _, period := range []string{"", dayPeriod, hourPeriod} {
|
||||||
|
require.NoError(t, w.WritePeriod(
|
||||||
|
delivery.ExportArchivedEvent{EventID: "in-" + period}, 0, period,
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
export, err := delivery.OpenArchiveExport(
|
current := path
|
||||||
t.Context(), path, archiveTestLogger(),
|
export := newExport(t, path, &sync.Mutex{}, func() (string, error) {
|
||||||
)
|
return current, nil
|
||||||
require.NoError(t, err)
|
})
|
||||||
|
|
||||||
defer func() { require.NoError(t, export.Close()) }()
|
|
||||||
|
|
||||||
require.NoError(t, w.Rename("archive-new.db"))
|
require.NoError(t, w.Rename("archive-new.db"))
|
||||||
require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "after"}, 0))
|
|
||||||
require.NoFileExists(t, path)
|
current = filepath.Join(dir, "archive-new.db")
|
||||||
|
|
||||||
|
removeArchiveFiles(t, periodPath(current, dayPeriod))
|
||||||
|
|
||||||
assert.Equal(t,
|
assert.Equal(t,
|
||||||
[]string{openedEventID}, exportedEventIDs(t, writeExport(t, export)),
|
[]string{"in-", "in-" + hourPeriod},
|
||||||
|
exportedEventIDs(t, writeExport(t, export)),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -306,7 +346,7 @@ func TestArchiveExport_SurvivesRename(t *testing.T) {
|
|||||||
// day, and proves the export holds every row: the file without a
|
// day, and proves the export holds every row: the file without a
|
||||||
// period first, then the others oldest period first, each row from a
|
// period first, then the others oldest period first, each row from a
|
||||||
// file named for a period carrying that period. A file made after the
|
// file named for a period carrying that period. A file made after the
|
||||||
// export opened is not in it.
|
// export was listed is not in it.
|
||||||
func TestArchiveExport_EveryFileOldestFirst(t *testing.T) {
|
func TestArchiveExport_EveryFileOldestFirst(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
@@ -320,12 +360,7 @@ func TestArchiveExport_EveryFileOldestFirst(t *testing.T) {
|
|||||||
))
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
export, err := delivery.OpenArchiveExport(
|
export := listExport(t, path)
|
||||||
t.Context(), path, archiveTestLogger(),
|
|
||||||
)
|
|
||||||
require.NoError(t, err)
|
|
||||||
|
|
||||||
defer func() { require.NoError(t, export.Close()) }()
|
|
||||||
|
|
||||||
require.NoError(t, w.WritePeriod(
|
require.NoError(t, w.WritePeriod(
|
||||||
delivery.ExportArchivedEvent{EventID: "later"}, 0, "2026-03-06",
|
delivery.ExportArchivedEvent{EventID: "later"}, 0, "2026-03-06",
|
||||||
@@ -351,6 +386,72 @@ func TestArchiveExport_EveryFileOldestFirst(t *testing.T) {
|
|||||||
"a row from the file without a period has no period")
|
"a row from the file without a period has no period")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// openFilesPeak is an io.Writer that discards what it is given and
|
||||||
|
// records the most archive files in dir the process had open at any
|
||||||
|
// write, as /proc/self/fd lists the files a process has open.
|
||||||
|
type openFilesPeak struct {
|
||||||
|
dir string
|
||||||
|
max int
|
||||||
|
}
|
||||||
|
|
||||||
|
func (p *openFilesPeak) Write(b []byte) (int, error) {
|
||||||
|
fds, err := os.ReadDir("/proc/self/fd")
|
||||||
|
if err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
|
||||||
|
open := map[string]bool{}
|
||||||
|
|
||||||
|
for _, fd := range fds {
|
||||||
|
file, err := os.Readlink(filepath.Join("/proc/self/fd", fd.Name()))
|
||||||
|
if err == nil && filepath.Dir(file) == p.dir &&
|
||||||
|
strings.HasSuffix(file, ".db") {
|
||||||
|
open[file] = true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
p.max = max(p.max, len(open))
|
||||||
|
|
||||||
|
return len(b), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestArchiveExport_OneFileOpenAtATime exports a target with a file for
|
||||||
|
// each of 24 hours and proves the export never had more than one of
|
||||||
|
// them open, and had one open while it wrote. Each file holds a row of
|
||||||
|
// 48 KiB of random base64, which gzip shrinks little, so the export
|
||||||
|
// writes output while it reads each file.
|
||||||
|
func TestArchiveExport_OneFileOpenAtATime(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
if runtime.GOOS != "linux" {
|
||||||
|
t.Skip("only Linux lists a process's open files in /proc/self/fd")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Readlink gives each open file's path with no symbolic link in it.
|
||||||
|
dir, err := filepath.EvalSymlinks(t.TempDir())
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
path := filepath.Join(dir, "archive-wh.db")
|
||||||
|
w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0)
|
||||||
|
random := make([]byte, 36<<10)
|
||||||
|
|
||||||
|
for hour := range 24 {
|
||||||
|
_, _ = rand.Read(random)
|
||||||
|
|
||||||
|
require.NoError(t, w.WritePeriod(delivery.ExportArchivedEvent{
|
||||||
|
Body: base64.StdEncoding.EncodeToString(random),
|
||||||
|
}, 0, fmt.Sprintf("2026-10-01-%02d", hour)))
|
||||||
|
}
|
||||||
|
|
||||||
|
// The writer's own handle on the last file is not the export's.
|
||||||
|
w.Evict()
|
||||||
|
|
||||||
|
peak := &openFilesPeak{dir: dir}
|
||||||
|
|
||||||
|
require.NoError(t, writeExportTo(t, listExport(t, path), peak))
|
||||||
|
assert.Equal(t, 1, peak.max)
|
||||||
|
}
|
||||||
|
|
||||||
// heapPeak is an io.Writer that discards what it is given and records
|
// 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
|
// the largest heap it saw at a write. It collects garbage before each
|
||||||
// reading, so the heap it reads is what is still held.
|
// reading, so the heap it reads is what is still held.
|
||||||
@@ -388,12 +489,7 @@ func exportHeapGrowth(t *testing.T, rows, bodySize int) uint64 {
|
|||||||
}, 0))
|
}, 0))
|
||||||
}
|
}
|
||||||
|
|
||||||
export, err := delivery.OpenArchiveExport(
|
export := listExport(t, path)
|
||||||
t.Context(), path, archiveTestLogger(),
|
|
||||||
)
|
|
||||||
require.NoError(t, err)
|
|
||||||
|
|
||||||
defer func() { require.NoError(t, export.Close()) }()
|
|
||||||
|
|
||||||
runtime.GC()
|
runtime.GC()
|
||||||
|
|
||||||
|
|||||||
@@ -721,3 +721,41 @@ func TestArchiveWriter_RenameMovesBackOnFailure(t *testing.T) {
|
|||||||
assert.NoFileExists(t, filepath.Join(dir, newName))
|
assert.NoFileExists(t, filepath.Join(dir, newName))
|
||||||
assert.Equal(t, oldPath, w.Path())
|
assert.Equal(t, oldPath, w.Path())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestArchiveWriter_RenameMovesBackEveryFile renames a target with a
|
||||||
|
// file without a period and a file for a month, and proves that when
|
||||||
|
// the month's file fails to move, the file already moved is moved back:
|
||||||
|
// both files are under the old name with their rows, and nothing is
|
||||||
|
// under the new name. The new name is 251 bytes, so the file without a
|
||||||
|
// period, with its -wal and -shm, can take it, but the month's file,
|
||||||
|
// eight bytes longer, cannot.
|
||||||
|
func TestArchiveWriter_RenameMovesBackEveryFile(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
dir := t.TempDir()
|
||||||
|
oldPath := filepath.Join(dir, "archive-old.db")
|
||||||
|
monthPath := filepath.Join(dir, "archive-old-2026-03.db")
|
||||||
|
w := delivery.NewExportArchiveWriter(oldPath, archiveTestLogger(), 0)
|
||||||
|
|
||||||
|
require.NoError(t, w.WritePeriod(
|
||||||
|
delivery.ExportArchivedEvent{EventID: "in-none"}, 0, "",
|
||||||
|
))
|
||||||
|
require.NoError(t, w.WritePeriod(
|
||||||
|
delivery.ExportArchivedEvent{EventID: "in-month"}, 0, "2026-03",
|
||||||
|
))
|
||||||
|
|
||||||
|
require.Error(t, w.Rename(strings.Repeat("a", 248)+".db"))
|
||||||
|
|
||||||
|
assert.Equal(t, []string{"in-none"}, archivedEventIDs(t, oldPath))
|
||||||
|
assert.Equal(t, []string{"in-month"}, archivedEventIDs(t, monthPath))
|
||||||
|
|
||||||
|
entries, err := os.ReadDir(dir)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
for _, entry := range entries {
|
||||||
|
assert.True(t, strings.HasPrefix(entry.Name(), "archive-old"),
|
||||||
|
"%s is not under the old name", entry.Name())
|
||||||
|
}
|
||||||
|
|
||||||
|
assert.Equal(t, oldPath, w.Path())
|
||||||
|
}
|
||||||
|
|||||||
@@ -234,7 +234,7 @@ func (c *httpCore) handleRetry(
|
|||||||
database.DeliveryStatusRetrying,
|
database.DeliveryStatusRetrying,
|
||||||
)
|
)
|
||||||
|
|
||||||
backoff := calcBackoff(attemptNum)
|
backoff := Backoff(attemptNum)
|
||||||
|
|
||||||
retryTask := *task
|
retryTask := *task
|
||||||
retryTask.AttemptNum = attemptNum + 1
|
retryTask.AttemptNum = attemptNum + 1
|
||||||
@@ -301,7 +301,7 @@ func (c *httpCore) remainingBackoff(
|
|||||||
return 0
|
return 0
|
||||||
}
|
}
|
||||||
|
|
||||||
backoff := calcBackoff(attemptNum)
|
backoff := Backoff(attemptNum)
|
||||||
elapsed := time.Since(lastResult.CreatedAt)
|
elapsed := time.Since(lastResult.CreatedAt)
|
||||||
remaining := backoff - elapsed
|
remaining := backoff - elapsed
|
||||||
|
|
||||||
@@ -326,12 +326,14 @@ func (c *httpCore) backoffElapsed(
|
|||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
backoff := calcBackoff(attemptNum)
|
backoff := Backoff(attemptNum)
|
||||||
|
|
||||||
return time.Since(lastResult.CreatedAt) >= backoff
|
return time.Since(lastResult.CreatedAt) >= backoff
|
||||||
}
|
}
|
||||||
|
|
||||||
func calcBackoff(attemptNum int) time.Duration {
|
// Backoff is how long an http or slack target with retries waits after
|
||||||
|
// a delivery's failed attempt attemptNum before trying it again.
|
||||||
|
func Backoff(attemptNum int) time.Duration {
|
||||||
shift := max(attemptNum-1, 0)
|
shift := max(attemptNum-1, 0)
|
||||||
shift = min(shift, maxBackoffShift)
|
shift = min(shift, maxBackoffShift)
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,89 @@
|
|||||||
|
package handlers_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"net/http"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
"gorm.io/gorm/clause"
|
||||||
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
|
)
|
||||||
|
|
||||||
|
// TestDeliveryAttempts_ReadInTheTargetTypesOwnTerms proves, on the
|
||||||
|
// event's page and in the event log, that an http or slack attempt
|
||||||
|
// shows its status as before, while a database or log attempt, which
|
||||||
|
// sends no HTTP request, says what it did and shows no status.
|
||||||
|
func TestDeliveryAttempts_ReadInTheTargetTypesOwnTerms(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
cases := []struct {
|
||||||
|
targetType database.TargetType
|
||||||
|
success bool
|
||||||
|
statusCode int
|
||||||
|
errText string
|
||||||
|
outcome string
|
||||||
|
status string // "" when the attempt must show no status
|
||||||
|
}{
|
||||||
|
{
|
||||||
|
database.TargetTypeHTTP, false, 0, "",
|
||||||
|
"failure", "Status: — (no response)",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
database.TargetTypeSlack, true, http.StatusOK, "",
|
||||||
|
"success", "Status: 200",
|
||||||
|
},
|
||||||
|
{database.TargetTypeDatabase, true, 0, "", "archived", ""},
|
||||||
|
{
|
||||||
|
database.TargetTypeDatabase, false, 0,
|
||||||
|
"opening archive database: disk full", "failure", "",
|
||||||
|
},
|
||||||
|
{database.TargetTypeLog, true, 0, "", "written to the log", ""},
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, tc := range cases {
|
||||||
|
t.Run(string(tc.targetType), func(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
f := newRecentEventsFixture(t)
|
||||||
|
target := seedTarget(t, f.db, f.webhook.ID, tc.targetType)
|
||||||
|
event := f.event(t, contentTypeJSON, "{}", time.Now())
|
||||||
|
dlv := f.delivery(
|
||||||
|
t, event, target.ID, database.DeliveryStatusDelivered,
|
||||||
|
)
|
||||||
|
|
||||||
|
require.NoError(t, f.webhookDB.Omit(clause.Associations).Create(
|
||||||
|
&database.DeliveryResult{
|
||||||
|
DeliveryID: dlv.ID,
|
||||||
|
AttemptNum: 1,
|
||||||
|
Success: tc.success,
|
||||||
|
StatusCode: tc.statusCode,
|
||||||
|
Error: tc.errText,
|
||||||
|
},
|
||||||
|
).Error)
|
||||||
|
|
||||||
|
w := serveEventPage(t, f.h, f.sess, f.webhook.ID, event.ID)
|
||||||
|
require.Equal(t, http.StatusOK, w.Code)
|
||||||
|
|
||||||
|
pages := []string{
|
||||||
|
w.Body.String(),
|
||||||
|
renderSourceLogsPage(t, f.h, f.sess, f.webhook.ID),
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, page := range pages {
|
||||||
|
assert.Contains(t, page, ">"+tc.outcome+"</span>")
|
||||||
|
|
||||||
|
if tc.errText != "" {
|
||||||
|
assert.Contains(t, page, "Error: "+tc.errText)
|
||||||
|
}
|
||||||
|
|
||||||
|
if tc.status == "" {
|
||||||
|
assert.NotContains(t, page, "Status:")
|
||||||
|
} else {
|
||||||
|
assert.Contains(t, page, tc.status)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -21,7 +21,9 @@ const (
|
|||||||
// replayTargetDeleted reports a target that once existed and has
|
// replayTargetDeleted reports a target that once existed and has
|
||||||
// since been deleted. Deletes are soft and deliveries carry no
|
// since been deleted. Deletes are soft and deliveries carry no
|
||||||
// foreign key to the target row, so the history survives its
|
// foreign key to the target row, so the history survives its
|
||||||
// target and this is the ordinary case for an old event.
|
// target and this is the ordinary case for an old event. The
|
||||||
|
// event log shows no Replay button for such a delivery, so only
|
||||||
|
// a page loaded before the delete reaches this.
|
||||||
replayTargetDeleted noticeCode = "replay-target-deleted"
|
replayTargetDeleted noticeCode = "replay-target-deleted"
|
||||||
|
|
||||||
// replayTargetMissing reports a target id that names no row at
|
// replayTargetMissing reports a target id that names no row at
|
||||||
|
|||||||
@@ -513,7 +513,11 @@ func TestHandleSourceLogs_RendersReplayControlAndBanner(t *testing.T) {
|
|||||||
)
|
)
|
||||||
|
|
||||||
assert.Contains(t, refused, "alert-error")
|
assert.Contains(t, refused, "alert-error")
|
||||||
assert.Contains(t, refused, "has been deleted")
|
assert.Contains(
|
||||||
|
t, refused,
|
||||||
|
"has been deleted. Use Resubmit to send the event "+
|
||||||
|
"to the webhook",
|
||||||
|
)
|
||||||
|
|
||||||
// An outcome code nobody issued renders no banner at all.
|
// An outcome code nobody issued renders no banner at all.
|
||||||
unknown := renderSourceLogsPageWithQuery(
|
unknown := renderSourceLogsPageWithQuery(
|
||||||
|
|||||||
@@ -1,6 +1,8 @@
|
|||||||
package handlers
|
package handlers
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"time"
|
||||||
|
|
||||||
"sneak.berlin/go/webhooker/internal/delivery"
|
"sneak.berlin/go/webhooker/internal/delivery"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -24,8 +26,8 @@ const maxRenderedResponseBytes = 4096
|
|||||||
// bytes rather than characters, and they make SQLite do the
|
// bytes rather than characters, and they make SQLite do the
|
||||||
// cut, so an oversized stored response never becomes a Go
|
// cut, so an oversized stored response never becomes a Go
|
||||||
// string at all.
|
// string at all.
|
||||||
const deliveryResultColumns = "delivery_id, attempt_num, success, " +
|
const deliveryResultColumns = "delivery_id, attempt_num, created_at, " +
|
||||||
"status_code, error, duration, " +
|
"success, status_code, error, duration, " +
|
||||||
"substr(cast(response_body as blob), 1, ?) AS response_body, " +
|
"substr(cast(response_body as blob), 1, ?) AS response_body, " +
|
||||||
"length(cast(response_body as blob)) AS response_bytes"
|
"length(cast(response_body as blob)) AS response_bytes"
|
||||||
|
|
||||||
@@ -100,6 +102,7 @@ func (v DeliveryResultView) HasStatusCode() bool {
|
|||||||
type deliveryResultRow struct {
|
type deliveryResultRow struct {
|
||||||
DeliveryID string
|
DeliveryID string
|
||||||
AttemptNum int
|
AttemptNum int
|
||||||
|
CreatedAt time.Time
|
||||||
Success bool
|
Success bool
|
||||||
StatusCode int
|
StatusCode int
|
||||||
Error string
|
Error string
|
||||||
|
|||||||
@@ -57,19 +57,20 @@ var errVerificationBusy = errors.New(
|
|||||||
type HandlersParams struct {
|
type HandlersParams struct {
|
||||||
fx.In
|
fx.In
|
||||||
|
|
||||||
Logger *logger.Logger
|
Logger *logger.Logger
|
||||||
Globals *globals.Globals
|
Globals *globals.Globals
|
||||||
Config *config.Config
|
Config *config.Config
|
||||||
Database *database.Database
|
Database *database.Database
|
||||||
WebhookDBMgr *database.WebhookDBManager
|
WebhookDBMgr *database.WebhookDBManager
|
||||||
Healthcheck *healthcheck.Healthcheck
|
Healthcheck *healthcheck.Healthcheck
|
||||||
Session *session.Session
|
Session *session.Session
|
||||||
Middleware *middleware.Middleware
|
Middleware *middleware.Middleware
|
||||||
Notifier delivery.Notifier
|
Notifier delivery.Notifier
|
||||||
Archives delivery.Archives
|
Archives delivery.Archives
|
||||||
SSRFGuard *delivery.Guard
|
CircuitBreakers delivery.CircuitBreakers
|
||||||
Metrics *metrics.Set
|
SSRFGuard *delivery.Guard
|
||||||
Registry *prometheus.Registry
|
Metrics *metrics.Set
|
||||||
|
Registry *prometheus.Registry
|
||||||
}
|
}
|
||||||
|
|
||||||
// Handlers provides HTTP handler methods for all application
|
// Handlers provides HTTP handler methods for all application
|
||||||
@@ -84,6 +85,7 @@ type Handlers struct {
|
|||||||
mw *middleware.Middleware
|
mw *middleware.Middleware
|
||||||
notifier delivery.Notifier
|
notifier delivery.Notifier
|
||||||
archives delivery.Archives
|
archives delivery.Archives
|
||||||
|
breakers delivery.CircuitBreakers
|
||||||
mtr *metrics.Set
|
mtr *metrics.Set
|
||||||
templates map[string]*template.Template
|
templates map[string]*template.Template
|
||||||
|
|
||||||
@@ -98,7 +100,9 @@ type Handlers struct {
|
|||||||
// 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. An archive download holds it while
|
||||||
// it reads the stored names and opens the file they give.
|
// it reads the stored names and lists the files they give, and
|
||||||
|
// again for each file while it finds the file under the names
|
||||||
|
// stored then and opens it.
|
||||||
renameMu sync.Mutex
|
renameMu sync.Mutex
|
||||||
|
|
||||||
// dummyVerifications counts the equivalent-cost verifications
|
// dummyVerifications counts the equivalent-cost verifications
|
||||||
@@ -148,6 +152,7 @@ func New(
|
|||||||
s.mw = params.Middleware
|
s.mw = params.Middleware
|
||||||
s.notifier = params.Notifier
|
s.notifier = params.Notifier
|
||||||
s.archives = params.Archives
|
s.archives = params.Archives
|
||||||
|
s.breakers = params.CircuitBreakers
|
||||||
s.mtr = params.Metrics
|
s.mtr = params.Metrics
|
||||||
s.ssrf = params.SSRFGuard
|
s.ssrf = params.SSRFGuard
|
||||||
|
|
||||||
|
|||||||
@@ -9,6 +9,7 @@ import (
|
|||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
@@ -181,6 +182,40 @@ func (r *recordingArchives) Renames() []archiveRename {
|
|||||||
return out
|
return out
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// testCircuitBreakers is a delivery.CircuitBreakers that reports, for
|
||||||
|
// each target, the circuit state and cooldown a test gave it with Set,
|
||||||
|
// and a closed breaker for any other target.
|
||||||
|
type testCircuitBreakers struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
states map[string]delivery.CircuitState
|
||||||
|
cooldowns map[string]time.Duration
|
||||||
|
}
|
||||||
|
|
||||||
|
// Set makes the target's breaker read as state, with cooldown left.
|
||||||
|
func (b *testCircuitBreakers) Set(
|
||||||
|
targetID string, state delivery.CircuitState, cooldown time.Duration,
|
||||||
|
) {
|
||||||
|
b.mu.Lock()
|
||||||
|
defer b.mu.Unlock()
|
||||||
|
|
||||||
|
if b.states == nil {
|
||||||
|
b.states = map[string]delivery.CircuitState{}
|
||||||
|
b.cooldowns = map[string]time.Duration{}
|
||||||
|
}
|
||||||
|
|
||||||
|
b.states[targetID] = state
|
||||||
|
b.cooldowns[targetID] = cooldown
|
||||||
|
}
|
||||||
|
|
||||||
|
func (b *testCircuitBreakers) StateAndCooldown(
|
||||||
|
targetID string,
|
||||||
|
) (delivery.CircuitState, time.Duration) {
|
||||||
|
b.mu.Lock()
|
||||||
|
defer b.mu.Unlock()
|
||||||
|
|
||||||
|
return b.states[targetID], b.cooldowns[targetID]
|
||||||
|
}
|
||||||
|
|
||||||
// newTestApp returns an app whose RequireStart fails the test when
|
// newTestApp returns an app whose RequireStart fails the test when
|
||||||
// starting takes longer than fx's default start timeout of 15s. That
|
// starting takes longer than fx's default start timeout of 15s. That
|
||||||
// limit catches a start that hangs, not a busy host: measured with make
|
// limit catches a start that hangs, not a busy host: measured with make
|
||||||
@@ -231,6 +266,12 @@ func newTestAppWithConfig(
|
|||||||
func(r *recordingArchives) delivery.Archives {
|
func(r *recordingArchives) delivery.Archives {
|
||||||
return r
|
return r
|
||||||
},
|
},
|
||||||
|
func() *testCircuitBreakers {
|
||||||
|
return &testCircuitBreakers{}
|
||||||
|
},
|
||||||
|
func(b *testCircuitBreakers) delivery.CircuitBreakers {
|
||||||
|
return b
|
||||||
|
},
|
||||||
metrics.NewRegistry,
|
metrics.NewRegistry,
|
||||||
metrics.New,
|
metrics.New,
|
||||||
middleware.New,
|
middleware.New,
|
||||||
|
|||||||
@@ -66,7 +66,8 @@ func noticeFor(r *http.Request) *notice {
|
|||||||
},
|
},
|
||||||
replayTargetDeleted: {
|
replayTargetDeleted: {
|
||||||
Text: "Not replayed: the target this delivery was for " +
|
Text: "Not replayed: the target this delivery was for " +
|
||||||
"has been deleted. Recreate the target, then replay.",
|
"has been deleted. Use Resubmit to send the event " +
|
||||||
|
"to the webhook's currently active targets.",
|
||||||
Failed: true,
|
Failed: true,
|
||||||
},
|
},
|
||||||
replayTargetMissing: {
|
replayTargetMissing: {
|
||||||
|
|||||||
@@ -94,6 +94,44 @@ func TestHandleSourceLogs_NamesDeletedTarget(t *testing.T) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestHandleSourceLogs_OffersNoReplayForDeletedTarget proves a
|
||||||
|
// finished delivery offers Replay while its target lives and not
|
||||||
|
// once the target is deleted. A replay to a deleted target is always
|
||||||
|
// refused, and recreating the target makes a new one that the old
|
||||||
|
// delivery does not name.
|
||||||
|
func TestHandleSourceLogs_OffersNoReplayForDeletedTarget(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
var (
|
||||||
|
h *handlers.Handlers
|
||||||
|
sess *session.Session
|
||||||
|
db *database.Database
|
||||||
|
dbMgr *database.WebhookDBManager
|
||||||
|
)
|
||||||
|
|
||||||
|
app := newTestApp(t, &h, &sess, &db, &dbMgr)
|
||||||
|
app.RequireStart()
|
||||||
|
|
||||||
|
t.Cleanup(app.RequireStop)
|
||||||
|
|
||||||
|
wh := seedWebhook(t, db)
|
||||||
|
tgt := seedTarget(t, db, wh.ID, database.TargetTypeLog)
|
||||||
|
|
||||||
|
_, failed := seedFailedDelivery(t, dbMgr, wh.ID, tgt.ID)
|
||||||
|
replayForm := `action="/hook/` + wh.ID + `/deliveries/` +
|
||||||
|
failed.ID + `/replay"`
|
||||||
|
|
||||||
|
before := renderSourceLogsPage(t, h, sess, wh.ID)
|
||||||
|
assert.Contains(t, before, replayForm)
|
||||||
|
|
||||||
|
deleteTargetThroughHandler(t, h, sess, wh.ID, tgt.ID)
|
||||||
|
|
||||||
|
after := renderSourceLogsPage(t, h, sess, wh.ID)
|
||||||
|
assert.NotContains(t, after, replayForm)
|
||||||
|
assert.NotContains(t, after, ">Replay<")
|
||||||
|
assert.Contains(t, after, tgt.Name+deletedMarker)
|
||||||
|
}
|
||||||
|
|
||||||
// TestHandleSourceLogs_MasksDeletedTargetConfig proves that
|
// TestHandleSourceLogs_MasksDeletedTargetConfig proves that
|
||||||
// naming a deleted target does not widen what the page shows of
|
// naming a deleted target does not widen what the page shows of
|
||||||
// it: its stored configuration stays masked by exactly the rules
|
// it: its stored configuration stays masked by exactly the rules
|
||||||
|
|||||||
@@ -109,6 +109,10 @@ type DeliveryView struct {
|
|||||||
// the middle of Results. The page must show it, or the
|
// the middle of Results. The page must show it, or the
|
||||||
// bound would hide history rather than fold it.
|
// bound would hide history rather than fold it.
|
||||||
AttemptsOmitted int
|
AttemptsOmitted int
|
||||||
|
|
||||||
|
// Paused is set while the delivery is retrying and its
|
||||||
|
// target's circuit breaker is open, and nil otherwise.
|
||||||
|
Paused *PausedView
|
||||||
}
|
}
|
||||||
|
|
||||||
// eventLogTarget is what the event log needs to know about
|
// eventLogTarget is what the event log needs to know about
|
||||||
@@ -1297,7 +1301,7 @@ func (h *Handlers) eventLogViews(
|
|||||||
}
|
}
|
||||||
|
|
||||||
for i := range rows {
|
for i := range rows {
|
||||||
result[i].Deliveries = newDeliveryViews(
|
result[i].Deliveries = h.newDeliveryViews(
|
||||||
eventDeliveries[i], targetMap, attempts,
|
eventDeliveries[i], targetMap, attempts,
|
||||||
)
|
)
|
||||||
result[i].ResubmitCount = resubmits[rows[i].ID]
|
result[i].ResubmitCount = resubmits[rows[i].ID]
|
||||||
@@ -1421,8 +1425,9 @@ func (h *Handlers) loadDeliveryResults(
|
|||||||
|
|
||||||
// newDeliveryViews projects deliveries for rendering,
|
// newDeliveryViews projects deliveries for rendering,
|
||||||
// resolving each one's target to its display-safe view and
|
// resolving each one's target to its display-safe view and
|
||||||
// each one's attempts through that target's redactor.
|
// each one's attempts through that target's redactor. A
|
||||||
func newDeliveryViews(
|
// retrying delivery also reads its target's circuit breaker.
|
||||||
|
func (h *Handlers) newDeliveryViews(
|
||||||
deliveries []database.Delivery,
|
deliveries []database.Delivery,
|
||||||
targetMap map[string]eventLogTarget,
|
targetMap map[string]eventLogTarget,
|
||||||
attempts map[string][]deliveryResultRow,
|
attempts map[string][]deliveryResultRow,
|
||||||
@@ -1445,6 +1450,12 @@ func newDeliveryViews(
|
|||||||
AttemptCount: len(rows),
|
AttemptCount: len(rows),
|
||||||
AttemptsOmitted: omitted,
|
AttemptsOmitted: omitted,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if deliveries[i].Status == database.DeliveryStatusRetrying {
|
||||||
|
views[i].Paused = h.deliveryPausedView(
|
||||||
|
deliveries[i].TargetID, rows,
|
||||||
|
)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return views
|
return views
|
||||||
|
|||||||
@@ -27,13 +27,11 @@ func (h *Handlers) HandleTargetDownload() http.HandlerFunc {
|
|||||||
return func(w http.ResponseWriter, r *http.Request) {
|
return func(w http.ResponseWriter, r *http.Request) {
|
||||||
ctx := context.WithoutCancel(r.Context())
|
ctx := context.WithoutCancel(r.Context())
|
||||||
|
|
||||||
webhook, target, export, ok := h.openTargetArchive(ctx, w, r)
|
webhook, target, export, ok := h.listTargetArchive(w, r)
|
||||||
if !ok {
|
if !ok {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
defer func() { _ = export.Close() }()
|
|
||||||
|
|
||||||
now := time.Now()
|
now := time.Now()
|
||||||
|
|
||||||
w.Header().Set("Content-Type", "application/gzip")
|
w.Header().Set("Content-Type", "application/gzip")
|
||||||
@@ -81,17 +79,15 @@ func (d downloadWriter) Write(b []byte) (int, error) {
|
|||||||
return d.w.Write(b)
|
return d.w.Write(b)
|
||||||
}
|
}
|
||||||
|
|
||||||
// openTargetArchive opens the archive of the request's database target
|
// listTargetArchive lists the archive files of the request's database
|
||||||
// for export, with its reads under ctx. It reports false once it has
|
// target for export. It reports false once it has written the response.
|
||||||
// written the response.
|
|
||||||
//
|
//
|
||||||
// It holds renameMu, which every archive rename runs under, while it
|
// It holds renameMu, which every archive rename runs under, while it
|
||||||
// reads the stored names and opens the files, so the files it opens
|
// reads the stored names and lists the files, so the files it lists
|
||||||
// are the ones those names give. It lets go before the export is
|
// are the ones those names give. It lets go before the export is
|
||||||
// streamed: once the files are open, a rename does not affect the
|
// streamed, which takes renameMu again for each file only while it
|
||||||
// export.
|
// finds the file under the names stored then and opens it.
|
||||||
func (h *Handlers) openTargetArchive(
|
func (h *Handlers) listTargetArchive(
|
||||||
ctx context.Context,
|
|
||||||
w http.ResponseWriter,
|
w http.ResponseWriter,
|
||||||
r *http.Request,
|
r *http.Request,
|
||||||
) (database.Webhook, *database.Target, *delivery.ArchiveExport, bool) {
|
) (database.Webhook, *database.Target, *delivery.ArchiveExport, bool) {
|
||||||
@@ -109,14 +105,43 @@ func (h *Handlers) openTargetArchive(
|
|||||||
return database.Webhook{}, nil, nil, false
|
return database.Webhook{}, nil, nil, false
|
||||||
}
|
}
|
||||||
|
|
||||||
export, err := delivery.OpenArchiveExport(
|
export, err := delivery.NewArchiveExport(
|
||||||
ctx, delivery.ArchivePath(h.dbMgr, &webhook, target), h.log,
|
delivery.ArchivePath(h.dbMgr, &webhook, target),
|
||||||
|
&h.renameMu,
|
||||||
|
func() (string, error) {
|
||||||
|
return h.storedArchivePath(webhook.ID, target.ID)
|
||||||
|
},
|
||||||
|
h.log,
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
h.serverError(w, r, "failed to open archive for export", err)
|
h.serverError(w, r, "failed to list archive for export", err)
|
||||||
|
|
||||||
return database.Webhook{}, nil, nil, false
|
return database.Webhook{}, nil, nil, false
|
||||||
}
|
}
|
||||||
|
|
||||||
return webhook, target, export, true
|
return webhook, target, export, true
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// storedArchivePath returns the path delivery.ArchivePath gives a
|
||||||
|
// database target under the names stored for it and its webhook now.
|
||||||
|
// Its caller holds renameMu. A webhook or target deleted since is still
|
||||||
|
// found, since deleting one leaves its archive files under their names.
|
||||||
|
func (h *Handlers) storedArchivePath(
|
||||||
|
webhookID, targetID string,
|
||||||
|
) (string, error) {
|
||||||
|
var webhook database.Webhook
|
||||||
|
|
||||||
|
err := h.db.DB().Unscoped().First(&webhook, "id = ?", webhookID).Error
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
|
||||||
|
var target database.Target
|
||||||
|
|
||||||
|
err = h.db.DB().Unscoped().First(&target, "id = ?", targetID).Error
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
|
||||||
|
return delivery.ArchivePath(h.dbMgr, &webhook, &target), nil
|
||||||
|
}
|
||||||
|
|||||||
@@ -13,6 +13,8 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"net/url"
|
"net/url"
|
||||||
|
"os"
|
||||||
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
@@ -85,7 +87,7 @@ func TestHandleTargetDownload(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// TestHandleTargetDownload_WaitsForRename proves a download reads the
|
// TestHandleTargetDownload_WaitsForRename proves a download reads the
|
||||||
// target's names and opens its archive under the lock a rename holds:
|
// target's names and lists its archive under the lock a rename holds:
|
||||||
// started while an edit is renaming the archive, it waits, and is
|
// started while an edit is renaming the archive, it waits, and is
|
||||||
// named for the target's new name.
|
// named for the target's new name.
|
||||||
func TestHandleTargetDownload_WaitsForRename(t *testing.T) {
|
func TestHandleTargetDownload_WaitsForRename(t *testing.T) {
|
||||||
@@ -148,18 +150,17 @@ func (s *stalledWriter) Write(b []byte) (int, error) {
|
|||||||
return s.ResponseRecorder.Write(b)
|
return s.ResponseRecorder.Write(b)
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestHandleTargetDownload_StreamsWithoutTheLock proves a download
|
// startStalledDownload starts a download of the target and returns once
|
||||||
// lets go of the rename lock once its archive is open: while the
|
// it is stalled at its first write, which comes before it opens any
|
||||||
// download is stalled writing, an edit can still rename the target.
|
// archive file. Closing the writer's resume lets it go on; the returned
|
||||||
func TestHandleTargetDownload_StreamsWithoutTheLock(t *testing.T) {
|
// channel is closed when it has finished.
|
||||||
t.Parallel()
|
func startStalledDownload(
|
||||||
|
t *testing.T, env *sourceTestEnv, webhookID, targetID string,
|
||||||
env := setupSourceTest(t)
|
) (*stalledWriter, <-chan struct{}) {
|
||||||
wh := seedWebhookWithRetention(t, env.db, 7)
|
t.Helper()
|
||||||
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
|
|
||||||
|
|
||||||
req := httptest.NewRequestWithContext(
|
req := httptest.NewRequestWithContext(
|
||||||
t.Context(), http.MethodGet, downloadPath(wh.ID, archive.ID), nil,
|
t.Context(), http.MethodGet, downloadPath(webhookID, targetID), nil,
|
||||||
)
|
)
|
||||||
for _, c := range env.cookies {
|
for _, c := range env.cookies {
|
||||||
req.AddCookie(c)
|
req.AddCookie(c)
|
||||||
@@ -179,6 +180,21 @@ func TestHandleTargetDownload_StreamsWithoutTheLock(t *testing.T) {
|
|||||||
|
|
||||||
<-sw.writing
|
<-sw.writing
|
||||||
|
|
||||||
|
return sw, downloaded
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestHandleTargetDownload_StreamsWithoutTheLock proves a download
|
||||||
|
// lets go of the rename lock once it has listed its archive: 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)
|
||||||
|
|
||||||
|
sw, downloaded := startStalledDownload(t, env, wh.ID, archive.ID)
|
||||||
|
|
||||||
edited := make(chan *httptest.ResponseRecorder, 1)
|
edited := make(chan *httptest.ResponseRecorder, 1)
|
||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
@@ -197,6 +213,62 @@ func TestHandleTargetDownload_StreamsWithoutTheLock(t *testing.T) {
|
|||||||
assert.Equal(t, http.StatusOK, sw.Code)
|
assert.Equal(t, http.StatusOK, sw.Code)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestHandleTargetDownload_FindsFilesAfterRename proves a download
|
||||||
|
// finds each of the target's files again under the names stored when
|
||||||
|
// it reaches the file: the target is renamed while the download is
|
||||||
|
// stalled before it has opened any file, and the rows of both its files
|
||||||
|
// are in the download.
|
||||||
|
func TestHandleTargetDownload_FindsFilesAfterRename(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
env := setupSourceTest(t)
|
||||||
|
wh := seedWebhookWithRetention(t, env.db, 7)
|
||||||
|
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
|
||||||
|
oldPath := delivery.ArchivePath(env.dbMgr, &wh, archive)
|
||||||
|
|
||||||
|
month := func(path string) string {
|
||||||
|
return strings.TrimSuffix(path, ".db") + "-2026-10.db"
|
||||||
|
}
|
||||||
|
|
||||||
|
seedArchive(t, oldPath, 1, 16)
|
||||||
|
seedArchive(t, month(oldPath), 1, 16)
|
||||||
|
|
||||||
|
sw, downloaded := startStalledDownload(t, env, wh.ID, archive.ID)
|
||||||
|
|
||||||
|
require.Equal(t,
|
||||||
|
http.StatusSeeOther, renameTarget(env, wh.ID, archive.ID).Code,
|
||||||
|
)
|
||||||
|
|
||||||
|
// The test's archives record a rename without moving any file, so
|
||||||
|
// the files are moved here, as the delivery engine moves them.
|
||||||
|
var renamed database.Target
|
||||||
|
|
||||||
|
require.NoError(t, env.db.DB().First(&renamed, "id = ?", archive.ID).Error)
|
||||||
|
|
||||||
|
newPath := delivery.ArchivePath(env.dbMgr, &wh, &renamed)
|
||||||
|
|
||||||
|
require.NoError(t, os.Rename(oldPath, newPath))
|
||||||
|
require.NoError(t, os.Rename(month(oldPath), month(newPath)))
|
||||||
|
|
||||||
|
close(sw.resume)
|
||||||
|
<-downloaded
|
||||||
|
require.Equal(t, http.StatusOK, sw.Code)
|
||||||
|
|
||||||
|
zr, err := gzip.NewReader(sw.Body)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
var (
|
||||||
|
got map[string]json.RawMessage
|
||||||
|
events []map[string]any
|
||||||
|
)
|
||||||
|
|
||||||
|
require.NoError(t, json.NewDecoder(zr).Decode(&got))
|
||||||
|
require.NoError(t, json.Unmarshal(got["archived_events"], &events))
|
||||||
|
require.Len(t, events, 2)
|
||||||
|
assert.NotContains(t, events[0], "period")
|
||||||
|
assert.Equal(t, "2026-10", events[1]["period"])
|
||||||
|
}
|
||||||
|
|
||||||
// seedArchive writes rows to the archive file at path, each with a
|
// seedArchive writes rows to the archive file at path, each with a
|
||||||
// body of bodySize random bytes, which do not compress. Its table has
|
// body of bodySize random bytes, which do not compress. Its table has
|
||||||
// only the columns the test fills; an export writes the others empty.
|
// only the columns the test fills; an export writes the others empty.
|
||||||
|
|||||||
@@ -25,6 +25,84 @@ type TargetRowView struct {
|
|||||||
// Archive is a database target's archive files, and nil for a
|
// Archive is a database target's archive files, and nil for a
|
||||||
// target of any other type.
|
// target of any other type.
|
||||||
Archive *ArchiveFileView
|
Archive *ArchiveFileView
|
||||||
|
|
||||||
|
// Paused is set while the target's circuit breaker is turning its
|
||||||
|
// deliveries away, and nil otherwise.
|
||||||
|
Paused *PausedView
|
||||||
|
}
|
||||||
|
|
||||||
|
// PausedView is a target's circuit breaker turning deliveries away.
|
||||||
|
// While the breaker is open, Until is a time in UTC, and Relative how
|
||||||
|
// long that is from now: on the target's row, when the cooldown ends;
|
||||||
|
// on a delivery, the earliest it can be tried next. While it is
|
||||||
|
// half-open both are empty: the cooldown has ended, and the target's
|
||||||
|
// deliveries are held while one delivery tests whether the target has
|
||||||
|
// recovered.
|
||||||
|
type PausedView struct {
|
||||||
|
Until string
|
||||||
|
Relative string
|
||||||
|
}
|
||||||
|
|
||||||
|
// pausedView reads the target's circuit breaker for its row, and
|
||||||
|
// returns nil when the breaker lets the target's deliveries through.
|
||||||
|
func (h *Handlers) pausedView(targetID string) *PausedView {
|
||||||
|
state, cooldown := h.breakers.StateAndCooldown(targetID)
|
||||||
|
|
||||||
|
switch {
|
||||||
|
case state == delivery.CircuitHalfOpen:
|
||||||
|
return &PausedView{}
|
||||||
|
case state == delivery.CircuitOpen && cooldown > 0:
|
||||||
|
return newPausedView(time.Now().Add(cooldown))
|
||||||
|
default:
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// deliveryPausedView reads the circuit breaker of a retrying delivery's
|
||||||
|
// target. While it is open, it says the earliest the delivery can be
|
||||||
|
// tried next: the later of the cooldown's end and the end of the
|
||||||
|
// delivery's own backoff after its last attempt. It is only the
|
||||||
|
// earliest: when the cooldown ends, one of the target's waiting
|
||||||
|
// deliveries is sent to test it while the others wait at least one more
|
||||||
|
// cooldown. Otherwise it returns nil, half-open included, since the
|
||||||
|
// delivery may then be the one being sent to test the target.
|
||||||
|
func (h *Handlers) deliveryPausedView(
|
||||||
|
targetID string, attempts []deliveryResultRow,
|
||||||
|
) *PausedView {
|
||||||
|
state, cooldown := h.breakers.StateAndCooldown(targetID)
|
||||||
|
if state != delivery.CircuitOpen || cooldown <= 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
next := time.Now().Add(cooldown)
|
||||||
|
|
||||||
|
if len(attempts) > 0 {
|
||||||
|
last := attempts[len(attempts)-1]
|
||||||
|
|
||||||
|
backoffEnd := last.CreatedAt.Add(delivery.Backoff(last.AttemptNum))
|
||||||
|
if backoffEnd.After(next) {
|
||||||
|
next = backoffEnd
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return newPausedView(next)
|
||||||
|
}
|
||||||
|
|
||||||
|
// newPausedView is a PausedView of deliveries paused until the given
|
||||||
|
// time. A time not on the current UTC day is written with its date, as
|
||||||
|
// the event log writes its times.
|
||||||
|
func newPausedView(until time.Time) *PausedView {
|
||||||
|
until = until.UTC()
|
||||||
|
|
||||||
|
layout := time.TimeOnly
|
||||||
|
if until.Format(time.DateOnly) != time.Now().UTC().Format(time.DateOnly) {
|
||||||
|
layout = time.DateTime
|
||||||
|
}
|
||||||
|
|
||||||
|
return &PausedView{
|
||||||
|
Until: until.Format(layout) + " UTC",
|
||||||
|
Relative: humanize.Time(until),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// TargetDeliveries is how many of a target's deliveries became
|
// TargetDeliveries is how many of a target's deliveries became
|
||||||
@@ -92,6 +170,8 @@ func (h *Handlers) targetRows(
|
|||||||
if targets[i].Type == database.TargetTypeDatabase {
|
if targets[i].Type == database.TargetTypeDatabase {
|
||||||
rows[i].Archive = h.archiveFileView(webhook, &targets[i], now)
|
rows[i].Archive = h.archiveFileView(webhook, &targets[i], now)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
rows[i].Paused = h.pausedView(targets[i].ID)
|
||||||
}
|
}
|
||||||
|
|
||||||
return rows
|
return rows
|
||||||
|
|||||||
@@ -0,0 +1,219 @@
|
|||||||
|
package handlers_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"net/http"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
"gorm.io/gorm/clause"
|
||||||
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
|
"sneak.berlin/go/webhooker/internal/delivery"
|
||||||
|
"sneak.berlin/go/webhooker/internal/handlers"
|
||||||
|
"sneak.berlin/go/webhooker/internal/session"
|
||||||
|
)
|
||||||
|
|
||||||
|
// cooldownEnds is how the pages write the end of a paused target's
|
||||||
|
// breaker's cooldown: the time in UTC, with its date when that falls on
|
||||||
|
// another UTC day, then how long that is from now.
|
||||||
|
const cooldownEnds = `(\d{4}-\d\d-\d\d )?\d\d:\d\d:\d\d UTC ` +
|
||||||
|
`\(\d+ seconds from now\)`
|
||||||
|
|
||||||
|
// TestPausedTarget_ShownUntilBreakerCloses takes an http target's
|
||||||
|
// circuit breaker from open through half-open to closed.
|
||||||
|
//
|
||||||
|
// Open, the target's row on the webhook page says its deliveries are
|
||||||
|
// paused and until when, and each retrying delivery says it is waiting
|
||||||
|
// and why in the event log and on the event's page, with the earliest
|
||||||
|
// it can be tried next: the later of the cooldown's end and the end of
|
||||||
|
// its own backoff, with the date when that is another UTC day.
|
||||||
|
// Half-open, the row says deliveries are held while one delivery tests
|
||||||
|
// the target, with no time, and no delivery says it is waiting. Closed,
|
||||||
|
// the pages say neither. The delivered delivery and the log target are
|
||||||
|
// shown as before throughout.
|
||||||
|
func TestPausedTarget_ShownUntilBreakerCloses(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
var (
|
||||||
|
h *handlers.Handlers
|
||||||
|
sess *session.Session
|
||||||
|
db *database.Database
|
||||||
|
dbMgr *database.WebhookDBManager
|
||||||
|
breakers *testCircuitBreakers
|
||||||
|
)
|
||||||
|
|
||||||
|
app := newTestApp(t, &h, &sess, &db, &dbMgr, &breakers)
|
||||||
|
app.RequireStart()
|
||||||
|
|
||||||
|
t.Cleanup(app.RequireStop)
|
||||||
|
|
||||||
|
wh := seedWebhook(t, db)
|
||||||
|
target := seedTarget(t, db, wh.ID, database.TargetTypeHTTP)
|
||||||
|
seedTarget(t, db, wh.ID, database.TargetTypeLog)
|
||||||
|
|
||||||
|
retrying := seedStoredEvent(t, dbMgr, wh.ID, `{"n":1}`)
|
||||||
|
addDelivery(t, dbMgr, wh.ID, retrying.ID, target.ID,
|
||||||
|
database.DeliveryStatusRetrying)
|
||||||
|
|
||||||
|
delivered := seedStoredEvent(t, dbMgr, wh.ID, `{"n":2}`)
|
||||||
|
addDelivery(t, dbMgr, wh.ID, delivered.ID, target.ID,
|
||||||
|
database.DeliveryStatusDelivered)
|
||||||
|
|
||||||
|
// This delivery's 18th attempt failed a minute ago, so its own
|
||||||
|
// backoff ends over a day from now: long after the cooldown, and on
|
||||||
|
// another UTC day, so the page shows the date.
|
||||||
|
backedOff := seedStoredEvent(t, dbMgr, wh.ID, `{"n":3}`)
|
||||||
|
backedOffID := addDelivery(t, dbMgr, wh.ID, backedOff.ID, target.ID,
|
||||||
|
database.DeliveryStatusRetrying)
|
||||||
|
|
||||||
|
failedAt := time.Now().Add(-time.Minute).Truncate(time.Second)
|
||||||
|
addFailedAttempt(t, dbMgr, wh.ID, backedOffID, 18, failedAt)
|
||||||
|
|
||||||
|
backoffEnds := failedAt.Add(delivery.Backoff(18)).UTC().
|
||||||
|
Format("2006-01-02 15:04:05") + " UTC (1 day from now)"
|
||||||
|
|
||||||
|
const waiting = "waiting: target paused after repeated failures, " +
|
||||||
|
"next try no earlier than "
|
||||||
|
|
||||||
|
breakers.Set(target.ID, delivery.CircuitOpen, 30*time.Second)
|
||||||
|
|
||||||
|
list := targetList(t, renderSourceDetailPage(t, h, sess, wh.ID))
|
||||||
|
assert.Regexp(t, "t-http http Active Edit Deactivate Delete "+
|
||||||
|
"Deliveries Paused: after repeated failures, until "+cooldownEnds+
|
||||||
|
", then one waiting delivery is sent to test the target while "+
|
||||||
|
"the others wait at least one more cooldown", list)
|
||||||
|
assert.Equal(t, 1, strings.Count(list, "Paused"))
|
||||||
|
|
||||||
|
log := renderSourceLogsPage(t, h, sess, wh.ID)
|
||||||
|
assert.Equal(t, 2, strings.Count(log, "t-http: waiting"))
|
||||||
|
assert.Contains(t, log, "t-http: delivered")
|
||||||
|
assert.Regexp(t, waiting+cooldownEnds, log)
|
||||||
|
assert.Contains(t, log, waiting+backoffEnds)
|
||||||
|
assert.NotContains(t, log, "retrying")
|
||||||
|
|
||||||
|
page := eventPage(t, h, sess, wh.ID, retrying.ID)
|
||||||
|
assert.Regexp(t, waiting+cooldownEnds, page)
|
||||||
|
assert.NotContains(t, page, "retrying")
|
||||||
|
|
||||||
|
page = eventPage(t, h, sess, wh.ID, backedOff.ID)
|
||||||
|
assert.Contains(t, page, waiting+backoffEnds)
|
||||||
|
assert.NotContains(t, page, "seconds from now")
|
||||||
|
|
||||||
|
breakers.Set(target.ID, delivery.CircuitHalfOpen, 0)
|
||||||
|
|
||||||
|
list = targetList(t, renderSourceDetailPage(t, h, sess, wh.ID))
|
||||||
|
assert.Contains(t, list, "t-http http Active Edit Deactivate Delete "+
|
||||||
|
"Deliveries Paused: held while one delivery tests whether the "+
|
||||||
|
"target has recovered")
|
||||||
|
// Not the whole list: the add target form above the rows says UTC.
|
||||||
|
assert.NotContains(t, targetRow(list, "t-http", "t-log"), "UTC")
|
||||||
|
|
||||||
|
assertRetryingNotWaiting(t, h, sess, wh.ID, retrying, backedOff)
|
||||||
|
|
||||||
|
breakers.Set(target.ID, delivery.CircuitClosed, 0)
|
||||||
|
|
||||||
|
list = targetList(t, renderSourceDetailPage(t, h, sess, wh.ID))
|
||||||
|
assert.NotContains(t, list, "Paused")
|
||||||
|
|
||||||
|
assertRetryingNotWaiting(t, h, sess, wh.ID, retrying, backedOff)
|
||||||
|
}
|
||||||
|
|
||||||
|
// targetRow returns the row of the target named name in a targetList:
|
||||||
|
// from its name to the name of the target listed after it, next.
|
||||||
|
func targetRow(list, name, next string) string {
|
||||||
|
_, row, _ := strings.Cut(list, name+" ")
|
||||||
|
row, _, _ = strings.Cut(row, next+" ")
|
||||||
|
|
||||||
|
return row
|
||||||
|
}
|
||||||
|
|
||||||
|
// assertRetryingNotWaiting checks that the event log and each event's
|
||||||
|
// page show the http target's delivery of the event as retrying, and
|
||||||
|
// none of them as waiting.
|
||||||
|
func assertRetryingNotWaiting(
|
||||||
|
t *testing.T,
|
||||||
|
h *handlers.Handlers,
|
||||||
|
sess *session.Session,
|
||||||
|
webhookID string,
|
||||||
|
events ...*database.Event,
|
||||||
|
) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
log := renderSourceLogsPage(t, h, sess, webhookID)
|
||||||
|
assert.Equal(t, len(events), strings.Count(log, "t-http: retrying"))
|
||||||
|
assert.NotContains(t, log, "waiting")
|
||||||
|
|
||||||
|
for _, event := range events {
|
||||||
|
page := eventPage(t, h, sess, webhookID, event.ID)
|
||||||
|
assert.Contains(t, page, ">retrying</span>")
|
||||||
|
assert.NotContains(t, page, "waiting")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// addDelivery records a delivery of the event to the target, with the
|
||||||
|
// given status, in the webhook's own database, and returns its ID.
|
||||||
|
func addDelivery(
|
||||||
|
t *testing.T,
|
||||||
|
dbMgr *database.WebhookDBManager,
|
||||||
|
webhookID, eventID, targetID string,
|
||||||
|
status database.DeliveryStatus,
|
||||||
|
) string {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
webhookDB, err := dbMgr.GetDB(webhookID)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
dlv := &database.Delivery{
|
||||||
|
EventID: eventID,
|
||||||
|
TargetID: targetID,
|
||||||
|
Status: status,
|
||||||
|
}
|
||||||
|
|
||||||
|
require.NoError(t, webhookDB.Omit(clause.Associations).Create(
|
||||||
|
dlv,
|
||||||
|
).Error)
|
||||||
|
|
||||||
|
return dlv.ID
|
||||||
|
}
|
||||||
|
|
||||||
|
// addFailedAttempt records the delivery's failed attempt attemptNum,
|
||||||
|
// made at the given time.
|
||||||
|
func addFailedAttempt(
|
||||||
|
t *testing.T,
|
||||||
|
dbMgr *database.WebhookDBManager,
|
||||||
|
webhookID, deliveryID string,
|
||||||
|
attemptNum int,
|
||||||
|
at time.Time,
|
||||||
|
) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
webhookDB, err := dbMgr.GetDB(webhookID)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
require.NoError(t, webhookDB.Omit(clause.Associations).Create(
|
||||||
|
&database.DeliveryResult{
|
||||||
|
BaseModel: database.BaseModel{CreatedAt: at},
|
||||||
|
DeliveryID: deliveryID,
|
||||||
|
AttemptNum: attemptNum,
|
||||||
|
Error: "connection refused",
|
||||||
|
},
|
||||||
|
).Error)
|
||||||
|
}
|
||||||
|
|
||||||
|
// eventPage runs the real event page handler and returns the
|
||||||
|
// rendered HTML.
|
||||||
|
func eventPage(
|
||||||
|
t *testing.T,
|
||||||
|
h *handlers.Handlers,
|
||||||
|
sess *session.Session,
|
||||||
|
webhookID, eventID string,
|
||||||
|
) string {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
w := serveEventPage(t, h, sess, webhookID, eventID)
|
||||||
|
require.Equal(t, http.StatusOK, w.Code)
|
||||||
|
|
||||||
|
return w.Body.String()
|
||||||
|
}
|
||||||
@@ -11,6 +11,7 @@ import (
|
|||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
@@ -142,6 +143,14 @@ func (n *noopArchives) Rename(_, _, _ string) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type noopCircuitBreakers struct{}
|
||||||
|
|
||||||
|
func (n *noopCircuitBreakers) StateAndCooldown(
|
||||||
|
string,
|
||||||
|
) (delivery.CircuitState, time.Duration) {
|
||||||
|
return delivery.CircuitClosed, 0
|
||||||
|
}
|
||||||
|
|
||||||
// newServerApp starts the real login path against dir: the handlers,
|
// newServerApp starts the real login path against dir: the handlers,
|
||||||
// the middleware that bounds password verification, the session store
|
// the middleware that bounds password verification, the session store
|
||||||
// and the database, exactly as internal/handlers builds them.
|
// and the database, exactly as internal/handlers builds them.
|
||||||
@@ -174,6 +183,9 @@ func newServerApp(
|
|||||||
session.New,
|
session.New,
|
||||||
func() delivery.Notifier { return &noopNotifier{} },
|
func() delivery.Notifier { return &noopNotifier{} },
|
||||||
func() delivery.Archives { return &noopArchives{} },
|
func() delivery.Archives { return &noopArchives{} },
|
||||||
|
func() delivery.CircuitBreakers {
|
||||||
|
return &noopCircuitBreakers{}
|
||||||
|
},
|
||||||
metrics.NewRegistry,
|
metrics.NewRegistry,
|
||||||
metrics.New,
|
metrics.New,
|
||||||
middleware.New,
|
middleware.New,
|
||||||
|
|||||||
@@ -18,10 +18,13 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/chromedp/cdproto/browser"
|
"github.com/chromedp/cdproto/browser"
|
||||||
|
"github.com/chromedp/cdproto/dom"
|
||||||
|
"github.com/chromedp/cdproto/input"
|
||||||
"github.com/chromedp/cdproto/log"
|
"github.com/chromedp/cdproto/log"
|
||||||
"github.com/chromedp/cdproto/network"
|
"github.com/chromedp/cdproto/network"
|
||||||
"github.com/chromedp/cdproto/runtime"
|
"github.com/chromedp/cdproto/runtime"
|
||||||
"github.com/chromedp/chromedp"
|
"github.com/chromedp/chromedp"
|
||||||
|
"github.com/chromedp/chromedp/kb"
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
"gorm.io/gorm/clause"
|
"gorm.io/gorm/clause"
|
||||||
@@ -40,6 +43,16 @@ const (
|
|||||||
phoneWidth = 390
|
phoneWidth = 390
|
||||||
phoneHeight = 844
|
phoneHeight = 844
|
||||||
|
|
||||||
|
// A window short enough that the event log scrolls with its last
|
||||||
|
// event expanded, and tall enough to show all of that event.
|
||||||
|
shortWidth = 1024
|
||||||
|
shortHeight = 450
|
||||||
|
|
||||||
|
// A person's double- or triple-click: each press is held a tenth of a
|
||||||
|
// second, and the next press comes a quarter second after the release.
|
||||||
|
pressHeld = 100 * time.Millisecond
|
||||||
|
betweenClicks = 250 * time.Millisecond
|
||||||
|
|
||||||
// olderBody is the body of the event received before the newest.
|
// olderBody is the body of the event received before the newest.
|
||||||
olderBody = "the older event"
|
olderBody = "the older event"
|
||||||
)
|
)
|
||||||
@@ -59,7 +72,7 @@ func TestAlpineRunsUnderTheSecurityPolicy(t *testing.T) {
|
|||||||
t.Cleanup(srv.Close)
|
t.Cleanup(srv.Close)
|
||||||
|
|
||||||
userID, _ := env.seedUser(t, "browser", "browser-password")
|
userID, _ := env.seedUser(t, "browser", "browser-password")
|
||||||
webhook, event, target := seedBrowserWebhook(t, env, userID)
|
webhook, older, event, target := seedBrowserWebhook(t, env, userID)
|
||||||
|
|
||||||
require.NoError(t, chromedp.Run(
|
require.NoError(t, chromedp.Run(
|
||||||
ctx, setCookies(srv.URL, env.authCookies(t, userID, "browser")),
|
ctx, setCookies(srv.URL, env.authCookies(t, userID, "browser")),
|
||||||
@@ -80,10 +93,10 @@ func TestAlpineRunsUnderTheSecurityPolicy(t *testing.T) {
|
|||||||
checkCopy(ctx, t, page)
|
checkCopy(ctx, t, page)
|
||||||
checkEntrypointEdit(ctx, t, page, page+"/events")
|
checkEntrypointEdit(ctx, t, page, page+"/events")
|
||||||
checkRecentEvents(ctx, t, page)
|
checkRecentEvents(ctx, t, page)
|
||||||
checkEventLog(ctx, t, page+"/events", event.ID, target.Name)
|
|
||||||
checkArchiveChoice(ctx, t, srv.URL+"/hooks/new", page)
|
checkArchiveChoice(ctx, t, srv.URL+"/hooks/new", page)
|
||||||
checkNewWebhookTargets(ctx, t, env, srv.URL+"/hooks/new")
|
checkNewWebhookTargets(ctx, t, env, srv.URL+"/hooks/new")
|
||||||
checkRefusedNewWebhook(ctx, t, srv.URL+"/hooks/new")
|
checkRefusedNewWebhook(ctx, t, srv.URL+"/hooks/new")
|
||||||
|
checkEventLog(ctx, t, page+"/events", event.ID, older.ID, target.Name)
|
||||||
checkMobileMenu(ctx, t, page)
|
checkMobileMenu(ctx, t, page)
|
||||||
|
|
||||||
assert.Empty(t, problems(), "the browser reported problems")
|
assert.Empty(t, problems(), "the browser reported problems")
|
||||||
@@ -91,11 +104,11 @@ func TestAlpineRunsUnderTheSecurityPolicy(t *testing.T) {
|
|||||||
|
|
||||||
// seedBrowserWebhook seeds the webhook the browser test loads, owned by
|
// seedBrowserWebhook seeds the webhook the browser test loads, owned by
|
||||||
// userID: an entrypoint, two events, and a target whose delivery of the
|
// userID: an entrypoint, two events, and a target whose delivery of the
|
||||||
// newer event failed once with a 502. It returns the webhook, the newer
|
// newer event failed once with a 502. It returns the webhook, the older
|
||||||
// event and the target.
|
// and the newer event, and the target.
|
||||||
func seedBrowserWebhook(
|
func seedBrowserWebhook(
|
||||||
t *testing.T, env *testEnv, userID string,
|
t *testing.T, env *testEnv, userID string,
|
||||||
) (*database.Webhook, *database.Event, *database.Target) {
|
) (*database.Webhook, *database.Event, *database.Event, *database.Target) {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
webhook := env.seedWebhook(t, userID)
|
webhook := env.seedWebhook(t, userID)
|
||||||
@@ -106,7 +119,7 @@ func seedBrowserWebhook(
|
|||||||
Active: true,
|
Active: true,
|
||||||
},
|
},
|
||||||
).Error)
|
).Error)
|
||||||
env.seedEvent(t, webhook.ID, olderBody)
|
older := env.seedEvent(t, webhook.ID, olderBody)
|
||||||
event := env.seedEvent(t, webhook.ID, `{"hello":"browser"}`)
|
event := env.seedEvent(t, webhook.ID, `{"hello":"browser"}`)
|
||||||
target := env.seedTarget(t, webhook.ID)
|
target := env.seedTarget(t, webhook.ID)
|
||||||
dlv := env.seedFailedDelivery(t, webhook.ID, event.ID, target.ID)
|
dlv := env.seedFailedDelivery(t, webhook.ID, event.ID, target.ID)
|
||||||
@@ -121,7 +134,7 @@ func seedBrowserWebhook(
|
|||||||
},
|
},
|
||||||
).Error)
|
).Error)
|
||||||
|
|
||||||
return webhook, event, target
|
return webhook, older, event, target
|
||||||
}
|
}
|
||||||
|
|
||||||
// startBrowser starts a headless browser for one test. It returns the
|
// startBrowser starts a headless browser for one test. It returns the
|
||||||
@@ -766,18 +779,30 @@ func checkRecentEvents(ctx context.Context, t *testing.T, url string) {
|
|||||||
"the event's own page does not show its body")
|
"the event's own page does not show its body")
|
||||||
}
|
}
|
||||||
|
|
||||||
// checkEventLog loads the event log and checks that clicking an event's
|
// checkEventLog loads the event log and checks an event's row. Clicking
|
||||||
// row expands it, that in there clicking its delivery shows the
|
// its ID expands the event, and in there clicking its delivery shows the
|
||||||
// delivery's attempts and clicking again hides them, and that clicking
|
// delivery's attempts and clicking again hides them. Clicking the row's
|
||||||
// the event's row again collapses it.
|
// caret collapses the event, clicking it again expands it, and clicking
|
||||||
|
// the ID again collapses it. While the event is expanded the row says so
|
||||||
|
// and its caret is turned up, and while it is collapsed neither. It then
|
||||||
|
// runs checkEventSelection on the log's last event, lastEventID, and
|
||||||
|
// checkEventKeyboard on eventID.
|
||||||
func checkEventLog(
|
func checkEventLog(
|
||||||
ctx context.Context, t *testing.T, url, eventID, targetName string,
|
ctx context.Context,
|
||||||
|
t *testing.T,
|
||||||
|
url, eventID, lastEventID, targetName string,
|
||||||
) {
|
) {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
// The event's row shows its ID, and its Resubmit form is in the part
|
// The event's row shows its ID and ends with its caret, which turns
|
||||||
// that expands. The delivery's row there shows the target's name.
|
// up with Tailwind's rotate-180 class, and its Resubmit form is in
|
||||||
eventRow := `//span[text()="` + eventID + `"]`
|
// the part that expands. The delivery's row there shows the target's
|
||||||
|
// name.
|
||||||
|
id := `//span[text()="` + eventID + `"]`
|
||||||
|
row := id + `/ancestor::div[@role="button"]`
|
||||||
|
caret := row + `//*[local-name()="svg"]`
|
||||||
|
caretUp := caret + `[contains(@class, "rotate-180")]`
|
||||||
|
caretDown := caret + `[not(contains(@class, "rotate-180"))]`
|
||||||
expanded := `form[action$="/` + eventID + `/resubmit"]`
|
expanded := `form[action$="/` + eventID + `/resubmit"]`
|
||||||
deliveryRow := `//span[text()="` + targetName + `"]`
|
deliveryRow := `//span[text()="` + targetName + `"]`
|
||||||
attempt := `//span[text()="Attempt 1"]`
|
attempt := `//span[text()="Attempt 1"]`
|
||||||
@@ -786,8 +811,13 @@ func checkEventLog(
|
|||||||
|
|
||||||
assert.True(t, hidden(ctx, expanded), "the event starts expanded")
|
assert.True(t, hidden(ctx, expanded), "the event starts expanded")
|
||||||
|
|
||||||
click(ctx, t, eventRow)
|
click(ctx, t, id)
|
||||||
assert.True(t, shown(ctx, expanded), "clicking the event does not expand it")
|
assert.True(t, shown(ctx, expanded),
|
||||||
|
"clicking the event's ID does not expand it")
|
||||||
|
assert.True(t, shown(ctx, row+`[@aria-expanded="true"]`),
|
||||||
|
"the expanded event's row does not say it is expanded")
|
||||||
|
assert.True(t, shown(ctx, caretUp),
|
||||||
|
"the expanded event's caret does not turn up")
|
||||||
|
|
||||||
assert.True(t, hidden(ctx, attempt), "the delivery's attempts start shown")
|
assert.True(t, hidden(ctx, attempt), "the delivery's attempts start shown")
|
||||||
|
|
||||||
@@ -799,9 +829,232 @@ func checkEventLog(
|
|||||||
assert.True(t, hidden(ctx, attempt),
|
assert.True(t, hidden(ctx, attempt),
|
||||||
"clicking the delivery again does not hide its attempts")
|
"clicking the delivery again does not hide its attempts")
|
||||||
|
|
||||||
click(ctx, t, eventRow)
|
click(ctx, t, caret)
|
||||||
assert.True(t, hidden(ctx, expanded),
|
assert.True(t, hidden(ctx, expanded),
|
||||||
"clicking the event again does not collapse it")
|
"clicking the caret does not collapse the event")
|
||||||
|
assert.True(t, shown(ctx, row+`[@aria-expanded="false"]`),
|
||||||
|
"the collapsed event's row does not say it is collapsed")
|
||||||
|
assert.True(t, shown(ctx, caretDown),
|
||||||
|
"the collapsed event's caret stays turned up")
|
||||||
|
|
||||||
|
click(ctx, t, caret)
|
||||||
|
assert.True(t, shown(ctx, expanded),
|
||||||
|
"clicking the caret again does not expand the event")
|
||||||
|
|
||||||
|
click(ctx, t, id)
|
||||||
|
assert.True(t, hidden(ctx, expanded),
|
||||||
|
"clicking the event's ID again does not collapse it")
|
||||||
|
|
||||||
|
checkEventSelection(ctx, t, url, lastEventID)
|
||||||
|
checkEventKeyboard(ctx, t, url, eventID)
|
||||||
|
}
|
||||||
|
|
||||||
|
// checkEventSelection loads the event log in a short window and checks
|
||||||
|
// that selecting the ID of its last event, eventID, with the mouse leaves
|
||||||
|
// the event as it was, and that its caret toggles it at once. Dragging
|
||||||
|
// over the ID leaves the event collapsed, and the caret's click expands
|
||||||
|
// it at once. With the page then scrolled to its end, a double-click on
|
||||||
|
// the ID that goes on to drag along it, and a triple-click on it, each
|
||||||
|
// leave the event expanded and select that ID. Had a click there
|
||||||
|
// collapsed the event, the page would have got shorter and moved under
|
||||||
|
// the pointer before the next click.
|
||||||
|
func checkEventSelection(
|
||||||
|
ctx context.Context, t *testing.T, url, eventID string,
|
||||||
|
) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
id := `//span[text()="` + eventID + `"]`
|
||||||
|
row := id + `/ancestor::div[@role="button"]`
|
||||||
|
caret := row + `//*[local-name()="svg"]`
|
||||||
|
expanded := `form[action$="/` + eventID + `/resubmit"]`
|
||||||
|
|
||||||
|
var (
|
||||||
|
selected, state string
|
||||||
|
hasState bool
|
||||||
|
scrolled float64
|
||||||
|
)
|
||||||
|
|
||||||
|
// What is selected, and whether the event's row says it is expanded.
|
||||||
|
read := chromedp.Tasks{
|
||||||
|
chromedp.Evaluate(`window.getSelection().toString()`, &selected),
|
||||||
|
chromedp.AttributeValue(
|
||||||
|
row, "aria-expanded", &state, &hasState, chromedp.BySearch,
|
||||||
|
),
|
||||||
|
}
|
||||||
|
|
||||||
|
// A click on the ID toggles the event half a second after it, so a
|
||||||
|
// check that selecting the ID did not toggle it waits a second first.
|
||||||
|
settle := chromedp.Sleep(time.Second)
|
||||||
|
|
||||||
|
// The double-click and the triple-click each start with nothing
|
||||||
|
// selected, so that their first click waits to toggle the event.
|
||||||
|
clearSelection := chromedp.Evaluate(
|
||||||
|
`window.getSelection().removeAllRanges()`, nil,
|
||||||
|
)
|
||||||
|
|
||||||
|
require.NoError(t, chromedp.Run(
|
||||||
|
ctx, chromedp.EmulateViewport(shortWidth, shortHeight), loadPage(url),
|
||||||
|
))
|
||||||
|
|
||||||
|
selectText(ctx, t, id)
|
||||||
|
require.NoError(t, chromedp.Run(ctx, settle, read))
|
||||||
|
assert.Equal(t, eventID, selected, "the event's ID cannot be selected")
|
||||||
|
require.True(t, hasState, "the event's row does not say if it is expanded")
|
||||||
|
assert.Equal(t, "false", state, "selecting the event's ID expands it")
|
||||||
|
|
||||||
|
click(ctx, t, caret)
|
||||||
|
require.NoError(t, chromedp.Run(ctx, read))
|
||||||
|
assert.Equal(t, "true", state,
|
||||||
|
"clicking the caret does not expand the event at once")
|
||||||
|
|
||||||
|
// The row says it is expanded before its expanded part is shown, so
|
||||||
|
// the scroll waits for that part, to end at the expanded page's end.
|
||||||
|
require.True(t, shown(ctx, expanded),
|
||||||
|
"clicking the caret does not show the event's expanded part")
|
||||||
|
|
||||||
|
require.NoError(t, chromedp.Run(ctx, chromedp.Evaluate(
|
||||||
|
`window.scrollTo(0, document.body.scrollHeight); window.scrollY`,
|
||||||
|
&scrolled,
|
||||||
|
)))
|
||||||
|
require.Positive(t, scrolled, "the event log does not scroll")
|
||||||
|
|
||||||
|
require.NoError(t, chromedp.Run(ctx, clearSelection))
|
||||||
|
doubleClickAndDrag(ctx, t, id)
|
||||||
|
require.NoError(t, chromedp.Run(ctx, settle, read))
|
||||||
|
assert.Contains(t, selected, eventID,
|
||||||
|
"a double-click and drag does not select the event's ID")
|
||||||
|
assert.Equal(t, "true", state,
|
||||||
|
"a double-click and drag over the event's ID collapses it")
|
||||||
|
|
||||||
|
require.NoError(t, chromedp.Run(ctx, clearSelection))
|
||||||
|
tripleClick(ctx, t, id)
|
||||||
|
require.NoError(t, chromedp.Run(ctx, settle, read))
|
||||||
|
assert.Contains(t, selected, eventID,
|
||||||
|
"a triple-click does not select the event's ID")
|
||||||
|
assert.Equal(t, "true", state,
|
||||||
|
"a triple-click selecting the event's ID collapses it")
|
||||||
|
}
|
||||||
|
|
||||||
|
// checkEventKeyboard loads the event log and checks that Tab from the
|
||||||
|
// page's Back link reaches the event's row, the first after it, and that
|
||||||
|
// Enter then expands the event and Space collapses it.
|
||||||
|
func checkEventKeyboard(
|
||||||
|
ctx context.Context, t *testing.T, url, eventID string,
|
||||||
|
) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
back := `//a[contains(text(), "Back to")]`
|
||||||
|
expanded := `form[action$="/` + eventID + `/resubmit"]`
|
||||||
|
|
||||||
|
var focused string
|
||||||
|
|
||||||
|
require.NoError(t, chromedp.Run(
|
||||||
|
ctx,
|
||||||
|
loadPage(url),
|
||||||
|
chromedp.Focus(back, chromedp.BySearch),
|
||||||
|
chromedp.KeyEvent(kb.Tab),
|
||||||
|
chromedp.Evaluate(`document.activeElement.textContent`, &focused),
|
||||||
|
))
|
||||||
|
require.Contains(t, focused, eventID,
|
||||||
|
"Tab from the Back link does not reach the event's row")
|
||||||
|
|
||||||
|
require.NoError(t, chromedp.Run(ctx, chromedp.KeyEvent(kb.Enter)))
|
||||||
|
assert.True(t, shown(ctx, expanded), "Enter does not expand the event")
|
||||||
|
|
||||||
|
require.NoError(t, chromedp.Run(ctx, chromedp.KeyEvent(" ")))
|
||||||
|
assert.True(t, hidden(ctx, expanded), "Space does not collapse the event")
|
||||||
|
}
|
||||||
|
|
||||||
|
// selectText selects the text of the element matching an XPath
|
||||||
|
// expression as a person does with the mouse: pressing the button at the
|
||||||
|
// text's start, moving to its end and releasing it there.
|
||||||
|
func selectText(ctx context.Context, t *testing.T, xpath string) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
left, right, y := textEnds(ctx, t, xpath)
|
||||||
|
|
||||||
|
require.NoError(t, chromedp.Run(
|
||||||
|
ctx, press(left, y, 1), drag(right, y), release(right, y, 1),
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
|
// doubleClickAndDrag double-clicks the start of the text of the element
|
||||||
|
// matching an XPath expression, which selects its first word, and keeps
|
||||||
|
// the button down to drag to the text's end, which selects it word by
|
||||||
|
// word. It holds the button for a second, longer than a single click on
|
||||||
|
// an event's row waits before it toggles the event.
|
||||||
|
func doubleClickAndDrag(ctx context.Context, t *testing.T, xpath string) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
left, right, y := textEnds(ctx, t, xpath)
|
||||||
|
|
||||||
|
require.NoError(t, chromedp.Run(
|
||||||
|
ctx,
|
||||||
|
press(left, y, 1), chromedp.Sleep(pressHeld), release(left, y, 1),
|
||||||
|
chromedp.Sleep(betweenClicks),
|
||||||
|
press(left, y, 2), drag(right, y), chromedp.Sleep(time.Second),
|
||||||
|
release(right, y, 2),
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
|
// tripleClick clicks three times in the middle of the text of the
|
||||||
|
// element matching an XPath expression, as a person does to select a
|
||||||
|
// whole line of text. The browser selects a word on the second click and
|
||||||
|
// the whole paragraph on the third.
|
||||||
|
func tripleClick(ctx context.Context, t *testing.T, xpath string) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
left, right, y := textEnds(ctx, t, xpath)
|
||||||
|
x := (left + right) / 2
|
||||||
|
|
||||||
|
require.NoError(t, chromedp.Run(
|
||||||
|
ctx,
|
||||||
|
press(x, y, 1), chromedp.Sleep(pressHeld), release(x, y, 1),
|
||||||
|
chromedp.Sleep(betweenClicks),
|
||||||
|
press(x, y, 2), chromedp.Sleep(pressHeld), release(x, y, 2),
|
||||||
|
chromedp.Sleep(betweenClicks),
|
||||||
|
press(x, y, 3), chromedp.Sleep(pressHeld), release(x, y, 3),
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
|
// textEnds returns where on screen the text of the element matching an
|
||||||
|
// XPath expression starts and ends, just inside its left and right
|
||||||
|
// edges, and the height of its middle: in that order, the x of its
|
||||||
|
// start, the x of its end, and the y of both.
|
||||||
|
func textEnds(
|
||||||
|
ctx context.Context, t *testing.T, xpath string,
|
||||||
|
) (float64, float64, float64) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
var box *dom.BoxModel
|
||||||
|
|
||||||
|
require.NoError(t, chromedp.Run(
|
||||||
|
ctx, chromedp.Dimensions(xpath, &box, chromedp.BySearch),
|
||||||
|
))
|
||||||
|
|
||||||
|
// The content box's corners, clockwise from its top left.
|
||||||
|
return box.Content[0] + 1, box.Content[2] - 1,
|
||||||
|
(box.Content[1] + box.Content[5]) / 2
|
||||||
|
}
|
||||||
|
|
||||||
|
// press presses the left mouse button at x, y, as the nth click of a
|
||||||
|
// double- or triple-click.
|
||||||
|
func press(x, y float64, nth int64) *input.DispatchMouseEventParams {
|
||||||
|
return input.DispatchMouseEvent(input.MousePressed, x, y).
|
||||||
|
WithButton(input.Left).WithButtons(1).WithClickCount(nth)
|
||||||
|
}
|
||||||
|
|
||||||
|
// drag moves the pointer to x, y with the left mouse button down.
|
||||||
|
func drag(x, y float64) *input.DispatchMouseEventParams {
|
||||||
|
return input.DispatchMouseEvent(input.MouseMoved, x, y).
|
||||||
|
WithButton(input.Left).WithButtons(1)
|
||||||
|
}
|
||||||
|
|
||||||
|
// release releases the left mouse button at x, y, as the nth click of a
|
||||||
|
// double- or triple-click.
|
||||||
|
func release(x, y float64, nth int64) *input.DispatchMouseEventParams {
|
||||||
|
return input.DispatchMouseEvent(input.MouseReleased, x, y).
|
||||||
|
WithButton(input.Left).WithClickCount(nth)
|
||||||
}
|
}
|
||||||
|
|
||||||
// The parts of the new webhook page the checks below find and click.
|
// The parts of the new webhook page the checks below find and click.
|
||||||
|
|||||||
@@ -62,6 +62,17 @@ func (e *noopArchives) Rename(_, _, _ string) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// noopCircuitBreakers satisfies handlers.New's
|
||||||
|
// delivery.CircuitBreakers dependency with no target's deliveries
|
||||||
|
// paused.
|
||||||
|
type noopCircuitBreakers struct{}
|
||||||
|
|
||||||
|
func (b *noopCircuitBreakers) StateAndCooldown(
|
||||||
|
string,
|
||||||
|
) (delivery.CircuitState, time.Duration) {
|
||||||
|
return delivery.CircuitClosed, 0
|
||||||
|
}
|
||||||
|
|
||||||
// testEnv is the real router from routes.go plus the collaborators
|
// testEnv is the real router from routes.go plus the collaborators
|
||||||
// tests need to seed users and forge sessions.
|
// tests need to seed users and forge sessions.
|
||||||
type testEnv struct {
|
type testEnv struct {
|
||||||
@@ -126,6 +137,9 @@ func newTestEnvWithConfig(
|
|||||||
session.New,
|
session.New,
|
||||||
func() delivery.Notifier { return &noopNotifier{} },
|
func() delivery.Notifier { return &noopNotifier{} },
|
||||||
func() delivery.Archives { return &noopArchives{} },
|
func() delivery.Archives { return &noopArchives{} },
|
||||||
|
func() delivery.CircuitBreakers {
|
||||||
|
return &noopCircuitBreakers{}
|
||||||
|
},
|
||||||
metrics.NewRegistry,
|
metrics.NewRegistry,
|
||||||
metrics.New,
|
metrics.New,
|
||||||
middleware.New,
|
middleware.New,
|
||||||
|
|||||||
@@ -77,12 +77,36 @@ document.addEventListener("alpine:init", function () {
|
|||||||
window.Alpine.data("collapsible", function () {
|
window.Alpine.data("collapsible", function () {
|
||||||
return {
|
return {
|
||||||
open: false,
|
open: false,
|
||||||
|
// The timer of the toggle a single click is waiting to make.
|
||||||
|
pendingToggle: null,
|
||||||
init() {
|
init() {
|
||||||
this.open = this.$root.hasAttribute("data-open");
|
this.open = this.$root.hasAttribute("data-open");
|
||||||
},
|
},
|
||||||
toggle() {
|
toggle() {
|
||||||
this.open = !this.open;
|
this.open = !this.open;
|
||||||
},
|
},
|
||||||
|
// Toggles on a click, except one that selects text, such as
|
||||||
|
// selecting an event's ID to copy it. A single click toggles
|
||||||
|
// only after 500 ms, the usual double-click interval, and the
|
||||||
|
// second press of a double- or triple-click cancels that (see
|
||||||
|
// cancelPendingToggle), so nothing moves under the pointer
|
||||||
|
// while it selects text.
|
||||||
|
toggleUnlessSelecting(event) {
|
||||||
|
if (
|
||||||
|
event.detail === 1 &&
|
||||||
|
window.getSelection().toString() === ""
|
||||||
|
) {
|
||||||
|
this.pendingToggle = setTimeout(() => this.toggle(), 500);
|
||||||
|
}
|
||||||
|
},
|
||||||
|
// Runs when the mouse button goes down, so the second press
|
||||||
|
// of a double- or triple-click cancels the toggle its first
|
||||||
|
// click is waiting to make, however long that press lasts.
|
||||||
|
cancelPendingToggle(event) {
|
||||||
|
if (event.detail > 1) {
|
||||||
|
clearTimeout(this.pendingToggle);
|
||||||
|
}
|
||||||
|
},
|
||||||
get closed() {
|
get closed() {
|
||||||
return !this.open;
|
return !this.open;
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -5,11 +5,15 @@
|
|||||||
<p class="text-xs text-gray-500">{{.AttemptsOmitted}} attempt{{if ne .AttemptsOmitted 1}}s{{end}} omitted between the first and last shown.</p>
|
<p class="text-xs text-gray-500">{{.AttemptsOmitted}} attempt{{if ne .AttemptsOmitted 1}}s{{end}} omitted between the first and last shown.</p>
|
||||||
{{end}}
|
{{end}}
|
||||||
{{range .Results}}
|
{{range .Results}}
|
||||||
|
{{/* Inside this loop, dot is one attempt and $ the whole delivery. */}}
|
||||||
<div class="rounded-md bg-white border border-gray-200 p-2">
|
<div class="rounded-md bg-white border border-gray-200 p-2">
|
||||||
<div class="flex flex-wrap items-center gap-3 text-xs">
|
<div class="flex flex-wrap items-center gap-3 text-xs">
|
||||||
<span class="text-gray-500">Attempt {{.AttemptNum}}</span>
|
<span class="text-gray-500">Attempt {{.AttemptNum}}</span>
|
||||||
<span class="{{if .Success}}text-green-600{{else}}text-red-600{{end}}">{{if .Success}}success{{else}}failure{{end}}</span>
|
<span class="{{if .Success}}text-green-600{{else}}text-red-600{{end}}">{{if not .Success}}failure{{else if eq $.Target.Type "database"}}archived{{else if eq $.Target.Type "log"}}written to the log{{else}}success{{end}}</span>
|
||||||
|
{{/* A database or log target sends no HTTP request, so its attempts have no status code. */}}
|
||||||
|
{{if not (eq $.Target.Type "database" "log")}}
|
||||||
<span class="text-gray-500">Status: {{if .HasStatusCode}}{{.StatusCode}}{{else}}— (no response){{end}}</span>
|
<span class="text-gray-500">Status: {{if .HasStatusCode}}{{.StatusCode}}{{else}}— (no response){{end}}</span>
|
||||||
|
{{end}}
|
||||||
<span class="text-gray-500">Duration: {{.DurationMS}} ms</span>
|
<span class="text-gray-500">Duration: {{.DurationMS}} ms</span>
|
||||||
</div>
|
</div>
|
||||||
{{if .Error}}
|
{{if .Error}}
|
||||||
|
|||||||
@@ -69,7 +69,7 @@
|
|||||||
<div class="flex flex-wrap items-center justify-between gap-3">
|
<div class="flex flex-wrap items-center justify-between gap-3">
|
||||||
<span class="text-sm text-gray-700">{{.Target.DisplayName}}</span>
|
<span class="text-sm text-gray-700">{{.Target.DisplayName}}</span>
|
||||||
<span class="flex flex-wrap items-center gap-3">
|
<span class="flex flex-wrap items-center gap-3">
|
||||||
<span class="text-xs {{if eq .Status "delivered"}}text-green-600{{else if eq .Status "failed"}}text-red-600{{else if eq .Status "retrying"}}text-yellow-600{{else}}text-gray-400{{end}}">{{.Status}}</span>
|
<span class="text-xs {{if eq .Status "delivered"}}text-green-600{{else if eq .Status "failed"}}text-red-600{{else if eq .Status "retrying"}}text-yellow-600{{else}}text-gray-400{{end}}">{{with .Paused}}waiting: target paused after repeated failures, next try no earlier than {{.Until}} ({{.Relative}}){{else}}{{.Status}}{{end}}</span>
|
||||||
<span class="text-xs text-gray-400">{{.AttemptCount}} attempt{{if ne .AttemptCount 1}}s{{end}}</span>
|
<span class="text-xs text-gray-400">{{.AttemptCount}} attempt{{if ne .AttemptCount 1}}s{{end}}</span>
|
||||||
</span>
|
</span>
|
||||||
</div>
|
</div>
|
||||||
|
|||||||
@@ -269,6 +269,12 @@
|
|||||||
</form>
|
</form>
|
||||||
</div>
|
</div>
|
||||||
</div>
|
</div>
|
||||||
|
{{with .Paused}}
|
||||||
|
<div class="text-xs text-yellow-600 mt-1">
|
||||||
|
<span class="font-medium">Deliveries Paused:</span>
|
||||||
|
<span>{{if .Until}}after repeated failures, until {{.Until}} ({{.Relative}}), then one waiting delivery is sent to test the target while the others wait at least one more cooldown{{else}}held while one delivery tests whether the target has recovered{{end}}</span>
|
||||||
|
</div>
|
||||||
|
{{end}}
|
||||||
{{range .Config}}
|
{{range .Config}}
|
||||||
<div class="text-xs text-gray-500 break-all mt-1">
|
<div class="text-xs text-gray-500 break-all mt-1">
|
||||||
<span class="font-medium text-gray-700">{{.Label}}:</span>
|
<span class="font-medium text-gray-700">{{.Label}}:</span>
|
||||||
|
|||||||
@@ -16,7 +16,8 @@
|
|||||||
<div class="divide-y divide-gray-100">
|
<div class="divide-y divide-gray-100">
|
||||||
{{range .Events}}
|
{{range .Events}}
|
||||||
<div class="p-4" x-data="collapsible">
|
<div class="p-4" x-data="collapsible">
|
||||||
<button type="button" class="btn-small w-full flex flex-wrap justify-between gap-2 text-left" @click="toggle">
|
<!-- Not a button element: browsers do not let a button's text be selected, and an event's ID must be. -->
|
||||||
|
<div role="button" tabindex="0" class="btn-small w-full flex flex-wrap justify-between gap-2" :aria-expanded="open" @mousedown="cancelPendingToggle" @click="toggleUnlessSelecting" @keydown.enter.prevent="toggle" @keydown.space.prevent="toggle">
|
||||||
<span class="flex flex-wrap items-center gap-3">
|
<span class="flex flex-wrap items-center gap-3">
|
||||||
<span class="badge-info">{{.Method}}</span>
|
<span class="badge-info">{{.Method}}</span>
|
||||||
<span class="text-sm font-mono text-gray-700">{{.ID}}</span>
|
<span class="text-sm font-mono text-gray-700">{{.ID}}</span>
|
||||||
@@ -31,15 +32,16 @@
|
|||||||
<span class="flex flex-wrap items-center gap-4">
|
<span class="flex flex-wrap items-center gap-4">
|
||||||
{{range .Deliveries}}
|
{{range .Deliveries}}
|
||||||
<span class="text-xs {{if eq .Status "delivered"}}text-green-600{{else if eq .Status "failed"}}text-red-600{{else if eq .Status "retrying"}}text-yellow-600{{else}}text-gray-400{{end}}">
|
<span class="text-xs {{if eq .Status "delivered"}}text-green-600{{else if eq .Status "failed"}}text-red-600{{else if eq .Status "retrying"}}text-yellow-600{{else}}text-gray-400{{end}}">
|
||||||
{{.Target.DisplayName}}: {{.Status}}
|
{{.Target.DisplayName}}: {{if .Paused}}waiting{{else}}{{.Status}}{{end}}
|
||||||
</span>
|
</span>
|
||||||
{{end}}
|
{{end}}
|
||||||
<span class="text-xs text-gray-400">{{.CreatedAt.Format "2006-01-02 15:04:05"}}</span>
|
<span class="text-xs text-gray-400">{{.CreatedAt.Format "2006-01-02 15:04:05"}}</span>
|
||||||
<svg class="w-4 h-4 text-gray-400 transition-transform" :class="caretClass" fill="none" stroke="currentColor" viewBox="0 0 24 24">
|
<!-- The caret has no text to select, so a click on it toggles at once. -->
|
||||||
|
<svg class="w-4 h-4 text-gray-400 transition-transform" :class="caretClass" @click.stop="toggle" fill="none" stroke="currentColor" viewBox="0 0 24 24">
|
||||||
<path stroke-linecap="round" stroke-linejoin="round" stroke-width="2" d="M19 9l-7 7-7-7"/>
|
<path stroke-linecap="round" stroke-linejoin="round" stroke-width="2" d="M19 9l-7 7-7-7"/>
|
||||||
</svg>
|
</svg>
|
||||||
</span>
|
</span>
|
||||||
</button>
|
</div>
|
||||||
|
|
||||||
<div x-show="open" x-cloak class="mt-3 p-3 bg-gray-50 rounded-md">
|
<div x-show="open" x-cloak class="mt-3 p-3 bg-gray-50 rounded-md">
|
||||||
<div class="mb-3 flex flex-wrap items-center justify-between gap-2">
|
<div class="mb-3 flex flex-wrap items-center justify-between gap-2">
|
||||||
@@ -65,7 +67,7 @@
|
|||||||
<button type="button" class="btn-small flex-1 flex-wrap justify-between gap-2 text-left" @click="toggle">
|
<button type="button" class="btn-small flex-1 flex-wrap justify-between gap-2 text-left" @click="toggle">
|
||||||
<span class="flex flex-wrap items-center gap-3">
|
<span class="flex flex-wrap items-center gap-3">
|
||||||
<span class="text-sm text-gray-700">{{.Target.DisplayName}}</span>
|
<span class="text-sm text-gray-700">{{.Target.DisplayName}}</span>
|
||||||
<span class="text-xs {{if eq .Status "delivered"}}text-green-600{{else if eq .Status "failed"}}text-red-600{{else if eq .Status "retrying"}}text-yellow-600{{else}}text-gray-400{{end}}">{{.Status}}</span>
|
<span class="text-xs {{if eq .Status "delivered"}}text-green-600{{else if eq .Status "failed"}}text-red-600{{else if eq .Status "retrying"}}text-yellow-600{{else}}text-gray-400{{end}}">{{with .Paused}}waiting: target paused after repeated failures, next try no earlier than {{.Until}} ({{.Relative}}){{else}}{{.Status}}{{end}}</span>
|
||||||
</span>
|
</span>
|
||||||
<span class="flex flex-wrap items-center gap-3">
|
<span class="flex flex-wrap items-center gap-3">
|
||||||
<span class="text-xs text-gray-400">{{.AttemptCount}} attempt{{if ne .AttemptCount 1}}s{{end}}</span>
|
<span class="text-xs text-gray-400">{{.AttemptCount}} attempt{{if ne .AttemptCount 1}}s{{end}}</span>
|
||||||
@@ -74,7 +76,7 @@
|
|||||||
</svg>
|
</svg>
|
||||||
</span>
|
</span>
|
||||||
</button>
|
</button>
|
||||||
{{if .Status.Terminal}}
|
{{if and .Status.Terminal (not .Target.Deleted)}}
|
||||||
<form method="POST" action="/hook/{{$.Webhook.ID}}/deliveries/{{.ID}}/replay" class="inline">
|
<form method="POST" action="/hook/{{$.Webhook.ID}}/deliveries/{{.ID}}/replay" class="inline">
|
||||||
<input type="hidden" name="csrf_token" value="{{$.CSRFToken}}">
|
<input type="hidden" name="csrf_token" value="{{$.CSRFToken}}">
|
||||||
<input type="hidden" name="page" value="{{$.Page}}">
|
<input type="hidden" name="page" value="{{$.Page}}">
|
||||||
|
|||||||
Reference in New Issue
Block a user