From e3c7b263f17f93ee3b0dcffbffce88b6be701d1e Mon Sep 17 00:00:00 2001 From: tycho garen Date: Fri, 24 Sep 2021 16:08:33 -0400 Subject: [PATCH] retry params fetch --- internal/statesync/reactor.go | 2 +- internal/statesync/reactor_test.go | 2 +- internal/statesync/stateprovider.go | 83 +++++++++++++++++------------ 3 files changed, 50 insertions(+), 37 deletions(-) diff --git a/internal/statesync/reactor.go b/internal/statesync/reactor.go index 58f4033d8..c66b77793 100644 --- a/internal/statesync/reactor.go +++ b/internal/statesync/reactor.go @@ -1004,7 +1004,7 @@ func (r *Reactor) fetchLightBlock(height uint64) (*types.LightBlock, error) { func (r *Reactor) waitForEnoughPeers(ctx context.Context, numPeers int) error { startAt := time.Now() - t := time.NewTicker(200 * time.Millisecond) + t := time.NewTicker(100 * time.Millisecond) defer t.Stop() for r.peers.Len() < numPeers { select { diff --git a/internal/statesync/reactor_test.go b/internal/statesync/reactor_test.go index 41dcf3d2d..22a56ce40 100644 --- a/internal/statesync/reactor_test.go +++ b/internal/statesync/reactor_test.go @@ -525,7 +525,7 @@ func TestReactor_StateProviderP2P(t *testing.T) { rts.reactor.cfg.UseP2P = true rts.reactor.cfg.TrustHeight = 1 rts.reactor.cfg.TrustHash = fmt.Sprintf("%X", chain[1].Hash()) - ctx, cancel := context.WithCancel(context.Background()) + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) defer cancel() rts.reactor.mtx.Lock() diff --git a/internal/statesync/stateprovider.go b/internal/statesync/stateprovider.go index 9eb776e28..ddbca57c9 100644 --- a/internal/statesync/stateprovider.go +++ b/internal/statesync/stateprovider.go @@ -337,44 +337,57 @@ func (s *stateProviderP2P) addProvider(p lightprovider.Provider) { } } -// consensusParams sends out a request for consensus params blocking until one is returned. -// If it fails to get a valid set of consensus params from any of the providers it returns an error. +// consensusParams sends out a request for consensus params blocking +// until one is returned. +// +// If it fails to get a valid set of consensus params from any of the +// providers it returns an error; however, it will retry indefinitely +// (with backoff) until the context is canceled. func (s *stateProviderP2P) consensusParams(ctx context.Context, height int64) (types.ConsensusParams, error) { - for _, provider := range s.lc.Witnesses() { - p, ok := provider.(*BlockProvider) - if !ok { - panic("expected p2p state provider to use p2p block providers") - } - - // extract the nodeID of the provider - peer, err := types.NewNodeID(p.String()) - if err != nil { - return types.ConsensusParams{}, fmt.Errorf("invalid provider (%s) node id: %w", p.String(), err) - } - - select { - case s.paramsSendCh <- p2p.Envelope{ - To: peer, - Message: &ssproto.ParamsRequest{ - Height: uint64(height), - }, - }: - case <-ctx.Done(): - return types.ConsensusParams{}, ctx.Err() - } - - select { - // if we get no response from this provider we move on to the next one - case <-time.After(consensusParamsResponseTimeout): - continue - case <-ctx.Done(): - return types.ConsensusParams{}, ctx.Err() - case params, ok := <-s.paramsRecvCh: + iterCount := 0 + for { + for _, provider := range s.lc.Witnesses() { + p, ok := provider.(*BlockProvider) if !ok { - return types.ConsensusParams{}, errors.New("params channel closed") + panic("expected p2p state provider to use p2p block providers") } - return params, nil + + // extract the nodeID of the provider + peer, err := types.NewNodeID(p.String()) + if err != nil { + return types.ConsensusParams{}, fmt.Errorf("invalid provider (%s) node id: %w", p.String(), err) + } + + select { + case s.paramsSendCh <- p2p.Envelope{ + To: peer, + Message: &ssproto.ParamsRequest{ + Height: uint64(height), + }, + }: + case <-ctx.Done(): + return types.ConsensusParams{}, ctx.Err() + } + + select { + // if we get no response from this provider we move on to the next one + case <-time.After(consensusParamsResponseTimeout): + continue + case <-ctx.Done(): + return types.ConsensusParams{}, ctx.Err() + case params, ok := <-s.paramsRecvCh: + if !ok { + return types.ConsensusParams{}, errors.New("params channel closed") + } + return params, nil + } + } + iterCount++ + + select { + case <-ctx.Done(): + return types.ConsensusParams{}, ctx.Err() + case <-time.After(time.Duration(iterCount) * consensusParamsResponseTimeout): } } - return types.ConsensusParams{}, errors.New("unable to fetch consensus params from connected providers") }