Merge branch 'master' into marko/remove-apphash

This commit is contained in:
Marko
2021-04-27 09:48:34 +00:00
committed by GitHub
112 changed files with 4146 additions and 1368 deletions
+67 -32
View File
@@ -5,6 +5,7 @@ import (
"context"
"errors"
"fmt"
"math"
"net"
"net/http"
_ "net/http/pprof" // nolint: gosec // securely exposed on separate, optional port
@@ -43,9 +44,8 @@ import (
"github.com/tendermint/tendermint/state/indexer"
blockidxkv "github.com/tendermint/tendermint/state/indexer/block/kv"
blockidxnull "github.com/tendermint/tendermint/state/indexer/block/null"
"github.com/tendermint/tendermint/state/txindex"
"github.com/tendermint/tendermint/state/txindex/kv"
"github.com/tendermint/tendermint/state/txindex/null"
"github.com/tendermint/tendermint/state/indexer/tx/kv"
"github.com/tendermint/tendermint/state/indexer/tx/null"
"github.com/tendermint/tendermint/statesync"
"github.com/tendermint/tendermint/store"
"github.com/tendermint/tendermint/types"
@@ -95,6 +95,7 @@ func DefaultNewNode(config *cfg.Config, logger log.Logger) (*Node, error) {
}
if config.Mode == cfg.ModeSeed {
return NewSeedNode(config,
DefaultDBProvider,
nodeKey,
DefaultGenesisDocProviderFunc(config),
logger,
@@ -224,9 +225,9 @@ type Node struct {
evidencePool *evidence.Pool // tracking evidence
proxyApp proxy.AppConns // connection to the application
rpcListeners []net.Listener // rpc servers
txIndexer txindex.TxIndexer
txIndexer indexer.TxIndexer
blockIndexer indexer.BlockIndexer
indexerService *txindex.IndexerService
indexerService *indexer.Service
prometheusSrv *http.Server
}
@@ -239,10 +240,6 @@ func initDBs(config *cfg.Config, dbProvider DBProvider) (blockStore *store.Block
blockStore = store.NewBlockStore(blockStoreDB)
stateDB, err = dbProvider(&DBContext{"state", config})
if err != nil {
return
}
return
}
@@ -269,10 +266,10 @@ func createAndStartIndexerService(
dbProvider DBProvider,
eventBus *types.EventBus,
logger log.Logger,
) (*txindex.IndexerService, txindex.TxIndexer, indexer.BlockIndexer, error) {
) (*indexer.Service, indexer.TxIndexer, indexer.BlockIndexer, error) {
var (
txIndexer txindex.TxIndexer
txIndexer indexer.TxIndexer
blockIndexer indexer.BlockIndexer
)
@@ -290,7 +287,7 @@ func createAndStartIndexerService(
blockIndexer = &blockidxnull.BlockerIndexer{}
}
indexerService := txindex.NewIndexerService(txIndexer, blockIndexer, eventBus)
indexerService := indexer.NewIndexerService(txIndexer, blockIndexer, eventBus)
indexerService.SetLogger(logger.With("module", "txindex"))
if err := indexerService.Start(); err != nil {
@@ -387,7 +384,7 @@ func createMempoolReactor(
peerUpdates *p2p.PeerUpdates
)
if config.P2P.UseNewP2P {
if config.P2P.DisableLegacy {
channels = makeChannelsFromShims(router, channelShims)
peerUpdates = peerManager.Subscribe()
} else {
@@ -438,7 +435,7 @@ func createEvidenceReactor(
peerUpdates *p2p.PeerUpdates
)
if config.P2P.UseNewP2P {
if config.P2P.DisableLegacy {
channels = makeChannelsFromShims(router, evidence.ChannelShims)
peerUpdates = peerManager.Subscribe()
} else {
@@ -479,7 +476,7 @@ func createBlockchainReactor(
peerUpdates *p2p.PeerUpdates
)
if config.P2P.UseNewP2P {
if config.P2P.DisableLegacy {
channels = makeChannelsFromShims(router, bcv0.ChannelShims)
peerUpdates = peerManager.Subscribe()
} else {
@@ -545,7 +542,7 @@ func createConsensusReactor(
peerUpdates *p2p.PeerUpdates
)
if config.P2P.UseNewP2P {
if config.P2P.DisableLegacy {
channels = makeChannelsFromShims(router, cs.ChannelShims)
peerUpdates = peerManager.Subscribe()
} else {
@@ -583,9 +580,35 @@ func createTransport(logger log.Logger, config *cfg.Config) *p2p.MConnTransport
)
}
func createPeerManager(config *cfg.Config, p2pLogger log.Logger, nodeID p2p.NodeID) (*p2p.PeerManager, error) {
func createPeerManager(
config *cfg.Config,
dbProvider DBProvider,
p2pLogger log.Logger,
nodeID p2p.NodeID) (*p2p.PeerManager, error) {
var maxConns uint16
switch {
case config.P2P.MaxConnections > 0:
maxConns = config.P2P.MaxConnections
case config.P2P.MaxNumInboundPeers > 0 && config.P2P.MaxNumOutboundPeers > 0:
x := config.P2P.MaxNumInboundPeers + config.P2P.MaxNumOutboundPeers
if x > math.MaxUint16 {
return nil, fmt.Errorf(
"max inbound peers (%d) + max outbound peers (%d) exceeds maximum (%d)",
config.P2P.MaxNumInboundPeers,
config.P2P.MaxNumOutboundPeers,
math.MaxUint16,
)
}
maxConns = uint16(x)
default:
maxConns = 64
}
options := p2p.PeerManagerOptions{
MaxConnected: 64,
MaxConnected: maxConns,
MaxConnectedUpgrade: 4,
MaxPeers: 1000,
MinRetryTime: 100 * time.Millisecond,
@@ -605,13 +628,25 @@ func createPeerManager(config *cfg.Config, p2pLogger log.Logger, nodeID p2p.Node
options.PersistentPeers = append(options.PersistentPeers, address.NodeID)
}
peerManager, err := p2p.NewPeerManager(nodeID, dbm.NewMemDB(), options)
for _, p := range splitAndTrimEmpty(config.P2P.BootstrapPeers, ",", " ") {
address, err := p2p.ParseNodeAddress(p)
if err != nil {
return nil, fmt.Errorf("invalid peer address %q: %w", p, err)
}
peers = append(peers, address)
}
peerDB, err := dbProvider(&DBContext{"peerstore", config})
if err != nil {
return nil, err
}
peerManager, err := p2p.NewPeerManager(nodeID, peerDB, options)
if err != nil {
return nil, fmt.Errorf("failed to create peer manager: %w", err)
}
for _, peer := range peers {
if err := peerManager.Add(peer); err != nil {
if _, err := peerManager.Add(peer); err != nil {
return nil, fmt.Errorf("failed to add peer %q: %w", peer, err)
}
}
@@ -853,6 +888,7 @@ func startStateSync(ssR *statesync.Reactor, bcR fastSyncReactor, conR *cs.Reacto
// NewSeedNode returns a new seed node, containing only p2p, pex reactor
func NewSeedNode(config *cfg.Config,
dbProvider DBProvider,
nodeKey p2p.NodeKey,
genesisDocProvider GenesisDocProvider,
logger log.Logger,
@@ -897,7 +933,7 @@ func NewSeedNode(config *cfg.Config,
return nil, fmt.Errorf("could not create addrbook: %w", err)
}
peerManager, err := createPeerManager(config, p2pLogger, nodeKey.ID)
peerManager, err := createPeerManager(config, dbProvider, p2pLogger, nodeKey.ID)
if err != nil {
return nil, fmt.Errorf("failed to create peer manager: %w", err)
}
@@ -1058,7 +1094,7 @@ func NewNode(config *cfg.Config,
p2pLogger := logger.With("module", "p2p")
transport := createTransport(p2pLogger, config)
peerManager, err := createPeerManager(config, p2pLogger, nodeKey.ID)
peerManager, err := createPeerManager(config, dbProvider, p2pLogger, nodeKey.ID)
if err != nil {
return nil, fmt.Errorf("failed to create peer manager: %w", err)
}
@@ -1138,7 +1174,7 @@ func NewNode(config *cfg.Config,
stateSyncReactorShim = p2p.NewReactorShim(logger.With("module", "statesync"), "StateSyncShim", statesync.ChannelShims)
if config.P2P.UseNewP2P {
if config.P2P.DisableLegacy {
channels = makeChannelsFromShims(router, statesync.ChannelShims)
peerUpdates = peerManager.Subscribe()
} else {
@@ -1305,9 +1341,9 @@ func (n *Node) OnStart() error {
n.isListening = true
n.Logger.Info("p2p service", "legacy_enabled", !n.config.P2P.UseNewP2P)
n.Logger.Info("p2p service", "legacy_enabled", !n.config.P2P.DisableLegacy)
if n.config.P2P.UseNewP2P {
if n.config.P2P.DisableLegacy {
err = n.router.Start()
} else {
err = n.sw.Start()
@@ -1345,7 +1381,7 @@ func (n *Node) OnStart() error {
}
}
if n.config.P2P.UseNewP2P && n.pexReactorV2 != nil {
if n.config.P2P.DisableLegacy && n.pexReactorV2 != nil {
if err := n.pexReactorV2.Start(); err != nil {
return err
}
@@ -1388,7 +1424,6 @@ func (n *Node) OnStop() {
}
if n.config.Mode != cfg.ModeSeed {
// now stop the reactors
if n.config.FastSync.Version == "v0" {
// Stop the real blockchain reactor separately since the switch uses the shim.
@@ -1418,13 +1453,13 @@ func (n *Node) OnStop() {
}
}
if n.config.P2P.UseNewP2P && n.pexReactorV2 != nil {
if n.config.P2P.DisableLegacy && n.pexReactorV2 != nil {
if err := n.pexReactorV2.Stop(); err != nil {
n.Logger.Error("failed to stop the PEX v2 reactor", "err", err)
}
}
if n.config.P2P.UseNewP2P {
if n.config.P2P.DisableLegacy {
if err := n.router.Stop(); err != nil {
n.Logger.Error("failed to stop router", "err", err)
}
@@ -1708,7 +1743,7 @@ func (n *Node) Config() *cfg.Config {
}
// TxIndexer returns the Node's TxIndexer.
func (n *Node) TxIndexer() txindex.TxIndexer {
func (n *Node) TxIndexer() indexer.TxIndexer {
return n.txIndexer
}
@@ -1732,7 +1767,7 @@ func (n *Node) NodeInfo() p2p.NodeInfo {
func makeNodeInfo(
config *cfg.Config,
nodeKey p2p.NodeKey,
txIndexer txindex.TxIndexer,
txIndexer indexer.TxIndexer,
genDoc *types.GenesisDoc,
state sm.State,
) (p2p.NodeInfo, error) {
@@ -1960,7 +1995,7 @@ func getRouterConfig(conf *cfg.Config, proxyApp proxy.AppConns) p2p.RouterOption
}
if conf.P2P.MaxNumInboundPeers > 0 {
opts.MaxIncommingConnectionsPerIP = uint(conf.P2P.MaxNumInboundPeers)
opts.MaxIncomingConnectionAttempts = conf.P2P.MaxIncomingConnectionAttempts
}
if conf.FilterPeers && proxyApp != nil {
+1
View File
@@ -533,6 +533,7 @@ func TestNodeNewSeedNode(t *testing.T) {
require.NoError(t, err)
n, err := NewSeedNode(config,
DefaultDBProvider,
nodeKey,
DefaultGenesisDocProviderFunc(config),
log.TestingLogger(),