Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
42946711b1 | ||
|
|
d2ecb83923 | ||
|
|
643077021d | ||
|
|
22fa502638 |
@@ -1700,7 +1700,6 @@ events should be forwarded.
|
|||||||
| `active` | boolean | Whether deliveries are enabled (default: true) |
|
| `active` | boolean | Whether deliveries are enabled (default: true) |
|
||||||
| `config` | JSON text | Type-specific configuration |
|
| `config` | JSON text | Type-specific configuration |
|
||||||
| `max_retries` | integer | Total delivery attempts for `http` and `slack` targets, not retries on top of the first: 0 is a single fire-and-forget attempt with no retries and no circuit breaker, and a value of N makes N attempts in all, with exponential backoff and a per-target circuit breaker. Ignored by `database` and `log` targets |
|
| `max_retries` | integer | Total delivery attempts for `http` and `slack` targets, not retries on top of the first: 0 is a single fire-and-forget attempt with no retries and no circuit breaker, and a value of N makes N attempts in all, with exponential backoff and a per-target circuit breaker. Ignored by `database` and `log` targets |
|
||||||
| `max_queue_size` | integer | Stored and shown on the target's detail view, but not enforced anywhere yet: nothing in the delivery engine consults it. Queue depth is set by the two fixed 10,000-entry channels |
|
|
||||||
|
|
||||||
**Relations:** Belongs to Webhook. Has many Deliveries.
|
**Relations:** Belongs to Webhook. Has many Deliveries.
|
||||||
|
|
||||||
@@ -1858,7 +1857,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
|
||||||
@@ -2368,6 +2371,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(
|
||||||
|
|||||||
@@ -31,8 +31,7 @@ type Target struct {
|
|||||||
|
|
||||||
// For HTTP targets (max_retries=0 means fire-and-forget,
|
// For HTTP targets (max_retries=0 means fire-and-forget,
|
||||||
// >0 enables retries with backoff)
|
// >0 enables retries with backoff)
|
||||||
MaxRetries int `json:"maxRetries,omitempty"`
|
MaxRetries int `json:"maxRetries,omitempty"`
|
||||||
MaxQueueSize int `json:"maxQueueSize,omitempty"`
|
|
||||||
|
|
||||||
// Relations. No model marshals the record it belongs to:
|
// Relations. No model marshals the record it belongs to:
|
||||||
// Webhook.Targets leads back here, and the JSON could loop.
|
// Webhook.Targets leads back here, and the JSON could loop.
|
||||||
|
|||||||
@@ -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.
|
||||||
@@ -300,6 +311,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{
|
||||||
|
|||||||
@@ -173,13 +173,6 @@ func httpConfigFields(t *database.Target) []ConfigField {
|
|||||||
|
|
||||||
fields = append(fields, maxRetriesField(t))
|
fields = append(fields, maxRetriesField(t))
|
||||||
|
|
||||||
if t.MaxQueueSize > 0 {
|
|
||||||
fields = append(fields, ConfigField{
|
|
||||||
Label: "Max Queue Size",
|
|
||||||
Value: strconv.Itoa(t.MaxQueueSize),
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
return fields
|
return fields
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -197,14 +197,12 @@ func TestNewTargetViews_Slack(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// TestNewTargetViews_SlackRetries proves a Slack target shows
|
// TestNewTargetViews_SlackRetries proves a Slack target shows
|
||||||
// its retry count the same way an HTTP target does, and no
|
// its retry count the same way an HTTP target does.
|
||||||
// queue size even when one is stored: delivery never reads it.
|
|
||||||
func TestNewTargetViews_SlackRetries(t *testing.T) {
|
func TestNewTargetViews_SlackRetries(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
target := slackTarget()
|
target := slackTarget()
|
||||||
target.MaxRetries = 2
|
target.MaxRetries = 2
|
||||||
target.MaxQueueSize = 100
|
|
||||||
|
|
||||||
view := viewFor(t, target)
|
view := viewFor(t, target)
|
||||||
|
|
||||||
@@ -226,8 +224,7 @@ func TestNewTargetViews_HTTP(t *testing.T) {
|
|||||||
Config: `{"url":"` + viewExampleHook + `",` +
|
Config: `{"url":"` + viewExampleHook + `",` +
|
||||||
`"timeout":30,` +
|
`"timeout":30,` +
|
||||||
`"headers":{"Authorization":"Bearer sekrit"}}`,
|
`"headers":{"Authorization":"Bearer sekrit"}}`,
|
||||||
MaxRetries: 5,
|
MaxRetries: 5,
|
||||||
MaxQueueSize: 100,
|
|
||||||
})
|
})
|
||||||
|
|
||||||
fields := fieldMap(view.Config)
|
fields := fieldMap(view.Config)
|
||||||
@@ -239,7 +236,6 @@ func TestNewTargetViews_HTTP(t *testing.T) {
|
|||||||
"Timeout": "30s",
|
"Timeout": "30s",
|
||||||
"Headers": "1 configured",
|
"Headers": "1 configured",
|
||||||
viewMaxRetries: "5",
|
viewMaxRetries: "5",
|
||||||
"Max Queue Size": "100",
|
|
||||||
},
|
},
|
||||||
fields,
|
fields,
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|
||||||
@@ -148,6 +150,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
|
||||||
@@ -1290,7 +1294,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]
|
||||||
@@ -1414,8 +1418,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,
|
||||||
@@ -1438,6 +1443,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
|
||||||
|
|||||||
@@ -24,6 +24,84 @@ type TargetRowView struct {
|
|||||||
// Archive is a database target's archive file, and nil for a target
|
// Archive is a database target's archive file, and nil for a target
|
||||||
// of any other type.
|
// 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
|
||||||
@@ -84,6 +162,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])
|
rows[i].Archive = h.archiveFileView(webhook, &targets[i])
|
||||||
}
|
}
|
||||||
|
|
||||||
|
rows[i].Paused = h.pausedView(targets[i].ID)
|
||||||
}
|
}
|
||||||
|
|
||||||
return rows
|
return rows
|
||||||
|
|||||||
@@ -0,0 +1,209 @@
|
|||||||
|
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")
|
||||||
|
assert.NotContains(t, list, "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)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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,
|
||||||
|
|||||||
@@ -852,6 +852,7 @@ func checkEventSelection(
|
|||||||
id := `//span[text()="` + eventID + `"]`
|
id := `//span[text()="` + eventID + `"]`
|
||||||
row := id + `/ancestor::div[@role="button"]`
|
row := id + `/ancestor::div[@role="button"]`
|
||||||
caret := row + `//*[local-name()="svg"]`
|
caret := row + `//*[local-name()="svg"]`
|
||||||
|
expanded := `form[action$="/` + eventID + `/resubmit"]`
|
||||||
|
|
||||||
var (
|
var (
|
||||||
selected, state string
|
selected, state string
|
||||||
@@ -892,6 +893,11 @@ func checkEventSelection(
|
|||||||
assert.Equal(t, "true", state,
|
assert.Equal(t, "true", state,
|
||||||
"clicking the caret does not expand the event at once")
|
"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(
|
require.NoError(t, chromedp.Run(ctx, chromedp.Evaluate(
|
||||||
`window.scrollTo(0, document.body.scrollHeight); window.scrollY`,
|
`window.scrollTo(0, document.body.scrollHeight); window.scrollY`,
|
||||||
&scrolled,
|
&scrolled,
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
@@ -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>
|
||||||
|
|||||||
@@ -255,6 +255,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>
|
||||||
|
|||||||
@@ -32,7 +32,7 @@
|
|||||||
<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>
|
||||||
@@ -67,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>
|
||||||
@@ -76,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