handler.go 38 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263
  1. // Copyright 2016 The go-ethereum Authors
  2. // This file is part of the go-ethereum library.
  3. //
  4. // The go-ethereum library is free software: you can redistribute it and/or modify
  5. // it under the terms of the GNU Lesser General Public License as published by
  6. // the Free Software Foundation, either version 3 of the License, or
  7. // (at your option) any later version.
  8. //
  9. // The go-ethereum library is distributed in the hope that it will be useful,
  10. // but WITHOUT ANY WARRANTY; without even the implied warranty of
  11. // MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  12. // GNU Lesser General Public License for more details.
  13. //
  14. // You should have received a copy of the GNU Lesser General Public License
  15. // along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
  16. // Package les implements the Light Ethereum Subprotocol.
  17. package les
  18. import (
  19. "encoding/binary"
  20. "errors"
  21. "fmt"
  22. "math/big"
  23. "net"
  24. "sync"
  25. "time"
  26. "github.com/ethereum/go-ethereum/common"
  27. "github.com/ethereum/go-ethereum/consensus"
  28. "github.com/ethereum/go-ethereum/core"
  29. "github.com/ethereum/go-ethereum/core/rawdb"
  30. "github.com/ethereum/go-ethereum/core/state"
  31. "github.com/ethereum/go-ethereum/core/types"
  32. "github.com/ethereum/go-ethereum/eth/downloader"
  33. "github.com/ethereum/go-ethereum/ethdb"
  34. "github.com/ethereum/go-ethereum/event"
  35. "github.com/ethereum/go-ethereum/light"
  36. "github.com/ethereum/go-ethereum/log"
  37. "github.com/ethereum/go-ethereum/p2p"
  38. "github.com/ethereum/go-ethereum/p2p/discover"
  39. "github.com/ethereum/go-ethereum/p2p/discv5"
  40. "github.com/ethereum/go-ethereum/params"
  41. "github.com/ethereum/go-ethereum/rlp"
  42. "github.com/ethereum/go-ethereum/trie"
  43. )
  44. const (
  45. softResponseLimit = 2 * 1024 * 1024 // Target maximum size of returned blocks, headers or node data.
  46. estHeaderRlpSize = 500 // Approximate size of an RLP encoded block header
  47. ethVersion = 63 // equivalent eth version for the downloader
  48. MaxHeaderFetch = 192 // Amount of block headers to be fetched per retrieval request
  49. MaxBodyFetch = 32 // Amount of block bodies to be fetched per retrieval request
  50. MaxReceiptFetch = 128 // Amount of transaction receipts to allow fetching per request
  51. MaxCodeFetch = 64 // Amount of contract codes to allow fetching per request
  52. MaxProofsFetch = 64 // Amount of merkle proofs to be fetched per retrieval request
  53. MaxHelperTrieProofsFetch = 64 // Amount of merkle proofs to be fetched per retrieval request
  54. MaxTxSend = 64 // Amount of transactions to be send per request
  55. MaxTxStatus = 256 // Amount of transactions to queried per request
  56. disableClientRemovePeer = false
  57. )
  58. // errIncompatibleConfig is returned if the requested protocols and configs are
  59. // not compatible (low protocol version restrictions and high requirements).
  60. var errIncompatibleConfig = errors.New("incompatible configuration")
  61. func errResp(code errCode, format string, v ...interface{}) error {
  62. return fmt.Errorf("%v - %v", code, fmt.Sprintf(format, v...))
  63. }
  64. type BlockChain interface {
  65. Config() *params.ChainConfig
  66. HasHeader(hash common.Hash, number uint64) bool
  67. GetHeader(hash common.Hash, number uint64) *types.Header
  68. GetHeaderByHash(hash common.Hash) *types.Header
  69. CurrentHeader() *types.Header
  70. GetTd(hash common.Hash, number uint64) *big.Int
  71. State() (*state.StateDB, error)
  72. InsertHeaderChain(chain []*types.Header, checkFreq int) (int, error)
  73. Rollback(chain []common.Hash)
  74. GetHeaderByNumber(number uint64) *types.Header
  75. GetBlockHashesFromHash(hash common.Hash, max uint64) []common.Hash
  76. Genesis() *types.Block
  77. SubscribeChainHeadEvent(ch chan<- core.ChainHeadEvent) event.Subscription
  78. }
  79. type txPool interface {
  80. AddRemotes(txs []*types.Transaction) []error
  81. Status(hashes []common.Hash) []core.TxStatus
  82. }
  83. type ProtocolManager struct {
  84. lightSync bool
  85. txpool txPool
  86. txrelay *LesTxRelay
  87. networkId uint64
  88. chainConfig *params.ChainConfig
  89. blockchain BlockChain
  90. chainDb ethdb.Database
  91. odr *LesOdr
  92. server *LesServer
  93. serverPool *serverPool
  94. lesTopic discv5.Topic
  95. reqDist *requestDistributor
  96. retriever *retrieveManager
  97. downloader *downloader.Downloader
  98. fetcher *lightFetcher
  99. peers *peerSet
  100. maxPeers int
  101. SubProtocols []p2p.Protocol
  102. eventMux *event.TypeMux
  103. // channels for fetcher, syncer, txsyncLoop
  104. newPeerCh chan *peer
  105. quitSync chan struct{}
  106. noMorePeers chan struct{}
  107. // wait group is used for graceful shutdowns during downloading
  108. // and processing
  109. wg *sync.WaitGroup
  110. }
  111. // NewProtocolManager returns a new ethereum sub protocol manager. The Ethereum sub protocol manages peers capable
  112. // with the ethereum network.
  113. func NewProtocolManager(chainConfig *params.ChainConfig, lightSync bool, protocolVersions []uint, networkId uint64, mux *event.TypeMux, engine consensus.Engine, peers *peerSet, blockchain BlockChain, txpool txPool, chainDb ethdb.Database, odr *LesOdr, txrelay *LesTxRelay, quitSync chan struct{}, wg *sync.WaitGroup) (*ProtocolManager, error) {
  114. // Create the protocol manager with the base fields
  115. manager := &ProtocolManager{
  116. lightSync: lightSync,
  117. eventMux: mux,
  118. blockchain: blockchain,
  119. chainConfig: chainConfig,
  120. chainDb: chainDb,
  121. odr: odr,
  122. networkId: networkId,
  123. txpool: txpool,
  124. txrelay: txrelay,
  125. peers: peers,
  126. newPeerCh: make(chan *peer),
  127. quitSync: quitSync,
  128. wg: wg,
  129. noMorePeers: make(chan struct{}),
  130. }
  131. if odr != nil {
  132. manager.retriever = odr.retriever
  133. manager.reqDist = odr.retriever.dist
  134. }
  135. // Initiate a sub-protocol for every implemented version we can handle
  136. manager.SubProtocols = make([]p2p.Protocol, 0, len(protocolVersions))
  137. for _, version := range protocolVersions {
  138. // Compatible, initialize the sub-protocol
  139. version := version // Closure for the run
  140. manager.SubProtocols = append(manager.SubProtocols, p2p.Protocol{
  141. Name: "les",
  142. Version: version,
  143. Length: ProtocolLengths[version],
  144. Run: func(p *p2p.Peer, rw p2p.MsgReadWriter) error {
  145. var entry *poolEntry
  146. peer := manager.newPeer(int(version), networkId, p, rw)
  147. if manager.serverPool != nil {
  148. addr := p.RemoteAddr().(*net.TCPAddr)
  149. entry = manager.serverPool.connect(peer, addr.IP, uint16(addr.Port))
  150. }
  151. peer.poolEntry = entry
  152. select {
  153. case manager.newPeerCh <- peer:
  154. manager.wg.Add(1)
  155. defer manager.wg.Done()
  156. err := manager.handle(peer)
  157. if entry != nil {
  158. manager.serverPool.disconnect(entry)
  159. }
  160. return err
  161. case <-manager.quitSync:
  162. if entry != nil {
  163. manager.serverPool.disconnect(entry)
  164. }
  165. return p2p.DiscQuitting
  166. }
  167. },
  168. NodeInfo: func() interface{} {
  169. return manager.NodeInfo()
  170. },
  171. PeerInfo: func(id discover.NodeID) interface{} {
  172. if p := manager.peers.Peer(fmt.Sprintf("%x", id[:8])); p != nil {
  173. return p.Info()
  174. }
  175. return nil
  176. },
  177. })
  178. }
  179. if len(manager.SubProtocols) == 0 {
  180. return nil, errIncompatibleConfig
  181. }
  182. removePeer := manager.removePeer
  183. if disableClientRemovePeer {
  184. removePeer = func(id string) {}
  185. }
  186. if lightSync {
  187. manager.downloader = downloader.New(downloader.LightSync, chainDb, manager.eventMux, nil, blockchain, removePeer)
  188. manager.peers.notify((*downloaderPeerNotify)(manager))
  189. manager.fetcher = newLightFetcher(manager)
  190. }
  191. return manager, nil
  192. }
  193. // removePeer initiates disconnection from a peer by removing it from the peer set
  194. func (pm *ProtocolManager) removePeer(id string) {
  195. pm.peers.Unregister(id)
  196. }
  197. func (pm *ProtocolManager) Start(maxPeers int) {
  198. pm.maxPeers = maxPeers
  199. if pm.lightSync {
  200. go pm.syncer()
  201. } else {
  202. go func() {
  203. for range pm.newPeerCh {
  204. }
  205. }()
  206. }
  207. }
  208. func (pm *ProtocolManager) Stop() {
  209. // Showing a log message. During download / process this could actually
  210. // take between 5 to 10 seconds and therefor feedback is required.
  211. log.Info("Stopping light Ethereum protocol")
  212. // Quit the sync loop.
  213. // After this send has completed, no new peers will be accepted.
  214. pm.noMorePeers <- struct{}{}
  215. close(pm.quitSync) // quits syncer, fetcher
  216. // Disconnect existing sessions.
  217. // This also closes the gate for any new registrations on the peer set.
  218. // sessions which are already established but not added to pm.peers yet
  219. // will exit when they try to register.
  220. pm.peers.Close()
  221. // Wait for any process action
  222. pm.wg.Wait()
  223. log.Info("Light Ethereum protocol stopped")
  224. }
  225. func (pm *ProtocolManager) newPeer(pv int, nv uint64, p *p2p.Peer, rw p2p.MsgReadWriter) *peer {
  226. return newPeer(pv, nv, p, newMeteredMsgWriter(rw))
  227. }
  228. // handle is the callback invoked to manage the life cycle of a les peer. When
  229. // this function terminates, the peer is disconnected.
  230. func (pm *ProtocolManager) handle(p *peer) error {
  231. // Ignore maxPeers if this is a trusted peer
  232. if pm.peers.Len() >= pm.maxPeers && !p.Peer.Info().Network.Trusted {
  233. return p2p.DiscTooManyPeers
  234. }
  235. p.Log().Debug("Light Ethereum peer connected", "name", p.Name())
  236. // Execute the LES handshake
  237. var (
  238. genesis = pm.blockchain.Genesis()
  239. head = pm.blockchain.CurrentHeader()
  240. hash = head.Hash()
  241. number = head.Number.Uint64()
  242. td = pm.blockchain.GetTd(hash, number)
  243. )
  244. if err := p.Handshake(td, hash, number, genesis.Hash(), pm.server); err != nil {
  245. p.Log().Debug("Light Ethereum handshake failed", "err", err)
  246. return err
  247. }
  248. if rw, ok := p.rw.(*meteredMsgReadWriter); ok {
  249. rw.Init(p.version)
  250. }
  251. // Register the peer locally
  252. if err := pm.peers.Register(p); err != nil {
  253. p.Log().Error("Light Ethereum peer registration failed", "err", err)
  254. return err
  255. }
  256. defer func() {
  257. if pm.server != nil && pm.server.fcManager != nil && p.fcClient != nil {
  258. p.fcClient.Remove(pm.server.fcManager)
  259. }
  260. pm.removePeer(p.id)
  261. }()
  262. // Register the peer in the downloader. If the downloader considers it banned, we disconnect
  263. if pm.lightSync {
  264. p.lock.Lock()
  265. head := p.headInfo
  266. p.lock.Unlock()
  267. if pm.fetcher != nil {
  268. pm.fetcher.announce(p, head)
  269. }
  270. if p.poolEntry != nil {
  271. pm.serverPool.registered(p.poolEntry)
  272. }
  273. }
  274. stop := make(chan struct{})
  275. defer close(stop)
  276. go func() {
  277. // new block announce loop
  278. for {
  279. select {
  280. case announce := <-p.announceChn:
  281. p.SendAnnounce(announce)
  282. case <-stop:
  283. return
  284. }
  285. }
  286. }()
  287. // main loop. handle incoming messages.
  288. for {
  289. if err := pm.handleMsg(p); err != nil {
  290. p.Log().Debug("Light Ethereum message handling failed", "err", err)
  291. return err
  292. }
  293. }
  294. }
  295. var reqList = []uint64{GetBlockHeadersMsg, GetBlockBodiesMsg, GetCodeMsg, GetReceiptsMsg, GetProofsV1Msg, SendTxMsg, SendTxV2Msg, GetTxStatusMsg, GetHeaderProofsMsg, GetProofsV2Msg, GetHelperTrieProofsMsg}
  296. // handleMsg is invoked whenever an inbound message is received from a remote
  297. // peer. The remote connection is torn down upon returning any error.
  298. func (pm *ProtocolManager) handleMsg(p *peer) error {
  299. // Read the next message from the remote peer, and ensure it's fully consumed
  300. msg, err := p.rw.ReadMsg()
  301. if err != nil {
  302. return err
  303. }
  304. p.Log().Trace("Light Ethereum message arrived", "code", msg.Code, "bytes", msg.Size)
  305. costs := p.fcCosts[msg.Code]
  306. reject := func(reqCnt, maxCnt uint64) bool {
  307. if p.fcClient == nil || reqCnt > maxCnt {
  308. return true
  309. }
  310. bufValue, _ := p.fcClient.AcceptRequest()
  311. cost := costs.baseCost + reqCnt*costs.reqCost
  312. if cost > pm.server.defParams.BufLimit {
  313. cost = pm.server.defParams.BufLimit
  314. }
  315. if cost > bufValue {
  316. recharge := time.Duration((cost - bufValue) * 1000000 / pm.server.defParams.MinRecharge)
  317. p.Log().Error("Request came too early", "recharge", common.PrettyDuration(recharge))
  318. return true
  319. }
  320. return false
  321. }
  322. if msg.Size > ProtocolMaxMsgSize {
  323. return errResp(ErrMsgTooLarge, "%v > %v", msg.Size, ProtocolMaxMsgSize)
  324. }
  325. defer msg.Discard()
  326. var deliverMsg *Msg
  327. // Handle the message depending on its contents
  328. switch msg.Code {
  329. case StatusMsg:
  330. p.Log().Trace("Received status message")
  331. // Status messages should never arrive after the handshake
  332. return errResp(ErrExtraStatusMsg, "uncontrolled status message")
  333. // Block header query, collect the requested headers and reply
  334. case AnnounceMsg:
  335. p.Log().Trace("Received announce message")
  336. if p.requestAnnounceType == announceTypeNone {
  337. return errResp(ErrUnexpectedResponse, "")
  338. }
  339. var req announceData
  340. if err := msg.Decode(&req); err != nil {
  341. return errResp(ErrDecode, "%v: %v", msg, err)
  342. }
  343. if p.requestAnnounceType == announceTypeSigned {
  344. if err := req.checkSignature(p.pubKey); err != nil {
  345. p.Log().Trace("Invalid announcement signature", "err", err)
  346. return err
  347. }
  348. p.Log().Trace("Valid announcement signature")
  349. }
  350. p.Log().Trace("Announce message content", "number", req.Number, "hash", req.Hash, "td", req.Td, "reorg", req.ReorgDepth)
  351. if pm.fetcher != nil {
  352. pm.fetcher.announce(p, &req)
  353. }
  354. case GetBlockHeadersMsg:
  355. p.Log().Trace("Received block header request")
  356. // Decode the complex header query
  357. var req struct {
  358. ReqID uint64
  359. Query getBlockHeadersData
  360. }
  361. if err := msg.Decode(&req); err != nil {
  362. return errResp(ErrDecode, "%v: %v", msg, err)
  363. }
  364. query := req.Query
  365. if reject(query.Amount, MaxHeaderFetch) {
  366. return errResp(ErrRequestRejected, "")
  367. }
  368. hashMode := query.Origin.Hash != (common.Hash{})
  369. // Gather headers until the fetch or network limits is reached
  370. var (
  371. bytes common.StorageSize
  372. headers []*types.Header
  373. unknown bool
  374. )
  375. for !unknown && len(headers) < int(query.Amount) && bytes < softResponseLimit {
  376. // Retrieve the next header satisfying the query
  377. var origin *types.Header
  378. if hashMode {
  379. origin = pm.blockchain.GetHeaderByHash(query.Origin.Hash)
  380. } else {
  381. origin = pm.blockchain.GetHeaderByNumber(query.Origin.Number)
  382. }
  383. if origin == nil {
  384. break
  385. }
  386. number := origin.Number.Uint64()
  387. headers = append(headers, origin)
  388. bytes += estHeaderRlpSize
  389. // Advance to the next header of the query
  390. switch {
  391. case query.Origin.Hash != (common.Hash{}) && query.Reverse:
  392. // Hash based traversal towards the genesis block
  393. for i := 0; i < int(query.Skip)+1; i++ {
  394. if header := pm.blockchain.GetHeader(query.Origin.Hash, number); header != nil {
  395. query.Origin.Hash = header.ParentHash
  396. number--
  397. } else {
  398. unknown = true
  399. break
  400. }
  401. }
  402. case query.Origin.Hash != (common.Hash{}) && !query.Reverse:
  403. // Hash based traversal towards the leaf block
  404. if header := pm.blockchain.GetHeaderByNumber(origin.Number.Uint64() + query.Skip + 1); header != nil {
  405. if pm.blockchain.GetBlockHashesFromHash(header.Hash(), query.Skip+1)[query.Skip] == query.Origin.Hash {
  406. query.Origin.Hash = header.Hash()
  407. } else {
  408. unknown = true
  409. }
  410. } else {
  411. unknown = true
  412. }
  413. case query.Reverse:
  414. // Number based traversal towards the genesis block
  415. if query.Origin.Number >= query.Skip+1 {
  416. query.Origin.Number -= query.Skip + 1
  417. } else {
  418. unknown = true
  419. }
  420. case !query.Reverse:
  421. // Number based traversal towards the leaf block
  422. query.Origin.Number += query.Skip + 1
  423. }
  424. }
  425. bv, rcost := p.fcClient.RequestProcessed(costs.baseCost + query.Amount*costs.reqCost)
  426. pm.server.fcCostStats.update(msg.Code, query.Amount, rcost)
  427. return p.SendBlockHeaders(req.ReqID, bv, headers)
  428. case BlockHeadersMsg:
  429. if pm.downloader == nil {
  430. return errResp(ErrUnexpectedResponse, "")
  431. }
  432. p.Log().Trace("Received block header response message")
  433. // A batch of headers arrived to one of our previous requests
  434. var resp struct {
  435. ReqID, BV uint64
  436. Headers []*types.Header
  437. }
  438. if err := msg.Decode(&resp); err != nil {
  439. return errResp(ErrDecode, "msg %v: %v", msg, err)
  440. }
  441. p.fcServer.GotReply(resp.ReqID, resp.BV)
  442. if pm.fetcher != nil && pm.fetcher.requestedID(resp.ReqID) {
  443. pm.fetcher.deliverHeaders(p, resp.ReqID, resp.Headers)
  444. } else {
  445. err := pm.downloader.DeliverHeaders(p.id, resp.Headers)
  446. if err != nil {
  447. log.Debug(fmt.Sprint(err))
  448. }
  449. }
  450. case GetBlockBodiesMsg:
  451. p.Log().Trace("Received block bodies request")
  452. // Decode the retrieval message
  453. var req struct {
  454. ReqID uint64
  455. Hashes []common.Hash
  456. }
  457. if err := msg.Decode(&req); err != nil {
  458. return errResp(ErrDecode, "msg %v: %v", msg, err)
  459. }
  460. // Gather blocks until the fetch or network limits is reached
  461. var (
  462. bytes int
  463. bodies []rlp.RawValue
  464. )
  465. reqCnt := len(req.Hashes)
  466. if reject(uint64(reqCnt), MaxBodyFetch) {
  467. return errResp(ErrRequestRejected, "")
  468. }
  469. for _, hash := range req.Hashes {
  470. if bytes >= softResponseLimit {
  471. break
  472. }
  473. // Retrieve the requested block body, stopping if enough was found
  474. if number := rawdb.ReadHeaderNumber(pm.chainDb, hash); number != nil {
  475. if data := rawdb.ReadBodyRLP(pm.chainDb, hash, *number); len(data) != 0 {
  476. bodies = append(bodies, data)
  477. bytes += len(data)
  478. }
  479. }
  480. }
  481. bv, rcost := p.fcClient.RequestProcessed(costs.baseCost + uint64(reqCnt)*costs.reqCost)
  482. pm.server.fcCostStats.update(msg.Code, uint64(reqCnt), rcost)
  483. return p.SendBlockBodiesRLP(req.ReqID, bv, bodies)
  484. case BlockBodiesMsg:
  485. if pm.odr == nil {
  486. return errResp(ErrUnexpectedResponse, "")
  487. }
  488. p.Log().Trace("Received block bodies response")
  489. // A batch of block bodies arrived to one of our previous requests
  490. var resp struct {
  491. ReqID, BV uint64
  492. Data []*types.Body
  493. }
  494. if err := msg.Decode(&resp); err != nil {
  495. return errResp(ErrDecode, "msg %v: %v", msg, err)
  496. }
  497. p.fcServer.GotReply(resp.ReqID, resp.BV)
  498. deliverMsg = &Msg{
  499. MsgType: MsgBlockBodies,
  500. ReqID: resp.ReqID,
  501. Obj: resp.Data,
  502. }
  503. case GetCodeMsg:
  504. p.Log().Trace("Received code request")
  505. // Decode the retrieval message
  506. var req struct {
  507. ReqID uint64
  508. Reqs []CodeReq
  509. }
  510. if err := msg.Decode(&req); err != nil {
  511. return errResp(ErrDecode, "msg %v: %v", msg, err)
  512. }
  513. // Gather state data until the fetch or network limits is reached
  514. var (
  515. bytes int
  516. data [][]byte
  517. )
  518. reqCnt := len(req.Reqs)
  519. if reject(uint64(reqCnt), MaxCodeFetch) {
  520. return errResp(ErrRequestRejected, "")
  521. }
  522. for _, req := range req.Reqs {
  523. // Retrieve the requested state entry, stopping if enough was found
  524. if number := rawdb.ReadHeaderNumber(pm.chainDb, req.BHash); number != nil {
  525. if header := rawdb.ReadHeader(pm.chainDb, req.BHash, *number); header != nil {
  526. statedb, err := pm.blockchain.State()
  527. if err != nil {
  528. continue
  529. }
  530. account, err := pm.getAccount(statedb, header.Root, common.BytesToHash(req.AccKey))
  531. if err != nil {
  532. continue
  533. }
  534. code, _ := statedb.Database().TrieDB().Node(common.BytesToHash(account.CodeHash))
  535. data = append(data, code)
  536. if bytes += len(code); bytes >= softResponseLimit {
  537. break
  538. }
  539. }
  540. }
  541. }
  542. bv, rcost := p.fcClient.RequestProcessed(costs.baseCost + uint64(reqCnt)*costs.reqCost)
  543. pm.server.fcCostStats.update(msg.Code, uint64(reqCnt), rcost)
  544. return p.SendCode(req.ReqID, bv, data)
  545. case CodeMsg:
  546. if pm.odr == nil {
  547. return errResp(ErrUnexpectedResponse, "")
  548. }
  549. p.Log().Trace("Received code response")
  550. // A batch of node state data arrived to one of our previous requests
  551. var resp struct {
  552. ReqID, BV uint64
  553. Data [][]byte
  554. }
  555. if err := msg.Decode(&resp); err != nil {
  556. return errResp(ErrDecode, "msg %v: %v", msg, err)
  557. }
  558. p.fcServer.GotReply(resp.ReqID, resp.BV)
  559. deliverMsg = &Msg{
  560. MsgType: MsgCode,
  561. ReqID: resp.ReqID,
  562. Obj: resp.Data,
  563. }
  564. case GetReceiptsMsg:
  565. p.Log().Trace("Received receipts request")
  566. // Decode the retrieval message
  567. var req struct {
  568. ReqID uint64
  569. Hashes []common.Hash
  570. }
  571. if err := msg.Decode(&req); err != nil {
  572. return errResp(ErrDecode, "msg %v: %v", msg, err)
  573. }
  574. // Gather state data until the fetch or network limits is reached
  575. var (
  576. bytes int
  577. receipts []rlp.RawValue
  578. )
  579. reqCnt := len(req.Hashes)
  580. if reject(uint64(reqCnt), MaxReceiptFetch) {
  581. return errResp(ErrRequestRejected, "")
  582. }
  583. for _, hash := range req.Hashes {
  584. if bytes >= softResponseLimit {
  585. break
  586. }
  587. // Retrieve the requested block's receipts, skipping if unknown to us
  588. var results types.Receipts
  589. if number := rawdb.ReadHeaderNumber(pm.chainDb, hash); number != nil {
  590. results = rawdb.ReadReceipts(pm.chainDb, hash, *number)
  591. }
  592. if results == nil {
  593. if header := pm.blockchain.GetHeaderByHash(hash); header == nil || header.ReceiptHash != types.EmptyRootHash {
  594. continue
  595. }
  596. }
  597. // If known, encode and queue for response packet
  598. if encoded, err := rlp.EncodeToBytes(results); err != nil {
  599. log.Error("Failed to encode receipt", "err", err)
  600. } else {
  601. receipts = append(receipts, encoded)
  602. bytes += len(encoded)
  603. }
  604. }
  605. bv, rcost := p.fcClient.RequestProcessed(costs.baseCost + uint64(reqCnt)*costs.reqCost)
  606. pm.server.fcCostStats.update(msg.Code, uint64(reqCnt), rcost)
  607. return p.SendReceiptsRLP(req.ReqID, bv, receipts)
  608. case ReceiptsMsg:
  609. if pm.odr == nil {
  610. return errResp(ErrUnexpectedResponse, "")
  611. }
  612. p.Log().Trace("Received receipts response")
  613. // A batch of receipts arrived to one of our previous requests
  614. var resp struct {
  615. ReqID, BV uint64
  616. Receipts []types.Receipts
  617. }
  618. if err := msg.Decode(&resp); err != nil {
  619. return errResp(ErrDecode, "msg %v: %v", msg, err)
  620. }
  621. p.fcServer.GotReply(resp.ReqID, resp.BV)
  622. deliverMsg = &Msg{
  623. MsgType: MsgReceipts,
  624. ReqID: resp.ReqID,
  625. Obj: resp.Receipts,
  626. }
  627. case GetProofsV1Msg:
  628. p.Log().Trace("Received proofs request")
  629. // Decode the retrieval message
  630. var req struct {
  631. ReqID uint64
  632. Reqs []ProofReq
  633. }
  634. if err := msg.Decode(&req); err != nil {
  635. return errResp(ErrDecode, "msg %v: %v", msg, err)
  636. }
  637. // Gather state data until the fetch or network limits is reached
  638. var (
  639. bytes int
  640. proofs proofsData
  641. )
  642. reqCnt := len(req.Reqs)
  643. if reject(uint64(reqCnt), MaxProofsFetch) {
  644. return errResp(ErrRequestRejected, "")
  645. }
  646. for _, req := range req.Reqs {
  647. // Retrieve the requested state entry, stopping if enough was found
  648. if number := rawdb.ReadHeaderNumber(pm.chainDb, req.BHash); number != nil {
  649. if header := rawdb.ReadHeader(pm.chainDb, req.BHash, *number); header != nil {
  650. statedb, err := pm.blockchain.State()
  651. if err != nil {
  652. continue
  653. }
  654. var trie state.Trie
  655. if len(req.AccKey) > 0 {
  656. account, err := pm.getAccount(statedb, header.Root, common.BytesToHash(req.AccKey))
  657. if err != nil {
  658. continue
  659. }
  660. trie, _ = statedb.Database().OpenStorageTrie(common.BytesToHash(req.AccKey), account.Root)
  661. } else {
  662. trie, _ = statedb.Database().OpenTrie(header.Root)
  663. }
  664. if trie != nil {
  665. var proof light.NodeList
  666. trie.Prove(req.Key, 0, &proof)
  667. proofs = append(proofs, proof)
  668. if bytes += proof.DataSize(); bytes >= softResponseLimit {
  669. break
  670. }
  671. }
  672. }
  673. }
  674. }
  675. bv, rcost := p.fcClient.RequestProcessed(costs.baseCost + uint64(reqCnt)*costs.reqCost)
  676. pm.server.fcCostStats.update(msg.Code, uint64(reqCnt), rcost)
  677. return p.SendProofs(req.ReqID, bv, proofs)
  678. case GetProofsV2Msg:
  679. p.Log().Trace("Received les/2 proofs request")
  680. // Decode the retrieval message
  681. var req struct {
  682. ReqID uint64
  683. Reqs []ProofReq
  684. }
  685. if err := msg.Decode(&req); err != nil {
  686. return errResp(ErrDecode, "msg %v: %v", msg, err)
  687. }
  688. // Gather state data until the fetch or network limits is reached
  689. var (
  690. lastBHash common.Hash
  691. statedb *state.StateDB
  692. root common.Hash
  693. )
  694. reqCnt := len(req.Reqs)
  695. if reject(uint64(reqCnt), MaxProofsFetch) {
  696. return errResp(ErrRequestRejected, "")
  697. }
  698. nodes := light.NewNodeSet()
  699. for _, req := range req.Reqs {
  700. // Look up the state belonging to the request
  701. if statedb == nil || req.BHash != lastBHash {
  702. statedb, root, lastBHash = nil, common.Hash{}, req.BHash
  703. if number := rawdb.ReadHeaderNumber(pm.chainDb, req.BHash); number != nil {
  704. if header := rawdb.ReadHeader(pm.chainDb, req.BHash, *number); header != nil {
  705. statedb, _ = pm.blockchain.State()
  706. root = header.Root
  707. }
  708. }
  709. }
  710. if statedb == nil {
  711. continue
  712. }
  713. // Pull the account or storage trie of the request
  714. var trie state.Trie
  715. if len(req.AccKey) > 0 {
  716. account, err := pm.getAccount(statedb, root, common.BytesToHash(req.AccKey))
  717. if err != nil {
  718. continue
  719. }
  720. trie, _ = statedb.Database().OpenStorageTrie(common.BytesToHash(req.AccKey), account.Root)
  721. } else {
  722. trie, _ = statedb.Database().OpenTrie(root)
  723. }
  724. if trie == nil {
  725. continue
  726. }
  727. // Prove the user's request from the account or stroage trie
  728. trie.Prove(req.Key, req.FromLevel, nodes)
  729. if nodes.DataSize() >= softResponseLimit {
  730. break
  731. }
  732. }
  733. bv, rcost := p.fcClient.RequestProcessed(costs.baseCost + uint64(reqCnt)*costs.reqCost)
  734. pm.server.fcCostStats.update(msg.Code, uint64(reqCnt), rcost)
  735. return p.SendProofsV2(req.ReqID, bv, nodes.NodeList())
  736. case ProofsV1Msg:
  737. if pm.odr == nil {
  738. return errResp(ErrUnexpectedResponse, "")
  739. }
  740. p.Log().Trace("Received proofs response")
  741. // A batch of merkle proofs arrived to one of our previous requests
  742. var resp struct {
  743. ReqID, BV uint64
  744. Data []light.NodeList
  745. }
  746. if err := msg.Decode(&resp); err != nil {
  747. return errResp(ErrDecode, "msg %v: %v", msg, err)
  748. }
  749. p.fcServer.GotReply(resp.ReqID, resp.BV)
  750. deliverMsg = &Msg{
  751. MsgType: MsgProofsV1,
  752. ReqID: resp.ReqID,
  753. Obj: resp.Data,
  754. }
  755. case ProofsV2Msg:
  756. if pm.odr == nil {
  757. return errResp(ErrUnexpectedResponse, "")
  758. }
  759. p.Log().Trace("Received les/2 proofs response")
  760. // A batch of merkle proofs arrived to one of our previous requests
  761. var resp struct {
  762. ReqID, BV uint64
  763. Data light.NodeList
  764. }
  765. if err := msg.Decode(&resp); err != nil {
  766. return errResp(ErrDecode, "msg %v: %v", msg, err)
  767. }
  768. p.fcServer.GotReply(resp.ReqID, resp.BV)
  769. deliverMsg = &Msg{
  770. MsgType: MsgProofsV2,
  771. ReqID: resp.ReqID,
  772. Obj: resp.Data,
  773. }
  774. case GetHeaderProofsMsg:
  775. p.Log().Trace("Received headers proof request")
  776. // Decode the retrieval message
  777. var req struct {
  778. ReqID uint64
  779. Reqs []ChtReq
  780. }
  781. if err := msg.Decode(&req); err != nil {
  782. return errResp(ErrDecode, "msg %v: %v", msg, err)
  783. }
  784. // Gather state data until the fetch or network limits is reached
  785. var (
  786. bytes int
  787. proofs []ChtResp
  788. )
  789. reqCnt := len(req.Reqs)
  790. if reject(uint64(reqCnt), MaxHelperTrieProofsFetch) {
  791. return errResp(ErrRequestRejected, "")
  792. }
  793. trieDb := trie.NewDatabase(ethdb.NewTable(pm.chainDb, light.ChtTablePrefix))
  794. for _, req := range req.Reqs {
  795. if header := pm.blockchain.GetHeaderByNumber(req.BlockNum); header != nil {
  796. sectionHead := rawdb.ReadCanonicalHash(pm.chainDb, req.ChtNum*light.CHTFrequencyServer-1)
  797. if root := light.GetChtRoot(pm.chainDb, req.ChtNum-1, sectionHead); root != (common.Hash{}) {
  798. trie, err := trie.New(root, trieDb)
  799. if err != nil {
  800. continue
  801. }
  802. var encNumber [8]byte
  803. binary.BigEndian.PutUint64(encNumber[:], req.BlockNum)
  804. var proof light.NodeList
  805. trie.Prove(encNumber[:], 0, &proof)
  806. proofs = append(proofs, ChtResp{Header: header, Proof: proof})
  807. if bytes += proof.DataSize() + estHeaderRlpSize; bytes >= softResponseLimit {
  808. break
  809. }
  810. }
  811. }
  812. }
  813. bv, rcost := p.fcClient.RequestProcessed(costs.baseCost + uint64(reqCnt)*costs.reqCost)
  814. pm.server.fcCostStats.update(msg.Code, uint64(reqCnt), rcost)
  815. return p.SendHeaderProofs(req.ReqID, bv, proofs)
  816. case GetHelperTrieProofsMsg:
  817. p.Log().Trace("Received helper trie proof request")
  818. // Decode the retrieval message
  819. var req struct {
  820. ReqID uint64
  821. Reqs []HelperTrieReq
  822. }
  823. if err := msg.Decode(&req); err != nil {
  824. return errResp(ErrDecode, "msg %v: %v", msg, err)
  825. }
  826. // Gather state data until the fetch or network limits is reached
  827. var (
  828. auxBytes int
  829. auxData [][]byte
  830. )
  831. reqCnt := len(req.Reqs)
  832. if reject(uint64(reqCnt), MaxHelperTrieProofsFetch) {
  833. return errResp(ErrRequestRejected, "")
  834. }
  835. var (
  836. lastIdx uint64
  837. lastType uint
  838. root common.Hash
  839. auxTrie *trie.Trie
  840. )
  841. nodes := light.NewNodeSet()
  842. for _, req := range req.Reqs {
  843. if auxTrie == nil || req.Type != lastType || req.TrieIdx != lastIdx {
  844. auxTrie, lastType, lastIdx = nil, req.Type, req.TrieIdx
  845. var prefix string
  846. if root, prefix = pm.getHelperTrie(req.Type, req.TrieIdx); root != (common.Hash{}) {
  847. auxTrie, _ = trie.New(root, trie.NewDatabase(ethdb.NewTable(pm.chainDb, prefix)))
  848. }
  849. }
  850. if req.AuxReq == auxRoot {
  851. var data []byte
  852. if root != (common.Hash{}) {
  853. data = root[:]
  854. }
  855. auxData = append(auxData, data)
  856. auxBytes += len(data)
  857. } else {
  858. if auxTrie != nil {
  859. auxTrie.Prove(req.Key, req.FromLevel, nodes)
  860. }
  861. if req.AuxReq != 0 {
  862. data := pm.getHelperTrieAuxData(req)
  863. auxData = append(auxData, data)
  864. auxBytes += len(data)
  865. }
  866. }
  867. if nodes.DataSize()+auxBytes >= softResponseLimit {
  868. break
  869. }
  870. }
  871. bv, rcost := p.fcClient.RequestProcessed(costs.baseCost + uint64(reqCnt)*costs.reqCost)
  872. pm.server.fcCostStats.update(msg.Code, uint64(reqCnt), rcost)
  873. return p.SendHelperTrieProofs(req.ReqID, bv, HelperTrieResps{Proofs: nodes.NodeList(), AuxData: auxData})
  874. case HeaderProofsMsg:
  875. if pm.odr == nil {
  876. return errResp(ErrUnexpectedResponse, "")
  877. }
  878. p.Log().Trace("Received headers proof response")
  879. var resp struct {
  880. ReqID, BV uint64
  881. Data []ChtResp
  882. }
  883. if err := msg.Decode(&resp); err != nil {
  884. return errResp(ErrDecode, "msg %v: %v", msg, err)
  885. }
  886. p.fcServer.GotReply(resp.ReqID, resp.BV)
  887. deliverMsg = &Msg{
  888. MsgType: MsgHeaderProofs,
  889. ReqID: resp.ReqID,
  890. Obj: resp.Data,
  891. }
  892. case HelperTrieProofsMsg:
  893. if pm.odr == nil {
  894. return errResp(ErrUnexpectedResponse, "")
  895. }
  896. p.Log().Trace("Received helper trie proof response")
  897. var resp struct {
  898. ReqID, BV uint64
  899. Data HelperTrieResps
  900. }
  901. if err := msg.Decode(&resp); err != nil {
  902. return errResp(ErrDecode, "msg %v: %v", msg, err)
  903. }
  904. p.fcServer.GotReply(resp.ReqID, resp.BV)
  905. deliverMsg = &Msg{
  906. MsgType: MsgHelperTrieProofs,
  907. ReqID: resp.ReqID,
  908. Obj: resp.Data,
  909. }
  910. case SendTxMsg:
  911. if pm.txpool == nil {
  912. return errResp(ErrRequestRejected, "")
  913. }
  914. // Transactions arrived, parse all of them and deliver to the pool
  915. var txs []*types.Transaction
  916. if err := msg.Decode(&txs); err != nil {
  917. return errResp(ErrDecode, "msg %v: %v", msg, err)
  918. }
  919. reqCnt := len(txs)
  920. if reject(uint64(reqCnt), MaxTxSend) {
  921. return errResp(ErrRequestRejected, "")
  922. }
  923. pm.txpool.AddRemotes(txs)
  924. _, rcost := p.fcClient.RequestProcessed(costs.baseCost + uint64(reqCnt)*costs.reqCost)
  925. pm.server.fcCostStats.update(msg.Code, uint64(reqCnt), rcost)
  926. case SendTxV2Msg:
  927. if pm.txpool == nil {
  928. return errResp(ErrRequestRejected, "")
  929. }
  930. // Transactions arrived, parse all of them and deliver to the pool
  931. var req struct {
  932. ReqID uint64
  933. Txs []*types.Transaction
  934. }
  935. if err := msg.Decode(&req); err != nil {
  936. return errResp(ErrDecode, "msg %v: %v", msg, err)
  937. }
  938. reqCnt := len(req.Txs)
  939. if reject(uint64(reqCnt), MaxTxSend) {
  940. return errResp(ErrRequestRejected, "")
  941. }
  942. hashes := make([]common.Hash, len(req.Txs))
  943. for i, tx := range req.Txs {
  944. hashes[i] = tx.Hash()
  945. }
  946. stats := pm.txStatus(hashes)
  947. for i, stat := range stats {
  948. if stat.Status == core.TxStatusUnknown {
  949. if errs := pm.txpool.AddRemotes([]*types.Transaction{req.Txs[i]}); errs[0] != nil {
  950. stats[i].Error = errs[0].Error()
  951. continue
  952. }
  953. stats[i] = pm.txStatus([]common.Hash{hashes[i]})[0]
  954. }
  955. }
  956. bv, rcost := p.fcClient.RequestProcessed(costs.baseCost + uint64(reqCnt)*costs.reqCost)
  957. pm.server.fcCostStats.update(msg.Code, uint64(reqCnt), rcost)
  958. return p.SendTxStatus(req.ReqID, bv, stats)
  959. case GetTxStatusMsg:
  960. if pm.txpool == nil {
  961. return errResp(ErrUnexpectedResponse, "")
  962. }
  963. // Transactions arrived, parse all of them and deliver to the pool
  964. var req struct {
  965. ReqID uint64
  966. Hashes []common.Hash
  967. }
  968. if err := msg.Decode(&req); err != nil {
  969. return errResp(ErrDecode, "msg %v: %v", msg, err)
  970. }
  971. reqCnt := len(req.Hashes)
  972. if reject(uint64(reqCnt), MaxTxStatus) {
  973. return errResp(ErrRequestRejected, "")
  974. }
  975. bv, rcost := p.fcClient.RequestProcessed(costs.baseCost + uint64(reqCnt)*costs.reqCost)
  976. pm.server.fcCostStats.update(msg.Code, uint64(reqCnt), rcost)
  977. return p.SendTxStatus(req.ReqID, bv, pm.txStatus(req.Hashes))
  978. case TxStatusMsg:
  979. if pm.odr == nil {
  980. return errResp(ErrUnexpectedResponse, "")
  981. }
  982. p.Log().Trace("Received tx status response")
  983. var resp struct {
  984. ReqID, BV uint64
  985. Status []txStatus
  986. }
  987. if err := msg.Decode(&resp); err != nil {
  988. return errResp(ErrDecode, "msg %v: %v", msg, err)
  989. }
  990. p.fcServer.GotReply(resp.ReqID, resp.BV)
  991. default:
  992. p.Log().Trace("Received unknown message", "code", msg.Code)
  993. return errResp(ErrInvalidMsgCode, "%v", msg.Code)
  994. }
  995. if deliverMsg != nil {
  996. err := pm.retriever.deliver(p, deliverMsg)
  997. if err != nil {
  998. p.responseErrors++
  999. if p.responseErrors > maxResponseErrors {
  1000. return err
  1001. }
  1002. }
  1003. }
  1004. return nil
  1005. }
  1006. // getAccount retrieves an account from the state based at root.
  1007. func (pm *ProtocolManager) getAccount(statedb *state.StateDB, root, hash common.Hash) (state.Account, error) {
  1008. trie, err := trie.New(root, statedb.Database().TrieDB())
  1009. if err != nil {
  1010. return state.Account{}, err
  1011. }
  1012. blob, err := trie.TryGet(hash[:])
  1013. if err != nil {
  1014. return state.Account{}, err
  1015. }
  1016. var account state.Account
  1017. if err = rlp.DecodeBytes(blob, &account); err != nil {
  1018. return state.Account{}, err
  1019. }
  1020. return account, nil
  1021. }
  1022. // getHelperTrie returns the post-processed trie root for the given trie ID and section index
  1023. func (pm *ProtocolManager) getHelperTrie(id uint, idx uint64) (common.Hash, string) {
  1024. switch id {
  1025. case htCanonical:
  1026. sectionHead := rawdb.ReadCanonicalHash(pm.chainDb, (idx+1)*light.CHTFrequencyClient-1)
  1027. return light.GetChtV2Root(pm.chainDb, idx, sectionHead), light.ChtTablePrefix
  1028. case htBloomBits:
  1029. sectionHead := rawdb.ReadCanonicalHash(pm.chainDb, (idx+1)*light.BloomTrieFrequency-1)
  1030. return light.GetBloomTrieRoot(pm.chainDb, idx, sectionHead), light.BloomTrieTablePrefix
  1031. }
  1032. return common.Hash{}, ""
  1033. }
  1034. // getHelperTrieAuxData returns requested auxiliary data for the given HelperTrie request
  1035. func (pm *ProtocolManager) getHelperTrieAuxData(req HelperTrieReq) []byte {
  1036. switch {
  1037. case req.Type == htCanonical && req.AuxReq == auxHeader && len(req.Key) == 8:
  1038. blockNum := binary.BigEndian.Uint64(req.Key)
  1039. hash := rawdb.ReadCanonicalHash(pm.chainDb, blockNum)
  1040. return rawdb.ReadHeaderRLP(pm.chainDb, hash, blockNum)
  1041. }
  1042. return nil
  1043. }
  1044. func (pm *ProtocolManager) txStatus(hashes []common.Hash) []txStatus {
  1045. stats := make([]txStatus, len(hashes))
  1046. for i, stat := range pm.txpool.Status(hashes) {
  1047. // Save the status we've got from the transaction pool
  1048. stats[i].Status = stat
  1049. // If the transaction is unknown to the pool, try looking it up locally
  1050. if stat == core.TxStatusUnknown {
  1051. if block, number, index := rawdb.ReadTxLookupEntry(pm.chainDb, hashes[i]); block != (common.Hash{}) {
  1052. stats[i].Status = core.TxStatusIncluded
  1053. stats[i].Lookup = &rawdb.TxLookupEntry{BlockHash: block, BlockIndex: number, Index: index}
  1054. }
  1055. }
  1056. }
  1057. return stats
  1058. }
  1059. // NodeInfo represents a short summary of the Ethereum sub-protocol metadata
  1060. // known about the host peer.
  1061. type NodeInfo struct {
  1062. Network uint64 `json:"network"` // Ethereum network ID (1=Frontier, 2=Morden, Ropsten=3, Rinkeby=4)
  1063. Difficulty *big.Int `json:"difficulty"` // Total difficulty of the host's blockchain
  1064. Genesis common.Hash `json:"genesis"` // SHA3 hash of the host's genesis block
  1065. Config *params.ChainConfig `json:"config"` // Chain configuration for the fork rules
  1066. Head common.Hash `json:"head"` // SHA3 hash of the host's best owned block
  1067. }
  1068. // NodeInfo retrieves some protocol metadata about the running host node.
  1069. func (self *ProtocolManager) NodeInfo() *NodeInfo {
  1070. head := self.blockchain.CurrentHeader()
  1071. hash := head.Hash()
  1072. return &NodeInfo{
  1073. Network: self.networkId,
  1074. Difficulty: self.blockchain.GetTd(hash, head.Number.Uint64()),
  1075. Genesis: self.blockchain.Genesis().Hash(),
  1076. Config: self.blockchain.Config(),
  1077. Head: hash,
  1078. }
  1079. }
  1080. // downloaderPeerNotify implements peerSetNotify
  1081. type downloaderPeerNotify ProtocolManager
  1082. type peerConnection struct {
  1083. manager *ProtocolManager
  1084. peer *peer
  1085. }
  1086. func (pc *peerConnection) Head() (common.Hash, *big.Int) {
  1087. return pc.peer.HeadAndTd()
  1088. }
  1089. func (pc *peerConnection) RequestHeadersByHash(origin common.Hash, amount int, skip int, reverse bool) error {
  1090. reqID := genReqID()
  1091. rq := &distReq{
  1092. getCost: func(dp distPeer) uint64 {
  1093. peer := dp.(*peer)
  1094. return peer.GetRequestCost(GetBlockHeadersMsg, amount)
  1095. },
  1096. canSend: func(dp distPeer) bool {
  1097. return dp.(*peer) == pc.peer
  1098. },
  1099. request: func(dp distPeer) func() {
  1100. peer := dp.(*peer)
  1101. cost := peer.GetRequestCost(GetBlockHeadersMsg, amount)
  1102. peer.fcServer.QueueRequest(reqID, cost)
  1103. return func() { peer.RequestHeadersByHash(reqID, cost, origin, amount, skip, reverse) }
  1104. },
  1105. }
  1106. _, ok := <-pc.manager.reqDist.queue(rq)
  1107. if !ok {
  1108. return ErrNoPeers
  1109. }
  1110. return nil
  1111. }
  1112. func (pc *peerConnection) RequestHeadersByNumber(origin uint64, amount int, skip int, reverse bool) error {
  1113. reqID := genReqID()
  1114. rq := &distReq{
  1115. getCost: func(dp distPeer) uint64 {
  1116. peer := dp.(*peer)
  1117. return peer.GetRequestCost(GetBlockHeadersMsg, amount)
  1118. },
  1119. canSend: func(dp distPeer) bool {
  1120. return dp.(*peer) == pc.peer
  1121. },
  1122. request: func(dp distPeer) func() {
  1123. peer := dp.(*peer)
  1124. cost := peer.GetRequestCost(GetBlockHeadersMsg, amount)
  1125. peer.fcServer.QueueRequest(reqID, cost)
  1126. return func() { peer.RequestHeadersByNumber(reqID, cost, origin, amount, skip, reverse) }
  1127. },
  1128. }
  1129. _, ok := <-pc.manager.reqDist.queue(rq)
  1130. if !ok {
  1131. return ErrNoPeers
  1132. }
  1133. return nil
  1134. }
  1135. func (d *downloaderPeerNotify) registerPeer(p *peer) {
  1136. pm := (*ProtocolManager)(d)
  1137. pc := &peerConnection{
  1138. manager: pm,
  1139. peer: p,
  1140. }
  1141. pm.downloader.RegisterLightPeer(p.id, ethVersion, pc)
  1142. }
  1143. func (d *downloaderPeerNotify) unregisterPeer(p *peer) {
  1144. pm := (*ProtocolManager)(d)
  1145. pm.downloader.UnregisterPeer(p.id)
  1146. }