check / check (push) Waiting to run
The HTTP drain at shutdown waited up to ShutdownTimeout regardless of how much of the stop budget earlier hooks had used, so a slow archive sweeper or retention reaper could eat the reserve the hooks after the server need, and the database close was skipped. The drain now waits at most the shorter of ShutdownTimeout and what is left of the budget less TailHookReserve, as the Sentry flush already does. The reserve is documented as derived from the two timeouts. Tests cover earlier hooks having spent part of the budget, on a clock that host speed cannot move, and pin that a drain on the full budget gets all of ShutdownTimeout. Model: opus-5-5
292 lines
9.2 KiB
Go
292 lines
9.2 KiB
Go
package main
|
|
|
|
import (
|
|
"bytes"
|
|
"encoding/json"
|
|
"io"
|
|
"log/slog"
|
|
"net"
|
|
"os"
|
|
"strconv"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
"sneak.berlin/go/webhooker/internal/config"
|
|
"sneak.berlin/go/webhooker/internal/datadir"
|
|
"sneak.berlin/go/webhooker/internal/resetpw"
|
|
"sneak.berlin/go/webhooker/internal/server"
|
|
)
|
|
|
|
// dockerStopGrace is Docker's default `docker stop` grace period.
|
|
// The Dockerfile sets no STOPSIGNAL or grace override, so this is
|
|
// the deadline the container is actually held to, and the fx stop
|
|
// timeout has to fit inside it with room for signal delivery and
|
|
// process exit.
|
|
const dockerStopGrace = 10 * time.Second
|
|
|
|
// TestNewApp_StopTimeout pins the fx stop timeout. Without the
|
|
// explicit fx.StopTimeout option the app reads fx's 15s
|
|
// DefaultTimeout, which exceeds dockerStopGrace: the container is
|
|
// SIGKILLed before the bound fires and every shutdown hook bounded
|
|
// by it — including the operator-facing timeout log — becomes
|
|
// unreachable in the image this repo produces.
|
|
//
|
|
// fx.New applies options before it executes invokes, so the timeout
|
|
// is set whether or not the graph itself can be constructed here.
|
|
func TestNewApp_StopTimeout(t *testing.T) {
|
|
config.ClearEnvForTest(t)
|
|
t.Setenv("DATA_DIR", t.TempDir())
|
|
|
|
got := newApp().StopTimeout()
|
|
|
|
require.Equal(t, stopTimeout, got)
|
|
require.Less(t, got, dockerStopGrace)
|
|
}
|
|
|
|
// freePort returns a loopback TCP port that was free a moment ago, by
|
|
// taking one and releasing it.
|
|
func freePort(t *testing.T) int {
|
|
t.Helper()
|
|
|
|
var listenCfg net.ListenConfig
|
|
|
|
l, err := listenCfg.Listen(t.Context(), "tcp", "127.0.0.1:0")
|
|
require.NoError(t, err)
|
|
|
|
addr, ok := l.Addr().(*net.TCPAddr)
|
|
require.True(t, ok, "listener is not TCP")
|
|
require.NoError(t, l.Close())
|
|
|
|
return addr.Port
|
|
}
|
|
|
|
// TestNewApp_SendsFxEventsToTheLogger starts and stops the app main
|
|
// runs, with DEBUG=true, and reads back what reached the service's
|
|
// logger. fx's own events must arrive there as structured records:
|
|
// the start at INFO, and at DEBUG the records of how the graph was
|
|
// built.
|
|
//
|
|
// fx holds its events back until its logger is built and then replays
|
|
// them all at once, so the earliest of them arriving shows the replay
|
|
// ran at DEBUG: that globals.New was provided, which fx records before
|
|
// anything is built, and the run of logger.New, which happens before
|
|
// the configuration sets the level.
|
|
func TestNewApp_SendsFxEventsToTheLogger(t *testing.T) {
|
|
config.ClearEnvForTest(t)
|
|
t.Setenv("DATA_DIR", t.TempDir())
|
|
t.Setenv("PORT", strconv.Itoa(freePort(t)))
|
|
t.Setenv("DEBUG", "true")
|
|
|
|
// internal/logger writes to whatever os.Stdout is when it builds
|
|
// its handler. A file is not a terminal, so that handler is the
|
|
// JSON one the service uses in production.
|
|
out, err := os.CreateTemp(t.TempDir(), "stdout")
|
|
require.NoError(t, err)
|
|
|
|
stdout := os.Stdout
|
|
os.Stdout = out
|
|
|
|
t.Cleanup(func() {
|
|
os.Stdout = stdout
|
|
_ = out.Close()
|
|
})
|
|
|
|
app := newApp()
|
|
require.NoError(t, app.Start(t.Context()))
|
|
require.NoError(t, app.Stop(t.Context()))
|
|
|
|
_, err = out.Seek(0, io.SeekStart)
|
|
require.NoError(t, err)
|
|
|
|
written, err := io.ReadAll(out)
|
|
require.NoError(t, err)
|
|
|
|
type record struct {
|
|
Level string `json:"level"`
|
|
Msg string `json:"msg"`
|
|
Name string `json:"name"`
|
|
Constructor string `json:"constructor"`
|
|
}
|
|
|
|
var records []record
|
|
|
|
for line := range strings.Lines(string(written)) {
|
|
var r record
|
|
|
|
// The first-boot banner is plain text, not a record.
|
|
if json.Unmarshal([]byte(line), &r) == nil {
|
|
records = append(records, r)
|
|
}
|
|
}
|
|
|
|
const pkg = "sneak.berlin/go/webhooker/internal/"
|
|
|
|
info := slog.LevelInfo.String()
|
|
debug := slog.LevelDebug.String()
|
|
|
|
assert.Contains(t, records, record{Level: info, Msg: "started"})
|
|
assert.Contains(t, records, record{
|
|
Level: debug, Msg: "provided", Constructor: pkg + "globals.New()",
|
|
})
|
|
assert.Contains(t, records, record{
|
|
Level: debug, Msg: "run", Name: pkg + "logger.New()",
|
|
})
|
|
assert.Contains(t, records, record{Level: debug, Msg: "invoking"})
|
|
assert.Contains(t, records, record{
|
|
Level: debug, Msg: "initialized custom fxevent.Logger",
|
|
})
|
|
}
|
|
|
|
// TestRunRefusesLockedDataDir pins what an operator's second start
|
|
// does. The entry point must refuse before it builds the fx graph —
|
|
// nothing may open a database in a DATA_DIR another process holds —
|
|
// and must exit non-zero with a message naming the directory rather
|
|
// than starting a second delivery engine over the same rows.
|
|
//
|
|
// flock(2) locks descriptors independently, so holding the lock here
|
|
// is the same denial a separate process gets; internal/datadir pins
|
|
// that property and covers the real two-process case.
|
|
func TestRunRefusesLockedDataDir(t *testing.T) {
|
|
dir := t.TempDir()
|
|
t.Setenv("DATA_DIR", dir)
|
|
|
|
lock, err := datadir.Acquire(dir)
|
|
require.NoError(t, err)
|
|
|
|
defer func() { _ = lock.Release() }()
|
|
|
|
var stderr bytes.Buffer
|
|
|
|
code := run(&stderr)
|
|
|
|
require.Equal(
|
|
t, 1, code, "a second instance must exit non-zero",
|
|
)
|
|
assert.Contains(
|
|
t, stderr.String(), dir,
|
|
"the refusal must name the directory",
|
|
)
|
|
assert.Contains(t, stderr.String(), "another instance")
|
|
}
|
|
|
|
// TestDispatch_NoArgumentsRunsTheServer pins the routing of a bare
|
|
// invocation, which is what the image's CMD and every deployment use.
|
|
// Adding subcommands must not move the server off the empty argument
|
|
// list, and must not move the DATA_DIR lock: this asserts the refusal
|
|
// arrives with no fx graph built, exactly as run does on its own.
|
|
func TestDispatch_NoArgumentsRunsTheServer(t *testing.T) {
|
|
dir := t.TempDir()
|
|
t.Setenv("DATA_DIR", dir)
|
|
|
|
lock, err := datadir.Acquire(dir)
|
|
require.NoError(t, err)
|
|
|
|
defer func() { _ = lock.Release() }()
|
|
|
|
var stdout, stderr bytes.Buffer
|
|
|
|
code := dispatch(nil, strings.NewReader(""), &stdout, &stderr)
|
|
|
|
require.Equal(t, 1, code)
|
|
assert.Contains(t, stderr.String(), "another instance")
|
|
}
|
|
|
|
// TestDispatch_UnknownSubcommand keeps a mistyped subcommand from
|
|
// starting a server. Anything else would have `webhooker resetpww`
|
|
// silently take the DATA_DIR lock and serve.
|
|
func TestDispatch_UnknownSubcommand(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
var stdout, stderr bytes.Buffer
|
|
|
|
code := dispatch(
|
|
[]string{"resetpww", "admin"},
|
|
strings.NewReader(""), &stdout, &stderr,
|
|
)
|
|
|
|
require.Equal(t, 2, code)
|
|
assert.Contains(t, stderr.String(), "unknown subcommand")
|
|
assert.Contains(
|
|
t, stderr.String(), resetpw.Name,
|
|
"the usage must name the subcommand that does exist",
|
|
)
|
|
}
|
|
|
|
// TestDispatch_Help answers on standard output with a zero status, so
|
|
// `webhooker help` is usable in a pipe.
|
|
func TestDispatch_Help(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
var stdout, stderr bytes.Buffer
|
|
|
|
code := dispatch(
|
|
[]string{helpCommand}, strings.NewReader(""), &stdout, &stderr,
|
|
)
|
|
|
|
require.Equal(t, 0, code)
|
|
assert.Empty(t, stderr.String())
|
|
assert.Contains(t, stdout.String(), resetpw.Name)
|
|
}
|
|
|
|
// tailHeadroom is the slack the fx stop budget must keep beyond the
|
|
// server stop hook. The hooks that run after the server — the
|
|
// delivery engine, the healthcheck, the webhook DB manager and the
|
|
// database close — are microsecond-scale in normal operation, so
|
|
// this is generous for them.
|
|
const tailHeadroom = 2 * time.Second
|
|
|
|
// TestStopTimeout_LeavesHeadroomForTailHooks pins the relationship
|
|
// between the server's stop hook and the fx stop budget. fx bounds
|
|
// the whole stop sequence, and returns without running its
|
|
// remaining hooks once the stop context has expired. If the hook
|
|
// could use the entire budget, every later hook — the database close
|
|
// included — would be skipped in exactly the case where the drain
|
|
// mattered.
|
|
//
|
|
// The hook is not just the HTTP drain: a Sentry flush follows it in
|
|
// the same hook, and sentry.Flush honours no context, so both halves
|
|
// have to be counted. The sweep walks every drain length the hook
|
|
// can produce, since a shorter drain leaves the flush more room and
|
|
// the worst case is not necessarily at either extreme.
|
|
//
|
|
// Nor does the hook start on a full budget: the ArchiveSweeper and
|
|
// RetentionReaper hooks run before it, and whatever they spent is
|
|
// gone. The outer sweep walks every amount they can spend. Once they
|
|
// have eaten into the headroom themselves, the hook must spend
|
|
// nothing of what is left. A drain that starts on the full budget
|
|
// must still get all of ShutdownTimeout, so a smaller stopTimeout
|
|
// cannot silently shorten every drain.
|
|
//
|
|
// Shrinking either budget, or unbounding the drain or the flush
|
|
// again, must fail here rather than silently recreating a hook that
|
|
// swallows the whole sequence.
|
|
func TestStopTimeout_LeavesHeadroomForTailHooks(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
require.Less(t, server.ShutdownTimeout, stopTimeout)
|
|
require.Equal(
|
|
t, server.ShutdownTimeout, server.DrainBudget(stopTimeout),
|
|
"a drain that starts on the full stop budget is cut short",
|
|
)
|
|
|
|
const step = 10 * time.Millisecond
|
|
|
|
for spent := time.Duration(0); spent <= stopTimeout; spent += step {
|
|
remaining := stopTimeout - spent
|
|
longest := max(server.DrainBudget(remaining), 0)
|
|
|
|
for drain := time.Duration(0); drain <= longest; drain += step {
|
|
hook := drain + server.SentryFlushBudget(remaining-drain)
|
|
|
|
require.GreaterOrEqual(
|
|
t, remaining-hook, min(remaining, tailHeadroom),
|
|
"a %s drain after %s of earlier hooks leaves "+
|
|
"the tail hooks short", drain, spent,
|
|
)
|
|
}
|
|
}
|
|
}
|