mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-31 00:53:46 +00:00
feat:add find content
Signed-off-by: grapebaba <281165273@qq.com>
This commit is contained in:
parent
adfd669e40
commit
5944d6589c
5 changed files with 261 additions and 93 deletions
13
go.mod
13
go.mod
|
|
@ -54,6 +54,7 @@ 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-20231030043430-a1331c25fa98
|
||||
github.com/peterh/liner v1.1.1-0.20190123174540-a2c9a5303de7
|
||||
github.com/protolambda/bls12-381-util v0.0.0-20220416220906-d8552aa452c7
|
||||
github.com/rs/cors v1.7.0
|
||||
|
|
@ -67,11 +68,11 @@ require (
|
|||
go.uber.org/automaxprocs v1.5.2
|
||||
golang.org/x/crypto v0.14.0
|
||||
golang.org/x/exp v0.0.0-20230905200255-921286631fa9
|
||||
golang.org/x/sync v0.3.0
|
||||
golang.org/x/sync v0.4.0
|
||||
golang.org/x/sys v0.13.0
|
||||
golang.org/x/text v0.13.0
|
||||
golang.org/x/time v0.3.0
|
||||
golang.org/x/tools v0.13.0
|
||||
golang.org/x/tools v0.14.0
|
||||
gopkg.in/natefinch/lumberjack.v2 v2.0.0
|
||||
gopkg.in/yaml.v3 v3.0.1
|
||||
)
|
||||
|
|
@ -128,7 +129,6 @@ 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-20231024092003-1dd76611b5f2 // 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
|
||||
|
|
@ -141,10 +141,9 @@ require (
|
|||
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.6.0 // indirect
|
||||
go.uber.org/zap v1.19.0 // indirect
|
||||
golang.org/x/mod v0.12.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
|
||||
gopkg.in/yaml.v2 v2.4.0 // indirect
|
||||
|
|
|
|||
31
go.sum
31
go.sum
|
|
@ -93,7 +93,6 @@ github.com/aws/aws-sdk-go-v2/service/sts v1.23.2/go.mod h1:Eows6e1uQEsc4ZaHANmsP
|
|||
github.com/aws/smithy-go v1.15.0 h1:PS/durmlzvAFpQHDs4wi4sNNP9ExsqZh6IlfdHXgKK8=
|
||||
github.com/aws/smithy-go v1.15.0/go.mod h1:Tg+OJXh4MB2R/uN61Ko2f6hTZwB/ZYGOtib8J3gBHzA=
|
||||
github.com/aymerick/raymond v2.0.3-0.20180322193309-b565731e1464+incompatible/go.mod h1:osfaiScAUVup+UC9Nfq76eWqDhXlp+4UYaA8uhTBO6g=
|
||||
github.com/benbjohnson/clock v1.1.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA=
|
||||
github.com/beorn7/perks v0.0.0-20180321164747-3a771d992973/go.mod h1:Dwedo/Wpr24TaqPxmxbtue+5NUziq4I4S80YR8gNf3Q=
|
||||
github.com/beorn7/perks v1.0.0/go.mod h1:KWe93zE9D1o94FZ5RNwFwVgaQK1VOXiVxmqh+CedLV8=
|
||||
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
|
||||
|
|
@ -496,8 +495,8 @@ 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-20231024092003-1dd76611b5f2 h1:1H9unjvDxqDVTF9rh99sfolMtp8zYCUNAq+aoeGCxLA=
|
||||
github.com/optimism-java/utp-go v0.0.0-20231024092003-1dd76611b5f2/go.mod h1:ohCuwoc66lfiNpo2Wk22zV07IbB0gte8+TYiclv9an4=
|
||||
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/pelletier/go-toml v1.2.0/go.mod h1:5z9KED0ma1S8pY6P1sdut58dfprrGBbd/94hg7ilaic=
|
||||
github.com/peterh/liner v1.1.1-0.20190123174540-a2c9a5303de7 h1:oYW+YCJ1pachXTQmzR3rNLYGGz4g/UgFcjb28p/viDM=
|
||||
github.com/peterh/liner v1.1.1-0.20190123174540-a2c9a5303de7/go.mod h1:CRroGNssyjTd/qIG2FyxByd2S8JEAZXBl4qUrZf8GS0=
|
||||
|
|
@ -617,15 +616,13 @@ go.opencensus.io v0.22.0/go.mod h1:+kGneAE2xo2IficOXnaByMWTGM9T73dGwxeWcUqIpI8=
|
|||
go.opencensus.io v0.22.2/go.mod h1:yxeiOL68Rb0Xd1ddK5vPZ/oVn4vY4Ynel7k9FzqtOIw=
|
||||
go.opencensus.io v0.22.3/go.mod h1:yxeiOL68Rb0Xd1ddK5vPZ/oVn4vY4Ynel7k9FzqtOIw=
|
||||
go.opencensus.io v0.22.4/go.mod h1:yxeiOL68Rb0Xd1ddK5vPZ/oVn4vY4Ynel7k9FzqtOIw=
|
||||
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/zap v1.19.0 h1:mZQZefskPPCMIBCSEH0v2/iUqqLrYtaeqwD6FUGUnFE=
|
||||
go.uber.org/zap v1.19.0/go.mod h1:xg/QME4nWcxGxrpdeYfq7UvYrLh66cuVKdrbD1XF/NI=
|
||||
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.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=
|
||||
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
|
||||
|
|
@ -672,8 +669,8 @@ 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=
|
||||
golang.org/x/net v0.0.0-20180906233101-161cd47e91fd/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
|
||||
|
|
@ -734,8 +731,8 @@ golang.org/x/sync v0.0.0-20201207232520-09787c993a3a/go.mod h1:RxMgew5VJxzue5/jJ
|
|||
golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.3.0 h1:ftCYgMx6zT/asHUrPw8BLLscYtGznsLAnjq5RH9P66E=
|
||||
golang.org/x/sync v0.3.0/go.mod h1:FU7BRWz2tNW+3quACPkgCx/L+uEAv1htQ0V83Z9Rj+Y=
|
||||
golang.org/x/sync v0.4.0 h1:zxkM55ReGkDlKSM+Fu41A+zmbZuaPVbGMzvvdUPznYQ=
|
||||
golang.org/x/sync v0.4.0/go.mod h1:FU7BRWz2tNW+3quACPkgCx/L+uEAv1htQ0V83Z9Rj+Y=
|
||||
golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
golang.org/x/sys v0.0.0-20180905080454-ebe1bf3edb33/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
golang.org/x/sys v0.0.0-20180909124046-d0be0721c37e/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
|
|
@ -841,7 +838,6 @@ golang.org/x/tools v0.0.0-20190628153133-6cdbf07be9d0/go.mod h1:/rFqwRUd4F7ZHNgw
|
|||
golang.org/x/tools v0.0.0-20190816200558-6889da9d5479/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
|
||||
golang.org/x/tools v0.0.0-20190911174233-4f2ddba30aff/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
|
||||
golang.org/x/tools v0.0.0-20191012152004-8de300cfc20a/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
|
||||
golang.org/x/tools v0.0.0-20191108193012-7d206e10da11/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
|
||||
golang.org/x/tools v0.0.0-20191113191852-77e3bb0ad9e7/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
|
||||
golang.org/x/tools v0.0.0-20191115202509-3a792d9c32b2/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
|
||||
golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
|
||||
|
|
@ -870,8 +866,8 @@ 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=
|
||||
golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||
|
|
@ -980,7 +976,6 @@ gopkg.in/yaml.v2 v2.3.0/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
|||
gopkg.in/yaml.v2 v2.4.0 h1:D8xgwECY7CYvx+Y2n4sBz93Jn9JRvxdiyyo8CTfuKaY=
|
||||
gopkg.in/yaml.v2 v2.4.0/go.mod h1:RDklbk79AGWmwhnvt/jBztapEOGDOx6ZbXqjP6csGnQ=
|
||||
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
|
||||
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
gotest.tools/v3 v3.5.1 h1:EENdUnS3pdur5nybKYIh2Vfgc8IUNBjxDPSjtiJcOzU=
|
||||
|
|
|
|||
|
|
@ -7,11 +7,13 @@ import (
|
|||
"encoding/binary"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"sort"
|
||||
"time"
|
||||
|
||||
"github.com/VictoriaMetrics/fastcache"
|
||||
"github.com/ethereum/go-ethereum/crypto"
|
||||
"github.com/ethereum/go-ethereum/log"
|
||||
"github.com/ethereum/go-ethereum/p2p/discover/portalwire"
|
||||
"github.com/ethereum/go-ethereum/p2p/enode"
|
||||
|
|
@ -43,9 +45,11 @@ const (
|
|||
|
||||
portalFindnodesResultLimit = 32
|
||||
|
||||
defaultUTPAcceptTimeout = 15 * time.Second
|
||||
defaultUTPConnectTimeout = 15 * time.Second
|
||||
|
||||
defaultUTPWriteTimeout = 60 * time.Second
|
||||
|
||||
defaultUTPReadTimeout = 60 * time.Second
|
||||
)
|
||||
|
||||
type PortalProtocolConfig struct {
|
||||
|
|
@ -94,7 +98,7 @@ type PortalProtocol struct {
|
|||
storage Storage
|
||||
}
|
||||
|
||||
func NewPortalProtocol(config *PortalProtocolConfig, protocolId string, privateKey *ecdsa.PrivateKey) (*PortalProtocol, error) {
|
||||
func NewPortalProtocol(config *PortalProtocolConfig, protocolId string, privateKey *ecdsa.PrivateKey, storage Storage) (*PortalProtocol, error) {
|
||||
nodeDB, err := enode.OpenDB(config.NodeDBPath)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -117,6 +121,7 @@ func NewPortalProtocol(config *PortalProtocolConfig, protocolId string, privateK
|
|||
cancelCloseCtx: cancelCloseCtx,
|
||||
localNode: localNode,
|
||||
validSchemes: enode.ValidSchemes,
|
||||
storage: storage,
|
||||
}
|
||||
|
||||
return protocol, nil
|
||||
|
|
@ -161,9 +166,8 @@ func (p *PortalProtocol) setupUDPListening() (*net.UDPConn, error) {
|
|||
p.utpPackets = make(chan *utp.UdpMessage, 10)
|
||||
p.utp, err = utp.ListenUTPOptions("utp", (*utp.Addr)(laddr), utp.WithCustomHandler(
|
||||
func(buf []byte, addr *net.UDPAddr) (int, error) {
|
||||
var a [32]byte
|
||||
// todo need to find enode.ID by addr
|
||||
_, err := p.DiscV5.TalkRequestToID(a, addr, portalwire.UTPNetwork, buf)
|
||||
id := crypto.Keccak256([]byte(addr.String()))
|
||||
_, err := p.DiscV5.TalkRequestToID(enode.ID(id), addr, portalwire.UTPNetwork, buf)
|
||||
return 0, err
|
||||
},
|
||||
func() ([]byte, *net.UDPAddr, error) {
|
||||
|
|
@ -211,6 +215,44 @@ func (p *PortalProtocol) setupDiscV5AndTable() error {
|
|||
return nil
|
||||
}
|
||||
|
||||
func (p *PortalProtocol) ping(node *enode.Node) (uint64, error) {
|
||||
enrSeq := p.DiscV5.LocalNode().Seq()
|
||||
radiusBytes, err := p.nodeRadius.MarshalSSZ()
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
customPayload := &portalwire.PingPongCustomData{
|
||||
Radius: radiusBytes,
|
||||
}
|
||||
|
||||
customPayloadBytes, err := customPayload.MarshalSSZ()
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
pingRequest := &portalwire.Ping{
|
||||
EnrSeq: enrSeq,
|
||||
CustomPayload: customPayloadBytes,
|
||||
}
|
||||
|
||||
p.log.Trace("Sending ping request", "protocol", p.protocolId, "source", p.Self().ID(), "target", node.ID(), "ping", pingRequest)
|
||||
pingRequestBytes, err := pingRequest.MarshalSSZ()
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
talkRequestBytes := make([]byte, 0, len(pingRequestBytes)+1)
|
||||
talkRequestBytes = append(talkRequestBytes, portalwire.PING)
|
||||
talkRequestBytes = append(talkRequestBytes, pingRequestBytes...)
|
||||
|
||||
talkResp, err := p.DiscV5.TalkRequest(node, p.protocolId, talkRequestBytes)
|
||||
|
||||
if err != nil {
|
||||
p.replaceNode(node)
|
||||
}
|
||||
return p.processPong(node, talkResp)
|
||||
}
|
||||
|
||||
func (p *PortalProtocol) findNodes(node *enode.Node, distances []uint) ([]*enode.Node, error) {
|
||||
distancesBytes := make([][2]byte, len(distances))
|
||||
for i, distance := range distances {
|
||||
|
|
@ -241,44 +283,148 @@ func (p *PortalProtocol) findNodes(node *enode.Node, distances []uint) ([]*enode
|
|||
return p.processNodes(node, talkResp, distances)
|
||||
}
|
||||
|
||||
func (p *PortalProtocol) processNodes(target *enode.Node, resp []byte, distances []uint) ([]*enode.Node, error) {
|
||||
var (
|
||||
nodes []*enode.Node
|
||||
seen = make(map[enode.ID]struct{})
|
||||
err error
|
||||
verified = 0
|
||||
)
|
||||
func (p *PortalProtocol) findContent(node *enode.Node, contentKey []byte) (byte, interface{}, error) {
|
||||
findContent := &portalwire.FindContent{
|
||||
ContentKey: contentKey,
|
||||
}
|
||||
|
||||
p.log.Trace("Sending find content request", "id", node.ID(), "findContent", findContent)
|
||||
findContentBytes, err := findContent.MarshalSSZ()
|
||||
if err != nil {
|
||||
p.log.Error("failed to marshal find content request", "err", err)
|
||||
return 0xff, nil, err
|
||||
}
|
||||
|
||||
talkRequestBytes := make([]byte, 0, len(findContentBytes)+1)
|
||||
talkRequestBytes = append(talkRequestBytes, portalwire.FINDCONTENT)
|
||||
talkRequestBytes = append(talkRequestBytes, findContentBytes...)
|
||||
|
||||
talkResp, err := p.DiscV5.TalkRequest(node, p.protocolId, talkRequestBytes)
|
||||
if err != nil {
|
||||
p.log.Error("failed to send find content request", "err", err)
|
||||
return 0xff, nil, err
|
||||
}
|
||||
|
||||
return p.processContent(node, talkResp)
|
||||
}
|
||||
|
||||
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")
|
||||
}
|
||||
|
||||
switch resp[1] {
|
||||
case portalwire.ContentRawSelector:
|
||||
content := &portalwire.Content{}
|
||||
err := content.UnmarshalSSZ(resp[2:])
|
||||
if err != nil {
|
||||
return 0xff, nil, err
|
||||
}
|
||||
|
||||
p.log.Trace("Received content response", "id", target.ID(), "content", content)
|
||||
return resp[1], content.Content, nil
|
||||
case portalwire.ContentConnIdSelector:
|
||||
connIdMsg := &portalwire.ConnectionId{}
|
||||
err := connIdMsg.UnmarshalSSZ(resp[2:])
|
||||
if err != nil {
|
||||
return 0xff, nil, err
|
||||
}
|
||||
|
||||
p.log.Trace("Received content response", "id", target.ID(), "connIdMsg", connIdMsg)
|
||||
rctx, rcancel := context.WithTimeout(context.Background(), defaultUTPConnectTimeout)
|
||||
laddr := p.utp.Addr().(*utp.Addr)
|
||||
raddr := &utp.Addr{IP: target.IP(), Port: target.UDP()}
|
||||
connId := binary.BigEndian.Uint16(connIdMsg.Id[:])
|
||||
conn, err := utp.DialUTPOptions("utp", laddr, raddr, utp.WithContext(rctx), utp.WithConnId(uint32(connId)))
|
||||
if err != nil {
|
||||
rcancel()
|
||||
return 0xff, nil, err
|
||||
}
|
||||
|
||||
err = conn.SetReadDeadline(time.Now().Add(defaultUTPReadTimeout))
|
||||
if err != nil {
|
||||
rcancel()
|
||||
return 0xff, nil, err
|
||||
}
|
||||
// Read ALL the data from the connection until EOF and return it
|
||||
data := make([]byte, 0)
|
||||
for {
|
||||
buf := make([]byte, 1024)
|
||||
var n int
|
||||
n, err = conn.Read(buf)
|
||||
if err != nil {
|
||||
rcancel()
|
||||
if errors.Is(err, io.EOF) {
|
||||
p.log.Trace("Received content response", "id", target.ID(), "data", data, "size", n)
|
||||
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]...)
|
||||
}
|
||||
case portalwire.ContentEnrsSelector:
|
||||
enrs := &portalwire.Enrs{}
|
||||
err := enrs.UnmarshalSSZ(resp[2:])
|
||||
|
||||
if err != nil {
|
||||
return 0xff, nil, err
|
||||
}
|
||||
|
||||
p.log.Trace("Received content response", "id", target.ID(), "enrs", enrs)
|
||||
|
||||
nodes := p.filterNodes(target, enrs.Enrs, nil)
|
||||
return resp[1], nodes, nil
|
||||
default:
|
||||
return 0xff, nil, fmt.Errorf("invalid content response")
|
||||
}
|
||||
}
|
||||
|
||||
func (p *PortalProtocol) processNodes(target *enode.Node, resp []byte, distances []uint) ([]*enode.Node, error) {
|
||||
if resp[0] != portalwire.NODES {
|
||||
return nil, fmt.Errorf("invalid nodes response")
|
||||
}
|
||||
|
||||
nodesResp := &portalwire.Nodes{}
|
||||
err = nodesResp.UnmarshalSSZ(resp[1:])
|
||||
err := nodesResp.UnmarshalSSZ(resp[1:])
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
p.table.addVerifiedNode(wrapNode(target))
|
||||
var n *enode.Node
|
||||
for _, b := range nodesResp.Enrs {
|
||||
nodes := p.filterNodes(target, nodesResp.Enrs, distances)
|
||||
|
||||
return nodes, nil
|
||||
}
|
||||
|
||||
func (p *PortalProtocol) filterNodes(target *enode.Node, enrs [][]byte, distances []uint) []*enode.Node {
|
||||
var (
|
||||
nodes []*enode.Node
|
||||
seen = make(map[enode.ID]struct{})
|
||||
err error
|
||||
verified = 0
|
||||
n *enode.Node
|
||||
)
|
||||
|
||||
for _, b := range enrs {
|
||||
record := &enr.Record{}
|
||||
err = rlp.DecodeBytes(b, record)
|
||||
if err != nil {
|
||||
p.log.Debug("Invalid record in nodes response", "id", target.ID(), "err", err)
|
||||
p.log.Error("Invalid record in nodes response", "id", target.ID(), "err", err)
|
||||
continue
|
||||
}
|
||||
n, err = p.verifyResponseNode(target, record, distances, seen)
|
||||
if err != nil {
|
||||
p.log.Debug("Invalid record in nodes response", "id", target.ID(), "err", err)
|
||||
p.log.Error("Invalid record in nodes response", "id", target.ID(), "err", err)
|
||||
continue
|
||||
}
|
||||
verified++
|
||||
nodes = append(nodes, n)
|
||||
}
|
||||
|
||||
p.log.Trace("Received nodes response", "id", target.ID(), "total", nodesResp.Total, "verified", verified, "nodes", nodes)
|
||||
return nodes, nil
|
||||
p.log.Trace("Received nodes response", "id", target.ID(), "total", len(enrs), "verified", verified, "nodes", nodes)
|
||||
return nodes
|
||||
}
|
||||
|
||||
func (p *PortalProtocol) processPong(target *enode.Node, resp []byte) (uint64, error) {
|
||||
|
|
@ -309,16 +455,16 @@ func (p *PortalProtocol) processPong(target *enode.Node, resp []byte) (uint64, e
|
|||
}
|
||||
|
||||
func (p *PortalProtocol) handleUtpTalkRequest(id enode.ID, addr *net.UDPAddr, msg []byte) []byte {
|
||||
if node := p.DiscV5.getNode(id); node != nil {
|
||||
p.table.addSeenNode(wrapNode(node))
|
||||
if n := p.DiscV5.getNode(id); n != nil {
|
||||
p.table.addSeenNode(wrapNode(n))
|
||||
}
|
||||
p.utpPackets <- &utp.UdpMessage{Buf: msg, Addr: addr}
|
||||
return []byte("")
|
||||
}
|
||||
|
||||
func (p *PortalProtocol) handleTalkRequest(id enode.ID, addr *net.UDPAddr, msg []byte) []byte {
|
||||
if node := p.DiscV5.getNode(id); node != nil {
|
||||
p.table.addSeenNode(wrapNode(node))
|
||||
if n := p.DiscV5.getNode(id); n != nil {
|
||||
p.table.addSeenNode(wrapNode(n))
|
||||
}
|
||||
|
||||
msgCode := msg[0]
|
||||
|
|
@ -459,16 +605,16 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque
|
|||
|
||||
contentId := p.storage.ContentId(request.ContentKey)
|
||||
if contentId == nil {
|
||||
return nil, fmt.Errorf("content not found")
|
||||
return nil, ContentNotFound
|
||||
}
|
||||
|
||||
var content []byte
|
||||
content, err = p.storage.Get(request.ContentKey, contentId)
|
||||
if err != nil {
|
||||
if err != nil && !errors.Is(err, ContentNotFound) {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if content == nil {
|
||||
if errors.Is(err, ContentNotFound) {
|
||||
closestNodes := p.findNodesCloseToContent(contentId)
|
||||
for i, n := range closestNodes {
|
||||
if n.ID() == id {
|
||||
|
|
@ -500,9 +646,21 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque
|
|||
|
||||
return talkRespBytes, nil
|
||||
} else if len(content) <= maxPayloadSize {
|
||||
contentMsgBytes := make([]byte, 0, len(content)+1)
|
||||
rawContentMsg := &portalwire.Content{
|
||||
Content: content,
|
||||
}
|
||||
|
||||
p.log.Trace("Sending raw content response", "protocol", p.protocolId, "source", addr, "content", rawContentMsg)
|
||||
|
||||
var rawContentMsgBytes []byte
|
||||
rawContentMsgBytes, err = rawContentMsg.MarshalSSZ()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
contentMsgBytes := make([]byte, 0, len(rawContentMsgBytes)+1)
|
||||
contentMsgBytes = append(contentMsgBytes, portalwire.ContentRawSelector)
|
||||
contentMsgBytes = append(contentMsgBytes, content...)
|
||||
contentMsgBytes = append(contentMsgBytes, rawContentMsgBytes...)
|
||||
|
||||
talkRespBytes := make([]byte, 0, len(contentMsgBytes)+1)
|
||||
talkRespBytes = append(talkRespBytes, portalwire.CONTENT)
|
||||
|
|
@ -515,7 +673,7 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque
|
|||
connIdSend := connId.SendId()
|
||||
|
||||
go func() {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), defaultUTPAcceptTimeout)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), defaultUTPConnectTimeout)
|
||||
var conn *utp.Conn
|
||||
conn, err = p.utp.AcceptUTPContext(ctx, connIdSend)
|
||||
if err != nil {
|
||||
|
|
@ -604,44 +762,6 @@ func (p *PortalProtocol) verifyResponseNode(sender *enode.Node, r *enr.Record, d
|
|||
return n, nil
|
||||
}
|
||||
|
||||
func (p *PortalProtocol) ping(node *enode.Node) (uint64, error) {
|
||||
enrSeq := p.DiscV5.LocalNode().Seq()
|
||||
radiusBytes, err := p.nodeRadius.MarshalSSZ()
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
customPayload := &portalwire.PingPongCustomData{
|
||||
Radius: radiusBytes,
|
||||
}
|
||||
|
||||
customPayloadBytes, err := customPayload.MarshalSSZ()
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
pingRequest := &portalwire.Ping{
|
||||
EnrSeq: enrSeq,
|
||||
CustomPayload: customPayloadBytes,
|
||||
}
|
||||
|
||||
p.log.Trace("Sending ping request", "protocol", p.protocolId, "source", p.Self().ID(), "target", node.ID(), "ping", pingRequest)
|
||||
pingRequestBytes, err := pingRequest.MarshalSSZ()
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
talkRequestBytes := make([]byte, 0, len(pingRequestBytes)+1)
|
||||
talkRequestBytes = append(talkRequestBytes, portalwire.PING)
|
||||
talkRequestBytes = append(talkRequestBytes, pingRequestBytes...)
|
||||
|
||||
talkResp, err := p.DiscV5.TalkRequest(node, p.protocolId, talkRequestBytes)
|
||||
|
||||
if err != nil {
|
||||
p.replaceNode(node)
|
||||
}
|
||||
return p.processPong(node, talkResp)
|
||||
}
|
||||
|
||||
func (p *PortalProtocol) replaceNode(node *enode.Node) {
|
||||
p.table.mutex.Lock()
|
||||
defer p.table.mutex.Unlock()
|
||||
|
|
|
|||
|
|
@ -1,10 +1,12 @@
|
|||
package discover
|
||||
|
||||
import (
|
||||
"crypto/rand"
|
||||
"fmt"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/ethereum/go-ethereum/crypto"
|
||||
"github.com/ethereum/go-ethereum/internal/testlog"
|
||||
"github.com/ethereum/go-ethereum/log"
|
||||
"github.com/ethereum/go-ethereum/p2p/discover/portalwire"
|
||||
|
|
@ -13,6 +15,26 @@ import (
|
|||
"golang.org/x/exp/slices"
|
||||
)
|
||||
|
||||
type MockStorage struct {
|
||||
db map[string][]byte
|
||||
}
|
||||
|
||||
func (m *MockStorage) ContentId(contentKey []byte) []byte {
|
||||
return crypto.Keccak256(contentKey)
|
||||
}
|
||||
|
||||
func (m *MockStorage) Get(contentKey []byte, contentId []byte) ([]byte, error) {
|
||||
if content, ok := m.db[string(contentId)]; ok {
|
||||
return content, nil
|
||||
}
|
||||
return nil, ContentNotFound
|
||||
}
|
||||
|
||||
func (m *MockStorage) Put(contentKey []byte, content []byte) error {
|
||||
m.db[string(m.ContentId(contentKey))] = content
|
||||
return nil
|
||||
}
|
||||
|
||||
func setupLocalPortalNode(addr string, bootNodes []*enode.Node) (*PortalProtocol, error) {
|
||||
conf := DefaultPortalProtocolConfig()
|
||||
if addr != "" {
|
||||
|
|
@ -21,7 +43,7 @@ func setupLocalPortalNode(addr string, bootNodes []*enode.Node) (*PortalProtocol
|
|||
if bootNodes != nil {
|
||||
conf.BootstrapNodes = bootNodes
|
||||
}
|
||||
portalProtocol, err := NewPortalProtocol(conf, portalwire.HistoryNetwork, newkey())
|
||||
portalProtocol, err := NewPortalProtocol(conf, portalwire.HistoryNetwork, newkey(), &MockStorage{db: make(map[string][]byte)})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -76,4 +98,32 @@ func TestPortalWireProtocol(t *testing.T) {
|
|||
slices.ContainsFunc(node3.table.Nodes(), func(n *enode.Node) bool {
|
||||
return n.ID() == node2.localNode.Node().ID()
|
||||
})
|
||||
|
||||
err = node1.storage.Put([]byte("test_key"), []byte("test_value"))
|
||||
assert.NoError(t, err)
|
||||
|
||||
flag, content, err := node2.findContent(node1.localNode.Node(), []byte("test_key"))
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, portalwire.ContentRawSelector, flag)
|
||||
assert.Equal(t, []byte("test_value"), content)
|
||||
|
||||
flag, content, err = node2.findContent(node3.localNode.Node(), []byte("test_key"))
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, portalwire.ContentEnrsSelector, flag)
|
||||
assert.Equal(t, 1, len(content.([]*enode.Node)))
|
||||
assert.Equal(t, node1.localNode.Node().ID(), content.([]*enode.Node)[0].ID())
|
||||
|
||||
// create a byte slice of length 1199 and fill it with random data
|
||||
// this will be used as a test content
|
||||
largeTestContent := make([]byte, 1199)
|
||||
_, err = rand.Read(largeTestContent)
|
||||
assert.NoError(t, err)
|
||||
|
||||
err = node1.storage.Put([]byte("large_test_key"), largeTestContent)
|
||||
assert.NoError(t, err)
|
||||
|
||||
//flag, content, err = node2.findContent(node1.localNode.Node(), []byte("large_test_key"))
|
||||
//assert.NoError(t, err)
|
||||
//assert.Equal(t, portalwire.ContentConnIdSelector, flag)
|
||||
//assert.Equal(t, largeTestContent, content)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,5 +1,9 @@
|
|||
package discover
|
||||
|
||||
import "fmt"
|
||||
|
||||
var ContentNotFound = fmt.Errorf("content not found")
|
||||
|
||||
type Storage interface {
|
||||
ContentId(contentKey []byte) []byte
|
||||
|
||||
|
|
|
|||
Loading…
Reference in a new issue