mirror of
https://gitlab.com/pulsechaincom/prysm-pulse.git
synced 2025-01-11 04:00:05 +00:00
918129cf36
* refactor initialization to blocking startup method * require genesisSetter in blockchain, fix tests * work-around gazelle weirdness * fix dep gazelle ignores * only call SetGenesis once * fix typo * validator test setup and fix to return right error * move waitForChainStart to Start * wire up sync Service.genesisWaiter * fix p2p genesisWaiter plumbing * remove extra clock type, integrate into genesis and rename * use time.Now when no Nower is specified * remove unused ClockSetter * simplify rpc context checking * fix typo * use clock everywhere in sync; [32]byte val root * don't use DeepEqual to compare [32]byte and []byte * don't use clock in init sync, not wired up yet * use clock waiter in blockchain as well * use cancelable contexts in tests with goroutines * missed a reference to WithClockSetter * Update beacon-chain/startup/genesis.go Co-authored-by: Radosław Kapka <rkapka@wp.pl> * Update beacon-chain/blockchain/service_test.go Co-authored-by: Radosław Kapka <rkapka@wp.pl> * more clear docs * doc for NewClock * move clock typedef to more logical file name * adding documentation * gaz * fixes for capella * reducing test raciness * fix races in committee cache tests * lint * add tests on Duration slot math helper * startup package test coverage * fix bad merge * set non-zero genesis time in tests that call Start * happy deepsource, happy me-epsource * replace Synced event with channel * remove unused error * remove accidental wip commit * gaz! * remove unused event constants * remove sync statefeed subscription to fix deadlock * remove state notifier * fix build --------- Co-authored-by: Kasey Kirkham <kasey@users.noreply.github.com> Co-authored-by: Radosław Kapka <rkapka@wp.pl> Co-authored-by: prylabs-bulldozer[bot] <58059840+prylabs-bulldozer[bot]@users.noreply.github.com> Co-authored-by: nisdas <nishdas93@gmail.com>
229 lines
6.6 KiB
Go
229 lines
6.6 KiB
Go
package sync
|
|
|
|
import (
|
|
"context"
|
|
"testing"
|
|
"time"
|
|
|
|
gcache "github.com/patrickmn/go-cache"
|
|
"github.com/prysmaticlabs/prysm/v4/async/abool"
|
|
mockChain "github.com/prysmaticlabs/prysm/v4/beacon-chain/blockchain/testing"
|
|
"github.com/prysmaticlabs/prysm/v4/beacon-chain/core/feed"
|
|
dbTest "github.com/prysmaticlabs/prysm/v4/beacon-chain/db/testing"
|
|
p2ptest "github.com/prysmaticlabs/prysm/v4/beacon-chain/p2p/testing"
|
|
"github.com/prysmaticlabs/prysm/v4/beacon-chain/startup"
|
|
state_native "github.com/prysmaticlabs/prysm/v4/beacon-chain/state/state-native"
|
|
mockSync "github.com/prysmaticlabs/prysm/v4/beacon-chain/sync/initial-sync/testing"
|
|
"github.com/prysmaticlabs/prysm/v4/crypto/bls"
|
|
"github.com/prysmaticlabs/prysm/v4/encoding/bytesutil"
|
|
ethpb "github.com/prysmaticlabs/prysm/v4/proto/prysm/v1alpha1"
|
|
"github.com/prysmaticlabs/prysm/v4/testing/assert"
|
|
"github.com/prysmaticlabs/prysm/v4/testing/require"
|
|
"github.com/prysmaticlabs/prysm/v4/testing/util"
|
|
)
|
|
|
|
func TestService_StatusZeroEpoch(t *testing.T) {
|
|
bState, err := state_native.InitializeFromProtoPhase0(ðpb.BeaconState{Slot: 0})
|
|
require.NoError(t, err)
|
|
chain := &mockChain.ChainService{
|
|
Genesis: time.Now(),
|
|
State: bState,
|
|
}
|
|
r := &Service{
|
|
cfg: &config{
|
|
p2p: p2ptest.NewTestP2P(t),
|
|
initialSync: new(mockSync.Sync),
|
|
chain: chain,
|
|
clock: startup.NewClock(chain.Genesis, chain.ValidatorsRoot),
|
|
},
|
|
chainStarted: abool.New(),
|
|
}
|
|
r.chainStarted.Set()
|
|
|
|
assert.NoError(t, r.Status(), "Wanted non failing status")
|
|
}
|
|
|
|
func TestSyncHandlers_WaitToSync(t *testing.T) {
|
|
p2p := p2ptest.NewTestP2P(t)
|
|
chainService := &mockChain.ChainService{
|
|
Genesis: time.Now(),
|
|
ValidatorsRoot: [32]byte{'A'},
|
|
}
|
|
gs := startup.NewClockSynchronizer()
|
|
r := Service{
|
|
ctx: context.Background(),
|
|
cfg: &config{
|
|
p2p: p2p,
|
|
chain: chainService,
|
|
initialSync: &mockSync.Sync{IsSyncing: false},
|
|
},
|
|
chainStarted: abool.New(),
|
|
clockWaiter: gs,
|
|
}
|
|
|
|
topic := "/eth2/%x/beacon_block"
|
|
go r.registerHandlers()
|
|
go r.waitForChainStart()
|
|
time.Sleep(100 * time.Millisecond)
|
|
|
|
var vr [32]byte
|
|
require.NoError(t, gs.SetClock(startup.NewClock(time.Now(), vr)))
|
|
b := []byte("sk")
|
|
b32 := bytesutil.ToBytes32(b)
|
|
sk, err := bls.SecretKeyFromBytes(b32[:])
|
|
require.NoError(t, err)
|
|
|
|
msg := util.NewBeaconBlock()
|
|
msg.Block.ParentRoot = util.Random32Bytes(t)
|
|
msg.Signature = sk.Sign([]byte("data")).Marshal()
|
|
p2p.ReceivePubSub(topic, msg)
|
|
// wait for chainstart to be sent
|
|
time.Sleep(400 * time.Millisecond)
|
|
require.Equal(t, true, r.chainStarted.IsSet(), "Did not receive chain start event.")
|
|
}
|
|
|
|
func TestSyncHandlers_WaitForChainStart(t *testing.T) {
|
|
p2p := p2ptest.NewTestP2P(t)
|
|
chainService := &mockChain.ChainService{
|
|
Genesis: time.Now(),
|
|
ValidatorsRoot: [32]byte{'A'},
|
|
}
|
|
gs := startup.NewClockSynchronizer()
|
|
r := Service{
|
|
ctx: context.Background(),
|
|
cfg: &config{
|
|
p2p: p2p,
|
|
chain: chainService,
|
|
initialSync: &mockSync.Sync{IsSyncing: false},
|
|
},
|
|
chainStarted: abool.New(),
|
|
slotToPendingBlocks: gcache.New(time.Second, 2*time.Second),
|
|
clockWaiter: gs,
|
|
}
|
|
|
|
go r.registerHandlers()
|
|
var vr [32]byte
|
|
require.NoError(t, gs.SetClock(startup.NewClock(time.Now(), vr)))
|
|
r.waitForChainStart()
|
|
|
|
require.Equal(t, true, r.chainStarted.IsSet(), "Did not receive chain start event.")
|
|
}
|
|
|
|
func TestSyncHandlers_WaitTillSynced(t *testing.T) {
|
|
p2p := p2ptest.NewTestP2P(t)
|
|
chainService := &mockChain.ChainService{
|
|
Genesis: time.Now(),
|
|
ValidatorsRoot: [32]byte{'A'},
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
gs := startup.NewClockSynchronizer()
|
|
r := Service{
|
|
ctx: ctx,
|
|
cfg: &config{
|
|
p2p: p2p,
|
|
beaconDB: dbTest.SetupDB(t),
|
|
chain: chainService,
|
|
blockNotifier: chainService.BlockNotifier(),
|
|
initialSync: &mockSync.Sync{IsSyncing: false},
|
|
},
|
|
chainStarted: abool.New(),
|
|
subHandler: newSubTopicHandler(),
|
|
clockWaiter: gs,
|
|
initialSyncComplete: make(chan struct{}),
|
|
}
|
|
r.initCaches()
|
|
|
|
syncCompleteCh := make(chan bool)
|
|
go func() {
|
|
r.registerHandlers()
|
|
syncCompleteCh <- true
|
|
}()
|
|
var vr [32]byte
|
|
require.NoError(t, gs.SetClock(startup.NewClock(time.Now(), vr)))
|
|
r.waitForChainStart()
|
|
require.Equal(t, true, r.chainStarted.IsSet(), "Did not receive chain start event.")
|
|
|
|
blockChan := make(chan *feed.Event, 1)
|
|
sub := r.cfg.blockNotifier.BlockFeed().Subscribe(blockChan)
|
|
defer sub.Unsubscribe()
|
|
|
|
b := []byte("sk")
|
|
b32 := bytesutil.ToBytes32(b)
|
|
sk, err := bls.SecretKeyFromBytes(b32[:])
|
|
require.NoError(t, err)
|
|
msg := util.NewBeaconBlock()
|
|
msg.Block.ParentRoot = util.Random32Bytes(t)
|
|
msg.Signature = sk.Sign([]byte("data")).Marshal()
|
|
p2p.Digest, err = r.currentForkDigest()
|
|
require.NoError(t, err)
|
|
|
|
// Save block into DB so that validateBeaconBlockPubSub() process gets short cut.
|
|
util.SaveBlock(t, ctx, r.cfg.beaconDB, msg)
|
|
|
|
topic := "/eth2/%x/beacon_block"
|
|
p2p.ReceivePubSub(topic, msg)
|
|
assert.Equal(t, 0, len(blockChan), "block was received by sync service despite not being fully synced")
|
|
|
|
close(r.initialSyncComplete)
|
|
<-syncCompleteCh
|
|
|
|
p2p.ReceivePubSub(topic, msg)
|
|
|
|
select {
|
|
case <-blockChan:
|
|
case <-ctx.Done():
|
|
}
|
|
assert.NoError(t, ctx.Err())
|
|
}
|
|
|
|
func TestSyncService_StopCleanly(t *testing.T) {
|
|
p2p := p2ptest.NewTestP2P(t)
|
|
chainService := &mockChain.ChainService{
|
|
Genesis: time.Now(),
|
|
ValidatorsRoot: [32]byte{'A'},
|
|
}
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
gs := startup.NewClockSynchronizer()
|
|
r := Service{
|
|
ctx: ctx,
|
|
cancel: cancel,
|
|
cfg: &config{
|
|
p2p: p2p,
|
|
chain: chainService,
|
|
initialSync: &mockSync.Sync{IsSyncing: false},
|
|
},
|
|
chainStarted: abool.New(),
|
|
subHandler: newSubTopicHandler(),
|
|
clockWaiter: gs,
|
|
initialSyncComplete: make(chan struct{}),
|
|
}
|
|
|
|
go r.registerHandlers()
|
|
var vr [32]byte
|
|
require.NoError(t, gs.SetClock(startup.NewClock(time.Now(), vr)))
|
|
r.waitForChainStart()
|
|
|
|
var err error
|
|
p2p.Digest, err = r.currentForkDigest()
|
|
require.NoError(t, err)
|
|
|
|
// wait for chainstart to be sent
|
|
time.Sleep(2 * time.Second)
|
|
require.Equal(t, true, r.chainStarted.IsSet(), "Did not receive chain start event.")
|
|
|
|
close(r.initialSyncComplete)
|
|
time.Sleep(1 * time.Second)
|
|
|
|
require.NotEqual(t, 0, len(r.cfg.p2p.PubSub().GetTopics()))
|
|
require.NotEqual(t, 0, len(r.cfg.p2p.Host().Mux().Protocols()))
|
|
|
|
// Both pubsub and rpc topcis should be unsubscribed.
|
|
require.NoError(t, r.Stop())
|
|
|
|
// Sleep to allow pubsub topics to be deregistered.
|
|
time.Sleep(1 * time.Second)
|
|
require.Equal(t, 0, len(r.cfg.p2p.PubSub().GetTopics()))
|
|
require.Equal(t, 0, len(r.cfg.p2p.Host().Mux().Protocols()))
|
|
}
|