mirror of
https://github.com/tendermint/tendermint.git
synced 2026-09-19 22:44:24 +00:00
Merge branch 'master' into wb/abci-prepare-proposal-synchronize
This commit is contained in:
@@ -69,6 +69,9 @@ func configSetup(t *testing.T) *config.Config {
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() { os.RemoveAll(configByzantineTest.RootDir) })
|
||||
|
||||
walDir := filepath.Dir(cfg.Consensus.WalFile())
|
||||
ensureDir(t, walDir, 0700)
|
||||
|
||||
return cfg
|
||||
}
|
||||
|
||||
@@ -787,6 +790,7 @@ func makeConsensusState(
|
||||
configOpts ...func(*config.Config),
|
||||
) ([]*State, cleanupFunc) {
|
||||
t.Helper()
|
||||
tempDir := t.TempDir()
|
||||
|
||||
valSet, privVals := factory.ValidatorSet(ctx, t, nValidators, 30)
|
||||
genDoc := factory.GenesisDoc(cfg, time.Now(), valSet.Validators, nil)
|
||||
@@ -801,7 +805,7 @@ func makeConsensusState(
|
||||
blockStore := store.NewBlockStore(dbm.NewMemDB()) // each state needs its own db
|
||||
state, err := sm.MakeGenesisState(genDoc)
|
||||
require.NoError(t, err)
|
||||
thisConfig, err := ResetConfig(t.TempDir(), fmt.Sprintf("%s_%d", testName, i))
|
||||
thisConfig, err := ResetConfig(tempDir, fmt.Sprintf("%s_%d", testName, i))
|
||||
require.NoError(t, err)
|
||||
|
||||
configRootDirs = append(configRootDirs, thisConfig.RootDir)
|
||||
@@ -810,7 +814,8 @@ func makeConsensusState(
|
||||
opt(thisConfig)
|
||||
}
|
||||
|
||||
ensureDir(t, filepath.Dir(thisConfig.Consensus.WalFile()), 0700) // dir for wal
|
||||
walDir := filepath.Dir(thisConfig.Consensus.WalFile())
|
||||
ensureDir(t, walDir, 0700)
|
||||
|
||||
app := kvstore.NewApplication()
|
||||
closeFuncs = append(closeFuncs, app.Close)
|
||||
|
||||
@@ -5,7 +5,6 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"path"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -466,7 +465,6 @@ func TestReactorWithEvidence(t *testing.T) {
|
||||
|
||||
defer os.RemoveAll(thisConfig.RootDir)
|
||||
|
||||
ensureDir(t, path.Dir(thisConfig.Consensus.WalFile()), 0700) // dir for wal
|
||||
app := kvstore.NewApplication()
|
||||
vals := types.TM2PB.ValidatorUpdates(state.Validators)
|
||||
app.InitChain(abci.RequestInitChain{Validators: vals})
|
||||
@@ -564,7 +562,6 @@ func TestReactorCreatesBlockWhenEmptyBlocksFalse(t *testing.T) {
|
||||
c.Consensus.CreateEmptyBlocks = false
|
||||
},
|
||||
)
|
||||
|
||||
t.Cleanup(cleanup)
|
||||
|
||||
rts := setup(ctx, t, n, states, 100) // buffer must be large enough to not deadlock
|
||||
|
||||
+25
-12
@@ -20,6 +20,7 @@ import (
|
||||
cstypes "github.com/tendermint/tendermint/internal/consensus/types"
|
||||
"github.com/tendermint/tendermint/internal/eventbus"
|
||||
"github.com/tendermint/tendermint/internal/jsontypes"
|
||||
"github.com/tendermint/tendermint/internal/libs/autofile"
|
||||
sm "github.com/tendermint/tendermint/internal/state"
|
||||
tmevents "github.com/tendermint/tendermint/libs/events"
|
||||
"github.com/tendermint/tendermint/libs/log"
|
||||
@@ -869,15 +870,27 @@ func (cs *State) receiveRoutine(ctx context.Context, maxSteps int) {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
cs.logger.Error("CONSENSUS FAILURE!!!", "err", r, "stack", string(debug.Stack()))
|
||||
// stop gracefully
|
||||
//
|
||||
// NOTE: We most probably shouldn't be running any further when there is
|
||||
// some unexpected panic. Some unknown error happened, and so we don't
|
||||
// know if that will result in the validator signing an invalid thing. It
|
||||
// might be worthwhile to explore a mechanism for manual resuming via
|
||||
// some console or secure RPC system, but for now, halting the chain upon
|
||||
// unexpected consensus bugs sounds like the better option.
|
||||
|
||||
// Make a best-effort attempt to close the WAL, but otherwise do not
|
||||
// attempt to gracefully terminate. Once consensus has irrecoverably
|
||||
// failed, any additional progress we permit the node to make may
|
||||
// complicate diagnosing and recovering from the failure.
|
||||
onExit(cs)
|
||||
|
||||
// Re-panic to ensure the node terminates.
|
||||
//
|
||||
// TODO(creachadair): In ordinary operation, the WAL autofile should
|
||||
// never be closed. This only happens during shutdown and production
|
||||
// nodes usually halt by panicking. Many existing tests, however,
|
||||
// assume a clean shutdown is possible. Prior to #8111, we were
|
||||
// swallowing the panic in receiveRoutine, making that appear to
|
||||
// work. Filtering this specific error is slightly risky, but should
|
||||
// affect only unit tests. In any case, not re-panicking here only
|
||||
// preserves the pre-existing behavior for this one error type.
|
||||
if err, ok := r.(error); ok && errors.Is(err, autofile.ErrAutoFileClosed) {
|
||||
return
|
||||
}
|
||||
panic(r)
|
||||
}
|
||||
}()
|
||||
|
||||
@@ -906,8 +919,8 @@ func (cs *State) receiveRoutine(ctx context.Context, maxSteps int) {
|
||||
case mi := <-cs.internalMsgQueue:
|
||||
err := cs.wal.WriteSync(mi) // NOTE: fsync
|
||||
if err != nil {
|
||||
panic(fmt.Sprintf(
|
||||
"failed to write %v msg to consensus WAL due to %v; check your file system and restart the node",
|
||||
panic(fmt.Errorf(
|
||||
"failed to write %v msg to consensus WAL due to %w; check your file system and restart the node",
|
||||
mi, err,
|
||||
))
|
||||
}
|
||||
@@ -1906,8 +1919,8 @@ func (cs *State) finalizeCommit(ctx context.Context, height int64) {
|
||||
// restart).
|
||||
endMsg := EndHeightMessage{height}
|
||||
if err := cs.wal.WriteSync(endMsg); err != nil { // NOTE: fsync
|
||||
panic(fmt.Sprintf(
|
||||
"failed to write %v msg to consensus WAL due to %v; check your file system and restart the node",
|
||||
panic(fmt.Errorf(
|
||||
"failed to write %v msg to consensus WAL due to %w; check your file system and restart the node",
|
||||
endMsg, err,
|
||||
))
|
||||
}
|
||||
|
||||
@@ -41,9 +41,9 @@ const (
|
||||
autoFilePerms = os.FileMode(0600)
|
||||
)
|
||||
|
||||
// errAutoFileClosed is reported when operations attempt to use an autofile
|
||||
// ErrAutoFileClosed is reported when operations attempt to use an autofile
|
||||
// after it has been closed.
|
||||
var errAutoFileClosed = errors.New("autofile is closed")
|
||||
var ErrAutoFileClosed = errors.New("autofile is closed")
|
||||
|
||||
// AutoFile automatically closes and re-opens file for writing. The file is
|
||||
// automatically setup to close itself every 1s and upon receiving SIGHUP.
|
||||
@@ -155,7 +155,7 @@ func (af *AutoFile) Write(b []byte) (n int, err error) {
|
||||
af.mtx.Lock()
|
||||
defer af.mtx.Unlock()
|
||||
if af.closed {
|
||||
return 0, fmt.Errorf("write: %w", errAutoFileClosed)
|
||||
return 0, fmt.Errorf("write: %w", ErrAutoFileClosed)
|
||||
}
|
||||
|
||||
if af.file == nil {
|
||||
@@ -174,7 +174,7 @@ func (af *AutoFile) Write(b []byte) (n int, err error) {
|
||||
func (af *AutoFile) Sync() error {
|
||||
return af.withLock(func() error {
|
||||
if af.closed {
|
||||
return fmt.Errorf("sync: %w", errAutoFileClosed)
|
||||
return fmt.Errorf("sync: %w", ErrAutoFileClosed)
|
||||
} else if af.file == nil {
|
||||
return nil // nothing to sync
|
||||
}
|
||||
@@ -207,7 +207,7 @@ func (af *AutoFile) Size() (int64, error) {
|
||||
af.mtx.Lock()
|
||||
defer af.mtx.Unlock()
|
||||
if af.closed {
|
||||
return 0, fmt.Errorf("size: %w", errAutoFileClosed)
|
||||
return 0, fmt.Errorf("size: %w", ErrAutoFileClosed)
|
||||
}
|
||||
|
||||
if af.file == nil {
|
||||
|
||||
+96
-137
@@ -3,14 +3,12 @@ package pex
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"runtime/debug"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/tendermint/tendermint/internal/p2p"
|
||||
"github.com/tendermint/tendermint/internal/p2p/conn"
|
||||
"github.com/tendermint/tendermint/libs/log"
|
||||
tmmath "github.com/tendermint/tendermint/libs/math"
|
||||
"github.com/tendermint/tendermint/libs/service"
|
||||
protop2p "github.com/tendermint/tendermint/proto/tendermint/p2p"
|
||||
"github.com/tendermint/tendermint/types"
|
||||
@@ -42,7 +40,7 @@ const (
|
||||
minReceiveRequestInterval = 100 * time.Millisecond
|
||||
|
||||
// the maximum amount of addresses that can be included in a response
|
||||
maxAddresses uint16 = 100
|
||||
maxAddresses = 100
|
||||
|
||||
// How long to wait when there are no peers available before trying again
|
||||
noAvailablePeersWaitPeriod = 1 * time.Second
|
||||
@@ -100,15 +98,8 @@ type Reactor struct {
|
||||
// minReceiveRequestInterval).
|
||||
lastReceivedRequests map[types.NodeID]time.Time
|
||||
|
||||
// keep track of how many new peers to existing peers we have received to
|
||||
// extrapolate the size of the network
|
||||
newPeers uint32
|
||||
totalPeers uint32
|
||||
|
||||
// discoveryRatio is the inverse ratio of new peers to old peers squared.
|
||||
// This is multiplied by the minimum duration to calculate how long to wait
|
||||
// between each request.
|
||||
discoveryRatio float32
|
||||
// the total number of unique peers added
|
||||
totalPeers int
|
||||
}
|
||||
|
||||
// NewReactor returns a reference to a new reactor.
|
||||
@@ -156,16 +147,6 @@ func (r *Reactor) OnStop() {}
|
||||
// processPexCh implements a blocking event loop where we listen for p2p
|
||||
// Envelope messages from the pexCh.
|
||||
func (r *Reactor) processPexCh(ctx context.Context) {
|
||||
timer := time.NewTimer(0)
|
||||
defer timer.Stop()
|
||||
|
||||
r.mtx.Lock()
|
||||
var (
|
||||
duration = r.calculateNextRequestTime()
|
||||
err error
|
||||
)
|
||||
r.mtx.Unlock()
|
||||
|
||||
incoming := make(chan *p2p.Envelope)
|
||||
go func() {
|
||||
defer close(incoming)
|
||||
@@ -179,36 +160,51 @@ func (r *Reactor) processPexCh(ctx context.Context) {
|
||||
}
|
||||
}()
|
||||
|
||||
// Initially, we will request peers quickly to bootstrap. This duration
|
||||
// will be adjusted upward as knowledge of the network grows.
|
||||
var nextPeerRequest = minReceiveRequestInterval
|
||||
|
||||
timer := time.NewTimer(0)
|
||||
defer timer.Stop()
|
||||
|
||||
for {
|
||||
timer.Reset(duration)
|
||||
timer.Reset(nextPeerRequest)
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
|
||||
// outbound requests for new peers
|
||||
case <-timer.C:
|
||||
duration, err = r.sendRequestForPeers(ctx)
|
||||
if err != nil {
|
||||
// Send a request for more peer addresses.
|
||||
if err := r.sendRequestForPeers(ctx); err != nil {
|
||||
return
|
||||
// TODO(creachadair): Do we really want to stop processing the PEX
|
||||
// channel just because of an error here?
|
||||
}
|
||||
// inbound requests for new peers or responses to requests sent by this
|
||||
// reactor
|
||||
|
||||
// Note we do not update the poll timer upon making a request, only
|
||||
// when we receive an update that updates our priors.
|
||||
|
||||
case envelope, ok := <-incoming:
|
||||
if !ok {
|
||||
return
|
||||
return // channel closed
|
||||
}
|
||||
duration, err = r.handleMessage(ctx, r.pexCh.ID, envelope)
|
||||
|
||||
// A request from another peer, or a response to one of our requests.
|
||||
dur, err := r.handlePexMessage(ctx, envelope)
|
||||
if err != nil {
|
||||
r.logger.Error("failed to process message", "ch_id", r.pexCh.ID, "envelope", envelope, "err", err)
|
||||
r.logger.Error("failed to process message",
|
||||
"ch_id", r.pexCh.ID, "envelope", envelope, "err", err)
|
||||
if serr := r.pexCh.SendError(ctx, p2p.PeerError{
|
||||
NodeID: envelope.From,
|
||||
Err: err,
|
||||
}); serr != nil {
|
||||
return
|
||||
}
|
||||
} else if dur != 0 {
|
||||
// We got a useful result; update the poll timer.
|
||||
nextPeerRequest = dur
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -228,19 +224,20 @@ func (r *Reactor) processPeerUpdates(ctx context.Context) {
|
||||
}
|
||||
|
||||
// handlePexMessage handles envelopes sent from peers on the PexChannel.
|
||||
// If an update was received, a new polling interval is returned; otherwise the
|
||||
// duration is 0.
|
||||
func (r *Reactor) handlePexMessage(ctx context.Context, envelope *p2p.Envelope) (time.Duration, error) {
|
||||
logger := r.logger.With("peer", envelope.From)
|
||||
|
||||
switch msg := envelope.Message.(type) {
|
||||
case *protop2p.PexRequest:
|
||||
// check if the peer hasn't sent a prior request too close to this one
|
||||
// in time
|
||||
// Verify that this peer hasn't sent us another request too recently.
|
||||
if err := r.markPeerRequest(envelope.From); err != nil {
|
||||
return time.Minute, err
|
||||
return 0, err
|
||||
}
|
||||
|
||||
// request peers from the peer manager and parse the NodeAddresses into
|
||||
// URL strings
|
||||
// Fetch peers from the peer manager, convert NodeAddresses into URL
|
||||
// strings, and send them back to the caller.
|
||||
nodeAddresses := r.peerManager.Advertise(envelope.From, maxAddresses)
|
||||
pexAddresses := make([]protop2p.PexAddress, len(nodeAddresses))
|
||||
for idx, addr := range nodeAddresses {
|
||||
@@ -248,28 +245,24 @@ func (r *Reactor) handlePexMessage(ctx context.Context, envelope *p2p.Envelope)
|
||||
URL: addr.String(),
|
||||
}
|
||||
}
|
||||
if err := r.pexCh.Send(ctx, p2p.Envelope{
|
||||
return 0, r.pexCh.Send(ctx, p2p.Envelope{
|
||||
To: envelope.From,
|
||||
Message: &protop2p.PexResponse{Addresses: pexAddresses},
|
||||
}); err != nil {
|
||||
})
|
||||
|
||||
case *protop2p.PexResponse:
|
||||
// Verify that this response corresponds to one of our pending requests.
|
||||
if err := r.markPeerResponse(envelope.From); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
return time.Second, nil
|
||||
case *protop2p.PexResponse:
|
||||
// check if the response matches a request that was made to that peer
|
||||
if err := r.markPeerResponse(envelope.From); err != nil {
|
||||
return time.Minute, err
|
||||
}
|
||||
|
||||
// check the size of the response
|
||||
if len(msg.Addresses) > int(maxAddresses) {
|
||||
return 10 * time.Minute, fmt.Errorf("peer sent too many addresses (max: %d, got: %d)",
|
||||
maxAddresses,
|
||||
len(msg.Addresses),
|
||||
)
|
||||
// Verify that the response does not exceed the safety limit.
|
||||
if len(msg.Addresses) > maxAddresses {
|
||||
return 0, fmt.Errorf("peer sent too many addresses (%d > maxiumum %d)",
|
||||
len(msg.Addresses), maxAddresses)
|
||||
}
|
||||
|
||||
var numAdded int
|
||||
for _, pexAddress := range msg.Addresses {
|
||||
peerAddress, err := p2p.ParseNodeAddress(pexAddress.URL)
|
||||
if err != nil {
|
||||
@@ -278,47 +271,21 @@ func (r *Reactor) handlePexMessage(ctx context.Context, envelope *p2p.Envelope)
|
||||
added, err := r.peerManager.Add(peerAddress)
|
||||
if err != nil {
|
||||
logger.Error("failed to add PEX address", "address", peerAddress, "err", err)
|
||||
continue
|
||||
}
|
||||
if added {
|
||||
r.newPeers++
|
||||
numAdded++
|
||||
logger.Debug("added PEX address", "address", peerAddress)
|
||||
}
|
||||
r.totalPeers++
|
||||
}
|
||||
|
||||
return 10 * time.Minute, nil
|
||||
return r.calculateNextRequestTime(numAdded), nil
|
||||
|
||||
default:
|
||||
return time.Second, fmt.Errorf("received unknown message: %T", msg)
|
||||
return 0, fmt.Errorf("received unknown message: %T", msg)
|
||||
}
|
||||
}
|
||||
|
||||
// handleMessage handles an Envelope sent from a peer on a specific p2p Channel.
|
||||
// It will handle errors and any possible panics gracefully. A caller can handle
|
||||
// any error returned by sending a PeerError on the respective channel.
|
||||
func (r *Reactor) handleMessage(ctx context.Context, chID p2p.ChannelID, envelope *p2p.Envelope) (duration time.Duration, err error) {
|
||||
defer func() {
|
||||
if e := recover(); e != nil {
|
||||
err = fmt.Errorf("panic in processing message: %v", e)
|
||||
r.logger.Error(
|
||||
"recovering from processing message panic",
|
||||
"err", err,
|
||||
"stack", string(debug.Stack()),
|
||||
)
|
||||
}
|
||||
}()
|
||||
|
||||
r.logger.Debug("received PEX message", "peer", envelope.From)
|
||||
|
||||
switch chID {
|
||||
case p2p.ChannelID(PexChannel):
|
||||
duration, err = r.handlePexMessage(ctx, envelope)
|
||||
default:
|
||||
err = fmt.Errorf("unknown channel ID (%d) for envelope (%v)", chID, envelope)
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
// processPeerUpdate processes a PeerUpdate. For added peers, PeerStatusUp, we
|
||||
// send a request for addresses.
|
||||
func (r *Reactor) processPeerUpdate(peerUpdate p2p.PeerUpdate) {
|
||||
@@ -338,95 +305,87 @@ func (r *Reactor) processPeerUpdate(peerUpdate p2p.PeerUpdate) {
|
||||
}
|
||||
}
|
||||
|
||||
// sendRequestForPeers pops the first peerID off the list and sends the
|
||||
// peer a request for more peer addresses. The function then moves the
|
||||
// peer into the requestsSent bucket and calculates when the next request
|
||||
// time should be
|
||||
func (r *Reactor) sendRequestForPeers(ctx context.Context) (time.Duration, error) {
|
||||
// sendRequestForPeers chooses a peer from the set of available peers and sends
|
||||
// that peer a request for more peer addresses. The chosen peer is moved into
|
||||
// the requestsSent bucket so that we will not attempt to contact them again
|
||||
// until they've replied or updated.
|
||||
func (r *Reactor) sendRequestForPeers(ctx context.Context) error {
|
||||
r.mtx.Lock()
|
||||
defer r.mtx.Unlock()
|
||||
if len(r.availablePeers) == 0 {
|
||||
// no peers are available
|
||||
r.logger.Debug("no available peers to send request to, waiting...")
|
||||
return noAvailablePeersWaitPeriod, nil
|
||||
r.logger.Debug("no available peers to send a PEX request to (retrying)")
|
||||
return nil
|
||||
}
|
||||
var peerID types.NodeID
|
||||
|
||||
// use range to get a random peer.
|
||||
// Select an arbitrary peer from the available set.
|
||||
var peerID types.NodeID
|
||||
for peerID = range r.availablePeers {
|
||||
break
|
||||
}
|
||||
|
||||
// send out the pex request
|
||||
if err := r.pexCh.Send(ctx, p2p.Envelope{
|
||||
To: peerID,
|
||||
Message: &protop2p.PexRequest{},
|
||||
}); err != nil {
|
||||
return 0, err
|
||||
return err
|
||||
}
|
||||
|
||||
// remove the peer from the abvailable peers list and mark it in the requestsSent map
|
||||
// Move the peer from available to pending.
|
||||
delete(r.availablePeers, peerID)
|
||||
r.requestsSent[peerID] = struct{}{}
|
||||
|
||||
dur := r.calculateNextRequestTime()
|
||||
r.logger.Debug("peer request sent", "next_request_time", dur)
|
||||
return dur, nil
|
||||
return nil
|
||||
}
|
||||
|
||||
// calculateNextRequestTime implements something of a proportional controller
|
||||
// to estimate how often the reactor should be requesting new peer addresses.
|
||||
// The dependent variable in this calculation is the ratio of new peers to
|
||||
// all peers that the reactor receives. The interval is thus calculated as the
|
||||
// inverse squared. In the beginning, all peers should be new peers.
|
||||
// We expect this ratio to be near 1 and thus the interval to be as short
|
||||
// as possible. As the node becomes more familiar with the network the ratio of
|
||||
// new nodes will plummet to a very small number, meaning the interval expands
|
||||
// to its upper bound.
|
||||
// calculateNextRequestTime selects how long we should wait before attempting
|
||||
// to send out another request for peer addresses.
|
||||
//
|
||||
// CONTRACT: The caller must hold r.mtx exclusively when calling this method.
|
||||
func (r *Reactor) calculateNextRequestTime() time.Duration {
|
||||
// check if the peer store is full. If so then there is no need
|
||||
// to send peer requests too often
|
||||
// This implements a simplified proportional control mechanism to poll more
|
||||
// often when our knowledge of the network is incomplete, and less often as our
|
||||
// knowledge grows. To estimate our knowledge of the network, we use the
|
||||
// fraction of "new" peers (addresses we have not previously seen) to the total
|
||||
// so far observed. When we first join the network, this fraction will be close
|
||||
// to 1, meaning most new peers are "new" to us, and as we discover more peers,
|
||||
// the fraction will go toward zero.
|
||||
//
|
||||
// The minimum interval will be minReceiveRequestInterval to ensure we will not
|
||||
// request from any peer more often than we would allow them to do from us.
|
||||
func (r *Reactor) calculateNextRequestTime(added int) time.Duration {
|
||||
r.mtx.Lock()
|
||||
defer r.mtx.Unlock()
|
||||
|
||||
r.totalPeers += added
|
||||
|
||||
// If the peer store is nearly full, wait the maximum interval.
|
||||
if ratio := r.peerManager.PeerRatio(); ratio >= 0.95 {
|
||||
r.logger.Debug("peer manager near full ratio, sleeping...",
|
||||
r.logger.Debug("Peer manager is nearly full",
|
||||
"sleep_period", fullCapacityInterval, "ratio", ratio)
|
||||
return fullCapacityInterval
|
||||
}
|
||||
|
||||
// baseTime represents the shortest interval that we can send peer requests
|
||||
// in. For example if we have 10 peers and we can't send a message to the
|
||||
// same peer every 500ms, then we can send a request every 50ms. In practice
|
||||
// we use a safety margin of 2, ergo 100ms
|
||||
peers := tmmath.MinInt(len(r.availablePeers), 50)
|
||||
baseTime := minReceiveRequestInterval
|
||||
if peers > 0 {
|
||||
baseTime = minReceiveRequestInterval * 2 / time.Duration(peers)
|
||||
// If there are no available peers to query, poll less aggressively.
|
||||
if len(r.availablePeers) == 0 {
|
||||
r.logger.Debug("No available peers to send a PEX request",
|
||||
"sleep_period", noAvailablePeersWaitPeriod)
|
||||
return noAvailablePeersWaitPeriod
|
||||
}
|
||||
|
||||
if r.totalPeers > 0 || r.discoveryRatio == 0 {
|
||||
// find the ratio of new peers. NOTE: We add 1 to both sides to avoid
|
||||
// divide by zero problems
|
||||
ratio := float32(r.totalPeers+1) / float32(r.newPeers+1)
|
||||
// square the ratio in order to get non linear time intervals
|
||||
// NOTE: The longest possible interval for a network with 100 or more peers
|
||||
// where a node is connected to 50 of them is 2 minutes.
|
||||
r.discoveryRatio = ratio * ratio
|
||||
r.newPeers = 0
|
||||
r.totalPeers = 0
|
||||
}
|
||||
// NOTE: As ratio is always >= 1, discovery ratio is >= 1. Therefore we don't need to worry
|
||||
// about the next request time being less than the minimum time
|
||||
return baseTime * time.Duration(r.discoveryRatio)
|
||||
// Reaching here, there are available peers to query and the peer store
|
||||
// still has space. Estimate our knowledge of the network from the latest
|
||||
// update and choose a new interval.
|
||||
base := float64(minReceiveRequestInterval) / float64(len(r.availablePeers))
|
||||
multiplier := float64(r.totalPeers+1) / float64(added+1) // +1 to avert zero division
|
||||
return time.Duration(base*multiplier*multiplier) + minReceiveRequestInterval
|
||||
}
|
||||
|
||||
func (r *Reactor) markPeerRequest(peer types.NodeID) error {
|
||||
r.mtx.Lock()
|
||||
defer r.mtx.Unlock()
|
||||
if lastRequestTime, ok := r.lastReceivedRequests[peer]; ok {
|
||||
if time.Now().Before(lastRequestTime.Add(minReceiveRequestInterval)) {
|
||||
return fmt.Errorf("peer sent a request too close after a prior one. Minimum interval: %v",
|
||||
minReceiveRequestInterval)
|
||||
if d := time.Since(lastRequestTime); d < minReceiveRequestInterval {
|
||||
return fmt.Errorf("peer %v sent PEX request too soon (%v < minimum %v)",
|
||||
peer, d, minReceiveRequestInterval)
|
||||
}
|
||||
}
|
||||
r.lastReceivedRequests[peer] = time.Now()
|
||||
|
||||
@@ -96,7 +96,7 @@ func TestReactorSendsRequestsTooOften(t *testing.T) {
|
||||
peerErr := <-r.pexErrCh
|
||||
require.Error(t, peerErr.Err)
|
||||
require.Empty(t, r.pexOutCh)
|
||||
require.Contains(t, peerErr.Err.Error(), "peer sent a request too close after a prior one")
|
||||
require.Contains(t, peerErr.Err.Error(), "sent PEX request too soon")
|
||||
require.Equal(t, badNode, peerErr.NodeID)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user