Compare commits
2
Commits
6b0fdc477a
..
next
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
fc43f893a5 | ||
|
|
ae06f7e3a1 |
@@ -182,6 +182,27 @@ dnswatcher exposes a lightweight HTTP API for operational visibility:
|
|||||||
| `GET /api/v1/status` | Current monitoring state |
|
| `GET /api/v1/status` | Current monitoring state |
|
||||||
| `GET /metrics` | Prometheus metrics (optional) |
|
| `GET /metrics` | Prometheus metrics (optional) |
|
||||||
|
|
||||||
|
#### Server timeouts
|
||||||
|
|
||||||
|
The HTTP server sets all four socket-level timeouts. These are compile-time
|
||||||
|
constants in `internal/server/server.go`, not configurable via environment
|
||||||
|
variables.
|
||||||
|
|
||||||
|
| Timeout | Value | Purpose |
|
||||||
|
|---------------------|-------|-----------------------------------------------|
|
||||||
|
| `ReadHeaderTimeout` | 10s | Bounds the request header read (slowloris) |
|
||||||
|
| `ReadTimeout` | 15s | Bounds the whole request read, headers + body |
|
||||||
|
| `WriteTimeout` | 75s | Bounds handler execution plus response flush |
|
||||||
|
| `IdleTimeout` | 120s | Reaps idle keep-alive connections |
|
||||||
|
|
||||||
|
These are distinct from the 60s per-request handler budget applied by
|
||||||
|
`chimw.Timeout` in `internal/server/routes.go`, which cancels the request
|
||||||
|
context but does not touch the socket. `WriteTimeout` is deliberately
|
||||||
|
larger than that budget: the write deadline is armed once request headers
|
||||||
|
are read, so a smaller value would sever the connection before a handler
|
||||||
|
using its full budget could respond. `IdleTimeout` exceeds common
|
||||||
|
Prometheus scrape intervals so the scraper reuses its connection.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## Architecture
|
## Architecture
|
||||||
@@ -218,7 +239,8 @@ internal/
|
|||||||
- **Structured logging**: All logs use `log/slog` with JSON output in
|
- **Structured logging**: All logs use `log/slog` with JSON output in
|
||||||
production (TTY detection for development).
|
production (TTY detection for development).
|
||||||
- **Graceful shutdown**: All background goroutines respect context
|
- **Graceful shutdown**: All background goroutines respect context
|
||||||
cancellation and the fx lifecycle.
|
cancellation and the fx lifecycle. In-flight notification deliveries
|
||||||
|
are drained on shutdown, bounded by the shutdown timeout.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -463,8 +485,14 @@ docker run -d \
|
|||||||
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, complete in-flight
|
5. **Shutdown**: Persist final state to disk, wait for in-flight
|
||||||
notifications, stop gracefully.
|
notification deliveries to complete, stop gracefully. The wait is
|
||||||
|
bounded by the fx shutdown timeout (15s by default): deliveries still
|
||||||
|
retrying against an unreachable endpoint when that expires are
|
||||||
|
abandoned, and the number abandoned is logged at warn level rather
|
||||||
|
than dropped silently. Notifications generated after shutdown has
|
||||||
|
begun are refused and logged, so a late burst cannot extend the
|
||||||
|
shutdown.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
|
|||||||
+4
-58
@@ -2,10 +2,8 @@
|
|||||||
|
|
||||||
## DNS Resolution Tests
|
## DNS Resolution Tests
|
||||||
|
|
||||||
All tests that involve DNS resolution — in every package, including
|
All resolver tests **MUST** use live queries against real DNS servers.
|
||||||
consumers of the resolver such as the watcher — **MUST** use live
|
No mocking of the DNS client layer is permitted.
|
||||||
queries against real DNS servers. No mocking, faking, or stubbing of
|
|
||||||
DNS at any layer is permitted.
|
|
||||||
|
|
||||||
### Rationale
|
### Rationale
|
||||||
|
|
||||||
@@ -14,8 +12,6 @@ the full delegation chain. Mocked responses cannot faithfully represent
|
|||||||
the variety of real-world DNS behavior (truncation, referrals, glue
|
the variety of real-world DNS behavior (truncation, referrals, glue
|
||||||
records, DNSSEC, varied response times, EDNS, etc.). Testing against
|
records, DNSSEC, varied response times, EDNS, etc.). Testing against
|
||||||
real servers ensures the resolver works correctly in production.
|
real servers ensures the resolver works correctly in production.
|
||||||
Robustness comes from handling real-world DNS behavior with tolerant
|
|
||||||
assertions and sensible timeouts, not from mocks.
|
|
||||||
|
|
||||||
### Constraints
|
### Constraints
|
||||||
|
|
||||||
@@ -28,64 +24,14 @@ assertions and sensible timeouts, not from mocks.
|
|||||||
- Flaky failures from transient network issues are acceptable and
|
- Flaky failures from transient network issues are acceptable and
|
||||||
should be investigated as potential resolver bugs, not papered over
|
should be investigated as potential resolver bugs, not papered over
|
||||||
with mocks or skip flags
|
with mocks or skip flags
|
||||||
- Watcher change-detection tests seed a synthetic *previous state*
|
|
||||||
and compare it against fresh live lookups; the DNS side is never
|
|
||||||
faked
|
|
||||||
- Live query concurrency is bounded per package (`liveGate` in
|
|
||||||
`internal/resolver`, `liveWatcherGate` in `internal/watcher`) so
|
|
||||||
parallel tests do not burst at the root servers
|
|
||||||
- Those gates are package-scoped and therefore per test binary, so
|
|
||||||
`script/test` also passes `-p 1`: with both live-DNS packages
|
|
||||||
running at once the gates sum instead of holding, and the resolver's
|
|
||||||
per-attempt deadlines start expiring
|
|
||||||
|
|
||||||
### Transport failures: loopback nameservers, not mocks
|
|
||||||
|
|
||||||
The resolver classifies a nameserver that stays silent as
|
|
||||||
`StatusTimeout` and one that answers SERVFAIL as `StatusError`. The
|
|
||||||
public network cannot be made to produce either on demand — a
|
|
||||||
black-holed address is only black-holed on some networks, and build
|
|
||||||
environments that transparently intercept UDP/53 answer it locally —
|
|
||||||
so a test built on a chosen remote address asserts on the network it
|
|
||||||
happens to run on rather than on the resolver.
|
|
||||||
|
|
||||||
`internal/resolver/transport_test.go` binds a real nameserver on
|
|
||||||
`127.0.0.1` instead and points the query at it.
|
|
||||||
`nameserverAddr` dials an address that already carries a port as
|
|
||||||
written, so no production behaviour is bypassed to arrange this.
|
|
||||||
|
|
||||||
**This is permitted, and it is not a mock.** The rule above bans
|
|
||||||
substituting `DNSClient` or any other DNS abstraction, which lets the
|
|
||||||
code under test skip DNS and hands it a manufactured verdict. A
|
|
||||||
loopback nameserver does the opposite: the resolver dials a real
|
|
||||||
socket, writes a real query with the real `miekg/dns` client, and
|
|
||||||
applies its real deadline and its real classification logic to what
|
|
||||||
comes back. Choosing which nameserver a live query is sent to is not
|
|
||||||
faking DNS — the resolver is aimed at a nameserver of the caller's
|
|
||||||
choosing in production too.
|
|
||||||
|
|
||||||
The distinction to hold on to: **substituting the client is banned;
|
|
||||||
choosing the server is not.** A test that reaches for a fake
|
|
||||||
`DNSClient` to force a classification is still forbidden, no matter
|
|
||||||
how awkward the alternative looks.
|
|
||||||
|
|
||||||
Such a test must stay cheap. The resolver asks for eight record
|
|
||||||
types and retries each once, so a nameserver silent on every type
|
|
||||||
costs sixteen query timeouts. `TestQueryNameserverIP_Timeout` is
|
|
||||||
silent on `A` alone and answers the rest, which is all the resolver
|
|
||||||
needs to classify the response and keeps the test to two.
|
|
||||||
|
|
||||||
### What NOT to do
|
### What NOT to do
|
||||||
|
|
||||||
- **Do not mock `DNSClient`**, the watcher's `DNSResolver` interface,
|
- **Do not mock `DNSClient`** for resolver tests (the mock constructor
|
||||||
or any other DNS abstraction — in any package, for any reason
|
exists for unit-testing other packages that consume the resolver)
|
||||||
- **Do not add `-short` flags** to skip slow tests
|
- **Do not add `-short` flags** to skip slow tests
|
||||||
- **Do not increase `-timeout`** to hide hanging queries
|
- **Do not increase `-timeout`** to hide hanging queries
|
||||||
- **Do not remove `-count=1` from `script/test`** — Go's test cache
|
- **Do not remove `-count=1` from `script/test`** — Go's test cache
|
||||||
replays a previous run's output without querying anything, so a
|
replays a previous run's output without querying anything, so a
|
||||||
cached pass is not evidence that live resolution works
|
cached pass is not evidence that live resolution works
|
||||||
- **Do not remove `-p 1` from `script/test`** — running the live-DNS
|
|
||||||
packages in parallel oversubscribes live DNS past what their gates
|
|
||||||
bound, which surfaces as unrelated resolver tests failing on
|
|
||||||
expired deadlines
|
|
||||||
- **Do not modify linter configuration** to suppress findings
|
- **Do not modify linter configuration** to suppress findings
|
||||||
|
|||||||
@@ -14,10 +14,7 @@ pre-1.0. No git tags. Core resolver work in flight on feature/resolver
|
|||||||
(dirty: internal/resolver/resolver_test.go). Local checkout has diverged
|
(dirty: internal/resolver/resolver_test.go). Local checkout has diverged
|
||||||
from origin: origin/main is 8 commits ahead (watcher orchestrator,
|
from origin: origin/main is 8 commits ahead (watcher orchestrator,
|
||||||
unified TARGETS) and origin/feature/resolver already contains the full
|
unified TARGETS) and origin/feature/resolver already contains the full
|
||||||
iterative resolver implementation. DNS mocking is banned in this repo
|
iterative resolver implementation with hermetic mocked tests.
|
||||||
(see `TESTING.md`): all tests use live DNS only. The hermetic mocked
|
|
||||||
tests previously noted on `feature/resolver` are gone from its current
|
|
||||||
tip, which carries a live-DNS suite against `*.dns.sneak.cloud`.
|
|
||||||
|
|
||||||
# Next Step
|
# Next Step
|
||||||
|
|
||||||
@@ -26,22 +23,6 @@ Rationale, Design, TODO, License, Author) if any are still missing.
|
|||||||
|
|
||||||
# Completed Steps
|
# Completed Steps
|
||||||
|
|
||||||
- 2026-09-03: restored the transport-failure classification coverage
|
|
||||||
the DNS-mock removal had dropped, and tightened two over-tolerant
|
|
||||||
watcher assertions. `internal/resolver/transport_test.go` covers
|
|
||||||
`StatusTimeout`, `StatusError` and the connection-refused path by
|
|
||||||
binding real nameservers on loopback rather than by mocking
|
|
||||||
`DNSClient` or by aiming a query at a remote address and hoping the
|
|
||||||
network black-holes it; `queryDNS` now honours a port already
|
|
||||||
present in a nameserver address, which is what lets a query be
|
|
||||||
aimed at one. The watcher's `assertStatePopulated` and
|
|
||||||
`TestDomainPortAndTLSChecks` now assert that port and certificate
|
|
||||||
state, and the arguments the port and TLS checkers were called
|
|
||||||
with, match the addresses live DNS returned — previously they
|
|
||||||
asserted only that those sets were non-empty, which would not have
|
|
||||||
caught resolving the wrong addresses. `TESTING.md` records why a
|
|
||||||
loopback nameserver is not a mock.
|
|
||||||
|
|
||||||
- 2026-08-10: comment-only corrections to `script/bootstrap`,
|
- 2026-08-10: comment-only corrections to `script/bootstrap`,
|
||||||
`script/cibuild`, and `Dockerfile.lint`. The `goimports` pin in
|
`script/cibuild`, and `Dockerfile.lint`. The `goimports` pin in
|
||||||
`script/bootstrap` was justified by a claim that `script/fmt-check`
|
`script/bootstrap` was justified by a claim that `script/fmt-check`
|
||||||
@@ -111,10 +92,23 @@ Rationale, Design, TODO, License, Author) if any are still missing.
|
|||||||
`make check` into `script/lint`. `golangci-lint config verify` is
|
`make check` into `script/lint`. `golangci-lint config verify` is
|
||||||
deliberately omitted: it fetches its schema over an unpinned live
|
deliberately omitted: it fetches its schema over an unpinned live
|
||||||
HTTPS call
|
HTTPS call
|
||||||
- 2026-08-07: DNS mocking removed from the entire test suite; watcher
|
- 2026-08-09: in-flight notification deliveries are now drained at
|
||||||
tests now drive the real iterative resolver against live DNS and
|
shutdown (#106): `notify.New` registers an fx `OnStop` hook that waits
|
||||||
`TESTING.md` bans DNS mocks in every package (`remove-dns-mocking`
|
on a `sync.WaitGroup` of tracked delivery goroutines, bounded by the
|
||||||
branch)
|
`OnStop` context; on expiry the outstanding count is logged at warn
|
||||||
|
level and parked retry backoffs are released instead of being dropped
|
||||||
|
silently, and deliveries submitted after the drain begins are refused
|
||||||
|
so shutdown cannot be extended indefinitely; an `OnStop` context that
|
||||||
|
is already expired on entry with nothing outstanding drains quietly
|
||||||
|
rather than warning about deliveries that were never abandoned
|
||||||
|
- 2026-08-09: `http.Server` now sets all four socket-level timeouts
|
||||||
|
(`ReadTimeout` 15s, `ReadHeaderTimeout` 10s, `WriteTimeout` 75s,
|
||||||
|
`IdleTimeout` 120s) as named constants in `internal/server/server.go`,
|
||||||
|
closing the slowloris / unreaped-keep-alive exposure required by
|
||||||
|
`REPO_POLICIES.md` before 1.0; `WriteTimeout` is deliberately greater
|
||||||
|
than the 60s `chimw.Timeout` handler budget so that budget stays
|
||||||
|
reachable, and tests in `internal/server` pin both the non-zero
|
||||||
|
values and that relationship (#99)
|
||||||
- 2026-08-07: golangci-lint bumped to v2.12.2 (commit-pinned installs
|
- 2026-08-07: golangci-lint bumped to v2.12.2 (commit-pinned installs
|
||||||
in `Dockerfile` and `script/bootstrap`); `.golangci.yml` set to the
|
in `Dockerfile` and `script/bootstrap`); `.golangci.yml` set to the
|
||||||
org-standard v2-schema config used across the org's repos
|
org-standard v2-schema config used across the org's repos
|
||||||
@@ -126,8 +120,7 @@ Rationale, Design, TODO, License, Author) if any are still missing.
|
|||||||
- 2026-07-07 Adopted scripts-to-rule-them-all: `script/` entrypoints,
|
- 2026-07-07 Adopted scripts-to-rule-them-all: `script/` entrypoints,
|
||||||
Makefile shims, README Entrypoints section
|
Makefile shims, README Entrypoints section
|
||||||
- 2026-02-20: iterative DNS resolver implemented; tests made hermetic
|
- 2026-02-20: iterative DNS resolver implemented; tests made hermetic
|
||||||
with mocked DNS (origin/feature/resolver, unmerged; superseded — DNS
|
with mocked DNS (origin/feature/resolver, unmerged)
|
||||||
mocking is banned, see `TESTING.md`)
|
|
||||||
- 2026-02-20: CI actions and go install refs pinned to commit SHAs;
|
- 2026-02-20: CI actions and go install refs pinned to commit SHAs;
|
||||||
Gitea Actions workflow for make check (origin/ci/make-check, unmerged)
|
Gitea Actions workflow for make check (origin/ci/make-check, unmerged)
|
||||||
- 2026-02-20: watcher monitoring orchestrator merged to main (#8)
|
- 2026-02-20: watcher monitoring orchestrator merged to main (#8)
|
||||||
@@ -152,9 +145,8 @@ Branch reconciliation:
|
|||||||
- Sync local checkout with origin: local main is 8 commits behind
|
- Sync local checkout with origin: local main is 8 commits behind
|
||||||
origin/main; local feature/resolver has diverged from
|
origin/main; local feature/resolver has diverged from
|
||||||
origin/feature/resolver, which already implements the resolver
|
origin/feature/resolver, which already implements the resolver
|
||||||
- Merge in-flight branches to main once green: feature/resolver
|
- Merge in-flight branches to main once green: feature/resolver,
|
||||||
(confirm its tests remain live-DNS — DNS mocking is banned, see
|
ci/make-check, feature/portcheck-implementation,
|
||||||
`TESTING.md`), ci/make-check, feature/portcheck-implementation,
|
|
||||||
feature/tlscheck-implementation
|
feature/tlscheck-implementation
|
||||||
|
|
||||||
Resolver (plan from untracked TODO.md; largely implemented on
|
Resolver (plan from untracked TODO.md; largely implemented on
|
||||||
@@ -238,7 +230,6 @@ Infrastructure notes (from untracked TODO.md):
|
|||||||
- Module path sneak.berlin/go/dnswatcher differs from the git.eeqj.de
|
- Module path sneak.berlin/go/dnswatcher differs from the git.eeqj.de
|
||||||
remote intentionally; do not "fix" it
|
remote intentionally; do not "fix" it
|
||||||
- Dependencies: github.com/miekg/dns, golang.org/x/net/publicsuffix
|
- Dependencies: github.com/miekg/dns, golang.org/x/net/publicsuffix
|
||||||
- Resolver tests originally used live DNS against `*.dns.sneak.cloud`
|
- Resolver tests originally used live DNS against *.dns.sneak.cloud
|
||||||
(required records documented in the test file header); `main` now
|
(required records documented in the test file header); origin now has
|
||||||
tests against live public DNS. DNS mocking is banned (see
|
mocked hermetic tests, keep them hermetic
|
||||||
`TESTING.md`); never reintroduce hermetic mocked DNS tests
|
|
||||||
|
|||||||
@@ -32,11 +32,27 @@ func NewRequestForTest(
|
|||||||
// NewTestService creates a Service suitable for unit testing.
|
// NewTestService creates a Service suitable for unit testing.
|
||||||
// It discards log output and uses the given transport.
|
// It discards log output and uses the given transport.
|
||||||
func NewTestService(transport http.RoundTripper) *Service {
|
func NewTestService(transport http.RoundTripper) *Service {
|
||||||
return &Service{
|
return newService(slog.New(slog.DiscardHandler), transport)
|
||||||
log: slog.New(slog.DiscardHandler),
|
}
|
||||||
transport: transport,
|
|
||||||
history: NewAlertHistory(),
|
// NewTestServiceWithLogger creates a Service that writes to the
|
||||||
}
|
// given handler, so tests can assert on emitted log records.
|
||||||
|
func NewTestServiceWithLogger(
|
||||||
|
transport http.RoundTripper,
|
||||||
|
handler slog.Handler,
|
||||||
|
) *Service {
|
||||||
|
return newService(slog.New(handler), transport)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Drain exports drain for testing.
|
||||||
|
func (svc *Service) Drain(ctx context.Context) {
|
||||||
|
svc.drain(ctx)
|
||||||
|
}
|
||||||
|
|
||||||
|
// OutstandingDeliveries reports how many delivery goroutines
|
||||||
|
// are currently tracked as in flight.
|
||||||
|
func (svc *Service) OutstandingDeliveries() int64 {
|
||||||
|
return svc.outstanding.Load()
|
||||||
}
|
}
|
||||||
|
|
||||||
// SetNtfyURL sets the ntfy URL on a Service for testing.
|
// SetNtfyURL sets the ntfy URL on a Service for testing.
|
||||||
|
|||||||
+73
-56
@@ -12,6 +12,8 @@ import (
|
|||||||
"log/slog"
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/url"
|
"net/url"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"go.uber.org/fx"
|
"go.uber.org/fx"
|
||||||
@@ -115,19 +117,41 @@ type Service struct {
|
|||||||
history *AlertHistory
|
history *AlertHistory
|
||||||
retryConfig RetryConfig
|
retryConfig RetryConfig
|
||||||
sleepFn func(time.Duration) <-chan time.Time
|
sleepFn func(time.Duration) <-chan time.Time
|
||||||
|
|
||||||
|
// Shutdown draining state. drainMu guards draining and
|
||||||
|
// serialises it against the counter increment in
|
||||||
|
// startDelivery; inFlight tracks the delivery goroutines
|
||||||
|
// themselves and outstanding mirrors its count so a timed
|
||||||
|
// out drain can report how many were abandoned.
|
||||||
|
drainMu sync.Mutex
|
||||||
|
draining bool
|
||||||
|
inFlight sync.WaitGroup
|
||||||
|
outstanding atomic.Int64
|
||||||
|
abandon chan struct{}
|
||||||
|
abandonOnce sync.Once
|
||||||
|
}
|
||||||
|
|
||||||
|
// newService builds a Service with the fields every Service
|
||||||
|
// needs regardless of how it was constructed.
|
||||||
|
func newService(
|
||||||
|
log *slog.Logger,
|
||||||
|
transport http.RoundTripper,
|
||||||
|
) *Service {
|
||||||
|
return &Service{
|
||||||
|
log: log,
|
||||||
|
transport: transport,
|
||||||
|
history: NewAlertHistory(),
|
||||||
|
abandon: make(chan struct{}),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// New creates a new notify Service.
|
// New creates a new notify Service.
|
||||||
func New(
|
func New(
|
||||||
_ fx.Lifecycle,
|
lifecycle fx.Lifecycle,
|
||||||
params Params,
|
params Params,
|
||||||
) (*Service, error) {
|
) (*Service, error) {
|
||||||
svc := &Service{
|
svc := newService(params.Logger.Get(), http.DefaultTransport)
|
||||||
log: params.Logger.Get(),
|
svc.config = params.Config
|
||||||
transport: http.DefaultTransport,
|
|
||||||
config: params.Config,
|
|
||||||
history: NewAlertHistory(),
|
|
||||||
}
|
|
||||||
|
|
||||||
if params.Config.NtfyTopic != "" {
|
if params.Config.NtfyTopic != "" {
|
||||||
u, err := ValidateWebhookURL(
|
u, err := ValidateWebhookURL(
|
||||||
@@ -168,6 +192,14 @@ func New(
|
|||||||
svc.mattermostWebhookURL = u
|
svc.mattermostWebhookURL = u
|
||||||
}
|
}
|
||||||
|
|
||||||
|
lifecycle.Append(fx.Hook{
|
||||||
|
OnStop: func(ctx context.Context) error {
|
||||||
|
svc.drain(ctx)
|
||||||
|
|
||||||
|
return nil
|
||||||
|
},
|
||||||
|
})
|
||||||
|
|
||||||
return svc, nil
|
return svc, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -194,6 +226,32 @@ func (svc *Service) SendNotification(
|
|||||||
svc.dispatchMattermost(ctx, title, message, priority)
|
svc.dispatchMattermost(ctx, title, message, priority)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// dispatch delivers a notification to one endpoint on a
|
||||||
|
// tracked background goroutine.
|
||||||
|
//
|
||||||
|
// The delivery context is detached from ctx with
|
||||||
|
// context.WithoutCancel so that a cancelled caller does not
|
||||||
|
// kill a delivery already under way; the shutdown drain, not
|
||||||
|
// the caller, decides how long deliveries may keep running.
|
||||||
|
func (svc *Service) dispatch(
|
||||||
|
ctx context.Context,
|
||||||
|
endpoint string,
|
||||||
|
send func(context.Context) error,
|
||||||
|
) {
|
||||||
|
notifyCtx := context.WithoutCancel(ctx)
|
||||||
|
|
||||||
|
svc.startDelivery(endpoint, func() {
|
||||||
|
err := svc.deliverWithRetry(notifyCtx, endpoint, send)
|
||||||
|
if err != nil {
|
||||||
|
svc.log.Error(
|
||||||
|
"failed to send notification after retries",
|
||||||
|
"endpoint", endpoint,
|
||||||
|
"error", err,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
func (svc *Service) dispatchNtfy(
|
func (svc *Service) dispatchNtfy(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
title, message, priority string,
|
title, message, priority string,
|
||||||
@@ -202,26 +260,11 @@ func (svc *Service) dispatchNtfy(
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
go func() {
|
svc.dispatch(ctx, "ntfy", func(c context.Context) error {
|
||||||
notifyCtx := context.WithoutCancel(ctx)
|
|
||||||
|
|
||||||
err := svc.deliverWithRetry(
|
|
||||||
notifyCtx, "ntfy",
|
|
||||||
func(c context.Context) error {
|
|
||||||
return svc.sendNtfy(
|
return svc.sendNtfy(
|
||||||
c, svc.ntfyURL,
|
c, svc.ntfyURL, title, message, priority,
|
||||||
title, message, priority,
|
|
||||||
)
|
)
|
||||||
},
|
})
|
||||||
)
|
|
||||||
if err != nil {
|
|
||||||
svc.log.Error(
|
|
||||||
"failed to send ntfy notification "+
|
|
||||||
"after retries",
|
|
||||||
"error", err,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (svc *Service) dispatchSlack(
|
func (svc *Service) dispatchSlack(
|
||||||
@@ -232,26 +275,11 @@ func (svc *Service) dispatchSlack(
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
go func() {
|
svc.dispatch(ctx, "slack", func(c context.Context) error {
|
||||||
notifyCtx := context.WithoutCancel(ctx)
|
|
||||||
|
|
||||||
err := svc.deliverWithRetry(
|
|
||||||
notifyCtx, "slack",
|
|
||||||
func(c context.Context) error {
|
|
||||||
return svc.sendSlack(
|
return svc.sendSlack(
|
||||||
c, svc.slackWebhookURL,
|
c, svc.slackWebhookURL, title, message, priority,
|
||||||
title, message, priority,
|
|
||||||
)
|
)
|
||||||
},
|
})
|
||||||
)
|
|
||||||
if err != nil {
|
|
||||||
svc.log.Error(
|
|
||||||
"failed to send slack notification "+
|
|
||||||
"after retries",
|
|
||||||
"error", err,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (svc *Service) dispatchMattermost(
|
func (svc *Service) dispatchMattermost(
|
||||||
@@ -262,11 +290,8 @@ func (svc *Service) dispatchMattermost(
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
go func() {
|
svc.dispatch(
|
||||||
notifyCtx := context.WithoutCancel(ctx)
|
ctx, "mattermost",
|
||||||
|
|
||||||
err := svc.deliverWithRetry(
|
|
||||||
notifyCtx, "mattermost",
|
|
||||||
func(c context.Context) error {
|
func(c context.Context) error {
|
||||||
return svc.sendSlack(
|
return svc.sendSlack(
|
||||||
c, svc.mattermostWebhookURL,
|
c, svc.mattermostWebhookURL,
|
||||||
@@ -274,14 +299,6 @@ func (svc *Service) dispatchMattermost(
|
|||||||
)
|
)
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
if err != nil {
|
|
||||||
svc.log.Error(
|
|
||||||
"failed to send mattermost notification "+
|
|
||||||
"after retries",
|
|
||||||
"error", err,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (svc *Service) sendNtfy(
|
func (svc *Service) sendNtfy(
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package notify
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"fmt"
|
||||||
"math"
|
"math"
|
||||||
"math/rand/v2"
|
"math/rand/v2"
|
||||||
"time"
|
"time"
|
||||||
@@ -121,6 +122,14 @@ func (svc *Service) deliverWithRetry(
|
|||||||
select {
|
select {
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
return ctx.Err()
|
return ctx.Err()
|
||||||
|
case <-svc.abandon:
|
||||||
|
// Shutdown drained past its deadline; stop
|
||||||
|
// sleeping rather than outlive the process.
|
||||||
|
// A nil channel (Service built without a
|
||||||
|
// constructor) simply never fires.
|
||||||
|
return fmt.Errorf(
|
||||||
|
"%w: %s", ErrDeliveryAbandoned, endpoint,
|
||||||
|
)
|
||||||
case <-svc.sleepFunc(delay):
|
case <-svc.sleepFunc(delay):
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,119 @@
|
|||||||
|
package notify
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
)
|
||||||
|
|
||||||
|
// ErrDeliveryAbandoned is returned by a retry loop that was
|
||||||
|
// cut short because shutdown drained past its deadline.
|
||||||
|
var ErrDeliveryAbandoned = errors.New(
|
||||||
|
"notification delivery abandoned at shutdown",
|
||||||
|
)
|
||||||
|
|
||||||
|
// startDelivery runs fn on its own goroutine while tracking it,
|
||||||
|
// so that drain can wait for it during shutdown.
|
||||||
|
//
|
||||||
|
// The WaitGroup counter is incremented here, on the caller's
|
||||||
|
// goroutine, before the worker exists: incrementing it inside
|
||||||
|
// the worker would race with drain's Wait and could let
|
||||||
|
// shutdown sail past a delivery that had not started yet.
|
||||||
|
//
|
||||||
|
// Once draining has begun the delivery is refused outright
|
||||||
|
// rather than queued, so a steady stream of newly submitted
|
||||||
|
// notifications cannot keep extending the drain.
|
||||||
|
func (svc *Service) startDelivery(endpoint string, fn func()) {
|
||||||
|
svc.drainMu.Lock()
|
||||||
|
|
||||||
|
if svc.draining {
|
||||||
|
svc.drainMu.Unlock()
|
||||||
|
|
||||||
|
svc.log.Warn(
|
||||||
|
"notification not dispatched: shutdown in progress",
|
||||||
|
"endpoint", endpoint,
|
||||||
|
)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
svc.outstanding.Add(1)
|
||||||
|
|
||||||
|
// WaitGroup.Go increments the counter synchronously, here,
|
||||||
|
// and only then starts the goroutine.
|
||||||
|
svc.inFlight.Go(func() {
|
||||||
|
// Runs before the WaitGroup counter is decremented, so
|
||||||
|
// a drain that times out reports an accurate count.
|
||||||
|
defer svc.outstanding.Add(-1)
|
||||||
|
|
||||||
|
fn()
|
||||||
|
})
|
||||||
|
|
||||||
|
svc.drainMu.Unlock()
|
||||||
|
}
|
||||||
|
|
||||||
|
// drain waits for in-flight notification deliveries to finish.
|
||||||
|
//
|
||||||
|
// It first stops accepting new deliveries, then waits until
|
||||||
|
// either every outstanding delivery has completed or ctx
|
||||||
|
// expires — whichever comes first. ctx is the context fx
|
||||||
|
// passes to the OnStop hook, so a permanently dead webhook
|
||||||
|
// cannot hang shutdown indefinitely.
|
||||||
|
//
|
||||||
|
// When the deadline arrives with deliveries still outstanding,
|
||||||
|
// the count is logged at warn level and the abandon channel is
|
||||||
|
// closed, which releases any retry loop sleeping in backoff.
|
||||||
|
// Deliveries already inside an HTTP round trip are bounded by
|
||||||
|
// the existing httpClientTimeout instead.
|
||||||
|
//
|
||||||
|
// A ctx that is already expired on entry is not by itself cause
|
||||||
|
// for alarm: if nothing is outstanding there is nothing to
|
||||||
|
// abandon, and the drain says so at debug level rather than
|
||||||
|
// warning about deliveries that do not exist.
|
||||||
|
func (svc *Service) drain(ctx context.Context) {
|
||||||
|
svc.drainMu.Lock()
|
||||||
|
svc.draining = true
|
||||||
|
svc.drainMu.Unlock()
|
||||||
|
|
||||||
|
done := make(chan struct{})
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
svc.inFlight.Wait()
|
||||||
|
close(done)
|
||||||
|
}()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-done:
|
||||||
|
svc.log.Debug(
|
||||||
|
"all in-flight notifications completed",
|
||||||
|
)
|
||||||
|
case <-ctx.Done():
|
||||||
|
// outstanding is decremented before the WaitGroup
|
||||||
|
// counter, and startDelivery can no longer add to it
|
||||||
|
// now that draining is set, so a zero here means every
|
||||||
|
// delivery really did finish. ctx expiring in that
|
||||||
|
// state (an OnStop context that was already cancelled
|
||||||
|
// on entry is the usual way) abandons nothing, so it
|
||||||
|
// must not close abandon or warn about it.
|
||||||
|
abandoned := svc.outstanding.Load()
|
||||||
|
if abandoned == 0 {
|
||||||
|
svc.log.Debug(
|
||||||
|
"all in-flight notifications completed",
|
||||||
|
)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
svc.abandonOnce.Do(func() {
|
||||||
|
if svc.abandon != nil {
|
||||||
|
close(svc.abandon)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
svc.log.Warn(
|
||||||
|
"shutdown deadline reached with notifications "+
|
||||||
|
"still in flight; abandoning them",
|
||||||
|
"abandoned", abandoned,
|
||||||
|
"error", ctx.Err(),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,531 @@
|
|||||||
|
package notify_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"log/slog"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"net/url"
|
||||||
|
"strings"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"go.uber.org/fx"
|
||||||
|
|
||||||
|
"sneak.berlin/go/dnswatcher/internal/config"
|
||||||
|
"sneak.berlin/go/dnswatcher/internal/globals"
|
||||||
|
"sneak.berlin/go/dnswatcher/internal/logger"
|
||||||
|
"sneak.berlin/go/dnswatcher/internal/notify"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Timings used by the drain tests. They stay in the same
|
||||||
|
// 10-100ms band as the retry tests so the suite never waits on
|
||||||
|
// a real backoff delay.
|
||||||
|
const (
|
||||||
|
// inFlightHold is how long a delivery is kept mid-request
|
||||||
|
// before the handler is released.
|
||||||
|
inFlightHold = 30 * time.Millisecond
|
||||||
|
|
||||||
|
// drainDeadline bounds a drain that is expected to time
|
||||||
|
// out.
|
||||||
|
drainDeadline = 50 * time.Millisecond
|
||||||
|
|
||||||
|
// drainSlack is the upper bound on how long a bounded
|
||||||
|
// drain may take; generous enough for a loaded CI box,
|
||||||
|
// still far below the 20s test ceiling.
|
||||||
|
drainSlack = 2 * time.Second
|
||||||
|
|
||||||
|
// settleDelay is how long to wait before asserting that
|
||||||
|
// something did *not* happen.
|
||||||
|
settleDelay = 50 * time.Millisecond
|
||||||
|
|
||||||
|
// idleDrainBound is the upper bound on a drain that has
|
||||||
|
// nothing in flight. It is deliberately far above the cost
|
||||||
|
// of the goroutine hop through inFlight.Wait() — which
|
||||||
|
// reached 57ms on a loaded box under -race with the package's
|
||||||
|
// parallel tests — and far below drainSlack, the deadline
|
||||||
|
// such a drain is given. A drain that blocked until its
|
||||||
|
// deadline instead of returning on the WaitGroup therefore
|
||||||
|
// still fails this bound, but scheduling delay alone cannot.
|
||||||
|
idleDrainBound = 500 * time.Millisecond
|
||||||
|
)
|
||||||
|
|
||||||
|
// syncBuffer is an io.Writer safe for concurrent use, so log
|
||||||
|
// output written from delivery goroutines can be inspected.
|
||||||
|
type syncBuffer struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
buf bytes.Buffer
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sb *syncBuffer) Write(p []byte) (int, error) {
|
||||||
|
sb.mu.Lock()
|
||||||
|
defer sb.mu.Unlock()
|
||||||
|
|
||||||
|
return sb.buf.Write(p) //nolint:wrapcheck // test helper
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sb *syncBuffer) String() string {
|
||||||
|
sb.mu.Lock()
|
||||||
|
defer sb.mu.Unlock()
|
||||||
|
|
||||||
|
return sb.buf.String()
|
||||||
|
}
|
||||||
|
|
||||||
|
// newLoggingService returns a Service writing JSON logs into
|
||||||
|
// the returned buffer.
|
||||||
|
func newLoggingService(
|
||||||
|
transport http.RoundTripper,
|
||||||
|
) (*notify.Service, *syncBuffer) {
|
||||||
|
logs := &syncBuffer{}
|
||||||
|
handler := slog.NewJSONHandler(logs, nil)
|
||||||
|
|
||||||
|
return notify.NewTestServiceWithLogger(transport, handler),
|
||||||
|
logs
|
||||||
|
}
|
||||||
|
|
||||||
|
// blockingNtfyServer returns a server whose handler signals on
|
||||||
|
// entered, waits for release, and then responds 200.
|
||||||
|
func blockingNtfyServer(
|
||||||
|
entered chan<- struct{},
|
||||||
|
release <-chan struct{},
|
||||||
|
served *atomic.Bool,
|
||||||
|
) *httptest.Server {
|
||||||
|
var once sync.Once
|
||||||
|
|
||||||
|
return httptest.NewServer(
|
||||||
|
http.HandlerFunc(
|
||||||
|
func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
once.Do(func() { close(entered) })
|
||||||
|
<-release
|
||||||
|
|
||||||
|
served.Store(true)
|
||||||
|
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestDrainWaitsForInFlightDelivery verifies that a delivery
|
||||||
|
// already under way when shutdown starts is allowed to finish.
|
||||||
|
func TestDrainWaitsForInFlightDelivery(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
var served atomic.Bool
|
||||||
|
|
||||||
|
entered := make(chan struct{})
|
||||||
|
release := make(chan struct{})
|
||||||
|
|
||||||
|
srv := blockingNtfyServer(entered, release, &served)
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
topicURL, _ := url.Parse(srv.URL)
|
||||||
|
|
||||||
|
svc := notify.NewTestService(http.DefaultTransport)
|
||||||
|
svc.SetNtfyURL(topicURL)
|
||||||
|
|
||||||
|
svc.SendNotification(
|
||||||
|
context.Background(), "t", "m", prioInfo,
|
||||||
|
)
|
||||||
|
|
||||||
|
// Make sure the delivery really is mid-request before the
|
||||||
|
// drain begins.
|
||||||
|
select {
|
||||||
|
case <-entered:
|
||||||
|
case <-time.After(drainSlack):
|
||||||
|
t.Fatal("delivery never reached the endpoint")
|
||||||
|
}
|
||||||
|
|
||||||
|
// As in TestDrainBoundedByContextDeadline: start is captured
|
||||||
|
// before the clock it is compared against, here the timer
|
||||||
|
// holding the delivery open, so elapsed covers the whole hold
|
||||||
|
// and the lower bound cannot come out short from scheduling
|
||||||
|
// delay alone.
|
||||||
|
start := time.Now()
|
||||||
|
|
||||||
|
timer := time.AfterFunc(inFlightHold, func() {
|
||||||
|
close(release)
|
||||||
|
})
|
||||||
|
defer timer.Stop()
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(
|
||||||
|
context.Background(), drainSlack,
|
||||||
|
)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
svc.Drain(ctx)
|
||||||
|
|
||||||
|
elapsed := time.Since(start)
|
||||||
|
|
||||||
|
if !served.Load() {
|
||||||
|
t.Error(
|
||||||
|
"drain returned before the in-flight delivery " +
|
||||||
|
"completed",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
if elapsed < inFlightHold {
|
||||||
|
t.Errorf(
|
||||||
|
"drain took %v, want at least %v",
|
||||||
|
elapsed, inFlightHold,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
if got := svc.OutstandingDeliveries(); got != 0 {
|
||||||
|
t.Errorf("outstanding deliveries = %d, want 0", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// neverFires returns a channel that never delivers, standing in
|
||||||
|
// for a long backoff sleep without actually sleeping.
|
||||||
|
func neverFires(_ time.Duration) <-chan time.Time {
|
||||||
|
return make(chan time.Time)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestDrainBoundedByContextDeadline verifies that a delivery
|
||||||
|
// stuck retrying against a dead endpoint does not hold shutdown
|
||||||
|
// past the OnStop context deadline, and that the abandoned
|
||||||
|
// deliveries are logged at warn level rather than dropped
|
||||||
|
// silently.
|
||||||
|
func TestDrainBoundedByContextDeadline(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
var requests atomic.Int64
|
||||||
|
|
||||||
|
srv := httptest.NewServer(
|
||||||
|
http.HandlerFunc(
|
||||||
|
func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
requests.Add(1)
|
||||||
|
|
||||||
|
w.WriteHeader(http.StatusInternalServerError)
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
topicURL, _ := url.Parse(srv.URL)
|
||||||
|
|
||||||
|
svc, logs := newLoggingService(http.DefaultTransport)
|
||||||
|
svc.SetNtfyURL(topicURL)
|
||||||
|
// Never let the backoff sleep complete: the delivery is
|
||||||
|
// parked in its retry wait until shutdown releases it.
|
||||||
|
svc.SetSleepFunc(neverFires)
|
||||||
|
svc.SetRetryConfig(notify.RetryConfig{
|
||||||
|
MaxRetries: 5,
|
||||||
|
BaseDelay: time.Hour,
|
||||||
|
MaxDelay: time.Hour,
|
||||||
|
})
|
||||||
|
|
||||||
|
svc.SendNotification(
|
||||||
|
context.Background(), "t", "m", prioError,
|
||||||
|
)
|
||||||
|
|
||||||
|
waitForCondition(t, func() bool {
|
||||||
|
return requests.Load() >= 1 &&
|
||||||
|
svc.OutstandingDeliveries() == 1
|
||||||
|
})
|
||||||
|
|
||||||
|
// start must be captured *before* the deadline clock starts,
|
||||||
|
// so that the measured interval is a superset of the deadline
|
||||||
|
// interval. Capturing it after context.WithTimeout would
|
||||||
|
// make elapsed structurally smaller than drainDeadline and
|
||||||
|
// the lower bound below unfalsifiable-by-luck: it would fail
|
||||||
|
// whenever the two statements were separated by any
|
||||||
|
// scheduling delay, and pass otherwise, regardless of what
|
||||||
|
// the drain did.
|
||||||
|
start := time.Now()
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(
|
||||||
|
context.Background(), drainDeadline,
|
||||||
|
)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
// The upper bound is enforced by a watchdog rather than by
|
||||||
|
// measuring after the fact: a drain that is not bounded at
|
||||||
|
// all never returns here (the delivery is parked in a backoff
|
||||||
|
// that never fires), so an unbounded drain must fail this
|
||||||
|
// test promptly instead of hanging the package until the test
|
||||||
|
// binary's 30s timeout.
|
||||||
|
returned := make(chan struct{})
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
defer close(returned)
|
||||||
|
|
||||||
|
svc.Drain(ctx)
|
||||||
|
}()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-returned:
|
||||||
|
case <-time.After(drainSlack):
|
||||||
|
t.Fatalf(
|
||||||
|
"drain did not return within %v; its %v deadline "+
|
||||||
|
"did not bound it",
|
||||||
|
drainSlack, drainDeadline,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// The lower bound is the real assertion: the drain must have
|
||||||
|
// waited for its whole deadline rather than giving up on the
|
||||||
|
// outstanding delivery early. With start captured above, an
|
||||||
|
// early return is the only thing that can make it fail.
|
||||||
|
if elapsed := time.Since(start); elapsed < drainDeadline {
|
||||||
|
t.Errorf(
|
||||||
|
"drain returned after %v, before its %v deadline",
|
||||||
|
elapsed, drainDeadline,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
assertAbandonLogged(t, logs.String())
|
||||||
|
|
||||||
|
// The abandoned delivery must stop retrying rather than
|
||||||
|
// outlive the drain.
|
||||||
|
waitForCondition(t, func() bool {
|
||||||
|
return svc.OutstandingDeliveries() == 0
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
// assertAbandonLogged checks that the drain logged the
|
||||||
|
// abandoned deliveries at warn level with a count.
|
||||||
|
func assertAbandonLogged(t *testing.T, output string) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
if !strings.Contains(output, `"level":"WARN"`) {
|
||||||
|
t.Errorf(
|
||||||
|
"abandoned deliveries not logged at warn level; "+
|
||||||
|
"log output: %s",
|
||||||
|
output,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
if !strings.Contains(output, `"abandoned":1`) {
|
||||||
|
t.Errorf(
|
||||||
|
"abandoned delivery count not logged; "+
|
||||||
|
"log output: %s",
|
||||||
|
output,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestDrainRefusesNewDeliveries verifies that notifications
|
||||||
|
// submitted after the drain has begun are refused and logged,
|
||||||
|
// so a stream of new work cannot extend shutdown indefinitely.
|
||||||
|
func TestDrainRefusesNewDeliveries(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
var requests atomic.Int64
|
||||||
|
|
||||||
|
srv := httptest.NewServer(
|
||||||
|
http.HandlerFunc(
|
||||||
|
func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
requests.Add(1)
|
||||||
|
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
target, _ := url.Parse(srv.URL)
|
||||||
|
|
||||||
|
svc, logs := newLoggingService(http.DefaultTransport)
|
||||||
|
svc.SetNtfyURL(target)
|
||||||
|
svc.SetSlackWebhookURL(target)
|
||||||
|
svc.SetMattermostWebhookURL(target)
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(
|
||||||
|
context.Background(), drainSlack,
|
||||||
|
)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
// Nothing is in flight, so this returns immediately and
|
||||||
|
// leaves the service refusing further deliveries.
|
||||||
|
svc.Drain(ctx)
|
||||||
|
|
||||||
|
for range 3 {
|
||||||
|
svc.SendNotification(
|
||||||
|
context.Background(), "t", "m", prioInfo,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
time.Sleep(settleDelay)
|
||||||
|
|
||||||
|
if got := requests.Load(); got != 0 {
|
||||||
|
t.Errorf(
|
||||||
|
"%d requests reached the endpoint after drain, "+
|
||||||
|
"want 0",
|
||||||
|
got,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
if got := svc.OutstandingDeliveries(); got != 0 {
|
||||||
|
t.Errorf("outstanding deliveries = %d, want 0", got)
|
||||||
|
}
|
||||||
|
|
||||||
|
output := logs.String()
|
||||||
|
if !strings.Contains(output, "shutdown in progress") {
|
||||||
|
t.Errorf(
|
||||||
|
"refused deliveries not logged; log output: %s",
|
||||||
|
output,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// recordingLifecycle is a minimal fx.Lifecycle that records the
|
||||||
|
// hooks appended to it, so the wiring done by notify.New can be
|
||||||
|
// inspected without standing up a whole fx application.
|
||||||
|
type recordingLifecycle struct {
|
||||||
|
hooks []fx.Hook
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *recordingLifecycle) Append(hook fx.Hook) {
|
||||||
|
l.hooks = append(l.hooks, hook)
|
||||||
|
}
|
||||||
|
|
||||||
|
// newNotifyService builds a Service through the real
|
||||||
|
// constructor, wired to the given lifecycle.
|
||||||
|
func newNotifyService(
|
||||||
|
t *testing.T,
|
||||||
|
lifecycle fx.Lifecycle,
|
||||||
|
ntfyTopic string,
|
||||||
|
) *notify.Service {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
g, err := globals.New(nil)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("globals.New: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
log, err := logger.New(nil, logger.Params{Globals: g})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("logger.New: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
svc, err := notify.New(lifecycle, notify.Params{
|
||||||
|
Logger: log,
|
||||||
|
Config: &config.Config{NtfyTopic: ntfyTopic},
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("notify.New: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return svc
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestNewRegistersDrainingStopHook verifies that notify.New
|
||||||
|
// wires an OnStop hook into the fx lifecycle and that the hook
|
||||||
|
// waits for in-flight deliveries.
|
||||||
|
func TestNewRegistersDrainingStopHook(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
var served atomic.Bool
|
||||||
|
|
||||||
|
entered := make(chan struct{})
|
||||||
|
release := make(chan struct{})
|
||||||
|
|
||||||
|
srv := blockingNtfyServer(entered, release, &served)
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
lifecycle := &recordingLifecycle{}
|
||||||
|
svc := newNotifyService(t, lifecycle, srv.URL)
|
||||||
|
|
||||||
|
if len(lifecycle.hooks) != 1 {
|
||||||
|
t.Fatalf(
|
||||||
|
"appended %d lifecycle hooks, want 1",
|
||||||
|
len(lifecycle.hooks),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
stop := lifecycle.hooks[0].OnStop
|
||||||
|
if stop == nil {
|
||||||
|
t.Fatal("lifecycle hook has no OnStop function")
|
||||||
|
}
|
||||||
|
|
||||||
|
svc.SendNotification(
|
||||||
|
context.Background(), "t", "m", prioInfo,
|
||||||
|
)
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-entered:
|
||||||
|
case <-time.After(drainSlack):
|
||||||
|
t.Fatal("delivery never reached the endpoint")
|
||||||
|
}
|
||||||
|
|
||||||
|
timer := time.AfterFunc(inFlightHold, func() {
|
||||||
|
close(release)
|
||||||
|
})
|
||||||
|
defer timer.Stop()
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(
|
||||||
|
context.Background(), drainSlack,
|
||||||
|
)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
err := stop(ctx)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("OnStop returned error: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if !served.Load() {
|
||||||
|
t.Error(
|
||||||
|
"OnStop returned before the in-flight delivery " +
|
||||||
|
"completed",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestDrainWithoutDeliveriesReturnsImmediately verifies the
|
||||||
|
// common case: nothing in flight, shutdown is not delayed.
|
||||||
|
func TestDrainWithoutDeliveriesReturnsImmediately(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
svc := notify.NewTestService(http.DefaultTransport)
|
||||||
|
|
||||||
|
// Captured before the deadline clock, as elsewhere in this
|
||||||
|
// file; for an upper bound that is the conservative
|
||||||
|
// direction, since the measured interval can then only be
|
||||||
|
// longer than the drain itself.
|
||||||
|
start := time.Now()
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(
|
||||||
|
context.Background(), drainSlack,
|
||||||
|
)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
svc.Drain(ctx)
|
||||||
|
|
||||||
|
if elapsed := time.Since(start); elapsed > idleDrainBound {
|
||||||
|
t.Errorf(
|
||||||
|
"drain of an idle service took %v, want well "+
|
||||||
|
"under its %v deadline",
|
||||||
|
elapsed, drainSlack,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestDrainWithCancelledContextDoesNotWarn verifies that an
|
||||||
|
// OnStop context that is already dead on entry does not produce
|
||||||
|
// an "abandoning them" warning when there was nothing in flight
|
||||||
|
// to abandon. The expired context wins the select immediately,
|
||||||
|
// so only the outstanding count can tell the difference between
|
||||||
|
// a genuine timeout and a shutdown that had simply already run
|
||||||
|
// out of time with no work left.
|
||||||
|
func TestDrainWithCancelledContextDoesNotWarn(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
svc, logs := newLoggingService(http.DefaultTransport)
|
||||||
|
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
cancel()
|
||||||
|
|
||||||
|
svc.Drain(ctx)
|
||||||
|
|
||||||
|
if output := logs.String(); strings.Contains(
|
||||||
|
output, `"level":"WARN"`,
|
||||||
|
) {
|
||||||
|
t.Errorf(
|
||||||
|
"drain with nothing in flight warned about "+
|
||||||
|
"abandoned deliveries; log output: %s",
|
||||||
|
output,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -7,8 +7,8 @@ import (
|
|||||||
"github.com/miekg/dns"
|
"github.com/miekg/dns"
|
||||||
)
|
)
|
||||||
|
|
||||||
// DNSClient abstracts DNS wire-protocol exchanges over a single
|
// DNSClient abstracts DNS wire-protocol exchanges so the resolver
|
||||||
// transport, letting the resolver switch between UDP and TCP.
|
// can be tested without hitting real nameservers.
|
||||||
type DNSClient interface {
|
type DNSClient interface {
|
||||||
ExchangeContext(
|
ExchangeContext(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
|
|||||||
@@ -20,10 +20,6 @@ const (
|
|||||||
minDomainLabels = 2
|
minDomainLabels = 2
|
||||||
)
|
)
|
||||||
|
|
||||||
// defaultDNSPort is the port a nameserver is assumed to listen on
|
|
||||||
// when its address does not carry one.
|
|
||||||
const defaultDNSPort = "53"
|
|
||||||
|
|
||||||
// ErrRefused is returned when a DNS server refuses a query.
|
// ErrRefused is returned when a DNS server refuses a query.
|
||||||
var ErrRefused = errors.New("dns query refused")
|
var ErrRefused = errors.New("dns query refused")
|
||||||
|
|
||||||
@@ -110,19 +106,6 @@ func (r *Resolver) retryTCP(
|
|||||||
return resp
|
return resp
|
||||||
}
|
}
|
||||||
|
|
||||||
// nameserverAddr renders a nameserver address for dialling. A bare
|
|
||||||
// address — the normal case, and what a delegation's glue records
|
|
||||||
// carry — is given the default DNS port. An address that already
|
|
||||||
// specifies a port is dialled as written, which is what makes a
|
|
||||||
// nameserver listening somewhere other than 53 reachable.
|
|
||||||
func nameserverAddr(nsIP string) string {
|
|
||||||
if _, _, err := net.SplitHostPort(nsIP); err == nil {
|
|
||||||
return nsIP
|
|
||||||
}
|
|
||||||
|
|
||||||
return net.JoinHostPort(nsIP, defaultDNSPort)
|
|
||||||
}
|
|
||||||
|
|
||||||
// queryDNS sends a DNS query to a specific server IP.
|
// queryDNS sends a DNS query to a specific server IP.
|
||||||
// Tries non-recursive first, falls back to recursive on
|
// Tries non-recursive first, falls back to recursive on
|
||||||
// REFUSED (handles DNS interception environments).
|
// REFUSED (handles DNS interception environments).
|
||||||
@@ -137,7 +120,7 @@ func (r *Resolver) queryDNS(
|
|||||||
}
|
}
|
||||||
|
|
||||||
name = dns.Fqdn(name)
|
name = dns.Fqdn(name)
|
||||||
addr := nameserverAddr(serverIP)
|
addr := net.JoinHostPort(serverIP, "53")
|
||||||
|
|
||||||
msg := new(dns.Msg)
|
msg := new(dns.Msg)
|
||||||
msg.SetQuestion(name, qtype)
|
msg.SetQuestion(name, qtype)
|
||||||
|
|||||||
@@ -67,4 +67,17 @@ func NewFromLogger(log *slog.Logger) *Resolver {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// NewFromLoggerWithClient creates a Resolver with a custom DNS
|
||||||
|
// client, useful for testing with mock DNS responses.
|
||||||
|
func NewFromLoggerWithClient(
|
||||||
|
log *slog.Logger,
|
||||||
|
client DNSClient,
|
||||||
|
) *Resolver {
|
||||||
|
return &Resolver{
|
||||||
|
log: log,
|
||||||
|
client: client,
|
||||||
|
tcp: client,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Method implementations are in iterative.go.
|
// Method implementations are in iterative.go.
|
||||||
|
|||||||
@@ -8,7 +8,9 @@ import (
|
|||||||
"sort"
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/miekg/dns"
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
|
|
||||||
@@ -517,9 +519,58 @@ func TestQueryAllNameservers_ContextCanceled(t *testing.T) {
|
|||||||
assert.Error(t, err)
|
assert.Error(t, err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Transport-failure classification (StatusTimeout / StatusError)
|
// ----------------------------------------------------------------
|
||||||
// is covered in transport_test.go, against real nameservers bound on
|
// Timeout tests
|
||||||
// loopback.
|
// ----------------------------------------------------------------
|
||||||
|
|
||||||
|
func TestQueryNameserverIP_Timeout(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
log := slog.New(slog.NewTextHandler(
|
||||||
|
os.Stderr,
|
||||||
|
&slog.HandlerOptions{Level: slog.LevelDebug},
|
||||||
|
))
|
||||||
|
|
||||||
|
r := resolver.NewFromLoggerWithClient(
|
||||||
|
log, &timeoutClient{},
|
||||||
|
)
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(
|
||||||
|
context.Background(), 10*time.Second,
|
||||||
|
)
|
||||||
|
t.Cleanup(cancel)
|
||||||
|
|
||||||
|
// Query any IP — the client always returns a timeout error.
|
||||||
|
resp, err := r.QueryNameserverIP(
|
||||||
|
ctx, "unreachable.test.", "192.0.2.1",
|
||||||
|
"example.com",
|
||||||
|
)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
assert.Equal(t, resolver.StatusTimeout, resp.Status)
|
||||||
|
assert.NotEmpty(t, resp.Error)
|
||||||
|
}
|
||||||
|
|
||||||
|
// timeoutClient simulates DNS timeout errors for testing.
|
||||||
|
type timeoutClient struct{}
|
||||||
|
|
||||||
|
func (c *timeoutClient) ExchangeContext(
|
||||||
|
_ context.Context,
|
||||||
|
_ *dns.Msg,
|
||||||
|
_ string,
|
||||||
|
) (*dns.Msg, time.Duration, error) {
|
||||||
|
return nil, 0, &net.OpError{
|
||||||
|
Op: "read",
|
||||||
|
Net: "udp",
|
||||||
|
Err: &timeoutError{},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type timeoutError struct{}
|
||||||
|
|
||||||
|
func (e *timeoutError) Error() string { return "i/o timeout" }
|
||||||
|
func (e *timeoutError) Timeout() bool { return true }
|
||||||
|
func (e *timeoutError) Temporary() bool { return true }
|
||||||
|
|
||||||
func TestResolveIPAddresses_ContextCanceled(t *testing.T) {
|
func TestResolveIPAddresses_ContextCanceled(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|||||||
@@ -1,253 +0,0 @@
|
|||||||
package resolver_test
|
|
||||||
|
|
||||||
// Transport-failure classification tests.
|
|
||||||
//
|
|
||||||
// These are live tests, not mocks. Nothing here substitutes the
|
|
||||||
// resolver's DNSClient: the resolver dials a real UDP socket, writes
|
|
||||||
// a real DNS query with the real miekg/dns client, and applies its
|
|
||||||
// real deadline and its real classification logic to what comes
|
|
||||||
// back. The only thing under test control is which address the query
|
|
||||||
// is sent to, and what — if anything — is listening there.
|
|
||||||
//
|
|
||||||
// That distinction is what the no-DNS-mocks rule in TESTING.md is
|
|
||||||
// about. A fake DNSClient lets the code under test skip DNS entirely
|
|
||||||
// and hands it a manufactured verdict; a nameserver bound on
|
|
||||||
// loopback makes it speak DNS for real and earn one. Pointing a live
|
|
||||||
// query at a nameserver of the test's choosing is no more a mock
|
|
||||||
// than pointing it at a.root-servers.net.
|
|
||||||
//
|
|
||||||
// The public network cannot produce these outcomes on demand. A
|
|
||||||
// black-holed address is not black-holed everywhere — build
|
|
||||||
// environments that intercept UDP/53 answer it locally — so a test
|
|
||||||
// built on one asserts on the network it happens to run on rather
|
|
||||||
// than on the resolver. A loopback nameserver is deterministic
|
|
||||||
// everywhere, and it is fast, because the test picks the deadline.
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"net"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/miekg/dns"
|
|
||||||
"github.com/stretchr/testify/assert"
|
|
||||||
"github.com/stretchr/testify/require"
|
|
||||||
|
|
||||||
"sneak.berlin/go/dnswatcher/internal/resolver"
|
|
||||||
)
|
|
||||||
|
|
||||||
const (
|
|
||||||
// transportBudget is the wall time the timeout test must stay
|
|
||||||
// under. The resolver asks a nameserver for eight record types
|
|
||||||
// in turn and retries each one once, so a nameserver silent on
|
|
||||||
// every type would cost sixteen query timeouts. The test's
|
|
||||||
// nameserver is silent on exactly one type, which costs two,
|
|
||||||
// and this budget fails loudly if that ever stops being true.
|
|
||||||
transportBudget = 8 * time.Second
|
|
||||||
|
|
||||||
// transportDeadline is the caller deadline the tests run
|
|
||||||
// under. It is generous on purpose: these tests are about the
|
|
||||||
// resolver classifying a nameserver's behaviour, so the
|
|
||||||
// caller's deadline must never be the thing that expires.
|
|
||||||
transportDeadline = 30 * time.Second
|
|
||||||
|
|
||||||
// silentNS and failingNS are the nameserver names reported
|
|
||||||
// back in NameserverResponse.Nameserver. They are .test names
|
|
||||||
// (RFC 6761) and are never resolved: the tests address the
|
|
||||||
// nameserver by its socket address.
|
|
||||||
silentNS = "silent.ns.test."
|
|
||||||
failingNS = "servfail.ns.test."
|
|
||||||
|
|
||||||
// transportHostname is the name queried. Nothing resolves it;
|
|
||||||
// the point is entirely how the nameserver behaves.
|
|
||||||
transportHostname = "example.com"
|
|
||||||
)
|
|
||||||
|
|
||||||
// startNameserver binds a real UDP nameserver on loopback and serves
|
|
||||||
// every datagram it receives with handle, which returns the reply to
|
|
||||||
// send or nil to stay silent. It returns the "host:port" address to
|
|
||||||
// aim a query at, and stops the server when the test ends.
|
|
||||||
func startNameserver(
|
|
||||||
t *testing.T,
|
|
||||||
handle func(query *dns.Msg) *dns.Msg,
|
|
||||||
) string {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
conn, err := net.ListenPacket("udp", "127.0.0.1:0")
|
|
||||||
require.NoError(t, err, "binding loopback nameserver")
|
|
||||||
|
|
||||||
stopped := make(chan struct{})
|
|
||||||
|
|
||||||
t.Cleanup(func() {
|
|
||||||
_ = conn.Close()
|
|
||||||
<-stopped
|
|
||||||
})
|
|
||||||
|
|
||||||
go serveNameserver(conn, handle, stopped)
|
|
||||||
|
|
||||||
return conn.LocalAddr().String()
|
|
||||||
}
|
|
||||||
|
|
||||||
// serveNameserver reads queries until conn is closed, replying with
|
|
||||||
// whatever handle produces.
|
|
||||||
func serveNameserver(
|
|
||||||
conn net.PacketConn,
|
|
||||||
handle func(query *dns.Msg) *dns.Msg,
|
|
||||||
stopped chan<- struct{},
|
|
||||||
) {
|
|
||||||
defer close(stopped)
|
|
||||||
|
|
||||||
buf := make([]byte, dns.MaxMsgSize)
|
|
||||||
|
|
||||||
for {
|
|
||||||
n, from, err := conn.ReadFrom(buf)
|
|
||||||
if err != nil {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
query := new(dns.Msg)
|
|
||||||
if query.Unpack(buf[:n]) != nil {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
reply := handle(query)
|
|
||||||
if reply == nil {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
wire, err := reply.Pack()
|
|
||||||
if err != nil {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
if _, err := conn.WriteTo(wire, from); err != nil {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// unservedAddr returns a loopback address with nothing listening on
|
|
||||||
// it, by binding a port and releasing it again.
|
|
||||||
func unservedAddr(t *testing.T) string {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
conn, err := net.ListenPacket("udp", "127.0.0.1:0")
|
|
||||||
require.NoError(t, err, "binding loopback port")
|
|
||||||
|
|
||||||
addr := conn.LocalAddr().String()
|
|
||||||
require.NoError(t, conn.Close(), "releasing loopback port")
|
|
||||||
|
|
||||||
return addr
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestQueryNameserverIP_Timeout covers the StatusTimeout branch: a
|
|
||||||
// nameserver that takes the query and never answers.
|
|
||||||
func TestQueryNameserverIP_Timeout(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
// A real nameserver that drops A queries and answers every
|
|
||||||
// other type. Silence on one type is all the resolver needs to
|
|
||||||
// classify the response as a timeout, and it keeps the test
|
|
||||||
// two query timeouts long instead of sixteen.
|
|
||||||
addr := startNameserver(t, func(query *dns.Msg) *dns.Msg {
|
|
||||||
if len(query.Question) > 0 &&
|
|
||||||
query.Question[0].Qtype == dns.TypeA {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
reply := new(dns.Msg)
|
|
||||||
reply.SetReply(query)
|
|
||||||
|
|
||||||
return reply
|
|
||||||
})
|
|
||||||
|
|
||||||
r := newTestResolver(t)
|
|
||||||
|
|
||||||
ctx, cancel := context.WithTimeout(
|
|
||||||
context.Background(), transportDeadline,
|
|
||||||
)
|
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
start := time.Now()
|
|
||||||
|
|
||||||
resp, err := r.QueryNameserverIP(
|
|
||||||
ctx, silentNS, addr, transportHostname,
|
|
||||||
)
|
|
||||||
|
|
||||||
elapsed := time.Since(start)
|
|
||||||
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NotNil(t, resp)
|
|
||||||
|
|
||||||
assert.Equal(t, resolver.StatusTimeout, resp.Status)
|
|
||||||
assert.Equal(t, "all queries timed out", resp.Error)
|
|
||||||
assert.Empty(t, resp.Records)
|
|
||||||
assert.Equal(t, silentNS, resp.Nameserver)
|
|
||||||
|
|
||||||
assert.Less(
|
|
||||||
t, elapsed, transportBudget,
|
|
||||||
"one silent record type must cost one query's retries, "+
|
|
||||||
"not every record type's",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestQueryNameserverIP_ServFail covers the StatusError branch: a
|
|
||||||
// nameserver that answers, and answers SERVFAIL.
|
|
||||||
func TestQueryNameserverIP_ServFail(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
addr := startNameserver(t, func(query *dns.Msg) *dns.Msg {
|
|
||||||
reply := new(dns.Msg)
|
|
||||||
reply.SetRcode(query, dns.RcodeServerFailure)
|
|
||||||
|
|
||||||
return reply
|
|
||||||
})
|
|
||||||
|
|
||||||
r := newTestResolver(t)
|
|
||||||
|
|
||||||
ctx, cancel := context.WithTimeout(
|
|
||||||
context.Background(), transportDeadline,
|
|
||||||
)
|
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
resp, err := r.QueryNameserverIP(
|
|
||||||
ctx, failingNS, addr, transportHostname,
|
|
||||||
)
|
|
||||||
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NotNil(t, resp)
|
|
||||||
|
|
||||||
assert.Equal(t, resolver.StatusError, resp.Status)
|
|
||||||
assert.Equal(t, "server returned SERVFAIL", resp.Error)
|
|
||||||
assert.Empty(t, resp.Records)
|
|
||||||
assert.Equal(t, failingNS, resp.Nameserver)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestQueryNameserverIP_NoListener pins the third transport outcome:
|
|
||||||
// a refused datagram is not a timeout. The socket fails immediately
|
|
||||||
// with ECONNREFUSED rather than going quiet, so isTimeout is false,
|
|
||||||
// no failure flag is set, and the response classifies as NoData.
|
|
||||||
// Asserting it here is what stops that path being mistaken for the
|
|
||||||
// timeout path, in either direction.
|
|
||||||
func TestQueryNameserverIP_NoListener(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
addr := unservedAddr(t)
|
|
||||||
|
|
||||||
r := newTestResolver(t)
|
|
||||||
|
|
||||||
ctx, cancel := context.WithTimeout(
|
|
||||||
context.Background(), transportDeadline,
|
|
||||||
)
|
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
resp, err := r.QueryNameserverIP(
|
|
||||||
ctx, silentNS, addr, transportHostname,
|
|
||||||
)
|
|
||||||
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NotNil(t, resp)
|
|
||||||
|
|
||||||
assert.Equal(t, resolver.StatusNoData, resp.Status)
|
|
||||||
assert.Empty(t, resp.Records)
|
|
||||||
}
|
|
||||||
@@ -0,0 +1,19 @@
|
|||||||
|
package server
|
||||||
|
|
||||||
|
import (
|
||||||
|
"net/http"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// NewHTTPServer exports newHTTPServer for testing.
|
||||||
|
func NewHTTPServer(
|
||||||
|
listenAddr string,
|
||||||
|
handler http.Handler,
|
||||||
|
) *http.Server {
|
||||||
|
return newHTTPServer(listenAddr, handler)
|
||||||
|
}
|
||||||
|
|
||||||
|
// RequestTimeout exports the handler execution budget applied by
|
||||||
|
// chimw.Timeout in SetupRoutes, so tests can assert the relationship
|
||||||
|
// between it and the server's WriteTimeout.
|
||||||
|
const RequestTimeout time.Duration = requestTimeout
|
||||||
@@ -33,8 +33,52 @@ type Params struct {
|
|||||||
// shutdownTimeout is how long to wait for graceful shutdown.
|
// shutdownTimeout is how long to wait for graceful shutdown.
|
||||||
const shutdownTimeout = 30 * time.Second
|
const shutdownTimeout = 30 * time.Second
|
||||||
|
|
||||||
// readHeaderTimeout is the max duration for reading request headers.
|
// Socket-level timeouts for the HTTP server.
|
||||||
const readHeaderTimeout = 10 * time.Second
|
//
|
||||||
|
// These bound time spent on the connection itself and are a distinct
|
||||||
|
// control from the per-request handler budget enforced by
|
||||||
|
// chimw.Timeout(requestTimeout) in routes.go: that one cancels the
|
||||||
|
// request context after requestTimeout but never touches the socket,
|
||||||
|
// so without the values below a peer can hold a connection open
|
||||||
|
// forever (slowloris, unreaped keep-alives).
|
||||||
|
//
|
||||||
|
// The one hard constraint between the two controls is
|
||||||
|
// writeTimeout > requestTimeout. net/http arms the write deadline
|
||||||
|
// once the request headers have been read, so on a plaintext
|
||||||
|
// connection it covers handler execution AND the response flush. If
|
||||||
|
// writeTimeout were <= requestTimeout the server would sever the
|
||||||
|
// connection before a handler that legitimately consumed its full
|
||||||
|
// budget could emit anything, making the 60s budget unreachable in
|
||||||
|
// practice. The margin between them is the response-flush allowance.
|
||||||
|
//
|
||||||
|
// The only clients of this service are browsers loading the dashboard
|
||||||
|
// and a Prometheus scraper; the values are sized for those.
|
||||||
|
const (
|
||||||
|
// readHeaderTimeout is the max duration for reading request
|
||||||
|
// headers.
|
||||||
|
readHeaderTimeout = 10 * time.Second
|
||||||
|
|
||||||
|
// readTimeout bounds reading the entire request, headers plus
|
||||||
|
// body. Every route here is a GET with no body, so this only
|
||||||
|
// ever needs to cover headers; the extra 5s over
|
||||||
|
// readHeaderTimeout is slack, not a real allowance, and keeps a
|
||||||
|
// body dribbled one byte at a time from holding the read side
|
||||||
|
// open indefinitely.
|
||||||
|
readTimeout = 15 * time.Second
|
||||||
|
|
||||||
|
// writeTimeout must exceed the requestTimeout handler budget
|
||||||
|
// (60s) per the note above. The 15s difference is the allowance
|
||||||
|
// for flushing a completed response to a slow client.
|
||||||
|
writeTimeout = 75 * time.Second
|
||||||
|
|
||||||
|
// idleTimeout reaps keep-alive connections between requests. It
|
||||||
|
// is deliberately longer than the common Prometheus scrape
|
||||||
|
// intervals (15s/30s/60s) so the scraper reuses its connection
|
||||||
|
// rather than reconnecting every cycle, while a browser tab
|
||||||
|
// left open on the dashboard stops occupying a connection
|
||||||
|
// within two minutes of going quiet.
|
||||||
|
idleTimeout = 120 * time.Second
|
||||||
|
)
|
||||||
|
|
||||||
// Server is the HTTP server.
|
// Server is the HTTP server.
|
||||||
type Server struct {
|
type Server struct {
|
||||||
@@ -76,16 +120,29 @@ func New(
|
|||||||
return srv, nil
|
return srv, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// newHTTPServer builds the listening http.Server with every
|
||||||
|
// socket-level timeout set. All four are set deliberately: a zero
|
||||||
|
// value in net/http means "no limit", not "some default".
|
||||||
|
func newHTTPServer(
|
||||||
|
listenAddr string,
|
||||||
|
handler http.Handler,
|
||||||
|
) *http.Server {
|
||||||
|
return &http.Server{
|
||||||
|
Addr: listenAddr,
|
||||||
|
Handler: handler,
|
||||||
|
ReadTimeout: readTimeout,
|
||||||
|
ReadHeaderTimeout: readHeaderTimeout,
|
||||||
|
WriteTimeout: writeTimeout,
|
||||||
|
IdleTimeout: idleTimeout,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Run starts the HTTP server.
|
// Run starts the HTTP server.
|
||||||
func (s *Server) Run() {
|
func (s *Server) Run() {
|
||||||
s.SetupRoutes()
|
s.SetupRoutes()
|
||||||
|
|
||||||
listenAddr := fmt.Sprintf(":%d", s.port)
|
listenAddr := fmt.Sprintf(":%d", s.port)
|
||||||
s.httpServer = &http.Server{
|
s.httpServer = newHTTPServer(listenAddr, s)
|
||||||
Addr: listenAddr,
|
|
||||||
Handler: s,
|
|
||||||
ReadHeaderTimeout: readHeaderTimeout,
|
|
||||||
}
|
|
||||||
|
|
||||||
s.log.Info("http server starting", "addr", listenAddr)
|
s.log.Info("http server starting", "addr", listenAddr)
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,112 @@
|
|||||||
|
package server_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"net/http"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"sneak.berlin/go/dnswatcher/internal/server"
|
||||||
|
)
|
||||||
|
|
||||||
|
// noopHandler stands in for the router; newHTTPServer only stores it.
|
||||||
|
func noopHandler() http.Handler {
|
||||||
|
return http.HandlerFunc(
|
||||||
|
func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
},
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestHTTPServerTimeoutsAreSet asserts that every socket-level
|
||||||
|
// timeout is configured. A zero value in net/http means "no limit",
|
||||||
|
// so a refactor that silently drops one of these reintroduces the
|
||||||
|
// slowloris / unreaped-keep-alive exposure this guards against.
|
||||||
|
//
|
||||||
|
// The assertions are on the configured field values only; nothing
|
||||||
|
// here measures elapsed time, so the test cannot flake on timing.
|
||||||
|
func TestHTTPServerTimeoutsAreSet(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
srv := server.NewHTTPServer(":8080", noopHandler())
|
||||||
|
|
||||||
|
if srv.ReadTimeout <= 0 {
|
||||||
|
t.Errorf(
|
||||||
|
"ReadTimeout must be non-zero, got %v",
|
||||||
|
srv.ReadTimeout,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
if srv.ReadHeaderTimeout <= 0 {
|
||||||
|
t.Errorf(
|
||||||
|
"ReadHeaderTimeout must be non-zero, got %v",
|
||||||
|
srv.ReadHeaderTimeout,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
if srv.WriteTimeout <= 0 {
|
||||||
|
t.Errorf(
|
||||||
|
"WriteTimeout must be non-zero, got %v",
|
||||||
|
srv.WriteTimeout,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
if srv.IdleTimeout <= 0 {
|
||||||
|
t.Errorf(
|
||||||
|
"IdleTimeout must be non-zero, got %v",
|
||||||
|
srv.IdleTimeout,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestWriteTimeoutExceedsHandlerBudget pins the one relationship the
|
||||||
|
// values must satisfy. net/http arms the write deadline once request
|
||||||
|
// headers are read, so it covers handler execution plus the response
|
||||||
|
// flush. If WriteTimeout were not greater than the chimw.Timeout
|
||||||
|
// handler budget, the connection would be severed before a handler
|
||||||
|
// that used its full budget could respond, making that budget
|
||||||
|
// unreachable.
|
||||||
|
func TestWriteTimeoutExceedsHandlerBudget(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
srv := server.NewHTTPServer(":8080", noopHandler())
|
||||||
|
|
||||||
|
if srv.WriteTimeout <= server.RequestTimeout {
|
||||||
|
t.Errorf(
|
||||||
|
"WriteTimeout (%v) must exceed handler budget (%v)",
|
||||||
|
srv.WriteTimeout,
|
||||||
|
server.RequestTimeout,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestReadTimeoutCoversHeaderTimeout asserts the read deadline for
|
||||||
|
// the whole request is at least as long as the header-only deadline;
|
||||||
|
// a smaller ReadTimeout would make ReadHeaderTimeout unreachable.
|
||||||
|
func TestReadTimeoutCoversHeaderTimeout(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
srv := server.NewHTTPServer(":8080", noopHandler())
|
||||||
|
|
||||||
|
if srv.ReadTimeout < srv.ReadHeaderTimeout {
|
||||||
|
t.Errorf(
|
||||||
|
"ReadTimeout (%v) must be >= ReadHeaderTimeout (%v)",
|
||||||
|
srv.ReadTimeout,
|
||||||
|
srv.ReadHeaderTimeout,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestHTTPServerAddrAndHandler covers the rest of the constructor so
|
||||||
|
// a future edit cannot drop the listen address or the handler.
|
||||||
|
func TestHTTPServerAddrAndHandler(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
srv := server.NewHTTPServer(":9999", noopHandler())
|
||||||
|
|
||||||
|
if srv.Addr != ":9999" {
|
||||||
|
t.Errorf("Addr = %q, want %q", srv.Addr, ":9999")
|
||||||
|
}
|
||||||
|
|
||||||
|
if srv.Handler == nil {
|
||||||
|
t.Error("Handler must not be nil")
|
||||||
|
}
|
||||||
|
}
|
||||||
+493
-647
File diff suppressed because it is too large
Load Diff
+2
-13
@@ -19,26 +19,15 @@
|
|||||||
#
|
#
|
||||||
# -timeout 90s is a deliberate backstop above the 60s hard cap on
|
# -timeout 90s is a deliberate backstop above the 60s hard cap on
|
||||||
# suite duration. Do not lower it.
|
# suite duration. Do not lower it.
|
||||||
#
|
|
||||||
# -p 1 runs one test package at a time, and is load-bearing. Live DNS
|
|
||||||
# is a resource outside the process: the concurrency gates that keep
|
|
||||||
# this suite from bursting at the root servers
|
|
||||||
# (internal/resolver/livedns_test.go, internal/watcher/watcher_test.go)
|
|
||||||
# are package-scoped, so each one only bounds its own test binary. Go
|
|
||||||
# runs package binaries in parallel by default, so with both live-DNS
|
|
||||||
# packages in flight at once their gates sum instead of holding, the
|
|
||||||
# root and TLD servers rate-limit the excess, and the resolver
|
|
||||||
# package's per-attempt deadlines expire. Serialising packages is what
|
|
||||||
# makes each gate authoritative while its package runs.
|
|
||||||
set -eu
|
set -eu
|
||||||
|
|
||||||
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
||||||
|
|
||||||
main() {
|
main() {
|
||||||
cd "$ROOT"
|
cd "$ROOT"
|
||||||
go test -count=1 -p 1 -race -timeout 90s -cover ./... || {
|
go test -count=1 -race -timeout 90s -cover ./... || {
|
||||||
echo "--- Rerunning with -v for details ---" >&2
|
echo "--- Rerunning with -v for details ---" >&2
|
||||||
go test -count=1 -p 1 -race -timeout 90s -v ./... || true
|
go test -count=1 -race -timeout 90s -v ./... || true
|
||||||
exit 1
|
exit 1
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user