Signed-off-by: Chen Kai <281165273grape@gmail.com>
This commit is contained in:
Chen Kai 2024-05-27 14:48:30 +08:00
parent 3ab51ce84f
commit 526ee8ad9e
3 changed files with 39 additions and 56 deletions

4
go.mod
View file

@ -55,7 +55,7 @@ require (
github.com/mattn/go-sqlite3 v1.14.18 github.com/mattn/go-sqlite3 v1.14.18
github.com/naoina/toml v0.1.2-0.20170918210437-9fafd6967416 github.com/naoina/toml v0.1.2-0.20170918210437-9fafd6967416
github.com/olekukonko/tablewriter v0.0.5 github.com/olekukonko/tablewriter v0.0.5
github.com/optimism-java/utp-go v0.0.0-20240518144144-6560912a0d99 github.com/optimism-java/utp-go v0.0.0-20240309041853-b6b3a0dea581
github.com/peterh/liner v1.1.1-0.20190123174540-a2c9a5303de7 github.com/peterh/liner v1.1.1-0.20190123174540-a2c9a5303de7
github.com/protolambda/bls12-381-util v0.1.0 github.com/protolambda/bls12-381-util v0.1.0
github.com/protolambda/zrnt v0.32.2 github.com/protolambda/zrnt v0.32.2
@ -75,7 +75,7 @@ require (
golang.org/x/crypto v0.22.0 golang.org/x/crypto v0.22.0
golang.org/x/exp v0.0.0-20231110203233-9a3e6036ecaa golang.org/x/exp v0.0.0-20231110203233-9a3e6036ecaa
golang.org/x/sync v0.7.0 golang.org/x/sync v0.7.0
golang.org/x/sys v0.20.0 golang.org/x/sys v0.19.0
golang.org/x/text v0.14.0 golang.org/x/text v0.14.0
golang.org/x/time v0.5.0 golang.org/x/time v0.5.0
golang.org/x/tools v0.20.0 golang.org/x/tools v0.20.0

4
go.sum
View file

@ -431,8 +431,6 @@ github.com/opentracing/opentracing-go v1.1.0 h1:pWlfV3Bxv7k65HYwkikxat0+s3pV4bsq
github.com/opentracing/opentracing-go v1.1.0/go.mod h1:UkNAQd3GIcIGf0SeVgPpRdFStlNbqXla1AfSYxPUl2o= github.com/opentracing/opentracing-go v1.1.0/go.mod h1:UkNAQd3GIcIGf0SeVgPpRdFStlNbqXla1AfSYxPUl2o=
github.com/optimism-java/utp-go v0.0.0-20240309041853-b6b3a0dea581 h1:ZxgrtI0xIw+clB32iDDDWaiTcCizTeN7rNyzH9YorPI= github.com/optimism-java/utp-go v0.0.0-20240309041853-b6b3a0dea581 h1:ZxgrtI0xIw+clB32iDDDWaiTcCizTeN7rNyzH9YorPI=
github.com/optimism-java/utp-go v0.0.0-20240309041853-b6b3a0dea581/go.mod h1:DZ0jYzLzt4ZsCmhI/iqYgGFoNx45OfpEoKzXB8HVALQ= github.com/optimism-java/utp-go v0.0.0-20240309041853-b6b3a0dea581/go.mod h1:DZ0jYzLzt4ZsCmhI/iqYgGFoNx45OfpEoKzXB8HVALQ=
github.com/optimism-java/utp-go v0.0.0-20240518144144-6560912a0d99 h1:8NEQQ8KNNUASMBB0OdnfYuxnOIEHJm1NcDiPnMb2Kvk=
github.com/optimism-java/utp-go v0.0.0-20240518144144-6560912a0d99/go.mod h1:DZ0jYzLzt4ZsCmhI/iqYgGFoNx45OfpEoKzXB8HVALQ=
github.com/optimism-java/zrnt v0.32.4-0.20240415084906-d9dbf06b32f7 h1:ZTQWXQ8xblCRUXhZs3h5qrBMSAHe8iNH7BG7a7IVFlI= github.com/optimism-java/zrnt v0.32.4-0.20240415084906-d9dbf06b32f7 h1:ZTQWXQ8xblCRUXhZs3h5qrBMSAHe8iNH7BG7a7IVFlI=
github.com/optimism-java/zrnt v0.32.4-0.20240415084906-d9dbf06b32f7/go.mod h1:A0fezkp9Tt3GBLATSPIbuY4ywYESyAuc/FFmPKg8Lqs= github.com/optimism-java/zrnt v0.32.4-0.20240415084906-d9dbf06b32f7/go.mod h1:A0fezkp9Tt3GBLATSPIbuY4ywYESyAuc/FFmPKg8Lqs=
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 h1:oYW+YCJ1pachXTQmzR3rNLYGGz4g/UgFcjb28p/viDM=
@ -711,8 +709,6 @@ 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.11.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.19.0 h1:q5f1RH2jigJ1MoAWp2KTp3gm5zAGFUTarQZ5U386+4o= golang.org/x/sys v0.19.0 h1:q5f1RH2jigJ1MoAWp2KTp3gm5zAGFUTarQZ5U386+4o=
golang.org/x/sys v0.19.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.19.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/sys v0.20.0 h1:Od9JTbYCk261bKm4M/mw7AklTlFYIa0bIp9BgSm1S8Y=
golang.org/x/sys v0.20.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/term v0.0.0-20201117132131-f5c789dd3221/go.mod h1:Nr5EML6q2oocZ2LXRh80K7BxOlk5/8JxuGnuhpl+muw= golang.org/x/term v0.0.0-20201117132131-f5c789dd3221/go.mod h1:Nr5EML6q2oocZ2LXRh80K7BxOlk5/8JxuGnuhpl+muw=
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8=

View file

@ -19,6 +19,7 @@ import (
"time" "time"
"github.com/ethereum/go-ethereum/common/hexutil" "github.com/ethereum/go-ethereum/common/hexutil"
"github.com/ethereum/go-ethereum/p2p/discover/v5wire"
"github.com/VictoriaMetrics/fastcache" "github.com/VictoriaMetrics/fastcache"
"github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/log"
@ -159,9 +160,7 @@ func DefaultPortalProtocolConfig() *PortalProtocolConfig {
} }
type PortalProtocol struct { type PortalProtocol struct {
table *Table table *Table
cachedIdsLock sync.Mutex
cachedIds map[string]enode.ID
protocolId string protocolId string
protocolName string protocolName string
@ -201,7 +200,6 @@ func NewPortalProtocol(config *PortalProtocolConfig, protocolId string, privateK
protocolName := portalwire.NetworkNameMap[protocolId] protocolName := portalwire.NetworkNameMap[protocolId]
protocol := &PortalProtocol{ protocol := &PortalProtocol{
cachedIds: make(map[string]enode.ID),
protocolId: protocolId, protocolId: protocolId,
protocolName: protocolName, protocolName: protocolName,
ListenAddr: config.ListenAddr, ListenAddr: config.ListenAddr,
@ -296,22 +294,39 @@ func (p *PortalProtocol) setupUDPListening() error {
var err error var err error
p.packetRouter = utp.NewPacketRouter( p.packetRouter = utp.NewPacketRouter(
func(buf []byte, addr *net.UDPAddr) (int, error) { func(buf []byte, addr *net.UDPAddr) (int, error) {
p.Log.Info("will send to target data", "ip", addr.IP.To4().String(), "port", addr.Port, "bufLength", len(buf)) nodes := p.table.Nodes()
var target *enode.Node
for _, n := range nodes {
if addr.Port != n.UDP() {
continue
}
if addr.IP != nil && addr.IP.To4().String() == n.IP().To4().String() {
target = n
p.cachedIdsLock.Lock() break
defer p.cachedIdsLock.Unlock() }
if id, ok := p.cachedIds[addr.String()]; ok { if addr.IP == nil {
_, err := p.DiscV5.TalkRequestToID(id, addr, string(portalwire.UTPNetwork), buf) nodeIp := n.IP().To4().String()
return len(buf), err if nodeIp == "127.0.0.1" || nodeIp == "0.0.0.0" {
} else { target = n
p.Log.Warn("not found target node info", "ip", addr.IP.To4().String(), "port", addr.Port, "bufLength", len(buf)) break
return 0, fmt.Errorf("not found target node id") }
}
} }
if target == nil {
p.Log.Warn("not fount target node info", "ip", addr.IP.To4().String(), "port", addr.Port, "bufLength", len(buf))
return 0, fmt.Errorf("not found target node info")
}
p.Log.Trace("send to target data", "ip", addr.IP.To4().String(), "port", addr.Port, "bufLength", len(buf))
req := &v5wire.TalkRequest{Protocol: string(portalwire.UTPNetwork), Message: buf}
p.DiscV5.sendFromAnotherThread(target.ID(), addr, req)
return len(buf), err
}) })
ctx := context.Background() ctx := context.Background()
var logger *zap.Logger var logger *zap.Logger
if p.Log.Enabled(ctx, log.LevelDebug) || p.Log.Enabled(ctx, log.LevelTrace) { if p.Log.Enabled(ctx, log.LevelDebug) || p.Log.Enabled(ctx, log.LevelTrace) {
logger, err = zap.NewDevelopmentConfig().Build() logger, err = zap.NewDevelopmentConfig().Build()
} else { } else {
@ -355,23 +370,6 @@ func (p *PortalProtocol) setupDiscV5AndTable() error {
return nil return nil
} }
func (p *PortalProtocol) putCacheNodeId(node *enode.Node) {
p.cachedIdsLock.Lock()
defer p.cachedIdsLock.Unlock()
addr := &net.UDPAddr{IP: node.IP(), Port: node.UDP()}
if _, ok := p.cachedIds[addr.String()]; !ok {
p.cachedIds[addr.String()] = node.ID()
}
}
func (p *PortalProtocol) putCacheId(id enode.ID, addr *net.UDPAddr) {
p.cachedIdsLock.Lock()
defer p.cachedIdsLock.Unlock()
if _, ok := p.cachedIds[addr.String()]; !ok {
p.cachedIds[addr.String()] = id
}
}
func (p *PortalProtocol) ping(node *enode.Node) (uint64, error) { func (p *PortalProtocol) ping(node *enode.Node) (uint64, error) {
pong, err := p.pingInner(node) pong, err := p.pingInner(node)
if err != nil { if err != nil {
@ -515,9 +513,6 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request *
return nil, fmt.Errorf("invalid accept response") return nil, fmt.Errorf("invalid accept response")
} }
p.Log.Info("will process Offer", "id", target.ID(), "ip", target.IP().To4().String(), "port", target.UDP())
p.putCacheNodeId(target)
accept := &portalwire.Accept{} accept := &portalwire.Accept{}
err = accept.UnmarshalSSZ(resp[1:]) err = accept.UnmarshalSSZ(resp[1:])
if err != nil { if err != nil {
@ -586,8 +581,8 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request *
connctx, conncancel := context.WithTimeout(ctx, defaultUTPConnectTimeout) connctx, conncancel := context.WithTimeout(ctx, defaultUTPConnectTimeout)
laddr := p.utp.Addr().(*utp.Addr) laddr := p.utp.Addr().(*utp.Addr)
raddr := &utp.Addr{IP: target.IP(), Port: target.UDP()} raddr := &utp.Addr{IP: target.IP(), Port: target.UDP()}
p.Log.Info("will connect to: ", "addr", raddr.String(), "connId", connId)
conn, err = utp.DialUTPOptions("utp", laddr, raddr, utp.WithContext(connctx), utp.WithSocketManager(p.utpSm), utp.WithConnId(uint32(connId))) conn, err = utp.DialUTPOptions("utp", laddr, raddr, utp.WithContext(connctx), utp.WithSocketManager(p.utpSm), utp.WithConnId(uint32(connId)))
p.Log.Info("will connect to: ", "addr", raddr.String(), "connId", connId)
if err != nil { if err != nil {
conncancel() conncancel()
p.Log.Error("failed to dial utp connection", "err", err) p.Log.Error("failed to dial utp connection", "err", err)
@ -641,9 +636,6 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte,
return 0xff, nil, fmt.Errorf("invalid content response") return 0xff, nil, fmt.Errorf("invalid content response")
} }
p.Log.Info("will process content", "id", target.ID(), "ip", target.IP().To4().String(), "port", target.UDP())
p.putCacheNodeId(target)
switch resp[1] { switch resp[1] {
case portalwire.ContentRawSelector: case portalwire.ContentRawSelector:
content := &portalwire.Content{} content := &portalwire.Content{}
@ -668,8 +660,8 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte,
laddr := p.utp.Addr().(*utp.Addr) laddr := p.utp.Addr().(*utp.Addr)
raddr := &utp.Addr{IP: target.IP(), Port: target.UDP()} raddr := &utp.Addr{IP: target.IP(), Port: target.UDP()}
connId := binary.BigEndian.Uint16(connIdMsg.Id[:]) connId := binary.BigEndian.Uint16(connIdMsg.Id[:])
p.Log.Info("will connect to: ", "addr", raddr.String(), "connId", connId)
conn, err := utp.DialUTPOptions("utp", laddr, raddr, utp.WithContext(connctx), utp.WithSocketManager(p.utpSm), utp.WithConnId(uint32(connId))) conn, err := utp.DialUTPOptions("utp", laddr, raddr, utp.WithContext(connctx), utp.WithSocketManager(p.utpSm), utp.WithConnId(uint32(connId)))
p.Log.Info("will connect to: ", "addr", raddr.String(), "connId", connId)
if err != nil { if err != nil {
conncancel() conncancel()
return 0xff, nil, err return 0xff, nil, err
@ -795,18 +787,16 @@ func (p *PortalProtocol) handleUtpTalkRequest(id enode.ID, addr *net.UDPAddr, ms
if n := p.DiscV5.getNode(id); n != nil { if n := p.DiscV5.getNode(id); n != nil {
p.table.addSeenNode(wrapNode(n)) p.table.addSeenNode(wrapNode(n))
} }
p.putCacheId(id, addr)
p.Log.Trace("receive utp data", "addr", addr, "msg-length", len(msg)) p.Log.Trace("receive utp data", "addr", addr, "msg-length", len(msg))
p.packetRouter.ReceiveMessage(msg, addr) p.packetRouter.ReceiveMessage(msg, addr)
return []byte("") return []byte("")
} }
func (p *PortalProtocol) handleTalkRequest(id enode.ID, addr *net.UDPAddr, msg []byte) []byte { func (p *PortalProtocol) handleTalkRequest(id enode.ID, addr *net.UDPAddr, msg []byte) []byte {
p.Log.Trace("handleTalkRequest", "id", id, "addr", addr)
if n := p.DiscV5.getNode(id); n != nil { if n := p.DiscV5.getNode(id); n != nil {
p.table.addSeenNode(wrapNode(n)) p.table.addSeenNode(wrapNode(n))
} }
p.putCacheId(id, addr)
msgCode := msg[0] msgCode := msg[0]
@ -971,8 +961,6 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque
return nil, err return nil, err
} }
p.putCacheId(id, addr)
if errors.Is(err, ContentNotFound) { if errors.Is(err, ContentNotFound) {
closestNodes := p.findNodesCloseToContent(contentId, portalFindnodesResultLimit) closestNodes := p.findNodesCloseToContent(contentId, portalFindnodesResultLimit)
for i, n := range closestNodes { for i, n := range closestNodes {
@ -1042,13 +1030,14 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque
default: default:
ctx, cancel := context.WithTimeout(bctx, defaultUTPConnectTimeout) ctx, cancel := context.WithTimeout(bctx, defaultUTPConnectTimeout)
var conn *utp.Conn var conn *utp.Conn
p.Log.Debug("will accept find content conn from: ", "source", addr, "connId", connId)
conn, err = p.utp.AcceptUTPContext(ctx, connIdSend) conn, err = p.utp.AcceptUTPContext(ctx, connIdSend)
p.Log.Info("will accept from: ", "source", addr, "connId", connId)
if err != nil { if err != nil {
p.Log.Error("failed to accept utp connection for handle find content", "connId", connIdSend, "err", err) p.Log.Error("failed to accept utp connection", "connId", connIdSend, "err", err)
cancel() cancel()
return return
} }
p.Log.Info("")
cancel() cancel()
err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout)) err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout))
@ -1149,8 +1138,6 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po
} }
} }
p.putCacheId(id, addr)
idBuffer := make([]byte, 2) idBuffer := make([]byte, 2)
if contentKeyBitlist.Count() != 0 { if contentKeyBitlist.Count() != 0 {
connId := p.connIdGen.GenCid(id, false) connId := p.connIdGen.GenCid(id, false)
@ -1164,10 +1151,10 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po
default: default:
ctx, cancel := context.WithTimeout(bctx, defaultUTPConnectTimeout) ctx, cancel := context.WithTimeout(bctx, defaultUTPConnectTimeout)
var conn *utp.Conn var conn *utp.Conn
p.Log.Debug("will accept offer conn from: ", "source", addr, "connId", connId)
conn, err = p.utp.AcceptUTPContext(ctx, connIdSend) conn, err = p.utp.AcceptUTPContext(ctx, connIdSend)
p.Log.Info("will accept from: ", "source", addr, "connId", connId)
if err != nil { if err != nil {
p.Log.Error("failed to accept utp connection for handle offer", "connId", connIdSend, "err", err) p.Log.Error("failed to accept utp connection", "connId", connIdSend, "err", err)
cancel() cancel()
return return
} }