Merge pull request #2 from optimism-java/feature_offer

Feature offer
This commit is contained in:
Chen Kai 2023-11-26 20:59:14 +08:00 committed by GitHub
commit 632c815ffa
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
14 changed files with 495 additions and 168 deletions

8
go.mod
View file

@ -54,8 +54,10 @@ require (
github.com/mattn/go-isatty v0.0.17
github.com/naoina/toml v0.1.2-0.20170918210437-9fafd6967416
github.com/olekukonko/tablewriter v0.0.5
github.com/optimism-java/utp-go v0.0.0-20231114114639-92925ba7e35e
github.com/peterh/liner v1.1.1-0.20190123174540-a2c9a5303de7
github.com/protolambda/bls12-381-util v0.0.0-20220416220906-d8552aa452c7
github.com/prysmaticlabs/go-bitfield v0.0.0-20210809151128-385d8c5e3fb7
github.com/rs/cors v1.7.0
github.com/shirou/gopsutil v3.21.4-0.20210419000835-c7a38de76ee5+incompatible
github.com/status-im/keycard-go v0.2.0
@ -65,6 +67,7 @@ require (
github.com/tyler-smith/go-bip39 v1.1.0
github.com/urfave/cli/v2 v2.25.7
go.uber.org/automaxprocs v1.5.2
go.uber.org/zap v1.26.0
golang.org/x/crypto v0.14.0
golang.org/x/exp v0.0.0-20230905200255-921286631fa9
golang.org/x/sync v0.4.0
@ -127,23 +130,20 @@ require (
github.com/mmcloughlin/addchain v0.4.0 // indirect
github.com/naoina/go-stringutil v0.1.0 // indirect
github.com/opentracing/opentracing-go v1.1.0 // indirect
github.com/optimism-java/utp-go v0.0.0-20231114114639-92925ba7e35e // indirect
github.com/pkg/errors v0.9.1 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect
github.com/prometheus/client_golang v1.12.0 // indirect
github.com/prometheus/client_model v0.2.1-0.20210607210712-147c58e9608a // indirect
github.com/prometheus/common v0.32.1 // indirect
github.com/prometheus/procfs v0.7.3 // indirect
github.com/prysmaticlabs/go-bitfield v0.0.0-20210809151128-385d8c5e3fb7 // indirect
github.com/rivo/uniseg v0.2.0 // indirect
github.com/rogpeppe/go-internal v1.9.0 // indirect
github.com/russross/blackfriday/v2 v2.1.0 // indirect
github.com/tetratelabs/wabin v0.0.0-20230304001439-f6f874872834 // indirect
github.com/tklauser/go-sysconf v0.3.12 // indirect
github.com/tklauser/numcpus v0.6.1 // indirect
github.com/xrash/smetrics v0.0.0-20201216005158-039620a65673 // indirect
go.uber.org/atomic v1.7.0 // indirect
go.uber.org/multierr v1.11.0 // indirect
go.uber.org/zap v1.26.0 // indirect
golang.org/x/mod v0.13.0 // indirect
golang.org/x/net v0.17.0 // indirect
google.golang.org/protobuf v1.27.1 // indirect

22
go.sum
View file

@ -497,10 +497,6 @@ github.com/onsi/gomega v1.10.1 h1:o0+MgICZLuZ7xjH7Vx6zS/zcu93/BEp1VwkIW1mEXCE=
github.com/onsi/gomega v1.10.1/go.mod h1:iN09h71vgCQne3DLsj+A5owkum+a2tYe+TOCB1ybHNo=
github.com/opentracing/opentracing-go v1.1.0 h1:pWlfV3Bxv7k65HYwkikxat0+s3pV4bsqf19k25Ur8rU=
github.com/opentracing/opentracing-go v1.1.0/go.mod h1:UkNAQd3GIcIGf0SeVgPpRdFStlNbqXla1AfSYxPUl2o=
github.com/optimism-java/utp-go v0.0.0-20231030043430-a1331c25fa98 h1:uxUbd8LFc24XetNFjTu9Kp9MqF2zKF92UMbcDuPxYZ8=
github.com/optimism-java/utp-go v0.0.0-20231030043430-a1331c25fa98/go.mod h1:DZ0jYzLzt4ZsCmhI/iqYgGFoNx45OfpEoKzXB8HVALQ=
github.com/optimism-java/utp-go v0.0.0-20231111152515-b2c1e9aba225 h1:UUVmsVAv/4v0TMW3AFPwxAkuyNjKsAWVyqxfR10gMGE=
github.com/optimism-java/utp-go v0.0.0-20231111152515-b2c1e9aba225/go.mod h1:DZ0jYzLzt4ZsCmhI/iqYgGFoNx45OfpEoKzXB8HVALQ=
github.com/optimism-java/utp-go v0.0.0-20231114114639-92925ba7e35e h1:61Mw2nE4trMg/Ze/0oFD5o5yKog2UuLUlbUnBS0lCh8=
github.com/optimism-java/utp-go v0.0.0-20231114114639-92925ba7e35e/go.mod h1:DZ0jYzLzt4ZsCmhI/iqYgGFoNx45OfpEoKzXB8HVALQ=
github.com/pelletier/go-toml v1.2.0/go.mod h1:5z9KED0ma1S8pY6P1sdut58dfprrGBbd/94hg7ilaic=
@ -587,6 +583,8 @@ github.com/supranational/blst v0.3.11 h1:LyU6FolezeWAhvQk0k6O/d49jqgO52MSDDfYgbe
github.com/supranational/blst v0.3.11/go.mod h1:jZJtfjgudtNl4en1tzwPIV3KjUnQUvG3/j+w+fVonLw=
github.com/syndtr/goleveldb v1.0.1-0.20210819022825-2ae1ddf74ef7 h1:epCh84lMvA70Z7CTTCmYQn2CKbY8j86K7/FAIr141uY=
github.com/syndtr/goleveldb v1.0.1-0.20210819022825-2ae1ddf74ef7/go.mod h1:q4W45IWZaF22tdD+VEXcAWRA037jwmWEB5VWYORlTpc=
github.com/tetratelabs/wabin v0.0.0-20230304001439-f6f874872834 h1:ZF+QBjOI+tILZjBaFj3HgFonKXUcwgJ4djLb6i42S3Q=
github.com/tetratelabs/wabin v0.0.0-20230304001439-f6f874872834/go.mod h1:m9ymHTgNSEjuxvw8E7WWe4Pl4hZQHXONY8wE6dMLaRk=
github.com/tklauser/go-sysconf v0.3.12 h1:0QaGUFOdQaIVdPgfITYzaTegZvdCjmYO52cSFAEVmqU=
github.com/tklauser/go-sysconf v0.3.12/go.mod h1:Ho14jnntGE1fpdOqQEEaiKRpvIavV0hSfmBq8nJbHYI=
github.com/tklauser/numcpus v0.6.1 h1:ng9scYS7az0Bk4OZLvrNXNSAO2Pxr1XXRAPyjhIx+Fk=
@ -628,12 +626,10 @@ go.uber.org/atomic v1.7.0 h1:ADUqmZGgLDDfbSL9ZmPxKTybcoEYHgpYfELNoN+7hsw=
go.uber.org/atomic v1.7.0/go.mod h1:fEN4uk6kAWBTFdckzkM89CLk9XfWZrxpCo0nPH17wJc=
go.uber.org/automaxprocs v1.5.2 h1:2LxUOGiR3O6tw8ui5sZa2LAaHnsviZdVOUZw4fvbnME=
go.uber.org/automaxprocs v1.5.2/go.mod h1:eRbA25aqJrxAbsLO0xy5jVwPt7FQnRgjW+efnwa1WM0=
go.uber.org/goleak v1.1.10/go.mod h1:8a7PlsEVH3e/a/GLqe5IIrQx6GzcnRmZEufDUTk4A7A=
go.uber.org/multierr v1.6.0 h1:y6IPFStTAIT5Ytl7/XYmHvzXQ7S3g/IeZW9hyZ5thw4=
go.uber.org/multierr v1.6.0/go.mod h1:cdWPpRnG4AhwMwsgIHip0KRBQjJy5kYEpYjJxpXp9iU=
go.uber.org/goleak v1.2.0 h1:xqgm/S+aQvhWFTtR0XK3Jvg7z8kGV8P4X14IzwN3Eqk=
go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0=
go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y=
go.uber.org/zap v1.19.0 h1:mZQZefskPPCMIBCSEH0v2/iUqqLrYtaeqwD6FUGUnFE=
go.uber.org/zap v1.19.0/go.mod h1:xg/QME4nWcxGxrpdeYfq7UvYrLh66cuVKdrbD1XF/NI=
go.uber.org/zap v1.26.0 h1:sI7k6L95XOKS281NhVKOFCUNIvv9e0w4BF8N3u+tCRo=
go.uber.org/zap v1.26.0/go.mod h1:dtElttAiwGvoJ/vj4IwHBS/gXsEu/pZ50mUIRWuG0so=
golang.org/x/crypto v0.0.0-20180904163835-0709b304e793/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4=
golang.org/x/crypto v0.0.0-20181203042331-505ab145d0a9/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4=
@ -681,8 +677,7 @@ golang.org/x/mod v0.1.1-0.20191107180719-034126e5016b/go.mod h1:QqPTAvyqsEbceGzB
golang.org/x/mod v0.2.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA=
golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA=
golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4=
golang.org/x/mod v0.12.0 h1:rmsUpXtvNzj340zd98LZ4KntptpfRHwpFOHG188oHXc=
golang.org/x/mod v0.12.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
golang.org/x/mod v0.13.0 h1:I/DsJXRlw/8l/0c24sM9yb0T4z9liZTduXvdAWYiysY=
golang.org/x/mod v0.13.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c=
golang.org/x/net v0.0.0-20180724234803-3673e40ba225/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
golang.org/x/net v0.0.0-20180826012351-8a410e7b638d/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
@ -807,8 +802,6 @@ golang.org/x/sys v0.0.0-20220908164124-27713097b956/go.mod h1:oPkhp1MJrh7nUepCBc
golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.11.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.13.0 h1:Af8nKPmuFypiUBjVoU9V20FiaFXOcuZI21p0ycVYYGE=
golang.org/x/sys v0.13.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.14.0 h1:Vz7Qs629MkJkGyHxUlRHizWJRG2j8fbQKjELVSNhy7Q=
golang.org/x/sys v0.14.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/term v0.0.0-20201117132131-f5c789dd3221/go.mod h1:Nr5EML6q2oocZ2LXRh80K7BxOlk5/8JxuGnuhpl+muw=
@ -880,8 +873,7 @@ golang.org/x/tools v0.0.0-20200804011535-6c149bb5ef0d/go.mod h1:njjCfa9FT2d7l9Bc
golang.org/x/tools v0.0.0-20200825202427-b303f430e36d/go.mod h1:njjCfa9FT2d7l9Bc6FUM5FLjQPp3cFF28FI3qnDFljA=
golang.org/x/tools v0.0.0-20210106214847-113979e3529a/go.mod h1:emZCQorbCU4vsT4fOWvOPXz4eW1wZW4PmDk9uLelYpA=
golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc=
golang.org/x/tools v0.13.0 h1:Iey4qkscZuv0VvIt8E0neZjtPVQFSc870HQ448QgEmQ=
golang.org/x/tools v0.13.0/go.mod h1:HvlwmtVNQAhOuCjW7xxvovg8wbNq7LwfXh/k7wXUl58=
golang.org/x/tools v0.14.0 h1:jvNa2pY0M4r62jkRQ6RwEZZyPcymeL9XZMLBbV7U2nc=
golang.org/x/tools v0.14.0/go.mod h1:uYBEerGOWcJyEORxN+Ek8+TT266gXkNlHdJBwexUsBg=
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=

View file

@ -49,7 +49,7 @@ type Config struct {
// Node table configuration:
Bootnodes []*enode.Node // list of bootstrap nodes
PingInterval time.Duration // speed of node liveness check
PingInterval time.Duration // speed of Node liveness check
RefreshInterval time.Duration // used in bucket refresh
// The options below are useful in very specific cases, like in unit tests.

View file

@ -149,7 +149,7 @@ func (it *lookup) query(n *node, reply chan<- []*node) {
} else if len(r) == 0 {
fails++
it.tab.db.UpdateFindFails(n.ID(), n.IP(), fails)
// Remove the node from the local table if it fails to return anything useful too
// Remove the Node from the local table if it fails to return anything useful too
// many times, but only if there are enough other nodes in the bucket.
dropped := false
if fails >= maxFindnodeFailures && it.tab.bucketLen(n.ID()) >= bucketSize/2 {
@ -187,7 +187,7 @@ func newLookupIterator(ctx context.Context, next lookupFunc) *lookupIterator {
return &lookupIterator{ctx: ctx, cancel: cancel, nextLookup: next}
}
// Node returns the current node.
// Node returns the current Node.
func (it *lookupIterator) Node() *enode.Node {
if len(it.buffer) == 0 {
return nil
@ -195,9 +195,9 @@ func (it *lookupIterator) Node() *enode.Node {
return unwrapNode(it.buffer[0])
}
// Next moves to the next node.
// Next moves to the next Node.
func (it *lookupIterator) Next() bool {
// Consume next node in buffer.
// Consume next Node in buffer.
if len(it.buffer) > 0 {
it.buffer = it.buffer[1:]
}

View file

@ -1,6 +1,7 @@
package discover
import (
"bytes"
"context"
"crypto/ecdsa"
crand "crypto/rand"
@ -13,6 +14,7 @@ import (
"sort"
"time"
"github.com/tetratelabs/wabin/leb128"
"go.uber.org/zap"
"github.com/VictoriaMetrics/fastcache"
@ -55,6 +57,35 @@ const (
defaultUTPReadTimeout = 60 * time.Second
)
const (
TransientOfferRequestKind byte = 0x01
PersistOfferRequestKind byte = 0x02
)
type ContentElement struct {
Node enode.ID
ContentKeys [][]byte
Contents [][]byte
}
type ContentEntry struct {
ContentKey []byte
Content []byte
}
type TransientOfferRequest struct {
Contents []*ContentEntry
}
type PersistOfferRequest struct {
ContentKeys [][]byte
}
type OfferRequest struct {
Kind byte
Request interface{}
}
type PortalProtocolConfig struct {
BootstrapNodes []*enode.Node
@ -100,9 +131,11 @@ type PortalProtocol struct {
closeCtx context.Context
cancelCloseCtx context.CancelFunc
storage Storage
contentQueue chan *ContentElement
}
func NewPortalProtocol(config *PortalProtocolConfig, protocolId string, privateKey *ecdsa.PrivateKey, storage Storage) (*PortalProtocol, error) {
func NewPortalProtocol(config *PortalProtocolConfig, protocolId string, privateKey *ecdsa.PrivateKey, storage Storage, contentQueue chan *ContentElement) (*PortalProtocol, error) {
nodeDB, err := enode.OpenDB(config.NodeDBPath)
if err != nil {
return nil, err
@ -126,6 +159,7 @@ func NewPortalProtocol(config *PortalProtocolConfig, protocolId string, privateK
localNode: localNode,
validSchemes: enode.ValidSchemes,
storage: storage,
contentQueue: contentQueue,
}
return protocol, nil
@ -331,6 +365,152 @@ func (p *PortalProtocol) findContent(node *enode.Node, contentKey []byte) (byte,
return p.processContent(node, talkResp)
}
func (p *PortalProtocol) offer(node *enode.Node, offerRequest *OfferRequest) ([]byte, error) {
contentKeys := getContentKeys(offerRequest)
offer := &portalwire.Offer{
ContentKeys: contentKeys,
}
p.log.Trace("Sending offer request", "offer", offer)
offerBytes, err := offer.MarshalSSZ()
if err != nil {
p.log.Error("failed to marshal offer request", "err", err)
return nil, err
}
talkRequestBytes := make([]byte, 0, len(offerBytes)+1)
talkRequestBytes = append(talkRequestBytes, portalwire.OFFER)
talkRequestBytes = append(talkRequestBytes, offerBytes...)
talkResp, err := p.DiscV5.TalkRequest(node, p.protocolId, talkRequestBytes)
if err != nil {
p.log.Error("failed to send offer request", "err", err)
return nil, err
}
return p.processOffer(node, talkResp, offerRequest)
}
func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request *OfferRequest) ([]byte, error) {
var err error
if resp[0] != portalwire.ACCEPT {
return nil, fmt.Errorf("invalid accept response")
}
accept := &portalwire.Accept{}
err = accept.UnmarshalSSZ(resp[1:])
if err != nil {
return nil, err
}
p.log.Trace("Received accept response", "id", target.ID(), "accept", accept)
var contentKeyLen int
if request.Kind == TransientOfferRequestKind {
contentKeyLen = len(request.Request.(*TransientOfferRequest).Contents)
} else {
contentKeyLen = len(request.Request.(*PersistOfferRequest).ContentKeys)
}
contentKeyBitlist := bitfield.Bitlist(accept.ContentKeys)
if int(contentKeyBitlist.Count()) != contentKeyLen {
return nil, fmt.Errorf("accepted content key bitlist has invalid size, expected %d, got %d", contentKeyLen, contentKeyBitlist.Len())
}
if contentKeyBitlist.Count() == 0 {
return nil, nil
}
connId := binary.BigEndian.Uint16(accept.ConnectionId[:])
go func(ctx context.Context) {
var conn net.Conn
for {
select {
case <-ctx.Done():
return
default:
contents := make([][]byte, 0, contentKeyBitlist.Count())
var content []byte
if request.Kind == TransientOfferRequestKind {
for _, index := range contentKeyBitlist.BitIndices() {
content = request.Request.(*TransientOfferRequest).Contents[index].Content
contents = append(contents, content)
}
} else {
for _, index := range contentKeyBitlist.BitIndices() {
contentKey := request.Request.(*PersistOfferRequest).ContentKeys[index]
contentId := p.storage.ContentId(contentKey)
if contentId != nil {
content, err = p.storage.Get(contentKey, contentId)
if err != nil {
p.log.Error("failed to get content from storage", "err", err)
contents = append(contents, []byte{})
} else {
contents = append(contents, content)
}
} else {
contents = append(contents, []byte{})
}
}
}
var contentsPayload []byte
contentsPayload, err = encodeContents(contents)
if err != nil {
p.log.Error("failed to encode contents", "err", err)
return
}
connctx, conncancel := context.WithTimeout(ctx, defaultUTPConnectTimeout)
laddr := p.utp.Addr().(*utp.Addr)
raddr := &utp.Addr{IP: target.IP(), Port: target.UDP()}
conn, err = utp.DialUTPOptions("utp", laddr, raddr, utp.WithContext(connctx), utp.WithSocketManager(p.utpSm), utp.WithConnId(uint32(connId)))
if err != nil {
conncancel()
p.log.Error("failed to dial utp connection", "err", err)
return
}
conncancel()
err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout))
if err != nil {
p.log.Error("failed to set write deadline", "err", err)
err = conn.Close()
if err != nil {
p.log.Error("failed to close utp connection", "err", err)
return
}
return
}
var written int
written, err = conn.Write(contentsPayload)
if err != nil {
p.log.Error("failed to write to utp connection", "err", err)
err = conn.Close()
if err != nil {
p.log.Error("failed to close utp connection", "err", err)
return
}
return
}
p.log.Trace("Sent content response", "id", target.ID(), "contents", contents, "size", written)
err = conn.Close()
if err != nil {
p.log.Error("failed to close utp connection", "err", err)
return
}
return
}
}
}(p.closeCtx)
return accept.ContentKeys, nil
}
func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, interface{}, error) {
if resp[0] != portalwire.CONTENT {
return 0xff, nil, fmt.Errorf("invalid content response")
@ -354,7 +534,7 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte,
}
p.log.Trace("Received content response", "id", target.ID(), "connIdMsg", connIdMsg)
connctx, conncancel := context.WithTimeout(context.Background(), defaultUTPConnectTimeout)
connctx, conncancel := context.WithTimeout(p.closeCtx, defaultUTPConnectTimeout)
laddr := p.utp.Addr().(*utp.Addr)
raddr := &utp.Addr{IP: target.IP(), Port: target.UDP()}
connId := binary.BigEndian.Uint16(connIdMsg.Id[:])
@ -373,18 +553,18 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte,
data := make([]byte, 0)
buf := make([]byte, 1024)
for {
var n int
n, err = conn.Read(buf)
var read int
read, err = conn.Read(buf)
if err != nil {
if errors.Is(err, io.EOF) {
p.log.Trace("Received content response", "id", target.ID(), "data", data, "size", n)
p.log.Trace("Received content response", "id", target.ID(), "data", data, "size", read)
return resp[1], data, nil
}
p.log.Error("failed to read from utp connection", "err", err)
return 0xff, nil, err
}
data = append(data, buf[:n]...)
data = append(data, buf[:read]...)
}
case portalwire.ContentEnrsSelector:
enrs := &portalwire.Enrs{}
@ -712,34 +892,56 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque
connId := connIdGen.GenCid(id, false)
connIdSend := connId.SendId()
go func() {
ctx, cancel := context.WithTimeout(context.Background(), defaultUTPConnectTimeout)
var conn *utp.Conn
conn, err = p.utp.AcceptUTPContext(ctx, connIdSend)
defer func(conn *utp.Conn) {
err = conn.Close()
if err != nil {
p.log.Error("failed to close utp connection", "err", err)
}
}(conn)
if err != nil {
p.log.Error("failed to accept utp connection", "connId", connIdSend, "err", err)
cancel()
return
}
cancel()
go func(bctx context.Context) {
for {
select {
case <-bctx.Done():
return
default:
ctx, cancel := context.WithTimeout(bctx, defaultUTPConnectTimeout)
var conn *utp.Conn
conn, err = p.utp.AcceptUTPContext(ctx, connIdSend)
if err != nil {
p.log.Error("failed to accept utp connection", "connId", connIdSend, "err", err)
cancel()
return
}
cancel()
wctx, wcancel := context.WithTimeout(context.Background(), defaultUTPWriteTimeout)
var n int
n, err = conn.WriteContext(wctx, content)
if err != nil {
p.log.Error("failed to write content to utp connection", "err", err)
wcancel()
return
err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout))
if err != nil {
p.log.Error("failed to set write deadline", "err", err)
err = conn.Close()
if err != nil {
p.log.Error("failed to close utp connection", "err", err)
return
}
return
}
var n int
n, err = conn.Write(content)
if err != nil {
p.log.Error("failed to write content to utp connection", "err", err)
err = conn.Close()
if err != nil {
p.log.Error("failed to close utp connection", "err", err)
return
}
return
}
err = conn.Close()
if err != nil {
p.log.Error("failed to close utp connection", "err", err)
return
}
p.log.Trace("wrote content size to utp connection", "n", n)
return
}
}
wcancel()
p.log.Trace("wrote content size to utp connection", "n", n)
}()
}(p.closeCtx)
idBuffer := make([]byte, 2)
binary.BigEndian.PutUint16(idBuffer, uint16(connIdSend))
@ -768,14 +970,34 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque
func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *portalwire.Offer) ([]byte, error) {
var err error
contentKeyBitsets := bitfield.NewBitlist(uint64(len(request.ContentKeys)))
contentKeyBitlist := bitfield.NewBitlist(uint64(len(request.ContentKeys)))
if len(p.contentQueue) >= cap(p.contentQueue) {
acceptMsg := &portalwire.Accept{
ConnectionId: []byte{0, 0},
ContentKeys: []byte(contentKeyBitlist),
}
p.log.Trace("Sending accept response", "protocol", p.protocolId, "source", addr, "accept", acceptMsg)
var acceptMsgBytes []byte
acceptMsgBytes, err = acceptMsg.MarshalSSZ()
if err != nil {
return nil, err
}
talkRespBytes := make([]byte, 0, len(acceptMsgBytes)+1)
talkRespBytes = append(talkRespBytes, portalwire.ACCEPT)
talkRespBytes = append(talkRespBytes, acceptMsgBytes...)
return talkRespBytes, nil
}
contentKeys := make([][]byte, 0)
for i, contentKey := range request.ContentKeys {
contentId := p.storage.ContentId(contentKey)
if contentId != nil {
if p.inRange(p.Self().ID(), p.nodeRadius, contentId) {
if inRange(p.Self().ID(), p.nodeRadius, contentId) {
if _, err = p.storage.Get(contentKey, contentId); err != nil {
contentKeyBitsets.SetBitAt(uint64(i), true)
contentKeyBitlist.SetBitAt(uint64(i), true)
contentKeys = append(contentKeys, contentKey)
}
}
@ -785,53 +1007,60 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po
}
idBuffer := make([]byte, 2)
if contentKeyBitsets.Count() != 0 {
if contentKeyBitlist.Count() != 0 {
connIdGen := utp.NewConnIdGenerator()
connId := connIdGen.GenCid(id, false)
connIdSend := connId.SendId()
go func() {
ctx, cancel := context.WithTimeout(context.Background(), defaultUTPConnectTimeout)
var conn *utp.Conn
conn, err = p.utp.AcceptUTPContext(ctx, connIdSend)
defer func(conn *utp.Conn) {
err = conn.Close()
if err != nil {
p.log.Error("failed to close utp connection", "err", err)
}
}(conn)
if err != nil {
p.log.Error("failed to accept utp connection", "connId", connIdSend, "err", err)
cancel()
return
}
cancel()
err = conn.SetReadDeadline(time.Now().Add(defaultUTPReadTimeout))
if err != nil {
p.log.Error("failed to set read deadline", "err", err)
return
}
// Read ALL the data from the connection until EOF and return it
data := make([]byte, 0)
buf := make([]byte, 1024)
go func(bctx context.Context) {
for {
var n int
n, err = conn.Read(buf)
if err != nil {
if errors.Is(err, io.EOF) {
p.log.Trace("Received content response", "id", id, "data", data, "size", n)
break
select {
case <-bctx.Done():
return
default:
ctx, cancel := context.WithTimeout(bctx, defaultUTPConnectTimeout)
var conn *utp.Conn
conn, err = p.utp.AcceptUTPContext(ctx, connIdSend)
if err != nil {
p.log.Error("failed to accept utp connection", "connId", connIdSend, "err", err)
cancel()
return
}
cancel()
err = conn.SetReadDeadline(time.Now().Add(defaultUTPReadTimeout))
if err != nil {
p.log.Error("failed to set read deadline", "err", err)
return
}
// Read ALL the data from the connection until EOF and return it
data := make([]byte, 0)
buf := make([]byte, 1024)
for {
var n int
n, err = conn.Read(buf)
if err != nil {
if errors.Is(err, io.EOF) {
p.log.Trace("Received content response", "id", id, "data", data, "size", n)
break
}
p.log.Error("failed to read from utp connection", "err", err)
return
}
data = append(data, buf[:n]...)
}
err = p.handleOfferedContents(id, contentKeys, data)
if err != nil {
p.log.Error("failed to handle offered Contents", "err", err)
return
}
p.log.Error("failed to read from utp connection", "err", err)
return
}
data = append(data, buf[:n]...)
}
p.handleOfferedContents(id, contentKeys, data)
}()
}(p.closeCtx)
binary.BigEndian.PutUint16(idBuffer, uint16(connIdSend))
} else {
@ -840,7 +1069,7 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po
acceptMsg := &portalwire.Accept{
ConnectionId: idBuffer,
ContentKeys: contentKeyBitsets.BytesNoTrim(),
ContentKeys: []byte(contentKeyBitlist),
}
p.log.Trace("Sending accept response", "protocol", p.protocolId, "source", addr, "accept", acceptMsg)
@ -857,8 +1086,27 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po
return talkRespBytes, nil
}
func (p *PortalProtocol) handleOfferedContents(id enode.ID, keys [][]byte, data []byte) {
func (p *PortalProtocol) handleOfferedContents(id enode.ID, keys [][]byte, payload []byte) error {
contents, err := decodeContents(payload)
if err != nil {
return err
}
keyLen := len(keys)
contentLen := len(contents)
if keyLen != contentLen {
return fmt.Errorf("content keys len %d doesn't match content values len %d", keyLen, contentLen)
}
contentElement := &ContentElement{
Node: id,
ContentKeys: keys,
Contents: contents,
}
p.contentQueue <- contentElement
return nil
}
func (p *PortalProtocol) Self() *enode.Node {
@ -990,8 +1238,62 @@ func (p *PortalProtocol) findNodesCloseToContent(contentId []byte) []*enode.Node
return allNodes
}
func (p *PortalProtocol) inRange(nodeId enode.ID, nodeRadius *uint256.Int, contentId []byte) bool {
func inRange(nodeId enode.ID, nodeRadius *uint256.Int, contentId []byte) bool {
distance := enode.LogDist(nodeId, enode.ID(contentId))
disBig := new(big.Int).SetInt64(int64(distance))
return nodeRadius.CmpBig(disBig) > 0
}
func encodeContents(contents [][]byte) ([]byte, error) {
contentsBytes := make([]byte, 0)
for _, content := range contents {
contentLen := len(content)
contentLenBytes := leb128.EncodeUint32(uint32(contentLen))
contentsBytes = append(contentsBytes, contentLenBytes...)
contentsBytes = append(contentsBytes, content...)
}
return contentsBytes, nil
}
func decodeContents(payload []byte) ([][]byte, error) {
contents := make([][]byte, 0)
buffer := bytes.NewBuffer(payload)
for {
contentLen, contentLenLen, err := leb128.DecodeUint32(bytes.NewReader(buffer.Bytes()))
if err != nil {
if errors.Is(err, io.EOF) {
return contents, nil
}
return nil, err
}
buffer.Next(int(contentLenLen))
content := make([]byte, contentLen)
_, err = buffer.Read(content)
if err != nil {
if errors.Is(err, io.EOF) {
return contents, nil
}
return nil, err
}
contents = append(contents, content)
}
}
func getContentKeys(request *OfferRequest) [][]byte {
if request.Kind == TransientOfferRequestKind {
contentKeys := make([][]byte, 0)
contents := request.Request.(*TransientOfferRequest).Contents
for _, content := range contents {
contentKeys = append(contentKeys, content.ContentKey)
}
return contentKeys
} else {
return request.Request.(*PersistOfferRequest).ContentKeys
}
}

View file

@ -11,6 +11,7 @@ import (
"time"
"github.com/optimism-java/utp-go"
"github.com/prysmaticlabs/go-bitfield"
"github.com/ethereum/go-ethereum/crypto"
"github.com/ethereum/go-ethereum/internal/testlog"
@ -49,7 +50,9 @@ func setupLocalPortalNode(addr string, bootNodes []*enode.Node) (*PortalProtocol
if bootNodes != nil {
conf.BootstrapNodes = bootNodes
}
portalProtocol, err := NewPortalProtocol(conf, portalwire.HistoryNetwork, newkey(), &MockStorage{db: make(map[string][]byte)})
contentQueue := make(chan *ContentElement, 50)
portalProtocol, err := NewPortalProtocol(conf, portalwire.HistoryNetwork, newkey(), &MockStorage{db: make(map[string][]byte)}, contentQueue)
if err != nil {
return nil, err
}
@ -263,4 +266,34 @@ func TestPortalWireProtocol(t *testing.T) {
assert.NoError(t, err)
assert.Equal(t, largeTestContent, content)
assert.Equal(t, portalwire.ContentConnIdSelector, flag)
testEntry1 := &ContentEntry{
ContentKey: []byte("test_entry1"),
Content: []byte("test_entry1_content"),
}
testEntry2 := &ContentEntry{
ContentKey: []byte("test_entry2"),
Content: []byte("test_entry2_content"),
}
testTransientOfferRequest := &TransientOfferRequest{
Contents: []*ContentEntry{testEntry1, testEntry2},
}
offerRequest := &OfferRequest{
Kind: TransientOfferRequestKind,
Request: testTransientOfferRequest,
}
contentKeys, err := node1.offer(node3.localNode.Node(), offerRequest)
assert.Equal(t, uint64(2), bitfield.Bitlist(contentKeys).Count())
assert.NoError(t, err)
contentElement := <-node3.contentQueue
assert.Equal(t, node1.localNode.Node().ID(), contentElement.Node)
assert.Equal(t, testEntry1.ContentKey, contentElement.ContentKeys[0])
assert.Equal(t, testEntry1.Content, contentElement.Contents[0])
assert.Equal(t, testEntry2.ContentKey, contentElement.ContentKeys[1])
assert.Equal(t, testEntry2.Content, contentElement.Contents[1])
}

View file

@ -60,8 +60,8 @@ const (
seedMaxAge = 5 * 24 * time.Hour
)
// Table is the 'node table', a Kademlia-like index of neighbor nodes. The table keeps
// itself up-to-date by verifying the liveness of neighbors and requesting their node
// Table is the 'Node table', a Kademlia-like index of neighbor nodes. The table keeps
// itself up-to-date by verifying the liveness of neighbors and requesting their Node
// records when announcements of a new record version are received.
type Table struct {
mutex sync.Mutex // protects buckets, bucket content, nursery, rand
@ -206,10 +206,10 @@ func (tab *Table) setFallbackNodes(nodes []*enode.Node) error {
nursery := make([]*node, 0, len(nodes))
for _, n := range nodes {
if err := n.ValidateComplete(); err != nil {
return fmt.Errorf("bad bootstrap node %q: %v", n, err)
return fmt.Errorf("bad bootstrap Node %q: %v", n, err)
}
if tab.cfg.NetRestrict != nil && !tab.cfg.NetRestrict.Contains(n.IP()) {
tab.log.Error("Bootstrap node filtered by netrestrict", "id", n.ID(), "ip", n.IP())
tab.log.Error("Bootstrap Node filtered by netrestrict", "id", n.ID(), "ip", n.IP())
continue
}
nursery = append(nursery, wrapNode(n))
@ -331,7 +331,7 @@ func (tab *Table) loadSeedNodes() {
for i := range seeds {
seed := seeds[i]
age := log.Lazy{Fn: func() interface{} { return time.Since(tab.db.LastPongReceived(seed.ID(), seed.IP())) }}
tab.log.Trace("Found seed node in database", "id", seed.ID(), "addr", seed.addr(), "age", age)
tab.log.Trace("Found seed Node in database", "id", seed.ID(), "addr", seed.addr(), "age", age)
tab.addSeenNode(seed)
}
}
@ -347,10 +347,10 @@ func (tab *Table) doRevalidate(done chan<- struct{}) {
return
}
// Ping the selected node and wait for a pong.
// Ping the selected Node and wait for a pong.
remoteSeq, err := tab.net.ping(unwrapNode(last))
// Also fetch record if the node replied and returned a higher sequence number.
// Also fetch record if the Node replied and returned a higher sequence number.
if last.Seq() < remoteSeq {
n, err := tab.net.RequestENR(unwrapNode(last))
if err != nil {
@ -364,18 +364,18 @@ func (tab *Table) doRevalidate(done chan<- struct{}) {
defer tab.mutex.Unlock()
b := tab.buckets[bi]
if err == nil {
// The node responded, move it to the front.
// The Node responded, move it to the front.
last.livenessChecks++
tab.log.Debug("Revalidated node", "b", bi, "id", last.ID(), "checks", last.livenessChecks)
tab.log.Debug("Revalidated Node", "b", bi, "id", last.ID(), "checks", last.livenessChecks)
tab.bumpInBucket(b, last)
return
}
// No reply received, pick a replacement or delete the node if there aren't
// No reply received, pick a replacement or delete the Node if there aren't
// any replacements.
if r := tab.replace(b, last); r != nil {
tab.log.Debug("Replaced dead node", "b", bi, "id", last.ID(), "ip", last.IP(), "checks", last.livenessChecks, "r", r.ID(), "rip", r.IP())
tab.log.Debug("Replaced dead Node", "b", bi, "id", last.ID(), "ip", last.IP(), "checks", last.livenessChecks, "r", r.ID(), "rip", r.IP())
} else {
tab.log.Debug("Removed dead node", "b", bi, "id", last.ID(), "ip", last.IP(), "checks", last.livenessChecks)
tab.log.Debug("Removed dead Node", "b", bi, "id", last.ID(), "ip", last.IP(), "checks", last.livenessChecks)
}
}
@ -664,8 +664,8 @@ func (tab *Table) bumpInBucket(b *bucket, n *node) bool {
}
func (tab *Table) deleteInBucket(b *bucket, n *node) {
// Check if the node is actually in the bucket so the removed hook
// isn't called multiple times for the same node.
// Check if the Node is actually in the bucket so the removed hook
// isn't called multiple times for the same Node.
if !contains(b.entries, n.ID()) {
return
}

View file

@ -61,7 +61,7 @@ func testPingReplace(t *testing.T, newNodeIsResponding, lastInBucketIsResponding
pingSender := wrapNode(enode.NewV4(&pingKey.PublicKey, net.IP{127, 0, 0, 1}, 99, 99))
last := fillBucket(tab, pingSender)
// Add the sender as if it just pinged us. Revalidate should replace the last node in
// Add the sender as if it just pinged us. Revalidate should replace the last Node in
// its bucket if it is unresponsive. Revalidate again to ensure that
transport.dead[last.ID()] = !lastInBucketIsResponding
transport.dead[pingSender.ID()] = !newNodeIsResponding
@ -70,8 +70,8 @@ func testPingReplace(t *testing.T, newNodeIsResponding, lastInBucketIsResponding
tab.doRevalidate(make(chan struct{}, 1))
if !transport.pinged[last.ID()] {
// Oldest node in bucket is pinged to see whether it is still alive.
t.Error("table did not ping last node in bucket")
// Oldest Node in bucket is pinged to see whether it is still alive.
t.Error("table did not ping last Node in bucket")
}
tab.mutex.Lock()
@ -194,7 +194,7 @@ func TestTable_findnodeByID(t *testing.T) {
t.Parallel()
test := func(test *closeTest) bool {
// for any node table, Target and N
// for any Node table, Target and N
transport := newPingRecorder()
tab, db := newTestTable(transport)
defer db.Close()
@ -232,7 +232,7 @@ func TestTable_findnodeByID(t *testing.T) {
}
farthestResult := result[len(result)-1].ID()
if enode.DistCmp(test.Target, n.ID(), farthestResult) < 0 {
t.Errorf("table contains node that is closer to target but it's not in result")
t.Errorf("table contains Node that is closer to target but it's not in result")
t.Logf(" Target: %v", test.Target)
t.Logf(" Farthest Result: %v", farthestResult)
t.Logf(" ID: %v", n.ID())
@ -333,7 +333,7 @@ func TestTable_addSeenNode(t *testing.T) {
checkIPLimitInvariant(t, tab)
}
// This test checks that ENR updates happen during revalidation. If a node in the table
// This test checks that ENR updates happen during revalidation. If a Node in the table
// announces a new sequence number, the new record should be pulled.
func TestTable_revalidateSyncRecord(t *testing.T) {
transport := newPingRecorder()
@ -342,14 +342,14 @@ func TestTable_revalidateSyncRecord(t *testing.T) {
defer db.Close()
defer tab.close()
// Insert a node.
// Insert a Node.
var r enr.Record
r.Set(enr.IP(net.IP{127, 0, 0, 1}))
id := enode.ID{1}
n1 := wrapNode(enode.SignNull(&r, id))
tab.addSeenNode(n1)
// Update the node record.
// Update the Node record.
r.Set(enr.WithEntry("foo", "bar"))
n2 := enode.SignNull(&r, id)
transport.updateRecord(n2)

View file

@ -56,7 +56,7 @@ func nodeAtDistance(base enode.ID, ld int, ip net.IP) *node {
return wrapNode(enode.SignNull(&r, idAtDistance(base, ld)))
}
// nodesAtDistance creates n nodes for which enode.LogDist(base, node.ID()) == ld.
// nodesAtDistance creates n nodes for which enode.LogDist(base, Node.ID()) == ld.
func nodesAtDistance(base enode.ID, ld int, n int) []*enode.Node {
results := make([]*enode.Node, n)
for i := range results {

View file

@ -39,7 +39,7 @@ func TestUDPv4_Lookup(t *testing.T) {
t.Fatalf("lookup on empty table returned %d results: %#v", len(results), results)
}
// Seed table with initial node.
// Seed table with initial Node.
fillTable(test.table, []*node{wrapNode(lookupTestnet.node(256, 0))})
// Start the lookup.
@ -114,7 +114,7 @@ func TestUDPv4_LookupIteratorClose(t *testing.T) {
it := test.udp.RandomNodes()
if ok := it.Next(); !ok || it.Node() == nil {
t.Fatalf("iterator didn't return any node")
t.Fatalf("iterator didn't return any Node")
}
it.Close()
@ -122,7 +122,7 @@ func TestUDPv4_LookupIteratorClose(t *testing.T) {
ncalls := 0
for ; ncalls < 100 && it.Next(); ncalls++ {
if it.Node() == nil {
t.Error("iterator returned Node() == nil node after Next() == true")
t.Error("iterator returned Node() == nil Node after Next() == true")
}
}
t.Logf("iterator returned %d nodes after close", ncalls)
@ -130,7 +130,7 @@ func TestUDPv4_LookupIteratorClose(t *testing.T) {
t.Errorf("Next() == true after close and %d more calls", ncalls)
}
if n := it.Node(); n != nil {
t.Errorf("iterator returned non-nil node after close and %d more calls", ncalls)
t.Errorf("iterator returned non-nil Node after close and %d more calls", ncalls)
}
}

View file

@ -40,7 +40,7 @@ import (
var (
errExpired = errors.New("expired")
errUnsolicitedReply = errors.New("unsolicited reply")
errUnknownNode = errors.New("unknown node")
errUnknownNode = errors.New("unknown Node")
errTimeout = errors.New("RPC timeout")
errClockWarp = errors.New("reply deadline too far in the future")
errClosed = errors.New("socket closed")
@ -155,7 +155,7 @@ func ListenV4(c UDPConn, ln *enode.LocalNode, cfg Config) (*UDPv4, error) {
return t, nil
}
// Self returns the local node.
// Self returns the local Node.
func (t *UDPv4) Self() *enode.Node {
return t.localNode.Node()
}
@ -170,10 +170,10 @@ func (t *UDPv4) Close() {
})
}
// Resolve searches for a specific node with the given ID and tries to get the most recent
// version of the node record for it. It returns n if the node could not be resolved.
// Resolve searches for a specific Node with the given ID and tries to get the most recent
// version of the Node record for it. It returns n if the Node could not be resolved.
func (t *UDPv4) Resolve(n *enode.Node) *enode.Node {
// Try asking directly. This works if the node is still responding on the endpoint we have.
// Try asking directly. This works if the Node is still responding on the endpoint we have.
if rn, err := t.RequestENR(n); err == nil {
return rn
}
@ -206,7 +206,7 @@ func (t *UDPv4) ourEndpoint() v4wire.Endpoint {
return v4wire.NewEndpoint(a, uint16(n.TCP()))
}
// Ping sends a ping message to the given node.
// Ping sends a ping message to the given Node.
func (t *UDPv4) Ping(n *enode.Node) error {
_, err := t.ping(n)
return err
@ -311,7 +311,7 @@ func (t *UDPv4) findnode(toid enode.ID, toaddr *net.UDPAddr, target v4wire.Pubke
nreceived++
n, err := t.nodeFromRPC(toaddr, rn)
if err != nil {
t.log.Trace("Invalid neighbor node received", "ip", rn.IP, "addr", toaddr, "err", err)
t.log.Trace("Invalid neighbor Node received", "ip", rn.IP, "addr", toaddr, "err", err)
continue
}
nodes = append(nodes, n)
@ -322,9 +322,9 @@ func (t *UDPv4) findnode(toid enode.ID, toaddr *net.UDPAddr, target v4wire.Pubke
Target: target,
Expiration: uint64(time.Now().Add(expiration).Unix()),
})
// Ensure that callers don't see a timeout if the node actually responded. Since
// Ensure that callers don't see a timeout if the Node actually responded. Since
// findnode can receive more than one neighbors response, the reply matcher will be
// active until the remote node sends enough nodes. If the remote end doesn't have
// active until the remote Node sends enough nodes. If the remote end doesn't have
// enough nodes the reply matcher will time out waiting for the second reply, but
// there's no need for an error in that case.
err := <-rm.errc
@ -334,7 +334,7 @@ func (t *UDPv4) findnode(toid enode.ID, toaddr *net.UDPAddr, target v4wire.Pubke
return nodes, err
}
// RequestENR sends ENRRequest to the given node and waits for a response.
// RequestENR sends ENRRequest to the given Node and waits for a response.
func (t *UDPv4) RequestENR(n *enode.Node) (*enode.Node, error) {
addr := &net.UDPAddr{IP: n.IP(), Port: n.UDP()}
t.ensureBond(n.ID(), addr)
@ -675,7 +675,7 @@ func (t *UDPv4) handlePing(h *packetHandlerV4, from *net.UDPAddr, fromID enode.I
t.tab.addVerifiedNode(n)
}
// Update node database and endpoint predictor.
// Update Node database and endpoint predictor.
t.db.UpdateLastPingReceived(n.ID(), from.IP, time.Now())
t.localNode.UDPEndpointStatement(from, &net.UDPAddr{IP: req.To.IP, Port: int(req.To.UDP)})
}

View file

@ -271,7 +271,7 @@ func TestUDPv4_findnode(t *testing.T) {
}
fillTable(test.table, nodes.entries)
// ensure there's a bond with the test node,
// ensure there's a bond with the test Node,
// findnode won't be accepted otherwise.
remoteID := v4wire.EncodePubkey(&test.remotekey.PublicKey).ID()
test.table.db.UpdateLastPongReceived(remoteID, test.remoteaddr.IP, time.Now())
@ -290,7 +290,7 @@ func TestUDPv4_findnode(t *testing.T) {
t.Errorf("result mismatch at %d:\n got: %v\n want: %v", i, n, expected.entries[i])
}
if !live[n.ID.ID()] {
t.Errorf("result includes dead node %v", n.ID.ID())
t.Errorf("result includes dead Node %v", n.ID.ID())
}
}
})
@ -434,25 +434,25 @@ func TestUDPv4_successfulPing(t *testing.T) {
test.packetIn(nil, &v4wire.Pong{ReplyTok: hash, Expiration: futureExp})
})
// The node should be added to the table shortly after getting the
// The Node should be added to the table shortly after getting the
// pong packet.
select {
case n := <-added:
rid := encodePubkey(&test.remotekey.PublicKey).id()
if n.ID() != rid {
t.Errorf("node has wrong ID: got %v, want %v", n.ID(), rid)
t.Errorf("Node has wrong ID: got %v, want %v", n.ID(), rid)
}
if !n.IP().Equal(test.remoteaddr.IP) {
t.Errorf("node has wrong IP: got %v, want: %v", n.IP(), test.remoteaddr.IP)
t.Errorf("Node has wrong IP: got %v, want: %v", n.IP(), test.remoteaddr.IP)
}
if n.UDP() != test.remoteaddr.Port {
t.Errorf("node has wrong UDP port: got %v, want: %v", n.UDP(), test.remoteaddr.Port)
t.Errorf("Node has wrong UDP port: got %v, want: %v", n.UDP(), test.remoteaddr.Port)
}
if n.TCP() != int(testRemote.TCP) {
t.Errorf("node has wrong TCP port: got %v, want: %v", n.TCP(), testRemote.TCP)
t.Errorf("Node has wrong TCP port: got %v, want: %v", n.TCP(), testRemote.TCP)
}
case <-time.After(2 * time.Second):
t.Errorf("node was not added within 2 seconds")
t.Errorf("Node was not added within 2 seconds")
}
}
@ -489,7 +489,7 @@ func TestUDPv4_EIP868(t *testing.T) {
t.Fatalf("invalid record: %v", err)
}
if !reflect.DeepEqual(n, wantNode) {
t.Fatalf("wrong node in ENRResponse: %v", n)
t.Fatalf("wrong Node in ENRResponse: %v", n)
}
})
}
@ -525,7 +525,7 @@ func TestUDPv4_smallNetConvergence(t *testing.T) {
return
}
}
status <- fmt.Errorf("node %s didn't find all nodes", node.Self().ID().TerminalString())
status <- fmt.Errorf("Node %s didn't find all nodes", node.Self().ID().TerminalString())
}()
}
@ -555,7 +555,7 @@ func startLocalhostV4(t *testing.T, cfg Config) *UDPv4 {
db, _ := enode.OpenDB("")
ln := enode.NewLocalNode(db, cfg.PrivateKey)
// Prefix logs with node ID.
// Prefix logs with Node ID.
lprefix := fmt.Sprintf("(%s)", ln.ID().TerminalString())
lfmt := log.TerminalFormat(false)
cfg.Log = testlog.Logger(t, log.LvlTrace)

View file

@ -182,7 +182,7 @@ func newUDPv5(conn UDPConn, ln *enode.LocalNode, cfg Config) (*UDPv5, error) {
return t, nil
}
// Self returns the local node record.
// Self returns the local Node record.
func (t *UDPv5) Self() *enode.Node {
return t.localNode.Node()
}
@ -198,19 +198,19 @@ func (t *UDPv5) Close() {
})
}
// Ping sends a ping message to the given node.
// Ping sends a ping message to the given Node.
func (t *UDPv5) Ping(n *enode.Node) error {
_, err := t.ping(n)
return err
}
// Resolve searches for a specific node with the given ID and tries to get the most recent
// version of the node record for it. It returns n if the node could not be resolved.
// Resolve searches for a specific Node with the given ID and tries to get the most recent
// version of the Node record for it. It returns n if the Node could not be resolved.
func (t *UDPv5) Resolve(n *enode.Node) *enode.Node {
if intable := t.tab.getNode(n.ID()); intable != nil && intable.Seq() > n.Seq() {
n = intable
}
// Try asking directly. This works if the node is still responding on the endpoint we have.
// Try asking directly. This works if the Node is still responding on the endpoint we have.
if resp, err := t.RequestENR(n); err == nil {
return resp
}
@ -238,7 +238,7 @@ func (t *UDPv5) AllNodes() []*enode.Node {
return nodes
}
// LocalNode returns the current local node running the
// LocalNode returns the current local Node running the
// protocol.
func (t *UDPv5) LocalNode() *enode.LocalNode {
return t.localNode
@ -251,7 +251,7 @@ func (t *UDPv5) RegisterTalkHandler(protocol string, handler TalkRequestHandler)
t.talk.register(protocol, handler)
}
// TalkRequest sends a talk request to a node and waits for a response.
// TalkRequest sends a talk request to a Node and waits for a response.
func (t *UDPv5) TalkRequest(n *enode.Node, protocol string, request []byte) ([]byte, error) {
req := &v5wire.TalkRequest{Protocol: protocol, Message: request}
resp := t.callToNode(n, v5wire.TalkResponseMsg, req)
@ -264,7 +264,7 @@ func (t *UDPv5) TalkRequest(n *enode.Node, protocol string, request []byte) ([]b
}
}
// TalkRequestToID sends a talk request to a node and waits for a response.
// TalkRequestToID sends a talk request to a Node and waits for a response.
func (t *UDPv5) TalkRequestToID(id enode.ID, addr *net.UDPAddr, protocol string, request []byte) ([]byte, error) {
req := &v5wire.TalkRequest{Protocol: protocol, Message: request}
resp := t.callToID(id, addr, v5wire.TalkResponseMsg, req)
@ -828,7 +828,7 @@ func (t *UDPv5) matchWithCall(fromID enode.ID, nonce v5wire.Nonce) (*callV5, err
func (t *UDPv5) handlePing(p *v5wire.Ping, fromID enode.ID, fromAddr *net.UDPAddr) {
remoteIP := fromAddr.IP
// Handle IPv4 mapped IPv6 addresses in the
// event the local node is binded to an
// event the local Node is binded to an
// ipv6 interface.
if remoteIP.To4() != nil {
remoteIP = remoteIP.To4()

View file

@ -77,7 +77,7 @@ func startLocalhostV5(t *testing.T, cfg Config) *UDPv5 {
db, _ := enode.OpenDB("")
ln := enode.NewLocalNode(db, cfg.PrivateKey)
// Prefix logs with node ID.
// Prefix logs with Node ID.
lprefix := fmt.Sprintf("(%s)", ln.ID().TerminalString())
lfmt := log.TerminalFormat(false)
cfg.Log = testlog.Logger(t, log.LvlTrace)
@ -138,13 +138,13 @@ func TestUDPv5_unknownPacket(t *testing.T) {
}
}
// Unknown packet from unknown node.
// Unknown packet from unknown Node.
test.packetIn(&v5wire.Unknown{Nonce: nonce})
test.waitPacketOut(func(p *v5wire.Whoareyou, addr *net.UDPAddr, _ v5wire.Nonce) {
check(p, 0)
})
// Make node known.
// Make Node known.
n := test.getNode(test.remotekey, test.remoteaddr).Node()
test.table.addSeenNode(wrapNode(n))
@ -168,7 +168,7 @@ func TestUDPv5_findnodeHandling(t *testing.T) {
fillTable(test.table, wrapNodes(nodes249))
fillTable(test.table, wrapNodes(nodes248))
// Requesting with distance zero should return the node's own record.
// Requesting with distance zero should return the Node's own record.
test.packetIn(&v5wire.Findnode{ReqID: []byte{0}, Distances: []uint{0}})
test.expectNodes([]byte{0}, 1, []*enode.Node{test.udp.Self()})
@ -215,7 +215,7 @@ func (test *udpV5Test) expectNodes(wantReqID []byte, wantTotal uint8, wantNodes
n, _ := enode.New(enode.ValidSchemesForTesting, record)
want := nodeSet[n.ID()]
if want == nil {
test.t.Fatalf("unexpected node in response: %v", n)
test.t.Fatalf("unexpected Node in response: %v", n)
}
if !reflect.DeepEqual(record, want) {
test.t.Fatalf("wrong record in response: %v", n)
@ -592,7 +592,7 @@ func TestUDPv5_lookup(t *testing.T) {
}
}
// Seed table with initial node.
// Seed table with initial Node.
initialNode := lookupTestnet.node(256, 0)
fillTable(test.table, []*node{wrapNode(initialNode)})
@ -613,7 +613,7 @@ func TestUDPv5_lookup(t *testing.T) {
test.packetInFrom(key, to, &v5wire.Pong{ReqID: p.ReqID})
case *v5wire.Findnode:
if asked[recipient.ID()] {
t.Error("Asked node", recipient.ID(), "twice")
t.Error("Asked Node", recipient.ID(), "twice")
}
asked[recipient.ID()] = true
nodes := lookupTestnet.neighborsAtDistances(recipient, p.Distances, 16)
@ -630,7 +630,7 @@ func TestUDPv5_lookup(t *testing.T) {
checkLookupResults(t, lookupTestnet, results)
}
// This test checks the local node can be utilised to set key-values.
// This test checks the local Node can be utilised to set key-values.
func TestUDPv5_LocalNode(t *testing.T) {
t.Parallel()
var cfg Config
@ -638,7 +638,7 @@ func TestUDPv5_LocalNode(t *testing.T) {
defer node.Close()
localNd := node.LocalNode()
// set value in node's local record
// set value in Node's local record
testVal := [4]byte{'A', 'B', 'C', 'D'}
localNd.Set(enr.WithEntry("testing", &testVal))
@ -831,7 +831,7 @@ func (test *udpV5Test) waitPacketOut(validate interface{}) (closed bool) {
}
ln := test.nodesByIP[string(dgram.to.IP)]
if ln == nil {
test.t.Fatalf("attempt to send to non-existing node %v", &dgram.to)
test.t.Fatalf("attempt to send to non-existing Node %v", &dgram.to)
return false
}
codec := &testCodec{test: test, id: ln.ID()}