mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-19 10:22:23 +00:00
Add handler functions for graphene messages
This commit is contained in:
parent
34d9744461
commit
60e3fc31f2
3 changed files with 157 additions and 5 deletions
153
eth/handler.go
153
eth/handler.go
|
|
@ -25,6 +25,8 @@ import (
|
||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
"sort"
|
||||||
|
"encoding/binary"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/consensus"
|
"github.com/ethereum/go-ethereum/consensus"
|
||||||
|
|
@ -40,6 +42,9 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/p2p/discover"
|
"github.com/ethereum/go-ethereum/p2p/discover"
|
||||||
"github.com/ethereum/go-ethereum/params"
|
"github.com/ethereum/go-ethereum/params"
|
||||||
"github.com/ethereum/go-ethereum/rlp"
|
"github.com/ethereum/go-ethereum/rlp"
|
||||||
|
"github.com/sasha-s/go-IBLT"
|
||||||
|
"github.com/willf/bloom"
|
||||||
|
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
|
|
@ -55,6 +60,23 @@ var (
|
||||||
daoChallengeTimeout = 15 * time.Second // Time allowance for a node to reply to the DAO handshake challenge
|
daoChallengeTimeout = 15 * time.Second // Time allowance for a node to reply to the DAO handshake challenge
|
||||||
)
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
grapheneC = 8 * math.Pow(math.Log(2), 2)
|
||||||
|
grapheneTau = 37.5
|
||||||
|
grapheneIBLTLookupTable = [8][2]int {
|
||||||
|
// FIXME: Remove library requirement for powers of 2
|
||||||
|
// FIXME: Remove library requiremenet for min. 2 * nHashFns number of cells
|
||||||
|
{3, 8}, //According to formula: 1 item -> 3 cells
|
||||||
|
{8, 16}, //2 -> 16
|
||||||
|
{6, 32}, //3 -> 18
|
||||||
|
{7, 32}, //4 -> 21
|
||||||
|
{6, 32}, //5 -> 24
|
||||||
|
{6, 32}, //6 -> 24
|
||||||
|
{9, 32}, //7 -> 27
|
||||||
|
{7, 32}, //8 -> 28
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
// errIncompatibleConfig is returned if the requested protocols and configs are
|
// errIncompatibleConfig is returned if the requested protocols and configs are
|
||||||
// not compatible (low protocol version restrictions and high requirements).
|
// not compatible (low protocol version restrictions and high requirements).
|
||||||
var errIncompatibleConfig = errors.New("incompatible configuration")
|
var errIncompatibleConfig = errors.New("incompatible configuration")
|
||||||
|
|
@ -94,6 +116,9 @@ type ProtocolManager struct {
|
||||||
// wait group is used for graceful shutdowns during downloading
|
// wait group is used for graceful shutdowns during downloading
|
||||||
// and processing
|
// and processing
|
||||||
wg sync.WaitGroup
|
wg sync.WaitGroup
|
||||||
|
|
||||||
|
// switch to enable graphene
|
||||||
|
useGraphene bool
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewProtocolManager returns a new Ethereum sub protocol manager. The Ethereum sub protocol manages peers capable
|
// NewProtocolManager returns a new Ethereum sub protocol manager. The Ethereum sub protocol manages peers capable
|
||||||
|
|
@ -111,6 +136,7 @@ func NewProtocolManager(config *params.ChainConfig, mode downloader.SyncMode, ne
|
||||||
noMorePeers: make(chan struct{}),
|
noMorePeers: make(chan struct{}),
|
||||||
txsyncCh: make(chan *txsync),
|
txsyncCh: make(chan *txsync),
|
||||||
quitSync: make(chan struct{}),
|
quitSync: make(chan struct{}),
|
||||||
|
useGraphene: false,
|
||||||
}
|
}
|
||||||
// Figure out whether to allow fast sync or not
|
// Figure out whether to allow fast sync or not
|
||||||
if mode == downloader.FastSync && blockchain.CurrentBlock().NumberU64() > 0 {
|
if mode == downloader.FastSync && blockchain.CurrentBlock().NumberU64() > 0 {
|
||||||
|
|
@ -628,8 +654,18 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
for _, block := range unknown {
|
for _, block := range unknown {
|
||||||
|
if pm.useGraphene {
|
||||||
|
// FIXME: Should we consider Queued too?
|
||||||
|
pending, _ := pm.txpool.Pending()
|
||||||
|
nTxs := 0
|
||||||
|
for _, txs := range pending {
|
||||||
|
nTxs += len(txs)
|
||||||
|
}
|
||||||
|
p.RequestGraphene(block.Hash, nTxs)
|
||||||
|
} else {
|
||||||
pm.fetcher.Notify(p.id, block.Hash, block.Number, time.Now(), p.RequestOneHeader, p.RequestBodies)
|
pm.fetcher.Notify(p.id, block.Hash, block.Number, time.Now(), p.RequestOneHeader, p.RequestBodies)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
case msg.Code == NewBlockMsg:
|
case msg.Code == NewBlockMsg:
|
||||||
// Retrieve and decode the propagated block
|
// Retrieve and decode the propagated block
|
||||||
|
|
@ -682,6 +718,123 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
|
||||||
}
|
}
|
||||||
pm.txpool.AddRemotes(txs)
|
pm.txpool.AddRemotes(txs)
|
||||||
|
|
||||||
|
case msg.Code == GetTxMsg:
|
||||||
|
var hashes []common.Hash
|
||||||
|
if err := msg.Decode(&hashes); err != nil {
|
||||||
|
return errResp(ErrDecode, "msg %v: %v", msg, err)
|
||||||
|
}
|
||||||
|
var txs []*types.Transaction
|
||||||
|
for _, hash := range hashes {
|
||||||
|
log.Info("Retrieving transaction", "hash", hash)
|
||||||
|
// FIXME: this function has been removed, but we don't this msg currently
|
||||||
|
// txs = append(txs, pm.blockchain.GetTxByHash(hash))
|
||||||
|
}
|
||||||
|
if len(txs) > 0 {
|
||||||
|
p.SendTransactions(txs)
|
||||||
|
}
|
||||||
|
|
||||||
|
case msg.Code == GetGrapheneMsg:
|
||||||
|
var query getGrapheneData
|
||||||
|
msg.Decode(&query)
|
||||||
|
block := pm.blockchain.GetBlockByHash(query.Hash)
|
||||||
|
|
||||||
|
nPendingAtRecv := query.NTxs
|
||||||
|
txs := block.Transactions()
|
||||||
|
nTxsInBlock := uint(len(txs))
|
||||||
|
if nTxsInBlock == 0 {
|
||||||
|
log.Debug("GetGraphene for empty block", "hash", query.Hash)
|
||||||
|
// p.SendGraphene(query.Hash, empty, empty, 0, 0)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
expectedDiff := float64(nTxsInBlock) / (grapheneC * grapheneTau)
|
||||||
|
bloomFPR := expectedDiff / float64((nPendingAtRecv - nTxsInBlock) + 1)
|
||||||
|
bloomFilter := bloom.NewWithEstimates(nTxsInBlock, bloomFPR)
|
||||||
|
|
||||||
|
ibltRow := grapheneIBLTLookupTable[int(math.Ceil(expectedDiff))]
|
||||||
|
senderIBLT := iblt.New(ibltRow[0], ibltRow[1])
|
||||||
|
txIds := make([]string, nTxsInBlock)
|
||||||
|
txPositions := make(map[string]uint32)
|
||||||
|
|
||||||
|
for i, tx := range txs {
|
||||||
|
txHash := tx.Hash().Bytes()[:5]
|
||||||
|
txHashString := string(txHash[:])
|
||||||
|
bloomFilter.Add(txHash)
|
||||||
|
senderIBLT.Add(txHash)
|
||||||
|
txIds = append(txIds, txHashString)
|
||||||
|
txPositions[txHashString] = uint32(i)
|
||||||
|
}
|
||||||
|
|
||||||
|
sort.Strings(txIds)
|
||||||
|
indexArray := make([]byte, nTxsInBlock * 4)
|
||||||
|
for _, txId := range txIds {
|
||||||
|
bs := make([]byte, 4)
|
||||||
|
binary.LittleEndian.PutUint32(bs, txPositions[txId])
|
||||||
|
indexArray = append(indexArray, bs...)
|
||||||
|
}
|
||||||
|
|
||||||
|
ibltBinary, _ := senderIBLT.MarshalBinary()
|
||||||
|
bloomFilterBinary, _ := bloomFilter.GobEncode()
|
||||||
|
|
||||||
|
// include nPendingAtRecv in the response too so we don't have to
|
||||||
|
// calculate it again while recomputing FPR (which can't be serialized because it is a float)
|
||||||
|
p.SendGraphene(query.Hash, ibltBinary, bloomFilterBinary, uint(nPendingAtRecv), uint(nTxsInBlock), indexArray, block.Uncles())
|
||||||
|
|
||||||
|
case msg.Code == GrapheneMsg:
|
||||||
|
var body grapheneData
|
||||||
|
msg.Decode(&body)
|
||||||
|
|
||||||
|
// This recomputation is necessary because floats cannot be serialized
|
||||||
|
expectedDiff := float64(body.NTxs) / (grapheneC * grapheneTau)
|
||||||
|
bloomFPR := expectedDiff / float64((body.NPending - body.NTxs) + 1)
|
||||||
|
bloomFilter := bloom.NewWithEstimates(body.NTxs, bloomFPR)
|
||||||
|
bloomFilter.GobDecode(body.GrapheneBloom)
|
||||||
|
|
||||||
|
ibltRow := grapheneIBLTLookupTable[int(math.Ceil(expectedDiff))]
|
||||||
|
receivedIBLT := iblt.New(ibltRow[0], ibltRow[1])
|
||||||
|
receivedIBLT.UnmarshalBinary(body.GrapheneIBLT)
|
||||||
|
|
||||||
|
localIBLT := iblt.New(ibltRow[0], ibltRow[1])
|
||||||
|
|
||||||
|
nAddedToBloom := 0
|
||||||
|
nTested := 0
|
||||||
|
pending, _ := pm.txpool.Pending()
|
||||||
|
shortTxHashMap := make(map[string]*types.Transaction)
|
||||||
|
|
||||||
|
for _, txs := range pending {
|
||||||
|
for _, tx := range txs {
|
||||||
|
nTested += 1
|
||||||
|
txHash := tx.Hash().Bytes()[:5]
|
||||||
|
txHashString := string(txHash[:])
|
||||||
|
shortTxHashMap[txHashString] = tx
|
||||||
|
if bloomFilter.Test(txHash) {
|
||||||
|
nAddedToBloom += 1
|
||||||
|
localIBLT.Add(txHash)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
localIBLT.Sub(*receivedIBLT)
|
||||||
|
decodedIBLT, _ := localIBLT.Decode()
|
||||||
|
|
||||||
|
for _, txHash := range decodedIBLT.Added {
|
||||||
|
txHashString := string(txHash[:])
|
||||||
|
shortTxHashMap[txHashString] = nil
|
||||||
|
}
|
||||||
|
|
||||||
|
txIds := make([]string, body.NTxs)
|
||||||
|
for k, _ := range shortTxHashMap {
|
||||||
|
txIds = append(txIds, k)
|
||||||
|
}
|
||||||
|
sort.Strings(txIds)
|
||||||
|
|
||||||
|
txs := make([]*types.Transaction, body.NTxs)
|
||||||
|
for i, txId := range txIds {
|
||||||
|
index := binary.LittleEndian.Uint32(body.Indices[4*i:4*i+4])
|
||||||
|
txs[index] = shortTxHashMap[txId]
|
||||||
|
}
|
||||||
|
|
||||||
default:
|
default:
|
||||||
return errResp(ErrInvalidMsgCode, "%v", msg.Code)
|
return errResp(ErrInvalidMsgCode, "%v", msg.Code)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -340,9 +340,9 @@ func (p *peer) RequestGraphene(hash common.Hash, nTxs int) error {
|
||||||
}
|
}
|
||||||
|
|
||||||
// SendGraphene sends the graphene message in response to a RequestGraphene
|
// SendGraphene sends the graphene message in response to a RequestGraphene
|
||||||
func (p *peer) SendGraphene(hash common.Hash, i []byte, b []byte, ni uint, fpr uint, nTxs uint, indexArray []byte, uncles []*types.Header) error {
|
func (p *peer) SendGraphene(hash common.Hash, i []byte, b []byte, nPending uint, nTxs uint, indexArray []byte, uncles []*types.Header) error {
|
||||||
p.Log().Debug("Sending graphene", "block", hash)
|
p.Log().Debug("Sending graphene", "block", hash)
|
||||||
return p2p.Send(p.rw, GrapheneMsg, &grapheneData{Hash: hash, GrapheneIBLT: i, GrapheneBloom: b, NIBLT: ni, FPR: fpr, NTxs: nTxs, Indices: indexArray, Uncles: uncles})
|
return p2p.Send(p.rw, GrapheneMsg, &grapheneData{Hash: hash, GrapheneIBLT: i, GrapheneBloom: b, NPending: nPending, NTxs: nTxs, Indices: indexArray, Uncles: uncles})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -196,8 +196,7 @@ type getGrapheneData struct {
|
||||||
type grapheneData struct {
|
type grapheneData struct {
|
||||||
GrapheneIBLT []byte
|
GrapheneIBLT []byte
|
||||||
GrapheneBloom []byte
|
GrapheneBloom []byte
|
||||||
FPR uint
|
NPending uint
|
||||||
NIBLT uint
|
|
||||||
NTxs uint
|
NTxs uint
|
||||||
Hash common.Hash
|
Hash common.Hash
|
||||||
Indices []byte
|
Indices []byte
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue