mirror of
https://github.com/tendermint/tendermint.git
synced 2026-09-21 07:24:36 +00:00
Merge branch 'master' into wb/add-consensus-param-internal
This commit is contained in:
@@ -21,6 +21,7 @@ Special thanks to external contributors on this release:
|
||||
- [rpc] \#7982 Add new Events interface and deprecate Subscribe. (@creachadair)
|
||||
- [cli] \#8081 make the reset command safe to use by intoducing `reset-state` command. Fixed by \#8259. (@marbar3778, @cmwaters)
|
||||
- [config] \#8222 default indexer configuration to null. (@creachadair)
|
||||
- [rpc] \#8570 rework timeouts to be per-method instead of global. (@creachadair)
|
||||
|
||||
- Apps
|
||||
|
||||
|
||||
@@ -171,7 +171,8 @@ for applications built w/ Cosmos SDK).
|
||||
// If necessary adjust global WriteTimeout to ensure it's greater than
|
||||
// TimeoutBroadcastTxCommit.
|
||||
// See https://github.com/tendermint/tendermint/issues/3435
|
||||
if cfg.WriteTimeout <= conf.RPC.TimeoutBroadcastTxCommit {
|
||||
// Note we don't need to adjust anything if the timeout is already unlimited.
|
||||
if cfg.WriteTimeout > 0 && cfg.WriteTimeout <= conf.RPC.TimeoutBroadcastTxCommit {
|
||||
cfg.WriteTimeout = conf.RPC.TimeoutBroadcastTxCommit + 1*time.Second
|
||||
}
|
||||
|
||||
|
||||
+1
-1
@@ -523,7 +523,7 @@ func DefaultRPCConfig() *RPCConfig {
|
||||
MaxSubscriptionClients: 100,
|
||||
MaxSubscriptionsPerClient: 5,
|
||||
ExperimentalDisableWebsocket: false, // compatible with TM v0.35 and earlier
|
||||
EventLogWindowSize: 0, // disables /events RPC by default
|
||||
EventLogWindowSize: 30 * time.Second,
|
||||
EventLogMaxItems: 0,
|
||||
|
||||
TimeoutBroadcastTxCommit: 10 * time.Second,
|
||||
|
||||
@@ -16,7 +16,7 @@ require (
|
||||
github.com/gorilla/websocket v1.5.0
|
||||
github.com/grpc-ecosystem/go-grpc-middleware v1.3.0
|
||||
github.com/grpc-ecosystem/go-grpc-prometheus v1.2.0
|
||||
github.com/lib/pq v1.10.5
|
||||
github.com/lib/pq v1.10.6
|
||||
github.com/libp2p/go-buffer-pool v0.0.2
|
||||
github.com/mroth/weightedrand v0.4.1
|
||||
github.com/oasisprotocol/curve25519-voi v0.0.0-20210609091139-0a56a4bca00b
|
||||
|
||||
@@ -678,8 +678,8 @@ github.com/lib/pq v1.0.0/go.mod h1:5WUZQaWbwv1U+lTReE5YruASi9Al49XbQIvNi/34Woo=
|
||||
github.com/lib/pq v1.8.0/go.mod h1:AlVN5x4E4T544tWzH6hKfbfQvm3HdbOxrmggDNAPY9o=
|
||||
github.com/lib/pq v1.9.0/go.mod h1:AlVN5x4E4T544tWzH6hKfbfQvm3HdbOxrmggDNAPY9o=
|
||||
github.com/lib/pq v1.10.4/go.mod h1:AlVN5x4E4T544tWzH6hKfbfQvm3HdbOxrmggDNAPY9o=
|
||||
github.com/lib/pq v1.10.5 h1:J+gdV2cUmX7ZqL2B0lFcW0m+egaHC2V3lpO8nWxyYiQ=
|
||||
github.com/lib/pq v1.10.5/go.mod h1:AlVN5x4E4T544tWzH6hKfbfQvm3HdbOxrmggDNAPY9o=
|
||||
github.com/lib/pq v1.10.6 h1:jbk+ZieJ0D7EVGJYpL9QTz7/YW6UHbmdnZWYyK5cdBs=
|
||||
github.com/lib/pq v1.10.6/go.mod h1:AlVN5x4E4T544tWzH6hKfbfQvm3HdbOxrmggDNAPY9o=
|
||||
github.com/libp2p/go-buffer-pool v0.0.2 h1:QNK2iAFa8gjAe1SPz6mHSMuCcjs+X1wlHzeOSqcmlfs=
|
||||
github.com/libp2p/go-buffer-pool v0.0.2/go.mod h1:MvaB6xw5vOrDl8rYZGLFdKAuk/hRoRZd1Vi32+RXyFM=
|
||||
github.com/lufeee/execinquery v1.0.0 h1:1XUTuLIVPDlFvUU3LXmmZwHDsolsxXnY67lzhpeqe0I=
|
||||
|
||||
@@ -125,7 +125,8 @@ func serverRPCConfig(r *config.RPCConfig) *server.Config {
|
||||
// If necessary adjust global WriteTimeout to ensure it's greater than
|
||||
// TimeoutBroadcastTxCommit.
|
||||
// See https://github.com/tendermint/tendermint/issues/3435
|
||||
if cfg.WriteTimeout <= r.TimeoutBroadcastTxCommit {
|
||||
// Note we don't need to adjust anything if the timeout is already unlimited.
|
||||
if cfg.WriteTimeout > 0 && cfg.WriteTimeout <= r.TimeoutBroadcastTxCommit {
|
||||
cfg.WriteTimeout = r.TimeoutBroadcastTxCommit + 1*time.Second
|
||||
}
|
||||
return cfg
|
||||
|
||||
+10
-30
@@ -6,7 +6,6 @@ import (
|
||||
"fmt"
|
||||
"runtime/debug"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/tendermint/tendermint/config"
|
||||
"github.com/tendermint/tendermint/internal/libs/clist"
|
||||
@@ -22,13 +21,6 @@ var (
|
||||
_ p2p.Wrapper = (*protomem.Message)(nil)
|
||||
)
|
||||
|
||||
// PeerManager defines the interface contract required for getting necessary
|
||||
// peer information. This should eventually be replaced with a message-oriented
|
||||
// approach utilizing the p2p stack.
|
||||
type PeerManager interface {
|
||||
GetHeight(types.NodeID) int64
|
||||
}
|
||||
|
||||
// Reactor implements a service that contains mempool of txs that are broadcasted
|
||||
// amongst peers. It maintains a map from peer ID to counter, to prevent gossiping
|
||||
// txs to the peers you received it from.
|
||||
@@ -40,9 +32,8 @@ type Reactor struct {
|
||||
mempool *TxMempool
|
||||
ids *IDs
|
||||
|
||||
getPeerHeight func(types.NodeID) int64
|
||||
peerEvents p2p.PeerEventSubscriber
|
||||
chCreator p2p.ChannelCreator
|
||||
peerEvents p2p.PeerEventSubscriber
|
||||
chCreator p2p.ChannelCreator
|
||||
|
||||
// observePanic is a function for observing panics that were recovered in methods on
|
||||
// Reactor. observePanic is called with the recovered value.
|
||||
@@ -59,18 +50,16 @@ func NewReactor(
|
||||
txmp *TxMempool,
|
||||
chCreator p2p.ChannelCreator,
|
||||
peerEvents p2p.PeerEventSubscriber,
|
||||
getPeerHeight func(types.NodeID) int64,
|
||||
) *Reactor {
|
||||
r := &Reactor{
|
||||
logger: logger,
|
||||
cfg: cfg,
|
||||
mempool: txmp,
|
||||
ids: NewMempoolIDs(),
|
||||
chCreator: chCreator,
|
||||
peerEvents: peerEvents,
|
||||
getPeerHeight: getPeerHeight,
|
||||
peerRoutines: make(map[types.NodeID]context.CancelFunc),
|
||||
observePanic: defaultObservePanic,
|
||||
logger: logger,
|
||||
cfg: cfg,
|
||||
mempool: txmp,
|
||||
ids: NewMempoolIDs(),
|
||||
chCreator: chCreator,
|
||||
peerEvents: peerEvents,
|
||||
peerRoutines: make(map[types.NodeID]context.CancelFunc),
|
||||
observePanic: defaultObservePanic,
|
||||
}
|
||||
|
||||
r.BaseService = *service.NewBaseService(logger, "Mempool", r)
|
||||
@@ -327,15 +316,6 @@ func (r *Reactor) broadcastTxRoutine(ctx context.Context, peerID types.NodeID, m
|
||||
|
||||
memTx := nextGossipTx.Value.(*WrappedTx)
|
||||
|
||||
if r.getPeerHeight != nil {
|
||||
height := r.getPeerHeight(peerID)
|
||||
if height > 0 && height < memTx.height-1 {
|
||||
// allow for a lag of one block
|
||||
time.Sleep(PeerCatchupSleepIntervalMS * time.Millisecond)
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
// NOTE: Transaction batching was disabled due to:
|
||||
// https://github.com/tendermint/tendermint/issues/5796
|
||||
if ok := r.mempool.txStore.TxHasPeer(memTx.hash, peerMempoolID); !ok {
|
||||
|
||||
@@ -85,7 +85,6 @@ func setupReactors(ctx context.Context, t *testing.T, logger log.Logger, numNode
|
||||
mempool,
|
||||
chCreator,
|
||||
func(ctx context.Context) *p2p.PeerUpdates { return rts.peerUpdates[nodeID] },
|
||||
rts.network.Nodes[nodeID].PeerManager.GetHeight,
|
||||
)
|
||||
rts.nodes = append(rts.nodes, nodeID)
|
||||
|
||||
|
||||
@@ -1027,37 +1027,6 @@ func (m *PeerManager) retryDelay(failures uint32, persistent bool) time.Duration
|
||||
return delay
|
||||
}
|
||||
|
||||
// GetHeight returns a peer's height, as reported via SetHeight, or 0 if the
|
||||
// peer or height is unknown.
|
||||
//
|
||||
// FIXME: This is a temporary workaround to share state between the consensus
|
||||
// and mempool reactors, carried over from the legacy P2P stack. Reactors should
|
||||
// not have dependencies on each other, instead tracking this themselves.
|
||||
func (m *PeerManager) GetHeight(peerID types.NodeID) int64 {
|
||||
m.mtx.Lock()
|
||||
defer m.mtx.Unlock()
|
||||
|
||||
peer, _ := m.store.Get(peerID)
|
||||
return peer.Height
|
||||
}
|
||||
|
||||
// SetHeight stores a peer's height, making it available via GetHeight.
|
||||
//
|
||||
// FIXME: This is a temporary workaround to share state between the consensus
|
||||
// and mempool reactors, carried over from the legacy P2P stack. Reactors should
|
||||
// not have dependencies on each other, instead tracking this themselves.
|
||||
func (m *PeerManager) SetHeight(peerID types.NodeID, height int64) error {
|
||||
m.mtx.Lock()
|
||||
defer m.mtx.Unlock()
|
||||
|
||||
peer, ok := m.store.Get(peerID)
|
||||
if !ok {
|
||||
peer = m.newPeerInfo(peerID)
|
||||
}
|
||||
peer.Height = height
|
||||
return m.store.Set(peer)
|
||||
}
|
||||
|
||||
// peerStore stores information about peers. It is not thread-safe, assuming it
|
||||
// is only used by PeerManager which handles concurrency control. This allows
|
||||
// the manager to execute multiple operations atomically via its own mutex.
|
||||
|
||||
@@ -1868,38 +1868,3 @@ func TestPeerManager_Advertise_Self(t *testing.T) {
|
||||
self,
|
||||
}, peerManager.Advertise(dID, 100))
|
||||
}
|
||||
|
||||
func TestPeerManager_SetHeight_GetHeight(t *testing.T) {
|
||||
a := p2p.NodeAddress{Protocol: "memory", NodeID: types.NodeID(strings.Repeat("a", 40))}
|
||||
b := p2p.NodeAddress{Protocol: "memory", NodeID: types.NodeID(strings.Repeat("b", 40))}
|
||||
|
||||
db := dbm.NewMemDB()
|
||||
peerManager, err := p2p.NewPeerManager(selfID, db, p2p.PeerManagerOptions{})
|
||||
require.NoError(t, err)
|
||||
|
||||
// Getting a height should default to 0, for unknown peers and
|
||||
// for known peers without height.
|
||||
added, err := peerManager.Add(a)
|
||||
require.NoError(t, err)
|
||||
require.True(t, added)
|
||||
require.EqualValues(t, 0, peerManager.GetHeight(a.NodeID))
|
||||
require.EqualValues(t, 0, peerManager.GetHeight(b.NodeID))
|
||||
|
||||
// Setting a height should work for a known node.
|
||||
require.NoError(t, peerManager.SetHeight(a.NodeID, 3))
|
||||
require.EqualValues(t, 3, peerManager.GetHeight(a.NodeID))
|
||||
|
||||
// Setting a height should add an unknown node.
|
||||
require.Equal(t, []types.NodeID{a.NodeID}, peerManager.Peers())
|
||||
require.NoError(t, peerManager.SetHeight(b.NodeID, 7))
|
||||
require.EqualValues(t, 7, peerManager.GetHeight(b.NodeID))
|
||||
require.ElementsMatch(t, []types.NodeID{a.NodeID, b.NodeID}, peerManager.Peers())
|
||||
|
||||
// The heights should not be persisted.
|
||||
peerManager, err = p2p.NewPeerManager(selfID, db, p2p.PeerManagerOptions{})
|
||||
require.NoError(t, err)
|
||||
|
||||
require.ElementsMatch(t, []types.NodeID{a.NodeID, b.NodeID}, peerManager.Peers())
|
||||
require.Zero(t, peerManager.GetHeight(a.NodeID))
|
||||
require.Zero(t, peerManager.GetHeight(b.NodeID))
|
||||
}
|
||||
|
||||
@@ -236,7 +236,8 @@ func (env *Environment) StartService(ctx context.Context, conf *config.Config) (
|
||||
// If necessary adjust global WriteTimeout to ensure it's greater than
|
||||
// TimeoutBroadcastTxCommit.
|
||||
// See https://github.com/tendermint/tendermint/issues/3435
|
||||
if cfg.WriteTimeout <= conf.RPC.TimeoutBroadcastTxCommit {
|
||||
// Note we don't need to adjust anything if the timeout is already unlimited.
|
||||
if cfg.WriteTimeout > 0 && cfg.WriteTimeout <= conf.RPC.TimeoutBroadcastTxCommit {
|
||||
cfg.WriteTimeout = conf.RPC.TimeoutBroadcastTxCommit + 1*time.Second
|
||||
}
|
||||
|
||||
|
||||
@@ -28,7 +28,7 @@ func NewRoutesMap(svc RPCService, opts *RouteOptions) RoutesMap {
|
||||
out := RoutesMap{
|
||||
// Event subscription. Note that subscribe, unsubscribe, and
|
||||
// unsubscribe_all are only available via the websocket endpoint.
|
||||
"events": rpc.NewRPCFunc(svc.Events),
|
||||
"events": rpc.NewRPCFunc(svc.Events).Timeout(0),
|
||||
"subscribe": rpc.NewWSRPCFunc(svc.Subscribe),
|
||||
"unsubscribe": rpc.NewWSRPCFunc(svc.Unsubscribe),
|
||||
"unsubscribe_all": rpc.NewWSRPCFunc(svc.UnsubscribeAll),
|
||||
|
||||
+1
-1
@@ -266,7 +266,7 @@ func makeNode(
|
||||
node.evPool = evPool
|
||||
|
||||
mpReactor, mp := createMempoolReactor(logger, cfg, proxyApp, stateStore, nodeMetrics.mempool,
|
||||
peerManager.Subscribe, node.router.OpenChannel, peerManager.GetHeight)
|
||||
peerManager.Subscribe, node.router.OpenChannel)
|
||||
node.rpcEnv.Mempool = mp
|
||||
node.services = append(node.services, mpReactor)
|
||||
|
||||
|
||||
@@ -147,7 +147,6 @@ func createMempoolReactor(
|
||||
memplMetrics *mempool.Metrics,
|
||||
peerEvents p2p.PeerEventSubscriber,
|
||||
chCreator p2p.ChannelCreator,
|
||||
peerHeight func(types.NodeID) int64,
|
||||
) (service.Service, mempool.Mempool) {
|
||||
logger = logger.With("module", "mempool")
|
||||
|
||||
@@ -166,7 +165,6 @@ func createMempoolReactor(
|
||||
mp,
|
||||
chCreator,
|
||||
peerEvents,
|
||||
peerHeight,
|
||||
)
|
||||
|
||||
if cfg.Consensus.WaitForTxs() {
|
||||
|
||||
@@ -20,16 +20,27 @@ import (
|
||||
|
||||
// Config is a RPC server configuration.
|
||||
type Config struct {
|
||||
// see netutil.LimitListener
|
||||
// The maximum number of connections that will be accepted by the listener.
|
||||
// See https://godoc.org/golang.org/x/net/netutil#LimitListener
|
||||
MaxOpenConnections int
|
||||
// mirrors http.Server#ReadTimeout
|
||||
|
||||
// Used to set the HTTP server's per-request read timeout.
|
||||
// See https://godoc.org/net/http#Server.ReadTimeout
|
||||
ReadTimeout time.Duration
|
||||
// mirrors http.Server#WriteTimeout
|
||||
|
||||
// Used to set the HTTP server's per-request write timeout. Note that this
|
||||
// affects ALL methods on the server, so it should not be set too low. This
|
||||
// should be used as a safety valve, not a resource-control timeout.
|
||||
//
|
||||
// See https://godoc.org/net/http#Server.WriteTimeout
|
||||
WriteTimeout time.Duration
|
||||
// MaxBodyBytes controls the maximum number of bytes the
|
||||
// server will read parsing the request body.
|
||||
|
||||
// Controls the maximum number of bytes the server will read parsing the
|
||||
// request body.
|
||||
MaxBodyBytes int64
|
||||
// mirrors http.Server#MaxHeaderBytes
|
||||
|
||||
// Controls the maximum size of a request header.
|
||||
// See https://godoc.org/net/http#Server.MaxHeaderBytes
|
||||
MaxHeaderBytes int
|
||||
}
|
||||
|
||||
@@ -38,9 +49,9 @@ func DefaultConfig() *Config {
|
||||
return &Config{
|
||||
MaxOpenConnections: 0, // unlimited
|
||||
ReadTimeout: 10 * time.Second,
|
||||
WriteTimeout: 10 * time.Second,
|
||||
MaxBodyBytes: int64(1000000), // 1MB
|
||||
MaxHeaderBytes: 1 << 20, // same as the net/http default
|
||||
WriteTimeout: 0, // no default timeout
|
||||
MaxBodyBytes: 1000000, // 1MB
|
||||
MaxHeaderBytes: 1 << 20, // same as the net/http default
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -9,11 +9,16 @@ import (
|
||||
"net/http"
|
||||
"reflect"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/tendermint/tendermint/libs/log"
|
||||
rpctypes "github.com/tendermint/tendermint/rpc/jsonrpc/types"
|
||||
)
|
||||
|
||||
// DefaultRPCTimeout is the default context timeout for calls to any RPC method
|
||||
// that does not override it with a more specific timeout.
|
||||
const DefaultRPCTimeout = 60 * time.Second
|
||||
|
||||
// RegisterRPCFuncs adds a route to mux for each non-websocket function in the
|
||||
// funcMap, and also a root JSON-RPC POST handler.
|
||||
func RegisterRPCFuncs(mux *http.ServeMux, funcMap map[string]*RPCFunc, logger log.Logger) {
|
||||
@@ -32,11 +37,12 @@ func RegisterRPCFuncs(mux *http.ServeMux, funcMap map[string]*RPCFunc, logger lo
|
||||
|
||||
// RPCFunc contains the introspected type information for a function.
|
||||
type RPCFunc struct {
|
||||
f reflect.Value // underlying rpc function
|
||||
param reflect.Type // the parameter struct, or nil
|
||||
result reflect.Type // the non-error result type, or nil
|
||||
args []argInfo // names and type information (for URL decoding)
|
||||
ws bool // websocket only
|
||||
f reflect.Value // underlying rpc function
|
||||
param reflect.Type // the parameter struct, or nil
|
||||
result reflect.Type // the non-error result type, or nil
|
||||
args []argInfo // names and type information (for URL decoding)
|
||||
timeout time.Duration // default request timeout, 0 means none
|
||||
ws bool // websocket only
|
||||
}
|
||||
|
||||
// argInfo records the name of a field, along with a bit to tell whether the
|
||||
@@ -52,6 +58,12 @@ type argInfo struct {
|
||||
// with the resulting argument value. It reports an error if parameter parsing
|
||||
// fails, otherwise it returns the result from the wrapped function.
|
||||
func (rf *RPCFunc) Call(ctx context.Context, params json.RawMessage) (interface{}, error) {
|
||||
// If ctx has its own deadline we will respect it; otherwise use rf.timeout.
|
||||
if _, ok := ctx.Deadline(); !ok && rf.timeout > 0 {
|
||||
var cancel context.CancelFunc
|
||||
ctx, cancel = context.WithTimeout(ctx, rf.timeout)
|
||||
defer cancel()
|
||||
}
|
||||
args, err := rf.parseParams(ctx, params)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -74,6 +86,11 @@ func (rf *RPCFunc) Call(ctx context.Context, params json.RawMessage) (interface{
|
||||
return returns[0].Interface(), nil
|
||||
}
|
||||
|
||||
// Timeout updates rf to include a default timeout for calls to rf. This
|
||||
// timeout is used if one is not already provided on the request context.
|
||||
// Setting d == 0 means there will be no timeout. Returns rf to allow chaining.
|
||||
func (rf *RPCFunc) Timeout(d time.Duration) *RPCFunc { rf.timeout = d; return rf }
|
||||
|
||||
// parseParams parses the parameters of a JSON-RPC request and returns the
|
||||
// corresponding argument values. On success, the first argument value will be
|
||||
// the value of ctx.
|
||||
@@ -129,7 +146,9 @@ func (rf *RPCFunc) adjustParams(data []byte) (json.RawMessage, error) {
|
||||
// func(context.Context, *T) (R, error)
|
||||
//
|
||||
// for an arbitrary struct type T and type R. NewRPCFunc will panic if f does
|
||||
// not have one of these forms.
|
||||
// not have one of these forms. A newly-constructed RPCFunc has a default
|
||||
// timeout of DefaultRPCTimeout; use the Timeout method to adjust this as
|
||||
// needed.
|
||||
func NewRPCFunc(f interface{}) *RPCFunc {
|
||||
rf, err := newRPCFunc(f)
|
||||
if err != nil {
|
||||
@@ -215,10 +234,11 @@ func newRPCFunc(f interface{}) (*RPCFunc, error) {
|
||||
}
|
||||
|
||||
return &RPCFunc{
|
||||
f: fv,
|
||||
param: ptype,
|
||||
result: rtype,
|
||||
args: args,
|
||||
f: fv,
|
||||
param: ptype,
|
||||
result: rtype,
|
||||
args: args,
|
||||
timeout: DefaultRPCTimeout, // until overridden
|
||||
}, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -210,7 +210,8 @@ func startLightNode(ctx context.Context, logger log.Logger, cfg *Config) error {
|
||||
// If necessary adjust global WriteTimeout to ensure it's greater than
|
||||
// TimeoutBroadcastTxCommit.
|
||||
// See https://github.com/tendermint/tendermint/issues/3435
|
||||
if rpccfg.WriteTimeout <= tmcfg.RPC.TimeoutBroadcastTxCommit {
|
||||
// Note we don't need to adjust anything if the timeout is already unlimited.
|
||||
if rpccfg.WriteTimeout > 0 && rpccfg.WriteTimeout <= tmcfg.RPC.TimeoutBroadcastTxCommit {
|
||||
rpccfg.WriteTimeout = tmcfg.RPC.TimeoutBroadcastTxCommit + 1*time.Second
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user