1 Commits

Author SHA1 Message Date
clawbot
0384de4a7b Terminally fail retrying deliveries with a non-retry target type (closes #82)
All checks were successful
check / check (push) Successful in 4m4s
Restart recovery and the 60s retry sweep both looked an orphaned
`retrying` delivery's target up in the registry and silently returned
when it did not implement `rescheduler`. If a target's type was edited
from a retry type (`http`/`slack`) to a fire-and-forget type
(`database`/`log`) or an unknown one while a delivery was still
retrying, that delivery stayed `retrying` forever.

Both sites now hand the delivery to one shared helper,
`failUnretryableRetry`, which records a `DeliveryResult` naming the
current target type as the reason and marks the delivery `failed`. It
logs at warn, not error: this is operator-caused state, not a system
fault.

Re-dispatching under the new type was rejected as it would perform a
delivery the operator never asked for; the event itself stays in the
per-webhook event database, so manual redelivery can recover it
deliberately.

Fire-and-forget targets never set status `retrying` under normal
operation, so this path stays unreachable for them in practice.
2026-08-09 05:51:51 +00:00
17 changed files with 289 additions and 2305 deletions

View File

@@ -531,36 +531,6 @@ older than the expiry are pruned each time the archive is (re)opened. An
archive write failure is never silent success: the delivery records a
failed attempt with the error and is marked failed.
Because reopens only happen on writes, an archive belonging to a webhook
that has stopped receiving events would never be pruned. A background
**archive sweeper** closes that gap: on the same interval as the event
retention reaper (`RETENTION_SWEEP_INTERVAL`) it prunes every archive
whose database target declares a positive expiry, whether or not the
webhook is still receiving traffic. The sweep never creates an archive —
a webhook whose archive file does not yet exist is skipped, not
initialised — it takes the same per-webhook lock the write path uses, so
it can never interleave with a write, and it leaves the archive closed
afterwards so the move-the-file-away workflow keeps working. Archives
with no expiry, or the expiry `never`, are not touched by the sweep at
all.
Note that a webhook has one archive file but may carry more than one
`database` target, each with its own `expiry`. The shortest expiry
configured on any of them therefore governs the whole archive, and the
sweep applies it whether or not the webhook is still receiving events.
Configure a single `database` target per webhook unless you intend that.
Deleting a webhook releases its archive: the delivery engine's cached
archive writer is dropped and its file handle closed, so nothing lingers
after the webhook is gone. The archive **file itself is deliberately
left on disk**. Unlike the event database — per-webhook working storage
that is hard-deleted with the webhook — an archive is long-term storage
an operator may still want to keep or move away for offline retention,
and destroying it as a side effect of deleting a webhook would be
unrecoverable. Removing `archive-{webhookID}.db` is the operator's call.
Deleting a webhook's last `database` target releases the writer the same
way, and for the same reason leaves the file alone.
The **Slack target type** sends webhook events as formatted messages to
any Slack-compatible incoming webhook URL (works with Slack, Mattermost,
and other compatible services). Each message includes event metadata
@@ -666,6 +636,18 @@ This means:
durable fallback that ensures no retry is permanently lost, even under
extreme backpressure.
**Changing a target's type does not migrate in-flight deliveries.** Only
`http` and `slack` targets own durable retries; `database` and `log`
targets are fire-and-forget and never produce a `retrying` delivery. If a
target's `type` is edited from a retrying type to a non-retrying (or
unknown) one while one of its deliveries is still `retrying`, both
recovery paths above terminally mark that delivery `failed` and record a
`DeliveryResult` naming the current target type as the reason, logging it
at warn level. The delivery is not re-dispatched under the new type — the
operator never asked for that delivery — and the event itself remains
stored in the per-webhook event database, so it can be redelivered
manually.
### Circuit Breaker (HTTP Targets with Retries)
HTTP targets with `max_retries` > 0 are protected by a **per-target circuit breaker** that

15
TODO.md
View File

