4 Commits

Author SHA1 Message Date
clawbot
cd06bba034 notify: fix flaky drain timing assertions, drop false abandon warn (closes #106)
All checks were successful
check / check (push) Successful in 30s
The drain tests measured elapsed time from an instant captured after
the clock they compared it against had already started, so the lower
bounds were structurally unreachable and passed only when the gap
between the two statements rounded to zero. TestDrainBoundedByContext-
Deadline failed the Docker gate outright (49.9ms against its own 50ms
deadline) and roughly 1 run in 12 locally.

- TestDrainBoundedByContextDeadline: capture start before
  context.WithTimeout, so the measured interval is a superset of the
  deadline interval and only an early return can fail the lower bound.
  The upper bound moves to a watchdog around the drain, which turns an
  unbounded drain into a prompt failure instead of a package-timeout
  hang.
- TestDrainWaitsForInFlightDelivery: same ordering fix, ahead of the
  timer that releases the held delivery.
- TestDrainWithoutDeliveriesReturnsImmediately: its 50ms ceiling was
  under the observed cost of the goroutine hop through inFlight.Wait()
  on a loaded box (57ms), and failed once in 20 runs. It now bounds the
  idle drain at 500ms, still well under the 2s deadline a stalled drain
  would hit.
- drain: an OnStop context already expired on entry with nothing
  outstanding logged a WARN about abandoning deliveries with
  abandoned=0 and closed the abandon channel for no reason. The timeout
  branch now reports at debug level when the outstanding count is zero,
  and warns only when deliveries genuinely are abandoned.
  TestDrainWithCancelledContextDoesNotWarn covers it.

Verified: script/cibuild passes; 25 consecutive cache-bypassed
make test runs under -race, all clean; make check green at 5.2s.
Both corrected assertions were confirmed non-vacuous by temporarily
breaking drain and watching them fail.
2026-08-09 05:21:57 +00:00
clawbot
970ea9fae8 notify: drain in-flight deliveries at shutdown (closes #106)
All checks were successful
check / check (push) Successful in 33s
notify.New accepted an fx.Lifecycle and never used it, so the three
dispatch goroutines were untracked. context.WithoutCancel kept a
delivery alive past its caller's cancellation but made nothing wait
for it: the process could exit while a delivery was still in its
retry backoff (up to five attempts, 60s max delay), silently losing
exactly the alert most worth keeping.

Deliveries are now tracked in a sync.WaitGroup whose counter is
incremented on the dispatching goroutine before the worker starts,
and notify.New registers an OnStop hook that drains them. The drain
is bounded by the context fx passes to OnStop; when it expires with
work outstanding, the count is logged at warn level and parked retry
backoffs are released via an abandon channel so they stop retrying
rather than outliving the drain. Deliveries submitted after the
drain has begun are refused and logged, so a stream of new
notifications cannot extend shutdown indefinitely.

The three near-identical dispatchers now share one tracked dispatch
helper. Tests use httptest servers and the existing retry knobs
(SetRetryConfig/SetSleepFunc) so nothing waits on a real backoff.

README's shutdown claim is reworded to match the bounded semantics.
2026-08-09 05:03:23 +00:00
9347a2838b build: update golangci-lint to v2.12.2 with org-standard v2 config (#96)
All checks were successful
check / check (push) Successful in 4s
Updates golangci-lint to v2.12.2 and sets `.golangci.yml` to the org-standard v2-schema config already deployed across the org's repos. The config change is owner-authorized (see #96 (comment) and #96 (comment)); the same file is being landed as canonical via prompts PR #24 (sneak/prompts#24).

## Changes

- **Commit-pinned installs**: golangci-lint pinned to commit `c0d3ddc9cf3faa61a4e378e879ece580256d76e5` (v2.12.2, released 2026-05-06) in `Dockerfile` and `script/bootstrap`.
- **`.golangci.yml` set to the org-standard v2 config** (sha256 `021cc83f4e6fc7c31b95b34b846723dfcf20b66b7baeea1dc40406e643346bcb`), byte-identical to the file used across the org's other repos. Settings live under `linters.settings`, so the `lll`/`funlen`/`cyclop`/`dupl` thresholds are actually applied (under the old hybrid file, v2 silently ignored the top-level `linters-settings` block).
- **Lint fixes** required by the now-active thresholds:
  - `goconst`: shared constants for repeated status/priority/DNS-fixture strings in `internal/watcher/watcher.go` and the notify, state, and watcher tests
  - `dupl`: consolidated duplicated ntfy/slack HTTP-error tests and SendNotification endpoint-error tests behind shared helpers in `internal/notify/delivery_test.go`
  - `lll`: wrapped long test table entries and comments in `internal/config/classify_test.go`, `internal/notify/history_test.go`, `internal/state/state_test.go`, `internal/watcher/watcher_test.go`; shortened one inline nolint justification in `internal/notify/retry.go`
- **`TODO.md`**: Completed Steps entry updated in the same commit.
- Rebased onto current `main` (`f79cd98`); the branch is one clean commit.

## Notes

- v2.12 deprecates the `gomodguard` linter in favor of `gomodguard_v2`. The org-standard config does not disable the deprecated linter, so golangci-lint may emit an informational deprecation warning; this is accepted by the owner and does not affect the exit status (this exact config+code combination was CI-green at `dea7e44`).

## Verification

- `make check` exits 0 (fmt-check, tests, lint)
- `make lint`: 0 issues; no deprecation warning surfaced in the runs performed
- sha256 of `.golangci.yml` at HEAD verified equal to `021cc83f4e6fc7c31b95b34b846723dfcf20b66b7baeea1dc40406e643346bcb`

Co-authored-by: sneak <sneak@sneak.berlin>
Reviewed-on: #96
Co-authored-by: clawbot <clawbot@noreply.example.org>
Co-committed-by: clawbot <clawbot@noreply.example.org>
2026-08-07 23:15:47 +02:00
f79cd98107 docs: document the no-DNS-mocking policy in README (closes #94) (#95)
All checks were successful
check / check (push) Successful in 5s
Adds a prominent "No DNS mocking. Ever." section near the top of `README.md`, per owner policy (sneak, 2026-08-07):

- DNS is never mocked in this project — no mock resolvers, fake DNS servers, or stubbed lookups, in tests or anywhere else.
- Tests exercise real iterative resolution against live nameservers by design.
- Flaky live tests are fixed with robustness (retries, multiple nameservers, timeouts) or explicit opt-in gating decided by the owner — never with mocks.
- Contributions introducing DNS mocks will be rejected.

Markdown-only change; matches the README's existing tone and hard-wrap style. `script/fmt` covers Go only, so no formatter output applies to this file. Verified via `script/cibuild` (docker build runs `make check` with the pinned toolchain) — green. A direct local `make check` shows 21 pre-existing `goconst` lint findings that come from a newer local `golangci-lint` (v2.12.2 vs the pinned v2.10.1) and are unrelated to this change.

Related: #93 is being reframed under this policy.
Co-authored-by: sneak <sneak@sneak.berlin>
Reviewed-on: #95
Co-authored-by: clawbot <clawbot@noreply.example.org>
Co-committed-by: clawbot <clawbot@noreply.example.org>
2026-08-07 22:31:48 +02:00
16 changed files with 1132 additions and 364 deletions

View File

@@ -1,5 +1,9 @@
version: "2" version: "2"
# Config schema uses the golangci-lint v2 layout (settings live under
# linters.settings, not top-level linters-settings) so that the
# thresholds below are actually applied by golangci-lint >= v2.
run: run:
timeout: 5m timeout: 5m
modules-download-mode: readonly modules-download-mode: readonly
@@ -14,8 +18,7 @@ linters:
- wsl # Deprecated, replaced by wsl_v5 - wsl # Deprecated, replaced by wsl_v5
- wrapcheck # Too verbose for internal packages - wrapcheck # Too verbose for internal packages
- varnamelen # Short names like db, id are idiomatic Go - varnamelen # Short names like db, id are idiomatic Go
settings:
linters-settings:
lll: lll:
line-length: 88 line-length: 88
funlen: funlen:
@@ -27,6 +30,5 @@ linters-settings:
threshold: 100 threshold: 100
issues: issues:
exclude-use-default: false
max-issues-per-linter: 0 max-issues-per-linter: 0
max-same-issues: 0 max-same-issues: 0

View File

@@ -4,8 +4,8 @@ FROM golang@sha256:f6751d823c26342f9506c03797d2527668d095b0a15f1862cddb4d927a7a4
RUN apk add --no-cache git make gcc musl-dev binutils-gold RUN apk add --no-cache git make gcc musl-dev binutils-gold
# golangci-lint v2.10.1 # golangci-lint v2.12.2, 2026-08-07
RUN go install github.com/golangci/golangci-lint/v2/cmd/golangci-lint@5d1e709b7be35cb2025444e19de266b056b7b7ee RUN go install github.com/golangci/golangci-lint/v2/cmd/golangci-lint@c0d3ddc9cf3faa61a4e378e879ece580256d76e5
# goimports v0.42.0 # goimports v0.42.0
RUN go install golang.org/x/tools/cmd/goimports@009367f5c17a8d4c45a961a3a509277190a9a6f0 RUN go install golang.org/x/tools/cmd/goimports@009367f5c17a8d4c45a961a3a509277190a9a6f0

View File

@@ -17,6 +17,26 @@ without requiring an external database.
--- ---
## No DNS mocking. Ever.
**DNS is never mocked in this project — not in tests, not anywhere else.**
No mock resolvers, no fake DNS servers, no stubbed lookups.
dnswatcher's entire purpose is correct behavior against the real DNS.
Tests exercise real iterative resolution against live nameservers by
design; a test suite that passes against a mock proves nothing about the
one thing this program exists to do.
When live tests are flaky, that is a robustness problem, and it gets
fixed with robustness: retries with backoff, querying multiple
independent nameservers, longer timeouts — or explicit opt-in gating
decided by the project owner. Never with mocks.
Contributions that introduce mocked, faked, or stubbed DNS will be
rejected.
---
## Features ## Features
### DNS Domain Monitoring (Apex Domains) ### DNS Domain Monitoring (Apex Domains)
@@ -198,7 +218,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.
--- ---
@@ -432,8 +453,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.
--- ---

17
TODO.md
View File

@@ -25,6 +25,23 @@ confirm make check still passes.
# Completed Steps # Completed Steps
- 2026-08-09: in-flight notification deliveries are now drained at
shutdown (#106): `notify.New` registers an fx `OnStop` hook that waits
on a `sync.WaitGroup` of tracked delivery goroutines, bounded by the
`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-07: golangci-lint bumped to v2.12.2 (commit-pinned installs
in `Dockerfile` and `script/bootstrap`); `.golangci.yml` set to the
org-standard v2-schema config used across the org's repos
(owner-authorized; same file is being landed as canonical via prompts
PR #24), with settings under `linters.settings` so the
lll/funlen/cyclop/dupl thresholds apply; fixed the resulting
`goconst`, `dupl`, and `lll` findings; the informational `gomodguard`
deprecation warning under this config is accepted
- 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

View File

@@ -17,13 +17,33 @@ func TestClassifyDNSName(t *testing.T) {
}{ }{
{name: "apex domain simple", input: "example.com", want: config.DNSNameTypeDomain}, {name: "apex domain simple", input: "example.com", want: config.DNSNameTypeDomain},
{name: "hostname simple", input: "www.example.com", want: config.DNSNameTypeHostname}, {name: "hostname simple", input: "www.example.com", want: config.DNSNameTypeHostname},
{name: "apex domain multi-part TLD", input: "example.co.uk", want: config.DNSNameTypeDomain}, {
{name: "hostname multi-part TLD", input: "api.example.co.uk", want: config.DNSNameTypeHostname}, name: "apex domain multi-part TLD",
input: "example.co.uk",
want: config.DNSNameTypeDomain,
},
{
name: "hostname multi-part TLD",
input: "api.example.co.uk",
want: config.DNSNameTypeHostname,
},
{name: "public suffix itself", input: "co.uk", wantErr: true}, {name: "public suffix itself", input: "co.uk", wantErr: true},
{name: "empty string", input: "", wantErr: true}, {name: "empty string", input: "", wantErr: true},
{name: "deeply nested hostname", input: "a.b.c.example.com", want: config.DNSNameTypeHostname}, {
{name: "trailing dot stripped", input: "example.com.", want: config.DNSNameTypeDomain}, name: "deeply nested hostname",
{name: "uppercase normalized", input: "WWW.Example.COM", want: config.DNSNameTypeHostname}, input: "a.b.c.example.com",
want: config.DNSNameTypeHostname,
},
{
name: "trailing dot stripped",
input: "example.com.",
want: config.DNSNameTypeDomain,
},
{
name: "uppercase normalized",
input: "WWW.Example.COM",
want: config.DNSNameTypeHostname,
},
} }
for _, tt := range tests { for _, tt := range tests {

View File

@@ -25,6 +25,20 @@ const (
colorDefault = "#6c757d" colorDefault = "#6c757d"
) )
// Priority strings used across multiple tests.
const (
prioError = "error"
prioWarning = "warning"
prioSuccess = "success"
prioInfo = "info"
prioUnknown = "unknown"
prioDefault = "default"
prioUrgent = "urgent"
)
// testHost is the hostname used in request construction tests.
const testHost = "example.com"
// errSimulated is a static error for transport failures. // errSimulated is a static error for transport failures.
var errSimulated = errors.New("simulated transport failure") var errSimulated = errors.New("simulated transport failure")
@@ -101,13 +115,13 @@ func TestNtfyPriority(t *testing.T) {
input string input string
want string want string
}{ }{
{"error", "urgent"}, {prioError, prioUrgent},
{"warning", "high"}, {prioWarning, "high"},
{"success", "default"}, {prioSuccess, prioDefault},
{"info", "low"}, {prioInfo, "low"},
{"", "default"}, {"", prioDefault},
{"unknown", "default"}, {prioUnknown, prioDefault},
{"critical", "default"}, {"critical", prioDefault},
} }
for _, tc := range cases { for _, tc := range cases {
@@ -134,12 +148,12 @@ func TestSlackColor(t *testing.T) {
input string input string
want string want string
}{ }{
{"error", colorError}, {prioError, colorError},
{"warning", colorWarning}, {prioWarning, colorWarning},
{"success", colorSuccess}, {prioSuccess, colorSuccess},
{"info", colorInfo}, {prioInfo, colorInfo},
{"", colorDefault}, {"", colorDefault},
{"unknown", colorDefault}, {prioUnknown, colorDefault},
{"critical", colorDefault}, {"critical", colorDefault},
} }
@@ -165,7 +179,7 @@ func TestNewRequest(t *testing.T) {
target := &url.URL{ target := &url.URL{
Scheme: "https", Scheme: "https",
Host: "example.com", Host: testHost,
Path: "/webhook", Path: "/webhook",
} }
body := bytes.NewBufferString("hello") body := bytes.NewBufferString("hello")
@@ -187,9 +201,9 @@ func TestNewRequest(t *testing.T) {
) )
} }
if req.Host != "example.com" { if req.Host != testHost {
t.Errorf( t.Errorf(
"Host = %q, want %q", req.Host, "example.com", "Host = %q, want %q", req.Host, testHost,
) )
} }
@@ -217,7 +231,7 @@ func TestNewRequestPreservesContext(t *testing.T) {
ctxKey("k"), ctxKey("k"),
"v", "v",
) )
target := &url.URL{Scheme: "https", Host: "example.com"} target := &url.URL{Scheme: "https", Host: testHost}
req := notify.NewRequestForTest( req := notify.NewRequestForTest(
ctx, http.MethodGet, target, http.NoBody, ctx, http.MethodGet, target, http.NoBody,
@@ -289,10 +303,10 @@ func TestSendNtfyHeaders(t *testing.T) {
) )
} }
if captured.priority != "urgent" { if captured.priority != prioUrgent {
t.Errorf( t.Errorf(
"Priority header = %q, want %q", "Priority header = %q, want %q",
captured.priority, "urgent", captured.priority, prioUrgent,
) )
} }
@@ -311,10 +325,10 @@ func TestSendNtfyAllPriorities(t *testing.T) {
input string input string
want string want string
}{ }{
{"error", "urgent"}, {prioError, prioUrgent},
{"warning", "high"}, {prioWarning, "high"},
{"success", "default"}, {prioSuccess, prioDefault},
{"info", "low"}, {prioInfo, "low"},
} }
for _, tc := range priorities { for _, tc := range priorities {
@@ -356,56 +370,69 @@ func TestSendNtfyAllPriorities(t *testing.T) {
} }
} }
func TestSendNtfyClientError(t *testing.T) { // assertSendStatusError verifies that send returns an error
t.Parallel() // wrapping wantErr when the server responds with status.
func assertSendStatusError(
t *testing.T,
status int,
wantErr error,
send func(*notify.Service, *url.URL) error,
) {
t.Helper()
srv := httptest.NewServer( srv := httptest.NewServer(
http.HandlerFunc( http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) { func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusForbidden) w.WriteHeader(status)
}), }),
) )
defer srv.Close() defer srv.Close()
svc := notify.NewTestService(srv.Client().Transport) svc := notify.NewTestService(srv.Client().Transport)
topicURL, _ := url.Parse(srv.URL) target, _ := url.Parse(srv.URL)
err := svc.SendNtfy( err := send(svc, target)
context.Background(), topicURL, "t", "m", "info",
)
if err == nil { if err == nil {
t.Fatal("expected error for 403 response") t.Fatalf("expected error for %d response", status)
} }
if !errors.Is(err, notify.ErrNtfyFailed) { if !errors.Is(err, wantErr) {
t.Errorf("error = %v, want ErrNtfyFailed", err) t.Errorf("error = %v, want %v", err, wantErr)
} }
} }
func sendNtfyInfo(
svc *notify.Service, target *url.URL,
) error {
return svc.SendNtfy(
context.Background(), target, "t", "m", prioInfo,
)
}
func sendSlackInfo(
svc *notify.Service, target *url.URL,
) error {
return svc.SendSlack(
context.Background(), target, "t", "m", prioInfo,
)
}
func TestSendNtfyClientError(t *testing.T) {
t.Parallel()
assertSendStatusError(
t, http.StatusForbidden,
notify.ErrNtfyFailed, sendNtfyInfo,
)
}
func TestSendNtfyServerError(t *testing.T) { func TestSendNtfyServerError(t *testing.T) {
t.Parallel() t.Parallel()
srv := httptest.NewServer( assertSendStatusError(
http.HandlerFunc( t, http.StatusInternalServerError,
func(w http.ResponseWriter, _ *http.Request) { notify.ErrNtfyFailed, sendNtfyInfo,
w.WriteHeader(http.StatusInternalServerError)
}),
) )
defer srv.Close()
svc := notify.NewTestService(srv.Client().Transport)
topicURL, _ := url.Parse(srv.URL)
err := svc.SendNtfy(
context.Background(), topicURL, "t", "m", "info",
)
if err == nil {
t.Fatal("expected error for 500 response")
}
if !errors.Is(err, notify.ErrNtfyFailed) {
t.Errorf("error = %v, want ErrNtfyFailed", err)
}
} }
func TestSendNtfySuccess(t *testing.T) { func TestSendNtfySuccess(t *testing.T) {
@@ -550,11 +577,11 @@ func TestSendSlackAllColors(t *testing.T) {
priority string priority string
want string want string
}{ }{
{"error", colorError}, {prioError, colorError},
{"warning", colorWarning}, {prioWarning, colorWarning},
{"success", colorSuccess}, {prioSuccess, colorSuccess},
{"info", colorInfo}, {prioInfo, colorInfo},
{"unknown", colorDefault}, {prioUnknown, colorDefault},
} }
for _, tc := range colors { for _, tc := range colors {
@@ -606,53 +633,19 @@ func TestSendSlackAllColors(t *testing.T) {
func TestSendSlackClientError(t *testing.T) { func TestSendSlackClientError(t *testing.T) {
t.Parallel() t.Parallel()
srv := httptest.NewServer( assertSendStatusError(
http.HandlerFunc( t, http.StatusBadRequest,
func(w http.ResponseWriter, _ *http.Request) { notify.ErrSlackFailed, sendSlackInfo,
w.WriteHeader(http.StatusBadRequest)
}),
) )
defer srv.Close()
svc := notify.NewTestService(srv.Client().Transport)
webhookURL, _ := url.Parse(srv.URL)
err := svc.SendSlack(
context.Background(), webhookURL, "t", "m", "info",
)
if err == nil {
t.Fatal("expected error for 400 response")
}
if !errors.Is(err, notify.ErrSlackFailed) {
t.Errorf("error = %v, want ErrSlackFailed", err)
}
} }
func TestSendSlackServerError(t *testing.T) { func TestSendSlackServerError(t *testing.T) {
t.Parallel() t.Parallel()
srv := httptest.NewServer( assertSendStatusError(
http.HandlerFunc( t, http.StatusBadGateway,
func(w http.ResponseWriter, _ *http.Request) { notify.ErrSlackFailed, sendSlackInfo,
w.WriteHeader(http.StatusBadGateway)
}),
) )
defer srv.Close()
svc := notify.NewTestService(srv.Client().Transport)
webhookURL, _ := url.Parse(srv.URL)
err := svc.SendSlack(
context.Background(), webhookURL, "t", "m", "error",
)
if err == nil {
t.Fatal("expected error for 502 response")
}
if !errors.Is(err, notify.ErrSlackFailed) {
t.Errorf("error = %v, want ErrSlackFailed", err)
}
} }
func TestSendSlackNetworkError(t *testing.T) { func TestSendSlackNetworkError(t *testing.T) {
@@ -977,74 +970,62 @@ func TestSendNotificationMattermostOnly(t *testing.T) {
} }
} }
func TestSendNotificationNtfyError(t *testing.T) { // assertSendNotificationTolerates verifies SendNotification
t.Parallel() // neither panics nor blocks when the endpoint configured by
// setURL responds with status.
func assertSendNotificationTolerates(
t *testing.T,
status int,
priority string,
setURL func(*notify.Service, *url.URL),
) {
t.Helper()
srv := httptest.NewServer( srv := httptest.NewServer(
http.HandlerFunc( http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) { func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusInternalServerError) w.WriteHeader(status)
}), }),
) )
defer srv.Close() defer srv.Close()
ntfyURL, _ := url.Parse(srv.URL) target, _ := url.Parse(srv.URL)
svc := notify.NewTestService(http.DefaultTransport) svc := notify.NewTestService(http.DefaultTransport)
svc.SetNtfyURL(ntfyURL) setURL(svc, target)
// Should not panic or block.
svc.SendNotification( svc.SendNotification(
context.Background(), "t", "m", "error", context.Background(), "t", "m", priority,
) )
time.Sleep(100 * time.Millisecond) time.Sleep(100 * time.Millisecond)
} }
func TestSendNotificationNtfyError(t *testing.T) {
t.Parallel()
assertSendNotificationTolerates(
t, http.StatusInternalServerError, prioError,
(*notify.Service).SetNtfyURL,
)
}
func TestSendNotificationSlackError(t *testing.T) { func TestSendNotificationSlackError(t *testing.T) {
t.Parallel() t.Parallel()
srv := httptest.NewServer( assertSendNotificationTolerates(
http.HandlerFunc( t, http.StatusForbidden, prioError,
func(w http.ResponseWriter, _ *http.Request) { (*notify.Service).SetSlackWebhookURL,
w.WriteHeader(http.StatusForbidden)
}),
) )
defer srv.Close()
slackURL, _ := url.Parse(srv.URL)
svc := notify.NewTestService(http.DefaultTransport)
svc.SetSlackWebhookURL(slackURL)
svc.SendNotification(
context.Background(), "t", "m", "error",
)
time.Sleep(100 * time.Millisecond)
} }
func TestSendNotificationMattermostError(t *testing.T) { func TestSendNotificationMattermostError(t *testing.T) {
t.Parallel() t.Parallel()
srv := httptest.NewServer( assertSendNotificationTolerates(
http.HandlerFunc( t, http.StatusBadGateway, prioWarning,
func(w http.ResponseWriter, _ *http.Request) { (*notify.Service).SetMattermostWebhookURL,
w.WriteHeader(http.StatusBadGateway)
}),
) )
defer srv.Close()
mmURL, _ := url.Parse(srv.URL)
svc := notify.NewTestService(http.DefaultTransport)
svc.SetMattermostWebhookURL(mmURL)
svc.SendNotification(
context.Background(), "t", "m", "warning",
)
time.Sleep(100 * time.Millisecond)
} }
// ── SlackPayload JSON marshaling ────────────────────────── // ── SlackPayload JSON marshaling ──────────────────────────

View File

@@ -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.

View File

@@ -29,14 +29,14 @@ func TestAlertHistoryAddAndRecent(t *testing.T) {
Timestamp: now.Add(-2 * time.Minute), Timestamp: now.Add(-2 * time.Minute),
Title: "first", Title: "first",
Message: "msg1", Message: "msg1",
Priority: "info", Priority: prioInfo,
}) })
h.Add(notify.AlertEntry{ h.Add(notify.AlertEntry{
Timestamp: now.Add(-1 * time.Minute), Timestamp: now.Add(-1 * time.Minute),
Title: "second", Title: "second",
Message: "msg2", Message: "msg2",
Priority: "warning", Priority: prioWarning,
}) })
entries := h.Recent() entries := h.Recent()

View File

@@ -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(

View File

@@ -2,6 +2,7 @@ package notify
import ( import (
"context" "context"
"fmt"
"math" "math"
"math/rand/v2" "math/rand/v2"
"time" "time"
@@ -69,7 +70,7 @@ func (rc RetryConfig) backoff(attempt int) time.Duration {
lo := raw * (1 - jitterFraction) lo := raw * (1 - jitterFraction)
hi := raw * (1 + jitterFraction) hi := raw * (1 + jitterFraction)
jittered := lo + rand.Float64()*(hi-lo) //nolint:gosec // jitter does not need crypto/rand jittered := lo + rand.Float64()*(hi-lo) //nolint:gosec // jitter needs no crypto/rand
return time.Duration(jittered) return time.Duration(jittered)
} }
@@ -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):
} }
} }

119
internal/notify/shutdown.go Normal file
View File

@@ -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(),
)
}
}

View File

@@ -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,
)
}
}

View File

@@ -13,6 +13,16 @@ import (
const testHostname = "www.example.com" const testHostname = "www.example.com"
// Shared fixture values used across tests.
const (
testNS1 = "ns1.example.com."
testNS2 = "ns2.example.com."
testAltNS1 = "ns1.test.com."
testIPv4 = "93.184.216.34"
testIP = "1.2.3.4"
statusError = "error"
)
// populateState fills a State with representative test data across all categories. // populateState fills a State with representative test data across all categories.
func populateState(t *testing.T, s *state.State) { func populateState(t *testing.T, s *state.State) {
t.Helper() t.Helper()
@@ -20,7 +30,7 @@ func populateState(t *testing.T, s *state.State) {
now := time.Now().UTC().Truncate(time.Second) now := time.Now().UTC().Truncate(time.Second)
s.SetDomainState("example.com", &state.DomainState{ s.SetDomainState("example.com", &state.DomainState{
Nameservers: []string{"ns1.example.com.", "ns2.example.com."}, Nameservers: []string{testNS1, testNS2},
LastChecked: now, LastChecked: now,
}) })
@@ -31,17 +41,17 @@ func populateState(t *testing.T, s *state.State) {
s.SetHostnameState(testHostname, &state.HostnameState{ s.SetHostnameState(testHostname, &state.HostnameState{
RecordsByNameserver: map[string]*state.NameserverRecordState{ RecordsByNameserver: map[string]*state.NameserverRecordState{
"ns1.example.com.": { testNS1: {
Records: map[string][]string{ Records: map[string][]string{
"A": {"93.184.216.34"}, "A": {testIPv4},
"AAAA": {"2606:2800:220:1:248:1893:25c8:1946"}, "AAAA": {"2606:2800:220:1:248:1893:25c8:1946"},
}, },
Status: "ok", Status: "ok",
LastChecked: now, LastChecked: now,
}, },
"ns2.example.com.": { testNS2: {
Records: map[string][]string{ Records: map[string][]string{
"A": {"93.184.216.34"}, "A": {testIPv4},
}, },
Status: "ok", Status: "ok",
LastChecked: now, LastChecked: now,
@@ -152,13 +162,13 @@ func TestSaveLoadRoundTrip_Hostnames(t *testing.T) {
func verifyNS1Records(t *testing.T, hn *state.HostnameState) { func verifyNS1Records(t *testing.T, hn *state.HostnameState) {
t.Helper() t.Helper()
ns1, ok := hn.RecordsByNameserver["ns1.example.com."] ns1, ok := hn.RecordsByNameserver[testNS1]
if !ok { if !ok {
t.Fatal("missing nameserver ns1.example.com.") t.Fatal("missing nameserver ns1.example.com.")
} }
aRecords := ns1.Records["A"] aRecords := ns1.Records["A"]
if len(aRecords) != 1 || aRecords[0] != "93.184.216.34" { if len(aRecords) != 1 || aRecords[0] != testIPv4 {
t.Errorf("ns1 A records: got %v", aRecords) t.Errorf("ns1 A records: got %v", aRecords)
} }
@@ -213,7 +223,8 @@ func TestSaveLoadRoundTrip_Ports(t *testing.T) {
} }
} }
// TestSaveLoadRoundTrip_Certificates verifies certificate data survives a save/load cycle. // TestSaveLoadRoundTrip_Certificates verifies certificate data
// survives a save/load cycle.
func TestSaveLoadRoundTrip_Certificates(t *testing.T) { func TestSaveLoadRoundTrip_Certificates(t *testing.T) {
t.Parallel() t.Parallel()
@@ -653,7 +664,7 @@ func TestDomainState_GetSet(t *testing.T) {
now := time.Now().UTC().Truncate(time.Second) now := time.Now().UTC().Truncate(time.Second)
ds := &state.DomainState{ ds := &state.DomainState{
Nameservers: []string{"ns1.test.com."}, Nameservers: []string{testAltNS1},
LastChecked: now, LastChecked: now,
} }
@@ -664,7 +675,7 @@ func TestDomainState_GetSet(t *testing.T) {
t.Fatal("expected true for existing domain") t.Fatal("expected true for existing domain")
} }
if len(got.Nameservers) != 1 || got.Nameservers[0] != "ns1.test.com." { if len(got.Nameservers) != 1 || got.Nameservers[0] != testAltNS1 {
t.Errorf("nameservers: got %v", got.Nameservers) t.Errorf("nameservers: got %v", got.Nameservers)
} }
@@ -674,7 +685,7 @@ func TestDomainState_GetSet(t *testing.T) {
// Overwrite. // Overwrite.
ds2 := &state.DomainState{ ds2 := &state.DomainState{
Nameservers: []string{"ns1.test.com.", "ns2.test.com."}, Nameservers: []string{testAltNS1, "ns2.test.com."},
LastChecked: now.Add(time.Hour), LastChecked: now.Add(time.Hour),
} }
@@ -704,8 +715,8 @@ func TestHostnameState_GetSet(t *testing.T) {
now := time.Now().UTC().Truncate(time.Second) now := time.Now().UTC().Truncate(time.Second)
hs := &state.HostnameState{ hs := &state.HostnameState{
RecordsByNameserver: map[string]*state.NameserverRecordState{ RecordsByNameserver: map[string]*state.NameserverRecordState{
"ns1.example.com.": { testNS1: {
Records: map[string][]string{"A": {"1.2.3.4"}}, Records: map[string][]string{"A": {testIP}},
Status: "ok", Status: "ok",
LastChecked: now, LastChecked: now,
}, },
@@ -720,7 +731,7 @@ func TestHostnameState_GetSet(t *testing.T) {
t.Fatal("expected true for existing hostname") t.Fatal("expected true for existing hostname")
} }
nsState, ok := got.RecordsByNameserver["ns1.example.com."] nsState, ok := got.RecordsByNameserver[testNS1]
if !ok { if !ok {
t.Fatal("missing nameserver entry") t.Fatal("missing nameserver entry")
} }
@@ -730,7 +741,7 @@ func TestHostnameState_GetSet(t *testing.T) {
} }
aRecords := nsState.Records["A"] aRecords := nsState.Records["A"]
if len(aRecords) != 1 || aRecords[0] != "1.2.3.4" { if len(aRecords) != 1 || aRecords[0] != testIP {
t.Errorf("A records: got %v", aRecords) t.Errorf("A records: got %v", aRecords)
} }
} }
@@ -869,7 +880,7 @@ func TestCertificateState_ErrorField(t *testing.T) {
now := time.Now().UTC().Truncate(time.Second) now := time.Now().UTC().Truncate(time.Second)
cs := &state.CertificateState{ cs := &state.CertificateState{
Status: "error", Status: statusError,
Error: "connection refused", Error: "connection refused",
LastChecked: now, LastChecked: now,
} }
@@ -893,8 +904,8 @@ func TestCertificateState_ErrorField(t *testing.T) {
t.Fatal("missing certificate after load") t.Fatal("missing certificate after load")
} }
if got.Status != "error" { if got.Status != statusError {
t.Errorf("status: got %q, want %q", got.Status, "error") t.Errorf("status: got %q, want %q", got.Status, statusError)
} }
if got.Error != "connection refused" { if got.Error != "connection refused" {
@@ -912,9 +923,9 @@ func TestHostnameState_ErrorField(t *testing.T) {
now := time.Now().UTC().Truncate(time.Second) now := time.Now().UTC().Truncate(time.Second)
hs := &state.HostnameState{ hs := &state.HostnameState{
RecordsByNameserver: map[string]*state.NameserverRecordState{ RecordsByNameserver: map[string]*state.NameserverRecordState{
"ns1.example.com.": { testNS1: {
Records: nil, Records: nil,
Status: "error", Status: statusError,
Error: "SERVFAIL", Error: "SERVFAIL",
LastChecked: now, LastChecked: now,
}, },
@@ -941,9 +952,9 @@ func TestHostnameState_ErrorField(t *testing.T) {
t.Fatal("missing hostname after load") t.Fatal("missing hostname after load")
} }
nsState := got.RecordsByNameserver["ns1.example.com."] nsState := got.RecordsByNameserver[testNS1]
if nsState.Status != "error" { if nsState.Status != statusError {
t.Errorf("status: got %q, want %q", nsState.Status, "error") t.Errorf("status: got %q, want %q", nsState.Status, statusError)
} }
if nsState.Error != "SERVFAIL" { if nsState.Error != "SERVFAIL" {
@@ -1062,7 +1073,8 @@ func TestConcurrentGetSet(t *testing.T) {
wg.Wait() wg.Wait()
} }
// runConcurrentOps performs a series of get/set/delete operations for concurrency testing. // runConcurrentOps performs a series of get/set/delete
// operations for concurrency testing.
func runConcurrentOps(s *state.State, key string, now time.Time) { func runConcurrentOps(s *state.State, key string, now time.Time) {
const iterations = 50 const iterations = 50
@@ -1085,7 +1097,7 @@ func runConcurrentOps(s *state.State, key string, now time.Time) {
s.SetHostnameState(key+".example.com", &state.HostnameState{ s.SetHostnameState(key+".example.com", &state.HostnameState{
RecordsByNameserver: map[string]*state.NameserverRecordState{ RecordsByNameserver: map[string]*state.NameserverRecordState{
"ns1.test.": { "ns1.test.": {
Records: map[string][]string{"A": {"1.2.3.4"}}, Records: map[string][]string{"A": {testIP}},
Status: "ok", Status: "ok",
LastChecked: now, LastChecked: now,
}, },

View File

@@ -26,6 +26,12 @@ const tlsPort = 443
// hoursPerDay converts days to hours for duration calculations. // hoursPerDay converts days to hours for duration calculations.
const hoursPerDay = 24 const hoursPerDay = 24
// Status values recorded for nameserver and certificate checks.
const (
statusOK = "ok"
statusError = "error"
)
// Params contains dependencies for Watcher. // Params contains dependencies for Watcher.
type Params struct { type Params struct {
fx.In fx.In
@@ -344,7 +350,7 @@ func buildHostnameState(
for ns, recs := range records { for ns, recs := range records {
hs.RecordsByNameserver[ns] = &state.NameserverRecordState{ hs.RecordsByNameserver[ns] = &state.NameserverRecordState{
Records: recs, Records: recs,
Status: "ok", Status: statusOK,
LastChecked: now, LastChecked: now,
} }
} }
@@ -402,7 +408,7 @@ func (w *Watcher) detectNSDisappearances(
current map[string]map[string][]string, current map[string]map[string][]string,
) { ) {
for ns, prevNS := range prev.RecordsByNameserver { for ns, prevNS := range prev.RecordsByNameserver {
if _, ok := current[ns]; ok || prevNS.Status != "ok" { if _, ok := current[ns]; ok || prevNS.Status != statusOK {
continue continue
} }
@@ -421,7 +427,7 @@ func (w *Watcher) detectNSDisappearances(
for ns := range current { for ns := range current {
prevNS, ok := prev.RecordsByNameserver[ns] prevNS, ok := prev.RecordsByNameserver[ns]
if !ok || prevNS.Status != "error" { if !ok || prevNS.Status != statusError {
continue continue
} }
@@ -705,7 +711,7 @@ func (w *Watcher) handleTLSError(
now time.Time, now time.Time,
err error, err error,
) { ) {
if hasPrev && !w.firstRun && prev.Status == "ok" { if hasPrev && !w.firstRun && prev.Status == statusOK {
msg := fmt.Sprintf( msg := fmt.Sprintf(
"Host: %s\nIP: %s\nError: %s", "Host: %s\nIP: %s\nError: %s",
hostname, ip, err, hostname, ip, err,
@@ -721,7 +727,7 @@ func (w *Watcher) handleTLSError(
w.state.SetCertificateState( w.state.SetCertificateState(
certKey, &state.CertificateState{ certKey, &state.CertificateState{
Status: "error", Status: statusError,
Error: err.Error(), Error: err.Error(),
LastChecked: now, LastChecked: now,
}, },
@@ -748,7 +754,7 @@ func (w *Watcher) handleTLSSuccess(
Issuer: cert.Issuer, Issuer: cert.Issuer,
NotAfter: cert.NotAfter, NotAfter: cert.NotAfter,
SubjectAlternativeNames: cert.SubjectAlternativeNames, SubjectAlternativeNames: cert.SubjectAlternativeNames,
Status: "ok", Status: statusOK,
LastChecked: now, LastChecked: now,
}, },
) )
@@ -760,7 +766,7 @@ func (w *Watcher) detectTLSChanges(
prev *state.CertificateState, prev *state.CertificateState,
cert *tlscheck.CertificateInfo, cert *tlscheck.CertificateInfo,
) { ) {
if prev.Status == "error" { if prev.Status == statusError {
msg := fmt.Sprintf( msg := fmt.Sprintf(
"Host: %s\nIP: %s\nTLS recovered", "Host: %s\nIP: %s\nTLS recovered",
hostname, ip, hostname, ip,

View File

@@ -18,6 +18,17 @@ import (
// errNotFound is returned when mock data is missing. // errNotFound is returned when mock data is missing.
var errNotFound = errors.New("not found") var errNotFound = errors.New("not found")
// Fixture values shared across tests.
const (
testDomain = "example.com"
testHost = "www.example.com"
testNS1 = "ns1.example.com."
testNS2 = "ns2.example.com."
testIPv4 = "93.184.216.34"
testIP = "1.2.3.4"
testIssuer = "DigiCert"
)
// --- Mock implementations --- // --- Mock implementations ---
type mockResolver struct { type mockResolver struct {
@@ -256,8 +267,8 @@ func TestFirstRunBaseline(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Domains = []string{"example.com"} cfg.Domains = []string{testDomain}
cfg.Hostnames = []string{"www.example.com"} cfg.Hostnames = []string{testHost}
w, deps := newTestWatcher(t, cfg) w, deps := newTestWatcher(t, cfg)
setupBaselineMocks(deps) setupBaselineMocks(deps)
@@ -269,37 +280,37 @@ func TestFirstRunBaseline(t *testing.T) {
} }
func setupBaselineMocks(deps *testDeps) { func setupBaselineMocks(deps *testDeps) {
deps.resolver.nsRecords["example.com"] = []string{ deps.resolver.nsRecords[testDomain] = []string{
"ns1.example.com.", testNS1,
"ns2.example.com.", testNS2,
} }
deps.resolver.allRecords["example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testDomain] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"93.184.216.34"}}, testNS1: {"A": {testIPv4}},
"ns2.example.com.": {"A": {"93.184.216.34"}}, testNS2: {"A": {testIPv4}},
} }
deps.resolver.allRecords["www.example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testHost] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"93.184.216.34"}}, testNS1: {"A": {testIPv4}},
"ns2.example.com.": {"A": {"93.184.216.34"}}, testNS2: {"A": {testIPv4}},
} }
deps.resolver.ipAddresses["www.example.com"] = []string{ deps.resolver.ipAddresses[testHost] = []string{
"93.184.216.34", testIPv4,
} }
deps.portChecker.results["93.184.216.34:80"] = true deps.portChecker.results["93.184.216.34:80"] = true
deps.portChecker.results["93.184.216.34:443"] = true deps.portChecker.results["93.184.216.34:443"] = true
deps.tlsChecker.certs["93.184.216.34:www.example.com"] = &tlscheck.CertificateInfo{ deps.tlsChecker.certs["93.184.216.34:www.example.com"] = &tlscheck.CertificateInfo{
CommonName: "www.example.com", CommonName: testHost,
Issuer: "DigiCert", Issuer: testIssuer,
NotAfter: time.Now().Add(90 * 24 * time.Hour), NotAfter: time.Now().Add(90 * 24 * time.Hour),
SubjectAlternativeNames: []string{ SubjectAlternativeNames: []string{
"www.example.com", testHost,
}, },
} }
deps.tlsChecker.certs["93.184.216.34:example.com"] = &tlscheck.CertificateInfo{ deps.tlsChecker.certs["93.184.216.34:example.com"] = &tlscheck.CertificateInfo{
CommonName: "example.com", CommonName: testDomain,
Issuer: "DigiCert", Issuer: testIssuer,
NotAfter: time.Now().Add(90 * 24 * time.Hour), NotAfter: time.Now().Add(90 * 24 * time.Hour),
SubjectAlternativeNames: []string{ SubjectAlternativeNames: []string{
"example.com", testDomain,
}, },
} }
} }
@@ -348,24 +359,24 @@ func TestDomainPortAndTLSChecks(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Domains = []string{"example.com"} cfg.Domains = []string{testDomain}
w, deps := newTestWatcher(t, cfg) w, deps := newTestWatcher(t, cfg)
deps.resolver.nsRecords["example.com"] = []string{ deps.resolver.nsRecords[testDomain] = []string{
"ns1.example.com.", testNS1,
} }
deps.resolver.allRecords["example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testDomain] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"93.184.216.34"}}, testNS1: {"A": {testIPv4}},
} }
deps.portChecker.results["93.184.216.34:80"] = true deps.portChecker.results["93.184.216.34:80"] = true
deps.portChecker.results["93.184.216.34:443"] = true deps.portChecker.results["93.184.216.34:443"] = true
deps.tlsChecker.certs["93.184.216.34:example.com"] = &tlscheck.CertificateInfo{ deps.tlsChecker.certs["93.184.216.34:example.com"] = &tlscheck.CertificateInfo{
CommonName: "example.com", CommonName: testDomain,
Issuer: "DigiCert", Issuer: testIssuer,
NotAfter: time.Now().Add(90 * 24 * time.Hour), NotAfter: time.Now().Add(90 * 24 * time.Hour),
SubjectAlternativeNames: []string{ SubjectAlternativeNames: []string{
"example.com", testDomain,
}, },
} }
@@ -406,17 +417,17 @@ func TestNSChangeDetection(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Domains = []string{"example.com"} cfg.Domains = []string{testDomain}
w, deps := newTestWatcher(t, cfg) w, deps := newTestWatcher(t, cfg)
deps.resolver.nsRecords["example.com"] = []string{ deps.resolver.nsRecords[testDomain] = []string{
"ns1.example.com.", testNS1,
"ns2.example.com.", testNS2,
} }
deps.resolver.allRecords["example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testDomain] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"1.2.3.4"}}, testNS1: {"A": {testIP}},
"ns2.example.com.": {"A": {"1.2.3.4"}}, testNS2: {"A": {testIP}},
} }
deps.portChecker.results["1.2.3.4:80"] = false deps.portChecker.results["1.2.3.4:80"] = false
deps.portChecker.results["1.2.3.4:443"] = false deps.portChecker.results["1.2.3.4:443"] = false
@@ -425,13 +436,13 @@ func TestNSChangeDetection(t *testing.T) {
w.RunOnce(ctx) w.RunOnce(ctx)
deps.resolver.mu.Lock() deps.resolver.mu.Lock()
deps.resolver.nsRecords["example.com"] = []string{ deps.resolver.nsRecords[testDomain] = []string{
"ns1.example.com.", testNS1,
"ns3.example.com.", "ns3.example.com.",
} }
deps.resolver.allRecords["example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testDomain] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"1.2.3.4"}}, testNS1: {"A": {testIP}},
"ns3.example.com.": {"A": {"1.2.3.4"}}, "ns3.example.com.": {"A": {testIP}},
} }
deps.resolver.mu.Unlock() deps.resolver.mu.Unlock()
@@ -459,15 +470,15 @@ func TestRecordChangeDetection(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Hostnames = []string{"www.example.com"} cfg.Hostnames = []string{testHost}
w, deps := newTestWatcher(t, cfg) w, deps := newTestWatcher(t, cfg)
deps.resolver.allRecords["www.example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testHost] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"93.184.216.34"}}, testNS1: {"A": {testIPv4}},
} }
deps.resolver.ipAddresses["www.example.com"] = []string{ deps.resolver.ipAddresses[testHost] = []string{
"93.184.216.34", testIPv4,
} }
deps.portChecker.results["93.184.216.34:80"] = false deps.portChecker.results["93.184.216.34:80"] = false
deps.portChecker.results["93.184.216.34:443"] = false deps.portChecker.results["93.184.216.34:443"] = false
@@ -476,10 +487,10 @@ func TestRecordChangeDetection(t *testing.T) {
w.RunOnce(ctx) w.RunOnce(ctx)
deps.resolver.mu.Lock() deps.resolver.mu.Lock()
deps.resolver.allRecords["www.example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testHost] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"93.184.216.35"}}, testNS1: {"A": {"93.184.216.35"}},
} }
deps.resolver.ipAddresses["www.example.com"] = []string{ deps.resolver.ipAddresses[testHost] = []string{
"93.184.216.35", "93.184.216.35",
} }
deps.resolver.mu.Unlock() deps.resolver.mu.Unlock()
@@ -501,24 +512,24 @@ func TestPortStateChange(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Hostnames = []string{"www.example.com"} cfg.Hostnames = []string{testHost}
w, deps := newTestWatcher(t, cfg) w, deps := newTestWatcher(t, cfg)
deps.resolver.allRecords["www.example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testHost] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"1.2.3.4"}}, testNS1: {"A": {testIP}},
} }
deps.resolver.ipAddresses["www.example.com"] = []string{ deps.resolver.ipAddresses[testHost] = []string{
"1.2.3.4", testIP,
} }
deps.portChecker.results["1.2.3.4:80"] = true deps.portChecker.results["1.2.3.4:80"] = true
deps.portChecker.results["1.2.3.4:443"] = true deps.portChecker.results["1.2.3.4:443"] = true
deps.tlsChecker.certs["1.2.3.4:www.example.com"] = &tlscheck.CertificateInfo{ deps.tlsChecker.certs["1.2.3.4:www.example.com"] = &tlscheck.CertificateInfo{
CommonName: "www.example.com", CommonName: testHost,
Issuer: "DigiCert", Issuer: testIssuer,
NotAfter: time.Now().Add(90 * 24 * time.Hour), NotAfter: time.Now().Add(90 * 24 * time.Hour),
SubjectAlternativeNames: []string{ SubjectAlternativeNames: []string{
"www.example.com", testHost,
}, },
} }
@@ -541,24 +552,24 @@ func TestTLSExpiryWarning(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Hostnames = []string{"www.example.com"} cfg.Hostnames = []string{testHost}
w, deps := newTestWatcher(t, cfg) w, deps := newTestWatcher(t, cfg)
deps.resolver.allRecords["www.example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testHost] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"1.2.3.4"}}, testNS1: {"A": {testIP}},
} }
deps.resolver.ipAddresses["www.example.com"] = []string{ deps.resolver.ipAddresses[testHost] = []string{
"1.2.3.4", testIP,
} }
deps.portChecker.results["1.2.3.4:80"] = true deps.portChecker.results["1.2.3.4:80"] = true
deps.portChecker.results["1.2.3.4:443"] = true deps.portChecker.results["1.2.3.4:443"] = true
deps.tlsChecker.certs["1.2.3.4:www.example.com"] = &tlscheck.CertificateInfo{ deps.tlsChecker.certs["1.2.3.4:www.example.com"] = &tlscheck.CertificateInfo{
CommonName: "www.example.com", CommonName: testHost,
Issuer: "DigiCert", Issuer: testIssuer,
NotAfter: time.Now().Add(3 * 24 * time.Hour), NotAfter: time.Now().Add(3 * 24 * time.Hour),
SubjectAlternativeNames: []string{ SubjectAlternativeNames: []string{
"www.example.com", testHost,
}, },
} }
@@ -592,25 +603,25 @@ func TestTLSExpiryWarningDedup(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Hostnames = []string{"www.example.com"} cfg.Hostnames = []string{testHost}
cfg.TLSInterval = 24 * time.Hour cfg.TLSInterval = 24 * time.Hour
w, deps := newTestWatcher(t, cfg) w, deps := newTestWatcher(t, cfg)
deps.resolver.allRecords["www.example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testHost] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"1.2.3.4"}}, testNS1: {"A": {testIP}},
} }
deps.resolver.ipAddresses["www.example.com"] = []string{ deps.resolver.ipAddresses[testHost] = []string{
"1.2.3.4", testIP,
} }
deps.portChecker.results["1.2.3.4:80"] = true deps.portChecker.results["1.2.3.4:80"] = true
deps.portChecker.results["1.2.3.4:443"] = true deps.portChecker.results["1.2.3.4:443"] = true
deps.tlsChecker.certs["1.2.3.4:www.example.com"] = &tlscheck.CertificateInfo{ deps.tlsChecker.certs["1.2.3.4:www.example.com"] = &tlscheck.CertificateInfo{
CommonName: "www.example.com", CommonName: testHost,
Issuer: "DigiCert", Issuer: testIssuer,
NotAfter: time.Now().Add(3 * 24 * time.Hour), NotAfter: time.Now().Add(3 * 24 * time.Hour),
SubjectAlternativeNames: []string{ SubjectAlternativeNames: []string{
"www.example.com", testHost,
}, },
} }
@@ -647,17 +658,17 @@ func TestGracefulShutdown(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Domains = []string{"example.com"} cfg.Domains = []string{testDomain}
cfg.DNSInterval = 100 * time.Millisecond cfg.DNSInterval = 100 * time.Millisecond
cfg.TLSInterval = 100 * time.Millisecond cfg.TLSInterval = 100 * time.Millisecond
w, deps := newTestWatcher(t, cfg) w, deps := newTestWatcher(t, cfg)
deps.resolver.nsRecords["example.com"] = []string{ deps.resolver.nsRecords[testDomain] = []string{
"ns1.example.com.", testNS1,
} }
deps.resolver.allRecords["example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testDomain] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"1.2.3.4"}}, testNS1: {"A": {testIP}},
} }
deps.portChecker.results["1.2.3.4:80"] = false deps.portChecker.results["1.2.3.4:80"] = false
deps.portChecker.results["1.2.3.4:443"] = false deps.portChecker.results["1.2.3.4:443"] = false
@@ -687,13 +698,13 @@ func setupHostnameIP(
hostname, ip string, hostname, ip string,
) { ) {
deps.resolver.allRecords[hostname] = map[string]map[string][]string{ deps.resolver.allRecords[hostname] = map[string]map[string][]string{
"ns1.example.com.": {"A": {ip}}, testNS1: {"A": {ip}},
} }
deps.portChecker.results[ip+":80"] = true deps.portChecker.results[ip+":80"] = true
deps.portChecker.results[ip+":443"] = true deps.portChecker.results[ip+":443"] = true
deps.tlsChecker.certs[ip+":"+hostname] = &tlscheck.CertificateInfo{ deps.tlsChecker.certs[ip+":"+hostname] = &tlscheck.CertificateInfo{
CommonName: hostname, CommonName: hostname,
Issuer: "DigiCert", Issuer: testIssuer,
NotAfter: time.Now().Add(90 * 24 * time.Hour), NotAfter: time.Now().Add(90 * 24 * time.Hour),
SubjectAlternativeNames: []string{hostname}, SubjectAlternativeNames: []string{hostname},
} }
@@ -702,7 +713,7 @@ func setupHostnameIP(
func updateHostnameIP(deps *testDeps, hostname, ip string) { func updateHostnameIP(deps *testDeps, hostname, ip string) {
deps.resolver.mu.Lock() deps.resolver.mu.Lock()
deps.resolver.allRecords[hostname] = map[string]map[string][]string{ deps.resolver.allRecords[hostname] = map[string]map[string][]string{
"ns1.example.com.": {"A": {ip}}, testNS1: {"A": {ip}},
} }
deps.resolver.mu.Unlock() deps.resolver.mu.Unlock()
@@ -714,7 +725,7 @@ func updateHostnameIP(deps *testDeps, hostname, ip string) {
deps.tlsChecker.mu.Lock() deps.tlsChecker.mu.Lock()
deps.tlsChecker.certs[ip+":"+hostname] = &tlscheck.CertificateInfo{ deps.tlsChecker.certs[ip+":"+hostname] = &tlscheck.CertificateInfo{
CommonName: hostname, CommonName: hostname,
Issuer: "DigiCert", Issuer: testIssuer,
NotAfter: time.Now().Add(90 * 24 * time.Hour), NotAfter: time.Now().Add(90 * 24 * time.Hour),
SubjectAlternativeNames: []string{hostname}, SubjectAlternativeNames: []string{hostname},
} }
@@ -725,11 +736,11 @@ func TestDNSRunsBeforePortAndTLSChecks(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Hostnames = []string{"www.example.com"} cfg.Hostnames = []string{testHost}
w, deps := newTestWatcher(t, cfg) w, deps := newTestWatcher(t, cfg)
setupHostnameIP(deps, "www.example.com", "10.0.0.1") setupHostnameIP(deps, testHost, "10.0.0.1")
ctx := t.Context() ctx := t.Context()
w.RunOnce(ctx) w.RunOnce(ctx)
@@ -740,7 +751,7 @@ func TestDNSRunsBeforePortAndTLSChecks(t *testing.T) {
} }
// DNS changes to a new IP; port and TLS must pick it up. // DNS changes to a new IP; port and TLS must pick it up.
updateHostnameIP(deps, "www.example.com", "10.0.0.2") updateHostnameIP(deps, testHost, "10.0.0.2")
w.RunOnce(ctx) w.RunOnce(ctx)
@@ -760,8 +771,8 @@ func TestSendTestNotification_Enabled(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Domains = []string{"example.com"} cfg.Domains = []string{testDomain}
cfg.Hostnames = []string{"www.example.com"} cfg.Hostnames = []string{testHost}
cfg.SendTestNotification = true cfg.SendTestNotification = true
w, deps := newTestWatcher(t, cfg) w, deps := newTestWatcher(t, cfg)
@@ -786,8 +797,8 @@ func TestSendTestNotification_ViaRun(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Domains = []string{"example.com"} cfg.Domains = []string{testDomain}
cfg.Hostnames = []string{"www.example.com"} cfg.Hostnames = []string{testHost}
cfg.SendTestNotification = true cfg.SendTestNotification = true
cfg.DNSInterval = 24 * time.Hour cfg.DNSInterval = 24 * time.Hour
cfg.TLSInterval = 24 * time.Hour cfg.TLSInterval = 24 * time.Hour
@@ -833,8 +844,8 @@ func TestSendTestNotification_Disabled(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Domains = []string{"example.com"} cfg.Domains = []string{testDomain}
cfg.Hostnames = []string{"www.example.com"} cfg.Hostnames = []string{testHost}
cfg.SendTestNotification = false cfg.SendTestNotification = false
cfg.DNSInterval = 24 * time.Hour cfg.DNSInterval = 24 * time.Hour
cfg.TLSInterval = 24 * time.Hour cfg.TLSInterval = 24 * time.Hour
@@ -871,16 +882,16 @@ func TestNSFailureAndRecovery(t *testing.T) {
t.Parallel() t.Parallel()
cfg := defaultTestConfig(t) cfg := defaultTestConfig(t)
cfg.Hostnames = []string{"www.example.com"} cfg.Hostnames = []string{testHost}
w, deps := newTestWatcher(t, cfg) w, deps := newTestWatcher(t, cfg)
deps.resolver.allRecords["www.example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testHost] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"1.2.3.4"}}, testNS1: {"A": {testIP}},
"ns2.example.com.": {"A": {"1.2.3.4"}}, testNS2: {"A": {testIP}},
} }
deps.resolver.ipAddresses["www.example.com"] = []string{ deps.resolver.ipAddresses[testHost] = []string{
"1.2.3.4", testIP,
} }
deps.portChecker.results["1.2.3.4:80"] = false deps.portChecker.results["1.2.3.4:80"] = false
deps.portChecker.results["1.2.3.4:443"] = false deps.portChecker.results["1.2.3.4:443"] = false
@@ -890,8 +901,8 @@ func TestNSFailureAndRecovery(t *testing.T) {
w.RunOnce(ctx) w.RunOnce(ctx)
deps.resolver.mu.Lock() deps.resolver.mu.Lock()
deps.resolver.allRecords["www.example.com"] = map[string]map[string][]string{ deps.resolver.allRecords[testHost] = map[string]map[string][]string{
"ns1.example.com.": {"A": {"1.2.3.4"}}, testNS1: {"A": {testIP}},
} }
deps.resolver.mu.Unlock() deps.resolver.mu.Unlock()

View File

@@ -9,9 +9,9 @@ set -eu
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)" ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
# Pinned versions, 2026-07-07 (same pins as the Dockerfile) # Pinned versions, 2026-08-07 (same pins as the Dockerfile)
# golangci-lint v2.10.1 # golangci-lint v2.12.2
GOLANGCI_LINT_REF="github.com/golangci/golangci-lint/v2/cmd/golangci-lint@5d1e709b7be35cb2025444e19de266b056b7b7ee" GOLANGCI_LINT_REF="github.com/golangci/golangci-lint/v2/cmd/golangci-lint@c0d3ddc9cf3faa61a4e378e879ece580256d76e5"
# goimports v0.42.0 # goimports v0.42.0
GOIMPORTS_REF="golang.org/x/tools/cmd/goimports@009367f5c17a8d4c45a961a3a509277190a9a6f0" GOIMPORTS_REF="golang.org/x/tools/cmd/goimports@009367f5c17a8d4c45a961a3a509277190a9a6f0"