Compare commits
1 Commits
issue-102-
...
issue-123-
| Author | SHA1 | Date | |
|---|---|---|---|
| 4a91635b2a |
@@ -46,19 +46,8 @@ func (r *RetentionReaper) ExportStart() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ExportStop stops the reaper's background loop for tests.
|
// ExportStop stops the reaper's background loop for tests.
|
||||||
func (r *RetentionReaper) ExportStop(ctx context.Context) error {
|
func (r *RetentionReaper) ExportStop() {
|
||||||
return r.stop(ctx)
|
r.stop()
|
||||||
}
|
|
||||||
|
|
||||||
// ExportWedgeLoop adds a goroutine to the reaper's WaitGroup that
|
|
||||||
// never observes cancellation and returns only when release is
|
|
||||||
// closed. It stands in for a sweep stuck on a locked database.
|
|
||||||
func (r *RetentionReaper) ExportWedgeLoop(
|
|
||||||
release <-chan struct{},
|
|
||||||
) {
|
|
||||||
r.wg.Go(func() {
|
|
||||||
<-release
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// ExportSetInterval overrides the sweep interval for tests.
|
// ExportSetInterval overrides the sweep interval for tests.
|
||||||
|
|||||||
@@ -10,7 +10,6 @@ import (
|
|||||||
"go.uber.org/fx"
|
"go.uber.org/fx"
|
||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
"sneak.berlin/go/webhooker/internal/config"
|
"sneak.berlin/go/webhooker/internal/config"
|
||||||
"sneak.berlin/go/webhooker/internal/lifecycle"
|
|
||||||
"sneak.berlin/go/webhooker/internal/logger"
|
"sneak.berlin/go/webhooker/internal/logger"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -63,9 +62,8 @@ func NewRetentionReaper(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// registerHooks wires the reaper's start and stop into the fx
|
// registerHooks wires the reaper's start and stop into the fx
|
||||||
// lifecycle. The start hook's context is deliberately ignored (see
|
// lifecycle. The start hook's context is deliberately ignored: see
|
||||||
// start for why the sweep loop must not inherit it); the stop hook's
|
// start for why the sweep loop must not inherit it.
|
||||||
// context is honoured (see stop).
|
|
||||||
func (r *RetentionReaper) registerHooks(lc fx.Lifecycle) {
|
func (r *RetentionReaper) registerHooks(lc fx.Lifecycle) {
|
||||||
lc.Append(fx.Hook{
|
lc.Append(fx.Hook{
|
||||||
//nolint:contextcheck // Not inheriting the hook context is
|
//nolint:contextcheck // Not inheriting the hook context is
|
||||||
@@ -75,8 +73,10 @@ func (r *RetentionReaper) registerHooks(lc fx.Lifecycle) {
|
|||||||
|
|
||||||
return nil
|
return nil
|
||||||
},
|
},
|
||||||
OnStop: func(ctx context.Context) error {
|
OnStop: func(_ context.Context) error {
|
||||||
return r.stop(ctx)
|
r.stop()
|
||||||
|
|
||||||
|
return nil
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -105,27 +105,15 @@ func (r *RetentionReaper) start() {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// stop cancels the sweep loop's context and waits for it to
|
func (r *RetentionReaper) stop() {
|
||||||
// exit, bounded by the stop hook's context: a sweep wedged on a
|
|
||||||
// locked database must not hang the process past fx's stop
|
|
||||||
// timeout.
|
|
||||||
func (r *RetentionReaper) stop(ctx context.Context) error {
|
|
||||||
r.log.Info("retention reaper stopping")
|
r.log.Info("retention reaper stopping")
|
||||||
|
|
||||||
if r.cancel != nil {
|
if r.cancel != nil {
|
||||||
r.cancel()
|
r.cancel()
|
||||||
}
|
}
|
||||||
|
|
||||||
err := lifecycle.WaitForShutdown(
|
r.wg.Wait()
|
||||||
ctx, r.log, "retention reaper", &r.wg,
|
|
||||||
)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
r.log.Info("retention reaper stopped")
|
r.log.Info("retention reaper stopped")
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r *RetentionReaper) run(ctx context.Context) {
|
func (r *RetentionReaper) run(ctx context.Context) {
|
||||||
|
|||||||
@@ -26,13 +26,6 @@ const (
|
|||||||
// reaperTestRetentionDays is the retention policy the lifecycle
|
// reaperTestRetentionDays is the retention policy the lifecycle
|
||||||
// tests give their webhook.
|
// tests give their webhook.
|
||||||
reaperTestRetentionDays = 30
|
reaperTestRetentionDays = 30
|
||||||
|
|
||||||
// reaperWedgeStopTimeout is the stop timeout the wedged-shutdown
|
|
||||||
// test hands OnStop, standing in for fx's StopTimeout. The test
|
|
||||||
// asserts only that the hook returns at all, and allows it
|
|
||||||
// reaperStopTimeout — forty times this budget — to do so, so no
|
|
||||||
// assertion races the wall clock.
|
|
||||||
reaperWedgeStopTimeout = 250 * time.Millisecond
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// recordingLifecycle is a minimal fx.Lifecycle that records the
|
// recordingLifecycle is a minimal fx.Lifecycle that records the
|
||||||
@@ -214,59 +207,3 @@ func TestRetentionReaper_StopHookStopsLoop(t *testing.T) {
|
|||||||
"a stopped reaper must not sweep anything",
|
"a stopped reaper must not sweep anything",
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestRetentionReaper_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 sweep blocked on a locked SQLite database
|
|
||||||
// — 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 TestRetentionReaper_StopHookHonoursStopTimeout(
|
|
||||||
t *testing.T,
|
|
||||||
) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
env := setupRetentionTest(t)
|
|
||||||
|
|
||||||
env.reaper.ExportSetInterval(reaperTestInterval)
|
|
||||||
|
|
||||||
lc := startReaperViaHook(t, env.reaper)
|
|
||||||
|
|
||||||
release := make(chan struct{})
|
|
||||||
|
|
||||||
t.Cleanup(func() { close(release) })
|
|
||||||
|
|
||||||
env.reaper.ExportWedgeLoop(release)
|
|
||||||
|
|
||||||
stopCtx, cancel := context.WithTimeout(
|
|
||||||
context.Background(), reaperWedgeStopTimeout,
|
|
||||||
)
|
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
var stopErr error
|
|
||||||
|
|
||||||
stopped := make(chan struct{})
|
|
||||||
|
|
||||||
go func() {
|
|
||||||
defer close(stopped)
|
|
||||||
|
|
||||||
stopErr = lc.hooks[0].OnStop(stopCtx)
|
|
||||||
}()
|
|
||||||
|
|
||||||
select {
|
|
||||||
case <-stopped:
|
|
||||||
case <-time.After(reaperStopTimeout):
|
|
||||||
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, "retention reaper")
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -10,7 +10,6 @@ import (
|
|||||||
"go.uber.org/fx"
|
"go.uber.org/fx"
|
||||||
"sneak.berlin/go/webhooker/internal/config"
|
"sneak.berlin/go/webhooker/internal/config"
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
"sneak.berlin/go/webhooker/internal/lifecycle"
|
|
||||||
"sneak.berlin/go/webhooker/internal/logger"
|
"sneak.berlin/go/webhooker/internal/logger"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -68,9 +67,10 @@ func NewArchiveSweeper(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// registerHooks wires the sweeper's start and stop into the fx
|
// registerHooks wires the sweeper's start and stop into the fx
|
||||||
// lifecycle. The start hook's context is deliberately ignored
|
// lifecycle. Both hook contexts are deliberately ignored: see
|
||||||
// (see start for why the background loop must not inherit it);
|
// start for why the background loop must not inherit the start
|
||||||
// the stop hook's context is honoured (see stop).
|
// 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) {
|
func (s *ArchiveSweeper) registerHooks(lc fx.Lifecycle) {
|
||||||
lc.Append(fx.Hook{
|
lc.Append(fx.Hook{
|
||||||
//nolint:contextcheck // Not passing the hook context is
|
//nolint:contextcheck // Not passing the hook context is
|
||||||
@@ -80,8 +80,10 @@ func (s *ArchiveSweeper) registerHooks(lc fx.Lifecycle) {
|
|||||||
|
|
||||||
return nil
|
return nil
|
||||||
},
|
},
|
||||||
OnStop: func(ctx context.Context) error {
|
OnStop: func(_ context.Context) error {
|
||||||
return s.stop(ctx)
|
s.stop()
|
||||||
|
|
||||||
|
return nil
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -111,27 +113,15 @@ func (s *ArchiveSweeper) start() {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// stop cancels the sweep loop's context and waits for it to
|
func (s *ArchiveSweeper) stop() {
|
||||||
// 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")
|
s.log.Info("archive sweeper stopping")
|
||||||
|
|
||||||
if s.cancel != nil {
|
if s.cancel != nil {
|
||||||
s.cancel()
|
s.cancel()
|
||||||
}
|
}
|
||||||
|
|
||||||
err := lifecycle.WaitForShutdown(
|
s.wg.Wait()
|
||||||
ctx, s.log, "archive sweeper", &s.wg,
|
|
||||||
)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
s.log.Info("archive sweeper stopped")
|
s.log.Info("archive sweeper stopped")
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *ArchiveSweeper) run(ctx context.Context) {
|
func (s *ArchiveSweeper) run(ctx context.Context) {
|
||||||
|
|||||||
@@ -14,6 +14,7 @@ import (
|
|||||||
"github.com/google/uuid"
|
"github.com/google/uuid"
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
|
"go.uber.org/fx"
|
||||||
"gorm.io/driver/sqlite"
|
"gorm.io/driver/sqlite"
|
||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
"gorm.io/gorm/clause"
|
"gorm.io/gorm/clause"
|
||||||
@@ -225,6 +226,17 @@ func countArchivedRows(path string) (int64, error) {
|
|||||||
return count, nil
|
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
|
// TestArchiveSweeper_LoopOutlivesStartHookContext is the
|
||||||
// regression test for a sweeper that never swept. fx calls
|
// regression test for a sweeper that never swept. fx calls
|
||||||
// OnStart with a context carrying the application's start
|
// OnStart with a context carrying the application's start
|
||||||
@@ -258,7 +270,7 @@ func TestArchiveSweeper_LoopOutlivesStartHookContext(
|
|||||||
|
|
||||||
// Drive the genuine fx hooks the application registers,
|
// Drive the genuine fx hooks the application registers,
|
||||||
// rather than a test-only entry point.
|
// rather than a test-only entry point.
|
||||||
lc := &recordingLifecycle{}
|
lc := &captureLifecycle{}
|
||||||
env.sweeper.ExportRegisterHooks(lc)
|
env.sweeper.ExportRegisterHooks(lc)
|
||||||
require.Len(t, lc.hooks, 1)
|
require.Len(t, lc.hooks, 1)
|
||||||
|
|
||||||
@@ -912,36 +924,7 @@ func TestArchiveSweeper_StopsCleanly(t *testing.T) {
|
|||||||
env.sweeper.ExportSetInterval(time.Millisecond)
|
env.sweeper.ExportSetInterval(time.Millisecond)
|
||||||
env.sweeper.ExportStart()
|
env.sweeper.ExportStart()
|
||||||
|
|
||||||
// stop blocks on the loop's WaitGroup, so returning without
|
// stop blocks on the loop's WaitGroup, so returning at all
|
||||||
// error proves the loop observed the cancellation and exited
|
// proves the loop observed the cancellation and exited.
|
||||||
// well inside the stop context.
|
env.sweeper.ExportStop()
|
||||||
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,7 +13,6 @@ import (
|
|||||||
"go.uber.org/fx"
|
"go.uber.org/fx"
|
||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
"sneak.berlin/go/webhooker/internal/database"
|
||||||
"sneak.berlin/go/webhooker/internal/lifecycle"
|
|
||||||
"sneak.berlin/go/webhooker/internal/logger"
|
"sneak.berlin/go/webhooker/internal/logger"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -235,9 +234,8 @@ func (e *Engine) ScheduleRetry(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// registerHooks wires the engine's start and stop into the fx
|
// registerHooks wires the engine's start and stop into the fx
|
||||||
// lifecycle. The start hook's context is deliberately ignored
|
// lifecycle. The start hook's context is deliberately ignored:
|
||||||
// (see start for why the worker pool must not inherit it); the
|
// see start for why the worker pool must not inherit it.
|
||||||
// stop hook's context is honoured (see stop).
|
|
||||||
func (e *Engine) registerHooks(lc fx.Lifecycle) {
|
func (e *Engine) registerHooks(lc fx.Lifecycle) {
|
||||||
lc.Append(fx.Hook{
|
lc.Append(fx.Hook{
|
||||||
//nolint:contextcheck // Not inheriting the hook context
|
//nolint:contextcheck // Not inheriting the hook context
|
||||||
@@ -247,8 +245,10 @@ func (e *Engine) registerHooks(lc fx.Lifecycle) {
|
|||||||
|
|
||||||
return nil
|
return nil
|
||||||
},
|
},
|
||||||
OnStop: func(ctx context.Context) error {
|
OnStop: func(_ context.Context) error {
|
||||||
return e.stop(ctx)
|
e.stop()
|
||||||
|
|
||||||
|
return nil
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -289,26 +289,11 @@ func (e *Engine) start() {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// stop cancels the worker pool's context and waits for the pool
|
func (e *Engine) stop() {
|
||||||
// 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.log.Info("delivery engine stopping")
|
||||||
|
e.cancel()
|
||||||
if e.cancel != nil {
|
e.wg.Wait()
|
||||||
e.cancel()
|
|
||||||
}
|
|
||||||
|
|
||||||
err := lifecycle.WaitForShutdown(
|
|
||||||
ctx, e.log, "delivery engine", &e.wg,
|
|
||||||
)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
e.log.Info("delivery engine stopped")
|
e.log.Info("delivery engine stopped")
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (e *Engine) worker(ctx context.Context) {
|
func (e *Engine) worker(ctx context.Context) {
|
||||||
|
|||||||
@@ -501,7 +501,7 @@ func TestWorkerLifecycle_StartStop(t *testing.T) {
|
|||||||
|
|
||||||
iWaitForDelivered(t, s.WebhookDB, d.ID)
|
iWaitForDelivered(t, s.WebhookDB, d.ID)
|
||||||
|
|
||||||
require.NoError(t, s.Engine.ExportStop(context.Background()))
|
s.Engine.ExportStop()
|
||||||
}
|
}
|
||||||
|
|
||||||
// iWaitForDelivered polls until the delivery reaches the
|
// iWaitForDelivered polls until the delivery reaches the
|
||||||
@@ -567,7 +567,7 @@ func TestWorkerLifecycle_ProcessesRetryChannel(
|
|||||||
|
|
||||||
iWaitForDelivered(t, s.WebhookDB, d.ID)
|
iWaitForDelivered(t, s.WebhookDB, d.ID)
|
||||||
|
|
||||||
require.NoError(t, s.Engine.ExportStop(context.Background()))
|
s.Engine.ExportStop()
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- processDelivery: unknown target type ---
|
// --- processDelivery: unknown target type ---
|
||||||
|
|||||||
@@ -27,13 +27,6 @@ const (
|
|||||||
// and a ready deliveryCh are chosen between at random and a
|
// and a ready deliveryCh are chosen between at random and a
|
||||||
// doomed pool still delivers.
|
// doomed pool still delivers.
|
||||||
hookSettleDelay = 250 * time.Millisecond
|
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
|
// recordingLifecycle is a minimal fx.Lifecycle that records the
|
||||||
@@ -47,44 +40,6 @@ func (l *recordingLifecycle) Append(h fx.Hook) {
|
|||||||
l.hooks = append(l.hooks, h)
|
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
|
// startEngineViaHook drives the genuine fx hooks the application
|
||||||
// registers for the engine, handing OnStart a context that is
|
// registers for the engine, handing OnStart a context that is
|
||||||
// already done, and returns only once a pool that inherited that
|
// already done, and returns only once a pool that inherited that
|
||||||
@@ -242,30 +197,3 @@ func TestEngine_StopHookStopsWorkers(t *testing.T) {
|
|||||||
"a stopped engine must not deliver anything",
|
"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,19 +216,8 @@ func (e *Engine) ExportRegisterHooks(lc fx.Lifecycle) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ExportStop exposes stop for testing.
|
// ExportStop exposes stop for testing.
|
||||||
func (e *Engine) ExportStop(ctx context.Context) error {
|
func (e *Engine) ExportStop() {
|
||||||
return e.stop(ctx)
|
e.stop()
|
||||||
}
|
|
||||||
|
|
||||||
// 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.
|
// ExportDeliveryCh returns the delivery channel.
|
||||||
@@ -529,19 +518,8 @@ func (s *ArchiveSweeper) ExportRegisterHooks(lc fx.Lifecycle) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ExportStop stops the sweeper's background loop for tests.
|
// ExportStop stops the sweeper's background loop for tests.
|
||||||
func (s *ArchiveSweeper) ExportStop(ctx context.Context) error {
|
func (s *ArchiveSweeper) ExportStop() {
|
||||||
return s.stop(ctx)
|
s.stop()
|
||||||
}
|
|
||||||
|
|
||||||
// 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.
|
// ExportSetInterval overrides the sweep interval for tests.
|
||||||
|
|||||||
@@ -1,6 +1,19 @@
|
|||||||
package handlers
|
package handlers
|
||||||
|
|
||||||
import "net/http"
|
import (
|
||||||
|
"html/template"
|
||||||
|
"net/http"
|
||||||
|
)
|
||||||
|
|
||||||
|
// AddTemplateForTest registers a template under a page name so that
|
||||||
|
// the handlers_test package can drive the render path with a
|
||||||
|
// template of its own.
|
||||||
|
func (s *Handlers) AddTemplateForTest(
|
||||||
|
pageTemplate string,
|
||||||
|
tmpl *template.Template,
|
||||||
|
) {
|
||||||
|
s.templates[pageTemplate] = tmpl
|
||||||
|
}
|
||||||
|
|
||||||
// RenderTemplateForTest exposes renderTemplate for use in the
|
// RenderTemplateForTest exposes renderTemplate for use in the
|
||||||
// handlers_test package.
|
// handlers_test package.
|
||||||
|
|||||||
@@ -3,6 +3,7 @@
|
|||||||
package handlers
|
package handlers
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
@@ -224,13 +225,20 @@ func (s *Handlers) renderTemplate(
|
|||||||
s.executeTemplate(w, tmpl, wrapper)
|
s.executeTemplate(w, tmpl, wrapper)
|
||||||
}
|
}
|
||||||
|
|
||||||
// executeTemplate runs the template and handles errors.
|
// executeTemplate renders the template into a buffer and writes to
|
||||||
|
// the response only once rendering has fully succeeded. Executing
|
||||||
|
// straight into the ResponseWriter commits a partial body and a 200
|
||||||
|
// status before a mid-render error can be reported, leaving no way
|
||||||
|
// to serve a 500. These pages are small, so holding one in memory is
|
||||||
|
// the right trade.
|
||||||
func (s *Handlers) executeTemplate(
|
func (s *Handlers) executeTemplate(
|
||||||
w http.ResponseWriter,
|
w http.ResponseWriter,
|
||||||
tmpl *template.Template,
|
tmpl *template.Template,
|
||||||
data any,
|
data any,
|
||||||
) {
|
) {
|
||||||
err := tmpl.Execute(w, data)
|
var buf bytes.Buffer
|
||||||
|
|
||||||
|
err := tmpl.Execute(&buf, data)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.log.Error(
|
s.log.Error(
|
||||||
"failed to execute template", "error", err,
|
"failed to execute template", "error", err,
|
||||||
@@ -239,5 +247,16 @@ func (s *Handlers) executeTemplate(
|
|||||||
w, "Internal server error",
|
w, "Internal server error",
|
||||||
http.StatusInternalServerError,
|
http.StatusInternalServerError,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
w.Header().Set("Content-Type", "text/html; charset=utf-8")
|
||||||
|
|
||||||
|
_, err = buf.WriteTo(w)
|
||||||
|
if err != nil {
|
||||||
|
s.log.Error(
|
||||||
|
"failed to write rendered page", "error", err,
|
||||||
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2,6 +2,8 @@ package handlers_test
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"errors"
|
||||||
|
"html/template"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"sync"
|
"sync"
|
||||||
@@ -220,6 +222,68 @@ func TestRenderTemplate(t *testing.T) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// errMidRender is the failure a test template raises partway through
|
||||||
|
// rendering.
|
||||||
|
var errMidRender = errors.New("deliberate mid-render failure")
|
||||||
|
|
||||||
|
// midRenderFailure is template data whose first method renders and
|
||||||
|
// whose second fails, so the template aborts after output has
|
||||||
|
// already been produced.
|
||||||
|
type midRenderFailure struct{}
|
||||||
|
|
||||||
|
// Prefix is the output a streaming renderer would flush before the
|
||||||
|
// failure below aborts the template.
|
||||||
|
func (midRenderFailure) Prefix() string { return partialPageMarker }
|
||||||
|
|
||||||
|
// Boom aborts template execution.
|
||||||
|
func (midRenderFailure) Boom() (string, error) {
|
||||||
|
return "", errMidRender
|
||||||
|
}
|
||||||
|
|
||||||
|
// partialPageMarker is content the failing template emits before it
|
||||||
|
// aborts.
|
||||||
|
const partialPageMarker = "PARTIAL PAGE CONTENT"
|
||||||
|
|
||||||
|
// TestRenderTemplateMidRenderErrorSendsNoPartialBody proves the
|
||||||
|
// renderer does not commit output it cannot finish: a template that
|
||||||
|
// fails partway through must yield a 500 and a body carrying none of
|
||||||
|
// the content emitted before the failure. Against a renderer that
|
||||||
|
// executes straight into the ResponseWriter this fails on both
|
||||||
|
// counts, returning 200 with the prefix already flushed.
|
||||||
|
func TestRenderTemplateMidRenderErrorSendsNoPartialBody(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
var h *handlers.Handlers
|
||||||
|
|
||||||
|
app := newTestApp(t, &h)
|
||||||
|
app.RequireStart()
|
||||||
|
|
||||||
|
t.Cleanup(app.RequireStop)
|
||||||
|
|
||||||
|
h.AddTemplateForTest("failing.html", template.Must(
|
||||||
|
template.New("failing").Parse(
|
||||||
|
`{{.Data.Prefix}}{{.Data.Boom}}TAIL`,
|
||||||
|
),
|
||||||
|
))
|
||||||
|
|
||||||
|
req := httptest.NewRequestWithContext(
|
||||||
|
context.Background(), http.MethodGet, "/", nil)
|
||||||
|
w := httptest.NewRecorder()
|
||||||
|
|
||||||
|
h.RenderTemplateForTest(
|
||||||
|
w, req, "failing.html", midRenderFailure{},
|
||||||
|
)
|
||||||
|
|
||||||
|
assert.Equal(
|
||||||
|
t, http.StatusInternalServerError, w.Code,
|
||||||
|
"a failed render must report a 500",
|
||||||
|
)
|
||||||
|
assert.Equal(
|
||||||
|
t, "Internal server error\n", w.Body.String(),
|
||||||
|
"the response must carry no part of the aborted page",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
func TestBuildDatabaseTargetConfig_Valid(t *testing.T) {
|
func TestBuildDatabaseTargetConfig_Valid(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
|
|||||||
@@ -1,57 +0,0 @@
|
|||||||
// Package lifecycle holds helpers shared by the components that
|
|
||||||
// register fx start and stop hooks.
|
|
||||||
package lifecycle
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"fmt"
|
|
||||||
"log/slog"
|
|
||||||
"sync"
|
|
||||||
)
|
|
||||||
|
|
||||||
// WaitForShutdown waits for wg to drain, bounded by ctx.
|
|
||||||
//
|
|
||||||
// fx hands OnStop a context carrying the application's stop
|
|
||||||
// timeout. A bare wg.Wait() discards that deadline, so a single
|
|
||||||
// goroutine that never observes cancellation — a delivery target
|
|
||||||
// that never returns, a SQLite operation blocked on a lock —
|
|
||||||
// hangs the process forever instead of letting it exit when the
|
|
||||||
// timeout expires, which is exactly when a clean shutdown matters
|
|
||||||
// most.
|
|
||||||
//
|
|
||||||
// On timeout it logs at error naming component and returns an
|
|
||||||
// error: the goroutines are still running, and reporting success
|
|
||||||
// would hide an unclean shutdown from the operator. The waiting
|
|
||||||
// goroutine outlives this call and exits when (if) wg drains; it
|
|
||||||
// holds nothing but the channel it closes.
|
|
||||||
func WaitForShutdown(
|
|
||||||
ctx context.Context,
|
|
||||||
log *slog.Logger,
|
|
||||||
component string,
|
|
||||||
wg *sync.WaitGroup,
|
|
||||||
) error {
|
|
||||||
done := make(chan struct{})
|
|
||||||
|
|
||||||
go func() {
|
|
||||||
defer close(done)
|
|
||||||
|
|
||||||
wg.Wait()
|
|
||||||
}()
|
|
||||||
|
|
||||||
select {
|
|
||||||
case <-done:
|
|
||||||
return nil
|
|
||||||
case <-ctx.Done():
|
|
||||||
log.Error(
|
|
||||||
"shutdown timed out, goroutines still running",
|
|
||||||
"component", component,
|
|
||||||
"error", ctx.Err(),
|
|
||||||
)
|
|
||||||
|
|
||||||
return fmt.Errorf(
|
|
||||||
"%s: shutdown timed out, "+
|
|
||||||
"goroutines still running: %w",
|
|
||||||
component, ctx.Err(),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,62 +0,0 @@
|
|||||||
package lifecycle_test
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"log/slog"
|
|
||||||
"sync"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/stretchr/testify/require"
|
|
||||||
"sneak.berlin/go/webhooker/internal/lifecycle"
|
|
||||||
)
|
|
||||||
|
|
||||||
// waitTimeout is the stop budget the timeout case gives a
|
|
||||||
// goroutine that never returns. The test's own patience is the
|
|
||||||
// go test deadline, so the only thing this value affects is how
|
|
||||||
// long the case takes.
|
|
||||||
const waitTimeout = 100 * time.Millisecond
|
|
||||||
|
|
||||||
func discardLogger() *slog.Logger {
|
|
||||||
return slog.New(slog.DiscardHandler)
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestWaitForShutdown_DrainedGroup(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
var wg sync.WaitGroup
|
|
||||||
|
|
||||||
wg.Go(func() {})
|
|
||||||
|
|
||||||
require.NoError(
|
|
||||||
t,
|
|
||||||
lifecycle.WaitForShutdown(
|
|
||||||
context.Background(), discardLogger(),
|
|
||||||
"test component", &wg,
|
|
||||||
),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestWaitForShutdown_ContextExpires(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
release := make(chan struct{})
|
|
||||||
|
|
||||||
t.Cleanup(func() { close(release) })
|
|
||||||
|
|
||||||
var wg sync.WaitGroup
|
|
||||||
|
|
||||||
wg.Go(func() { <-release })
|
|
||||||
|
|
||||||
ctx, cancel := context.WithTimeout(
|
|
||||||
context.Background(), waitTimeout,
|
|
||||||
)
|
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
err := lifecycle.WaitForShutdown(
|
|
||||||
ctx, discardLogger(), "test component", &wg,
|
|
||||||
)
|
|
||||||
|
|
||||||
require.ErrorIs(t, err, context.DeadlineExceeded)
|
|
||||||
require.ErrorContains(t, err, "test component")
|
|
||||||
}
|
|
||||||
Reference in New Issue
Block a user