@@ -10,11 +10,9 @@
# Status
pre-1.0. No git tags exist. main (4f5ecb1) is a working webhook proxy
pre-1.0. No git tags exist. main (afe88c6) is a working webhook proxy
with auth, CSRF/SSRF protections, login rate limiting, Slack target,
event retention (#63), the database archiving target (#43), the admin
password change flow (#65), policy compliance (#6), and pinned lint
tooling (#55). Note: TODO.md was
policy compliance (#6), and pinned lint tooling (#55). Note: TODO.md was
deliberately deleted from this repo in f9a9569 (2026-03-01, #6); its
content was folded into the README TODO section, which this draft
reconstructs as of 2026-07-06.
@@ -30,11 +28,10 @@ databases currently grow without bound.
# Completed Steps
- 2026-08-09 Archive writer lifecycle (#89): deleting a webhook (or its
last `database` target) evicts the cached archive writer and closes
its handle while deliberately leaving `archive-{webhookID}.db` on
disk, and a new `ArchiveSweeper` prunes idle archives on the existing
`RETENTION_SWEEP_INTERVAL` without ever creating an archive file
- 2026-08-09 Restart recovery and the 60s retry sweep terminally fail an
orphaned `retrying` delivery whose target type no longer supports
retries, recording a `DeliveryResult` with the reason instead of
leaving the delivery stuck forever (#82)
- 2026-08-07 Update golangci-lint to v2.12.2 (Docker image digest in
`Dockerfile`, release-archive sha256 pins in `script/bootstrap`),
adopt the canonical `.golangci.yml` (v2 `linters.settings` layout so

View File

@@ -40,15 +40,9 @@ func main() {
handlers.New,
middleware.New,
delivery.New,
delivery.NewArchiveSweeper,
// Wire *delivery.Engine as delivery.Notifier so the
// webhook handler can notify the engine of new deliveries.
func(e *delivery.Engine) delivery.Notifier { return e },
// Wire *delivery.Engine as delivery.WebhookEvictor so
// deleting a webhook releases its archive writer.
func(e *delivery.Engine) delivery.WebhookEvictor {
return e
},
server.New,
),
fx.Invoke(
@@ -56,7 +50,6 @@ func main() {
*server.Server,
*delivery.Engine,
*database.RetentionReaper,
*delivery.ArchiveSweeper,
) {
},
),

View File

@@ -1,229 +0,0 @@
package delivery
import (
"context"
"errors"
"log/slog"
"sync"
"time"
"go.uber.org/fx"
"sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/logger"
)
// ArchiveSweeperParams holds the fx dependencies for the
// ArchiveSweeper.
type ArchiveSweeperParams struct {
fx.In
Config *config.Config
Database *database.Database
Engine *Engine
Logger *logger.Logger
}
// ArchiveSweeper periodically prunes expired rows from
// per-webhook archive databases whose database target carries a
// positive expiry.
//
// Without it, pruning happens only when an archive is
// (re)opened, and archives are only ever reopened by writes: an
// archive belonging to a webhook that has stopped receiving
// events would keep its expired rows forever. The sweep closes
// that gap without changing anything for archives whose expiry
// is unset or "never".
//
// It reuses Config.RetentionSweepInterval rather than
// introducing a second interval: this is a retention sweep with
// the same semantics as the event retention reaper.
type ArchiveSweeper struct {
db *database.Database
eng *Engine
log *slog.Logger
interval time.Duration
cancel context.CancelFunc
wg sync.WaitGroup
}
// NewArchiveSweeper creates the archive sweeper and registers
// its fx lifecycle hooks. The background sweep loop starts on
// OnStart and stops cleanly on OnStop via context cancellation.
func NewArchiveSweeper(
lc fx.Lifecycle,
params ArchiveSweeperParams,
) *ArchiveSweeper {
s := &ArchiveSweeper{
db: params.Database,
eng: params.Engine,
log: params.Logger.Get(),
interval: params.Config.RetentionSweepInterval,
}
s.registerHooks(lc)
return s
}
// registerHooks wires the sweeper's start and stop into the fx
// lifecycle. Both hook contexts are deliberately ignored: see
// start for why the background loop must not inherit the start
// hook's context, and stop for why shutdown blocks on the loop
// rather than on the stop hook's deadline.
func (s *ArchiveSweeper) registerHooks(lc fx.Lifecycle) {
lc.Append(fx.Hook{
//nolint:contextcheck // Not passing the hook context is
// the point: see start.
OnStart: func(_ context.Context) error {
s.start()
return nil
},
OnStop: func(_ context.Context) error {
s.stop()
return nil
},
})
}
// start launches the background sweep loop.
//
// The loop's context is derived from context.Background(), NOT
// from the fx OnStart hook context. The hook context carries
// fx's start timeout (15s by default), so a loop derived from it
// is cancelled 15 seconds after the application starts — long
// before the first tick under the default one-hour sweep
// interval, leaving a sweeper that never sweeps. A long-lived
// goroutine must outlive the startup phase, so its lifetime is
// bounded by OnStop instead: stop cancels this context and waits
// on the WaitGroup.
func (s *ArchiveSweeper) start() {
ctx, cancel := context.WithCancel(context.Background())
s.cancel = cancel
s.wg.Add(1)
go s.run(ctx)
s.log.Info(
"archive sweeper started",
"interval", s.interval.String(),
)
}
func (s *ArchiveSweeper) stop() {
s.log.Info("archive sweeper stopping")
if s.cancel != nil {
s.cancel()
}
s.wg.Wait()
s.log.Info("archive sweeper stopped")
}
func (s *ArchiveSweeper) run(ctx context.Context) {
defer s.wg.Done()
ticker := time.NewTicker(s.interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
s.sweep(ctx)
}
}
}
// sweep prunes every archive whose database target declares a
// positive expiry. Targets belonging to a deleted webhook are
// soft-deleted along with it, so GORM's default scope already
// excludes them.
//
// A failure for one webhook is logged and the sweep continues,
// matching how the write path already treats a prune error as
// non-fatal.
func (s *ArchiveSweeper) sweep(ctx context.Context) {
var targets []database.Target
err := s.db.DB().
Model(&database.Target{}).
Where("type = ?", database.TargetTypeDatabase).
Find(&targets).Error
if err != nil {
s.log.Error(
"archive sweep: failed to list database targets",
"error", err,
)
return
}
for i := range targets {
select {
case <-ctx.Done():
return
default:
}
s.sweepTarget(&targets[i])
}
}
// sweepTarget prunes the archive of a single database target.
// A missing, empty, or "never" expiry parses as a zero duration
// and is skipped entirely, so those archives keep exactly the
// behaviour they had before the sweep existed.
func (s *ArchiveSweeper) sweepTarget(target *database.Target) {
expiry, err := parseArchiveExpiry(target.Config)
if err != nil {
s.log.Error(
"archive sweep: invalid database target config",
"webhook_id", target.WebhookID,
"target_id", target.ID,
"error", err,
)
return
}
if expiry <= 0 {
return
}
if s.eng == nil || s.eng.dbTarget == nil {
return
}
err = s.eng.dbTarget.sweepWebhook(target.WebhookID, expiry)
if err == nil {
return
}
// A writer evicted underneath the sweep means the operator
// deleted the webhook (or its last database target) while the
// sweep was walking the target list. That is an ordinary
// interleaving, not a failure, so it must not produce an
// error line.
if errors.Is(err, errArchiveWriterEvicted) {
s.log.Debug(
"archive sweep: writer evicted mid-sweep",
"webhook_id", target.WebhookID,
"target_id", target.ID,
)
return
}
s.log.Error(
"archive sweep: failed to prune archive",
"webhook_id", target.WebhookID,
"target_id", target.ID,
"error", err,
)
}

View File

@@ -1,713 +0,0 @@
package delivery_test
import (
"context"
"database/sql"
"fmt"
"net/http"
"os"
"path/filepath"
"sync"
"testing"
"time"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.uber.org/fx"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"gorm.io/gorm/clause"
_ "modernc.org/sqlite" // Pure Go SQLite driver.
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
)
const (
// sweepRowOld and sweepRowNew are the event ids
// seedArchiveRows assigns to the first and second seeded
// rows.
sweepRowOld = "ev-0"
sweepRowNew = "ev-1"
// sweepConcurrentWrites is how many deliveries the
// concurrent write-plus-sweep test races against the sweep.
sweepConcurrentWrites = 20
)
// sweeperEnv bundles the pieces an archive sweep test drives:
// a main configuration database holding webhooks and targets, a
// delivery engine owning the archive writer registry, and the
// data directory the archive files live in.
type sweeperEnv struct {
sweeper *delivery.ArchiveSweeper
eng *delivery.Engine
mainDB *database.Database
dataDir string
}
func setupSweeperTest(t *testing.T) *sweeperEnv {
t.Helper()
dataDir := t.TempDir()
log := archiveTestLogger()
sqlDB, err := sql.Open(
"sqlite",
fmt.Sprintf(
"file:%s?mode=rwc",
filepath.Join(dataDir, "main.db"),
),
)
require.NoError(t, err)
t.Cleanup(func() { _ = sqlDB.Close() })
gdb, err := gorm.Open(
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
)
require.NoError(t, err)
mainDB := database.NewTestDatabase(gdb)
require.NoError(t, mainDB.Migrate())
eng := delivery.NewTestEngineWithDB(
mainDB,
database.NewTestWebhookDBManager(dataDir),
log,
&http.Client{Timeout: 5 * time.Second},
1,
)
return &sweeperEnv{
sweeper: delivery.NewTestArchiveSweeper(
mainDB, eng, log,
),
eng: eng,
mainDB: mainDB,
dataDir: dataDir,
}
}
// archivePath returns where the engine keeps a webhook's
// archive file.
func (env *sweeperEnv) archivePath(webhookID string) string {
return filepath.Join(
env.dataDir, fmt.Sprintf("archive-%s.db", webhookID),
)
}
// seedDatabaseTarget creates a webhook with one database target
// carrying the given target config JSON, and returns the
// webhook id.
func (env *sweeperEnv) seedDatabaseTarget(
t *testing.T, configJSON string,
) string {
t.Helper()
wh := &database.Webhook{
UserID: uuid.New().String(),
Name: "sweep-test",
}
require.NoError(
t,
env.mainDB.DB().
Omit(clause.Associations).
Create(wh).Error,
)
tgt := &database.Target{
WebhookID: wh.ID,
Name: "archive",
Type: database.TargetTypeDatabase,
Active: true,
Config: configJSON,
}
require.NoError(
t,
env.mainDB.DB().
Omit(clause.Associations).
Create(tgt).Error,
)
return wh.ID
}
// seedArchiveRows creates the archive file for a webhook and
// inserts one row per supplied archived-at timestamp, returning
// the archive path. The handle is closed before returning, so
// the archive is idle exactly as it would be with no traffic.
func (env *sweeperEnv) seedArchiveRows(
t *testing.T, webhookID string, archivedAt ...time.Time,
) string {
t.Helper()
path := env.archivePath(webhookID)
sqlDB, err := sql.Open(
"sqlite", fmt.Sprintf("file:%s?mode=rwc", path),
)
require.NoError(t, err)
gdb, err := gorm.Open(
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
)
require.NoError(t, err)
require.NoError(
t, gdb.AutoMigrate(&delivery.ExportArchivedEvent{}),
)
for i, at := range archivedAt {
row := delivery.ExportArchivedEvent{
EventID: fmt.Sprintf("ev-%d", i),
WebhookID: webhookID,
Method: http.MethodPost,
Body: `{"seeded":true}`,
ArchivedAt: at,
}
require.NoError(t, gdb.Create(&row).Error)
}
require.NoError(t, sqlDB.Close())
return path
}
// archivedEventIDs returns the event ids currently stored in an
// archive file, read through a separate read-only handle.
func archivedEventIDs(
t *testing.T, path string,
) []string {
t.Helper()
var rows []delivery.ExportArchivedEvent
rdb := openArchiveDBForRead(t, path)
require.NoError(t, rdb.Order("event_id").Find(&rows).Error)
ids := make([]string, 0, len(rows))
for i := range rows {
ids = append(ids, rows[i].EventID)
}
return ids
}
// countArchivedRows counts the rows in an archive file without
// asserting anything, so it is safe to poll from an
// assert.Eventually condition (which runs off the test
// goroutine, where testify assertions must not be used).
func countArchivedRows(path string) (int64, error) {
sqlDB, err := sql.Open(
"sqlite", fmt.Sprintf("file:%s?mode=ro", path),
)
if err != nil {
return 0, err
}
defer func() { _ = sqlDB.Close() }()
gdb, err := gorm.Open(
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
)
if err != nil {
return 0, err
}
var count int64
err = gdb.Model(&delivery.ExportArchivedEvent{}).
Count(&count).Error
if err != nil {
return 0, err
}
return count, nil
}
// captureLifecycle is a minimal fx.Lifecycle that records the
// hooks a component registers, so a test can invoke the real
// OnStart/OnStop functions with a context of its choosing.
type captureLifecycle struct {
hooks []fx.Hook
}
func (l *captureLifecycle) Append(h fx.Hook) {
l.hooks = append(l.hooks, h)
}
// TestArchiveSweeper_LoopOutlivesStartHookContext is the
// regression test for a sweeper that never swept. fx calls
// OnStart with a context carrying the application's start
// timeout (15 seconds by default), so a background loop whose
// context is derived from it is cancelled 15 seconds into the
// process — three quarters of an hour before the first tick
// under the default one-hour sweep interval.
//
// The hook context here is already cancelled, which is the same
// defect taken to its limit: a loop that inherits it never runs
// a single tick, while a correctly rooted loop keeps sweeping
// for as long as the process lives. Handing the hook a plain
// context.Background() would assert nothing at all.
func TestArchiveSweeper_LoopOutlivesStartHookContext(
t *testing.T,
) {
t.Parallel()
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
now := time.Now()
path := env.seedArchiveRows(
t, webhookID,
now.Add(-48*time.Hour),
now.Add(-time.Minute),
)
env.sweeper.ExportSetInterval(10 * time.Millisecond)
// Drive the genuine fx hooks the application registers,
// rather than a test-only entry point.
lc := &captureLifecycle{}
env.sweeper.ExportRegisterHooks(lc)
require.Len(t, lc.hooks, 1)
hookCtx, cancel := context.WithCancel(context.Background())
cancel()
require.NoError(t, lc.hooks[0].OnStart(hookCtx))
t.Cleanup(func() {
_ = lc.hooks[0].OnStop(context.Background())
})
assert.Eventually(
t,
func() bool {
count, err := countArchivedRows(path)
return err == nil && count == 1
},
5*time.Second,
10*time.Millisecond,
"the sweep loop must keep running after the start "+
"hook's context is done; it pruned nothing, so it "+
"inherited the hook context and died",
)
}
// TestArchiveSweep_DoesNotResurrectEvictedWriter covers the
// interleaving where a sweep tick has already listed a webhook's
// target when the webhook is deleted and its writer evicted. The
// sweep must not put a writer back into the registry: nothing
// would ever evict it again, which is precisely the leak this
// change exists to close.
func TestArchiveSweep_DoesNotResurrectEvictedWriter(
t *testing.T,
) {
t.Parallel()
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
)
// Prime the registry the way a delivery would, then evict as
// the deletion path does. The target row is deliberately left
// in place: this is the tick that listed the webhook before
// the deletion committed.
_, err := env.eng.ExportEnsureArchiveWriter(webhookID)
require.NoError(t, err)
env.eng.EvictWebhook(webhookID)
require.False(t, env.eng.ExportHasArchiveWriter(webhookID))
env.sweeper.ExportSweep(context.Background())
assert.False(
t, env.eng.ExportHasArchiveWriter(webhookID),
"a sweep must never re-register a writer for a webhook "+
"whose registry entry has already been released",
)
}
// TestArchiveSweep_LeavesNoRegistryEntry states the same
// invariant in its general form: sweeping an archive whose
// webhook has no cached writer must not leave one behind, so the
// registry keeps holding only writers a delivery created and an
// eviction can reach.
func TestArchiveSweep_LeavesNoRegistryEntry(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.seedArchiveRows(
t, webhookID,
time.Now().Add(-48*time.Hour),
time.Now().Add(-time.Minute),
)
require.False(t, env.eng.ExportHasArchiveWriter(webhookID))
env.sweeper.ExportSweep(context.Background())
assert.Equal(
t, []string{sweepRowNew}, archivedEventIDs(t, path),
"the sweep must still prune an idle archive",
)
assert.False(
t, env.eng.ExportHasArchiveWriter(webhookID),
"the sweep must release the registry entry it created",
)
}
// TestArchiveSweep_KeepsWriterAdoptedByDelivery is the other
// half of that invariant: an entry the sweep created but a
// delivery then claimed belongs to the registry and must survive
// the sweep, or the delivery would be left holding a detached
// writer with an open handle that no eviction can reach.
func TestArchiveSweep_KeepsWriterAdoptedByDelivery(
t *testing.T,
) {
t.Parallel()
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
)
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"n":1}`)
event.WebhookID = webhookID
d := seedDatabaseTargetDelivery(
t, webhookDB, event, `{"expiry":"1h"}`,
)
env.sweeper.ExportSweep(context.Background())
require.False(t, env.eng.ExportHasArchiveWriter(webhookID))
env.eng.ExportDeliverDatabase(webhookDB, d)
assert.True(
t, env.eng.ExportHasArchiveWriter(webhookID),
"a delivery's writer must stay registered",
)
env.sweeper.ExportSweep(context.Background())
assert.True(
t, env.eng.ExportHasArchiveWriter(webhookID),
"a sweep must not drop a writer a delivery owns",
)
}
// TestArchiveSweep_ContinuesAfterPerWebhookFailure proves a
// failure for one webhook does not abort the sweep for the
// others: an unparseable expiry and an unreadable archive both
// have to be logged and stepped over.
func TestArchiveSweep_ContinuesAfterPerWebhookFailure(
t *testing.T,
) {
t.Parallel()
env := setupSweeperTest(t)
// Seeded first so the sweep reaches them before the healthy
// webhook: targets come back in insertion order.
badConfigID := env.seedDatabaseTarget(t, `{"expiry":"!!!"}`)
env.seedArchiveRows(
t, badConfigID, time.Now().Add(-48*time.Hour),
)
corruptID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
require.NoError(t, os.WriteFile(
env.archivePath(corruptID),
[]byte("this is not a sqlite database"),
0o600,
))
healthyID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
healthyPath := env.seedArchiveRows(
t, healthyID,
time.Now().Add(-48*time.Hour),
time.Now().Add(-time.Minute),
)
env.sweeper.ExportSweep(context.Background())
assert.Equal(
t, []string{sweepRowNew},
archivedEventIDs(t, healthyPath),
"a failure for an earlier webhook must not stop the "+
"sweep from pruning the ones after it",
)
}
// TestArchiveSweep_OpenExistingDoesNotCreateFile pins the second
// of the two no-create guards. The first is the stat in
// sweepWebhook; this one is the SQLite open mode, which is what
// protects the window between that stat and the open. Flipping
// the sweep's mode to create-if-missing makes this fail.
func TestArchiveSweep_OpenExistingDoesNotCreateFile(
t *testing.T,
) {
t.Parallel()
dir := t.TempDir()
path := filepath.Join(dir, "archive-absent.db")
w := delivery.NewExportArchiveWriter(
path, archiveTestLogger(), 0,
)
err := w.OpenExisting(time.Hour)
require.Error(
t, err,
"opening a missing archive without create permission "+
"must fail rather than conjure the file",
)
for _, suffix := range archiveFileSuffixes() {
assert.NoFileExists(t, path+suffix)
}
}
// TestArchiveSweep_PrunesIdleArchive is the core regression
// test for this issue: an archive that receives no further
// writes must still lose its expired rows. Before the sweeper
// existed, pruning only ever ran on a write-triggered reopen,
// so an idle archive kept expired rows forever.
func TestArchiveSweep_PrunesIdleArchive(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
now := time.Now()
path := env.seedArchiveRows(
t, webhookID,
now.Add(-48*time.Hour),
now.Add(-time.Minute),
)
require.Equal(
t, []string{sweepRowOld, sweepRowNew},
archivedEventIDs(t, path),
)
env.sweeper.ExportSweep(context.Background())
assert.Equal(
t, []string{sweepRowNew}, archivedEventIDs(t, path),
"the sweep should prune rows older than the expiry "+
"from an idle archive and keep the rest",
)
}
// TestArchiveSweep_LeavesArchiveClosed proves the sweep does
// not hold the archive open afterwards, so an operator can
// still move the file away for offline retention.
func TestArchiveSweep_LeavesArchiveClosed(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
)
env.sweeper.ExportSweep(context.Background())
assert.False(
t, env.eng.ExportArchiveHandleOpen(webhookID),
"an idle archive must end the sweep closed",
)
}
// TestArchiveSweep_NeverExpiryUntouched proves the sweep is a
// no-op for the default retention policy, so archives with no
// expiry (or the literal "never") behave exactly as before.
func TestArchiveSweep_NeverExpiryUntouched(t *testing.T) {
t.Parallel()
for _, configJSON := range []string{
`{"expiry":"never"}`,
`{"expiry":""}`,
"",
} {
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, configJSON)
path := env.seedArchiveRows(
t, webhookID,
time.Now().Add(-10000*time.Hour),
)
env.sweeper.ExportSweep(context.Background())
assert.Equal(
t, []string{sweepRowOld}, archivedEventIDs(t, path),
"config %q must keep rows forever", configJSON,
)
}
}
// TestArchiveSweep_DoesNotCreateArchiveFile proves the sweep
// never conjures an archive: a webhook with a database target
// that has never received an event must still have no archive
// file (nor SQLite sidecar) after a sweep.
func TestArchiveSweep_DoesNotCreateArchiveFile(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.archivePath(webhookID)
require.NoFileExists(t, path)
env.sweeper.ExportSweep(context.Background())
for _, suffix := range archiveFileSuffixes() {
assert.NoFileExists(
t, path+suffix,
"the sweep must not create an archive file",
)
}
}
// TestArchiveSweep_DoesNotCreateAfterWriterExists covers the
// same guarantee once a writer is cached in the registry but
// the file itself is still absent (for instance because the
// operator moved the archive away).
func TestArchiveSweep_DoesNotCreateAfterWriterExists(
t *testing.T,
) {
t.Parallel()
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path, err := env.eng.ExportEnsureArchiveWriter(webhookID)
require.NoError(t, err)
require.NoFileExists(t, path)
env.sweeper.ExportSweep(context.Background())
assert.NoFileExists(t, path)
}
// TestArchiveSweep_SkipsDeletedWebhookTargets proves that the
// sweep ignores targets soft-deleted along with their webhook,
// so a deleted webhook's archive is never reopened.
func TestArchiveSweep_SkipsDeletedWebhookTargets(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
)
require.NoError(
t,
env.mainDB.DB().
Where("webhook_id = ?", webhookID).
Delete(&database.Target{}).Error,
)
env.sweeper.ExportSweep(context.Background())
assert.Equal(
t, []string{sweepRowOld}, archivedEventIDs(t, path),
"a deleted target's archive must be left alone",
)
}
// TestArchiveSweep_ConcurrentWrites proves the sweep serialises
// against writes through the per-webhook writer mutex. Run
// under -race, an unsynchronised sweep would be caught here.
func TestArchiveSweep_ConcurrentWrites(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
webhookDB := testWebhookDB(t)
// The deliveries are seeded up front, on the test's own
// goroutine: the seed helpers assert, and testify assertions
// must not run off the test goroutine.
deliveries := make(
[]*database.Delivery, 0, sweepConcurrentWrites,
)
for range sweepConcurrentWrites {
event := seedEvent(t, webhookDB, `{"n":1}`)
event.WebhookID = webhookID
deliveries = append(
deliveries,
seedDatabaseTargetDelivery(
t, webhookDB, event, `{"expiry":"1h"}`,
),
)
}
var wg sync.WaitGroup
wg.Add(2)
go func() {
defer wg.Done()
for _, d := range deliveries {
env.eng.ExportDeliverDatabase(webhookDB, d)
}
}()
go func() {
defer wg.Done()
for range sweepConcurrentWrites {
env.sweeper.ExportSweep(context.Background())
}
}()
wg.Wait()
assert.FileExists(t, env.archivePath(webhookID))
}
// TestArchiveSweeper_StopsCleanly proves the background loop
// exits on OnStop rather than leaking a goroutine.
func TestArchiveSweeper_StopsCleanly(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
)
env.sweeper.ExportSetInterval(time.Millisecond)
env.sweeper.ExportStart()
// stop blocks on the loop's WaitGroup, so returning at all
// proves the loop observed the cancellation and exited.
env.sweeper.ExportStop()
}

View File

@@ -94,23 +94,6 @@ type Notifier interface {
Notify(tasks []Task)
}
// WebhookEvictor releases the delivery engine's per-webhook
// state for a webhook that no longer needs it — currently the
// cached archive writer of the database target, whose open
// file handle would otherwise outlive the webhook.
//
// It is deliberately separate from Notifier and deliberately
// one method wide: archiving lifecycle is not notification, and
// a single-method interface keeps the handlers package free of
// any dependency on the engine's internals while staying
// trivially fakeable in tests.
//
// EvictWebhook never deletes an archive file. It is idempotent
// and is a no-op for a webhook with no engine state.
type WebhookEvictor interface {
EvictWebhook(webhookID string)
}
// EngineParams are the fx dependencies for the delivery
// engine.
type EngineParams struct {
@@ -144,10 +127,6 @@ type Engine struct {
// httpTarget is retained so tests can reach the HTTP
// target's shared client and circuit breakers.
httpTarget *httpTarget
// dbTarget is retained so the engine can reach the archive
// writer registry for webhook eviction and the idle sweep.
dbTarget *databaseTarget
}
// New creates and registers the delivery engine with the
@@ -203,19 +182,6 @@ func (e *Engine) Notify(tasks []Task) {
}
}
// EvictWebhook implements WebhookEvictor. It releases the
// engine's per-webhook archiving state: the database target's
// cached archive writer is dropped from the registry and its
// file handle closed. The archive file itself is left on disk
// — it is long-term storage the operator owns.
func (e *Engine) EvictWebhook(webhookID string) {
if e.dbTarget == nil {
return
}
e.dbTarget.evict(webhookID)
}
// 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.
@@ -487,8 +453,9 @@ func (e *Engine) recoverRetryingDeliveries(
// recoverSingleRetry hands an orphaned retrying delivery back
// to its target to recompute the remaining backoff, then
// reschedules it. Targets that do not own durable retries
// (fire-and-forget) never produce retrying deliveries, so
// they are skipped.
// (fire-and-forget) never produce retrying deliveries, so a
// delivery found in that state has had its target's type
// changed underneath it and is terminally failed.
func (e *Engine) recoverSingleRetry(
webhookDB *gorm.DB,
webhookID string,
@@ -509,6 +476,10 @@ func (e *Engine) recoverSingleRetry(
rs, ok := e.targets[target.Type].(rescheduler)
if !ok {
e.failUnretryableRetry(
webhookDB, webhookID, d, &target,
)
return
}
@@ -683,8 +654,8 @@ func (e *Engine) sweepWebhookRetries(
// sweepSingleRetry re-enqueues an orphaned retrying delivery
// whose backoff window has elapsed, delegating the backoff
// decision to the delivery's target. Targets that do not own
// durable retries are skipped.
// decision to the delivery's target. A delivery whose target
// no longer owns durable retries is terminally failed.
func (e *Engine) sweepSingleRetry(
webhookDB *gorm.DB,
webhookID string,
@@ -704,6 +675,10 @@ func (e *Engine) sweepSingleRetry(
rs, ok := e.targets[target.Type].(rescheduler)
if !ok {
e.failUnretryableRetry(
webhookDB, webhookID, d, &target,
)
return
}
@@ -744,6 +719,59 @@ func (e *Engine) sweepSingleRetry(
}
}
// failUnretryableRetry terminally fails an orphaned retrying
// delivery whose target type no longer supports retries. Both
// restart recovery and the periodic sweep call it, so the
// terminal transition exists once.
//
// This is only reachable when a target's type has been changed
// out from under an in-flight retrying delivery (or the type is
// unknown to the registry): fire-and-forget targets never set
// status retrying themselves. Re-dispatching under the new type
// would be a delivery the operator never asked for, and leaving
// the row retrying strands it forever, so the delivery is
// failed with a recorded reason and can be redelivered
// manually. Logged at warn, not error: this is operator-caused
// state, not a system fault.
func (e *Engine) failUnretryableRetry(
webhookDB *gorm.DB,
webhookID string,
d *database.Delivery,
target *database.Target,
) {
e.log.Warn(
"failing orphaned retrying delivery: target "+
"type no longer supports retries",
"webhook_id", webhookID,
"delivery_id", d.ID,
"target_id", target.ID,
"target_name", target.Name,
"target_type", target.Type,
)
reason := fmt.Sprintf(
"target type %q does not support retries; "+
"delivery was left retrying by a previous "+
"target type and has been failed terminally",
target.Type,
)
e.recordResult(
webhookDB,
d,
e.countAttempts(webhookDB, d.ID)+1,
false,
0,
"",
reason,
0,
)
e.updateDeliveryStatus(
webhookDB, d, database.DeliveryStatusFailed,
)
}
// processDelivery dispatches a delivery to the target that
// owns its type. Unknown target types fail the delivery.
func (e *Engine) processDelivery(

View File

@@ -748,6 +748,193 @@ func TestRecoverWebhookDeliveries_RetryingDeliveries(
case <-time.After(5 * time.Second):
t.Fatal("expected retry task from recovery")
}
// Regression guard: a target that still supports retries
// must be rescheduled, never terminally failed, and must
// not gain a synthetic result row.
iAssertStatus(
t, s.WebhookDB, d.ID,
database.DeliveryStatusRetrying,
)
assert.Len(t, iResults(t, s.WebhookDB, d.ID), 1)
}
// --- Retrying deliveries whose target type changed ---
// iSeedRetryingWithType seeds a retrying delivery with one
// recorded failed attempt against a target of the given type,
// standing in for a target whose type was edited in the main
// database while the delivery was still retrying.
func iSeedRetryingWithType(
t *testing.T,
s iSetup,
targetType database.TargetType,
) string {
t.Helper()
targetID := uuid.New().String()
iCreateTarget(t, s.MainDB, targetID,
s.WebhookID, "mutated-target", targetType,
iHTTPConfig("http://example.com/hook"), 5,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID,
`{"orphaned":"retry"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
iSeedFailedResult(t, s.WebhookDB, d.ID)
return d.ID
}
// iResults loads a delivery's results in attempt order.
func iResults(
t *testing.T, db *gorm.DB, deliveryID string,
) []database.DeliveryResult {
t.Helper()
var results []database.DeliveryResult
require.NoError(t, db.
Where("delivery_id = ?", deliveryID).
Order("attempt_num").
Find(&results).Error)
return results
}
// iAssertTerminallyFailed asserts the delivery ended failed
// with a result row recording why, and was not rescheduled.
func iAssertTerminallyFailed(
t *testing.T,
s iSetup,
deliveryID string,
targetType database.TargetType,
) {
t.Helper()
iAssertStatus(
t, s.WebhookDB, deliveryID,
database.DeliveryStatusFailed,
)
results := iResults(t, s.WebhookDB, deliveryID)
require.Len(t, results, 2)
last := results[1]
assert.False(t, last.Success)
assert.Equal(t, 2, last.AttemptNum)
assert.Contains(
t, last.Error, string(targetType),
)
assert.Contains(
t, last.Error, "does not support retries",
)
assert.Empty(t, s.Engine.ExportRetryCh())
}
func TestRecoverSingleRetry_TypeNoLongerRetries(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "mutated-type",
)
deliveryID := iSeedRetryingWithType(
t, s, database.TargetTypeLog,
)
s.Engine.ExportRecoverWebhookDeliveries(
context.Background(), s.WebhookID,
)
iAssertTerminallyFailed(
t, s, deliveryID, database.TargetTypeLog,
)
}
func TestSweepSingleRetry_TypeNoLongerRetries(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "mutated-type-sweep",
)
deliveryID := iSeedRetryingWithType(
t, s, database.TargetTypeDatabase,
)
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
iAssertTerminallyFailed(
t, s, deliveryID, database.TargetTypeDatabase,
)
}
func TestRecoverSingleRetry_UnknownTargetType(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "unknown-type",
)
unknown := database.TargetType("not-a-target-type")
deliveryID := iSeedRetryingWithType(t, s, unknown)
s.Engine.ExportRecoverWebhookDeliveries(
context.Background(), s.WebhookID,
)
iAssertTerminallyFailed(t, s, deliveryID, unknown)
}
func TestSweepSingleRetry_UnknownTargetType(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "unknown-type-sweep",
)
unknown := database.TargetType("not-a-target-type")
deliveryID := iSeedRetryingWithType(t, s, unknown)
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
iAssertTerminallyFailed(t, s, deliveryID, unknown)
}
// iSeedFailedResult creates a failed delivery result.

