What changed, and why it matters
This commit adds a new inbound connection rate-limiting and admission-control system to btcd. It is designed to reduce the risk of denial-of-service attacks where an attacker opens many incomplete handshake connections or forces the server to do expensive v2 transport cryptography for many sources at once. The change caps how many unfinished handshakes a single IP prefix can hold, limits the rate of CPU-heavy v2 responder handshakes both globally and per source, and reserves part of the peer budget for automatic outbound peers so inbound connections cannot starve them.
Treat as a hardening/DoS-prevention patch. Review the default constants (8 pending per /24 or /64, 20 sustained/40 burst global v2, 2 sustained/4 burst per source, 4 concurrent crypto slots) against expected node load and adversarial scenarios. Ensure the v2transport replace directive and dependency changes are intentional and that downstream consumers rebuild with the updated v2transport package. Run the new admission and server lifecycle tests.
Security signals we found
New DoS-mitigation admission control for inbound handshakes
Per-source-prefix limits on incomplete handshakes
Rate and concurrency limits on CPU-bound v2 responder handshake cryptography
Source-limit bypass restricted to whitelisted and loopback addresses
Handshake slot released exactly once via sync.Once on verack or disconnect
MaxInbound connection-manager budget added to reserve outbound peer capacity
No CVE, advisory, or vendor security disclosure present in commit or references
Evidence from the diff
The patch introduces internal/inbound/admission.go, which implements Admission with per-source-prefix incomplete-handshake accounting (IPv4 /24, IPv6 /64), a global token-bucket rate limiter for v2 responder handshakes, per-source token buckets, and a bounded channel for concurrent v2 crypto work. server.go wires this into inboundPeerConnected: it calls AcquireSource before constructing the peer, binds a V2Admission for inbound v2 responders, and releases the source slot exactly once on verack or disconnect. The connection manager now receives a MaxInbound value computed by maxInboundPeers, which reserves targetOutbound (up to 8) slots. Whitelisted and loopback sources bypass per-source limits but remain subject to global v2 and concurrency limits. Tests cover normalization, bypass behavior, concurrency, rate rollback on rejection, and lifecycle release.
Changed components
server.gopeer/peer.gointernal/inbound/admission.goconnmgr (MaxInbound integration)v2transport (responder handshake admission hook)Inspect captured patch +1001 / −18
diff --git a/config.go b/config.go
index 06a8798..3178557 100644
--- a/config.go
+++ b/config.go
@@ -129,7 +129,7 @@ type config struct {
Listeners []string `long:"listen" description:"Add an interface/port to listen for connections (default all interfaces port: 8333, testnet: 18333)"`
LogDir string `long:"logdir" description:"Directory to log output."`
MaxOrphanTxs int `long:"maxorphantx" description:"Max number of orphan transactions to keep in memory"`
- MaxPeers int `long:"maxpeers" description:"Max number of inbound and outbound peers"`
+ MaxPeers int `long:"maxpeers" description:"Max number of inbound and outbound peers. Up to 8 slots are reserved for automatic outbound peers; values of 8 or less disable inbound connections"`
MiningAddrs []string `long:"miningaddr" description:"Add the specified payment address to the list of addresses to use for generated blocks -- At least one address is required if the generate option is set"`
MinRelayTxFee float64 `long:"minrelaytxfee" description:"The minimum transaction fee in BTC/kB to be considered a non-zero fee."`
DisableBanning bool `long:"nobanning" description:"Disable banning of misbehaving peers"`
diff --git a/go.mod b/go.mod
index 78624c2..c707edb 100644
--- a/go.mod
+++ b/go.mod
@@ -42,6 +42,7 @@ require (
replace (
github.com/btcsuite/btcd/btcutil/v2 => ./btcutil
+ github.com/btcsuite/btcd/v2transport => ./v2transport
github.com/btcsuite/btcd/wire/v2 => ./wire
)
diff --git a/go.sum b/go.sum
index a7ea81c..f0c211f 100644
--- a/go.sum
+++ b/go.sum
@@ -4,18 +4,12 @@ github.com/btcsuite/btcd/address/v2 v2.0.0 h1:UVu8Hal6Siu4XastFe+JX5JkeBYONbDUIY
github.com/btcsuite/btcd/address/v2 v2.0.0/go.mod h1:htJK1AtaeK3bKNfZY63ep2oN8LbrI6qvmPGe1vekb3I=
github.com/btcsuite/btcd/btcec/v2 v2.5.0 h1:KioMXOWa76b86sTZZOmbzv/ldaQCmB8KFAyn5PbB8E8=
github.com/btcsuite/btcd/btcec/v2 v2.5.0/go.mod h1:+K/MYXcLBtHEQjRbjHuJChuybk4LCgjdjgRwil+e+Kk=
-github.com/btcsuite/btcd/btcutil/v2 v2.0.0 h1:77pgf/4tjWaSBLdos8yiWVWL3rSphxWNqkLwcyONExA=
-github.com/btcsuite/btcd/btcutil/v2 v2.0.0/go.mod h1:ZF8MMdsx1JGgvHJUanxbigekSO+8bN/ai34LBk/lg3c=
github.com/btcsuite/btcd/chaincfg/v2 v2.0.0 h1:M/RTtXfXA9odC1RUEOyZFXj/NXKVHPYZXVjb60xTOok=
github.com/btcsuite/btcd/chaincfg/v2 v2.0.0/go.mod h1:rHgHIXYYfn70m25a+BJ9f9z7VZAsTiDQGB2XYaippGQ=
github.com/btcsuite/btcd/chainhash/v2 v2.0.0 h1:PMLlSloHJuEeB80XG9EjpXWNEKAZAMLl6YHZ6YsEuoA=
github.com/btcsuite/btcd/chainhash/v2 v2.0.0/go.mod h1:mKxcZ7oGTXE7IRV+sS9hP4EVBwc/SzfNR+52IsOP9j8=
github.com/btcsuite/btcd/txscript/v2 v2.0.0 h1:pEmmHaC8eRx6KSB63zSVJD7qrit9/c9cLSrw++XrYP8=
github.com/btcsuite/btcd/txscript/v2 v2.0.0/go.mod h1:pZXabc11Xr9nz/18kXY3yErdAajYc3gi28Zqb3KqlFo=
-github.com/btcsuite/btcd/v2transport v1.0.1 h1:pIyyyBCPwd087K3Wdb/9tIvUubAQdzTJghjPgzTQVsE=
-github.com/btcsuite/btcd/v2transport v1.0.1/go.mod h1:N6H0HGSElVVJKntzaYHYVbW71DtWDLMw2yhwVRO3ZOE=
-github.com/btcsuite/btcd/wire/v2 v2.0.0 h1:mYSKzZZ0a1sK+aMhXzfDSVsSzRkWkU3x2U04TFRS2z8=
-github.com/btcsuite/btcd/wire/v2 v2.0.0/go.mod h1:bGxkPkk8IiDvUo1D96wE03llBIk7p2MdWYRyAQwLmqM=
github.com/btcsuite/btclog v1.0.0 h1:sEkpKJMmfGiyZjADwEIgB1NSwMyfdD1FB8v6+w1T0Ns=
github.com/btcsuite/btclog v1.0.0/go.mod h1:w7xnGOhwT3lmrS4H3b/D1XAXxvh+tbhUm8xeHN2y3TQ=
github.com/btcsuite/go-socks v0.0.0-20170105172521-4720035b7bfd h1:R/opQEbFEy9JGkIguV40SvRY1uliPX8ifOvi6ICsFCw=
diff --git a/integration/sync_race_test.go b/integration/sync_race_test.go
index 68cd5ba..2bf2fad 100644
--- a/integration/sync_race_test.go
+++ b/integration/sync_race_test.go
@@ -99,7 +99,12 @@ func fakePeerConn(nodeAddr string) error {
// sync, it was stuck with a dead sync peer (the sync manager still
// has a dead peer as sync peer, so it ignores the new live one).
func TestSyncManagerRaceCorruption(t *testing.T) {
- stressedHarness, err := rpctest.New(&chaincfg.SimNetParams, nil, nil, "")
+ // This test deliberately opens more concurrent peers than the default
+ // inbound limit. Raise the harness limit so admission does not dilute the
+ // sync-manager lifecycle stress this test is intended to apply.
+ stressedHarness, err := rpctest.New(
+ &chaincfg.SimNetParams, nil, []string{"--maxpeers=400"}, "",
+ )
require.NoError(t, err)
require.NoError(t, stressedHarness.SetUp(true, 0))
t.Cleanup(func() {
diff --git a/internal/inbound/admission.go b/internal/inbound/admission.go
new file mode 100644
index 0000000..acfe68d
--- /dev/null
+++ b/internal/inbound/admission.go
@@ -0,0 +1,379 @@
+package inbound
+
+import (
+ "errors"
+ "fmt"
+ "net"
+ "net/netip"
+ "sync"
+ "sync/atomic"
+ "time"
+
+ "github.com/decred/dcrd/lru"
+ "golang.org/x/time/rate"
+)
+
+const (
+ // inboundIPv4PrefixBits groups IPv4 sources by /24.
+ inboundIPv4PrefixBits = 24
+
+ // inboundIPv6PrefixBits groups IPv6 sources by /64.
+ inboundIPv6PrefixBits = 64
+
+ // defaultMaxPendingInboundHandshakes limits concurrent handshakes from a
+ // single normalized source prefix.
+ defaultMaxPendingInboundHandshakes = 8
+
+ // defaultV2HandshakeRate limits sustained inbound responder handshakes.
+ defaultV2HandshakeRate = 20
+
+ // defaultV2HandshakeBurst permits short bursts of responder handshakes.
+ defaultV2HandshakeBurst = 40
+
+ // defaultV2SourceHandshakeRate limits the sustained responder handshake
+ // rate from one normalized source prefix.
+ defaultV2SourceHandshakeRate = 2
+
+ // defaultV2SourceHandshakeBurst permits a short responder burst from one
+ // normalized source prefix.
+ defaultV2SourceHandshakeBurst = 4
+
+ // defaultV2SourceCacheSize bounds the number of source token buckets kept
+ // in memory.
+ defaultV2SourceCacheSize = 2048
+
+ // defaultV2HandshakeConcurrency limits concurrent responder cryptography.
+ defaultV2HandshakeConcurrency = 4
+)
+
+var (
+ // errInboundSourceLimit is returned when a source prefix already has the
+ // maximum number of incomplete handshakes.
+ errInboundSourceLimit = errors.New("inbound source handshake limit reached")
+
+ // errV2HandshakeRateLimit is returned when the responder handshake rate
+ // budget is exhausted.
+ errV2HandshakeRateLimit = errors.New("v2 handshake rate limit reached")
+
+ // errV2HandshakeSourceRateLimit is returned when one source prefix
+ // exhausts its responder handshake budget.
+ errV2HandshakeSourceRateLimit = errors.New("v2 source handshake rate limit reached")
+
+ // errV2HandshakeConcurrency is returned when all responder crypto slots
+ // are occupied.
+ errV2HandshakeConcurrency = errors.New("v2 handshake concurrency limit reached")
+)
+
+// admissionConfig defines the resource budgets used while accepting and
+// negotiating inbound connections.
+type admissionConfig struct {
+ maxPendingPerSource int
+ v2Rate rate.Limit
+ v2Burst int
+ v2SourceRate rate.Limit
+ v2SourceBurst int
+ v2SourceCacheSize uint
+ v2Concurrency int
+ now func() time.Time
+}
+
+// Admission tracks incomplete handshakes by normalized source prefix and
+// bounds the CPU-intensive portion of v2 responder handshakes.
+type Admission struct {
+ mu sync.Mutex
+ pendingBySource map[netip.Prefix]int
+ maxPendingSource int
+
+ v2Limiter *rate.Limiter
+ v2Slots chan struct{}
+ now func() time.Time
+
+ v2SourceMu sync.Mutex
+ v2SourceLimiters lru.KVCache
+ v2SourceRate rate.Limit
+ v2SourceBurst int
+
+ sourceRejected atomic.Uint64
+ v2Rejected atomic.Uint64
+ sourceLog rate.Sometimes
+ v2Log rate.Sometimes
+}
+
+// V2Admission binds the server-wide v2 admission policy to a single remote
+// address. The transport uses this value after it has classified the
+// connection as v2, but before it performs key generation or key agreement.
+type V2Admission struct {
+ admission *Admission
+ remote net.Addr
+ bypassSourceLimits bool
+}
+
+// Acquire reserves the v2 rate and concurrency budgets for the bound remote.
+func (a *V2Admission) Acquire() (func(), error) {
+ return a.admission.admitV2(a.remote, a.bypassSourceLimits)
+}
+
+// BindV2 binds the v2 admission policy to a remote address.
+func (a *Admission) BindV2(
+ remote net.Addr, bypassSourceLimits bool,
+) *V2Admission {
+
+ return &V2Admission{
+ admission: a,
+ remote: remote,
+ bypassSourceLimits: bypassSourceLimits,
+ }
+}
+
+// newAdmission creates an inbound handshake admission policy.
+func newAdmission(cfg admissionConfig) *Admission {
+
+ if cfg.now == nil {
+ cfg.now = time.Now
+ }
+
+ return &Admission{
+ pendingBySource: make(map[netip.Prefix]int),
+ maxPendingSource: cfg.maxPendingPerSource,
+ v2Limiter: rate.NewLimiter(cfg.v2Rate, cfg.v2Burst),
+ v2Slots: make(chan struct{}, cfg.v2Concurrency),
+ now: cfg.now,
+ v2SourceLimiters: lru.NewKVCache(cfg.v2SourceCacheSize),
+ v2SourceRate: cfg.v2SourceRate,
+ v2SourceBurst: cfg.v2SourceBurst,
+ sourceLog: rate.Sometimes{
+ First: 3,
+ Interval: 30 * time.Second,
+ },
+ v2Log: rate.Sometimes{
+ First: 3,
+ Interval: 30 * time.Second,
+ },
+ }
+}
+
+// New creates the production inbound admission policy.
+func New() *Admission {
+ return newAdmission(admissionConfig{
+ maxPendingPerSource: defaultMaxPendingInboundHandshakes,
+ v2Rate: rate.Limit(defaultV2HandshakeRate),
+ v2Burst: defaultV2HandshakeBurst,
+ v2SourceRate: rate.Limit(defaultV2SourceHandshakeRate),
+ v2SourceBurst: defaultV2SourceHandshakeBurst,
+ v2SourceCacheSize: defaultV2SourceCacheSize,
+ v2Concurrency: defaultV2HandshakeConcurrency,
+ })
+}
+
+// inboundSourceAddr returns the normalized IP for an inbound network address.
+// Ports and IPv6 zones do not affect the result.
+func inboundSourceAddr(addr net.Addr) (netip.Addr, error) {
+ if addr == nil {
+ return netip.Addr{}, errors.New("nil inbound address")
+ }
+
+ var ip netip.Addr
+ switch addr := addr.(type) {
+ case *net.TCPAddr:
+ var ok bool
+ ip, ok = netip.AddrFromSlice(addr.IP)
+ if !ok {
+ return netip.Addr{}, fmt.Errorf(
+ "invalid inbound TCP address: %v", addr,
+ )
+ }
+
+ default:
+ host, _, err := net.SplitHostPort(addr.String())
+ if err != nil {
+ return netip.Addr{}, fmt.Errorf(
+ "invalid inbound address %q: %w", addr.String(), err,
+ )
+ }
+
+ ip, err = netip.ParseAddr(host)
+ if err != nil {
+ return netip.Addr{}, fmt.Errorf(
+ "invalid inbound IP %q: %w", host, err,
+ )
+ }
+ }
+
+ if ip.Is6() {
+ ip = ip.WithZone("")
+ }
+ ip = ip.Unmap()
+ return ip, nil
+}
+
+// inboundSourcePrefix returns the normalized source prefix for an inbound
+// network address. Ports and IPv6 zones do not affect the result.
+func inboundSourcePrefix(addr net.Addr) (netip.Prefix, error) {
+ ip, err := inboundSourceAddr(addr)
+ if err != nil {
+ return netip.Prefix{}, err
+ }
+
+ bits := inboundIPv6PrefixBits
+ if ip.Is4() {
+ bits = inboundIPv4PrefixBits
+ }
+
+ return netip.PrefixFrom(ip, bits).Masked(), nil
+}
+
+// IsLoopback returns whether addr identifies a loopback source.
+func IsLoopback(addr net.Addr) bool {
+ ip, err := inboundSourceAddr(addr)
+ return err == nil && ip.IsLoopback()
+}
+
+// AcquireSource reserves one incomplete-handshake slot for addr. The returned
+// release function is safe to call more than once. Source limits are bypassed
+// when bypassSourceLimits is set, while the global connection limit remains
+// active at the listener boundary.
+func (a *Admission) AcquireSource(
+ addr net.Addr, bypassSourceLimits bool,
+) (func(), error) {
+
+ if bypassSourceLimits {
+ return func() {}, nil
+ }
+
+ prefix, err := inboundSourcePrefix(addr)
+ if err != nil {
+ return nil, err
+ }
+
+ a.mu.Lock()
+ if a.pendingBySource[prefix] >= a.maxPendingSource {
+ a.mu.Unlock()
+
+ rejected := a.sourceRejected.Add(1)
+ a.sourceLog.Do(func() {
+ log.Warnf("Inbound handshake source limit reached: "+
+ "rejected=%d source=%s", rejected, prefix)
+ })
+
+ return nil, errInboundSourceLimit
+ }
+ a.pendingBySource[prefix]++
+ a.mu.Unlock()
+
+ var once sync.Once
+ release := func() {
+ once.Do(func() {
+ a.mu.Lock()
+ defer a.mu.Unlock()
+
+ pending := a.pendingBySource[prefix] - 1
+ if pending == 0 {
+ delete(a.pendingBySource, prefix)
+ return
+ }
+
+ a.pendingBySource[prefix] = pending
+ })
+ }
+
+ return release, nil
+}
+
+// admitV2 reserves the rate and concurrency budgets for the CPU-bound portion
+// of an inbound v2 handshake. The returned release function only releases the
+// concurrency slot; consumed rate tokens are not returned.
+func (a *Admission) admitV2(
+ addr net.Addr, bypassSourceLimits bool,
+) (func(), error) {
+
+ now := a.now()
+ var sourceReservation *rate.Reservation
+ if !bypassSourceLimits {
+ prefix, err := inboundSourcePrefix(addr)
+ if err != nil {
+ return nil, err
+ }
+
+ var ok bool
+ sourceReservation, ok = a.reserveV2Source(prefix, now)
+ if !ok {
+ a.logV2Rejection(addr, "source-rate")
+ return nil, errV2HandshakeSourceRateLimit
+ }
+ }
+
+ globalReservation, ok := reserveImmediate(a.v2Limiter, now)
+ if !ok {
+ if sourceReservation != nil {
+ sourceReservation.CancelAt(now)
+ }
+ a.logV2Rejection(addr, "global-rate")
+ return nil, errV2HandshakeRateLimit
+ }
+
+ select {
+ case a.v2Slots <- struct{}{}:
+ var once sync.Once
+ return func() {
+ once.Do(func() {
+ <-a.v2Slots
+ })
+ }, nil
+
+ default:
+ globalReservation.CancelAt(now)
+ if sourceReservation != nil {
+ sourceReservation.CancelAt(now)
+ }
+ a.logV2Rejection(addr, "concurrency")
+ return nil, errV2HandshakeConcurrency
+ }
+}
+
+// reserveImmediate reserves one token only if the limiter permits immediate
+// admission. A delayed reservation is canceled before returning.
+func reserveImmediate(
+ limiter *rate.Limiter, now time.Time,
+) (*rate.Reservation, bool) {
+
+ reservation := limiter.ReserveN(now, 1)
+ if !reservation.OK() {
+ return nil, false
+ }
+ if reservation.DelayFrom(now) > 0 {
+ reservation.CancelAt(now)
+ return nil, false
+ }
+
+ return reservation, true
+}
+
+// reserveV2Source reserves one token for a normalized source prefix. The LRU
+// bounds retained source history while the global limiter remains the
+// authority when an evicted source returns.
+func (a *Admission) reserveV2Source(
+ prefix netip.Prefix, now time.Time,
+) (*rate.Reservation, bool) {
+
+ a.v2SourceMu.Lock()
+ value, ok := a.v2SourceLimiters.Lookup(prefix)
+ if !ok {
+ value = rate.NewLimiter(a.v2SourceRate, a.v2SourceBurst)
+ a.v2SourceLimiters.Add(prefix, value)
+ }
+ a.v2SourceMu.Unlock()
+
+ return reserveImmediate(value.(*rate.Limiter), now)
+}
+
+// logV2Rejection records and occasionally logs a rejected v2 handshake.
+func (a *Admission) logV2Rejection(
+ addr net.Addr, reason string,
+) {
+
+ rejected := a.v2Rejected.Add(1)
+ a.v2Log.Do(func() {
+ log.Warnf("Inbound v2 handshake limited: "+
+ "rejected=%d reason=%s remote=%s", rejected, reason, addr)
+ })
+}
diff --git a/internal/inbound/admission_test.go b/internal/inbound/admission_test.go
new file mode 100644
index 0000000..a29915b
--- /dev/null
+++ b/internal/inbound/admission_test.go
@@ -0,0 +1,371 @@
+package inbound
+
+import (
+ "errors"
+ "net"
+ "net/netip"
+ "testing"
+ "time"
+
+ "github.com/stretchr/testify/require"
+ "golang.org/x/time/rate"
+ "pgregory.net/rapid"
+)
+
+// stringAddr is a net.Addr with a caller-controlled string representation.
+type stringAddr string
+
+func (a stringAddr) Network() string { return "tcp" }
+func (a stringAddr) String() string { return string(a) }
+
+// TestInboundSourcePrefix verifies source normalization across address forms.
+func TestInboundSourcePrefix(t *testing.T) {
+ t.Parallel()
+
+ tests := []struct {
+ name string
+ addr net.Addr
+ want netip.Prefix
+ }{
+ {
+ name: "ipv4",
+ addr: &net.TCPAddr{IP: net.ParseIP("192.0.2.99"), Port: 1},
+ want: netip.MustParsePrefix("192.0.2.0/24"),
+ },
+ {
+ name: "ipv4 mapped",
+ addr: stringAddr("[::ffff:192.0.2.99]:8333"),
+ want: netip.MustParsePrefix("192.0.2.0/24"),
+ },
+ {
+ name: "ipv6",
+ addr: &net.TCPAddr{
+ IP: net.ParseIP("2001:db8:1:2:3:4:5:6"),
+ Port: 8333,
+ },
+ want: netip.MustParsePrefix("2001:db8:1:2::/64"),
+ },
+ {
+ name: "ipv6 zone",
+ addr: stringAddr("[fe80::1234%en0]:8333"),
+ want: netip.MustParsePrefix("fe80::/64"),
+ },
+ }
+
+ for _, test := range tests {
+ t.Run(test.name, func(t *testing.T) {
+ got, err := inboundSourcePrefix(test.addr)
+ require.NoError(t, err)
+ require.Equal(t, test.want, got)
+ })
+ }
+}
+
+// TestInboundSourcePrefixIgnoresPort verifies that ports cannot create new
+// source-accounting entries.
+func TestInboundSourcePrefixIgnoresPort(t *testing.T) {
+ t.Parallel()
+
+ rapid.Check(t, func(t *rapid.T) {
+ portA := rapid.IntRange(0, 65535).Draw(t, "port_a")
+ portB := rapid.IntRange(0, 65535).Draw(t, "port_b")
+ ip := net.IPv4(
+ rapid.Byte().Draw(t, "a"),
+ rapid.Byte().Draw(t, "b"),
+ rapid.Byte().Draw(t, "c"),
+ rapid.Byte().Draw(t, "d"),
+ )
+
+ prefixA, err := inboundSourcePrefix(&net.TCPAddr{
+ IP: ip, Port: portA,
+ })
+ require.NoError(t, err)
+ prefixB, err := inboundSourcePrefix(&net.TCPAddr{
+ IP: ip, Port: portB,
+ })
+ require.NoError(t, err)
+ require.Equal(t, prefixA, prefixB)
+ })
+}
+
+// TestIsLoopback verifies loopback detection across supported address forms.
+func TestIsLoopback(t *testing.T) {
+ t.Parallel()
+
+ tests := []struct {
+ name string
+ addr net.Addr
+ want bool
+ }{
+ {
+ name: "ipv4 loopback",
+ addr: &net.TCPAddr{
+ IP: net.ParseIP("127.0.0.2"), Port: 8333,
+ },
+ want: true,
+ },
+ {
+ name: "ipv6 loopback",
+ addr: stringAddr("[::1]:8333"),
+ want: true,
+ },
+ {
+ name: "public",
+ addr: stringAddr("192.0.2.1:8333"),
+ },
+ {
+ name: "malformed",
+ addr: stringAddr("attacker-controlled"),
+ },
+ }
+
+ for _, test := range tests {
+ t.Run(test.name, func(t *testing.T) {
+ require.Equal(t, test.want, IsLoopback(test.addr))
+ })
+ }
+}
+
+// TestInboundSourceAdmission verifies independent per-prefix limits, exact
+// release semantics, and map cleanup.
+func TestInboundSourceAdmission(t *testing.T) {
+ t.Parallel()
+
+ admission := newAdmission(admissionConfig{
+ maxPendingPerSource: 2,
+ v2Rate: rate.Inf,
+ v2SourceRate: rate.Inf,
+ v2SourceCacheSize: 16,
+ v2Concurrency: 1,
+ })
+
+ sourceA1 := &net.TCPAddr{IP: net.ParseIP("192.0.2.1"), Port: 1}
+ sourceA2 := &net.TCPAddr{IP: net.ParseIP("192.0.2.200"), Port: 2}
+ sourceB := &net.TCPAddr{IP: net.ParseIP("192.0.3.1"), Port: 3}
+
+ releaseA1, err := admission.AcquireSource(sourceA1, false)
+ require.NoError(t, err)
+ releaseA2, err := admission.AcquireSource(sourceA2, false)
+ require.NoError(t, err)
+
+ _, err = admission.AcquireSource(sourceA1, false)
+ require.ErrorIs(t, err, errInboundSourceLimit)
+
+ releaseB, err := admission.AcquireSource(sourceB, false)
+ require.NoError(t, err, "an independent prefix should remain admissible")
+
+ releaseA1()
+ releaseA1()
+ replacement, err := admission.AcquireSource(sourceA1, false)
+ require.NoError(t, err, "double release must only free one slot")
+
+ releaseA2()
+ replacement()
+ releaseB()
+
+ admission.mu.Lock()
+ defer admission.mu.Unlock()
+ require.Empty(t, admission.pendingBySource)
+}
+
+// TestInboundSourceAdmissionMalformedAddress verifies malformed addresses fail
+// closed without allocating map entries.
+func TestInboundSourceAdmissionMalformedAddress(t *testing.T) {
+ t.Parallel()
+
+ admission := newAdmission(admissionConfig{
+ maxPendingPerSource: 2,
+ v2Rate: rate.Inf,
+ v2SourceRate: rate.Inf,
+ v2SourceCacheSize: 16,
+ v2Concurrency: 1,
+ })
+
+ _, err := admission.AcquireSource(
+ stringAddr("attacker-controlled"), false,
+ )
+ require.Error(t, err)
+ require.Empty(t, admission.pendingBySource)
+}
+
+// TestBypassSourceLimits verifies trusted sources bypass only per-source
+// accounting while the global v2 rate and concurrency budgets remain active.
+func TestBypassSourceLimits(t *testing.T) {
+ t.Parallel()
+
+ now := time.Unix(1000, 0)
+ trusted := &net.TCPAddr{IP: net.ParseIP("192.0.2.1"), Port: 8333}
+ admission := newAdmission(admissionConfig{
+ maxPendingPerSource: 1,
+ v2Rate: 2,
+ v2Burst: 2,
+ v2SourceRate: 1,
+ v2SourceBurst: 1,
+ v2SourceCacheSize: 16,
+ v2Concurrency: 2,
+ now: func() time.Time { return now },
+ })
+
+ releaseSource1, err := admission.AcquireSource(trusted, true)
+ require.NoError(t, err)
+ releaseSource2, err := admission.AcquireSource(trusted, true)
+ require.NoError(t, err)
+ releaseSource1()
+ releaseSource1()
+ releaseSource2()
+ require.Empty(t, admission.pendingBySource)
+
+ releaseV2First, err := admission.BindV2(trusted, true).Acquire()
+ require.NoError(t, err)
+ releaseV2Second, err := admission.BindV2(trusted, true).Acquire()
+ require.NoError(t, err)
+
+ _, err = admission.BindV2(trusted, true).Acquire()
+ require.ErrorIs(t, err, errV2HandshakeRateLimit)
+
+ now = now.Add(500 * time.Millisecond)
+ _, err = admission.BindV2(trusted, true).Acquire()
+ require.ErrorIs(t, err, errV2HandshakeConcurrency)
+
+ releaseV2First()
+ releaseV2Second()
+}
+
+// TestV2HandshakeAdmission verifies deterministic token refill and concurrent
+// crypto-slot admission.
+func TestV2HandshakeAdmission(t *testing.T) {
+ t.Parallel()
+
+ now := time.Unix(1000, 0)
+ remote := &net.TCPAddr{IP: net.ParseIP("192.0.2.1"), Port: 8333}
+ admission := newAdmission(admissionConfig{
+ maxPendingPerSource: 1,
+ v2Rate: 2,
+ v2Burst: 2,
+ v2SourceRate: rate.Inf,
+ v2SourceCacheSize: 16,
+ v2Concurrency: 2,
+ now: func() time.Time { return now },
+ })
+
+ release1, err := admission.BindV2(remote, false).Acquire()
+ require.NoError(t, err)
+ release2, err := admission.BindV2(remote, false).Acquire()
+ require.NoError(t, err)
+
+ _, err = admission.BindV2(remote, false).Acquire()
+ require.ErrorIs(t, err, errV2HandshakeRateLimit)
+
+ release1()
+ release2()
+ now = now.Add(500 * time.Millisecond)
+
+ release3, err := admission.BindV2(remote, false).Acquire()
+ require.NoError(t, err, "one token should refill after half a second")
+ release3()
+
+ _, err = admission.BindV2(remote, false).Acquire()
+ require.ErrorIs(t, err, errV2HandshakeRateLimit)
+}
+
+// TestV2HandshakeConcurrency verifies that crypto slots are released exactly
+// once and do not wait when all slots are occupied.
+func TestV2HandshakeConcurrency(t *testing.T) {
+ t.Parallel()
+
+ now := time.Unix(1000, 0)
+ sourceA := &net.TCPAddr{IP: net.ParseIP("192.0.2.1"), Port: 1}
+ sourceB := &net.TCPAddr{IP: net.ParseIP("192.0.3.1"), Port: 2}
+ admission := newAdmission(admissionConfig{
+ maxPendingPerSource: 1,
+ v2Rate: 1,
+ v2Burst: 2,
+ v2SourceRate: 1,
+ v2SourceBurst: 1,
+ v2SourceCacheSize: 16,
+ v2Concurrency: 1,
+ now: func() time.Time { return now },
+ })
+
+ release, err := admission.BindV2(sourceA, false).Acquire()
+ require.NoError(t, err)
+ _, err = admission.BindV2(sourceB, false).Acquire()
+ require.True(t, errors.Is(err, errV2HandshakeConcurrency))
+
+ release()
+ release()
+ replacement, err := admission.BindV2(sourceB, false).Acquire()
+ require.NoError(t, err)
+ replacement()
+}
+
+// TestV2HandshakeGlobalRateRollback verifies a global rejection does not
+// consume the rejected source's independent token.
+func TestV2HandshakeGlobalRateRollback(t *testing.T) {
+ t.Parallel()
+
+ now := time.Unix(1000, 0)
+ admission := newAdmission(admissionConfig{
+ maxPendingPerSource: 1,
+ v2Rate: 1,
+ v2Burst: 1,
+ v2SourceRate: 0.01,
+ v2SourceBurst: 1,
+ v2SourceCacheSize: 16,
+ v2Concurrency: 1,
+ now: func() time.Time { return now },
+ })
+
+ sourceA := &net.TCPAddr{IP: net.ParseIP("192.0.2.1"), Port: 1}
+ sourceB := &net.TCPAddr{IP: net.ParseIP("192.0.3.1"), Port: 2}
+
+ release, err := admission.BindV2(sourceA, false).Acquire()
+ require.NoError(t, err)
+ release()
+
+ _, err = admission.BindV2(sourceB, false).Acquire()
+ require.ErrorIs(t, err, errV2HandshakeRateLimit)
+
+ now = now.Add(time.Second)
+ release, err = admission.BindV2(sourceB, false).Acquire()
+ require.NoError(t, err,
+ "a global rejection must return the source reservation")
+ release()
+}
+
+// TestV2HandshakeSourceRate verifies one source prefix cannot consume the
+// global responder budget or prevent another source from being admitted.
+func TestV2HandshakeSourceRate(t *testing.T) {
+ t.Parallel()
+
+ now := time.Unix(1000, 0)
+ admission := newAdmission(admissionConfig{
+ maxPendingPerSource: 1,
+ v2Rate: rate.Inf,
+ v2SourceRate: 1,
+ v2SourceBurst: 2,
+ v2SourceCacheSize: 16,
+ v2Concurrency: 1,
+ now: func() time.Time { return now },
+ })
+
+ sourceA := &net.TCPAddr{IP: net.ParseIP("192.0.2.1"), Port: 1}
+ sourceASamePrefix := &net.TCPAddr{
+ IP: net.ParseIP("192.0.2.200"), Port: 2,
+ }
+ sourceB := &net.TCPAddr{IP: net.ParseIP("192.0.3.1"), Port: 3}
+
+ for _, source := range []net.Addr{sourceA, sourceASamePrefix} {
+ release, err := admission.BindV2(source, false).Acquire()
+ require.NoError(t, err)
+ release()
+ }
+
+ _, err := admission.BindV2(sourceA, false).Acquire()
+ require.ErrorIs(t, err, errV2HandshakeSourceRateLimit)
+
+ release, err := admission.BindV2(sourceB, false).Acquire()
+ require.NoError(t, err,
+ "an independent source prefix should retain its own budget")
+ release()
+}
diff --git a/internal/inbound/log.go b/internal/inbound/log.go
new file mode 100644
index 0000000..1e0a411
--- /dev/null
+++ b/internal/inbound/log.go
@@ -0,0 +1,26 @@
+// Copyright (c) 2026 The btcsuite developers
+// Use of this source code is governed by an ISC
+// license that can be found in the LICENSE file.
+
+package inbound
+
+import "github.com/btcsuite/btclog"
+
+// log is initialized with no output filters. This means the package will not
+// perform any logging until the caller requests it.
+var log btclog.Logger
+
+// The default amount of logging is none.
+func init() {
+ DisableLog()
+}
+
+// DisableLog disables all package log output.
+func DisableLog() {
+ log = btclog.Disabled
+}
+
+// UseLogger uses the specified logger for package output.
+func UseLogger(logger btclog.Logger) {
+ log = logger
+}
diff --git a/log.go b/log.go
index 92a38eb..54dcb9f 100644
--- a/log.go
+++ b/log.go
@@ -15,6 +15,7 @@ import (
"github.com/btcsuite/btcd/blockchain/indexers"
"github.com/btcsuite/btcd/connmgr"
"github.com/btcsuite/btcd/database"
+ "github.com/btcsuite/btcd/internal/inbound"
"github.com/btcsuite/btcd/mempool"
"github.com/btcsuite/btcd/mining"
"github.com/btcsuite/btcd/mining/cpuminer"
@@ -78,6 +79,7 @@ func init() {
addrmgr.UseLogger(amgrLog)
connmgr.UseLogger(cmgrLog)
database.UseLogger(bcdbLog)
+ inbound.UseLogger(srvrLog)
blockchain.UseLogger(chanLog)
indexers.UseLogger(indxLog)
mining.UseLogger(minrLog)
diff --git a/peer/peer.go b/peer/peer.go
index 7124941..ca057b3 100644
--- a/peer/peer.go
+++ b/peer/peer.go
@@ -292,6 +292,11 @@ type Config struct {
// UsingV2Conn is defined if and only if we accept and attempt to make
// v2 connections.
UsingV2Conn bool
+
+ // V2HandshakeAdmission optionally reserves resources for the CPU-bound
+ // portion of an inbound v2 handshake. It is not used for outbound
+ // handshakes or v1 fallback.
+ V2HandshakeAdmission v2transport.HandshakeAdmission
}
// minUint32 is a helper function to return the minimum of two uint32s.
@@ -2553,7 +2558,16 @@ func newPeerBase(origCfg *Config, inbound bool) *Peer {
}
if p.cfg.UsingV2Conn && p.Services()&wire.SFNodeP2PV2 == wire.SFNodeP2PV2 {
- p.V2Transport = v2transport.NewPeer()
+ var options []v2transport.PeerOption
+ if inbound && p.cfg.V2HandshakeAdmission != nil {
+ options = append(options,
+ v2transport.WithResponderHandshakeAdmission(
+ p.cfg.V2HandshakeAdmission,
+ ),
+ )
+ }
+
+ p.V2Transport = v2transport.NewPeerWithOptions(options...)
} else {
// TODO: Hack, change.
p.cfg.UsingV2Conn = false
diff --git a/peer/peer_test.go b/peer/peer_test.go
index 4021672..2903cb9 100644
--- a/peer/peer_test.go
+++ b/peer/peer_test.go
@@ -10,6 +10,7 @@ import (
"io"
"net"
"strconv"
+ "sync/atomic"
"testing"
"time"
@@ -20,6 +21,15 @@ import (
"github.com/btcsuite/go-socks/socks"
)
+// testHandshakeAdmission adapts a function to the v2 handshake admission
+// interface used by the peer configuration.
+type testHandshakeAdmission func() (func(), error)
+
+// Acquire invokes the test admission function.
+func (a testHandshakeAdmission) Acquire() (func(), error) {
+ return a()
+}
+
// conn mocks a network connection by implementing the net.Conn interface. It
// is used to test peer connection without actually opening a network
// connection.
@@ -1202,3 +1212,65 @@ func TestSendAddrV2Handshake(t *testing.T) {
outPeer.WaitForDisconnect()
}
}
+
+// TestV2HandshakeAdmission verifies the responder admission hook is invoked
+// exactly once for an inbound v2 handshake and never for the initiator.
+func TestV2HandshakeAdmission(t *testing.T) {
+ verack := make(chan struct{}, 2)
+ var (
+ admissions atomic.Uint32
+ releases atomic.Uint32
+ )
+
+ newConfig := func() *peer.Config {
+ return &peer.Config{
+ Listeners: peer.MessageListeners{
+ OnVerAck: func(*peer.Peer, *wire.MsgVerAck) {
+ verack <- struct{}{}
+ },
+ },
+ AllowSelfConns: true,
+ ChainParams: &chaincfg.MainNetParams,
+ Services: wire.SFNodeNetwork | wire.SFNodeP2PV2,
+ UsingV2Conn: true,
+ TrickleInterval: time.Second,
+ V2HandshakeAdmission: testHandshakeAdmission(func() (func(), error) {
+ admissions.Add(1)
+ return func() {
+ releases.Add(1)
+ }, nil
+ }),
+ }
+ }
+
+ inbound := peer.NewInboundPeer(newConfig())
+ outbound, err := peer.NewOutboundPeer(newConfig(), "127.0.0.1:8333")
+ if err != nil {
+ t.Fatalf("NewOutboundPeer failed: %v", err)
+ }
+
+ if err := setupPeerConnection(inbound, outbound); err != nil {
+ t.Fatalf("setupPeerConnection failed: %v", err)
+ }
+ t.Cleanup(func() {
+ inbound.Disconnect()
+ outbound.Disconnect()
+ inbound.WaitForDisconnect()
+ outbound.WaitForDisconnect()
+ })
+
+ for i := 0; i < 2; i++ {
+ select {
+ case <-verack:
+ case <-time.After(5 * time.Second):
+ t.Fatal("timed out waiting for v2 verack")
+ }
+ }
+
+ if got := admissions.Load(); got != 1 {
+ t.Fatalf("admission invoked %d times, want 1", got)
+ }
+ if got := releases.Load(); got != 1 {
+ t.Fatalf("admission released %d times, want 1", got)
+ }
+}
diff --git a/server.go b/server.go
index b94abdc..a60875d 100644
--- a/server.go
+++ b/server.go
@@ -31,6 +31,7 @@ import (
"github.com/btcsuite/btcd/chainhash/v2"
"github.com/btcsuite/btcd/connmgr"
"github.com/btcsuite/btcd/database"
+ "github.com/btcsuite/btcd/internal/inbound"
"github.com/btcsuite/btcd/mempool"
"github.com/btcsuite/btcd/mining"
"github.com/btcsuite/btcd/mining/cpuminer"
@@ -73,6 +74,16 @@ var (
// zeroHash is the zero value hash (all zeros). It is defined as a convenience.
var zeroHash chainhash.Hash
+// maxInboundPeers returns the accepted inbound connection budget after
+// reserving capacity for automatic outbound peers.
+func maxInboundPeers(maxPeers, targetOutbound int) uint32 {
+ if maxPeers <= targetOutbound {
+ return 0
+ }
+
+ return uint32(maxPeers - targetOutbound)
+}
+
// onionAddr implements the net.Addr interface and represents a tor address.
type onionAddr struct {
addr string
@@ -237,6 +248,7 @@ type server struct {
cpuMiner *cpuminer.CPUMiner
modifyRebroadcastInv chan interface{}
p2pDowngrader *peer.P2PDowngrader
+ inboundAdmission *inbound.Admission
peerLifecycle chan peerLifecycleEvent
banPeers chan *serverPeer
query chan interface{}
@@ -296,12 +308,17 @@ type serverPeer struct {
addressesMtx sync.RWMutex
knownAddresses lru.Cache
banScore connmgr.DynamicBanScore
- quit chan struct{}
+ quit chan struct{}
// Closed by verAckOnce when OnVerAck fires.
verAckCh chan struct{}
verAckOnce sync.Once
+ // releaseInboundHandshake releases the source-prefix slot held while an
+ // inbound handshake is incomplete.
+ releaseInboundHandshake func()
+ releaseInboundHandshakeOnce sync.Once
+
// peerAdded is set by peerLifecycleHandler after a peerAdd event
// has been enqueued on s.peerLifecycle. handleDonePeerMsg reads
// this to decide whether to notify the sync manager and evict
@@ -331,6 +348,15 @@ func newServerPeer(s *server, isPersistent bool) *serverPeer {
}
}
+// releaseHandshake releases the source-prefix slot for an inbound handshake.
+func (sp *serverPeer) releaseHandshake() {
+ sp.releaseInboundHandshakeOnce.Do(func() {
+ if sp.releaseInboundHandshake != nil {
+ sp.releaseInboundHandshake()
+ }
+ })
+}
+
// newestBlock returns the current best block hash and height using the format
// required by the configuration for the peer package.
func (sp *serverPeer) newestBlock() (*chainhash.Hash, int32, error) {
@@ -574,7 +600,10 @@ func (sp *serverPeer) OnVersion(_ *peer.Peer, msg *wire.MsgVersion) *wire.MsgRej
// sync.Once guard ensures verAckCh is closed at most once even if
// OnVerAck is ever invoked more than once for a given peer.
func (sp *serverPeer) OnVerAck(_ *peer.Peer, _ *wire.MsgVerAck) {
- sp.verAckOnce.Do(func() { close(sp.verAckCh) })
+ sp.verAckOnce.Do(func() {
+ sp.releaseHandshake()
+ close(sp.verAckCh)
+ })
}
// OnMemPool is invoked when a peer receives a mempool bitcoin message.
@@ -2282,9 +2311,38 @@ func newPeerConfig(sp *serverPeer) *peer.Config {
// instance, associates it with the connection, and starts a goroutine to wait
// for disconnection.
func (s *server) inboundPeerConnected(conn net.Conn) {
+ remoteAddr := conn.RemoteAddr()
+ whitelisted := isWhitelisted(remoteAddr)
+
+ // Loopback includes onion peers forwarded into the listener. Bypass the
+ // per-source limits for these and configured whitelisted peers, while the
+ // global socket, v2 rate, and v2 concurrency limits remain active.
+ bypassSourceLimits := whitelisted || inbound.IsLoopback(remoteAddr)
+
+ var releaseHandshake func()
+ if s.inboundAdmission != nil {
+ var err error
+ releaseHandshake, err = s.inboundAdmission.AcquireSource(
+ remoteAddr, bypassSourceLimits,
+ )
+ if err != nil {
+ _ = conn.Close()
+ return
+ }
+ }
+
sp := newServerPeer(s, false)
- sp.isWhitelisted = isWhitelisted(conn.RemoteAddr())
- sp.Peer = peer.NewInboundPeer(newPeerConfig(sp))
+ sp.isWhitelisted = whitelisted
+ sp.releaseInboundHandshake = releaseHandshake
+
+ peerCfg := newPeerConfig(sp)
+ if s.inboundAdmission != nil {
+ peerCfg.V2HandshakeAdmission = s.inboundAdmission.BindV2(
+ remoteAddr, bypassSourceLimits,
+ )
+ }
+
+ sp.Peer = peer.NewInboundPeer(peerCfg)
sp.AssociateConnection(conn)
go s.peerLifecycleHandler(sp)
}
@@ -2347,9 +2405,10 @@ func (s *server) peerLifecycleHandler(sp *serverPeer) {
}
sp.peerAdded.Store(true)
- case <-sp.Peer.Done():
+ case <-sp.Done():
// Disconnected before verack; no peerAdd needed.
}
+ sp.releaseHandshake()
// Wait for full disconnect (may already be done).
sp.WaitForDisconnect()
@@ -2903,8 +2962,9 @@ func newServer(listenAddrs, agentBlacklist, agentWhitelist []string,
}
s := server{
- chainParams: chainParams,
- addrManager: amgr,
+ chainParams: chainParams,
+ addrManager: amgr,
+ inboundAdmission: inbound.New(),
// peerLifecycle is buffered for up to two events per peer
// (peerAdd followed by peerDone) so peerLifecycleHandler
@@ -3145,9 +3205,15 @@ func newServer(listenAddrs, agentBlacklist, agentWhitelist []string,
if cfg.MaxPeers < targetOutbound {
targetOutbound = cfg.MaxPeers
}
+ maxInbound := maxInboundPeers(cfg.MaxPeers, targetOutbound)
+ if maxInbound == 0 && len(listeners) > 0 {
+ srvrLog.Infof("Inbound connections disabled: maxpeers=%d, "+
+ "reserved-outbound=%d", cfg.MaxPeers, targetOutbound)
+ }
cmgr, err := connmgr.New(&connmgr.Config{
Listeners: listeners,
OnAccept: s.inboundPeerConnected,
+ MaxInbound: &maxInbound,
RetryDuration: connectionRetryInterval,
TargetOutbound: uint32(targetOutbound),
Dial: btcdDial,
diff --git a/server_test.go b/server_test.go
index 9cbfad3..b205e96 100644
--- a/server_test.go
+++ b/server_test.go
@@ -3,6 +3,7 @@ package main
import (
"os"
"path/filepath"
+ "sync/atomic"
"testing"
"time"
@@ -61,6 +62,10 @@ func TestOnVerAckDoubleCall(t *testing.T) {
t.Parallel()
_, sp := newTestServerPeer(t)
+ var releases atomic.Uint32
+ sp.releaseInboundHandshake = func() {
+ releases.Add(1)
+ }
sp.OnVerAck(nil, nil)
@@ -79,6 +84,54 @@ func TestOnVerAckDoubleCall(t *testing.T) {
default:
t.Fatal("verAckCh should still be closed after second OnVerAck call")
}
+ require.Equal(t, uint32(1), releases.Load(),
+ "the source-prefix slot must be released exactly once")
+}
+
+// TestHandshakeReleaseOnDisconnect verifies that a connection which
+// disconnects before verack releases its source-prefix slot.
+func TestHandshakeReleaseOnDisconnect(t *testing.T) {
+ t.Parallel()
+
+ s, sp := newTestServerPeer(t)
+ var releases atomic.Uint32
+ sp.releaseInboundHandshake = func() {
+ releases.Add(1)
+ }
+
+ sp.Disconnect()
+ go s.peerLifecycleHandler(sp)
+
+ event := recvLifecycleEvent(t, s.peerLifecycle)
+ require.Equal(t, peerDone, event.action)
+ require.Equal(t, uint32(1), releases.Load())
+}
+
+// TestMaxInboundPeers verifies that automatic outbound capacity is reserved
+// without underflow at small peer limits.
+func TestMaxInboundPeers(t *testing.T) {
+ t.Parallel()
+
+ tests := []struct {
+ name string
+ maxPeers int
+ targetOutbound int
+ want uint32
+ }{
+ {name: "zero", maxPeers: 0, targetOutbound: 0, want: 0},
+ {name: "all outbound", maxPeers: 8, targetOutbound: 8, want: 0},
+ {name: "below outbound", maxPeers: 7, targetOutbound: 8, want: 0},
+ {name: "one inbound", maxPeers: 9, targetOutbound: 8, want: 1},
+ {name: "default", maxPeers: 125, targetOutbound: 8, want: 117},
+ }
+
+ for _, test := range tests {
+ t.Run(test.name, func(t *testing.T) {
+ require.Equal(t, test.want, maxInboundPeers(
+ test.maxPeers, test.targetOutbound,
+ ))
+ })
+ }
}
// TestPeerLifecycleOrdering verifies that when verack arrives before
@@ -100,7 +153,7 @@ func TestPeerLifecycleOrdering(t *testing.T) {
require.Equal(t, sp, first.sp)
// Trigger disconnect after peerAdd is observed.
- sp.Peer.Disconnect()
+ sp.Disconnect()
second := recvLifecycleEvent(t, s.peerLifecycle)
require.Equal(t, peerDone, second.action,
@@ -123,7 +176,7 @@ func TestPeerLifecycleSimultaneousReady(t *testing.T) {
s, sp := newTestServerPeer(t)
close(sp.verAckCh)
- sp.Peer.Disconnect()
+ sp.Disconnect()
go s.peerLifecycleHandler(sp)
Why this scored 61/100
Community notes
Notes can correct, qualify, or add evidence to the AI analysis. Every note shown here has been validated by a human moderator.
The AI analysis stands alone for now. Submit a note if you can add evidence or important context.