From d66292df43284e030c2876de188c942f326964ab Mon Sep 17 00:00:00 2001 From: sneak Date: Fri, 2 Oct 2026 22:04:53 +0000 Subject: [PATCH] Show a target paused by its circuit breaker (closes #385) While an http or slack target's circuit breaker is open, the target's row on the webhook page says its deliveries are paused until the cooldown ends, in UTC and from now. Each of its retrying deliveries shows as waiting in the event log and on the event's page, until the later of the cooldown's end and the end of its own backoff. While the breaker is half-open, the row says deliveries are held while one delivery tests the target, with no time, and deliveries keep their plain status. The engine gains one read, StateAndCooldown(targetID), taking a breaker's state and remaining cooldown under one lock; the handlers reach it through a one-method interface wired like Archives. Model: opus-5-5 --- README.md | 10 ++ cmd/webhooker/main.go | 4 + internal/delivery/circuit_breaker.go | 14 ++ internal/delivery/engine.go | 41 ++++- internal/delivery/engine_test.go | 56 ++++++ internal/delivery/export_test.go | 8 + internal/delivery/target.go | 1 + internal/delivery/target_http.go | 10 +- internal/handlers/delivery_result_view.go | 7 +- internal/handlers/handlers.go | 29 +-- internal/handlers/handlers_test.go | 41 +++++ internal/handlers/source_management.go | 17 +- internal/handlers/target_list.go | 67 +++++++ internal/handlers/target_paused_test.go | 205 ++++++++++++++++++++++ internal/resetpw/resetpw_test.go | 12 ++ internal/server/routes_test.go | 14 ++ templates/event_detail.html | 2 +- templates/source_detail.html | 6 + templates/source_logs.html | 4 +- 19 files changed, 520 insertions(+), 28 deletions(-) create mode 100644 internal/handlers/target_paused_test.go diff --git a/README.md b/README.md index bc13f49..7fe2e5a 100644 --- a/README.md +++ b/README.md @@ -2367,6 +2367,16 @@ just delayed until the target is healthy again. A delivery already in `retrying` keeps that status without another database write each time 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 time it will be tried next: +the later of the cooldown's end and the end of its own backoff after its +last attempt. 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` serves one Prometheus registry behind basic auth (see diff --git a/cmd/webhooker/main.go b/cmd/webhooker/main.go index 966df32..9269e99 100644 --- a/cmd/webhooker/main.go +++ b/cmd/webhooker/main.go @@ -212,6 +212,10 @@ func newApp() *fx.App { // or renaming a webhook or target reaches its archive // files. 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, ), fx.Invoke( diff --git a/internal/delivery/circuit_breaker.go b/internal/delivery/circuit_breaker.go index afb77d4..f0ae8eb 100644 --- a/internal/delivery/circuit_breaker.go +++ b/internal/delivery/circuit_breaker.go @@ -102,6 +102,20 @@ func (cb *CircuitBreaker) CooldownRemaining() time.Duration { 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 // the circuit breaker to closed state. func (cb *CircuitBreaker) RecordSuccess() { diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index 9fe588f..17572b9 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -143,6 +143,15 @@ type Archives interface { 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 // engine. type EngineParams struct { @@ -186,9 +195,11 @@ type Engine struct { // targets maps each target type to its implementation. targets map[database.TargetType]Target - // httpTarget is retained so tests can reach the HTTP - // target's shared client and circuit breakers. - httpTarget *httpTarget + // httpTarget and slackTarget are retained so StateAndCooldown + // can read their circuit breakers, and so tests can reach the + // HTTP target's shared client. + httpTarget *httpTarget + slackTarget *slackTarget // dbTarget is retained so the engine can reach the archive // writer registry for eviction, renames and the idle sweep. @@ -300,6 +311,30 @@ func (e *Engine) Rename( return e.dbTarget.rename(targetID, webhookName, targetName) } +// StateAndCooldown implements CircuitBreakers. It is +// CircuitBreaker.StateAndCooldown for the target's breaker. While the +// breaker is open, the pages show the target's deliveries as paused +// until its cooldown ends; while it is half-open, they show them as +// held, with no time, while one delivery tests whether the target has +// recovered. A target with no breaker reads as closed, 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 // retry channel after delay. It implements the Scheduler // interface the targets use to own their durable retries. diff --git a/internal/delivery/engine_test.go b/internal/delivery/engine_test.go index c06a772..5f85ad8 100644 --- a/internal/delivery/engine_test.go +++ b/internal/delivery/engine_test.go @@ -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) { t.Parallel() diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index 40585aa..2a6a1e9 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -212,6 +212,14 @@ func (e *Engine) ExportSetCircuitBreaker( 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. func (e *Engine) ExportParseHTTPConfig( configJSON string, diff --git a/internal/delivery/target.go b/internal/delivery/target.go index 17c3bf5..5ccb7d4 100644 --- a/internal/delivery/target.go +++ b/internal/delivery/target.go @@ -105,6 +105,7 @@ func (e *Engine) initTargets(client *http.Client) { dbT := &databaseTarget{eng: e} e.httpTarget = httpT + e.slackTarget = slackT e.dbTarget = dbT e.targets = map[database.TargetType]Target{ diff --git a/internal/delivery/target_http.go b/internal/delivery/target_http.go index 1f699f8..896bd92 100644 --- a/internal/delivery/target_http.go +++ b/internal/delivery/target_http.go @@ -234,7 +234,7 @@ func (c *httpCore) handleRetry( database.DeliveryStatusRetrying, ) - backoff := calcBackoff(attemptNum) + backoff := Backoff(attemptNum) retryTask := *task retryTask.AttemptNum = attemptNum + 1 @@ -301,7 +301,7 @@ func (c *httpCore) remainingBackoff( return 0 } - backoff := calcBackoff(attemptNum) + backoff := Backoff(attemptNum) elapsed := time.Since(lastResult.CreatedAt) remaining := backoff - elapsed @@ -326,12 +326,14 @@ func (c *httpCore) backoffElapsed( return true } - backoff := calcBackoff(attemptNum) + backoff := Backoff(attemptNum) 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 = min(shift, maxBackoffShift) diff --git a/internal/handlers/delivery_result_view.go b/internal/handlers/delivery_result_view.go index e3bb27b..4b0e51f 100644 --- a/internal/handlers/delivery_result_view.go +++ b/internal/handlers/delivery_result_view.go @@ -1,6 +1,8 @@ package handlers import ( + "time" + "sneak.berlin/go/webhooker/internal/delivery" ) @@ -24,8 +26,8 @@ const maxRenderedResponseBytes = 4096 // bytes rather than characters, and they make SQLite do the // cut, so an oversized stored response never becomes a Go // string at all. -const deliveryResultColumns = "delivery_id, attempt_num, success, " + - "status_code, error, duration, " + +const deliveryResultColumns = "delivery_id, attempt_num, created_at, " + + "success, status_code, error, duration, " + "substr(cast(response_body as blob), 1, ?) AS response_body, " + "length(cast(response_body as blob)) AS response_bytes" @@ -100,6 +102,7 @@ func (v DeliveryResultView) HasStatusCode() bool { type deliveryResultRow struct { DeliveryID string AttemptNum int + CreatedAt time.Time Success bool StatusCode int Error string diff --git a/internal/handlers/handlers.go b/internal/handlers/handlers.go index 9decb63..88734a3 100644 --- a/internal/handlers/handlers.go +++ b/internal/handlers/handlers.go @@ -57,19 +57,20 @@ var errVerificationBusy = errors.New( type HandlersParams struct { fx.In - Logger *logger.Logger - Globals *globals.Globals - Config *config.Config - Database *database.Database - WebhookDBMgr *database.WebhookDBManager - Healthcheck *healthcheck.Healthcheck - Session *session.Session - Middleware *middleware.Middleware - Notifier delivery.Notifier - Archives delivery.Archives - SSRFGuard *delivery.Guard - Metrics *metrics.Set - Registry *prometheus.Registry + Logger *logger.Logger + Globals *globals.Globals + Config *config.Config + Database *database.Database + WebhookDBMgr *database.WebhookDBManager + Healthcheck *healthcheck.Healthcheck + Session *session.Session + Middleware *middleware.Middleware + Notifier delivery.Notifier + Archives delivery.Archives + CircuitBreakers delivery.CircuitBreakers + SSRFGuard *delivery.Guard + Metrics *metrics.Set + Registry *prometheus.Registry } // Handlers provides HTTP handler methods for all application @@ -84,6 +85,7 @@ type Handlers struct { mw *middleware.Middleware notifier delivery.Notifier archives delivery.Archives + breakers delivery.CircuitBreakers mtr *metrics.Set templates map[string]*template.Template @@ -145,6 +147,7 @@ func New( s.mw = params.Middleware s.notifier = params.Notifier s.archives = params.Archives + s.breakers = params.CircuitBreakers s.mtr = params.Metrics s.ssrf = params.SSRFGuard diff --git a/internal/handlers/handlers_test.go b/internal/handlers/handlers_test.go index bbf4f11..6af3d28 100644 --- a/internal/handlers/handlers_test.go +++ b/internal/handlers/handlers_test.go @@ -9,6 +9,7 @@ import ( "net/http/httptest" "sync" "testing" + "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -181,6 +182,40 @@ func (r *recordingArchives) Renames() []archiveRename { 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 // 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 @@ -231,6 +266,12 @@ func newTestAppWithConfig( func(r *recordingArchives) delivery.Archives { return r }, + func() *testCircuitBreakers { + return &testCircuitBreakers{} + }, + func(b *testCircuitBreakers) delivery.CircuitBreakers { + return b + }, metrics.NewRegistry, metrics.New, middleware.New, diff --git a/internal/handlers/source_management.go b/internal/handlers/source_management.go index 6c2af2e..01dcb50 100644 --- a/internal/handlers/source_management.go +++ b/internal/handlers/source_management.go @@ -109,6 +109,10 @@ type DeliveryView struct { // the middle of Results. The page must show it, or the // bound would hide history rather than fold it. 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 @@ -1290,7 +1294,7 @@ func (h *Handlers) eventLogViews( } for i := range rows { - result[i].Deliveries = newDeliveryViews( + result[i].Deliveries = h.newDeliveryViews( eventDeliveries[i], targetMap, attempts, ) result[i].ResubmitCount = resubmits[rows[i].ID] @@ -1414,8 +1418,9 @@ func (h *Handlers) loadDeliveryResults( // newDeliveryViews projects deliveries for rendering, // resolving each one's target to its display-safe view and -// each one's attempts through that target's redactor. -func newDeliveryViews( +// each one's attempts through that target's redactor. A +// retrying delivery also reads its target's circuit breaker. +func (h *Handlers) newDeliveryViews( deliveries []database.Delivery, targetMap map[string]eventLogTarget, attempts map[string][]deliveryResultRow, @@ -1438,6 +1443,12 @@ func newDeliveryViews( AttemptCount: len(rows), AttemptsOmitted: omitted, } + + if deliveries[i].Status == database.DeliveryStatusRetrying { + views[i].Paused = h.deliveryPausedView( + deliveries[i].TargetID, rows, + ) + } } return views diff --git a/internal/handlers/target_list.go b/internal/handlers/target_list.go index 035a575..3c8ac55 100644 --- a/internal/handlers/target_list.go +++ b/internal/handlers/target_list.go @@ -24,6 +24,71 @@ type TargetRowView struct { // Archive is a database target's archive file, and nil for a target // of any other type. 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 when they resume, in UTC, and +// Relative how long that is from now. 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 when the delivery will be tried +// next: the later of the cooldown's end and the end of the delivery's +// own backoff after its last attempt. 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 that resume at until. +func newPausedView(until time.Time) *PausedView { + return &PausedView{ + Until: until.UTC().Format(time.TimeOnly) + " UTC", + Relative: humanize.Time(until), + } } // TargetDeliveries is how many of a target's deliveries became @@ -84,6 +149,8 @@ func (h *Handlers) targetRows( if targets[i].Type == database.TargetTypeDatabase { rows[i].Archive = h.archiveFileView(webhook, &targets[i]) } + + rows[i].Paused = h.pausedView(targets[i].ID) } return rows diff --git a/internal/handlers/target_paused_test.go b/internal/handlers/target_paused_test.go new file mode 100644 index 0000000..adea510 --- /dev/null +++ b/internal/handlers/target_paused_test.go @@ -0,0 +1,205 @@ +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" +) + +// resumesAt is how the pages write when a paused target's deliveries +// resume at the end of its breaker's cooldown: the time in UTC, then +// how long that is from now. +const resumesAt = `\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, until the later of +// the cooldown's end and the end of its own backoff. 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 13th attempt failed a minute ago, so its own + // backoff ends over an hour from now, long after the cooldown. + 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, 13, failedAt) + + backoffEnds := failedAt.Add(delivery.Backoff(13)).UTC(). + Format(time.TimeOnly) + " UTC (1 hour from now)" + + const waiting = "waiting: target paused after repeated failures, " + + "resumes " + + 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 "+resumesAt, + 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+resumesAt, 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+resumesAt, 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") + 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() +} diff --git a/internal/resetpw/resetpw_test.go b/internal/resetpw/resetpw_test.go index 922a8ac..99ba5ed 100644 --- a/internal/resetpw/resetpw_test.go +++ b/internal/resetpw/resetpw_test.go @@ -11,6 +11,7 @@ import ( "path/filepath" "strings" "testing" + "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -142,6 +143,14 @@ func (n *noopArchives) Rename(_, _, _ string) error { 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, // the middleware that bounds password verification, the session store // and the database, exactly as internal/handlers builds them. @@ -174,6 +183,9 @@ func newServerApp( session.New, func() delivery.Notifier { return &noopNotifier{} }, func() delivery.Archives { return &noopArchives{} }, + func() delivery.CircuitBreakers { + return &noopCircuitBreakers{} + }, metrics.NewRegistry, metrics.New, middleware.New, diff --git a/internal/server/routes_test.go b/internal/server/routes_test.go index da41aa2..dac68de 100644 --- a/internal/server/routes_test.go +++ b/internal/server/routes_test.go @@ -62,6 +62,17 @@ func (e *noopArchives) Rename(_, _, _ string) error { 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 // tests need to seed users and forge sessions. type testEnv struct { @@ -126,6 +137,9 @@ func newTestEnvWithConfig( session.New, func() delivery.Notifier { return &noopNotifier{} }, func() delivery.Archives { return &noopArchives{} }, + func() delivery.CircuitBreakers { + return &noopCircuitBreakers{} + }, metrics.NewRegistry, metrics.New, middleware.New, diff --git a/templates/event_detail.html b/templates/event_detail.html index aec0280..03a44eb 100644 --- a/templates/event_detail.html +++ b/templates/event_detail.html @@ -69,7 +69,7 @@
{{.Target.DisplayName}} - {{.Status}} + {{with .Paused}}waiting: target paused after repeated failures, resumes {{.Until}} ({{.Relative}}){{else}}{{.Status}}{{end}} {{.AttemptCount}} attempt{{if ne .AttemptCount 1}}s{{end}}
diff --git a/templates/source_detail.html b/templates/source_detail.html index bdb181b..5ecc2ce 100644 --- a/templates/source_detail.html +++ b/templates/source_detail.html @@ -255,6 +255,12 @@ + {{with .Paused}} +
+ Deliveries Paused: + {{if .Until}}after repeated failures, until {{.Until}} ({{.Relative}}){{else}}held while one delivery tests whether the target has recovered{{end}} +
+ {{end}} {{range .Config}}
{{.Label}}: diff --git a/templates/source_logs.html b/templates/source_logs.html index d900661..d6a58a0 100644 --- a/templates/source_logs.html +++ b/templates/source_logs.html @@ -31,7 +31,7 @@ {{range .Deliveries}} - {{.Target.DisplayName}}: {{.Status}} + {{.Target.DisplayName}}: {{if .Paused}}waiting{{else}}{{.Status}}{{end}} {{end}} {{.CreatedAt.Format "2006-01-02 15:04:05"}} @@ -65,7 +65,7 @@