mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-20 10:52:25 +00:00
les: added test for StopMsg/ResumeMsg
This commit is contained in:
parent
3809ae416c
commit
b45ce52b6d
7 changed files with 75 additions and 22 deletions
|
|
@ -422,7 +422,7 @@ func (list RequestCostList) decode(protocolLength uint64) requestCostTable {
|
||||||
}
|
}
|
||||||
|
|
||||||
// testCostList returns a dummy request cost list used by tests
|
// testCostList returns a dummy request cost list used by tests
|
||||||
func testCostList() RequestCostList {
|
func testCostList(testCost uint64) RequestCostList {
|
||||||
cl := make(RequestCostList, len(reqAvgTimeCost))
|
cl := make(RequestCostList, len(reqAvgTimeCost))
|
||||||
var max uint64
|
var max uint64
|
||||||
for code := range reqAvgTimeCost {
|
for code := range reqAvgTimeCost {
|
||||||
|
|
@ -434,7 +434,7 @@ func testCostList() RequestCostList {
|
||||||
for code := uint64(0); code <= max; code++ {
|
for code := uint64(0); code <= max; code++ {
|
||||||
if _, ok := reqAvgTimeCost[code]; ok {
|
if _, ok := reqAvgTimeCost[code]; ok {
|
||||||
cl[i].MsgCode = code
|
cl[i].MsgCode = code
|
||||||
cl[i].BaseCost = 0
|
cl[i].BaseCost = testCost
|
||||||
cl[i].ReqCost = 0
|
cl[i].ReqCost = 0
|
||||||
i++
|
i++
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -258,9 +258,11 @@ func (cm *ClientManager) updateRecharge(now mclock.AbsTime) {
|
||||||
if sumRecharge > cm.totalRecharge {
|
if sumRecharge > cm.totalRecharge {
|
||||||
sumRecharge = cm.totalRecharge
|
sumRecharge = cm.totalRecharge
|
||||||
}
|
}
|
||||||
bonusRatio := cm.curve.ValueAt(sumRecharge) / float64(sumRecharge)
|
bonusRatio := float64(1)
|
||||||
if bonusRatio < 1 {
|
v := cm.curve.ValueAt(sumRecharge)
|
||||||
bonusRatio = 1
|
s := float64(sumRecharge)
|
||||||
|
if v > s && s > 0 {
|
||||||
|
bonusRatio = v / s
|
||||||
}
|
}
|
||||||
dt := now - lastUpdate
|
dt := now - lastUpdate
|
||||||
// fetch the client that finishes first
|
// fetch the client that finishes first
|
||||||
|
|
|
||||||
|
|
@ -24,6 +24,7 @@ import (
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/common/mclock"
|
||||||
"github.com/ethereum/go-ethereum/consensus/ethash"
|
"github.com/ethereum/go-ethereum/consensus/ethash"
|
||||||
"github.com/ethereum/go-ethereum/core"
|
"github.com/ethereum/go-ethereum/core"
|
||||||
"github.com/ethereum/go-ethereum/core/rawdb"
|
"github.com/ethereum/go-ethereum/core/rawdb"
|
||||||
|
|
@ -438,7 +439,7 @@ func TestTransactionStatusLes2(t *testing.T) {
|
||||||
config.Journal = ""
|
config.Journal = ""
|
||||||
txpool := core.NewTxPool(config, params.TestChainConfig, chain)
|
txpool := core.NewTxPool(config, params.TestChainConfig, chain)
|
||||||
pm.txpool = txpool
|
pm.txpool = txpool
|
||||||
peer, _ := newTestPeer(t, "peer", 2, pm, true)
|
peer, _ := newTestPeer(t, "peer", 2, pm, true, 0)
|
||||||
defer peer.close()
|
defer peer.close()
|
||||||
|
|
||||||
var reqID uint64
|
var reqID uint64
|
||||||
|
|
@ -519,3 +520,51 @@ func TestTransactionStatusLes2(t *testing.T) {
|
||||||
test(tx1, false, light.TxStatus{Status: core.TxStatusPending})
|
test(tx1, false, light.TxStatus{Status: core.TxStatusPending})
|
||||||
test(tx2, false, light.TxStatus{Status: core.TxStatusPending})
|
test(tx2, false, light.TxStatus{Status: core.TxStatusPending})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestStopResumeLes3(t *testing.T) {
|
||||||
|
db := rawdb.NewMemoryDatabase()
|
||||||
|
clock := &mclock.Simulated{}
|
||||||
|
testCost := testBufLimit / 10
|
||||||
|
pm, err := newTestProtocolManager(false, 0, nil, nil, nil, db, nil, testCost, clock)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Failed to create protocol manager: %v", err)
|
||||||
|
}
|
||||||
|
peer, _ := newTestPeer(t, "peer", 3, pm, true, testCost)
|
||||||
|
defer peer.close()
|
||||||
|
|
||||||
|
expBuf := testBufLimit
|
||||||
|
var reqID uint64
|
||||||
|
|
||||||
|
req := func() {
|
||||||
|
reqID++
|
||||||
|
sendRequest(peer.app, GetBlockHeadersMsg, reqID, testCost, &getBlockHeadersData{Origin: hashOrNumber{Hash: common.Hash{1}}, Amount: 1})
|
||||||
|
}
|
||||||
|
|
||||||
|
for i := 1; i <= 5; i++ {
|
||||||
|
// send requests while we still have enough buffer and expect a response
|
||||||
|
for expBuf >= testCost {
|
||||||
|
req()
|
||||||
|
expBuf -= testCost
|
||||||
|
if err := expectResponse(peer.app, BlockHeadersMsg, reqID, expBuf, nil); err != nil {
|
||||||
|
t.Errorf("expected response and failed: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// send some more requests in excess and expect a single StopMsg
|
||||||
|
c := i
|
||||||
|
for c > 0 {
|
||||||
|
req()
|
||||||
|
c--
|
||||||
|
}
|
||||||
|
if err := p2p.ExpectMsg(peer.app, StopMsg, nil); err != nil {
|
||||||
|
t.Errorf("expected StopMsg and failed: %v", err)
|
||||||
|
}
|
||||||
|
// wait until the buffer is recharged by half of the limit
|
||||||
|
wait := testBufLimit / testBufRecharge / 2
|
||||||
|
clock.Run(time.Millisecond * time.Duration(wait))
|
||||||
|
// expect a ResumeMsg with the partially recharged buffer value
|
||||||
|
expBuf += testBufRecharge * wait
|
||||||
|
if err := p2p.ExpectMsg(peer.app, ResumeMsg, expBuf); err != nil {
|
||||||
|
t.Errorf("expected ResumeMsg and failed: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -62,7 +62,8 @@ var (
|
||||||
testEventEmitterCode = common.Hex2Bytes("60606040523415600e57600080fd5b7f57050ab73f6b9ebdd9f76b8d4997793f48cf956e965ee070551b9ca0bb71584e60405160405180910390a160358060476000396000f3006060604052600080fd00a165627a7a723058203f727efcad8b5811f8cb1fc2620ce5e8c63570d697aef968172de296ea3994140029")
|
testEventEmitterCode = common.Hex2Bytes("60606040523415600e57600080fd5b7f57050ab73f6b9ebdd9f76b8d4997793f48cf956e965ee070551b9ca0bb71584e60405160405180910390a160358060476000396000f3006060604052600080fd00a165627a7a723058203f727efcad8b5811f8cb1fc2620ce5e8c63570d697aef968172de296ea3994140029")
|
||||||
testEventEmitterAddr common.Address
|
testEventEmitterAddr common.Address
|
||||||
|
|
||||||
testBufLimit = uint64(100)
|
testBufLimit = uint64(1000000)
|
||||||
|
testBufRecharge = uint64(1000)
|
||||||
)
|
)
|
||||||
|
|
||||||
/*
|
/*
|
||||||
|
|
@ -138,7 +139,7 @@ func testIndexers(db ethdb.Database, odr light.OdrBackend, iConfig *light.Indexe
|
||||||
// newTestProtocolManager creates a new protocol manager for testing purposes,
|
// newTestProtocolManager creates a new protocol manager for testing purposes,
|
||||||
// with the given number of blocks already known, potential notification
|
// with the given number of blocks already known, potential notification
|
||||||
// channels for different events and relative chain indexers array.
|
// channels for different events and relative chain indexers array.
|
||||||
func newTestProtocolManager(lightSync bool, blocks int, generator func(int, *core.BlockGen), odr *LesOdr, peers *peerSet, db ethdb.Database, ulcConfig *eth.ULCConfig) (*ProtocolManager, error) {
|
func newTestProtocolManager(lightSync bool, blocks int, generator func(int, *core.BlockGen), odr *LesOdr, peers *peerSet, db ethdb.Database, ulcConfig *eth.ULCConfig, testCost uint64, clock mclock.Clock) (*ProtocolManager, error) {
|
||||||
var (
|
var (
|
||||||
evmux = new(event.TypeMux)
|
evmux = new(event.TypeMux)
|
||||||
engine = ethash.NewFaker()
|
engine = ethash.NewFaker()
|
||||||
|
|
@ -182,10 +183,10 @@ func newTestProtocolManager(lightSync bool, blocks int, generator func(int, *cor
|
||||||
|
|
||||||
srv.defParams = flowcontrol.ServerParams{
|
srv.defParams = flowcontrol.ServerParams{
|
||||||
BufLimit: testBufLimit,
|
BufLimit: testBufLimit,
|
||||||
MinRecharge: 1,
|
MinRecharge: testBufRecharge,
|
||||||
}
|
}
|
||||||
|
srv.testCost = testCost
|
||||||
srv.fcManager = flowcontrol.NewClientManager(nil, &mclock.System{})
|
srv.fcManager = flowcontrol.NewClientManager(nil, clock)
|
||||||
}
|
}
|
||||||
pm.Start(1000)
|
pm.Start(1000)
|
||||||
return pm, nil
|
return pm, nil
|
||||||
|
|
@ -196,7 +197,7 @@ func newTestProtocolManager(lightSync bool, blocks int, generator func(int, *cor
|
||||||
// channels for different events and relative chain indexers array. In case of an error, the constructor force-
|
// channels for different events and relative chain indexers array. In case of an error, the constructor force-
|
||||||
// fails the test.
|
// fails the test.
|
||||||
func newTestProtocolManagerMust(t *testing.T, lightSync bool, blocks int, generator func(int, *core.BlockGen), odr *LesOdr, peers *peerSet, db ethdb.Database, ulcConfig *eth.ULCConfig) *ProtocolManager {
|
func newTestProtocolManagerMust(t *testing.T, lightSync bool, blocks int, generator func(int, *core.BlockGen), odr *LesOdr, peers *peerSet, db ethdb.Database, ulcConfig *eth.ULCConfig) *ProtocolManager {
|
||||||
pm, err := newTestProtocolManager(lightSync, blocks, generator, odr, peers, db, ulcConfig)
|
pm, err := newTestProtocolManager(lightSync, blocks, generator, odr, peers, db, ulcConfig, 0, &mclock.System{})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("Failed to create protocol manager: %v", err)
|
t.Fatalf("Failed to create protocol manager: %v", err)
|
||||||
}
|
}
|
||||||
|
|
@ -211,7 +212,7 @@ type testPeer struct {
|
||||||
}
|
}
|
||||||
|
|
||||||
// newTestPeer creates a new peer registered at the given protocol manager.
|
// newTestPeer creates a new peer registered at the given protocol manager.
|
||||||
func newTestPeer(t *testing.T, name string, version int, pm *ProtocolManager, shake bool) (*testPeer, <-chan error) {
|
func newTestPeer(t *testing.T, name string, version int, pm *ProtocolManager, shake bool, testCost uint64) (*testPeer, <-chan error) {
|
||||||
// Create a message pipe to communicate through
|
// Create a message pipe to communicate through
|
||||||
app, net := p2p.MsgPipe()
|
app, net := p2p.MsgPipe()
|
||||||
|
|
||||||
|
|
@ -243,7 +244,7 @@ func newTestPeer(t *testing.T, name string, version int, pm *ProtocolManager, sh
|
||||||
head = pm.blockchain.CurrentHeader()
|
head = pm.blockchain.CurrentHeader()
|
||||||
td = pm.blockchain.GetTd(head.Hash(), head.Number.Uint64())
|
td = pm.blockchain.GetTd(head.Hash(), head.Number.Uint64())
|
||||||
)
|
)
|
||||||
tp.handshake(t, td, head.Hash(), head.Number.Uint64(), genesis.Hash())
|
tp.handshake(t, td, head.Hash(), head.Number.Uint64(), genesis.Hash(), testCost)
|
||||||
}
|
}
|
||||||
return tp, errc
|
return tp, errc
|
||||||
}
|
}
|
||||||
|
|
@ -283,7 +284,7 @@ func newTestPeerPair(name string, version int, pm, pm2 *ProtocolManager) (*peer,
|
||||||
|
|
||||||
// handshake simulates a trivial handshake that expects the same state from the
|
// handshake simulates a trivial handshake that expects the same state from the
|
||||||
// remote side as we are simulating locally.
|
// remote side as we are simulating locally.
|
||||||
func (p *testPeer) handshake(t *testing.T, td *big.Int, head common.Hash, headNum uint64, genesis common.Hash) {
|
func (p *testPeer) handshake(t *testing.T, td *big.Int, head common.Hash, headNum uint64, genesis common.Hash, testCost uint64) {
|
||||||
var expList keyValueList
|
var expList keyValueList
|
||||||
expList = expList.add("protocolVersion", uint64(p.version))
|
expList = expList.add("protocolVersion", uint64(p.version))
|
||||||
expList = expList.add("networkId", uint64(NetworkId))
|
expList = expList.add("networkId", uint64(NetworkId))
|
||||||
|
|
@ -299,8 +300,8 @@ func (p *testPeer) handshake(t *testing.T, td *big.Int, head common.Hash, headNu
|
||||||
expList = expList.add("serveRecentState", uint64(core.TriesInMemory-4))
|
expList = expList.add("serveRecentState", uint64(core.TriesInMemory-4))
|
||||||
expList = expList.add("txRelay", nil)
|
expList = expList.add("txRelay", nil)
|
||||||
expList = expList.add("flowControl/BL", testBufLimit)
|
expList = expList.add("flowControl/BL", testBufLimit)
|
||||||
expList = expList.add("flowControl/MRR", uint64(1))
|
expList = expList.add("flowControl/MRR", testBufRecharge)
|
||||||
expList = expList.add("flowControl/MRC", testCostList())
|
expList = expList.add("flowControl/MRC", testCostList(testCost))
|
||||||
|
|
||||||
if err := p2p.ExpectMsg(p.app, StatusMsg, expList); err != nil {
|
if err := p2p.ExpectMsg(p.app, StatusMsg, expList); err != nil {
|
||||||
t.Fatalf("status recv: %v", err)
|
t.Fatalf("status recv: %v", err)
|
||||||
|
|
@ -311,7 +312,7 @@ func (p *testPeer) handshake(t *testing.T, td *big.Int, head common.Hash, headNu
|
||||||
|
|
||||||
p.fcParams = flowcontrol.ServerParams{
|
p.fcParams = flowcontrol.ServerParams{
|
||||||
BufLimit: testBufLimit,
|
BufLimit: testBufLimit,
|
||||||
MinRecharge: 1,
|
MinRecharge: testBufRecharge,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -340,7 +341,7 @@ func newServerEnv(t *testing.T, blocks int, protocol int, waitIndexers func(*cor
|
||||||
cIndexer, bIndexer, btIndexer := testIndexers(db, nil, light.TestServerIndexerConfig)
|
cIndexer, bIndexer, btIndexer := testIndexers(db, nil, light.TestServerIndexerConfig)
|
||||||
|
|
||||||
pm := newTestProtocolManagerMust(t, false, blocks, testChainGen, nil, nil, db, nil)
|
pm := newTestProtocolManagerMust(t, false, blocks, testChainGen, nil, nil, db, nil)
|
||||||
peer, _ := newTestPeer(t, "peer", protocol, pm, true)
|
peer, _ := newTestPeer(t, "peer", protocol, pm, true, 0)
|
||||||
|
|
||||||
cIndexer.Start(pm.blockchain.(*core.BlockChain))
|
cIndexer.Start(pm.blockchain.(*core.BlockChain))
|
||||||
bIndexer.Start(pm.blockchain.(*core.BlockChain))
|
bIndexer.Start(pm.blockchain.(*core.BlockChain))
|
||||||
|
|
|
||||||
|
|
@ -560,7 +560,7 @@ func (p *peer) Handshake(td *big.Int, head common.Hash, headNum uint64, genesis
|
||||||
if server.costTracker != nil {
|
if server.costTracker != nil {
|
||||||
costList = server.costTracker.makeCostList(server.costTracker.globalFactor())
|
costList = server.costTracker.makeCostList(server.costTracker.globalFactor())
|
||||||
} else {
|
} else {
|
||||||
costList = testCostList()
|
costList = testCostList(server.testCost)
|
||||||
}
|
}
|
||||||
send = send.add("flowControl/MRC", costList)
|
send = send.add("flowControl/MRC", costList)
|
||||||
p.fcCosts = costList.decode(ProtocolLengths[uint(p.version)])
|
p.fcCosts = costList.decode(ProtocolLengths[uint(p.version)])
|
||||||
|
|
|
||||||
|
|
@ -54,7 +54,7 @@ func TestPeerHandshakeSetAnnounceTypeToAnnounceTypeSignedForTrustedPeer(t *testi
|
||||||
l = l.add("txRelay", nil)
|
l = l.add("txRelay", nil)
|
||||||
l = l.add("flowControl/BL", uint64(0))
|
l = l.add("flowControl/BL", uint64(0))
|
||||||
l = l.add("flowControl/MRR", uint64(0))
|
l = l.add("flowControl/MRR", uint64(0))
|
||||||
l = l.add("flowControl/MRC", testCostList())
|
l = l.add("flowControl/MRC", testCostList(0))
|
||||||
|
|
||||||
return l
|
return l
|
||||||
},
|
},
|
||||||
|
|
@ -99,7 +99,7 @@ func TestPeerHandshakeAnnounceTypeSignedForTrustedPeersPeerNotInTrusted(t *testi
|
||||||
l = l.add("txRelay", nil)
|
l = l.add("txRelay", nil)
|
||||||
l = l.add("flowControl/BL", uint64(0))
|
l = l.add("flowControl/BL", uint64(0))
|
||||||
l = l.add("flowControl/MRR", uint64(0))
|
l = l.add("flowControl/MRR", uint64(0))
|
||||||
l = l.add("flowControl/MRC", testCostList())
|
l = l.add("flowControl/MRC", testCostList(0))
|
||||||
|
|
||||||
return l
|
return l
|
||||||
},
|
},
|
||||||
|
|
|
||||||
|
|
@ -53,6 +53,7 @@ type LesServer struct {
|
||||||
|
|
||||||
fcManager *flowcontrol.ClientManager // nil if our node is client only
|
fcManager *flowcontrol.ClientManager // nil if our node is client only
|
||||||
costTracker *costTracker
|
costTracker *costTracker
|
||||||
|
testCost uint64
|
||||||
defParams flowcontrol.ServerParams
|
defParams flowcontrol.ServerParams
|
||||||
lesTopics []discv5.Topic
|
lesTopics []discv5.Topic
|
||||||
privateKey *ecdsa.PrivateKey
|
privateKey *ecdsa.PrivateKey
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue