From 4d6c54dbf6d609a1a9ce0dea69472fb96727b4a2 Mon Sep 17 00:00:00 2001 From: zsfelfoldi Date: Thu, 23 Jun 2016 19:26:35 +0200 Subject: [PATCH] sync 2 --- les/handler.go | 18 +++++++++++++----- les/odr_test.go | 2 +- les/protocol.go | 3 ++- les/request_test.go | 2 +- les/sync.go | 21 ++++++++------------- 5 files changed, 25 insertions(+), 21 deletions(-) diff --git a/les/handler.go b/les/handler.go index 7aacdaa5e6..2189dcd471 100644 --- a/les/handler.go +++ b/les/handler.go @@ -95,7 +95,7 @@ type ProtocolManager struct { server *LesServer downloader *downloader.Downloader - //fetcher *fetcher.Fetcher + fetcher *lightFetcher peers *peerSet 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, nil, blockchain.CurrentHeader, nil, nil, nil, blockchain.GetTdByHash, 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 { @@ -359,6 +360,9 @@ func (pm *ProtocolManager) handleMsg(p *peer) error { return errResp(ErrDecode, "%v: %v", msg, err) } fmt.Println("RECEIVED", req[0].Number, req[0].Hash, req[0].Td) + for _, r := range req { + pm.fetcher.notify(p, r) + } case GetBlockHeadersMsg: 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) } p.fcServer.GotReply(resp.ReqID, resp.BV) - err := pm.downloader.DeliverHeaders(p.id, resp.Headers) - if err != nil { - glog.V(logger.Debug).Infoln(err) + if pm.fetcher.requestedID(resp.ReqID) { + pm.fetcher.deliverHeaders(resp.ReqID, resp.Headers) + } else { + err := pm.downloader.DeliverHeaders(p.id, resp.Headers) + if err != nil { + glog.V(logger.Debug).Infoln(err) + } } case GetBlockBodiesMsg: @@ -757,7 +765,7 @@ func (pm *ProtocolManager) broadcastBlockLoop() { fmt.Println("BROADCAST", number, hash, td) announce := newBlockHashesData{{Hash: hash, Number: number, Td: td}} for _, p := range peers { - go p.SendNewBlockHashes(announce) + p.SendNewBlockHashes(announce) } } case <-pm.quitSync: diff --git a/les/odr_test.go b/les/odr_test.go index 08d0c64656..b70d9f5f28 100644 --- a/les/odr_test.go +++ b/les/odr_test.go @@ -191,7 +191,7 @@ func testOdr(t *testing.T, protocol int, expFail uint64, fn odrTestFn) { t.Fatalf("peer 1 handshake error: %v", err) } - lpm.synchronise(lpeer, true) + lpm.synchronise(lpeer) test := func(expFail uint64) { for i := uint64(0); i <= pm.blockchain.CurrentHeader().GetNumberU64(); i++ { diff --git a/les/protocol.go b/les/protocol.go index ca4dd51b29..352833c1a0 100644 --- a/les/protocol.go +++ b/les/protocol.go @@ -124,7 +124,8 @@ type statusData struct { } // 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 Number uint64 // Number of one particular block being announced Td *big.Int // Total difficulty of one particular block being announced diff --git a/les/request_test.go b/les/request_test.go index 320580c6fc..0d12b835f7 100644 --- a/les/request_test.go +++ b/les/request_test.go @@ -63,7 +63,7 @@ func testAccess(t *testing.T, protocol int, fn accessTestFn) { t.Fatalf("peer 1 handshake error: %v", err) } - lpm.synchronise(lpeer, true) + lpm.synchronise(lpeer) test := func(expFail uint64) { for i := uint64(0); i <= pm.blockchain.CurrentHeader().GetNumberU64(); i++ { diff --git a/les/sync.go b/les/sync.go index e540d4289f..608e992393 100644 --- a/les/sync.go +++ b/les/sync.go @@ -44,11 +44,11 @@ func (pm *ProtocolManager) syncer() { if pm.peers.Len() < minDesiredPeerCount { break } - go pm.synchronise(pm.peers.BestPeer(), false) + go pm.synchronise(pm.peers.BestPeer()) case <-forceSync: // 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: return @@ -57,22 +57,17 @@ func (pm *ProtocolManager) syncer() { } // 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 if peer == nil { return } // Make sure the peer's TD is higher than our own. td := pm.blockchain.GetTdByHash(pm.blockchain.LastBlockHash()) - if peer.Td().Cmp(td) > 0 { - for { - if pm.downloader.Synchronise(peer.id, peer.Head(), peer.Td(), downloader.LightSync) != nil { - return - } - if exit { - return - } - time.Sleep(time.Second * 5) - } + if peer.Td().Cmp(td) <= 0 { + return + } + if pm.downloader.Synchronise(peer.id, peer.Head(), peer.Td(), downloader.LightSync) != nil { + return } }