mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-24 21:56:43 +00:00
sync 2
This commit is contained in:
parent
0eef519481
commit
4d6c54dbf6
5 changed files with 25 additions and 21 deletions
|
|
@ -95,7 +95,7 @@ type ProtocolManager struct {
|
||||||
server *LesServer
|
server *LesServer
|
||||||
|
|
||||||
downloader *downloader.Downloader
|
downloader *downloader.Downloader
|
||||||
//fetcher *fetcher.Fetcher
|
fetcher *lightFetcher
|
||||||
peers *peerSet
|
peers *peerSet
|
||||||
|
|
||||||
SubProtocols []p2p.Protocol
|
SubProtocols []p2p.Protocol
|
||||||
|
|
@ -171,6 +171,7 @@ func NewProtocolManager(chainConfig *core.ChainConfig, lightSync bool, networkId
|
||||||
manager.downloader = downloader.New(downloader.LightSync, chainDb, manager.eventMux, blockchain.HasHeader, nil, blockchain.GetHeaderByHash,
|
manager.downloader = downloader.New(downloader.LightSync, chainDb, manager.eventMux, blockchain.HasHeader, nil, blockchain.GetHeaderByHash,
|
||||||
nil, blockchain.CurrentHeader, nil, nil, nil, blockchain.GetTdByHash,
|
nil, blockchain.CurrentHeader, nil, nil, nil, blockchain.GetTdByHash,
|
||||||
blockchain.InsertHeaderChain, nil, nil, blockchain.Rollback, func(id string) {}) // manager.removePeer)
|
blockchain.InsertHeaderChain, nil, nil, blockchain.Rollback, func(id string) {}) // manager.removePeer)
|
||||||
|
manager.fetcher = newLightFetcher(odr, blockchain)
|
||||||
}
|
}
|
||||||
|
|
||||||
/*validator := func(block *types.Block, parent *types.Block) error {
|
/*validator := func(block *types.Block, parent *types.Block) error {
|
||||||
|
|
@ -359,6 +360,9 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
|
||||||
return errResp(ErrDecode, "%v: %v", msg, err)
|
return errResp(ErrDecode, "%v: %v", msg, err)
|
||||||
}
|
}
|
||||||
fmt.Println("RECEIVED", req[0].Number, req[0].Hash, req[0].Td)
|
fmt.Println("RECEIVED", req[0].Number, req[0].Hash, req[0].Td)
|
||||||
|
for _, r := range req {
|
||||||
|
pm.fetcher.notify(p, r)
|
||||||
|
}
|
||||||
|
|
||||||
case GetBlockHeadersMsg:
|
case GetBlockHeadersMsg:
|
||||||
glog.V(logger.Debug).Infof("LES: received GetBlockHeadersMsg from peer %v", p.id)
|
glog.V(logger.Debug).Infof("LES: received GetBlockHeadersMsg from peer %v", p.id)
|
||||||
|
|
@ -452,9 +456,13 @@ fmt.Println("RECEIVED", req[0].Number, req[0].Hash, req[0].Td)
|
||||||
return errResp(ErrDecode, "msg %v: %v", msg, err)
|
return errResp(ErrDecode, "msg %v: %v", msg, err)
|
||||||
}
|
}
|
||||||
p.fcServer.GotReply(resp.ReqID, resp.BV)
|
p.fcServer.GotReply(resp.ReqID, resp.BV)
|
||||||
err := pm.downloader.DeliverHeaders(p.id, resp.Headers)
|
if pm.fetcher.requestedID(resp.ReqID) {
|
||||||
if err != nil {
|
pm.fetcher.deliverHeaders(resp.ReqID, resp.Headers)
|
||||||
glog.V(logger.Debug).Infoln(err)
|
} else {
|
||||||
|
err := pm.downloader.DeliverHeaders(p.id, resp.Headers)
|
||||||
|
if err != nil {
|
||||||
|
glog.V(logger.Debug).Infoln(err)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
case GetBlockBodiesMsg:
|
case GetBlockBodiesMsg:
|
||||||
|
|
@ -757,7 +765,7 @@ func (pm *ProtocolManager) broadcastBlockLoop() {
|
||||||
fmt.Println("BROADCAST", number, hash, td)
|
fmt.Println("BROADCAST", number, hash, td)
|
||||||
announce := newBlockHashesData{{Hash: hash, Number: number, Td: td}}
|
announce := newBlockHashesData{{Hash: hash, Number: number, Td: td}}
|
||||||
for _, p := range peers {
|
for _, p := range peers {
|
||||||
go p.SendNewBlockHashes(announce)
|
p.SendNewBlockHashes(announce)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
case <-pm.quitSync:
|
case <-pm.quitSync:
|
||||||
|
|
|
||||||
|
|
@ -191,7 +191,7 @@ func testOdr(t *testing.T, protocol int, expFail uint64, fn odrTestFn) {
|
||||||
t.Fatalf("peer 1 handshake error: %v", err)
|
t.Fatalf("peer 1 handshake error: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
lpm.synchronise(lpeer, true)
|
lpm.synchronise(lpeer)
|
||||||
|
|
||||||
test := func(expFail uint64) {
|
test := func(expFail uint64) {
|
||||||
for i := uint64(0); i <= pm.blockchain.CurrentHeader().GetNumberU64(); i++ {
|
for i := uint64(0); i <= pm.blockchain.CurrentHeader().GetNumberU64(); i++ {
|
||||||
|
|
|
||||||
|
|
@ -124,7 +124,8 @@ type statusData struct {
|
||||||
}
|
}
|
||||||
|
|
||||||
// newBlockHashesData is the network packet for the block announcements.
|
// newBlockHashesData is the network packet for the block announcements.
|
||||||
type newBlockHashesData []struct {
|
type newBlockHashesData []newBlockHashData
|
||||||
|
type newBlockHashData struct {
|
||||||
Hash common.Hash // Hash of one particular block being announced
|
Hash common.Hash // Hash of one particular block being announced
|
||||||
Number uint64 // Number of one particular block being announced
|
Number uint64 // Number of one particular block being announced
|
||||||
Td *big.Int // Total difficulty of one particular block being announced
|
Td *big.Int // Total difficulty of one particular block being announced
|
||||||
|
|
|
||||||
|
|
@ -63,7 +63,7 @@ func testAccess(t *testing.T, protocol int, fn accessTestFn) {
|
||||||
t.Fatalf("peer 1 handshake error: %v", err)
|
t.Fatalf("peer 1 handshake error: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
lpm.synchronise(lpeer, true)
|
lpm.synchronise(lpeer)
|
||||||
|
|
||||||
test := func(expFail uint64) {
|
test := func(expFail uint64) {
|
||||||
for i := uint64(0); i <= pm.blockchain.CurrentHeader().GetNumberU64(); i++ {
|
for i := uint64(0); i <= pm.blockchain.CurrentHeader().GetNumberU64(); i++ {
|
||||||
|
|
|
||||||
21
les/sync.go
21
les/sync.go
|
|
@ -44,11 +44,11 @@ func (pm *ProtocolManager) syncer() {
|
||||||
if pm.peers.Len() < minDesiredPeerCount {
|
if pm.peers.Len() < minDesiredPeerCount {
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
go pm.synchronise(pm.peers.BestPeer(), false)
|
go pm.synchronise(pm.peers.BestPeer())
|
||||||
|
|
||||||
case <-forceSync:
|
case <-forceSync:
|
||||||
// Force a sync even if not enough peers are present
|
// Force a sync even if not enough peers are present
|
||||||
go pm.synchronise(pm.peers.BestPeer(), false)
|
go pm.synchronise(pm.peers.BestPeer())
|
||||||
|
|
||||||
case <-pm.noMorePeers:
|
case <-pm.noMorePeers:
|
||||||
return
|
return
|
||||||
|
|
@ -57,22 +57,17 @@ func (pm *ProtocolManager) syncer() {
|
||||||
}
|
}
|
||||||
|
|
||||||
// synchronise tries to sync up our local block chain with a remote peer.
|
// synchronise tries to sync up our local block chain with a remote peer.
|
||||||
func (pm *ProtocolManager) synchronise(peer *peer, exit bool) {
|
func (pm *ProtocolManager) synchronise(peer *peer) {
|
||||||
// Short circuit if no peers are available
|
// Short circuit if no peers are available
|
||||||
if peer == nil {
|
if peer == nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// Make sure the peer's TD is higher than our own.
|
// Make sure the peer's TD is higher than our own.
|
||||||
td := pm.blockchain.GetTdByHash(pm.blockchain.LastBlockHash())
|
td := pm.blockchain.GetTdByHash(pm.blockchain.LastBlockHash())
|
||||||
if peer.Td().Cmp(td) > 0 {
|
if peer.Td().Cmp(td) <= 0 {
|
||||||
for {
|
return
|
||||||
if pm.downloader.Synchronise(peer.id, peer.Head(), peer.Td(), downloader.LightSync) != nil {
|
}
|
||||||
return
|
if pm.downloader.Synchronise(peer.id, peer.Head(), peer.Td(), downloader.LightSync) != nil {
|
||||||
}
|
return
|
||||||
if exit {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
time.Sleep(time.Second * 5)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue