1 Commits

Author SHA1 Message Date
ff66ecc0c9 ci: re-run make check on every cibuild instead of serving it from the layer cache (closes #115)
All checks were successful
check / check (push) Successful in 43s
`script/cibuild` was plain `docker build .`. The Dockerfile does
`COPY . .` and then `RUN make check`, and Docker invalidates `COPY . .`
only on a content change, so on a byte-identical tree the check layer
was reused and the suite never ran. The script's header comment claimed
that a successful build implies all checks pass, which was false
whenever the cache was warm. Reproduced on this branch's parent: a
second consecutive run returned success in 283 ms with
`#13 [builder 9/10] RUN make check` reported `CACHED`.

That matters more here than in a typical repo. DNS is never mocked in
this repository, so the suite queries live DNS and its outcome varies
with real-world conditions; caching the verdict of a non-deterministic
check replays a stale result in exactly the case where re-running is
most valuable. It is also the gate every PR is verified through.

Fix: declare `ARG CHECK_EPOCH` immediately above the check step and
expand it into the command, with `script/cibuild` passing a fresh
`$(date +%s%N)` per invocation. A build argument's value participates in
the cache key of later instructions in the stage even when they do not
reference it, so a fresh value busts this layer either way; the value is
expanded into the command deliberately, which makes the invalidation a
property of the command string itself rather than of how a given builder
treats unreferenced args, and surfaces the epoch in the build log as a
diagnostic. Placing the ARG here and no earlier keeps the pinned
toolchain installs and `go mod download` above the invalidation line, so
only the check and the steps after it re-run. The epoch is nanosecond
granular so that two concurrent invocations starting in the same second
cannot share a value.

A plain `docker build` without the argument caches as before; nothing
outside the CI entrypoint changes behaviour.

Verified by experiment, not inspection:

- Two consecutive runs on an unchanged tree: 55.2 s and 42.2 s, both
  exit 0, with distinct epochs. The second run shows
  `RUN echo "check epoch: ..." && make check` executing for 36.0 s and
  216 passing tests across all eight packages, while `apk add`, both
  pinned `go install` steps, `go mod download`, `COPY go.mod go.sum` and
  `COPY . .` all report `CACHED`.
- Negative control: planted `internal/config/zz_negative_control_test.go`
  calling `t.Fatal("NEGATIVE-CONTROL-115: planted failure, cache did not
  serve this layer")`. The build failed in 24.7 s with exit 1, printing
  that exact message and `--- FAIL: TestNegativeControlIssue115`, and the
  check step exited with code 2. A cached layer cannot produce a failure
  predicted in advance, so this establishes the suite ran. The file was
  then removed, `git status` confirmed clean, and the tree built green
  again in 48.1 s.
- Total build time 42-55 s against the policy's 5-minute ceiling.
- `make check` green. No pin touched: the `golang` and `alpine` sha256
  digests, golangci-lint `c0d3ddc9`, and goimports `009367f5` are
  unchanged, and `.golangci.yml` still hashes to `021cc83f4e6f...`.
2026-08-09 06:22:28 +00:00
9 changed files with 113 additions and 786 deletions

View File

@@ -15,8 +15,25 @@ RUN go mod download
COPY . . COPY . .
# Run all checks - build fails if any check fails # Run all checks - build fails if any check fails.
RUN make check #
# CHECK_EPOCH is a cache-busting build argument. Without it, an
# unchanged tree leaves this layer's cache key identical and Docker
# serves the previous verdict instead of re-running the suite, so the
# build reports a green it did not earn. A build argument's value
# participates in the cache key of later instructions in the stage even
# when they do not reference it, so a fresh value busts this layer
# either way. It is expanded into the command deliberately: that makes
# the invalidation a property of the command string itself rather than
# of how a given builder treats unreferenced args, and it surfaces the
# epoch in the build log as a diagnostic.
#
# Placing the ARG here and nowhere earlier keeps everything above it
# (toolchain install, go mod download) cached, so only the check and the
# steps after it re-run. script/cibuild passes a fresh value per run; a
# plain `docker build` without it caches as before.
ARG CHECK_EPOCH
RUN echo "check epoch: ${CHECK_EPOCH}" && make check
# Build the binary # Build the binary
RUN make build RUN make build

View File

@@ -218,8 +218,7 @@ 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. In-flight notification deliveries cancellation and the fx lifecycle.
are drained on shutdown, bounded by the shutdown timeout.
--- ---
@@ -394,7 +393,10 @@ them. We provide:
- `script/check` — run test, lint, and fmt-check - `script/check` — run test, lint, and fmt-check
- `script/docker` — build the Docker image tagged via - `script/docker` — build the Docker image tagged via
`script/projectname` `script/projectname`
- `script/cibuild` — CI entrypoint: plain `docker build .` - `script/cibuild` — CI entrypoint: `docker build .` with a fresh
`CHECK_EPOCH` build argument, so the Dockerfile's `make check` layer
is never served from the cache and a green build always means the
checks ran on this invocation
- `script/precommit` — run by the git pre-commit hook; `go mod tidy` - `script/precommit` — run by the git pre-commit hook; `go mod tidy`
guard, then `script/check` guard, then `script/check`
- `script/install-precommit` — install the git pre-commit hook - `script/install-precommit` — install the git pre-commit hook
@@ -453,14 +455,8 @@ 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, wait for in-flight 5. **Shutdown**: Persist final state to disk, complete in-flight
notification deliveries to complete, stop gracefully. The wait is notifications, stop gracefully.
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.
--- ---

19
TODO.md
View File

@@ -25,15 +25,16 @@ confirm make check still passes.
# Completed Steps # Completed Steps
- 2026-08-09: in-flight notification deliveries are now drained at - 2026-08-09: `script/cibuild` can no longer report a green it did not
shutdown (#106): `notify.New` registers an fx `OnStop` hook that waits earn. The Dockerfile declares `ARG CHECK_EPOCH` immediately above the
on a `sync.WaitGroup` of tracked delivery goroutines, bounded by the check step and expands it into the `RUN` command, and `script/cibuild`
`OnStop` context; on expiry the outstanding count is logged at warn passes a fresh `$(date +%s%N)` per invocation, so the `make check`
level and parked retry backoffs are released instead of being dropped layer is always re-executed while the pinned toolchain install and
silently, and deliveries submitted after the drain begins are refused `go mod download` stay cached. Verified by experiment: before the fix
so shutdown cannot be extended indefinitely; an `OnStop` context that a second run on an unchanged tree returned in 283 ms with the check
is already expired on entry with nothing outstanding drains quietly layer `CACHED`; after it the check runs every time, and a deliberately
rather than warning about deliveries that were never abandoned planted always-failing test made the build fail with exactly that
test's message
- 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

View File

@@ -32,27 +32,11 @@ 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 newService(slog.New(slog.DiscardHandler), transport) return &Service{
} log: slog.New(slog.DiscardHandler),
transport: transport,
// NewTestServiceWithLogger creates a Service that writes to the history: NewAlertHistory(),
// 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

@@ -12,8 +12,6 @@ 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"
@@ -117,41 +115,19 @@ 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(
lifecycle fx.Lifecycle, _ fx.Lifecycle,
params Params, params Params,
) (*Service, error) { ) (*Service, error) {
svc := newService(params.Logger.Get(), http.DefaultTransport) svc := &Service{
svc.config = params.Config log: params.Logger.Get(),
transport: http.DefaultTransport,
config: params.Config,
history: NewAlertHistory(),
}
if params.Config.NtfyTopic != "" { if params.Config.NtfyTopic != "" {
u, err := ValidateWebhookURL( u, err := ValidateWebhookURL(
@@ -192,14 +168,6 @@ 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
} }
@@ -226,32 +194,6 @@ 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,
@@ -260,11 +202,26 @@ func (svc *Service) dispatchNtfy(
return return
} }
svc.dispatch(ctx, "ntfy", func(c context.Context) error { go func() {
return svc.sendNtfy( notifyCtx := context.WithoutCancel(ctx)
c, svc.ntfyURL, title, message, priority,
err := svc.deliverWithRetry(
notifyCtx, "ntfy",
func(c context.Context) error {
return svc.sendNtfy(
c, svc.ntfyURL,
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(
@@ -275,11 +232,26 @@ func (svc *Service) dispatchSlack(
return return
} }
svc.dispatch(ctx, "slack", func(c context.Context) error { go func() {
return svc.sendSlack( notifyCtx := context.WithoutCancel(ctx)
c, svc.slackWebhookURL, title, message, priority,
err := svc.deliverWithRetry(
notifyCtx, "slack",
func(c context.Context) error {
return svc.sendSlack(
c, svc.slackWebhookURL,
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(
@@ -290,15 +262,26 @@ func (svc *Service) dispatchMattermost(
return return
} }
svc.dispatch( go func() {
ctx, "mattermost", notifyCtx := context.WithoutCancel(ctx)
func(c context.Context) error {
return svc.sendSlack( err := svc.deliverWithRetry(
c, svc.mattermostWebhookURL, notifyCtx, "mattermost",
title, message, priority, func(c context.Context) error {
return svc.sendSlack(
c, svc.mattermostWebhookURL,
title, message, priority,
)
},
)
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,7 +2,6 @@ package notify
import ( import (
"context" "context"
"fmt"
"math" "math"
"math/rand/v2" "math/rand/v2"
"time" "time"
@@ -122,14 +121,6 @@ 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):
} }
} }

View File

@@ -1,119 +0,0 @@
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

@@ -1,531 +0,0 @@
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

@@ -1,13 +1,18 @@
#!/bin/sh #!/bin/sh
# script/cibuild: run the CI build. The Dockerfile runs make check, so # script/cibuild: run the CI build. The Dockerfile runs make check, and
# a successful build implies all checks pass. # the CHECK_EPOCH build argument below is fresh on every invocation, so
# the check layer is never served from the Docker layer cache: a
# successful build means the checks were executed and passed on this
# run, not on some earlier one. Only the check step and the steps after
# it are invalidated; the toolchain install and go mod download stay
# cached.
set -eu set -eu
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)" ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
main() { main() {
cd "$ROOT" cd "$ROOT"
docker build . docker build --build-arg CHECK_EPOCH="$(date +%s%N)" .
} }
main "$@" main "$@"