diff --git a/cmd/tendermint/commands/inspect.go b/cmd/tendermint/commands/inspect.go index 53106ba44..3cf33630b 100644 --- a/cmd/tendermint/commands/inspect.go +++ b/cmd/tendermint/commands/inspect.go @@ -38,13 +38,14 @@ func init() { } func runInspect(cmd *cobra.Command, args []string) error { - ctx, cancelFunc := context.WithCancel(context.Background()) + ctx, cancel := context.WithCancel(cmd.Context()) + defer cancel() c := make(chan os.Signal, 1) - signal.Notify(c, os.Interrupt, syscall.SIGTERM, syscall.SIGINT) + signal.Notify(c, syscall.SIGTERM, syscall.SIGINT) go func() { <-c - cancelFunc() + cancel() }() blockStoreDB, err := cfg.DefaultDBProvider(&cfg.DBContext{ID: "blockstore", Config: config}) @@ -66,10 +67,10 @@ func runInspect(cmd *cobra.Command, args []string) error { } stateStore := state.NewStore(stateDB) - d := inspect.New(config.RPC, blockStore, stateStore, sinks, logger) + ins := inspect.New(config.RPC, blockStore, stateStore, sinks, logger) logger.Info("starting inspect server") - if err := d.Run(ctx); err != nil { + if err := ins.Run(ctx); err != nil { logger.Error("error encountered while running inspect server", "err", err) return err } diff --git a/docs/tools/debugging/README.md b/docs/tools/debugging/README.md index 9c94f3c36..c2677db45 100644 --- a/docs/tools/debugging/README.md +++ b/docs/tools/debugging/README.md @@ -73,6 +73,7 @@ entire Tendermint process. While in this inconsistent state, a node running Tendermint's consensus engine will not start up. The `inspect` command runs only a subset of Tendermint's RPC endpoints for querying the block store and state store. +`inspect` allows operators to query a read-only view of the stage. `inspect` does not run the consensus engine at all and can therefore be used to debug processes that have crashed due to inconsistent state. diff --git a/go.mod b/go.mod index 84c6f43ac..d6ced2c95 100644 --- a/go.mod +++ b/go.mod @@ -37,6 +37,7 @@ require ( github.com/vektra/mockery/v2 v2.9.0 golang.org/x/crypto v0.0.0-20210314154223-e6e6c4f2bb5b golang.org/x/net v0.0.0-20210405180319-a5a99cb37ef4 + golang.org/x/sync v0.0.0-20210220032951-036812b2e83c // indirect google.golang.org/grpc v1.40.0 gopkg.in/check.v1 v1.0.0-20200902074654-038fdea0a05b // indirect pgregory.net/rapid v0.4.7 diff --git a/go.sum b/go.sum index 81abd8a00..4f314aefd 100644 --- a/go.sum +++ b/go.sum @@ -1071,6 +1071,7 @@ golang.org/x/sync v0.0.0-20200317015054-43a5402ce75a/go.mod h1:RxMgew5VJxzue5/jJ golang.org/x/sync v0.0.0-20200625203802-6e8e738ad208/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20201207232520-09787c993a3a/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20210220032951-036812b2e83c h1:5KslGYwFpkhGh+Q16bwMP3cOontH8FOep7tGV86Y7SQ= golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sys v0.0.0-20180823144017-11551d06cbcc/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= diff --git a/inspect/inspect.go b/inspect/inspect.go index 0cba9d27a..57b10c092 100644 --- a/inspect/inspect.go +++ b/inspect/inspect.go @@ -4,18 +4,20 @@ import ( "context" "errors" "net" - "sync" - cfg "github.com/tendermint/tendermint/config" - inspect_rpc "github.com/tendermint/tendermint/inspect/rpc" + "github.com/tendermint/tendermint/config" + "github.com/tendermint/tendermint/inspect/rpc" "github.com/tendermint/tendermint/libs/log" + "github.com/tendermint/tendermint/libs/pubsub" tmstrings "github.com/tendermint/tendermint/libs/strings" rpccore "github.com/tendermint/tendermint/rpc/core" - sm "github.com/tendermint/tendermint/state" + "github.com/tendermint/tendermint/state" "github.com/tendermint/tendermint/state/indexer" "github.com/tendermint/tendermint/state/indexer/sink" "github.com/tendermint/tendermint/store" "github.com/tendermint/tendermint/types" + + "golang.org/x/sync/errgroup" ) // Inspect manages an RPC service that exports methods to debug a failed node. @@ -27,148 +29,124 @@ import ( type Inspect struct { routes rpccore.RoutesMap - rpcConfig *cfg.RPCConfig + config *config.RPCConfig indexerService *indexer.Service eventBus *types.EventBus logger log.Logger } -// New constructs a new Inspect from the passed in parameters. +// New returns an Inspect that serves RPC on the given BlockStore and StateStore. /// //nolint:lll -func New(rpcConfig *cfg.RPCConfig, blockStore sm.BlockStore, stateStore sm.Store, eventSinks []indexer.EventSink, logger log.Logger) *Inspect { - routes := inspect_rpc.Routes(stateStore, blockStore, eventSinks) - eventBus := types.NewEventBus() - eventBus.SetLogger(logger.With("module", "events")) - indexerService := indexer.NewIndexerService(eventSinks, eventBus) - indexerService.SetLogger(logger.With("module", "txindex")) +func New(cfg *config.RPCConfig, bs state.BlockStore, ss state.Store, es []indexer.EventSink, logger log.Logger) *Inspect { + routes := rpc.Routes(ss, bs, es) + eb := types.NewEventBus() + eb.SetLogger(logger.With("module", "events")) + is := indexer.NewIndexerService(es, eb) + is.SetLogger(logger.With("module", "txindex")) return &Inspect{ routes: routes, - rpcConfig: rpcConfig, + config: cfg, logger: logger, - eventBus: eventBus, - indexerService: indexerService, + eventBus: eb, + indexerService: is, } } // NewFromConfig constructs an Inspect using the values defined in the passed in config. -func NewFromConfig(config *cfg.Config) (*Inspect, error) { - blockStoreDB, err := cfg.DefaultDBProvider(&cfg.DBContext{ID: "blockstore", Config: config}) +func NewFromConfig(cfg *config.Config) (*Inspect, error) { + bsDB, err := config.DefaultDBProvider(&config.DBContext{ID: "blockstore", Config: cfg}) if err != nil { return nil, err } - blockStore := store.NewBlockStore(blockStoreDB) - stateDB, err := cfg.DefaultDBProvider(&cfg.DBContext{ID: "statestore", Config: config}) + bs := store.NewBlockStore(bsDB) + sDB, err := config.DefaultDBProvider(&config.DBContext{ID: "statestore", Config: cfg}) if err != nil { return nil, err } - genDoc, err := types.GenesisDocFromFile(config.GenesisFile()) + genDoc, err := types.GenesisDocFromFile(cfg.GenesisFile()) if err != nil { return nil, err } - sinks, err := sink.EventSinksFromConfig(config, cfg.DefaultDBProvider, genDoc.ChainID) + sinks, err := sink.EventSinksFromConfig(cfg, config.DefaultDBProvider, genDoc.ChainID) if err != nil { return nil, err } - l := log.MustNewDefaultLogger(log.LogFormatPlain, log.LogLevelInfo, false) - stateStore := sm.NewStore(stateDB) - return New(config.RPC, blockStore, stateStore, sinks, l), nil -} - -// NewDefault constructs a new Inspect using the default values. -func NewDefault() (*Inspect, error) { - config := cfg.Config{ - BaseConfig: cfg.DefaultBaseConfig(), - RPC: cfg.DefaultRPCConfig(), - TxIndex: cfg.DefaultTxIndexConfig(), - } - return NewFromConfig(&config) + logger := log.MustNewDefaultLogger(log.LogFormatPlain, log.LogLevelInfo, false) + ss := state.NewStore(sDB) + return New(cfg.RPC, bs, ss, sinks, logger), nil } // Run starts the Inspect servers and blocks until the servers shut down. The passed // in context is used to control the lifecycle of the servers. -func (inspect *Inspect) Run(ctx context.Context) error { - err := inspect.eventBus.Start() +func (ins *Inspect) Run(ctx context.Context) error { + err := ins.eventBus.Start() if err != nil { return err } defer func() { - err := inspect.eventBus.Stop() + err := ins.eventBus.Stop() if err != nil { - inspect.logger.Error("event bus stopped with error", "err", err) + ins.logger.Error("event bus stopped with error", "err", err) } }() + g, tctx := errgroup.WithContext(ctx) + g.Go(func() error { + if !errors.Is(ins.indexerService.Start(), pubsub.ErrUnsubscribed) { + return err + } + return nil + }) + g.Go(func() error { + return startRPCServers(ctx, ins.config, ins.logger, ins.routes) + }) - err = inspect.indexerService.Start() + <-tctx.Done() + err = ins.indexerService.Stop() if err != nil { - return err + ins.logger.Error("event bus stopped with error", "err", err) } - defer func() { - err := inspect.indexerService.Stop() - if err != nil { - inspect.logger.Error("indexer stopped with error", "err", err) - } - }() - return startRPCServers(ctx, inspect.rpcConfig, inspect.logger, inspect.routes) + return g.Wait() } -func startRPCServers(ctx context.Context, rpcConfig *cfg.RPCConfig, logger log.Logger, routes rpccore.RoutesMap) error { - wg := &sync.WaitGroup{} - listenAddrs := tmstrings.SplitAndTrimEmpty(rpcConfig.ListenAddress, ",", " ") - rootHandler := inspect_rpc.Handler(rpcConfig, routes, logger) - errChan := make(chan error) +func startRPCServers(ctx context.Context, cfg *config.RPCConfig, logger log.Logger, routes rpccore.RoutesMap) error { + g, _ := errgroup.WithContext(ctx) + listenAddrs := tmstrings.SplitAndTrimEmpty(cfg.ListenAddress, ",", " ") + rh := rpc.Handler(cfg, routes, logger) for _, listenerAddr := range listenAddrs { - server := inspect_rpc.Server{ + server := rpc.Server{ Logger: logger, - Config: rpcConfig, - Handler: rootHandler, + Config: cfg, + Handler: rh, Addr: listenerAddr, } - if rpcConfig.IsTLSEnabled() { - keyFile := rpcConfig.KeyFile() - certFile := rpcConfig.CertFile() - wg.Add(1) - go func() { - defer wg.Done() + if cfg.IsTLSEnabled() { + keyFile := cfg.KeyFile() + certFile := cfg.CertFile() + g.Go(func() error { logger.Info("RPC HTTPS server starting", "address", listenerAddr, "certfile", certFile, "keyfile", keyFile) err := server.ListenAndServeTLS(ctx, certFile, keyFile) if !errors.Is(err, net.ErrClosed) { logger.Error("RPC HTTPS server stopped with error", "address", listenerAddr, "err", err) - errChan <- err - return + return err } logger.Info("RPC HTTPS server stopped", "address", listenerAddr) - }() + return nil + }) } else { - wg.Add(1) - go func() { - defer wg.Done() + g.Go(func() error { logger.Info("RPC HTTP server starting", "address", listenerAddr) err := server.ListenAndServe(ctx) if !errors.Is(err, net.ErrClosed) { logger.Error("RPC HTTP server stopped with error", "address", listenerAddr, "err", err) - errChan <- err - return + return err } logger.Info("RPC HTTP server stopped", "address", listenerAddr) - }() + return nil + }) } } - select { - case <-chanFromWG(wg): - return nil - case err := <-errChan: - return err - } -} - -func chanFromWG(wg *sync.WaitGroup) chan struct{} { - ch := make(chan struct{}) - go func() { - wg.Wait() - close(ch) - }() - return ch + return g.Wait() } diff --git a/inspect/rpc/rpc.go b/inspect/rpc/rpc.go index d9348fc88..2066f8bcd 100644 --- a/inspect/rpc/rpc.go +++ b/inspect/rpc/rpc.go @@ -6,11 +6,12 @@ import ( "time" "github.com/rs/cors" + "github.com/tendermint/tendermint/config" "github.com/tendermint/tendermint/libs/log" - tmpubsub "github.com/tendermint/tendermint/libs/pubsub" - rpccore "github.com/tendermint/tendermint/rpc/core" - rpcserver "github.com/tendermint/tendermint/rpc/jsonrpc/server" + "github.com/tendermint/tendermint/libs/pubsub" + "github.com/tendermint/tendermint/rpc/core" + "github.com/tendermint/tendermint/rpc/jsonrpc/server" "github.com/tendermint/tendermint/state" "github.com/tendermint/tendermint/state/indexer" "github.com/tendermint/tendermint/types" @@ -25,30 +26,30 @@ type Server struct { } // Routes returns the set of routes used by the Inspect server. -func Routes(store state.Store, blockStore state.BlockStore, eventSinks []indexer.EventSink) rpccore.RoutesMap { - env := &rpccore.Environment{ - EventSinks: eventSinks, - StateStore: store, - BlockStore: blockStore, +func Routes(s state.Store, bs state.BlockStore, es []indexer.EventSink) core.RoutesMap { + env := &core.Environment{ + EventSinks: es, + StateStore: s, + BlockStore: bs, } - return rpccore.RoutesMap{ - "blockchain": rpcserver.NewRPCFunc(env.BlockchainInfo, "minHeight,maxHeight", true), - "consensus_params": rpcserver.NewRPCFunc(env.ConsensusParams, "height", true), - "block": rpcserver.NewRPCFunc(env.Block, "height", true), - "block_by_hash": rpcserver.NewRPCFunc(env.BlockByHash, "hash", true), - "block_results": rpcserver.NewRPCFunc(env.BlockResults, "height", true), - "commit": rpcserver.NewRPCFunc(env.Commit, "height", true), - "validators": rpcserver.NewRPCFunc(env.Validators, "height,page,per_page", true), - "tx": rpcserver.NewRPCFunc(env.Tx, "hash,prove", true), - "tx_search": rpcserver.NewRPCFunc(env.TxSearch, "query,prove,page,per_page,order_by", false), - "block_search": rpcserver.NewRPCFunc(env.BlockSearch, "query,page,per_page,order_by", false), + return core.RoutesMap{ + "blockchain": server.NewRPCFunc(env.BlockchainInfo, "minHeight,maxHeight", true), + "consensus_params": server.NewRPCFunc(env.ConsensusParams, "height", true), + "block": server.NewRPCFunc(env.Block, "height", true), + "block_by_hash": server.NewRPCFunc(env.BlockByHash, "hash", true), + "block_results": server.NewRPCFunc(env.BlockResults, "height", true), + "commit": server.NewRPCFunc(env.Commit, "height", true), + "validators": server.NewRPCFunc(env.Validators, "height,page,per_page", true), + "tx": server.NewRPCFunc(env.Tx, "hash,prove", true), + "tx_search": server.NewRPCFunc(env.TxSearch, "query,prove,page,per_page,order_by", false), + "block_search": server.NewRPCFunc(env.BlockSearch, "query,page,per_page,order_by", false), } } // Handler returns the http.Handler configured for use with an Inspect server. Handler // registers the routes on the http.Handler and also registers the websocket handler // and the CORS handler if specified by the configuration options. -func Handler(rpcConfig *config.RPCConfig, routes rpccore.RoutesMap, logger log.Logger) http.Handler { +func Handler(rpcConfig *config.RPCConfig, routes core.RoutesMap, logger log.Logger) http.Handler { mux := http.NewServeMux() wmLogger := logger.With("protocol", "websocket") @@ -56,17 +57,17 @@ func Handler(rpcConfig *config.RPCConfig, routes rpccore.RoutesMap, logger log.L websocketDisconnectFn := func(remoteAddr string) { err := eventBus.UnsubscribeAll(context.Background(), remoteAddr) - if err != nil && err != tmpubsub.ErrSubscriptionNotFound { + if err != nil && err != pubsub.ErrSubscriptionNotFound { wmLogger.Error("Failed to unsubscribe addr from events", "addr", remoteAddr, "err", err) } } - wm := rpcserver.NewWebsocketManager(routes, - rpcserver.OnDisconnect(websocketDisconnectFn), - rpcserver.ReadLimit(rpcConfig.MaxBodyBytes)) + wm := server.NewWebsocketManager(routes, + server.OnDisconnect(websocketDisconnectFn), + server.ReadLimit(rpcConfig.MaxBodyBytes)) wm.SetLogger(wmLogger) mux.HandleFunc("/websocket", wm.WebsocketHandler) - rpcserver.RegisterRPCFuncs(mux, routes, logger) + server.RegisterRPCFuncs(mux, routes, logger) var rootHandler http.Handler = mux if rpcConfig.IsCorsEnabled() { rootHandler = addCORSHandler(rpcConfig, mux) @@ -87,7 +88,7 @@ func addCORSHandler(rpcConfig *config.RPCConfig, h http.Handler) http.Handler { // ListenAndServe listens on the address specified in srv.Addr and handles any // incoming requests over HTTP using the Inspect rpc handler specified on the server. func (srv *Server) ListenAndServe(ctx context.Context) error { - listener, err := rpcserver.Listen(srv.Addr, srv.Config.MaxOpenConnections) + listener, err := server.Listen(srv.Addr, srv.Config.MaxOpenConnections) if err != nil { return err } @@ -95,13 +96,13 @@ func (srv *Server) ListenAndServe(ctx context.Context) error { <-ctx.Done() listener.Close() }() - return rpcserver.Serve(listener, srv.Handler, srv.Logger, serverRPCConfig(srv.Config)) + return server.Serve(listener, srv.Handler, srv.Logger, serverRPCConfig(srv.Config)) } // ListenAndServeTLS listens on the address specified in srv.Addr. ListenAndServeTLS handles // incoming requests over HTTPS using the Inspect rpc handler specified on the server. func (srv *Server) ListenAndServeTLS(ctx context.Context, certFile, keyFile string) error { - listener, err := rpcserver.Listen(srv.Addr, srv.Config.MaxOpenConnections) + listener, err := server.Listen(srv.Addr, srv.Config.MaxOpenConnections) if err != nil { return err } @@ -109,11 +110,11 @@ func (srv *Server) ListenAndServeTLS(ctx context.Context, certFile, keyFile stri <-ctx.Done() listener.Close() }() - return rpcserver.ServeTLS(listener, srv.Handler, certFile, keyFile, srv.Logger, serverRPCConfig(srv.Config)) + return server.ServeTLS(listener, srv.Handler, certFile, keyFile, srv.Logger, serverRPCConfig(srv.Config)) } -func serverRPCConfig(r *config.RPCConfig) *rpcserver.Config { - cfg := rpcserver.DefaultConfig() +func serverRPCConfig(r *config.RPCConfig) *server.Config { + cfg := server.DefaultConfig() cfg.MaxBodyBytes = r.MaxBodyBytes cfg.MaxHeaderBytes = r.MaxHeaderBytes // If necessary adjust global WriteTimeout to ensure it's greater than diff --git a/node/node.go b/node/node.go index b6ea1e974..4d57f7bab 100644 --- a/node/node.go +++ b/node/node.go @@ -24,6 +24,7 @@ import ( "github.com/tendermint/tendermint/internal/statesync" "github.com/tendermint/tendermint/libs/log" tmnet "github.com/tendermint/tendermint/libs/net" + "github.com/tendermint/tendermint/libs/pubsub" tmpubsub "github.com/tendermint/tendermint/libs/pubsub" "github.com/tendermint/tendermint/libs/service" "github.com/tendermint/tendermint/libs/strings" @@ -161,11 +162,17 @@ func makeNode(config *cfg.Config, return nil, err } - indexerService, eventSinks, err := createAndStartIndexerService(config, dbProvider, eventBus, logger, genDoc.ChainID) + indexerService, eventSinks, err := createIndexerService(config, dbProvider, eventBus, logger, genDoc.ChainID) if err != nil { return nil, err } + go func() { + if !errors.Is(ins.indexerService.Start(), pubsub.ErrUnsubscribed) { + return err + } + }() + // If an address is provided, listen on the socket for a connection from an // external signing process. if config.PrivValidator.ListenAddr != "" { diff --git a/node/setup.go b/node/setup.go index ae706ee22..da11f85d1 100644 --- a/node/setup.go +++ b/node/setup.go @@ -68,7 +68,7 @@ func createAndStartEventBus(logger log.Logger) (*types.EventBus, error) { return eventBus, nil } -func createAndStartIndexerService( +func createIndexerService( config *cfg.Config, dbProvider cfg.DBProvider, eventBus *types.EventBus, @@ -81,6 +81,7 @@ func createAndStartIndexerService( } indexerService := indexer.NewIndexerService(eventSinks, eventBus) indexerService.SetLogger(logger.With("module", "txindex")) + return indexerService, eventSinks, nil if err := indexerService.Start(); err != nil { return nil, nil, err diff --git a/state/indexer/indexer_service.go b/state/indexer/indexer_service.go index 01d72bf9a..c1d006dce 100644 --- a/state/indexer/indexer_service.go +++ b/state/indexer/indexer_service.go @@ -52,55 +52,52 @@ func (is *Service) OnStart() error { return err } - go func() { - for { - select { - case <-is.doneChan: - return - case msg := <-blockHeadersSub.Out(): + for { + select { + case <-blockHeadersSub.Canceled(): + return blockHeadersSub.Err() + case msg := <-blockHeadersSub.Out(): - eventDataHeader := msg.Data().(types.EventDataNewBlockHeader) - height := eventDataHeader.Header.Height - batch := NewBatch(eventDataHeader.NumTxs) + eventDataHeader := msg.Data().(types.EventDataNewBlockHeader) + height := eventDataHeader.Header.Height + batch := NewBatch(eventDataHeader.NumTxs) - for i := int64(0); i < eventDataHeader.NumTxs; i++ { - msg2 := <-txsSub.Out() - txResult := msg2.Data().(types.EventDataTx).TxResult + for i := int64(0); i < eventDataHeader.NumTxs; i++ { + msg2 := <-txsSub.Out() + txResult := msg2.Data().(types.EventDataTx).TxResult - if err = batch.Add(&txResult); err != nil { - is.Logger.Error( - "failed to add tx to batch", - "height", height, - "index", txResult.Index, - "err", err, - ) - } + if err = batch.Add(&txResult); err != nil { + is.Logger.Error( + "failed to add tx to batch", + "height", height, + "index", txResult.Index, + "err", err, + ) + } + } + + if !IndexingEnabled(is.eventSinks) { + continue + } + + for _, sink := range is.eventSinks { + if err := sink.IndexBlockEvents(eventDataHeader); err != nil { + is.Logger.Error("failed to index block", "height", height, "err", err) + } else { + is.Logger.Debug("indexed block", "height", height, "sink", sink.Type()) } - if !IndexingEnabled(is.eventSinks) { - continue - } - - for _, sink := range is.eventSinks { - if err := sink.IndexBlockEvents(eventDataHeader); err != nil { - is.Logger.Error("failed to index block", "height", height, "err", err) + if len(batch.Ops) > 0 { + err := sink.IndexTxEvents(batch.Ops) + if err != nil { + is.Logger.Error("failed to index block txs", "height", height, "err", err) } else { - is.Logger.Debug("indexed block", "height", height, "sink", sink.Type()) - } - - if len(batch.Ops) > 0 { - err := sink.IndexTxEvents(batch.Ops) - if err != nil { - is.Logger.Error("failed to index block txs", "height", height, "err", err) - } else { - is.Logger.Debug("indexed txs", "height", height, "sink", sink.Type()) - } + is.Logger.Debug("indexed txs", "height", height, "sink", sink.Type()) } } } } - }() - return nil + } } // OnStop implements service.Service by unsubscribing from all transactions and @@ -115,7 +112,6 @@ func (is *Service) OnStop() { is.Logger.Error("failed to close eventsink", "eventsink", sink.Type(), "err", err) } } - close(is.doneChan) } // KVSinkEnabled returns the given eventSinks is containing KVEventSink.