Compare commits
1 Commits
issue-102-
...
f5bcfdccb1
| Author | SHA1 | Date | |
|---|---|---|---|
| f5bcfdccb1 |
24
README.md
24
README.md
@@ -100,10 +100,9 @@ TTY detection, and security headers are always applied.
|
|||||||
#### Trusted proxies
|
#### Trusted proxies
|
||||||
|
|
||||||
`TRUSTED_PROXIES` is a comma-separated list of CIDR blocks (a bare
|
`TRUSTED_PROXIES` is a comma-separated list of CIDR blocks (a bare
|
||||||
address such as `192.168.1.7` is accepted and treated as a single
|
address such as `10.0.0.1` is accepted and treated as a single host),
|
||||||
host), for example `192.168.1.7, 2001:db8::5`. It decides whose
|
for example `10.0.0.0/8, 192.168.1.7, 2001:db8::/32`. It decides whose
|
||||||
`X-Forwarded-For` header the rate limiters believe, so it should name
|
`X-Forwarded-For` header the rate limiters believe.
|
||||||
the addresses of your reverse proxies and nothing else.
|
|
||||||
|
|
||||||
`X-Forwarded-For` is honoured **only** when the connecting peer is
|
`X-Forwarded-For` is honoured **only** when the connecting peer is
|
||||||
inside one of these blocks; for every other peer the client identity is
|
inside one of these blocks; for every other peer the client identity is
|
||||||
@@ -126,8 +125,7 @@ hop that is not itself a trusted proxy is taken as the client. A hop
|
|||||||
that is not a bare IP address — `ip:port`, a bracketed IPv6 literal,
|
that is not a bare IP address — `ip:port`, a bracketed IPv6 literal,
|
||||||
the token `unknown` — ends the walk and the peer address is used
|
the token `unknown` — ends the walk and the peer address is used
|
||||||
instead, since past such an entry the chain is not the shape assumed
|
instead, since past such an entry the chain is not the shape assumed
|
||||||
here. The peer address is likewise used when the header is absent or
|
here.
|
||||||
every hop in it is a trusted proxy.
|
|
||||||
|
|
||||||
Two operator requirements follow:
|
Two operator requirements follow:
|
||||||
|
|
||||||
@@ -135,14 +133,12 @@ Two operator requirements follow:
|
|||||||
(nginx `$proxy_add_x_forwarded_for`, HAProxy `option forwardfor`,
|
(nginx `$proxy_add_x_forwarded_for`, HAProxy `option forwardfor`,
|
||||||
Caddy and AWS ALB by default), and must append a bare address with
|
Caddy and AWS ALB by default), and must append a bare address with
|
||||||
no port.
|
no port.
|
||||||
- List proxy hosts **only**. Any address inside `TRUSTED_PROXIES`
|
- Keep the list narrow. A client whose own address falls inside a
|
||||||
chooses its own rate-limit key: its `X-Forwarded-For` is walked, so
|
broad block such as `10.0.0.0/8` is treated as a proxy: its address
|
||||||
it can name a different address on every request to get a fresh
|
is skipped during the walk, so it shares the bucket of whatever lies
|
||||||
bucket each time, or name another client's address to drain that
|
further left rather than getting one of its own. That is safe — the
|
||||||
client's bucket. Never list a block that also covers clients — a
|
block is trusted by definition — but surprising if the block covers
|
||||||
broad `10.0.0.0/8` on a network where clients live in the same range
|
ordinary clients as well as proxies.
|
||||||
makes all three limits, including the unauthenticated webhook
|
|
||||||
receiver, silently bypassable by every client in the block.
|
|
||||||
|
|
||||||
Sessions are bounded by two independent clocks, and end at whichever
|
Sessions are bounded by two independent clocks, and end at whichever
|
||||||
one runs out first:
|
one runs out first:
|
||||||
|
|||||||
4
TODO.md
4
TODO.md
@@ -26,10 +26,6 @@ capability in the README rationale).
|
|||||||
|
|
||||||
# Completed Steps
|
# Completed Steps
|
||||||
|
|
||||||
- 2026-08-11 Web UI cleanup: nav terminology unified on Webhooks, the
|
|
||||||
Profile settings placeholder removed, a progressive-enhancement copy
|
|
||||||
button for the entrypoint URL, and retention form copy that states the
|
|
||||||
actual policy (deletion by the reaper, 0 retains forever) (#57)
|
|
||||||
- 2026-08-09 Inactivity-based session timeout: sliding idle expiry
|
- 2026-08-09 Inactivity-based session timeout: sliding idle expiry
|
||||||
(`SESSION_IDLE_TIMEOUT`, default `24h`) refreshed on authenticated
|
(`SESSION_IDLE_TIMEOUT`, default `24h`) refreshed on authenticated
|
||||||
requests, with the 7-day absolute cap kept as an independent
|
requests, with the 7-day absolute cap kept as an independent
|
||||||
|
|||||||
@@ -103,14 +103,11 @@ type Config struct {
|
|||||||
ReceiverRateLimit int
|
ReceiverRateLimit int
|
||||||
|
|
||||||
// TrustedProxies is the set of networks whose members are
|
// TrustedProxies is the set of networks whose members are
|
||||||
// allowed to speak for the client with X-Forwarded-For, the
|
// allowed to speak for the client with forwarded headers
|
||||||
// only forwarded header read. It is empty unless
|
// (X-Forwarded-For, X-Real-IP, True-Client-IP). It is empty
|
||||||
// TRUSTED_PROXIES is set, and empty means no peer is
|
// unless TRUSTED_PROXIES is set, and empty means no peer is
|
||||||
// trusted: forwarded headers are then ignored entirely and
|
// trusted: forwarded headers are then ignored entirely and
|
||||||
// clients are identified by the connection's own address.
|
// clients are identified by the connection's own address.
|
||||||
// Members can choose their own rate-limit key, so this must
|
|
||||||
// name proxy hosts only, never a block that also covers
|
|
||||||
// clients.
|
|
||||||
TrustedProxies []netip.Prefix
|
TrustedProxies []netip.Prefix
|
||||||
|
|
||||||
params *ConfigParams
|
params *ConfigParams
|
||||||
|
|||||||
@@ -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")
|
||||||
|
|
||||||
if e.cancel != nil {
|
|
||||||
e.cancel()
|
e.cancel()
|
||||||
}
|
e.wg.Wait()
|
||||||
|
|
||||||
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,299 +0,0 @@
|
|||||||
package handlers_test
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"net/http"
|
|
||||||
"net/http/httptest"
|
|
||||||
"strconv"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/stretchr/testify/assert"
|
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
|
||||||
"sneak.berlin/go/webhooker/internal/delivery"
|
|
||||||
"sneak.berlin/go/webhooker/internal/handlers"
|
|
||||||
"sneak.berlin/go/webhooker/internal/session"
|
|
||||||
)
|
|
||||||
|
|
||||||
// Template data keys the page templates read. The handlers package has
|
|
||||||
// its own unexported constants for these; this is the external test
|
|
||||||
// package, so it needs its own.
|
|
||||||
const (
|
|
||||||
dataKeyWebhook = "Webhook"
|
|
||||||
dataKeyError = "Error"
|
|
||||||
)
|
|
||||||
|
|
||||||
// testWebhookID is the identifier given to the webhook under test on
|
|
||||||
// pages that render one.
|
|
||||||
const testWebhookID = "wh-1"
|
|
||||||
|
|
||||||
// renderPage renders a page template through the real template set as
|
|
||||||
// an authenticated user and returns the resulting HTML.
|
|
||||||
func renderPage(
|
|
||||||
t *testing.T,
|
|
||||||
h *handlers.Handlers,
|
|
||||||
sess *session.Session,
|
|
||||||
page string,
|
|
||||||
data map[string]any,
|
|
||||||
) string {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
cookies := authenticatedCookies(t, sess, "test-user-id", "testuser")
|
|
||||||
|
|
||||||
req := httptest.NewRequestWithContext(
|
|
||||||
context.Background(), http.MethodGet, "/", nil,
|
|
||||||
)
|
|
||||||
for _, c := range cookies {
|
|
||||||
req.AddCookie(c)
|
|
||||||
}
|
|
||||||
|
|
||||||
w := httptest.NewRecorder()
|
|
||||||
h.RenderTemplateForTest(w, req, page, data)
|
|
||||||
|
|
||||||
return w.Body.String()
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestNavbarUsesWebhookTerminology pins the user-visible navigation
|
|
||||||
// label to "Webhooks". The /sources route is deliberately unchanged, so
|
|
||||||
// the assertion targets the link text rather than the href.
|
|
||||||
func TestNavbarUsesWebhookTerminology(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
var h *handlers.Handlers
|
|
||||||
|
|
||||||
var sess *session.Session
|
|
||||||
|
|
||||||
app := newTestApp(t, &h, &sess)
|
|
||||||
app.RequireStart()
|
|
||||||
|
|
||||||
t.Cleanup(app.RequireStop)
|
|
||||||
|
|
||||||
// One item, so the list body renders too: it calls
|
|
||||||
// WebhookListItem.RetentionLabel, promoted from the embedded
|
|
||||||
// Webhook and therefore a pointer method. An empty list would
|
|
||||||
// skip that call and hide a template error behind the
|
|
||||||
// navigation assertions below.
|
|
||||||
item := handlers.WebhookListItem{}
|
|
||||||
item.Name = "wh"
|
|
||||||
item.ID = testWebhookID
|
|
||||||
item.RetentionDays = 14
|
|
||||||
|
|
||||||
body := renderPage(t, h, sess, "sources_list.html", map[string]any{
|
|
||||||
"Webhooks": []handlers.WebhookListItem{item},
|
|
||||||
})
|
|
||||||
|
|
||||||
assert.Contains(t, body, "Retention: 14 days")
|
|
||||||
assert.Contains(t, body, `class="btn-text">Webhooks</a>`)
|
|
||||||
assert.Contains(
|
|
||||||
t, body, `class="btn-text w-full text-left">Webhooks</a>`,
|
|
||||||
)
|
|
||||||
assert.Contains(
|
|
||||||
t, body,
|
|
||||||
`<h1 class="text-2xl font-medium text-gray-900">Webhooks</h1>`,
|
|
||||||
)
|
|
||||||
assert.NotContains(
|
|
||||||
t, body, ">Sources<",
|
|
||||||
"no user-visible element may still be labelled Sources",
|
|
||||||
)
|
|
||||||
assert.Contains(
|
|
||||||
t, body, `href="/sources"`,
|
|
||||||
"the /sources route itself must not change",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestEditPageUsesWebhookTerminology pins the edit page's heading and
|
|
||||||
// its back link. The link's href still points at /source/{id}, which is
|
|
||||||
// intentional: only user-visible copy changes.
|
|
||||||
func TestEditPageUsesWebhookTerminology(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
var h *handlers.Handlers
|
|
||||||
|
|
||||||
var sess *session.Session
|
|
||||||
|
|
||||||
app := newTestApp(t, &h, &sess)
|
|
||||||
app.RequireStart()
|
|
||||||
|
|
||||||
t.Cleanup(app.RequireStop)
|
|
||||||
|
|
||||||
// The webhook goes in as a pointer because source_edit.html calls
|
|
||||||
// Webhook.RetentionLabel, a pointer method: a map element is not
|
|
||||||
// addressable, so a value here renders an error instead of the
|
|
||||||
// page.
|
|
||||||
webhook := &database.Webhook{Name: "wh", RetentionDays: 14}
|
|
||||||
webhook.ID = testWebhookID
|
|
||||||
|
|
||||||
body := renderPage(t, h, sess, "source_edit.html", map[string]any{
|
|
||||||
dataKeyWebhook: webhook,
|
|
||||||
dataKeyError: "",
|
|
||||||
})
|
|
||||||
|
|
||||||
assert.Contains(t, body, "Edit Webhook")
|
|
||||||
assert.NotContains(t, body, ">Sources<")
|
|
||||||
assert.Contains(t, body, `href="/source/wh-1"`)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestCreateFormRetentionCopyMatchesBehaviour pins the create form's
|
|
||||||
// retention copy to what the code does: the reaper permanently deletes
|
|
||||||
// events past the cutoff, an empty field falls back to
|
|
||||||
// DefaultRetentionDays, and 0 is rewritten to the retain-forever
|
|
||||||
// sentinel by Webhook.BeforeSave.
|
|
||||||
func TestCreateFormRetentionCopyMatchesBehaviour(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
var h *handlers.Handlers
|
|
||||||
|
|
||||||
var sess *session.Session
|
|
||||||
|
|
||||||
app := newTestApp(t, &h, &sess)
|
|
||||||
app.RequireStart()
|
|
||||||
|
|
||||||
t.Cleanup(app.RequireStop)
|
|
||||||
|
|
||||||
body := renderPage(t, h, sess, "sources_new.html", map[string]any{
|
|
||||||
"Name": "",
|
|
||||||
"Description": "",
|
|
||||||
"DefaultRetentionDays": database.DefaultRetentionDays,
|
|
||||||
dataKeyError: "",
|
|
||||||
})
|
|
||||||
|
|
||||||
assert.Contains(
|
|
||||||
t, body,
|
|
||||||
"permanently deletes events older than this",
|
|
||||||
"the form must say retention is enforced by deletion",
|
|
||||||
)
|
|
||||||
assert.Contains(t, body, "Enter 0 to retain events forever")
|
|
||||||
assert.Contains(
|
|
||||||
t, body,
|
|
||||||
"leave blank to use the default of "+
|
|
||||||
strconv.Itoa(database.DefaultRetentionDays)+" days",
|
|
||||||
"blank means the default, not forever",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestEditFormRetentionCopyMatchesBehaviour pins the edit form's
|
|
||||||
// retention copy, including that it states the stored policy via
|
|
||||||
// RetentionLabel and that an empty field leaves that policy unchanged
|
|
||||||
// rather than meaning forever.
|
|
||||||
func TestEditFormRetentionCopyMatchesBehaviour(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
var h *handlers.Handlers
|
|
||||||
|
|
||||||
var sess *session.Session
|
|
||||||
|
|
||||||
app := newTestApp(t, &h, &sess)
|
|
||||||
app.RequireStart()
|
|
||||||
|
|
||||||
t.Cleanup(app.RequireStop)
|
|
||||||
|
|
||||||
finite := &database.Webhook{Name: "wh", RetentionDays: 14}
|
|
||||||
finite.ID = testWebhookID
|
|
||||||
|
|
||||||
body := renderPage(t, h, sess, "source_edit.html", map[string]any{
|
|
||||||
dataKeyWebhook: finite,
|
|
||||||
dataKeyError: "",
|
|
||||||
})
|
|
||||||
|
|
||||||
assert.Contains(t, body, "Currently 14 days.")
|
|
||||||
assert.Contains(
|
|
||||||
t, body,
|
|
||||||
"permanently deletes events older than this",
|
|
||||||
)
|
|
||||||
assert.Contains(t, body, "Enter 0 to retain events forever")
|
|
||||||
assert.Contains(
|
|
||||||
t, body,
|
|
||||||
"leave blank to keep the current setting",
|
|
||||||
"blank means unchanged, not forever",
|
|
||||||
)
|
|
||||||
|
|
||||||
forever := &database.Webhook{
|
|
||||||
Name: "wh",
|
|
||||||
RetentionDays: database.RetentionForeverDays,
|
|
||||||
}
|
|
||||||
forever.ID = "wh-2"
|
|
||||||
|
|
||||||
foreverBody := renderPage(
|
|
||||||
t, h, sess, "source_edit.html", map[string]any{
|
|
||||||
dataKeyWebhook: forever,
|
|
||||||
dataKeyError: "",
|
|
||||||
},
|
|
||||||
)
|
|
||||||
|
|
||||||
assert.Contains(
|
|
||||||
t, foreverBody, "Currently forever.",
|
|
||||||
"a retain-forever webhook must not read as a day count",
|
|
||||||
)
|
|
||||||
assert.Contains(
|
|
||||||
t, foreverBody,
|
|
||||||
"No events are deleted while retention is set to forever",
|
|
||||||
)
|
|
||||||
assert.NotContains(
|
|
||||||
t, foreverBody,
|
|
||||||
"permanently deletes events older than this",
|
|
||||||
"the reaper skips retain-forever webhooks, so the form "+
|
|
||||||
"must not claim it deletes their events",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestEntrypointCopyButtonIsProgressiveEnhancement proves the copy
|
|
||||||
// affordance degrades: the button ships with the hidden attribute, so a
|
|
||||||
// browser that never runs app.js shows no dead control, and the URL is
|
|
||||||
// rendered as ordinary selectable text either way.
|
|
||||||
func TestEntrypointCopyButtonIsProgressiveEnhancement(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
var h *handlers.Handlers
|
|
||||||
|
|
||||||
var sess *session.Session
|
|
||||||
|
|
||||||
app := newTestApp(t, &h, &sess)
|
|
||||||
app.RequireStart()
|
|
||||||
|
|
||||||
t.Cleanup(app.RequireStop)
|
|
||||||
|
|
||||||
entrypoint := database.Entrypoint{Path: "abc123"}
|
|
||||||
entrypoint.ID = "ep-1"
|
|
||||||
|
|
||||||
// The webhook goes in as a pointer because source_detail.html
|
|
||||||
// calls Webhook.RetentionLabel, a pointer method: a map element
|
|
||||||
// is not addressable, so a value here aborts execution partway
|
|
||||||
// down the page, after the copy button has already been flushed
|
|
||||||
// to the response.
|
|
||||||
webhook := &database.Webhook{Name: "wh", RetentionDays: 14}
|
|
||||||
webhook.ID = testWebhookID
|
|
||||||
webhook.CreatedAt = time.Date(
|
|
||||||
2026, time.January, 2, 3, 4, 5, 0, time.UTC,
|
|
||||||
)
|
|
||||||
|
|
||||||
body := renderPage(t, h, sess, "source_detail.html", map[string]any{
|
|
||||||
dataKeyWebhook: webhook,
|
|
||||||
"Entrypoints": []database.Entrypoint{entrypoint},
|
|
||||||
// The handler passes delivery.NewTargetViews(targets), never
|
|
||||||
// raw targets, so the test data has to have that same shape.
|
|
||||||
"Targets": delivery.NewTargetViews(nil),
|
|
||||||
"Events": []database.Event{},
|
|
||||||
"BaseURL": "https://hooks.example.com",
|
|
||||||
})
|
|
||||||
|
|
||||||
assert.Contains(
|
|
||||||
t, body,
|
|
||||||
`<code id="entrypoint-url-ep-1"`,
|
|
||||||
)
|
|
||||||
assert.Contains(t, body, "https://hooks.example.com/webhook/abc123")
|
|
||||||
assert.Contains(
|
|
||||||
t, body,
|
|
||||||
`hidden data-copy-target="entrypoint-url-ep-1"`,
|
|
||||||
"the button must start hidden and be revealed by script",
|
|
||||||
)
|
|
||||||
|
|
||||||
// renderTemplate streams to the ResponseWriter, so an abort
|
|
||||||
// midway still leaves everything above it in the body. This pins
|
|
||||||
// content from the last line of the template, which is below the
|
|
||||||
// assertions above: without it, a page that renders the copy
|
|
||||||
// button and then 500s passes.
|
|
||||||
assert.Contains(
|
|
||||||
t, body, "Retention: 14 days",
|
|
||||||
"the page must render to completion, not abort partway",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
@@ -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")
|
|
||||||
}
|
|
||||||
@@ -1,60 +1,2 @@
|
|||||||
// Webhooker client-side JavaScript
|
// Webhooker client-side JavaScript
|
||||||
console.log("Webhooker loaded");
|
console.log("Webhooker loaded");
|
||||||
|
|
||||||
// Copy-to-clipboard, as progressive enhancement.
|
|
||||||
//
|
|
||||||
// Markup renders each copy button with the `hidden` attribute and a
|
|
||||||
// `data-copy-target` pointing at the id of the element holding the
|
|
||||||
// text. This script reveals a button only once it has both a resolvable
|
|
||||||
// target and a usable Clipboard API, so a browser without either shows
|
|
||||||
// no button at all and the text stays selectable.
|
|
||||||
(function () {
|
|
||||||
"use strict";
|
|
||||||
|
|
||||||
const revertDelayMs = 2000;
|
|
||||||
|
|
||||||
function flash(button, message) {
|
|
||||||
const original = button.getAttribute("data-copy-label");
|
|
||||||
button.textContent = message;
|
|
||||||
window.setTimeout(function () {
|
|
||||||
button.textContent = original;
|
|
||||||
}, revertDelayMs);
|
|
||||||
}
|
|
||||||
|
|
||||||
function wire(button) {
|
|
||||||
const target = document.getElementById(
|
|
||||||
button.getAttribute("data-copy-target")
|
|
||||||
);
|
|
||||||
if (!target) {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
button.setAttribute("data-copy-label", button.textContent);
|
|
||||||
button.addEventListener("click", function () {
|
|
||||||
navigator.clipboard.writeText(target.textContent.trim()).then(
|
|
||||||
function () {
|
|
||||||
flash(button, "Copied");
|
|
||||||
},
|
|
||||||
function () {
|
|
||||||
flash(button, "Copy failed");
|
|
||||||
}
|
|
||||||
);
|
|
||||||
});
|
|
||||||
button.removeAttribute("hidden");
|
|
||||||
}
|
|
||||||
|
|
||||||
function init() {
|
|
||||||
if (!navigator.clipboard || !navigator.clipboard.writeText) {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
const buttons = document.querySelectorAll("[data-copy-target]");
|
|
||||||
buttons.forEach(wire);
|
|
||||||
}
|
|
||||||
|
|
||||||
if (document.readyState === "loading") {
|
|
||||||
document.addEventListener("DOMContentLoaded", init);
|
|
||||||
} else {
|
|
||||||
init();
|
|
||||||
}
|
|
||||||
})();
|
|
||||||
|
|||||||
@@ -16,7 +16,7 @@
|
|||||||
<!-- Desktop navigation -->
|
<!-- Desktop navigation -->
|
||||||
<div class="hidden md:flex items-center gap-4">
|
<div class="hidden md:flex items-center gap-4">
|
||||||
{{if .User}}
|
{{if .User}}
|
||||||
<a href="/sources" class="btn-text">Webhooks</a>
|
<a href="/sources" class="btn-text">Sources</a>
|
||||||
<a href="/user/{{.User.Username}}" class="btn-text">
|
<a href="/user/{{.User.Username}}" class="btn-text">
|
||||||
<svg class="w-5 h-5 mr-1" fill="currentColor" viewBox="0 0 16 16">
|
<svg class="w-5 h-5 mr-1" fill="currentColor" viewBox="0 0 16 16">
|
||||||
<path d="M11 6a3 3 0 1 1-6 0 3 3 0 0 1 6 0z"/>
|
<path d="M11 6a3 3 0 1 1-6 0 3 3 0 0 1 6 0z"/>
|
||||||
@@ -38,7 +38,7 @@
|
|||||||
<div x-show="open" x-cloak x-transition class="md:hidden mt-4 pt-4 border-t border-gray-200">
|
<div x-show="open" x-cloak x-transition class="md:hidden mt-4 pt-4 border-t border-gray-200">
|
||||||
<div class="flex flex-col gap-2">
|
<div class="flex flex-col gap-2">
|
||||||
{{if .User}}
|
{{if .User}}
|
||||||
<a href="/sources" class="btn-text w-full text-left">Webhooks</a>
|
<a href="/sources" class="btn-text w-full text-left">Sources</a>
|
||||||
<a href="/user/{{.User.Username}}" class="btn-text w-full text-left">Profile</a>
|
<a href="/user/{{.User.Username}}" class="btn-text w-full text-left">Profile</a>
|
||||||
<form method="POST" action="/pages/logout">
|
<form method="POST" action="/pages/logout">
|
||||||
<input type="hidden" name="csrf_token" value="{{.CSRFToken}}">
|
<input type="hidden" name="csrf_token" value="{{.CSRFToken}}">
|
||||||
|
|||||||
@@ -34,6 +34,7 @@
|
|||||||
|
|
||||||
<hr class="border-gray-200 mb-6">
|
<hr class="border-gray-200 mb-6">
|
||||||
|
|
||||||
|
<div class="grid grid-cols-1 md:grid-cols-2 gap-8">
|
||||||
<div>
|
<div>
|
||||||
<h3 class="text-lg font-medium text-gray-900 mb-3">Account Information</h3>
|
<h3 class="text-lg font-medium text-gray-900 mb-3">Account Information</h3>
|
||||||
<dl class="space-y-3">
|
<dl class="space-y-3">
|
||||||
@@ -47,6 +48,11 @@
|
|||||||
</div>
|
</div>
|
||||||
</dl>
|
</dl>
|
||||||
</div>
|
</div>
|
||||||
|
<div>
|
||||||
|
<h3 class="text-lg font-medium text-gray-900 mb-3">Settings</h3>
|
||||||
|
<p class="text-sm text-gray-500">Profile settings and preferences will be available here.</p>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
<div class="card p-6 mt-6">
|
<div class="card p-6 mt-6">
|
||||||
|
|||||||
@@ -69,12 +69,7 @@
|
|||||||
</form>
|
</form>
|
||||||
</div>
|
</div>
|
||||||
</div>
|
</div>
|
||||||
<div class="flex items-start gap-2 mt-1">
|
<code class="text-xs text-gray-500 break-all block mt-1">{{$.BaseURL}}/webhook/{{.Path}}</code>
|
||||||
<code id="entrypoint-url-{{.ID}}" class="text-xs text-gray-500 break-all block flex-1">{{$.BaseURL}}/webhook/{{.Path}}</code>
|
|
||||||
<!-- Hidden until app.js reveals it; without the
|
|
||||||
script the URL above stays selectable. -->
|
|
||||||
<button type="button" hidden data-copy-target="entrypoint-url-{{.ID}}" class="text-xs text-gray-500 hover:text-primary-600">Copy</button>
|
|
||||||
</div>
|
|
||||||
</div>
|
</div>
|
||||||
{{else}}
|
{{else}}
|
||||||
<div class="p-4 text-sm text-gray-500">No entrypoints configured.</div>
|
<div class="p-4 text-sm text-gray-500">No entrypoints configured.</div>
|
||||||
|
|||||||
@@ -29,7 +29,7 @@
|
|||||||
<div class="form-group">
|
<div class="form-group">
|
||||||
<label for="retention_days" class="label">Retention (days)</label>
|
<label for="retention_days" class="label">Retention (days)</label>
|
||||||
<input type="number" id="retention_days" name="retention_days" value="{{.Webhook.RetentionDays}}" min="0" class="input">
|
<input type="number" id="retention_days" name="retention_days" value="{{.Webhook.RetentionDays}}" min="0" class="input">
|
||||||
<p class="text-xs text-gray-500 mt-1">Currently {{.Webhook.RetentionLabel}}.{{if .Webhook.RetainsForever}} No events are deleted while retention is set to forever.{{else}} A periodic cleanup permanently deletes events older than this, along with their delivery records.{{end}} Enter 0 to retain events forever; leave blank to keep the current setting.</p>
|
<p class="text-xs text-gray-500 mt-1">Currently {{.Webhook.RetentionLabel}}. Enter 0 to retain events forever.</p>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
<div class="flex gap-3">
|
<div class="flex gap-3">
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
{{template "base" .}}
|
{{template "base" .}}
|
||||||
|
|
||||||
{{define "title"}}Webhooks - Webhooker{{end}}
|
{{define "title"}}Sources - Webhooker{{end}}
|
||||||
|
|
||||||
{{define "content"}}
|
{{define "content"}}
|
||||||
<div class="max-w-6xl mx-auto px-6 py-8">
|
<div class="max-w-6xl mx-auto px-6 py-8">
|
||||||
|
|||||||
@@ -29,7 +29,7 @@
|
|||||||
<div class="form-group">
|
<div class="form-group">
|
||||||
<label for="retention_days" class="label">Retention (days)</label>
|
<label for="retention_days" class="label">Retention (days)</label>
|
||||||
<input type="number" id="retention_days" name="retention_days" value="{{.DefaultRetentionDays}}" min="0" class="input">
|
<input type="number" id="retention_days" name="retention_days" value="{{.DefaultRetentionDays}}" min="0" class="input">
|
||||||
<p class="text-xs text-gray-500 mt-1">A periodic cleanup permanently deletes events older than this, along with their delivery records. Enter 0 to retain events forever; leave blank to use the default of {{.DefaultRetentionDays}} days.</p>
|
<p class="text-xs text-gray-500 mt-1">How long to keep event data. Enter 0 to retain events forever.</p>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
<div class="flex gap-3">
|
<div class="flex gap-3">
|
||||||
|
|||||||
Reference in New Issue
Block a user