mirror of
https://github.com/Control-D-Inc/ctrld.git
synced 2026-07-29 01:18:48 +02:00
fix: back off unroutable IPv6 DoH upstream health-check spam
When IPv6 is available locally but the selected IPv6 DoH endpoint is unroutable (e.g. dialing [2606:1a40::22]:443 returns "no route to host" while IPv4 stays usable), ctrld re-bootstrapped and re-dialed the endpoint every ~2s. A weekend soak produced ~46.7k "no route to host" lines, with the dial/health-check loop dominating the log during bad windows. Add bounded backoff/suppression for network-unreachable endpoints at two levels: - ParallelDialer (internal/net): track dial addresses that fail with ENETUNREACH/EHOSTUNREACH and skip them for an exponentially growing, bounded window (5s -> 60s). A successful dial clears the entry immediately, so recovery is preserved when the route returns. When every candidate is suppressed the dial fails fast and quietly instead of hammering known-unroutable addresses. - Upstream recovery loop (cmd/cli): demote unreachable check failures to debug and back off the retry cadence (2s -> 60s) for an unreachable streak; any other failure resets to the base cadence. The new IsUnreachable classifier lives in internal/net and is reused by cmd/cli's errNetworkError, so the unreachable-errno matching has a single definition. Note the explicit winsock constants (10051/10065) are required on Windows: syscall.ENETUNREACH/EHOSTUNREACH are Go's portable "invented" values and never equal the raw WSA codes a failing connect surfaces. Suppression and backoff are always bounded, so IPv6 is never disabled until restart and recovers on its own once the route is back. Split-stack selection and the #549 macOS intercept recovery work are untouched. Adds unit tests for the classifier, the dialer's suppression tracker, and the recovery backoff schedule.
This commit is contained in:
@@ -1417,15 +1417,6 @@ func (p *prog) ensurePFAnchorActive() bool {
|
|||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
func (p *prog) scheduleDNSAfterVPNSettleRefresh(reason string, delay time.Duration) {
|
|
||||||
time.AfterFunc(delay, func() {
|
|
||||||
if p.dnsInterceptState == nil {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
p.refreshDNSAfterVPNSettle(reason)
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
func (p *prog) pfExecBackoffActive() bool {
|
func (p *prog) pfExecBackoffActive() bool {
|
||||||
until := p.pfExecBackoffUntil.Load()
|
until := p.pfExecBackoffUntil.Load()
|
||||||
if until == 0 {
|
if until == 0 {
|
||||||
@@ -1436,7 +1427,7 @@ func (p *prog) pfExecBackoffActive() bool {
|
|||||||
p.pfExecBackoffUntil.CompareAndSwap(until, 0)
|
p.pfExecBackoffUntil.CompareAndSwap(until, 0)
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
mainLog.Load().Debug().Dur("remaining", remaining).Msg("DNS intercept watchdog: suppressed during pf exec backoff")
|
mainLog.Load().Debug().Msgf("DNS intercept watchdog: suppressed during pf exec backoff (remaining: %s)", remaining)
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1446,8 +1437,7 @@ func (p *prog) pfBackoffResourceExhaustion(err error, output []byte, operation s
|
|||||||
}
|
}
|
||||||
until := time.Now().Add(pfExecFailureBackoff)
|
until := time.Now().Add(pfExecFailureBackoff)
|
||||||
p.pfExecBackoffUntil.Store(until.UnixMilli())
|
p.pfExecBackoffUntil.Store(until.UnixMilli())
|
||||||
mainLog.Load().Warn().Err(err).Dur("backoff", pfExecFailureBackoff).Str("operation", operation).
|
mainLog.Load().Warn().Err(err).Msgf("DNS intercept watchdog: backing off after local exec resource exhaustion (operation: %s, backoff: %s)", operation, pfExecFailureBackoff)
|
||||||
Msg("DNS intercept watchdog: backing off after local exec resource exhaustion")
|
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+26
-2
@@ -1945,6 +1945,15 @@ func (p *prog) checkUpstreamOnce(upstream string, uc *ctrld.UpstreamConfig) erro
|
|||||||
p.Debug().Err(err).Msgf("Upstream %s check failed after %v (WFP loopback protect active)", upstream, duration)
|
p.Debug().Err(err).Msgf("Upstream %s check failed after %v (WFP loopback protect active)", upstream, duration)
|
||||||
return errOsHealthcheckSuppressed
|
return errOsHealthcheckSuppressed
|
||||||
}
|
}
|
||||||
|
// A no-route/network-unreachable failure means the endpoint's address
|
||||||
|
// family is available locally but unroutable (e.g. an IPv6 DoH endpoint
|
||||||
|
// while IPv6 is up but has no route). These repeat until the route
|
||||||
|
// returns and are handled by bounded backoff in the recovery loop, so
|
||||||
|
// keep them at debug to avoid sustained error-log spam.
|
||||||
|
if ctrldnet.IsUnreachable(err) {
|
||||||
|
p.Debug().Err(err).Msgf("Upstream %s check failed after %v (network unreachable)", upstream, duration)
|
||||||
|
return err
|
||||||
|
}
|
||||||
p.Error().Err(err).Msgf("Upstream %s check failed after %v", upstream, duration)
|
p.Error().Err(err).Msgf("Upstream %s check failed after %v", upstream, duration)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -2223,6 +2232,7 @@ func (p *prog) waitForUpstreamRecovery(ctx context.Context, upstreams map[string
|
|||||||
defer wg.Done()
|
defer wg.Done()
|
||||||
p.Debug().Msgf("Starting recovery check loop for upstream: %s", name)
|
p.Debug().Msgf("Starting recovery check loop for upstream: %s", name)
|
||||||
attempts := 0
|
attempts := 0
|
||||||
|
unreachableStreak := 0
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-recoveryCtx.Done():
|
case <-recoveryCtx.Done():
|
||||||
@@ -2243,8 +2253,22 @@ func (p *prog) waitForUpstreamRecovery(ctx context.Context, upstreams map[string
|
|||||||
}
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
p.Debug().Msgf("Upstream %s check failed, sleeping before retry", name)
|
// Back off the retry cadence for an unroutable endpoint so a
|
||||||
if !sleepWithContext(recoveryCtx, checkUpstreamBackoffSleep) {
|
// host with IPv6 up but no route to the IPv6 DoH endpoint does
|
||||||
|
// not re-bootstrap/re-check every checkUpstreamBackoffSleep and
|
||||||
|
// spam the log. The backoff is bounded (checkUpstreamUnreachableBackoffMax)
|
||||||
|
// so the endpoint is still re-probed and recovers when the route
|
||||||
|
// returns; any other failure resets to the base cadence.
|
||||||
|
sleep := checkUpstreamBackoffSleep
|
||||||
|
if ctrldnet.IsUnreachable(err) {
|
||||||
|
unreachableStreak++
|
||||||
|
sleep = unreachableRecoveryBackoff(unreachableStreak)
|
||||||
|
p.Debug().Msgf("Upstream %s unreachable (streak %d), backing off %s before retry", name, unreachableStreak, sleep)
|
||||||
|
} else {
|
||||||
|
unreachableStreak = 0
|
||||||
|
p.Debug().Msgf("Upstream %s check failed, sleeping before retry", name)
|
||||||
|
}
|
||||||
|
if !sleepWithContext(recoveryCtx, sleep) {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+8
-7
@@ -35,6 +35,7 @@ import (
|
|||||||
"github.com/Control-D-Inc/ctrld/internal/controld"
|
"github.com/Control-D-Inc/ctrld/internal/controld"
|
||||||
"github.com/Control-D-Inc/ctrld/internal/dnscache"
|
"github.com/Control-D-Inc/ctrld/internal/dnscache"
|
||||||
"github.com/Control-D-Inc/ctrld/internal/firewall"
|
"github.com/Control-D-Inc/ctrld/internal/firewall"
|
||||||
|
ctrldnet "github.com/Control-D-Inc/ctrld/internal/net"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -1311,13 +1312,14 @@ func errAddrInUse(err error) bool {
|
|||||||
|
|
||||||
var _ = errAddrInUse
|
var _ = errAddrInUse
|
||||||
|
|
||||||
|
// The unreachable winsock errnos (ENETUNREACH/EHOSTUNREACH) are matched via
|
||||||
|
// ctrldnet.IsUnreachable, which owns their definitions.
|
||||||
|
//
|
||||||
// https://learn.microsoft.com/en-us/windows/win32/winsock/windows-sockets-error-codes-2
|
// https://learn.microsoft.com/en-us/windows/win32/winsock/windows-sockets-error-codes-2
|
||||||
var (
|
var (
|
||||||
windowsECONNREFUSED = syscall.Errno(10061)
|
windowsECONNREFUSED = syscall.Errno(10061)
|
||||||
windowsENETUNREACH = syscall.Errno(10051)
|
|
||||||
windowsEINVAL = syscall.Errno(10022)
|
windowsEINVAL = syscall.Errno(10022)
|
||||||
windowsEADDRINUSE = syscall.Errno(10048)
|
windowsEADDRINUSE = syscall.Errno(10048)
|
||||||
windowsEHOSTUNREACH = syscall.Errno(10065)
|
|
||||||
)
|
)
|
||||||
|
|
||||||
func errUrlNetworkError(err error) bool {
|
func errUrlNetworkError(err error) bool {
|
||||||
@@ -1334,15 +1336,14 @@ func errNetworkError(err error) bool {
|
|||||||
if opErr.Temporary() {
|
if opErr.Temporary() {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
if ctrldnet.IsUnreachable(err) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
switch {
|
switch {
|
||||||
case errors.Is(opErr.Err, syscall.ECONNREFUSED),
|
case errors.Is(opErr.Err, syscall.ECONNREFUSED),
|
||||||
errors.Is(opErr.Err, syscall.EINVAL),
|
errors.Is(opErr.Err, syscall.EINVAL),
|
||||||
errors.Is(opErr.Err, syscall.ENETUNREACH),
|
|
||||||
errors.Is(opErr.Err, syscall.EHOSTUNREACH),
|
|
||||||
errors.Is(opErr.Err, windowsENETUNREACH),
|
|
||||||
errors.Is(opErr.Err, windowsEINVAL),
|
errors.Is(opErr.Err, windowsEINVAL),
|
||||||
errors.Is(opErr.Err, windowsECONNREFUSED),
|
errors.Is(opErr.Err, windowsECONNREFUSED):
|
||||||
errors.Is(opErr.Err, windowsEHOSTUNREACH):
|
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -32,6 +32,15 @@ func TestSleepWithContext(t *testing.T) {
|
|||||||
assert.Less(t, time.Since(start), 100*time.Millisecond)
|
assert.Less(t, time.Since(start), 100*time.Millisecond)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestUnreachableRecoveryBackoff(t *testing.T) {
|
||||||
|
// Streak starts at the base cadence and doubles each attempt, capped at the max.
|
||||||
|
assert.Equal(t, checkUpstreamBackoffSleep, unreachableRecoveryBackoff(0))
|
||||||
|
assert.Equal(t, checkUpstreamBackoffSleep, unreachableRecoveryBackoff(1))
|
||||||
|
assert.Equal(t, 2*checkUpstreamBackoffSleep, unreachableRecoveryBackoff(2))
|
||||||
|
assert.Equal(t, 4*checkUpstreamBackoffSleep, unreachableRecoveryBackoff(3))
|
||||||
|
assert.Equal(t, checkUpstreamUnreachableBackoffMax, unreachableRecoveryBackoff(100))
|
||||||
|
}
|
||||||
|
|
||||||
func Test_prog_dnsWatchdogEnabled(t *testing.T) {
|
func Test_prog_dnsWatchdogEnabled(t *testing.T) {
|
||||||
p := &prog{cfg: &ctrld.Config{}}
|
p := &prog{cfg: &ctrld.Config{}}
|
||||||
|
|
||||||
|
|||||||
@@ -13,8 +13,27 @@ const (
|
|||||||
maxFailureRequest = 50
|
maxFailureRequest = 50
|
||||||
// checkUpstreamBackoffSleep is the time interval between each upstream checks.
|
// checkUpstreamBackoffSleep is the time interval between each upstream checks.
|
||||||
checkUpstreamBackoffSleep = 2 * time.Second
|
checkUpstreamBackoffSleep = 2 * time.Second
|
||||||
|
// checkUpstreamUnreachableBackoffMax caps the recovery retry interval for an
|
||||||
|
// endpoint that keeps failing with a network-unreachable error. It bounds
|
||||||
|
// the backoff so an unroutable endpoint is still re-probed periodically and
|
||||||
|
// recovers once the route returns.
|
||||||
|
checkUpstreamUnreachableBackoffMax = 60 * time.Second
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// unreachableRecoveryBackoff returns the retry interval for the given streak of
|
||||||
|
// consecutive network-unreachable failures. It starts at checkUpstreamBackoffSleep
|
||||||
|
// and doubles each attempt, capped at checkUpstreamUnreachableBackoffMax.
|
||||||
|
func unreachableRecoveryBackoff(streak int) time.Duration {
|
||||||
|
d := checkUpstreamBackoffSleep
|
||||||
|
for i := 1; i < streak; i++ {
|
||||||
|
d *= 2
|
||||||
|
if d >= checkUpstreamUnreachableBackoffMax {
|
||||||
|
return checkUpstreamUnreachableBackoffMax
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return d
|
||||||
|
}
|
||||||
|
|
||||||
// upstreamMonitor performs monitoring upstreams health.
|
// upstreamMonitor performs monitoring upstreams health.
|
||||||
type upstreamMonitor struct {
|
type upstreamMonitor struct {
|
||||||
cfg *ctrld.Config
|
cfg *ctrld.Config
|
||||||
|
|||||||
@@ -17,7 +17,7 @@ func withVPNDNSSettlingEnabled(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func TestVPNDNSRefreshSkipsConcurrentDuplicate(t *testing.T) {
|
func TestVPNDNSRefreshSkipsConcurrentDuplicate(t *testing.T) {
|
||||||
m := newVPNDNSManager(nil)
|
m := newVPNDNSManager(&mainLog, nil)
|
||||||
started := make(chan struct{})
|
started := make(chan struct{})
|
||||||
release := make(chan struct{})
|
release := make(chan struct{})
|
||||||
done := make(chan struct{})
|
done := make(chan struct{})
|
||||||
@@ -33,11 +33,11 @@ func TestVPNDNSRefreshSkipsConcurrentDuplicate(t *testing.T) {
|
|||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
defer close(done)
|
defer close(done)
|
||||||
m.Refresh(true)
|
m.Refresh(context.Background(), true)
|
||||||
}()
|
}()
|
||||||
|
|
||||||
<-started
|
<-started
|
||||||
m.Refresh(true)
|
m.Refresh(context.Background(), true)
|
||||||
close(release)
|
close(release)
|
||||||
<-done
|
<-done
|
||||||
|
|
||||||
|
|||||||
+135
-3
@@ -156,6 +156,101 @@ type parallelDialerResult struct {
|
|||||||
err error
|
err error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const (
|
||||||
|
// unreachableBackoffBase is the initial suppression window applied to a
|
||||||
|
// dial address after it returns a network-unreachable error (e.g.
|
||||||
|
// "connect: no route to host"). The window grows exponentially up to
|
||||||
|
// unreachableBackoffMax on repeated failures, and is cleared as soon as
|
||||||
|
// the address dials successfully.
|
||||||
|
unreachableBackoffBase = 5 * time.Second
|
||||||
|
// unreachableBackoffMax caps the suppression window so an address is
|
||||||
|
// always re-probed within a bounded interval, preserving recovery when
|
||||||
|
// the route comes back.
|
||||||
|
unreachableBackoffMax = 60 * time.Second
|
||||||
|
)
|
||||||
|
|
||||||
|
// Windows winsock codes for the unreachable errnos. A failing connect on
|
||||||
|
// Windows surfaces these raw WSA codes (WSAENETUNREACH/WSAEHOSTUNREACH), whereas
|
||||||
|
// syscall.ENETUNREACH/EHOSTUNREACH are Go's portable "invented" values
|
||||||
|
// (APPLICATION_ERROR + iota) that never equal them. Matching these explicitly is
|
||||||
|
// therefore required for the classifier to detect unreachable errors on Windows;
|
||||||
|
// errors.Is against the syscall.* constants alone would not.
|
||||||
|
//
|
||||||
|
// https://learn.microsoft.com/en-us/windows/win32/winsock/windows-sockets-error-codes-2
|
||||||
|
var (
|
||||||
|
windowsENETUNREACH = syscall.Errno(10051)
|
||||||
|
windowsEHOSTUNREACH = syscall.Errno(10065)
|
||||||
|
)
|
||||||
|
|
||||||
|
// IsUnreachable reports whether err indicates the destination network or host
|
||||||
|
// has no route (ENETUNREACH/EHOSTUNREACH). These are the errors produced when
|
||||||
|
// an endpoint's address family is available locally but unroutable, e.g. an
|
||||||
|
// IPv6 DoH endpoint while the host has IPv6 but no route to it.
|
||||||
|
func IsUnreachable(err error) bool {
|
||||||
|
if err == nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
var opErr *net.OpError
|
||||||
|
if errors.As(err, &opErr) {
|
||||||
|
return errors.Is(opErr.Err, syscall.ENETUNREACH) ||
|
||||||
|
errors.Is(opErr.Err, syscall.EHOSTUNREACH) ||
|
||||||
|
errors.Is(opErr.Err, windowsENETUNREACH) ||
|
||||||
|
errors.Is(opErr.Err, windowsEHOSTUNREACH)
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
type unreachableEntry struct {
|
||||||
|
until time.Time
|
||||||
|
backoff time.Duration
|
||||||
|
}
|
||||||
|
|
||||||
|
// unreachableTracker records dial addresses that recently failed with a
|
||||||
|
// network-unreachable error so ParallelDialer can temporarily stop hammering
|
||||||
|
// them. This prevents an unroutable endpoint from generating a sustained dial
|
||||||
|
// /health-check storm, while still re-probing each address once its bounded
|
||||||
|
// backoff window expires so genuine recovery is never permanently blocked.
|
||||||
|
type unreachableTracker struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
entries map[string]unreachableEntry
|
||||||
|
}
|
||||||
|
|
||||||
|
var unreachable = &unreachableTracker{entries: make(map[string]unreachableEntry)}
|
||||||
|
|
||||||
|
// suppressed reports whether addr is currently within its unreachable backoff
|
||||||
|
// window and should be skipped.
|
||||||
|
func (t *unreachableTracker) suppressed(addr string, now time.Time) bool {
|
||||||
|
t.mu.Lock()
|
||||||
|
defer t.mu.Unlock()
|
||||||
|
e, ok := t.entries[addr]
|
||||||
|
return ok && now.Before(e.until)
|
||||||
|
}
|
||||||
|
|
||||||
|
// markUnreachable extends the suppression window for addr using a bounded
|
||||||
|
// exponential backoff.
|
||||||
|
func (t *unreachableTracker) markUnreachable(addr string, now time.Time) {
|
||||||
|
t.mu.Lock()
|
||||||
|
defer t.mu.Unlock()
|
||||||
|
e := t.entries[addr]
|
||||||
|
if e.backoff == 0 {
|
||||||
|
e.backoff = unreachableBackoffBase
|
||||||
|
} else {
|
||||||
|
e.backoff *= 2
|
||||||
|
if e.backoff > unreachableBackoffMax {
|
||||||
|
e.backoff = unreachableBackoffMax
|
||||||
|
}
|
||||||
|
}
|
||||||
|
e.until = now.Add(e.backoff)
|
||||||
|
t.entries[addr] = e
|
||||||
|
}
|
||||||
|
|
||||||
|
// markReachable clears any suppression for addr after a successful dial.
|
||||||
|
func (t *unreachableTracker) markReachable(addr string) {
|
||||||
|
t.mu.Lock()
|
||||||
|
defer t.mu.Unlock()
|
||||||
|
delete(t.entries, addr)
|
||||||
|
}
|
||||||
|
|
||||||
type ParallelDialer struct {
|
type ParallelDialer struct {
|
||||||
net.Dialer
|
net.Dialer
|
||||||
}
|
}
|
||||||
@@ -164,26 +259,63 @@ func (d *ParallelDialer) DialContext(ctx context.Context, network string, addrs
|
|||||||
if len(addrs) == 0 {
|
if len(addrs) == 0 {
|
||||||
return nil, errors.New("empty addresses")
|
return nil, errors.New("empty addresses")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Skip addresses that recently returned a network-unreachable error so an
|
||||||
|
// unroutable endpoint (e.g. an IPv6 DoH address while the host has IPv6
|
||||||
|
// but no route to it) does not generate a sustained dial storm. Suppression
|
||||||
|
// is bounded: once an address's backoff window expires it is re-probed, so
|
||||||
|
// genuine recovery is preserved.
|
||||||
|
now := time.Now()
|
||||||
|
live := make([]string, 0, len(addrs))
|
||||||
|
var suppressed int
|
||||||
|
for _, addr := range addrs {
|
||||||
|
if unreachable.suppressed(addr, now) {
|
||||||
|
suppressed++
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
live = append(live, addr)
|
||||||
|
}
|
||||||
|
if len(live) == 0 {
|
||||||
|
// Every candidate is within its unreachable backoff window. Fail fast
|
||||||
|
// and quietly instead of re-dialing known-unroutable addresses; the
|
||||||
|
// windows expire and re-probe shortly, so recovery still happens.
|
||||||
|
logger.Debug("Skipping unreachable addresses, all in backoff", zap.Int("suppressed", suppressed))
|
||||||
|
// TODO: the errno here is hardcoded to EHOSTUNREACH, but the actual
|
||||||
|
// failure that triggered suppression may have been ENETUNREACH. This is
|
||||||
|
// harmless today (IsUnreachable treats both the same and nothing else
|
||||||
|
// inspects the errno), but if these errors are ever recorded/reported we
|
||||||
|
// should retain the real error in unreachableEntry and surface it here.
|
||||||
|
return nil, &net.OpError{Op: "dial", Net: network, Err: syscall.EHOSTUNREACH}
|
||||||
|
}
|
||||||
|
if suppressed > 0 {
|
||||||
|
logger.Debug("Skipping unreachable addresses in backoff", zap.Int("suppressed", suppressed))
|
||||||
|
}
|
||||||
|
|
||||||
ctx, cancel := context.WithCancel(ctx)
|
ctx, cancel := context.WithCancel(ctx)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
done := make(chan struct{})
|
done := make(chan struct{})
|
||||||
defer close(done)
|
defer close(done)
|
||||||
ch := make(chan *parallelDialerResult, len(addrs))
|
ch := make(chan *parallelDialerResult, len(live))
|
||||||
var wg sync.WaitGroup
|
var wg sync.WaitGroup
|
||||||
wg.Add(len(addrs))
|
wg.Add(len(live))
|
||||||
go func() {
|
go func() {
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
close(ch)
|
close(ch)
|
||||||
}()
|
}()
|
||||||
|
|
||||||
for _, addr := range addrs {
|
for _, addr := range live {
|
||||||
go func(addr string) {
|
go func(addr string) {
|
||||||
defer wg.Done()
|
defer wg.Done()
|
||||||
logger.Debug("Dialing to", zap.String("address", addr))
|
logger.Debug("Dialing to", zap.String("address", addr))
|
||||||
conn, err := d.Dialer.DialContext(ctx, network, addr)
|
conn, err := d.Dialer.DialContext(ctx, network, addr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Debug("Failed to dial", zap.String("address", addr), zap.Error(err))
|
logger.Debug("Failed to dial", zap.String("address", addr), zap.Error(err))
|
||||||
|
if IsUnreachable(err) {
|
||||||
|
unreachable.markUnreachable(addr, time.Now())
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
unreachable.markReachable(addr)
|
||||||
}
|
}
|
||||||
select {
|
select {
|
||||||
case ch <- ¶llelDialerResult{conn: conn, err: err}:
|
case ch <- ¶llelDialerResult{conn: conn, err: err}:
|
||||||
|
|||||||
@@ -2,10 +2,77 @@ package net
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"net"
|
||||||
|
"syscall"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
func TestIsUnreachable(t *testing.T) {
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
err error
|
||||||
|
want bool
|
||||||
|
}{
|
||||||
|
{"nil", nil, false},
|
||||||
|
{"enetunreach", &net.OpError{Op: "dial", Net: "tcp", Err: syscall.ENETUNREACH}, true},
|
||||||
|
{"ehostunreach", &net.OpError{Op: "dial", Net: "tcp", Err: syscall.EHOSTUNREACH}, true},
|
||||||
|
{"windows enetunreach", &net.OpError{Op: "dial", Net: "tcp", Err: windowsENETUNREACH}, true},
|
||||||
|
{"windows ehostunreach", &net.OpError{Op: "dial", Net: "tcp", Err: windowsEHOSTUNREACH}, true},
|
||||||
|
{"connection refused", &net.OpError{Op: "dial", Net: "tcp", Err: syscall.ECONNREFUSED}, false},
|
||||||
|
{"not an opError", syscall.ENETUNREACH, false},
|
||||||
|
}
|
||||||
|
for _, tc := range tests {
|
||||||
|
tc := tc
|
||||||
|
t.Run(tc.name, func(t *testing.T) {
|
||||||
|
if got := IsUnreachable(tc.err); got != tc.want {
|
||||||
|
t.Errorf("IsUnreachable(%v) = %v, want %v", tc.err, got, tc.want)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestUnreachableTracker(t *testing.T) {
|
||||||
|
tr := &unreachableTracker{entries: make(map[string]unreachableEntry)}
|
||||||
|
const addr = "[2606:1a40::22]:443"
|
||||||
|
now := time.Unix(0, 0)
|
||||||
|
|
||||||
|
// Not suppressed before any failure.
|
||||||
|
if tr.suppressed(addr, now) {
|
||||||
|
t.Fatal("addr suppressed before any failure")
|
||||||
|
}
|
||||||
|
|
||||||
|
// First failure suppresses for the base window.
|
||||||
|
tr.markUnreachable(addr, now)
|
||||||
|
if !tr.suppressed(addr, now.Add(unreachableBackoffBase-time.Millisecond)) {
|
||||||
|
t.Fatal("addr not suppressed within base backoff window")
|
||||||
|
}
|
||||||
|
if tr.suppressed(addr, now.Add(unreachableBackoffBase)) {
|
||||||
|
t.Fatal("addr still suppressed at end of base backoff window")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Backoff grows exponentially and is capped at the max.
|
||||||
|
tr.markUnreachable(addr, now)
|
||||||
|
if got := tr.entries[addr].backoff; got != 2*unreachableBackoffBase {
|
||||||
|
t.Fatalf("backoff after second failure = %v, want %v", got, 2*unreachableBackoffBase)
|
||||||
|
}
|
||||||
|
for i := 0; i < 10; i++ {
|
||||||
|
tr.markUnreachable(addr, now)
|
||||||
|
}
|
||||||
|
if got := tr.entries[addr].backoff; got != unreachableBackoffMax {
|
||||||
|
t.Fatalf("backoff not capped: got %v, want %v", got, unreachableBackoffMax)
|
||||||
|
}
|
||||||
|
|
||||||
|
// A successful dial clears suppression entirely.
|
||||||
|
tr.markReachable(addr)
|
||||||
|
if tr.suppressed(addr, now) {
|
||||||
|
t.Fatal("addr still suppressed after markReachable")
|
||||||
|
}
|
||||||
|
if _, ok := tr.entries[addr]; ok {
|
||||||
|
t.Fatal("entry not removed after markReachable")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestProbeStackTimeout(t *testing.T) {
|
func TestProbeStackTimeout(t *testing.T) {
|
||||||
done := make(chan struct{})
|
done := make(chan struct{})
|
||||||
started := make(chan struct{})
|
started := make(chan struct{})
|
||||||
|
|||||||
Reference in New Issue
Block a user