Bound shutdown hooks by their stop context (closes #102)
All checks were successful
check / check (push) Successful in 3m45s
All checks were successful
check / check (push) Successful in 3m45s
fx hands OnStop a context carrying the application's stop timeout, and the delivery engine, the retention reaper, and the archive sweeper all discarded it and called wg.Wait() bare. A worker wedged inside a delivery target that never returns, or a sweep blocked on a locked SQLite database, hung the process forever instead of letting it exit when the timeout expired. All three now wait through internal/lifecycle.WaitForShutdown, which selects the drained WaitGroup against the stop context and, on timeout, logs at error naming the component and returns an error rather than reporting a clean stop. Engine.stop also gains the cancel != nil guard its two mirrored components already had.
This commit is contained in:
@@ -10,6 +10,7 @@ import (
|
||||
"go.uber.org/fx"
|
||||
"sneak.berlin/go/webhooker/internal/config"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
"sneak.berlin/go/webhooker/internal/lifecycle"
|
||||
"sneak.berlin/go/webhooker/internal/logger"
|
||||
)
|
||||
|
||||
@@ -67,10 +68,9 @@ func NewArchiveSweeper(
|
||||
}
|
||||
|
||||
// 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.
|
||||
// lifecycle. The start hook's context is deliberately ignored
|
||||
// (see start for why the background loop must not inherit it);
|
||||
// the stop hook's context is honoured (see stop).
|
||||
func (s *ArchiveSweeper) registerHooks(lc fx.Lifecycle) {
|
||||
lc.Append(fx.Hook{
|
||||
//nolint:contextcheck // Not passing the hook context is
|
||||
@@ -80,10 +80,8 @@ func (s *ArchiveSweeper) registerHooks(lc fx.Lifecycle) {
|
||||
|
||||
return nil
|
||||
},
|
||||
OnStop: func(_ context.Context) error {
|
||||
s.stop()
|
||||
|
||||
return nil
|
||||
OnStop: func(ctx context.Context) error {
|
||||
return s.stop(ctx)
|
||||
},
|
||||
})
|
||||
}
|
||||
@@ -113,15 +111,27 @@ func (s *ArchiveSweeper) start() {
|
||||
)
|
||||
}
|
||||
|
||||
func (s *ArchiveSweeper) stop() {
|
||||
// stop cancels the sweep loop's context and waits for it to
|
||||
// exit, bounded by the stop hook's context: a prune wedged on a
|
||||
// locked archive must not hang the process past fx's stop
|
||||
// timeout.
|
||||
func (s *ArchiveSweeper) stop(ctx context.Context) error {
|
||||
s.log.Info("archive sweeper stopping")
|
||||
|
||||
if s.cancel != nil {
|
||||
s.cancel()
|
||||
}
|
||||
|
||||
s.wg.Wait()
|
||||
err := lifecycle.WaitForShutdown(
|
||||
ctx, s.log, "archive sweeper", &s.wg,
|
||||
)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
s.log.Info("archive sweeper stopped")
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *ArchiveSweeper) run(ctx context.Context) {
|
||||
|
||||
@@ -14,7 +14,6 @@ import (
|
||||
"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"
|
||||
@@ -226,17 +225,6 @@ func countArchivedRows(path string) (int64, error) {
|
||||
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
|
||||
@@ -270,7 +258,7 @@ func TestArchiveSweeper_LoopOutlivesStartHookContext(
|
||||
|
||||
// Drive the genuine fx hooks the application registers,
|
||||
// rather than a test-only entry point.
|
||||
lc := &captureLifecycle{}
|
||||
lc := &recordingLifecycle{}
|
||||
env.sweeper.ExportRegisterHooks(lc)
|
||||
require.Len(t, lc.hooks, 1)
|
||||
|
||||
@@ -924,7 +912,36 @@ func TestArchiveSweeper_StopsCleanly(t *testing.T) {
|
||||
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()
|
||||
// stop blocks on the loop's WaitGroup, so returning without
|
||||
// error proves the loop observed the cancellation and exited
|
||||
// well inside the stop context.
|
||||
require.NoError(
|
||||
t, env.sweeper.ExportStop(context.Background()),
|
||||
)
|
||||
}
|
||||
|
||||
// TestArchiveSweeper_StopHookHonoursStopTimeout is the sweeper's
|
||||
// half of the same shutdown defect the engine and the retention
|
||||
// reaper carried: an OnStop that discards its context and waits
|
||||
// on the WaitGroup bare hangs the process forever on a prune
|
||||
// wedged inside a locked archive.
|
||||
func TestArchiveSweeper_StopHookHonoursStopTimeout(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
env := setupSweeperTest(t)
|
||||
|
||||
lc := &recordingLifecycle{}
|
||||
env.sweeper.ExportRegisterHooks(lc)
|
||||
require.Len(t, lc.hooks, 1)
|
||||
require.NoError(t, lc.hooks[0].OnStart(context.Background()))
|
||||
|
||||
release := make(chan struct{})
|
||||
|
||||
t.Cleanup(func() { close(release) })
|
||||
|
||||
env.sweeper.ExportWedgeLoop(release)
|
||||
|
||||
requireStopHookExpires(t, lc.hooks[0], "archive sweeper")
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
"go.uber.org/fx"
|
||||
"gorm.io/gorm"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
"sneak.berlin/go/webhooker/internal/lifecycle"
|
||||
"sneak.berlin/go/webhooker/internal/logger"
|
||||
)
|
||||
|
||||
@@ -234,8 +235,9 @@ func (e *Engine) ScheduleRetry(
|
||||
}
|
||||
|
||||
// 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.
|
||||
// lifecycle. The start hook's context is deliberately ignored
|
||||
// (see start for why the worker pool must not inherit it); the
|
||||
// stop hook's context is honoured (see stop).
|
||||
func (e *Engine) registerHooks(lc fx.Lifecycle) {
|
||||
lc.Append(fx.Hook{
|
||||
//nolint:contextcheck // Not inheriting the hook context
|
||||
@@ -245,10 +247,8 @@ func (e *Engine) registerHooks(lc fx.Lifecycle) {
|
||||
|
||||
return nil
|
||||
},
|
||||
OnStop: func(_ context.Context) error {
|
||||
e.stop()
|
||||
|
||||
return nil
|
||||
OnStop: func(ctx context.Context) error {
|
||||
return e.stop(ctx)
|
||||
},
|
||||
})
|
||||
}
|
||||
@@ -289,11 +289,26 @@ func (e *Engine) start() {
|
||||
)
|
||||
}
|
||||
|
||||
func (e *Engine) stop() {
|
||||
// stop cancels the worker pool's context and waits for the pool
|
||||
// to drain, bounded by the stop hook's context: a wedged worker
|
||||
// must not hang the process past fx's stop timeout.
|
||||
func (e *Engine) stop(ctx context.Context) error {
|
||||
e.log.Info("delivery engine stopping")
|
||||
e.cancel()
|
||||
e.wg.Wait()
|
||||
|
||||
if e.cancel != nil {
|
||||
e.cancel()
|
||||
}
|
||||
|
||||
err := lifecycle.WaitForShutdown(
|
||||
ctx, e.log, "delivery engine", &e.wg,
|
||||
)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
e.log.Info("delivery engine stopped")
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (e *Engine) worker(ctx context.Context) {
|
||||
|
||||
@@ -501,7 +501,7 @@ func TestWorkerLifecycle_StartStop(t *testing.T) {
|
||||
|
||||
iWaitForDelivered(t, s.WebhookDB, d.ID)
|
||||
|
||||
s.Engine.ExportStop()
|
||||
require.NoError(t, s.Engine.ExportStop(context.Background()))
|
||||
}
|
||||
|
||||
// iWaitForDelivered polls until the delivery reaches the
|
||||
@@ -567,7 +567,7 @@ func TestWorkerLifecycle_ProcessesRetryChannel(
|
||||
|
||||
iWaitForDelivered(t, s.WebhookDB, d.ID)
|
||||
|
||||
s.Engine.ExportStop()
|
||||
require.NoError(t, s.Engine.ExportStop(context.Background()))
|
||||
}
|
||||
|
||||
// --- processDelivery: unknown target type ---
|
||||
|
||||
@@ -27,6 +27,13 @@ const (
|
||||
// and a ready deliveryCh are chosen between at random and a
|
||||
// doomed pool still delivers.
|
||||
hookSettleDelay = 250 * time.Millisecond
|
||||
|
||||
// wedgeStopTimeout is the stop timeout a wedged-shutdown test
|
||||
// hands OnStop, standing in for fx's StopTimeout. The test
|
||||
// asserts only that the hook returns at all, and allows it
|
||||
// hookStopTimeout — forty times this budget — to do so, so no
|
||||
// assertion here races the wall clock.
|
||||
wedgeStopTimeout = 250 * time.Millisecond
|
||||
)
|
||||
|
||||
// recordingLifecycle is a minimal fx.Lifecycle that records the
|
||||
@@ -40,6 +47,44 @@ func (l *recordingLifecycle) Append(h fx.Hook) {
|
||||
l.hooks = append(l.hooks, h)
|
||||
}
|
||||
|
||||
// requireStopHookExpires drives hook.OnStop with a stop context
|
||||
// that expires while a wedged goroutine is still running, and
|
||||
// requires the hook to return the deadline error naming
|
||||
// component instead of blocking on the WaitGroup forever.
|
||||
func requireStopHookExpires(
|
||||
t *testing.T, hook fx.Hook, component string,
|
||||
) {
|
||||
t.Helper()
|
||||
|
||||
stopCtx, cancel := context.WithTimeout(
|
||||
context.Background(), wedgeStopTimeout,
|
||||
)
|
||||
defer cancel()
|
||||
|
||||
var stopErr error
|
||||
|
||||
stopped := make(chan struct{})
|
||||
|
||||
go func() {
|
||||
defer close(stopped)
|
||||
|
||||
stopErr = hook.OnStop(stopCtx)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-stopped:
|
||||
case <-time.After(hookStopTimeout):
|
||||
t.Fatal(
|
||||
"OnStop did not return: it discarded the stop " +
|
||||
"context and is waiting on a wedged goroutine " +
|
||||
"that will never observe cancellation",
|
||||
)
|
||||
}
|
||||
|
||||
require.ErrorIs(t, stopErr, context.DeadlineExceeded)
|
||||
require.ErrorContains(t, stopErr, component)
|
||||
}
|
||||
|
||||
// startEngineViaHook drives the genuine fx hooks the application
|
||||
// registers for the engine, handing OnStart a context that is
|
||||
// already done, and returns only once a pool that inherited that
|
||||
@@ -197,3 +242,30 @@ func TestEngine_StopHookStopsWorkers(t *testing.T) {
|
||||
"a stopped engine must not deliver anything",
|
||||
)
|
||||
}
|
||||
|
||||
// TestEngine_StopHookHonoursStopTimeout is the regression test
|
||||
// for a shutdown that could never complete. fx hands OnStop a
|
||||
// context carrying the application's stop timeout; an OnStop
|
||||
// that discards it and calls wg.Wait() bare hangs the process
|
||||
// forever on a single worker stuck inside a delivery target that
|
||||
// never returns — precisely when a bounded shutdown matters
|
||||
// most.
|
||||
//
|
||||
// The wedged goroutine here never observes cancellation, so the
|
||||
// hook can only return by honouring its context, and it must say
|
||||
// so rather than reporting a clean stop.
|
||||
func TestEngine_StopHookHonoursStopTimeout(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
s := newISetup(t)
|
||||
|
||||
lc := startEngineViaHook(t, s.Engine)
|
||||
|
||||
release := make(chan struct{})
|
||||
|
||||
t.Cleanup(func() { close(release) })
|
||||
|
||||
s.Engine.ExportWedgeWorker(release)
|
||||
|
||||
requireStopHookExpires(t, lc.hooks[0], "delivery engine")
|
||||
}
|
||||
|
||||
@@ -216,8 +216,19 @@ func (e *Engine) ExportRegisterHooks(lc fx.Lifecycle) {
|
||||
}
|
||||
|
||||
// ExportStop exposes stop for testing.
|
||||
func (e *Engine) ExportStop() {
|
||||
e.stop()
|
||||
func (e *Engine) ExportStop(ctx context.Context) error {
|
||||
return e.stop(ctx)
|
||||
}
|
||||
|
||||
// ExportWedgeWorker adds a goroutine to the engine's WaitGroup
|
||||
// that never observes cancellation and returns only when release
|
||||
// is closed. It stands in for a worker stuck inside a delivery
|
||||
// target that never returns, which is the only way stop can be
|
||||
// made to outlast its context.
|
||||
func (e *Engine) ExportWedgeWorker(release <-chan struct{}) {
|
||||
e.wg.Go(func() {
|
||||
<-release
|
||||
})
|
||||
}
|
||||
|
||||
// ExportDeliveryCh returns the delivery channel.
|
||||
@@ -518,8 +529,19 @@ func (s *ArchiveSweeper) ExportRegisterHooks(lc fx.Lifecycle) {
|
||||
}
|
||||
|
||||
// ExportStop stops the sweeper's background loop for tests.
|
||||
func (s *ArchiveSweeper) ExportStop() {
|
||||
s.stop()
|
||||
func (s *ArchiveSweeper) ExportStop(ctx context.Context) error {
|
||||
return s.stop(ctx)
|
||||
}
|
||||
|
||||
// ExportWedgeLoop adds a goroutine to the sweeper's WaitGroup
|
||||
// that never observes cancellation and returns only when release
|
||||
// is closed. It stands in for a prune stuck on a locked archive.
|
||||
func (s *ArchiveSweeper) ExportWedgeLoop(
|
||||
release <-chan struct{},
|
||||
) {
|
||||
s.wg.Go(func() {
|
||||
<-release
|
||||
})
|
||||
}
|
||||
|
||||
// ExportSetInterval overrides the sweep interval for tests.
|
||||
|
||||
Reference in New Issue
Block a user