Add a statistics pane to the webhook page (closes #368)
check / check (push) Successful in 4m10s
check / check (push) Successful in 4m10s
Each webhook's event database keeps running totals: one row for its events, with when the newest arrived, and one row per target for its deliveries, delivered and failed, each with what retention removed. Every write to them shares the transaction of the rows it counts, and a delivery already delivered or failed is not settled again. Deliveries get a finished_at column; it and target_id end the status index, so each target's deliveries finished in a window come from one index-range query grouped by target. Retention deletes 1000 expired events per transaction, pausing 200 ms between them so other writers get in, and stops between them on shutdown. The pane is its own template, its figures in tables. The schema changes in place with nothing back-filled, so an existing database must be recreated. Model: opus-5-5
This commit is contained in:
@@ -0,0 +1,177 @@
|
||||
package delivery_test
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gorm.io/gorm"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
)
|
||||
|
||||
// targetTotals reads one target's totals from a webhook database, all
|
||||
// zero when it has no row.
|
||||
func targetTotals(
|
||||
t *testing.T, db *gorm.DB, targetID string,
|
||||
) database.TargetTotals {
|
||||
t.Helper()
|
||||
|
||||
var rows []database.TargetTotals
|
||||
|
||||
require.NoError(t, db.Where("target_id = ?", targetID).
|
||||
Find(&rows).Error)
|
||||
|
||||
if len(rows) == 0 {
|
||||
return database.TargetTotals{TargetID: targetID}
|
||||
}
|
||||
|
||||
return rows[0]
|
||||
}
|
||||
|
||||
// TestUpdateDeliveryStatus_FinishTimeAndTargetTotals pins what a status
|
||||
// write records for the webhook page's statistics: the time a delivery
|
||||
// finished, set only when it becomes delivered or failed, and one more
|
||||
// on its target's delivered or failed total.
|
||||
func TestUpdateDeliveryStatus_FinishTimeAndTargetTotals(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
tests := []struct {
|
||||
status database.DeliveryStatus
|
||||
finished bool
|
||||
delivered int64
|
||||
failed int64
|
||||
}{
|
||||
{database.DeliveryStatusRetrying, false, 0, 0},
|
||||
{database.DeliveryStatusDelivered, true, 1, 0},
|
||||
{database.DeliveryStatusFailed, true, 0, 1},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(string(tt.status), func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db := testWebhookDB(t)
|
||||
e := testEngine(t, 1)
|
||||
event := seedEvent(t, db, `{}`)
|
||||
targetID := uuid.New().String()
|
||||
d := seedDelivery(
|
||||
t, db, event.ID, targetID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
before := time.Now()
|
||||
|
||||
require.NoError(t, e.ExportUpdateDeliveryStatus(
|
||||
db, &d, tt.status,
|
||||
))
|
||||
|
||||
var stored database.Delivery
|
||||
|
||||
require.NoError(t, db.First(&stored, "id = ?", d.ID).Error)
|
||||
assert.Equal(t, tt.status, stored.Status)
|
||||
|
||||
if tt.finished {
|
||||
require.NotNil(t, stored.FinishedAt)
|
||||
assert.False(t, stored.FinishedAt.Before(before))
|
||||
} else {
|
||||
assert.Nil(t, stored.FinishedAt)
|
||||
}
|
||||
|
||||
assert.Equal(t, database.TargetTotals{
|
||||
TargetID: targetID,
|
||||
Delivered: tt.delivered,
|
||||
Failed: tt.failed,
|
||||
}, targetTotals(t, db, targetID))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestUpdateDeliveryStatus_DeletedDeliveryIsNotCounted covers a
|
||||
// delivery retention deleted while the engine still held it. Failing
|
||||
// it afterwards writes no row, so it adds no failure either: retention
|
||||
// has already counted what it removed.
|
||||
func TestUpdateDeliveryStatus_DeletedDeliveryIsNotCounted(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db := testWebhookDB(t)
|
||||
e := testEngine(t, 1)
|
||||
event := seedEvent(t, db, `{}`)
|
||||
targetID := uuid.New().String()
|
||||
d := seedDelivery(
|
||||
t, db, event.ID, targetID,
|
||||
database.DeliveryStatusRetrying,
|
||||
)
|
||||
|
||||
require.NoError(t, db.Unscoped().
|
||||
Delete(&database.Delivery{}, "id = ?", d.ID).Error)
|
||||
|
||||
require.NoError(t, e.ExportUpdateDeliveryStatus(
|
||||
db, &d, database.DeliveryStatusFailed,
|
||||
))
|
||||
|
||||
assert.Equal(t, database.TargetTotals{TargetID: targetID},
|
||||
targetTotals(t, db, targetID))
|
||||
}
|
||||
|
||||
// TestUpdateDeliveryStatus_FinishedDeliveryIsNotSettledAgain covers a
|
||||
// delivery settled a second time, as recovery can do when a worker has
|
||||
// settled it since recovery read it. Neither status writes over the
|
||||
// first, and the totals do not move.
|
||||
func TestUpdateDeliveryStatus_FinishedDeliveryIsNotSettledAgain(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
finished := []database.DeliveryStatus{
|
||||
database.DeliveryStatusDelivered,
|
||||
database.DeliveryStatusFailed,
|
||||
}
|
||||
|
||||
for _, first := range finished {
|
||||
t.Run(string(first), func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db := testWebhookDB(t)
|
||||
e := testEngine(t, 1)
|
||||
event := seedEvent(t, db, `{}`)
|
||||
targetID := uuid.New().String()
|
||||
d := seedDelivery(
|
||||
t, db, event.ID, targetID,
|
||||
database.DeliveryStatusRetrying,
|
||||
)
|
||||
|
||||
// The delivery as recovery read it, before the worker
|
||||
// settled it.
|
||||
readBefore := d
|
||||
|
||||
require.NoError(t, e.ExportUpdateDeliveryStatus(
|
||||
db, &d, first,
|
||||
))
|
||||
|
||||
var settled database.Delivery
|
||||
|
||||
require.NoError(t, db.First(&settled, "id = ?", d.ID).Error)
|
||||
require.NotNil(t, settled.FinishedAt)
|
||||
|
||||
totals := targetTotals(t, db, targetID)
|
||||
|
||||
for _, again := range finished {
|
||||
stale := readBefore
|
||||
|
||||
require.NoError(t, e.ExportUpdateDeliveryStatus(
|
||||
db, &stale, again,
|
||||
))
|
||||
}
|
||||
|
||||
var stored database.Delivery
|
||||
|
||||
require.NoError(t, db.First(&stored, "id = ?", d.ID).Error)
|
||||
assert.Equal(t, first, stored.Status)
|
||||
require.NotNil(t, stored.FinishedAt)
|
||||
assert.True(t, settled.FinishedAt.Equal(*stored.FinishedAt))
|
||||
assert.Equal(t, totals, targetTotals(t, db, targetID))
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -531,6 +531,11 @@ func (e *Engine) processRetryTask(
|
||||
return
|
||||
}
|
||||
|
||||
// Set before anything below can fail the delivery: the failure is
|
||||
// added to this target's totals.
|
||||
d.EventID = task.EventID
|
||||
d.TargetID = task.TargetID
|
||||
|
||||
if d.Status != database.DeliveryStatusRetrying {
|
||||
e.log.Debug(
|
||||
"skipping retry for delivery "+
|
||||
@@ -562,8 +567,6 @@ func (e *Engine) processRetryTask(
|
||||
}
|
||||
|
||||
target := buildTargetFromTask(task)
|
||||
d.EventID = task.EventID
|
||||
d.TargetID = task.TargetID
|
||||
d.Event = event
|
||||
d.Target = target
|
||||
|
||||
@@ -1554,8 +1557,9 @@ func (e *Engine) updateDeliveryStatus(
|
||||
targetType database.TargetType,
|
||||
status database.DeliveryStatus,
|
||||
) error {
|
||||
err := webhookDB.Model(d).
|
||||
Update("status", status).Error
|
||||
err := webhookDB.Transaction(func(tx *gorm.DB) error {
|
||||
return writeDeliveryStatus(tx, d, status)
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf(
|
||||
"updating delivery %s to status %s: %w",
|
||||
@@ -1574,6 +1578,43 @@ func (e *Engine) updateDeliveryStatus(
|
||||
return nil
|
||||
}
|
||||
|
||||
// writeDeliveryStatus writes a delivery's new status. A delivery that
|
||||
// becomes delivered or failed also gets the time it finished, and is
|
||||
// added to its target's delivered or failed total. That write changes
|
||||
// only a delivery not yet delivered or failed, and the total moves
|
||||
// only when it changed a row: retention may have deleted the delivery
|
||||
// while the engine was working on it, and a recovery path may settle
|
||||
// a delivery that a worker has already settled.
|
||||
func writeDeliveryStatus(
|
||||
tx *gorm.DB,
|
||||
d *database.Delivery,
|
||||
status database.DeliveryStatus,
|
||||
) error {
|
||||
if !status.Terminal() {
|
||||
return tx.Model(d).Update("status", status).Error
|
||||
}
|
||||
|
||||
res := tx.Model(d).
|
||||
Where("status NOT IN ?", []database.DeliveryStatus{
|
||||
database.DeliveryStatusDelivered,
|
||||
database.DeliveryStatusFailed,
|
||||
}).
|
||||
Updates(map[string]any{
|
||||
"status": status,
|
||||
"finished_at": time.Now(),
|
||||
})
|
||||
if res.Error != nil || res.RowsAffected == 0 {
|
||||
return res.Error
|
||||
}
|
||||
|
||||
add := database.TargetTotals{TargetID: d.TargetID, Delivered: 1}
|
||||
if status == database.DeliveryStatusFailed {
|
||||
add = database.TargetTotals{TargetID: d.TargetID, Failed: 1}
|
||||
}
|
||||
|
||||
return database.AddTargetTotals(tx, add)
|
||||
}
|
||||
|
||||
// settleStatus moves a delivery to its outcome status and reports a
|
||||
// failed write through bookkeepingFailed, which leaves the row
|
||||
// recoverable. It exists so the target call sites read as one
|
||||
|
||||
@@ -57,7 +57,10 @@ func testWebhookDB(t *testing.T) *gorm.DB {
|
||||
&database.Event{},
|
||||
&database.Delivery{},
|
||||
&database.DeliveryResult{},
|
||||
&database.EventTotals{},
|
||||
&database.TargetTotals{},
|
||||
))
|
||||
require.NoError(t, db.Create(&database.EventTotals{}).Error)
|
||||
|
||||
return db
|
||||
}
|
||||
|
||||
@@ -150,6 +150,16 @@ func (e *Engine) ExportDeliverSlack(
|
||||
)
|
||||
}
|
||||
|
||||
// ExportUpdateDeliveryStatus exposes updateDeliveryStatus. It passes no
|
||||
// target type, so no metric moves.
|
||||
func (e *Engine) ExportUpdateDeliveryStatus(
|
||||
webhookDB *gorm.DB,
|
||||
d *database.Delivery,
|
||||
status database.DeliveryStatus,
|
||||
) error {
|
||||
return e.updateDeliveryStatus(webhookDB, d, "", status)
|
||||
}
|
||||
|
||||
// ExportProcessNewTask exposes processNewTask.
|
||||
func (e *Engine) ExportProcessNewTask(
|
||||
ctx context.Context, task *Task,
|
||||
|
||||
@@ -418,6 +418,38 @@ func TestProcessRetryTask_TargetDeleted_MakesNoAttempt(
|
||||
assert.Zero(t, s.Engine.ExportInflightHeld())
|
||||
}
|
||||
|
||||
// TestProcessRetryTask_TargetDeleted_CountsFailureOnTarget verifies
|
||||
// that the failure of a retry abandoned because its target is gone is
|
||||
// added to that target's own totals, not to a row with no target.
|
||||
func TestProcessRetryTask_TargetDeleted_CountsFailureOnTarget(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
s := newISetup(t)
|
||||
|
||||
var hits atomic.Int64
|
||||
|
||||
task, targetID := tRetryChainSetup(
|
||||
t, s, "gone-counted", &hits,
|
||||
)
|
||||
|
||||
require.NoError(t, s.MainDB.Delete(
|
||||
&database.Target{}, "id = ?", targetID,
|
||||
).Error)
|
||||
|
||||
s.Engine.ExportProcessRetryTask(
|
||||
context.Background(), &task,
|
||||
)
|
||||
|
||||
var rows []database.TargetTotals
|
||||
|
||||
require.NoError(t, s.WebhookDB.Find(&rows).Error)
|
||||
assert.Equal(t, []database.TargetTotals{
|
||||
{TargetID: targetID, Failed: 1},
|
||||
}, rows)
|
||||
}
|
||||
|
||||
// TestProcessRetryTask_TargetPresent_StillDelivers is the guard's
|
||||
// mutation check: a liveness check that refused every retry would pass
|
||||
// the test above and break every retry there is.
|
||||
|
||||
Reference in New Issue
Block a user