1 Commits

Author SHA1 Message Date
clawbot
df1f76b006 Evict archive writers on deletion and sweep idle archives (closes #89)
All checks were successful
check / check (push) Successful in 2m56s
The per-webhook archiveWriter registry in the database delivery target
was never evicted, so a deleted webhook's writer -- and any archive
file handle open within its debounce window -- lingered for the process
lifetime. Separately, expiry pruning ran only when an archive was
(re)opened, and reopens only happen on writes, so an archive belonging
to a webhook that stopped receiving events kept its expired rows
forever.

Eviction: a new one-method delivery.WebhookEvictor interface (kept
separate from Notifier: archiving lifecycle is not notification) is
implemented by the Engine and injected into the handlers. Deleting a
webhook, or deleting its last database target, drops the writer from
the registry and closes its handle under the writer's own mutex, so
eviction can never race an in-flight write. An evicted writer refuses
further writes rather than reopening a file nothing holds. The archive
file is deliberately left on disk: it is long-term storage an operator
may want to keep or move away, and destroying it as a side effect of
deleting a webhook would be unrecoverable.

Idle sweep: a new ArchiveSweeper, modelled on the event RetentionReaper
(fx lifecycle hooks, cancellable context, WaitGroup, ticker loop),
prunes archives whose database target declares a positive expiry. It
reuses the existing RETENTION_SWEEP_INTERVAL rather than adding a
config key. It never creates an archive -- a missing file is skipped,
and the reopen uses SQLite mode=rw so the file cannot be conjured even
if it disappears mid-sweep -- routes the prune through the per-webhook
writer so its mutex orders the sweep against concurrent writes, and
leaves the archive closed so the move-the-file-away workflow keeps
working. A failure for one webhook is logged and the sweep continues.
Archives with no expiry or the expiry "never" are untouched.

The sweep loop's context is rooted at context.Background(), not at the
fx OnStart hook context. The hook context carries fx's 15 second start
timeout, so a loop derived from it is cancelled three quarters of an
hour before the first tick under the default one-hour interval, giving
a sweeper that never sweeps. OnStop still cancels the loop and waits on
the WaitGroup, so shutdown is unchanged.

The sweep also never leaves a registry entry behind. Reaching the
writer through the ordinary create-and-cache accessor would let a sweep
that raced a webhook deletion re-insert a writer for a webhook that no
longer exists, which nothing would ever evict again -- the very leak
this change closes. An entry the sweep has to create is marked
sweep-owned and released when the prune finishes, unless a delivery
claimed it meanwhile, in which case it belongs to the registry and an
eviction can still reach it. A writer evicted underneath a sweep is an
ordinary interleaving and is logged at debug, not error.
2026-08-09 03:05:13 +00:00
23 changed files with 2323 additions and 894 deletions

View File

@@ -92,27 +92,6 @@ TTY detection, and security headers are always applied.
| `METRICS_USERNAME` | Basic auth username for `/metrics` | `""` |
| `METRICS_PASSWORD` | Basic auth password for `/metrics` | `""` |
| `SENTRY_DSN` | Sentry error reporting DSN | `""` |
| `SESSION_IDLE_TIMEOUT` | Idle session timeout (Go duration) | `24h` |
Sessions are bounded by two independent clocks, and end at whichever
one runs out first:
- **Idle expiry** (`SESSION_IDLE_TIMEOUT`, default `24h`) is a sliding
window. Every authenticated request pushes it forward, so a session
in continuous use never hits it, while an abandoned one expires a day
after its last use. Set it to `0` to disable idle expiry entirely;
the absolute cap below still applies. A set-but-unparseable value
aborts startup rather than silently falling back to the default.
- **Absolute expiry** is a fixed 7 days from login. Activity does
**not** extend it: after a week, every session ends and the user
authenticates again.
Only requests that authenticate with the session count as activity, so
an unauthenticated request carrying the cookie cannot keep a session
alive. The idle timestamp is rewritten at most once per tenth of the
idle window rather than on every request, which means a session may
expire up to 10% early relative to the user's true last request, but
never late.
On first startup, webhooker automatically generates a cryptographically
secure session encryption key and stores it in the database. This key
@@ -552,6 +531,36 @@ older than the expiry are pruned each time the archive is (re)opened. An
archive write failure is never silent success: the delivery records a
failed attempt with the error and is marked failed.
Because reopens only happen on writes, an archive belonging to a webhook
that has stopped receiving events would never be pruned. A background
**archive sweeper** closes that gap: on the same interval as the event
retention reaper (`RETENTION_SWEEP_INTERVAL`) it prunes every archive
whose database target declares a positive expiry, whether or not the
webhook is still receiving traffic. The sweep never creates an archive —
a webhook whose archive file does not yet exist is skipped, not
initialised — it takes the same per-webhook lock the write path uses, so
it can never interleave with a write, and it leaves the archive closed
afterwards so the move-the-file-away workflow keeps working. Archives
with no expiry, or the expiry `never`, are not touched by the sweep at
all.
Note that a webhook has one archive file but may carry more than one
`database` target, each with its own `expiry`. The shortest expiry
configured on any of them therefore governs the whole archive, and the
sweep applies it whether or not the webhook is still receiving events.
Configure a single `database` target per webhook unless you intend that.
Deleting a webhook releases its archive: the delivery engine's cached
archive writer is dropped and its file handle closed, so nothing lingers
after the webhook is gone. The archive **file itself is deliberately
left on disk**. Unlike the event database — per-webhook working storage
that is hard-deleted with the webhook — an archive is long-term storage
an operator may still want to keep or move away for offline retention,
and destroying it as a side effect of deleting a webhook would be
unrecoverable. Removing `archive-{webhookID}.db` is the operator's call.
Deleting a webhook's last `database` target releases the writer the same
way, and for the same reason leaves the file alone.
The **Slack target type** sends webhook events as formatted messages to
any Slack-compatible incoming webhook URL (works with Slack, Mattermost,
and other compatible services). Each message includes event metadata

17
TODO.md
View File

@@ -10,9 +10,11 @@
# Status
pre-1.0. No git tags exist. main (afe88c6) is a working webhook proxy
pre-1.0. No git tags exist. main (4f5ecb1) is a working webhook proxy
with auth, CSRF/SSRF protections, login rate limiting, Slack target,
policy compliance (#6), and pinned lint tooling (#55). Note: TODO.md was
event retention (#63), the database archiving target (#43), the admin
password change flow (#65), policy compliance (#6), and pinned lint
tooling (#55). Note: TODO.md was
deliberately deleted from this repo in f9a9569 (2026-03-01, #6); its
content was folded into the README TODO section, which this draft
reconstructs as of 2026-07-06.
@@ -28,10 +30,11 @@ databases currently grow without bound.
# Completed Steps
- 2026-08-09 Inactivity-based session timeout: sliding idle expiry
(`SESSION_IDLE_TIMEOUT`, default `24h`) refreshed on authenticated
requests, with the 7-day absolute cap kept as an independent
backstop that activity never extends (#66)
- 2026-08-09 Archive writer lifecycle (#89): deleting a webhook (or its
last `database` target) evicts the cached archive writer and closes
its handle while deliberately leaving `archive-{webhookID}.db` on
disk, and a new `ArchiveSweeper` prunes idle archives on the existing
`RETENTION_SWEEP_INTERVAL` without ever creating an archive file
- 2026-08-07 Update golangci-lint to v2.12.2 (Docker image digest in
`Dockerfile`, release-archive sha256 pins in `script/bootstrap`),
adopt the canonical `.golangci.yml` (v2 `linters.settings` layout so
@@ -75,7 +78,7 @@ databases currently grow without bound.
- event redelivery endpoint
- OpenAPI specification
- Analytics dashboard: success rates, response times, volume
- A remember-me option at login
- Session expiration tuning and a remember-me option
- Password change and reset flow
- Later, nice to have
- email delivery target type

View File

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

View File

@@ -31,10 +31,6 @@ const (
// defaultRetentionSweepInterval is how often the retention
// reaper deletes events older than each webhook's RetentionDays.
defaultRetentionSweepInterval = time.Hour
// defaultSessionIdleTimeout is how long a session may go without
// authenticated activity before it expires.
defaultSessionIdleTimeout = 24 * time.Hour
)
// ErrInvalidEnvironment is returned when WEBHOOKER_ENVIRONMENT
@@ -64,10 +60,6 @@ type Config struct {
// RetentionSweepInterval is how often the retention reaper runs.
RetentionSweepInterval time.Duration
// SessionIdleTimeout is the sliding inactivity window after
// which a session expires. Non-positive disables idle expiry.
SessionIdleTimeout time.Duration
params *ConfigParams
log *slog.Logger
}
@@ -170,15 +162,6 @@ func New(lc fx.Lifecycle, params ConfigParams) (*Config, error) {
return nil, err
}
// Same fail-loud treatment for the session idle timeout.
sessionIdleTimeout, err := envDuration(
"SESSION_IDLE_TIMEOUT",
defaultSessionIdleTimeout,
)
if err != nil {
return nil, err
}
// Load configuration values from environment variables
s := &Config{
DataDir: envString("DATA_DIR"),
@@ -190,7 +173,6 @@ func New(lc fx.Lifecycle, params ConfigParams) (*Config, error) {
Port: envInt("PORT", defaultPort),
SentryDSN: envString("SENTRY_DSN"),
RetentionSweepInterval: retentionSweepInterval,
SessionIdleTimeout: sessionIdleTimeout,
log: log,
params: &params,
}

View File

@@ -163,7 +163,7 @@ func TestRetentionSweepInterval(t *testing.T) {
}
if tt.expectError {
expectStartupError(t)
testRetentionSweepIntervalError(t)
} else {
testRetentionSweepIntervalSuccess(t, tt.expected)
}
@@ -171,9 +171,7 @@ func TestRetentionSweepInterval(t *testing.T) {
}
}
// expectStartupError asserts that fx refuses to build the app,
// which is what a set-but-unparseable duration must cause.
func expectStartupError(t *testing.T) {
func testRetentionSweepIntervalError(t *testing.T) {
t.Helper()
var cfg *config.Config
@@ -217,82 +215,6 @@ func testRetentionSweepIntervalSuccess(
assert.Equal(t, expected, cfg.RetentionSweepInterval)
}
func TestSessionIdleTimeout(t *testing.T) {
tests := []struct {
name string
set bool
value string
expectError bool
expected time.Duration
}{
{
name: "unset uses default",
set: false,
expected: 24 * time.Hour,
},
{
name: "valid value is parsed",
set: true,
value: "30m",
expected: 30 * time.Minute,
},
{
name: "unparseable value fails startup",
set: true,
value: "not-a-duration",
expectError: true,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
// Cannot use t.Parallel() here because t.Setenv
// is incompatible with parallel subtests.
t.Setenv("WEBHOOKER_ENVIRONMENT", "dev")
if tt.set {
t.Setenv("SESSION_IDLE_TIMEOUT", tt.value)
} else {
require.NoError(t, os.Unsetenv(
"SESSION_IDLE_TIMEOUT",
))
}
if tt.expectError {
expectStartupError(t)
} else {
testSessionIdleTimeoutSuccess(t, tt.expected)
}
})
}
}
func testSessionIdleTimeoutSuccess(
t *testing.T,
expected time.Duration,
) {
t.Helper()
var cfg *config.Config
app := fxtest.New(
t,
fx.Provide(
globals.New,
logger.New,
config.New,
),
fx.Populate(&cfg),
)
require.NoError(t, app.Err())
app.RequireStart()
defer app.RequireStop()
assert.Equal(t, expected, cfg.SessionIdleTimeout)
}
func TestDefaultDataDir(t *testing.T) {
for _, env := range []string{"", "dev", "prod"} {
name := env

View File

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

View File

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

View File

@@ -94,6 +94,23 @@ type Notifier interface {
Notify(tasks []Task)
}
// WebhookEvictor releases the delivery engine's per-webhook
// state for a webhook that no longer needs it — currently the
// cached archive writer of the database target, whose open
// file handle would otherwise outlive the webhook.
//
// It is deliberately separate from Notifier and deliberately
// one method wide: archiving lifecycle is not notification, and
// a single-method interface keeps the handlers package free of
// any dependency on the engine's internals while staying
// trivially fakeable in tests.
//
// EvictWebhook never deletes an archive file. It is idempotent
// and is a no-op for a webhook with no engine state.
type WebhookEvictor interface {
EvictWebhook(webhookID string)
}
// EngineParams are the fx dependencies for the delivery
// engine.
type EngineParams struct {
@@ -127,6 +144,10 @@ type Engine struct {
// httpTarget is retained so tests can reach the HTTP
// target's shared client and circuit breakers.
httpTarget *httpTarget
// dbTarget is retained so the engine can reach the archive
// writer registry for webhook eviction and the idle sweep.
dbTarget *databaseTarget
}
// New creates and registers the delivery engine with the
@@ -182,6 +203,19 @@ func (e *Engine) Notify(tasks []Task) {
}
}
// EvictWebhook implements WebhookEvictor. It releases the
// engine's per-webhook archiving state: the database target's
// cached archive writer is dropped from the registry and its
// file handle closed. The archive file itself is left on disk
// — it is long-term storage the operator owns.
func (e *Engine) EvictWebhook(webhookID string) {
if e.dbTarget == nil {
return
}
e.dbTarget.evict(webhookID)
}
// ScheduleRetry schedules a task to be re-enqueued onto the
// retry channel after delay. It implements the Scheduler
// interface the targets use to own their durable retries.

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

@@ -186,10 +186,6 @@ func (s *Middleware) RequireAuth() func(http.Handler) http.Handler {
return
}
// IsAuthenticated also enforces both session expiry
// deadlines, so an idle-expired or absolutely-expired
// session lands here and is sent back to the login
// page.
if !s.session.IsAuthenticated(sess) {
s.log.Debug(
"auth middleware: unauthenticated request",
@@ -203,26 +199,6 @@ func (s *Middleware) RequireAuth() func(http.Handler) http.Handler {
return
}
// This request authenticated with the session, so it
// counts as activity: push the idle deadline forward.
// This is the only place sessions are refreshed, which
// is what keeps an unauthenticated request from
// extending someone else's session. Touch advances the
// idle clock only -- the absolute cap is untouched --
// and reports false when nothing changed, so most
// requests do not re-issue the cookie. Save before the
// handler runs, while the headers are still ours to
// write.
if s.session.Touch(sess) {
err = s.session.Save(r, w, sess)
if err != nil {
s.log.Error(
"auth middleware: failed to refresh session",
"error", err,
)
}
}
next.ServeHTTP(w, r)
})
}

View File

@@ -8,7 +8,6 @@ import (
"net/http/httptest"
"os"
"testing"
"time"
"github.com/gorilla/sessions"
"github.com/stretchr/testify/assert"
@@ -29,22 +28,6 @@ func testMiddleware(
) (*middleware.Middleware, *session.Session) {
t.Helper()
m, s, _ := testMiddlewareWithSessionClock(t, env, 0, nil)
return m, s
}
// testMiddlewareWithSessionClock is testMiddleware with a
// configurable session idle timeout and a manually advanced clock,
// for the session-expiry tests. A nil clock uses the real one.
func testMiddlewareWithSessionClock(
t *testing.T,
env string,
idleTimeout time.Duration,
clock *fakeClock,
) (*middleware.Middleware, *session.Session, *fakeClock) {
t.Helper()
log := slog.New(slog.NewTextHandler(
os.Stderr,
&slog.HandlerOptions{Level: slog.LevelDebug},
@@ -52,7 +35,6 @@ func testMiddlewareWithSessionClock(
cfg := &config.Config{
Environment: env,
SessionIdleTimeout: idleTimeout,
}
// Create a real session manager with a known key
@@ -71,40 +53,11 @@ func testMiddlewareWithSessionClock(
SameSite: http.SameSiteLaxMode,
}
var now func() time.Time
if clock != nil {
now = clock.Now
}
sessManager := session.NewForTest(store, cfg, log, key, now)
sessManager := session.NewForTest(store, cfg, log, key)
m := middleware.NewForTest(log, cfg, sessManager)
return m, sessManager, clock
}
// fakeClock is a manually advanced clock, so session expiry can be
// tested without sleeping.
type fakeClock struct {
t time.Time
}
func (c *fakeClock) Now() time.Time {
return c.t
}
func (c *fakeClock) Advance(d time.Duration) {
c.t = c.t.Add(d)
}
// newFakeClock returns a clock started at a fixed instant.
func newFakeClock() *fakeClock {
return &fakeClock{
t: time.Date(
2026, time.January, 2, 3, 4, 5, 0, time.UTC,
),
}
return m, sessManager
}
// --- Logging Middleware Tests ---
@@ -434,181 +387,6 @@ func TestRequireAuth_UnauthenticatedSession_RedirectsToLogin(
assert.Equal(t, "/pages/login", w.Header().Get("Location"))
}
// --- RequireAuth Session Expiry Tests ---
// loginCookies authenticates a new session and returns the cookies
// a browser would then send back.
func loginCookies(
t *testing.T,
sessManager *session.Session,
) []*http.Cookie {
t.Helper()
req := httptest.NewRequestWithContext(
context.Background(), http.MethodGet, "/login", nil)
w := httptest.NewRecorder()
sess, err := sessManager.Get(req)
require.NoError(t, err)
sessManager.SetUser(sess, "user-123", "testuser")
require.NoError(t, sessManager.Save(req, w, sess))
cookies := w.Result().Cookies()
require.NotEmpty(t, cookies, "session cookie should be set")
return cookies
}
// runAuthed sends a request carrying cookies through RequireAuth
// and reports whether the protected handler ran, plus the response.
func runAuthed(
t *testing.T,
m *middleware.Middleware,
cookies []*http.Cookie,
) (bool, *httptest.ResponseRecorder) {
t.Helper()
var called bool
handler := m.RequireAuth()(http.HandlerFunc(
func(_ http.ResponseWriter, _ *http.Request) {
called = true
},
))
req := httptest.NewRequestWithContext(
context.Background(),
http.MethodGet, "/dashboard", nil,
)
for _, c := range cookies {
req.AddCookie(c)
}
w := httptest.NewRecorder()
handler.ServeHTTP(w, req)
return called, w
}
// sessionCookies filters a response's cookies down to the session
// cookie, so tests can tell whether the session was re-issued.
func sessionCookies(
w *httptest.ResponseRecorder,
) []*http.Cookie {
var out []*http.Cookie
for _, c := range w.Result().Cookies() {
if c.Name == session.SessionName {
out = append(out, c)
}
}
return out
}
func TestRequireAuth_IdleExpiredSession_RedirectsToLogin(
t *testing.T,
) {
t.Parallel()
idle := time.Hour
m, sessManager, clock := testMiddlewareWithSessionClock(
t, config.EnvironmentDev, idle, newFakeClock(),
)
cookies := loginCookies(t, sessManager)
clock.Advance(idle)
called, w := runAuthed(t, m, cookies)
assert.False(
t, called,
"handler should not run for an idle-expired session",
)
assert.Equal(t, http.StatusSeeOther, w.Code)
assert.Equal(t, "/pages/login", w.Header().Get("Location"))
assert.Empty(
t, sessionCookies(w),
"an expired session must not be refreshed",
)
}
func TestRequireAuth_RefreshesIdleDeadlineOnActivity(
t *testing.T,
) {
t.Parallel()
idle := time.Hour
m, sessManager, clock := testMiddlewareWithSessionClock(
t, config.EnvironmentDev, idle, newFakeClock(),
)
cookies := loginCookies(t, sessManager)
// Activity halfway through the idle window.
clock.Advance(idle / 2)
called, w := runAuthed(t, m, cookies)
require.True(t, called, "handler should run while valid")
refreshed := sessionCookies(w)
require.NotEmpty(
t, refreshed,
"activity should re-issue the session cookie",
)
// Past the original deadline. The refreshed cookie is still
// good; the original one is not.
clock.Advance(idle - time.Second)
calledRefreshed, _ := runAuthed(t, m, refreshed)
assert.True(
t, calledRefreshed,
"refreshed session should outlive the original deadline",
)
calledStale, staleW := runAuthed(t, m, cookies)
assert.False(
t, calledStale,
"the pre-refresh cookie carries the old idle deadline",
)
assert.Equal(t, http.StatusSeeOther, staleW.Code)
}
func TestRequireAuth_UnauthenticatedRequestDoesNotRefresh(
t *testing.T,
) {
t.Parallel()
m, sessManager, _ := testMiddlewareWithSessionClock(
t, config.EnvironmentDev, time.Hour, newFakeClock(),
)
// A session cookie that exists but was never authenticated.
req := httptest.NewRequestWithContext(
context.Background(), http.MethodGet, "/setup", nil)
setupW := httptest.NewRecorder()
sess, err := sessManager.Get(req)
require.NoError(t, err)
require.NoError(t, sessManager.Save(req, setupW, sess))
cookies := setupW.Result().Cookies()
require.NotEmpty(t, cookies)
called, w := runAuthed(t, m, cookies)
assert.False(t, called)
assert.Empty(
t, sessionCookies(w),
"an unauthenticated request must not stamp the session",
)
}
// --- NoCache Middleware Tests ---
func TestNoCache_SetsHeaders(t *testing.T) {
@@ -701,7 +479,7 @@ func metricsAuthMiddleware(
store := sessions.NewCookieStore(key)
store.Options = &sessions.Options{Path: "/", MaxAge: 86400}
sessManager := session.NewForTest(store, cfg, log, key, nil)
sessManager := session.NewForTest(store, cfg, log, key)
return middleware.NewForTest(log, cfg, sessManager)
}

View File

@@ -10,7 +10,6 @@ import (
"log/slog"
"maps"
"net/http"
"time"
"github.com/gorilla/sessions"
"go.uber.org/fx"
@@ -33,18 +32,6 @@ const (
// status.
AuthenticatedKey = "authenticated"
// CreatedAtKey is the session key holding the Unix timestamp at
// which the session was authenticated. It anchors the ABSOLUTE
// expiry clock and is written exactly once, by SetUser. Nothing
// refreshes it: an absolute deadline that moved with activity
// would not be a cap at all.
CreatedAtKey = "created_at"
// LastSeenKey is the session key holding the Unix timestamp of
// the most recent authenticated request. It anchors the IDLE
// expiry clock and is pushed forward by Touch.
LastSeenKey = "last_seen"
// sessionKeyLength is the required length in bytes for the
// session authentication key.
sessionKeyLength = 32
@@ -54,19 +41,6 @@ const (
// secondsPerDay is the number of seconds in a day.
secondsPerDay = 86400
// sessionAbsoluteMaxAge is the hard upper bound on how long a
// session may live, measured from CreatedAtKey. Activity never
// extends it, so even a continuously used session ends here and
// the user has to authenticate again.
sessionAbsoluteMaxAge = sessionMaxAgeDays * secondsPerDay * time.Second
// idleRefreshDivisor rate-limits idle-deadline refreshes. Touch
// only rewrites LastSeenKey once the stored value is older than
// idleTimeout/idleRefreshDivisor, so an active session is
// re-saved at most this many times per idle window instead of
// once per request. See Touch for the tradeoff this buys.
idleRefreshDivisor = 10
)
// ErrSessionKeyLength is returned when the decoded session key
@@ -88,16 +62,6 @@ type Session struct {
key []byte // raw 32-byte auth key, also used for CSRF cookie signing
log *slog.Logger
config *config.Config
// idleTimeout is the sliding inactivity window. A session that
// sees no authenticated request within this window expires,
// independently of the absolute cap. Non-positive disables idle
// expiry and leaves sessionAbsoluteMaxAge as the only bound.
idleTimeout time.Duration
// now reads the current time. Injected so expiry can be tested
// without sleeping.
now func() time.Time
}
// New creates a new session manager. The cookie store is
@@ -111,8 +75,6 @@ func New(
s := &Session{
log: params.Logger.Get(),
config: params.Config,
idleTimeout: params.Config.SessionIdleTimeout,
now: time.Now,
}
lc.Append(fx.Hook{
@@ -187,98 +149,29 @@ func (s *Session) Save(
return sess.Save(r, w)
}
// SetUser sets the user information in the session. It starts both
// expiry clocks: CreatedAtKey (absolute, never refreshed again) and
// LastSeenKey (idle, refreshed by Touch).
// SetUser sets the user information in the session.
func (s *Session) SetUser(
sess *sessions.Session,
userID, username string,
) {
now := s.now().Unix()
sess.Values[UserIDKey] = userID
sess.Values[UsernameKey] = username
sess.Values[AuthenticatedKey] = true
sess.Values[CreatedAtKey] = now
sess.Values[LastSeenKey] = now
}
// ClearUser removes user information from the session, including
// both expiry timestamps.
// ClearUser removes user information from the session.
func (s *Session) ClearUser(sess *sessions.Session) {
delete(sess.Values, UserIDKey)
delete(sess.Values, UsernameKey)
delete(sess.Values, AuthenticatedKey)
delete(sess.Values, CreatedAtKey)
delete(sess.Values, LastSeenKey)
}
// sessionTime reads a Unix-second timestamp stored under key.
func sessionTime(
sess *sessions.Session,
key string,
) (time.Time, bool) {
secs, ok := sess.Values[key].(int64)
if !ok {
return time.Time{}, false
}
return time.Unix(secs, 0), true
}
// IsAuthenticated checks if the session has an authenticated user
// whose session has not passed either expiry deadline. Every
// authentication decision goes through here, so neither clock can
// be bypassed by a caller that forgets to check it.
// IsAuthenticated checks if the session has an authenticated
// user.
func (s *Session) IsAuthenticated(sess *sessions.Session) bool {
auth, ok := sess.Values[AuthenticatedKey].(bool)
if !ok || !auth {
return false
}
return !s.expired(sess)
}
// Touch records authenticated activity by pushing the IDLE deadline
// forward. It writes LastSeenKey only; CreatedAtKey is left alone so
// the absolute cap keeps counting down even for a user who never
// stops clicking.
//
// Callers must only invoke Touch for a request that authenticated
// with this session. Refreshing on an unauthenticated request would
// let anyone holding a stolen or abandoned cookie keep the session
// alive by polling a public endpoint. Touch enforces that itself by
// returning false for any session that is not currently
// authenticated and unexpired.
//
// To avoid re-encrypting and re-emitting the session cookie on every
// single request, the timestamp is advanced only once it is older
// than idleTimeout/idleRefreshDivisor. The tradeoff is that
// LastSeenKey lags real activity by up to that much, so a session
// can expire slightly early relative to the user's true last
// request -- never late.
//
// Touch reports whether it changed the session; only then does the
// caller need to save it.
func (s *Session) Touch(sess *sessions.Session) bool {
if s.idleTimeout <= 0 {
return false
}
if !s.IsAuthenticated(sess) {
return false
}
now := s.now()
lastSeen, ok := sessionTime(sess, LastSeenKey)
if ok && now.Sub(lastSeen) < s.idleTimeout/idleRefreshDivisor {
return false
}
sess.Values[LastSeenKey] = now.Unix()
return true
return ok && auth
}
// GetUserID retrieves the user ID from the session.
@@ -360,41 +253,3 @@ func (s *Session) Regenerate(
return newSess, nil
}
// expired reports whether the session has passed either of its two
// independent deadlines. They are deliberately kept apart:
//
// - the ABSOLUTE deadline is CreatedAtKey + sessionAbsoluteMaxAge.
// It is fixed at login and no amount of activity moves it.
// - the IDLE deadline is LastSeenKey + idleTimeout. Activity moves
// it forward via Touch.
//
// Whichever comes first ends the session.
//
// A session that claims to be authenticated but carries no
// timestamps predates this check; it is treated as expired so the
// user re-authenticates rather than being granted an unbounded
// session.
func (s *Session) expired(sess *sessions.Session) bool {
now := s.now()
createdAt, ok := sessionTime(sess, CreatedAtKey)
if !ok {
return true
}
if !now.Before(createdAt.Add(sessionAbsoluteMaxAge)) {
return true
}
if s.idleTimeout <= 0 {
return false
}
lastSeen, ok := sessionTime(sess, LastSeenKey)
if !ok {
return true
}
return !now.Before(lastSeen.Add(s.idleTimeout))
}

View File

@@ -7,7 +7,6 @@ import (
"net/http/httptest"
"os"
"testing"
"time"
"github.com/gorilla/sessions"
"github.com/stretchr/testify/assert"
@@ -18,47 +17,11 @@ import (
const testKeySize = 32
// testIdleTimeout is the idle window used by the expiry tests.
const testIdleTimeout = time.Hour
// testAbsoluteMaxAge restates the documented absolute session cap
// independently of the implementation constant.
const testAbsoluteMaxAge = 7 * 24 * time.Hour
// fakeClock is a manually advanced clock, so expiry can be tested
// without sleeping.
type fakeClock struct {
t time.Time
}
func (c *fakeClock) Now() time.Time {
return c.t
}
func (c *fakeClock) Advance(d time.Duration) {
c.t = c.t.Add(d)
}
// testSession creates a Session with a real cookie store and the
// real clock.
// testSession creates a Session with a real cookie store for
// testing.
func testSession(t *testing.T) *session.Session {
t.Helper()
s, _ := testSessionWithClock(t, testIdleTimeout, nil)
return s
}
// testSessionWithClock creates a Session with a real cookie store,
// the given idle timeout, and a manually advanced clock. Passing a
// nil clock uses the real one.
func testSessionWithClock(
t *testing.T,
idleTimeout time.Duration,
clock *fakeClock,
) (*session.Session, *fakeClock) {
t.Helper()
key := make([]byte, testKeySize)
for i := range key {
@@ -76,7 +39,6 @@ func testSessionWithClock(
cfg := &config.Config{
Environment: config.EnvironmentDev,
SessionIdleTimeout: idleTimeout,
}
log := slog.New(slog.NewTextHandler(
@@ -84,46 +46,7 @@ func testSessionWithClock(
&slog.HandlerOptions{Level: slog.LevelDebug},
))
var now func() time.Time
if clock != nil {
now = clock.Now
}
return session.NewForTest(store, cfg, log, key, now), clock
}
// newFakeClock returns a clock started at a fixed instant.
func newFakeClock() *fakeClock {
return &fakeClock{
t: time.Date(
2026, time.January, 2, 3, 4, 5, 0, time.UTC,
),
}
}
// authenticatedSession returns a fresh session that has just been
// logged in, along with its manager and clock.
func authenticatedSession(
t *testing.T,
idleTimeout time.Duration,
) (*session.Session, *sessions.Session, *fakeClock) {
t.Helper()
s, clock := testSessionWithClock(
t, idleTimeout, newFakeClock(),
)
req := httptest.NewRequestWithContext(
context.Background(), http.MethodGet, "/", nil)
sess, err := s.Get(req)
require.NoError(t, err)
s.SetUser(sess, "user-123", "alice")
require.True(t, s.IsAuthenticated(sess))
return s, sess, clock
return session.NewForTest(store, cfg, log, key)
}
// --- Get and Save Tests ---
@@ -507,263 +430,6 @@ func TestSessionConstants(t *testing.T) {
assert.Equal(t, "user_id", session.UserIDKey)
assert.Equal(t, "username", session.UsernameKey)
assert.Equal(t, "authenticated", session.AuthenticatedKey)
assert.Equal(t, "created_at", session.CreatedAtKey)
assert.Equal(t, "last_seen", session.LastSeenKey)
}
// --- Expiry Tests ---
func TestSetUser_StartsBothClocks(t *testing.T) {
t.Parallel()
_, sess, clock := authenticatedSession(t, testIdleTimeout)
assert.Equal(
t, clock.Now().Unix(), sess.Values[session.CreatedAtKey],
"SetUser should anchor the absolute clock",
)
assert.Equal(
t, clock.Now().Unix(), sess.Values[session.LastSeenKey],
"SetUser should anchor the idle clock",
)
}
func TestIsAuthenticated_WithinIdleWindow(t *testing.T) {
t.Parallel()
s, sess, clock := authenticatedSession(t, testIdleTimeout)
clock.Advance(testIdleTimeout - time.Second)
assert.True(
t, s.IsAuthenticated(sess),
"session should still be valid just inside the idle window",
)
}
func TestIsAuthenticated_IdleExpired(t *testing.T) {
t.Parallel()
s, sess, clock := authenticatedSession(t, testIdleTimeout)
clock.Advance(testIdleTimeout)
assert.False(
t, s.IsAuthenticated(sess),
"session should expire once the idle window lapses",
)
}
// TestTouch_DoesNotExtendAbsoluteCap is the regression test for the
// refresh-the-wrong-clock bug: a session that is used continuously
// must survive well past the idle window and still die at the
// absolute cap.
func TestTouch_DoesNotExtendAbsoluteCap(t *testing.T) {
t.Parallel()
s, sess, clock := authenticatedSession(t, testIdleTimeout)
createdAt := sess.Values[session.CreatedAtKey]
// Stay active: a request every half idle window, right up to
// the absolute cap.
step := testIdleTimeout / 2
steps := int(testAbsoluteMaxAge/step) - 1
for i := range steps {
clock.Advance(step)
s.Touch(sess)
require.True(
t, s.IsAuthenticated(sess),
"active session should survive the idle window "+
"(step %d of %d)", i+1, steps,
)
}
// One more step of activity takes the session to exactly the
// absolute cap, measured from login. Nothing that happened in
// the loop may have moved that deadline.
clock.Advance(step)
s.Touch(sess)
assert.False(
t, s.IsAuthenticated(sess),
"activity must not extend the absolute cap",
)
assert.Equal(
t, createdAt, sess.Values[session.CreatedAtKey],
"Touch must never rewrite the absolute-clock anchor",
)
}
func TestTouch_RefreshesIdleDeadline(t *testing.T) {
t.Parallel()
s, sess, clock := authenticatedSession(t, testIdleTimeout)
// Halfway through the window, activity happens.
clock.Advance(testIdleTimeout / 2)
assert.True(
t, s.Touch(sess),
"Touch should refresh once past the lazy-refresh threshold",
)
// Past the original deadline, but inside the refreshed one.
clock.Advance(testIdleTimeout - time.Second)
assert.True(
t, s.IsAuthenticated(sess),
"refreshed session should outlive the original deadline",
)
// And it still expires an idle window after that activity.
clock.Advance(time.Second)
assert.False(
t, s.IsAuthenticated(sess),
"refreshed session should expire one window after activity",
)
}
func TestTouch_LazyBelowRefreshThreshold(t *testing.T) {
t.Parallel()
s, sess, clock := authenticatedSession(t, testIdleTimeout)
before := sess.Values[session.LastSeenKey]
// A request arriving almost immediately is not worth a cookie
// rewrite.
clock.Advance(time.Second)
assert.False(
t, s.Touch(sess),
"Touch should not rewrite the session below the threshold",
)
assert.Equal(
t, before, sess.Values[session.LastSeenKey],
"last-seen should be unchanged below the threshold",
)
}
func TestTouch_UnauthenticatedSessionIsNotRefreshed(t *testing.T) {
t.Parallel()
s, clock := testSessionWithClock(
t, testIdleTimeout, newFakeClock(),
)
req := httptest.NewRequestWithContext(
context.Background(), http.MethodGet, "/", nil)
sess, err := s.Get(req)
require.NoError(t, err)
clock.Advance(testIdleTimeout / 2)
assert.False(
t, s.Touch(sess),
"an unauthenticated session must not be refreshed",
)
_, hasLastSeen := sess.Values[session.LastSeenKey]
assert.False(
t, hasLastSeen,
"Touch must not stamp an unauthenticated session",
)
}
func TestTouch_IdleExpiredSessionIsNotRevived(t *testing.T) {
t.Parallel()
s, sess, clock := authenticatedSession(t, testIdleTimeout)
clock.Advance(testIdleTimeout)
require.False(t, s.IsAuthenticated(sess))
assert.False(
t, s.Touch(sess),
"an already expired session must not be refreshed",
)
assert.False(
t, s.IsAuthenticated(sess),
"Touch must not revive an expired session",
)
}
func TestIsAuthenticated_MissingTimestamps(t *testing.T) {
t.Parallel()
s, _ := testSessionWithClock(
t, testIdleTimeout, newFakeClock(),
)
req := httptest.NewRequestWithContext(
context.Background(), http.MethodGet, "/", nil)
sess, err := s.Get(req)
require.NoError(t, err)
// A session from before idle expiry existed: authenticated,
// but with no timestamps. Fail closed.
sess.Values[session.AuthenticatedKey] = true
assert.False(
t, s.IsAuthenticated(sess),
"a session with no timestamps should be rejected",
)
}
func TestIsAuthenticated_MissingLastSeen(t *testing.T) {
t.Parallel()
s, sess, _ := authenticatedSession(t, testIdleTimeout)
delete(sess.Values, session.LastSeenKey)
assert.False(
t, s.IsAuthenticated(sess),
"a session with no idle anchor should be rejected",
)
}
func TestIdleTimeoutDisabled_AbsoluteCapStillApplies(t *testing.T) {
t.Parallel()
s, sess, clock := authenticatedSession(t, 0)
// Idle expiry is off, so an untouched session survives an
// arbitrary idle stretch.
clock.Advance(testAbsoluteMaxAge - time.Second)
assert.True(
t, s.IsAuthenticated(sess),
"idle expiry should be disabled by a non-positive timeout",
)
assert.False(
t, s.Touch(sess),
"Touch should be a no-op when idle expiry is disabled",
)
// The absolute cap still ends it.
clock.Advance(time.Second)
assert.False(
t, s.IsAuthenticated(sess),
"the absolute cap must still apply with idle expiry off",
)
}
func TestClearUser_RemovesTimestamps(t *testing.T) {
t.Parallel()
s, sess, _ := authenticatedSession(t, testIdleTimeout)
s.ClearUser(sess)
_, hasCreatedAt := sess.Values[session.CreatedAtKey]
assert.False(t, hasCreatedAt, "CreatedAtKey should be removed")
_, hasLastSeen := sess.Values[session.LastSeenKey]
assert.False(t, hasLastSeen, "LastSeenKey should be removed")
}
// --- Edge Cases ---

View File

@@ -2,7 +2,6 @@ package session
import (
"log/slog"
"time"
"github.com/gorilla/sessions"
"sneak.berlin/go/webhooker/internal/config"
@@ -13,28 +12,16 @@ import (
// middleware and handler tests to use real session functionality. The key
// parameter is the raw 32-byte authentication key used for session encryption
// and CSRF cookie signing.
//
// The idle timeout is taken from cfg.SessionIdleTimeout, exactly as in
// production. The now parameter supplies the clock used for expiry
// checks so tests can advance time without sleeping; pass nil for the
// real clock.
func NewForTest(
store *sessions.CookieStore,
cfg *config.Config,
log *slog.Logger,
key []byte,
now func() time.Time,
) *Session {
if now == nil {
now = time.Now
}
return &Session{
store: store,
key: key,
config: cfg,
log: log,
idleTimeout: cfg.SessionIdleTimeout,
now: now,
}
}