View File

@@ -7,17 +7,10 @@ import (
"net/http"
"time"
"go.uber.org/fx"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
)
// ErrExportArchiveWriterEvicted exposes the sentinel returned by
// an evicted archive writer. It carries the Err prefix rather
// than this file's usual Export one because it is a sentinel
// error.
var ErrExportArchiveWriterEvicted = errArchiveWriterEvicted
// Exported constants for test access.
const (
ExportDeliveryChannelSize = deliveryChannelSize
@@ -195,6 +188,13 @@ func (e *Engine) ExportRecoverInFlight(
e.recoverInFlight(ctx)
}
// ExportSweepWebhookRetries exposes sweepWebhookRetries.
func (e *Engine) ExportSweepWebhookRetries(
ctx context.Context, webhookID string,
) {
e.sweepWebhookRetries(ctx, webhookID)
}
// ExportStart exposes start for testing.
func (e *Engine) ExportStart(ctx context.Context) {
e.start(ctx)
@@ -335,151 +335,6 @@ func (e *ExportArchiveWriter) DB() *gorm.DB {
return e.w.db
}
// Path returns the archive file the writer owns.
func (e *ExportArchiveWriter) Path() string {
return e.w.path
}
// OpenExisting opens the archive without permitting creation,
// the way the idle sweep does.
func (e *ExportArchiveWriter) OpenExisting(
expiry time.Duration,
) error {
return e.w.openMode(archiveModeExisting, expiry)
}
// SweepExpired runs an idle sweep of the archive.
func (e *ExportArchiveWriter) SweepExpired(
expiry time.Duration,
) error {
return e.w.sweepExpired(expiry)
}
// Evict marks the writer evicted and closes its handle, exactly
// as leaving the registry does.
func (e *ExportArchiveWriter) Evict() {
e.w.evict()
}
// HandleOpen reports whether the writer currently holds an open
// archive handle.
func (e *ExportArchiveWriter) HandleOpen() bool {
e.w.mu.Lock()
defer e.w.mu.Unlock()
return e.w.db != nil
}
// ExportArchiveWriterFor returns the archive writer the registry
// currently caches for a webhook, or nil when none is cached. It
// never creates one, so a test can hold a reference to the very
// writer an eviction is about to detach.
func (e *Engine) ExportArchiveWriterFor(
webhookID string,
) *ExportArchiveWriter {
e.dbTarget.mu.Lock()
defer e.dbTarget.mu.Unlock()
w, ok := e.dbTarget.writers[webhookID]
if !ok {
return nil
}
return &ExportArchiveWriter{w: w}
}
// ExportHasArchiveWriter reports whether the database target
// currently caches an archive writer for a webhook.
func (e *Engine) ExportHasArchiveWriter(
webhookID string,
) bool {
e.dbTarget.mu.Lock()
defer e.dbTarget.mu.Unlock()
_, ok := e.dbTarget.writers[webhookID]
return ok
}
// ExportArchiveHandleOpen reports whether the cached archive
// writer for a webhook holds an open database handle. It
// returns false when no writer is cached.
func (e *Engine) ExportArchiveHandleOpen(
webhookID string,
) bool {
e.dbTarget.mu.Lock()
w, ok := e.dbTarget.writers[webhookID]
e.dbTarget.mu.Unlock()
if !ok {
return false
}
w.mu.Lock()
defer w.mu.Unlock()
return w.db != nil
}
// ExportEnsureArchiveWriter creates (if needed) and returns the
// archive file path of the cached writer for a webhook, so a
// test can prime the registry the way a delivery would.
func (e *Engine) ExportEnsureArchiveWriter(
webhookID string,
) (string, error) {
w, err := e.dbTarget.writerFor(webhookID)
if err != nil {
return "", err
}
return w.path, nil
}
// NewTestArchiveSweeper builds an ArchiveSweeper backed by the
// given main database and engine, without the fx lifecycle.
// Intended for tests.
func NewTestArchiveSweeper(
db *database.Database,
eng *Engine,
log *slog.Logger,
) *ArchiveSweeper {
return &ArchiveSweeper{
db: db,
eng: eng,
log: log,
interval: time.Hour,
}
}
// ExportSweep runs a single archive sweep synchronously for
// tests.
func (s *ArchiveSweeper) ExportSweep(ctx context.Context) {
s.sweep(ctx)
}
// ExportStart starts the sweeper's background loop for tests.
func (s *ArchiveSweeper) ExportStart() {
s.start()
}
// ExportRegisterHooks registers the sweeper's real fx lifecycle
// hooks on a lifecycle supplied by a test, so a test can drive
// the exact OnStart/OnStop functions the application runs and
// hand OnStart the kind of context fx actually supplies.
func (s *ArchiveSweeper) ExportRegisterHooks(lc fx.Lifecycle) {
s.registerHooks(lc)
}
// ExportStop stops the sweeper's background loop for tests.
func (s *ArchiveSweeper) ExportStop() {
s.stop()
}
// ExportSetInterval overrides the sweep interval for tests.
func (s *ArchiveSweeper) ExportSetInterval(d time.Duration) {
s.interval = d
}
// ExportParseArchiveExpiry exposes parseArchiveExpiry.
func ExportParseArchiveExpiry(
configJSON string,

View File

@@ -90,15 +90,12 @@ func (e *Engine) initTargets(client *http.Client) {
client: client,
}
dbT := &databaseTarget{eng: e}
e.httpTarget = httpT
e.dbTarget = dbT
e.targets = map[database.TargetType]Target{
database.TargetTypeHTTP: httpT,
database.TargetTypeSlack: slackT,
database.TargetTypeDatabase: dbT,
database.TargetTypeDatabase: &databaseTarget{eng: e},
database.TargetTypeLog: &logTarget{eng: e},
}
}

View File

@@ -5,7 +5,6 @@ import (
"fmt"
"path/filepath"
"sync"
"time"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
@@ -112,184 +111,27 @@ func (t *databaseTarget) archive(d *database.Delivery) error {
func (t *databaseTarget) writerFor(
webhookID string,
) (*archiveWriter, error) {
path, err := t.archivePath(webhookID)
if err != nil {
return nil, err
}
t.mu.Lock()
defer t.mu.Unlock()
if t.writers == nil {
t.writers = make(map[string]*archiveWriter)
}
w, ok := t.writers[webhookID]
if !ok {
w = newArchiveWriter(path, t.eng.log)
t.writers[webhookID] = w
}
// A delivery claims the entry: even if the idle sweep created
// it moments ago, it now belongs to the registry proper and
// the sweep must leave it in place when it finishes.
w.sweepOwned = false
return w, nil
}
// sweepWriterFor returns the archive writer the idle sweep should
// prune a webhook through, together with whether the sweep itself
// created the registry entry.
//
// The sweep must route its prune through the registered writer so
// the writer's mutex orders it against concurrent writes, but it
// must never leave a registry entry behind: a sweep that ran
// concurrently with the webhook's deletion would otherwise
// re-create an entry that nothing will ever evict again, which is
// exactly the leak eviction exists to prevent. An entry the sweep
// creates is therefore marked sweep-owned and handed back to
// releaseSweepWriter when the sweep is done.
func (t *databaseTarget) sweepWriterFor(
webhookID string,
) (*archiveWriter, bool, error) {
path, err := t.archivePath(webhookID)
if err != nil {
return nil, false, err
}
t.mu.Lock()
defer t.mu.Unlock()
if t.writers == nil {
t.writers = make(map[string]*archiveWriter)
}
w, ok := t.writers[webhookID]
if ok {
return w, false, nil
}
w = newArchiveWriter(path, t.eng.log)
w.sweepOwned = true
t.writers[webhookID] = w
return w, true, nil
}
// releaseSweepWriter drops a registry entry that the idle sweep
// created, so a sweep leaves the registry exactly as it found it.
//
// The entry is removed only if it is still the very writer the
// sweep installed and no delivery has claimed it in the meantime
// (writerFor clears sweepOwned when it hands a writer to the
// write path). Both conditions are evaluated under the registry
// lock, so an eviction that raced the sweep — which removes the
// entry outright — simply finds nothing left to do here, and a
// delivery that adopted the writer keeps a registered, evictable
// one.
func (t *databaseTarget) releaseSweepWriter(
webhookID string, w *archiveWriter,
) {
t.mu.Lock()
defer t.mu.Unlock()
cur, ok := t.writers[webhookID]
if !ok || cur != w || !cur.sweepOwned {
return
}
delete(t.writers, webhookID)
}
// archivePath returns the archive file path for a webhook: it
// lives beside the per-webhook event database in the data
// directory. It does not touch the filesystem.
func (t *databaseTarget) archivePath(
webhookID string,
) (string, error) {
if t.eng.dbManager == nil {
return "", errArchiveNoDataDir
return nil, errArchiveNoDataDir
}
dir := filepath.Dir(t.eng.dbManager.DBPath(webhookID))
return filepath.Join(
path := filepath.Join(
dir, fmt.Sprintf("archive-%s.db", webhookID),
), nil
}
)
// evict drops a webhook's archive writer from the registry and
// closes its handle, so a deleted webhook does not leave a
// writer (and an open archive handle within its debounce
// window) alive for the process lifetime.
//
// The map entry is removed under the registry lock, which is
// then released before the handle is closed under the writer's
// own lock: that ordering keeps the registry available to other
// webhooks while an in-flight write on this one drains, and
// closing under the writer's lock means eviction can never race
// a write.
//
// Eviction is idempotent and silent for a webhook with no
// writer, which is the common case: a webhook with no database
// target never creates one. It never deletes the archive file.
func (t *databaseTarget) evict(webhookID string) {
t.mu.Lock()
defer t.mu.Unlock()
if t.writers == nil {
t.writers = make(map[string]*archiveWriter)
}
w, ok := t.writers[webhookID]
if ok {
delete(t.writers, webhookID)
}
t.mu.Unlock()
if !ok {
return
w = newArchiveWriter(path, t.eng.log)
t.writers[webhookID] = w
}
w.evict()
t.eng.log.Info(
"evicted archive writer",
"webhook_id", webhookID,
"path", w.path,
)
}
// sweepWebhook prunes one webhook's archive of rows older than
// expiry, without requiring a write. It returns nil (nothing to
// do) when the archive file does not exist, so a sweep never
// creates an archive for a webhook that has a database target
// but has never received an event.
//
// It also never leaves a registry entry behind: an entry it had
// to create to reach the writer's mutex is released again once
// the prune is done, so a sweep racing a webhook deletion cannot
// resurrect the writer the eviction just dropped.
func (t *databaseTarget) sweepWebhook(
webhookID string, expiry time.Duration,
) error {
path, err := t.archivePath(webhookID)
if err != nil {
return err
}
// Check before taking a writer at all: a webhook whose
// archive has never been created gets no writer, no handle,
// and no file.
if !fileExists(path) {
return nil
}
w, created, err := t.sweepWriterFor(webhookID)
if err != nil {
return err
}
if created {
defer t.releaseSweepWriter(webhookID, w)
}
return w.sweepExpired(expiry)
return w, nil
}

View File

@@ -24,20 +24,6 @@ const archiveExpiryNever = "never"
// offline archiving, but never more than once per this window.
const archiveReopenDebounce = time.Second
const (
// archiveModeCreate is the SQLite URI mode used by the write
// path: open the archive file, creating it if missing, so a
// first write (or a write after the operator moved the file
// away) recreates it.
archiveModeCreate = "rwc"
// archiveModeExisting is the SQLite URI mode used by the idle
// sweep: open read-write but never create. A sweep must never
// conjure an empty archive file for a webhook that has a
// database target but has never received an event.
archiveModeExisting = "rw"
)
var (
// errArchiveMissingWebhookID is returned when an event to
// archive has no webhook id to key its archive file on.
@@ -58,15 +44,6 @@ var (
errArchiveExpiryNotPositive = errors.New(
"expiry must be a positive duration or \"never\"",
)
// errArchiveWriterEvicted is returned when a writer that has
// been evicted (its webhook was deleted, or its last database
// target was removed) is used again. An evicted writer is no
// longer in the registry, so reopening its file would leak a
// handle nothing owns.
errArchiveWriterEvicted = errors.New(
"archive writer has been evicted",
)
)
// databaseTargetConfig is the optional per-target JSON config
@@ -184,25 +161,6 @@ type archiveWriter struct {
db *gorm.DB
lastReopen time.Time
reopens int
// evicted marks a writer that has been removed from the
// per-webhook registry. Its handle is closed and it must
// never open the file again: nothing holds it any more, so a
// reopen would leak the handle for the process lifetime.
evicted bool
// sweepOwned marks a registry entry that the idle sweep
// created because no writer was cached for the webhook. The
// sweep removes such an entry again when it is done, so a
// sweep can never leave — or resurrect — a registry entry
// for a webhook that has been deleted. A delivery that adopts
// the writer clears the flag, handing the entry to the
// registry proper.
//
// Unlike every other field here it is guarded by
// databaseTarget.mu, not by this writer's mu: it describes the
// registry entry rather than the file.
sweepOwned bool
}
// newArchiveWriter builds an archiveWriter for a file path with
@@ -227,12 +185,6 @@ func (w *archiveWriter) write(
w.mu.Lock()
defer w.mu.Unlock()
if w.evicted {
return fmt.Errorf(
"%w: %s", errArchiveWriterEvicted, w.path,
)
}
if w.db == nil || !fileExists(w.path) {
err := w.reopen(expiry)
if err != nil {
@@ -260,19 +212,7 @@ func (w *archiveWriter) write(
// its schema, records the reopen time, and prunes expired rows
// when expiry is positive.
func (w *archiveWriter) open(expiry time.Duration) error {
return w.openMode(archiveModeCreate, expiry)
}
// openMode opens the archive file with the given SQLite URI
// mode, migrates its schema, records the reopen time, and
// prunes expired rows when expiry is positive. The write path
// passes archiveModeCreate so a missing file is recreated; the
// idle sweep passes archiveModeExisting so a missing file is an
// error rather than a newly conjured empty archive.
func (w *archiveWriter) openMode(
mode string, expiry time.Duration,
) error {
dbURL := fmt.Sprintf("file:%s?mode=%s", w.path, mode)
dbURL := fmt.Sprintf("file:%s?mode=rwc", w.path)
sqlDB, err := sql.Open("sqlite", dbURL)
if err != nil {
@@ -335,63 +275,6 @@ func (w *archiveWriter) close() {
w.db = nil
}
// sweepExpired prunes an archive that may have gone idle, with
// no write to trigger the usual on-reopen prune. It takes the
// writer's own mutex for the whole operation, so a sweep is
// ordered against concurrent writes rather than reaching around
// them to the file.
//
// It never creates the archive file: a missing file is skipped,
// and the reopen uses archiveModeExisting so SQLite itself
// refuses to create one if the file disappears between the
// check and the open.
//
// The archive is left CLOSED afterwards. An idle archive holding
// no handle is what keeps the operator's move-the-file-away
// workflow working; the next write reopens (and recreates) the
// file as it always has.
func (w *archiveWriter) sweepExpired(expiry time.Duration) error {
w.mu.Lock()
defer w.mu.Unlock()
if w.evicted {
return fmt.Errorf(
"%w: %s", errArchiveWriterEvicted, w.path,
)
}
if !fileExists(w.path) {
return nil
}
// Drop any live handle first so the prune runs against a
// freshly opened file, matching the write path's semantics.
w.close()
err := w.openMode(archiveModeExisting, expiry)
if err != nil {
return err
}
w.close()
return nil
}
// evict closes the writer's handle and marks it unusable. It is
// called when the writer leaves the registry, either because the
// webhook was deleted or because its last database target was
// removed. The archive FILE is deliberately left on disk: it is
// long-term storage an operator may still want.
func (w *archiveWriter) evict() {
w.mu.Lock()
defer w.mu.Unlock()
w.evicted = true
w.close()
}
// prune deletes archived rows older than expiry, measured from
// each row's archived time. It runs on every (re)open, and
// because the file is reopened after writes this keeps the

View File

@@ -1,363 +0,0 @@
package delivery_test
import (
"errors"
"fmt"
"net/http"
"os"
"path/filepath"
"sync"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
)
// evictTestEngine builds an engine backed by a temporary data
// directory and returns it along with that directory.
func evictTestEngine(t *testing.T) (*delivery.Engine, string) {
t.Helper()
dataDir := t.TempDir()
eng := delivery.NewTestEngineWithDB(
nil,
database.NewTestWebhookDBManager(dataDir),
archiveTestLogger(),
&http.Client{Timeout: 5 * time.Second},
1,
)
return eng, dataDir
}
// TestEvictWebhook_ClosesAndRemovesWriter proves that evicting
// a webhook drops its archive writer from the registry and
// closes the open archive handle, rather than leaving both
// alive for the process lifetime.
func TestEvictWebhook_ClosesAndRemovesWriter(t *testing.T) {
t.Parallel()
eng, dataDir := evictTestEngine(t)
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":true}`)
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
eng.ExportDeliverDatabase(webhookDB, d)
webhookID := event.WebhookID
require.True(
t, eng.ExportHasArchiveWriter(webhookID),
"a delivery should have cached an archive writer",
)
require.True(
t, eng.ExportArchiveHandleOpen(webhookID),
"the writer should hold an open handle after a write",
)
eng.EvictWebhook(webhookID)
assert.False(
t, eng.ExportHasArchiveWriter(webhookID),
"eviction should remove the registry entry",
)
assert.False(
t, eng.ExportArchiveHandleOpen(webhookID),
"eviction should close the archive handle",
)
archivePath := filepath.Join(
dataDir, fmt.Sprintf("archive-%s.db", webhookID),
)
assert.FileExists(
t, archivePath,
"eviction must not delete the archive file",
)
}
// TestEvictWebhook_UnknownWebhookIsNoOp proves eviction is safe
// for the common case of a webhook that never had a database
// target, and that repeating it does not panic.
func TestEvictWebhook_UnknownWebhookIsNoOp(t *testing.T) {
t.Parallel()
eng, _ := evictTestEngine(t)
assert.NotPanics(t, func() {
eng.EvictWebhook("no-such-webhook")
eng.EvictWebhook("no-such-webhook")
})
assert.False(
t, eng.ExportHasArchiveWriter("no-such-webhook"),
"eviction must not create a writer",
)
}
// evictTestRow builds an archive row for the eviction tests.
func evictTestRow(eventID string) delivery.ExportArchivedEvent {
return delivery.ExportArchivedEvent{
EventID: eventID,
WebhookID: "wh-evict",
Method: http.MethodPost,
Body: `{"seeded":true}`,
}
}
// TestEvictedWriter_WriteDoesNotReopenFile is the direct test of
// the evicted guard on the write path. A writer that has left
// the registry is held by nobody, so a handle it opened could
// never be closed again: it must refuse the write outright
// rather than recreate the archive behind the registry's back.
//
// The archive file is removed before the eviction, so an
// unguarded write is unmistakable — it recreates the file.
func TestEvictedWriter_WriteDoesNotReopenFile(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "archive-evicted.db")
w := delivery.NewExportArchiveWriter(
path, archiveTestLogger(), 0,
)
require.NoError(t, w.Write(evictTestRow("ev-1"), 0))
require.FileExists(t, path)
// The operator moves the archive away for offline retention,
// which the write path would ordinarily undo on the next
// write by recreating the file.
require.NoError(t, os.Remove(path))
w.Evict()
err := w.Write(evictTestRow("ev-2"), 0)
require.ErrorIs(
t, err, delivery.ErrExportArchiveWriterEvicted,
"an evicted writer must refuse writes",
)
assert.NoFileExists(
t, path,
"an evicted writer must not reopen (or recreate) the "+
"archive file",
)
assert.False(
t, w.HandleOpen(),
"an evicted writer must hold no handle",
)
}
// TestEvictedWriter_SweepDoesNotReopenFile is the same test for
// the sweep path: an idle sweep that reaches a writer already
// evicted underneath it must return the sentinel rather than
// reopen a file nothing owns.
func TestEvictedWriter_SweepDoesNotReopenFile(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "archive-evicted.db")
w := delivery.NewExportArchiveWriter(
path, archiveTestLogger(), 0,
)
require.NoError(t, w.Write(evictTestRow("ev-1"), 0))
require.FileExists(t, path)
w.Evict()
err := w.SweepExpired(time.Hour)
require.ErrorIs(
t, err, delivery.ErrExportArchiveWriterEvicted,
"an evicted writer must refuse an idle sweep",
)
assert.False(
t, w.HandleOpen(),
"a refused sweep must not leave a handle open",
)
}
// racingWrites drives a pack of goroutines writing to one
// archive writer until each is refused, so an eviction on the
// test goroutine has to take the writer's mutex away from writes
// that are already contending for it.
type racingWrites struct {
wg sync.WaitGroup
mu sync.Mutex
sawEvicted bool
otherErr error
started chan struct{}
}
// racingWriteGoroutines is how many goroutines contend for the
// writer's mutex while the eviction lands.
const racingWriteGoroutines = 4
// startRacingWrites launches the writing goroutines. Each writes
// in a loop and stops at its first error, recording whether that
// error was the eviction sentinel. The deadline is a backstop
// against a hang, not a timing assumption: the first write after
// the eviction is refused.
func startRacingWrites(
w *delivery.ExportArchiveWriter,
) *racingWrites {
r := &racingWrites{
started: make(chan struct{}, racingWriteGoroutines),
}
deadline := time.Now().Add(10 * time.Second)
r.wg.Add(racingWriteGoroutines)
for i := range racingWriteGoroutines {
go func() {
defer r.wg.Done()
first := true
for time.Now().Before(deadline) {
err := w.Write(
evictTestRow(fmt.Sprintf("ev-%d", i)), 0,
)
if first {
r.started <- struct{}{}
first = false
}
if err == nil {
continue
}
r.record(err)
return
}
}()
}
return r
}
// record classifies the error that stopped one goroutine.
func (r *racingWrites) record(err error) {
r.mu.Lock()
defer r.mu.Unlock()
if errors.Is(err, delivery.ErrExportArchiveWriterEvicted) {
r.sawEvicted = true
return
}
r.otherErr = err
}
// awaitFirstWrite blocks until at least one write has run, so
// the eviction that follows is a genuine race.
func (r *racingWrites) awaitFirstWrite() {
<-r.started
}
// wait joins the goroutines and reports whether any write was
// refused with the eviction sentinel, plus any unexpected error.
func (r *racingWrites) wait() (bool, error) {
r.wg.Wait()
r.mu.Lock()
defer r.mu.Unlock()
return r.sawEvicted, r.otherErr
}
// TestEvictWebhook_RacingWriteDoesNotReopenHandle exercises the
// interleaving the evicted flag exists for: writes already
// contending for the writer's mutex when the eviction takes it.
// The write that wins the mutex after the eviction must abandon
// its work rather than reopen the archive, leaving the writer
// permanently handle-free. Run under -race.
func TestEvictWebhook_RacingWriteDoesNotReopenHandle(
t *testing.T,
) {
t.Parallel()
eng, _ := evictTestEngine(t)
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":true}`)
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
// Prime the registry so the test can hold the very writer the
// eviction is about to detach.
eng.ExportDeliverDatabase(webhookDB, d)
w := eng.ExportArchiveWriterFor(event.WebhookID)
require.NotNil(t, w)
require.True(t, w.HandleOpen())
race := startRacingWrites(w)
// Evict only once writes are genuinely in flight, so the
// eviction has to contend for the writer's mutex.
race.awaitFirstWrite()
eng.EvictWebhook(event.WebhookID)
sawEvicted, otherErr := race.wait()
require.NoError(t, otherErr)
assert.True(
t, sawEvicted,
"a write after eviction must be refused",
)
assert.False(
t, w.HandleOpen(),
"no write may reopen the archive once the writer has "+
"been evicted",
)
assert.False(
t, eng.ExportHasArchiveWriter(event.WebhookID),
"the registry entry must stay gone",
)
}
// TestEvictWebhook_LaterDeliveryRecreatesWriter proves eviction
// does not break archiving for a webhook that is still alive: a
// subsequent delivery gets a brand new writer from the registry.
// It says nothing about the evicted writer itself — that is what
// TestEvictedWriter_WriteDoesNotReopenFile covers.
func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) {
t.Parallel()
eng, _ := evictTestEngine(t)
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":true}`)
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
eng.ExportDeliverDatabase(webhookDB, d)
require.True(
t, eng.ExportHasArchiveWriter(event.WebhookID),
)
eng.EvictWebhook(event.WebhookID)
// A fresh delivery for the same webhook gets a brand new
// writer from the registry, so archiving keeps working.
second := seedDatabaseTargetDelivery(
t, webhookDB, event, "",
)
eng.ExportDeliverDatabase(webhookDB, second)
assert.True(
t, eng.ExportHasArchiveWriter(event.WebhookID),
"a later delivery should recreate the writer",
)
}

View File

@@ -50,13 +50,6 @@ func openArchiveDBForRead(
return gdb
}
// archiveFileSuffixes returns the archive file itself and the
// SQLite sidecars that accompany an open database. A test that
// asserts no archive was created has to check all of them.
func archiveFileSuffixes() []string {
return []string{"", "-wal", "-shm"}
}
// removeArchiveFiles simulates an operator moving the archive
// away by deleting the SQLite file and its sidecar files.
func removeArchiveFiles(t *testing.T, path string) {

View File

@@ -51,7 +51,6 @@ type HandlersParams struct {
Healthcheck *healthcheck.Healthcheck
Session *session.Session
Notifier delivery.Notifier
Evictor delivery.WebhookEvictor
}
// Handlers provides HTTP handler methods for all application
@@ -64,7 +63,6 @@ type Handlers struct {
dbMgr *database.WebhookDBManager
session *session.Session
notifier delivery.Notifier
evictor delivery.WebhookEvictor
templates map[string]*template.Template
}
@@ -99,7 +97,6 @@ func New(
s.dbMgr = params.WebhookDBMgr
s.session = params.Session
s.notifier = params.Notifier
s.evictor = params.Evictor
// Parse all page templates once at startup
s.templates = map[string]*template.Template{

View File

@@ -4,7 +4,6 @@ import (
"context"
"net/http"
"net/http/httptest"
"sync"
"testing"
"github.com/stretchr/testify/assert"
@@ -25,32 +24,6 @@ type noopNotifier struct{}
func (n *noopNotifier) Notify([]delivery.Task) {}
// recordingEvictor is a delivery.WebhookEvictor that records
// the webhook ids it was asked to evict, so a test can prove
// that a deletion path reached the delivery engine.
type recordingEvictor struct {
mu sync.Mutex
evicted []string
}
func (r *recordingEvictor) EvictWebhook(webhookID string) {
r.mu.Lock()
defer r.mu.Unlock()
r.evicted = append(r.evicted, webhookID)
}
// Evicted returns a copy of the recorded webhook ids.
func (r *recordingEvictor) Evicted() []string {
r.mu.Lock()
defer r.mu.Unlock()
out := make([]string, len(r.evicted))
copy(out, r.evicted)
return out
}
func newTestApp(
t *testing.T,
targets ...any,
@@ -74,12 +47,6 @@ func newTestApp(
func() delivery.Notifier {
return &noopNotifier{}
},
func() *recordingEvictor {
return &recordingEvictor{}
},
func(r *recordingEvictor) delivery.WebhookEvictor {
return r
},
handlers.New,
),
fx.Populate(targets...),

View File

@@ -1,355 +0,0 @@
package handlers_test
import (
"context"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"testing"
"github.com/go-chi/chi"
"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/handlers"
"sneak.berlin/go/webhooker/internal/session"
)
const (
deleteTestUserID = "test-user-id"
deleteTestUsername = "testuser"
// paramSourceID and paramTargetID are the chi URL parameter
// names the deletion handlers read.
paramSourceID = "sourceID"
paramTargetID = "targetID"
)
// seedWebhook inserts a webhook owned by the test user and
// returns it.
func seedWebhook(
t *testing.T,
db *database.Database,
) *database.Webhook {
t.Helper()
wh := &database.Webhook{
UserID: deleteTestUserID,
Name: "delete-me",
}
require.NoError(
t,
db.DB().Omit(clause.Associations).Create(wh).Error,
)
return wh
}
// seedTarget inserts a target of the given type for a webhook
// and returns it.
func seedTarget(
t *testing.T,
db *database.Database,
webhookID string,
targetType database.TargetType,
) *database.Target {
t.Helper()
tgt := &database.Target{
WebhookID: webhookID,
Name: "t-" + string(targetType),
Type: targetType,
Active: true,
}
require.NoError(
t,
db.DB().Omit(clause.Associations).Create(tgt).Error,
)
return tgt
}
// archivePathFor returns the archive database path the
// delivery engine would use for a webhook: beside the webhook's
// event database in the data directory.
func archivePathFor(
t *testing.T,
mgr *database.WebhookDBManager,
webhookID string,
) string {
t.Helper()
return filepath.Join(
filepath.Dir(mgr.DBPath(webhookID)),
"archive-"+webhookID+".db",
)
}
// writeArchivePlaceholder creates a stand-in archive file so a
// test can assert the file survives webhook deletion.
func writeArchivePlaceholder(path string) error {
return os.WriteFile(path, []byte("archive"), 0o600)
}
// postRequest builds an authenticated POST request carrying the
// given chi URL parameters.
func postRequest(
path string,
cookies []*http.Cookie,
params map[string]string,
) *http.Request {
req := httptest.NewRequestWithContext(
context.Background(), http.MethodPost, path, nil,
)
for _, c := range cookies {
req.AddCookie(c)
}
rctx := chi.NewRouteContext()
for k, v := range params {
rctx.URLParams.Add(k, v)
}
return req.WithContext(
context.WithValue(req.Context(), chi.RouteCtxKey, rctx),
)
}
// TestHandleSourceDelete_EvictsArchiveWriter proves that
// deleting a webhook reaches the delivery engine and releases
// the webhook's archive writer, exercised through the real
// deletion handler rather than by calling the evictor directly.
func TestHandleSourceDelete_EvictsArchiveWriter(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
ev *recordingEvictor
)
app := newTestApp(t, &h, &sess, &db, &ev)
app.RequireStart()
t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db)
seedTarget(t, db, wh.ID, database.TargetTypeDatabase)
cookies := authenticatedCookies(
t, sess, deleteTestUserID, deleteTestUsername,
)
req := postRequest(
"/source/"+wh.ID+"/delete",
cookies,
map[string]string{paramSourceID: wh.ID},
)
w := httptest.NewRecorder()
h.HandleSourceDelete().ServeHTTP(w, req)
require.Equal(t, http.StatusSeeOther, w.Code)
assert.Equal(
t, []string{wh.ID}, ev.Evicted(),
"deleting a webhook should evict its archive writer",
)
}
// TestHandleSourceDelete_KeepsArchiveFile proves that deleting
// a webhook does not remove its archive database file: the
// archive is long-term storage the operator owns.
func TestHandleSourceDelete_KeepsArchiveFile(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
mgr *database.WebhookDBManager
)
app := newTestApp(t, &h, &sess, &db, &mgr)
app.RequireStart()
t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db)
// Place an archive file where the delivery engine would.
archivePath := archivePathFor(t, mgr, wh.ID)
require.NoError(
t,
writeArchivePlaceholder(archivePath),
)
cookies := authenticatedCookies(
t, sess, deleteTestUserID, deleteTestUsername,
)
req := postRequest(
"/source/"+wh.ID+"/delete",
cookies,
map[string]string{paramSourceID: wh.ID},
)
w := httptest.NewRecorder()
h.HandleSourceDelete().ServeHTTP(w, req)
require.Equal(t, http.StatusSeeOther, w.Code)
assert.FileExists(
t, archivePath,
"webhook deletion must not destroy the archive file",
)
}
// TestHandleTargetDelete_EvictsWhenLastDatabaseTargetGone
// proves that removing the last database target releases the
// archive writer.
func TestHandleTargetDelete_EvictsWhenLastDatabaseTargetGone(
t *testing.T,
) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
ev *recordingEvictor
)
app := newTestApp(t, &h, &sess, &db, &ev)
app.RequireStart()
t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db)
tgt := seedTarget(
t, db, wh.ID, database.TargetTypeDatabase,
)
cookies := authenticatedCookies(
t, sess, deleteTestUserID, deleteTestUsername,
)
req := postRequest(
"/source/"+wh.ID+"/targets/"+tgt.ID+"/delete",
cookies,
map[string]string{
paramSourceID: wh.ID,
paramTargetID: tgt.ID,
},
)
w := httptest.NewRecorder()
h.HandleTargetDelete().ServeHTTP(w, req)
require.Equal(t, http.StatusSeeOther, w.Code)
assert.Equal(
t, []string{wh.ID}, ev.Evicted(),
"removing the last database target should evict",
)
}
// TestHandleTargetDelete_KeepsWriterWhenDatabaseTargetRemains
// proves that deleting one of several database targets leaves
// the still-needed archive writer alone: the surviving target
// keeps archiving to the same file, so the writer must stay.
func TestHandleTargetDelete_KeepsWriterWhenDatabaseTargetRemains(
t *testing.T,
) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
ev *recordingEvictor
)
app := newTestApp(t, &h, &sess, &db, &ev)
app.RequireStart()
t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db)
doomed := seedTarget(
t, db, wh.ID, database.TargetTypeDatabase,
)
seedTarget(t, db, wh.ID, database.TargetTypeDatabase)
cookies := authenticatedCookies(
t, sess, deleteTestUserID, deleteTestUsername,
)
req := postRequest(
"/source/"+wh.ID+"/targets/"+doomed.ID+"/delete",
cookies,
map[string]string{
paramSourceID: wh.ID,
paramTargetID: doomed.ID,
},
)
w := httptest.NewRecorder()
h.HandleTargetDelete().ServeHTTP(w, req)
require.Equal(t, http.StatusSeeOther, w.Code)
assert.Empty(
t, ev.Evicted(),
"a second database target still needs the writer",
)
}
// TestHandleTargetDelete_KeepsWriterWhileDatabaseTargetRemains
// proves that deleting an unrelated target type leaves a
// still-needed archive writer alone.
func TestHandleTargetDelete_KeepsWriterWhileDatabaseTargetRemains(
t *testing.T,
) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
ev *recordingEvictor
)
app := newTestApp(t, &h, &sess, &db, &ev)
app.RequireStart()
t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db)
seedTarget(t, db, wh.ID, database.TargetTypeDatabase)
other := seedTarget(t, db, wh.ID, database.TargetTypeLog)
cookies := authenticatedCookies(
t, sess, deleteTestUserID, deleteTestUsername,
)
req := postRequest(
"/source/"+wh.ID+"/targets/"+other.ID+"/delete",
cookies,
map[string]string{
paramSourceID: wh.ID,
paramTargetID: other.ID,
},
)
w := httptest.NewRecorder()
h.HandleTargetDelete().ServeHTTP(w, req)
require.Equal(t, http.StatusSeeOther, w.Code)
assert.Empty(
t, ev.Evicted(),
"a surviving database target must keep its writer",
)
}

