watcher: save state when it stops and wait for that save (closes #114) #186
@@ -618,9 +618,10 @@ repository's `Dockerfile` and runs it. The app needs:
|
|||||||
from a previous cycle.
|
from a previous cycle.
|
||||||
4. **On change detection**: Send notifications to all configured
|
4. **On change detection**: Send notifications to all configured
|
||||||
endpoints, update in-memory state, persist to disk.
|
endpoints, update in-memory state, persist to disk.
|
||||||
5. **Shutdown**: Persist final state to disk, wait for in-flight
|
5. **Shutdown**: The watcher stops checking and saves the final state
|
||||||
notification deliveries to complete, stop gracefully. The wait is
|
to disk, and shutdown waits for that save before it goes on. Then it
|
||||||
bounded by the fx shutdown timeout (15s by default): deliveries still
|
waits for in-flight notification deliveries to complete. Both waits
|
||||||
|
share the fx shutdown timeout (15s by default): deliveries still
|
||||||
retrying against an unreachable endpoint when that expires are
|
retrying against an unreachable endpoint when that expires are
|
||||||
abandoned, and the number abandoned is logged at warn level rather
|
abandoned, and the number abandoned is logged at warn level rather
|
||||||
than dropped silently. Notifications generated after shutdown has
|
than dropped silently. Notifications generated after shutdown has
|
||||||
|
|||||||
@@ -19,6 +19,8 @@ nameserver IP address changes: https://git.eeqj.de/sneak/dnswatcher/issues/105
|
|||||||
|
|
||||||
# Completed Steps
|
# Completed Steps
|
||||||
|
|
||||||
|
- 2026-10-01: the watcher saves state when it stops, and shutdown waits for that
|
||||||
|
save, so it no longer relies on the state's own stop hook (closes #114).
|
||||||
- 2026-10-01: `DNSWATCHER_SENTRY_DSN` reports panics in HTTP handlers to Sentry,
|
- 2026-10-01: `DNSWATCHER_SENTRY_DSN` reports panics in HTTP handlers to Sentry,
|
||||||
and a DSN Sentry cannot parse stops startup (closes #107).
|
and a DSN Sentry cannot parse stops startup (closes #107).
|
||||||
- 2026-10-01: a port or TLS check that shutdown cuts short saves nothing and
|
- 2026-10-01: a port or TLS check that shutdown cuts short saves nothing and
|
||||||
@@ -109,7 +111,6 @@ nameserver IP address changes: https://git.eeqj.de/sneak/dnswatcher/issues/105
|
|||||||
https://git.eeqj.de/sneak/dnswatcher/issues/66
|
https://git.eeqj.de/sneak/dnswatcher/issues/66
|
||||||
- `goimports` in `make fmt-check`, Markdown formatting:
|
- `goimports` in `make fmt-check`, Markdown formatting:
|
||||||
https://git.eeqj.de/sneak/dnswatcher/issues/119
|
https://git.eeqj.de/sneak/dnswatcher/issues/119
|
||||||
- final state save at shutdown: https://git.eeqj.de/sneak/dnswatcher/issues/114
|
|
||||||
- README accuracy sweep: https://git.eeqj.de/sneak/dnswatcher/issues/108
|
- README accuracy sweep: https://git.eeqj.de/sneak/dnswatcher/issues/108
|
||||||
- README sections required by policy:
|
- README sections required by policy:
|
||||||
https://git.eeqj.de/sneak/dnswatcher/issues/173
|
https://git.eeqj.de/sneak/dnswatcher/issues/173
|
||||||
|
|||||||
+31
-13
@@ -56,6 +56,7 @@ type Watcher struct {
|
|||||||
tlsCheck TLSChecker
|
tlsCheck TLSChecker
|
||||||
notify Notifier
|
notify Notifier
|
||||||
cancel context.CancelFunc
|
cancel context.CancelFunc
|
||||||
|
done chan struct{} // closed when Run returns
|
||||||
firstRun bool
|
firstRun bool
|
||||||
expiryNotifiedMu sync.Mutex
|
expiryNotifiedMu sync.Mutex
|
||||||
expiryNotified map[string]time.Time
|
expiryNotified map[string]time.Time
|
||||||
@@ -79,31 +80,47 @@ func New(
|
|||||||
}
|
}
|
||||||
|
|
||||||
lifecycle.Append(fx.Hook{
|
lifecycle.Append(fx.Hook{
|
||||||
OnStart: func(_ context.Context) error {
|
OnStart: func(startCtx context.Context) error {
|
||||||
// Use context.Background() — the fx startup context
|
// The fx startup context expires after startup
|
||||||
// expires after startup completes, so deriving from it
|
// completes, so the watcher's context drops its
|
||||||
// would cancel the watcher immediately. The watcher's
|
// cancellation. The watcher's lifetime is controlled
|
||||||
// lifetime is controlled by w.cancel in OnStop.
|
// by w.cancel in OnStop.
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(
|
||||||
|
context.WithoutCancel(startCtx),
|
||||||
|
)
|
||||||
w.cancel = cancel
|
w.cancel = cancel
|
||||||
|
w.done = make(chan struct{})
|
||||||
|
|
||||||
go w.Run(ctx) //nolint:contextcheck // intentionally not derived from startCtx
|
go func() {
|
||||||
|
defer close(w.done)
|
||||||
|
|
||||||
|
w.Run(ctx)
|
||||||
|
}()
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
},
|
},
|
||||||
OnStop: func(_ context.Context) error {
|
OnStop: func(ctx context.Context) error {
|
||||||
if w.cancel != nil {
|
w.cancel()
|
||||||
w.cancel()
|
|
||||||
}
|
|
||||||
|
|
||||||
return nil
|
// Run saves state as it returns. Waiting for it here
|
||||||
|
// means the save is done before shutdown goes on.
|
||||||
|
select {
|
||||||
|
case <-w.done:
|
||||||
|
return nil
|
||||||
|
case <-ctx.Done():
|
||||||
|
return fmt.Errorf(
|
||||||
|
"waiting for the watcher to stop: %w",
|
||||||
|
ctx.Err(),
|
||||||
|
)
|
||||||
|
}
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
return w, nil
|
return w, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Run starts the monitoring loop with periodic scheduling.
|
// Run starts the monitoring loop with periodic scheduling. When ctx
|
||||||
|
// is cancelled, it saves state and returns.
|
||||||
func (w *Watcher) Run(ctx context.Context) {
|
func (w *Watcher) Run(ctx context.Context) {
|
||||||
w.log.Info(
|
w.log.Info(
|
||||||
"watcher starting",
|
"watcher starting",
|
||||||
@@ -125,6 +142,7 @@ func (w *Watcher) Run(ctx context.Context) {
|
|||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
|
w.saveState()
|
||||||
w.log.Info("watcher stopped")
|
w.log.Info("watcher stopped")
|
||||||
|
|
||||||
return
|
return
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
|
"os"
|
||||||
"slices"
|
"slices"
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
@@ -135,6 +136,7 @@ type testDeps struct {
|
|||||||
notifier *mockNotifier
|
notifier *mockNotifier
|
||||||
state *state.State
|
state *state.State
|
||||||
config *config.Config
|
config *config.Config
|
||||||
|
log *logger.Logger
|
||||||
}
|
}
|
||||||
|
|
||||||
func newTestWatcher(
|
func newTestWatcher(
|
||||||
@@ -143,6 +145,23 @@ func newTestWatcher(
|
|||||||
) (*watcher.Watcher, *testDeps) {
|
) (*watcher.Watcher, *testDeps) {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
|
deps := newTestDeps(t, cfg)
|
||||||
|
|
||||||
|
w := watcher.NewForTest(
|
||||||
|
deps.config,
|
||||||
|
deps.state,
|
||||||
|
resolver.NewFromLogger(slog.Default()),
|
||||||
|
deps.portChecker,
|
||||||
|
deps.tlsChecker,
|
||||||
|
deps.notifier,
|
||||||
|
)
|
||||||
|
|
||||||
|
return w, deps
|
||||||
|
}
|
||||||
|
|
||||||
|
func newTestDeps(t *testing.T, cfg *config.Config) *testDeps {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
deps := &testDeps{
|
deps := &testDeps{
|
||||||
portChecker: &mockPortChecker{},
|
portChecker: &mockPortChecker{},
|
||||||
tlsChecker: &mockTLSChecker{
|
tlsChecker: &mockTLSChecker{
|
||||||
@@ -157,30 +176,21 @@ func newTestWatcher(
|
|||||||
t.Fatalf("globals.New: %v", err)
|
t.Fatalf("globals.New: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
log, err := logger.New(nil, logger.Params{Globals: g})
|
deps.log, err = logger.New(nil, logger.Params{Globals: g})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("logger.New: %v", err)
|
t.Fatalf("logger.New: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// The watcher saves state after every check, into cfg.DataDir.
|
// The watcher saves state after every check, into cfg.DataDir.
|
||||||
deps.state, err = state.New(fxtest.NewLifecycle(t), state.Params{
|
deps.state, err = state.New(fxtest.NewLifecycle(t), state.Params{
|
||||||
Logger: log,
|
Logger: deps.log,
|
||||||
Config: cfg,
|
Config: cfg,
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("state.New: %v", err)
|
t.Fatalf("state.New: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
w := watcher.NewForTest(
|
return deps
|
||||||
deps.config,
|
|
||||||
deps.state,
|
|
||||||
resolver.NewFromLogger(slog.Default()),
|
|
||||||
deps.portChecker,
|
|
||||||
deps.tlsChecker,
|
|
||||||
deps.notifier,
|
|
||||||
)
|
|
||||||
|
|
||||||
return w, deps
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func defaultTestConfig(t *testing.T) *config.Config {
|
func defaultTestConfig(t *testing.T) *config.Config {
|
||||||
@@ -539,6 +549,85 @@ func TestGracefulShutdown(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestStopSavesState stops a watcher built by New the way fx stops it,
|
||||||
|
// and checks that a change made to the state after the last check is in
|
||||||
|
// the state file afterwards. The state's own stop hook never runs here,
|
||||||
|
// so only the watcher can have saved it. Nothing is configured to
|
||||||
|
// check, so no DNS is involved.
|
||||||
|
func TestStopSavesState(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
cfg := defaultTestConfig(t)
|
||||||
|
deps := newTestDeps(t, cfg)
|
||||||
|
lc := fxtest.NewLifecycle(t)
|
||||||
|
|
||||||
|
_, err := watcher.New(lc, watcher.Params{
|
||||||
|
Logger: deps.log,
|
||||||
|
Config: cfg,
|
||||||
|
State: deps.state,
|
||||||
|
Resolver: resolver.NewFromLogger(slog.Default()),
|
||||||
|
PortCheck: deps.portChecker,
|
||||||
|
TLSCheck: deps.tlsChecker,
|
||||||
|
Notify: deps.notifier,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("watcher.New: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
lc.RequireStart()
|
||||||
|
|
||||||
|
// The first check saves state once. Wait for that save before
|
||||||
|
// changing the state, so the change can reach the file only
|
||||||
|
// through the save made at stop.
|
||||||
|
deadline := time.Now().Add(5 * time.Second)
|
||||||
|
|
||||||
|
for {
|
||||||
|
_, err = os.Stat(cfg.StatePath())
|
||||||
|
if err == nil {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
|
||||||
|
if time.Now().After(deadline) {
|
||||||
|
t.Fatalf("the first check saved no state: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
time.Sleep(10 * time.Millisecond)
|
||||||
|
}
|
||||||
|
|
||||||
|
deps.state.SetDomainState(testDomain, &state.DomainState{
|
||||||
|
Nameservers: []string{oldNS1},
|
||||||
|
})
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
err = lc.Stop(ctx)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("stopping the watcher: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
saved, err := state.New(fxtest.NewLifecycle(t), state.Params{
|
||||||
|
Logger: deps.log,
|
||||||
|
Config: cfg,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("state.New: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = saved.Load()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("loading the state file: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
ds, ok := saved.GetDomainState(testDomain)
|
||||||
|
if !ok || !slices.Equal(ds.Nameservers, []string{oldNS1}) {
|
||||||
|
t.Errorf(
|
||||||
|
"state file after stop has %+v for %s, want nameservers %v",
|
||||||
|
ds, testDomain, []string{oldNS1},
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestDNSRunsBeforePortAndTLSChecks(t *testing.T) {
|
func TestDNSRunsBeforePortAndTLSChecks(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user