diff --git a/internal/database/export_test.go b/internal/database/export_test.go index 29321fe..fdf8804 100644 --- a/internal/database/export_test.go +++ b/internal/database/export_test.go @@ -5,6 +5,8 @@ import ( "log/slog" "os" "time" + + "go.uber.org/fx" ) // NewTestRetentionReaper builds a RetentionReaper backed by the given @@ -29,3 +31,26 @@ func NewTestRetentionReaper( func (r *RetentionReaper) ExportSweep(ctx context.Context) { r.sweep(ctx) } + +// ExportRegisterHooks registers the reaper'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 (r *RetentionReaper) ExportRegisterHooks(lc fx.Lifecycle) { + r.registerHooks(lc) +} + +// ExportStart starts the reaper's background loop for tests. +func (r *RetentionReaper) ExportStart() { + r.start() +} + +// ExportStop stops the reaper's background loop for tests. +func (r *RetentionReaper) ExportStop() { + r.stop() +} + +// ExportSetInterval overrides the sweep interval for tests. +func (r *RetentionReaper) ExportSetInterval(d time.Duration) { + r.interval = d +} diff --git a/internal/database/retention.go b/internal/database/retention.go index 23d516f..523817c 100644 --- a/internal/database/retention.go +++ b/internal/database/retention.go @@ -56,9 +56,20 @@ func NewRetentionReaper( interval: params.Config.RetentionSweepInterval, } + r.registerHooks(lc) + + return r +} + +// registerHooks wires the reaper's start and stop into the fx +// lifecycle. The start hook's context is deliberately ignored: see +// start for why the sweep loop must not inherit it. +func (r *RetentionReaper) registerHooks(lc fx.Lifecycle) { lc.Append(fx.Hook{ - OnStart: func(ctx context.Context) error { - r.start(ctx) + //nolint:contextcheck // Not inheriting the hook context is + // the point: see start. + OnStart: func(_ context.Context) error { + r.start() return nil }, @@ -68,12 +79,20 @@ func NewRetentionReaper( return nil }, }) - - return r } -func (r *RetentionReaper) start(ctx context.Context) { - ctx, cancel := context.WithCancel(ctx) +// 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) and is cancelled once the start phase +// completes, so a loop derived from it dies 45 minutes before its +// first tick under the default one-hour sweep interval, leaving a +// reaper that never reaps. 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 (r *RetentionReaper) start() { + ctx, cancel := context.WithCancel(context.Background()) r.cancel = cancel r.wg.Add(1) diff --git a/internal/database/retention_lifecycle_test.go b/internal/database/retention_lifecycle_test.go new file mode 100644 index 0000000..47e8d80 --- /dev/null +++ b/internal/database/retention_lifecycle_test.go @@ -0,0 +1,209 @@ +package database_test + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.uber.org/fx" + "gorm.io/gorm" + "sneak.berlin/go/webhooker/internal/database" +) + +const ( + // reaperTestInterval is the sweep interval a lifecycle test + // runs the reaper at, so a loop that survives startup produces + // an observable sweep quickly. + reaperTestInterval = 10 * time.Millisecond + + // reaperStopTimeout bounds how long a lifecycle test waits for + // the reaper's OnStop hook to return before declaring the + // shutdown hung. + reaperStopTimeout = 10 * time.Second + + // reaperTestRetentionDays is the retention policy the lifecycle + // tests give their webhook. + reaperTestRetentionDays = 30 +) + +// recordingLifecycle 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 recordingLifecycle struct { + hooks []fx.Hook +} + +func (l *recordingLifecycle) Append(h fx.Hook) { + l.hooks = append(l.hooks, h) +} + +// startReaperViaHook drives the genuine fx hooks the application +// registers for the reaper, handing OnStart a context that is +// already done. It returns the recorded lifecycle so the caller +// can drive OnStop too. +func startReaperViaHook( + t *testing.T, r *database.RetentionReaper, +) *recordingLifecycle { + t.Helper() + + lc := &recordingLifecycle{} + r.ExportRegisterHooks(lc) + require.Len(t, lc.hooks, 1) + + // fx hands OnStart a context carrying the application start + // timeout, and cancels it when the start phase ends. An + // already-cancelled context is that same defect taken to its + // limit, and unlike a plain context.Background() it actually + // distinguishes a correctly rooted loop from a broken one. + hookCtx, cancel := context.WithCancel(context.Background()) + cancel() + + require.NoError(t, lc.hooks[0].OnStart(hookCtx)) + + return lc +} + +// eventGone reports whether an event row has been removed. It +// takes no *testing.T because it is polled from an +// assert.Eventually condition, which runs off the test goroutine +// where testify assertions must not be used. +func eventGone(db *gorm.DB, eventID string) bool { + var n int64 + + err := db.Unscoped().Model(&database.Event{}). + Where("id = ?", eventID).Count(&n).Error + if err != nil { + return false + } + + return n == 0 +} + +// seedExpiredWebhook creates a webhook with a finite retention +// policy plus one long-expired event chain, and returns the +// webhook's database and the chain's event ID. +func seedExpiredWebhook( + t *testing.T, env *retentionTestEnv, +) (*gorm.DB, string) { + t.Helper() + + webhookID := createWebhook( + t, env.mainDB.DB(), reaperTestRetentionDays, + ) + + db, err := env.mgr.GetDB(webhookID) + require.NoError(t, err) + + chain := seedEventChain( + t, db, webhookID, + time.Now().Add(-365*24*time.Hour), + ) + + return db, chain.eventID +} + +// TestRetentionReaper_LoopOutlivesStartHookContext is the +// regression test for a reaper that never reaped. fx calls +// OnStart with a context carrying the application's start timeout +// (15s by default) and cancels it when the start phase ends, so a +// sweep loop rooted in it is dead three quarters of an hour +// before its first tick under the default one-hour interval, and +// per-webhook event databases grow without bound exactly as they +// did before retention existed. +// +// Driving OnStart with an already-cancelled context is that +// defect taken to its limit: a loop that inherits the hook +// context never ticks once, while a correctly rooted loop keeps +// sweeping for as long as the process lives. +func TestRetentionReaper_LoopOutlivesStartHookContext( + t *testing.T, +) { + t.Parallel() + + env := setupRetentionTest(t) + + db, eventID := seedExpiredWebhook(t, env) + + env.reaper.ExportSetInterval(reaperTestInterval) + + lc := startReaperViaHook(t, env.reaper) + t.Cleanup(func() { + _ = lc.hooks[0].OnStop(context.Background()) + }) + + assert.Eventually( + t, + func() bool { return eventGone(db, eventID) }, + 5*time.Second, + reaperTestInterval, + "the sweep loop must keep running after the start "+ + "hook's context is done; it reaped nothing, so it "+ + "inherited the hook context and died", + ) +} + +// TestRetentionReaper_StopHookStopsLoop proves the fix did not +// trade a startup bug for a shutdown hang: now that the sweep +// loop no longer observes the start hook's cancellation, OnStop +// is the only thing that can stop it, and it must both return +// promptly and actually leave the loop stopped. +func TestRetentionReaper_StopHookStopsLoop(t *testing.T) { + t.Parallel() + + env := setupRetentionTest(t) + + db, eventID := seedExpiredWebhook(t, env) + + env.reaper.ExportSetInterval(reaperTestInterval) + + lc := startReaperViaHook(t, env.reaper) + + // Let the loop prove it is running before stopping it, so a + // fast OnStop cannot pass by stopping something already dead. + require.Eventually( + t, + func() bool { return eventGone(db, eventID) }, + 5*time.Second, + reaperTestInterval, + ) + + var stopErr error + + stopped := make(chan struct{}) + + go func() { + defer close(stopped) + + // stop blocks on the loop's WaitGroup, so returning at all + // proves the goroutine observed the cancellation. + stopErr = lc.hooks[0].OnStop(context.Background()) + }() + + select { + case <-stopped: + case <-time.After(reaperStopTimeout): + t.Fatal( + "OnStop did not return: the retention reaper's " + + "WaitGroup is still waiting on a loop that never " + + "observed cancellation", + ) + } + + require.NoError(t, stopErr) + + // With the loop gone, a newly expired chain must survive. + survivor := seedEventChain( + t, db, "stopped-webhook", + time.Now().Add(-365*24*time.Hour), + ) + + time.Sleep(20 * reaperTestInterval) + + assert.False( + t, + eventGone(db, survivor.eventID), + "a stopped reaper must not sweep anything", + ) +} diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index 6011558..5f1a7f3 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -149,9 +149,20 @@ func New( Transport: NewSSRFSafeTransport(), }) + e.registerHooks(lc) + + return e +} + +// registerHooks wires the engine's start and stop into the fx +// lifecycle. The start hook's context is deliberately ignored: +// see start for why the worker pool must not inherit it. +func (e *Engine) registerHooks(lc fx.Lifecycle) { lc.Append(fx.Hook{ - OnStart: func(ctx context.Context) error { - e.start(ctx) + //nolint:contextcheck // Not inheriting the hook context + // is the point: see start. + OnStart: func(_ context.Context) error { + e.start() return nil }, @@ -161,8 +172,6 @@ func New( return nil }, }) - - return e } // Notify signals the delivery engine that new deliveries @@ -210,8 +219,20 @@ func (e *Engine) ScheduleRetry( }) } -func (e *Engine) start(ctx context.Context) { - ctx, cancel := context.WithCancel(ctx) +// start launches the worker pool, restart recovery, and the +// periodic retry sweep. +// +// Their context is derived from context.Background(), NOT from +// the fx OnStart hook context. The hook context carries fx's +// start timeout (15s by default) and is cancelled once the start +// phase completes, so goroutines derived from it stop a few +// seconds into the process: every worker would return and the +// engine would silently stop delivering webhooks entirely. 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 (e *Engine) start() { + ctx, cancel := context.WithCancel(context.Background()) e.cancel = cancel for range e.workers { diff --git a/internal/delivery/engine_integration_test.go b/internal/delivery/engine_integration_test.go index 6b07c43..f4ad67e 100644 --- a/internal/delivery/engine_integration_test.go +++ b/internal/delivery/engine_integration_test.go @@ -476,7 +476,7 @@ func TestWorkerLifecycle_StartStop(t *testing.T) { t.Parallel() s := newISetup(t) - s.Engine.ExportStart(context.Background()) + s.Engine.ExportStart() event := iSeedEvent( t, s.WebhookDB, s.WebhookID, @@ -558,7 +558,7 @@ func TestWorkerLifecycle_ProcessesRetryChannel( database.DeliveryStatusRetrying, ) - s.Engine.ExportStart(context.Background()) + s.Engine.ExportStart() bodyStr := event.Body cfg := iHTTPConfig(ts.URL) diff --git a/internal/delivery/engine_lifecycle_test.go b/internal/delivery/engine_lifecycle_test.go new file mode 100644 index 0000000..c2d9735 --- /dev/null +++ b/internal/delivery/engine_lifecycle_test.go @@ -0,0 +1,183 @@ +package delivery_test + +import ( + "context" + "testing" + "time" + + "github.com/google/uuid" + "github.com/stretchr/testify/require" + "go.uber.org/fx" + "sneak.berlin/go/webhooker/internal/database" + "sneak.berlin/go/webhooker/internal/delivery" +) + +// hookStopTimeout bounds how long a lifecycle test waits for the +// engine's OnStop hook to return before declaring the shutdown +// hung. +const hookStopTimeout = 10 * time.Second + +// recordingLifecycle 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 recordingLifecycle struct { + hooks []fx.Hook +} + +func (l *recordingLifecycle) Append(h fx.Hook) { + l.hooks = append(l.hooks, h) +} + +// startEngineViaHook drives the genuine fx hooks the application +// registers for the engine, handing OnStart a context that is +// already done. It returns the recorded lifecycle so the caller +// can drive OnStop too. +func startEngineViaHook( + t *testing.T, eng *delivery.Engine, +) *recordingLifecycle { + t.Helper() + + lc := &recordingLifecycle{} + eng.ExportRegisterHooks(lc) + require.Len(t, lc.hooks, 1) + + // fx hands OnStart a context carrying the application start + // timeout, and cancels it when the start phase ends. An + // already-cancelled context is that same defect taken to its + // limit, and unlike a plain context.Background() it actually + // distinguishes a correctly rooted loop from a broken one. + hookCtx, cancel := context.WithCancel(context.Background()) + cancel() + + require.NoError(t, lc.hooks[0].OnStart(hookCtx)) + + return lc +} + +// seedLogTask seeds a pending delivery for a log target and +// returns its ID together with the task that drives it. The log +// target needs no network, so a delivery completing proves only +// that a worker picked the task up. +func seedLogTask( + t *testing.T, s iSetup, +) (string, delivery.Task) { + t.Helper() + + event := iSeedEvent( + t, s.WebhookDB, s.WebhookID, + `{"lifecycle":"hook-context"}`, + ) + targetID := uuid.New().String() + + d := iSeedDelivery( + t, s.WebhookDB, event.ID, targetID, + database.DeliveryStatusPending, + ) + + bodyStr := event.Body + task := iTask( + d, event, s.WebhookID, targetID, + "hook-context-test", "", 0, 1, &bodyStr, + ) + task.TargetType = database.TargetTypeLog + + return d.ID, task +} + +// TestEngine_WorkersOutliveStartHookContext is the regression +// test for a delivery engine that stopped delivering roughly +// fifteen seconds after boot. fx calls OnStart with a context +// carrying the application's start timeout (15s by default) and +// cancels it when the start phase ends, so a worker pool rooted +// in it exits shortly after startup: the process keeps accepting +// and persisting events while nothing at all forwards them. +// +// Driving OnStart with an already-cancelled context is that +// defect taken to its limit. A pool that inherits the hook +// context never processes a single task; a correctly rooted pool +// keeps working for as long as the process lives. +func TestEngine_WorkersOutliveStartHookContext(t *testing.T) { + t.Parallel() + + s := newISetup(t) + + deliveryID, task := seedLogTask(t, s) + + lc := startEngineViaHook(t, s.Engine) + t.Cleanup(func() { + _ = lc.hooks[0].OnStop(context.Background()) + }) + + s.Engine.Notify([]delivery.Task{task}) + + iWaitForStatus( + t, s.WebhookDB, deliveryID, + database.DeliveryStatusDelivered, + ) +} + +// TestEngine_StopHookStopsWorkers proves the fix did not trade a +// startup bug for a shutdown hang: now that the worker pool no +// longer observes the start hook's cancellation, OnStop is the +// only thing that can stop it, and it must both return promptly +// and actually leave the pool drained. +func TestEngine_StopHookStopsWorkers(t *testing.T) { + t.Parallel() + + s := newISetup(t) + + lc := startEngineViaHook(t, s.Engine) + + // Let the pool prove it is running before stopping it, so a + // fast OnStop cannot pass by stopping something already dead. + firstID, firstTask := seedLogTask(t, s) + s.Engine.Notify([]delivery.Task{firstTask}) + iWaitForStatus( + t, s.WebhookDB, firstID, + database.DeliveryStatusDelivered, + ) + + var stopErr error + + stopped := make(chan struct{}) + + go func() { + defer close(stopped) + + // stop blocks on the workers' WaitGroup, so returning at + // all proves every goroutine observed the cancellation. + stopErr = lc.hooks[0].OnStop(context.Background()) + }() + + select { + case <-stopped: + case <-time.After(hookStopTimeout): + t.Fatal( + "OnStop did not return: the delivery engine's " + + "WaitGroup is still waiting on a goroutine that " + + "never observed cancellation", + ) + } + + require.NoError(t, stopErr) + + // With every worker gone, a freshly notified task must sit + // untouched in the queue rather than being delivered. + secondID, secondTask := seedLogTask(t, s) + s.Engine.Notify([]delivery.Task{secondTask}) + + time.Sleep(200 * time.Millisecond) + + var after database.Delivery + + require.NoError( + t, + s.WebhookDB.First(&after, "id = ?", secondID).Error, + ) + require.Equal( + t, + database.DeliveryStatusPending, + after.Status, + "a stopped engine must not deliver anything", + ) +} diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index 739eb03..f5b01c7 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -7,6 +7,7 @@ import ( "net/http" "time" + "go.uber.org/fx" "gorm.io/gorm" "sneak.berlin/go/webhooker/internal/database" ) @@ -189,8 +190,16 @@ func (e *Engine) ExportRecoverInFlight( } // ExportStart exposes start for testing. -func (e *Engine) ExportStart(ctx context.Context) { - e.start(ctx) +func (e *Engine) ExportStart() { + e.start() +} + +// ExportRegisterHooks registers the engine'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 (e *Engine) ExportRegisterHooks(lc fx.Lifecycle) { + e.registerHooks(lc) } // ExportStop exposes stop for testing.