View File

@@ -533,13 +533,6 @@ func (h *Handlers) deleteWebhookResources(
return
}
// Release the delivery engine's per-webhook archiving state
// so a deleted webhook's archive writer (and any handle open
// within its debounce window) does not linger for the
// process lifetime. The archive file itself is deliberately
// left on disk; see evictArchiveWriter.
h.evictArchiveWriter(webhook.ID)
err = h.dbMgr.DeleteDB(webhook.ID)
if err != nil {
h.log.Error(
@@ -558,64 +551,6 @@ func (h *Handlers) deleteWebhookResources(
http.Redirect(w, r, "/sources", http.StatusSeeOther)
}
// evictArchiveWriter asks the delivery engine to drop its
// cached archive writer for a webhook, closing the archive file
// handle.
//
// The archive database file is NOT deleted. Unlike the event
// database — which is per-webhook working storage and is
// hard-deleted with the webhook — an archive is explicitly
// long-term storage that an operator may want to keep or move
// away for offline retention. Destroying it as a side effect of
// deleting a webhook would be a surprising and unrecoverable
// data loss, so the file is left for the operator to handle.
func (h *Handlers) evictArchiveWriter(webhookID string) {
if h.evictor == nil {
return
}
h.evictor.EvictWebhook(webhookID)
}
// evictArchiveWriterIfUnused releases a webhook's archive
// writer once the webhook has no database target left to feed
// it.
//
// It is called after any child resource of a webhook is
// deleted, and is correct without knowing which kind was: it
// evicts only when no database target remains, so deleting one
// of several database targets — or deleting an unrelated
// target type — leaves a still-needed writer alone. When no
// database target ever existed there is no writer and eviction
// is a no-op. Soft-deleted targets are excluded by GORM's
// default scope, so the row just deleted is not counted.
func (h *Handlers) evictArchiveWriterIfUnused(webhookID string) {
var remaining int64
err := h.db.DB().
Model(&database.Target{}).
Where(
"webhook_id = ? AND type = ?",
webhookID, database.TargetTypeDatabase,
).
Count(&remaining).Error
if err != nil {
h.log.Error(
"failed to count remaining database targets",
"webhook_id", webhookID,
"error", err,
)
return
}
if remaining > 0 {
return
}
h.evictArchiveWriter(webhookID)
}
// HandleSourceLogs shows the request/response logs for a
// webhook.
func (h *Handlers) HandleSourceLogs() http.HandlerFunc {
@@ -1089,31 +1024,23 @@ func (h *Handlers) HandleEntrypointDelete() http.HandlerFunc {
return h.deleteChildResource(
"entrypointID", &database.Entrypoint{},
"failed to delete entrypoint",
nil,
)
}
// HandleTargetDelete handles deleting a target. Deleting the
// last database target of a webhook leaves its archive writer
// with nothing to write, so the writer is evicted and its
// handle closed; the archive file is left on disk.
// HandleTargetDelete handles deleting a target.
func (h *Handlers) HandleTargetDelete() http.HandlerFunc {
return h.deleteChildResource(
"targetID", &database.Target{},
"failed to delete target",
h.evictArchiveWriterIfUnused,
)
}
// deleteChildResource returns a handler that deletes a child
// resource (entrypoint or target) belonging to a webhook. The
// optional afterDelete hook runs with the webhook's id once the
// delete has succeeded, before the redirect.
// resource (entrypoint or target) belonging to a webhook.
func (h *Handlers) deleteChildResource(
idParam string,
model any,
errMsg string,
afterDelete func(webhookID string),
) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
userID, ok := h.getUserID(r)
@@ -1153,10 +1080,6 @@ func (h *Handlers) deleteChildResource(
return
}
if afterDelete != nil {
afterDelete(webhook.ID)
}
http.Redirect(
w, r,
"/source/"+webhook.ID,