handler.go 40 KB

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