From 64744ac20b8a21b429bb0edd33ba724a6c22d436 Mon Sep 17 00:00:00 2001 From: thinkAfCod Date: Sun, 13 Oct 2024 12:30:19 +0800 Subject: [PATCH 1/6] fix: cache node --- p2p/discover/portal_protocol.go | 44 ++++++++++++++------------------- p2p/discover/v5_udp.go | 24 ++++++++++++++++-- 2 files changed, 41 insertions(+), 27 deletions(-) diff --git a/p2p/discover/portal_protocol.go b/p2p/discover/portal_protocol.go index cceb256d1152..590d9aa29c2d 100644 --- a/p2p/discover/portal_protocol.go +++ b/p2p/discover/portal_protocol.go @@ -169,7 +169,7 @@ func DefaultPortalProtocolConfig() *PortalProtocolConfig { type PortalProtocol struct { table *Table cachedIdsLock sync.Mutex - cachedIds map[string]enode.ID + cachedIds map[string]*enode.Node protocolId string protocolName string @@ -213,7 +213,7 @@ func NewPortalProtocol(config *PortalProtocolConfig, protocolId portalwire.Proto closeCtx, cancelCloseCtx := context.WithCancel(context.Background()) protocol := &PortalProtocol{ - cachedIds: make(map[string]enode.ID), + cachedIds: make(map[string]*enode.Node), protocolId: string(protocolId), protocolName: protocolId.Name(), ListenAddr: config.ListenAddr, @@ -330,10 +330,10 @@ func (p *PortalProtocol) setupUDPListening() error { p.cachedIdsLock.Lock() defer p.cachedIdsLock.Unlock() - if id, ok := p.cachedIds[addr.String()]; ok { + if n, ok := p.cachedIds[addr.String()]; ok { //_, err := p.DiscV5.TalkRequestToID(id, addr, string(portalwire.UTPNetwork), buf) req := &v5wire.TalkRequest{Protocol: string(portalwire.Utp), Message: buf} - p.DiscV5.sendFromAnotherThread(id, netip.AddrPortFrom(netutil.IPToAddr(addr.IP), uint16(addr.Port)), req) + p.DiscV5.sendFromAnotherThreadWithNode(n, netip.AddrPortFrom(netutil.IPToAddr(addr.IP), uint16(addr.Port)), req) return len(buf), err } else { @@ -388,21 +388,22 @@ func (p *PortalProtocol) setupDiscV5AndTable() error { return nil } -func (p *PortalProtocol) putCacheNodeId(node *enode.Node) { +func (p *PortalProtocol) cacheNode(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() + p.cachedIds[addr.String()] = node } } -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) cacheNodeById(id enode.ID, addr *net.UDPAddr) { + go func() { + if _, ok := p.cachedIds[addr.String()]; !ok { + n := p.ResolveNodeId(id) + p.cacheNode(n) + } + }() } func (p *PortalProtocol) ping(node *enode.Node) (uint64, error) { @@ -566,7 +567,7 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request * } p.Log.Info("will process Offer", "id", target.ID(), "ip", target.IP().To4().String(), "port", target.UDP()) - p.putCacheNodeId(target) + p.cacheNode(target) accept := &portalwire.Accept{} err = accept.UnmarshalSSZ(resp[1:]) @@ -704,7 +705,7 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, } p.Log.Info("will process content", "id", target.ID(), "ip", target.IP().To4().String(), "port", target.UDP()) - p.putCacheNodeId(target) + p.cacheNode(target) switch resp[1] { case portalwire.ContentRawSelector: @@ -910,21 +911,14 @@ func (p *PortalProtocol) processPong(target *enode.Node, resp []byte) (*portalwi } func (p *PortalProtocol) handleUtpTalkRequest(id enode.ID, addr *net.UDPAddr, msg []byte) []byte { - if n := p.DiscV5.getNode(id); n != nil { - p.table.addInboundNode(n) - } - - p.putCacheId(id, addr) + p.cacheNodeById(id, addr) p.Log.Trace("receive utp data", "addr", addr, "msg-length", len(msg)) p.packetRouter.ReceiveMessage(msg, addr) return []byte("") } func (p *PortalProtocol) handleTalkRequest(id enode.ID, addr *net.UDPAddr, msg []byte) []byte { - if n := p.DiscV5.getNode(id); n != nil { - p.table.addInboundNode(n) - } - p.putCacheId(id, addr) + p.cacheNodeById(id, addr) msgCode := msg[0] @@ -1107,7 +1101,7 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque return nil, err } - p.putCacheId(id, addr) + p.cacheNodeById(id, addr) if errors.Is(err, ContentNotFound) { closestNodes := p.findNodesCloseToContent(contentId, portalFindnodesResultLimit) @@ -1303,7 +1297,7 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po } } - p.putCacheId(id, addr) + p.cacheNodeById(id, addr) idBuffer := make([]byte, 2) if contentKeyBitlist.Count() != 0 { diff --git a/p2p/discover/v5_udp.go b/p2p/discover/v5_udp.go index d4d9a054d9f1..fa043852ce28 100644 --- a/p2p/discover/v5_udp.go +++ b/p2p/discover/v5_udp.go @@ -102,6 +102,7 @@ type UDPv5 struct { type sendRequest struct { destID enode.ID + destNode *enode.Node destAddr netip.AddrPort msg v5wire.Packet } @@ -596,7 +597,19 @@ func (t *UDPv5) dispatch() { t.sendNextCall(c.id) case r := <-t.sendCh: - t.send(r.destID, r.destAddr, r.msg, nil) + c := &callV5{id: r.destID, addr: r.destAddr} + c.node = r.destNode + c.packet = r.msg + c.reqid = make([]byte, 8) + c.ch = make(chan v5wire.Packet, 1) + c.err = make(chan error, 1) + // Assign request ID. + crand.Read(c.reqid) + c.packet.SetRequestID(c.reqid) + newNonce, _ := t.send(c.id, c.addr, c.packet, nil) + c.nonce = newNonce + t.activeCallByAuth[newNonce] = c + t.startResponseTimeout(c) case p := <-t.packetInCh: t.handlePacket(p.Data, p.Addr) @@ -681,7 +694,14 @@ func (t *UDPv5) sendResponse(toID enode.ID, toAddr netip.AddrPort, packet v5wire func (t *UDPv5) sendFromAnotherThread(toID enode.ID, toAddr netip.AddrPort, packet v5wire.Packet) { select { - case t.sendCh <- sendRequest{toID, toAddr, packet}: + case t.sendCh <- sendRequest{toID, nil, toAddr, packet}: + case <-t.closeCtx.Done(): + } +} + +func (t *UDPv5) sendFromAnotherThreadWithNode(node *enode.Node, toAddr netip.AddrPort, packet v5wire.Packet) { + select { + case t.sendCh <- sendRequest{node.ID(), node, toAddr, packet}: case <-t.closeCtx.Done(): } } From befa40941e3278010ed13b78938ee2e33cca222f Mon Sep 17 00:00:00 2001 From: thinkAfCod Date: Sun, 13 Oct 2024 13:20:34 +0800 Subject: [PATCH 2/6] fix: log packet nonce --- p2p/discover/v5_udp.go | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/p2p/discover/v5_udp.go b/p2p/discover/v5_udp.go index fa043852ce28..87d5cb4882b3 100644 --- a/p2p/discover/v5_udp.go +++ b/p2p/discover/v5_udp.go @@ -713,6 +713,7 @@ func (t *UDPv5) send(toID enode.ID, toAddr netip.AddrPort, packet v5wire.Packet, t.logcontext = packet.AppendLogInfo(t.logcontext) enc, nonce, err := t.codec.Encode(toID, addr, packet, c) + t.logcontext = append(t.logcontext, "nonce", fmt.Sprintf("%x", nonce[:])) if err != nil { t.logcontext = append(t.logcontext, "err", err) t.log.Warn(">> "+packet.Name(), t.logcontext...) @@ -866,7 +867,7 @@ var ( func (t *UDPv5) handleWhoareyou(p *v5wire.Whoareyou, fromID enode.ID, fromAddr netip.AddrPort) { c, err := t.matchWithCall(fromID, p.Nonce) if err != nil { - t.log.Debug("Invalid "+p.Name(), "addr", fromAddr, "err", err) + t.log.Debug("Invalid "+p.Name(), "addr", fromAddr, "nonce", fmt.Sprintf("%x", p.Nonce[:]), "err", err) return } @@ -877,7 +878,7 @@ func (t *UDPv5) handleWhoareyou(p *v5wire.Whoareyou, fromID enode.ID, fromAddr n return } // Resend the call that was answered by WHOAREYOU. - t.log.Trace("<< "+p.Name(), "id", c.node.ID(), "addr", fromAddr) + t.log.Trace("<< "+p.Name(), "id", c.node.ID(), "addr", fromAddr, "nonce", fmt.Sprintf("%x", p.Nonce[:])) c.handshakeCount++ c.challenge = p p.Node = c.node From 3f0105712a54942aeb73efb692c5f2d5d4a45d66 Mon Sep 17 00:00:00 2001 From: thinkAfCod Date: Mon, 14 Oct 2024 11:20:37 +0800 Subject: [PATCH 3/6] fix: test case modify --- p2p/discover/portal_protocol_test.go | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/p2p/discover/portal_protocol_test.go b/p2p/discover/portal_protocol_test.go index 8b4732fd510f..efc49ae7c314 100644 --- a/p2p/discover/portal_protocol_test.go +++ b/p2p/discover/portal_protocol_test.go @@ -114,14 +114,14 @@ func TestPortalWireProtocolUdp(t *testing.T) { assert.NoError(t, err) time.Sleep(12 * time.Second) - node1.putCacheNodeId(node2.localNode.Node()) - node1.putCacheNodeId(node3.localNode.Node()) + node1.cacheNode(node2.localNode.Node()) + node1.cacheNode(node3.localNode.Node()) - node2.putCacheNodeId(node1.localNode.Node()) - node2.putCacheNodeId(node3.localNode.Node()) + node2.cacheNode(node1.localNode.Node()) + node2.cacheNode(node3.localNode.Node()) - node3.putCacheNodeId(node1.localNode.Node()) - node3.putCacheNodeId(node2.localNode.Node()) + node3.cacheNode(node1.localNode.Node()) + node3.cacheNode(node2.localNode.Node()) udpAddrStr1 := fmt.Sprintf("%s:%d", node1.localNode.Node().IP(), node1.localNode.Node().UDP()) udpAddrStr2 := fmt.Sprintf("%s:%d", node2.localNode.Node().IP(), node2.localNode.Node().UDP()) From 4ddb2f28c1d106d82d7b7ee81e49cf704603c4cb Mon Sep 17 00:00:00 2001 From: thinkAfCod Date: Wed, 16 Oct 2024 20:37:15 +0800 Subject: [PATCH 4/6] fix: not cache nil node --- p2p/discover/portal_protocol.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/p2p/discover/portal_protocol.go b/p2p/discover/portal_protocol.go index 590d9aa29c2d..66e170e70ca7 100644 --- a/p2p/discover/portal_protocol.go +++ b/p2p/discover/portal_protocol.go @@ -401,7 +401,9 @@ func (p *PortalProtocol) cacheNodeById(id enode.ID, addr *net.UDPAddr) { go func() { if _, ok := p.cachedIds[addr.String()]; !ok { n := p.ResolveNodeId(id) - p.cacheNode(n) + if n != nil { + p.cacheNode(n) + } } }() } From f09031f9e8016627f023a5f78e528c40dcfd45dd Mon Sep 17 00:00:00 2001 From: thinkAfCod Date: Thu, 17 Oct 2024 21:54:35 +0800 Subject: [PATCH 5/6] fix: check lint --- p2p/discover/portal_protocol.go | 29 +++++++++++++++---- p2p/discover/v5_udp.go | 15 ++++++---- portalnetwork/history/history_network_test.go | 5 ++++ portalnetwork/history/storage.go | 2 +- 4 files changed, 39 insertions(+), 12 deletions(-) diff --git a/p2p/discover/portal_protocol.go b/p2p/discover/portal_protocol.go index 66e170e70ca7..d9a39fa452a9 100644 --- a/p2p/discover/portal_protocol.go +++ b/p2p/discover/portal_protocol.go @@ -169,7 +169,8 @@ func DefaultPortalProtocolConfig() *PortalProtocolConfig { type PortalProtocol struct { table *Table cachedIdsLock sync.Mutex - cachedIds map[string]*enode.Node + cachedNodes map[string]*enode.Node + cachedIds map[string]enode.ID protocolId string protocolName string @@ -213,7 +214,8 @@ func NewPortalProtocol(config *PortalProtocolConfig, protocolId portalwire.Proto closeCtx, cancelCloseCtx := context.WithCancel(context.Background()) protocol := &PortalProtocol{ - cachedIds: make(map[string]*enode.Node), + cachedNodes: make(map[string]*enode.Node), + cachedIds: make(map[string]enode.ID), protocolId: string(protocolId), protocolName: protocolId.Name(), ListenAddr: config.ListenAddr, @@ -330,11 +332,17 @@ func (p *PortalProtocol) setupUDPListening() error { p.cachedIdsLock.Lock() defer p.cachedIdsLock.Unlock() - if n, ok := p.cachedIds[addr.String()]; ok { + if n, ok := p.cachedNodes[addr.String()]; ok { //_, err := p.DiscV5.TalkRequestToID(id, addr, string(portalwire.UTPNetwork), buf) req := &v5wire.TalkRequest{Protocol: string(portalwire.Utp), Message: buf} p.DiscV5.sendFromAnotherThreadWithNode(n, netip.AddrPortFrom(netutil.IPToAddr(addr.IP), uint16(addr.Port)), req) + return len(buf), err + } else if id, ok := p.cachedIds[addr.String()]; ok { + //_, err := p.DiscV5.TalkRequestToID(id, addr, string(portalwire.UTPNetwork), buf) + req := &v5wire.TalkRequest{Protocol: string(portalwire.Utp), Message: buf} + p.DiscV5.sendFromAnotherThread(id, netip.AddrPortFrom(netutil.IPToAddr(addr.IP), uint16(addr.Port)), req) + return len(buf), err } else { p.Log.Warn("not found target node info", "ip", addr.IP.To4().String(), "port", addr.Port, "bufLength", len(buf)) @@ -392,14 +400,23 @@ func (p *PortalProtocol) cacheNode(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 + if _, ok := p.cachedNodes[addr.String()]; !ok && node != nil { + p.cachedNodes[addr.String()] = node + } +} + +func (p *PortalProtocol) cacheNodeId(id enode.ID, addr *net.UDPAddr) { + p.cachedIdsLock.Lock() + defer p.cachedIdsLock.Unlock() + if (id != enode.ID{}) { + p.cachedIds[addr.String()] = id } } func (p *PortalProtocol) cacheNodeById(id enode.ID, addr *net.UDPAddr) { go func() { - if _, ok := p.cachedIds[addr.String()]; !ok { + p.cacheNodeId(id, addr) + if _, ok := p.cachedNodes[addr.String()]; !ok { n := p.ResolveNodeId(id) if n != nil { p.cacheNode(n) diff --git a/p2p/discover/v5_udp.go b/p2p/discover/v5_udp.go index 87d5cb4882b3..0422e9009dac 100644 --- a/p2p/discover/v5_udp.go +++ b/p2p/discover/v5_udp.go @@ -604,12 +604,17 @@ func (t *UDPv5) dispatch() { c.ch = make(chan v5wire.Packet, 1) c.err = make(chan error, 1) // Assign request ID. - crand.Read(c.reqid) - c.packet.SetRequestID(c.reqid) - newNonce, _ := t.send(c.id, c.addr, c.packet, nil) - c.nonce = newNonce - t.activeCallByAuth[newNonce] = c + if tq, ok := r.msg.(*v5wire.TalkRequest); ok { + if len(tq.ReqID) == 0 { + crand.Read(c.reqid) + c.packet.SetRequestID(c.reqid) + } + } + nonce, _ := t.send(c.id, c.addr, c.packet, nil) + c.nonce = nonce + t.activeCallByAuth[nonce] = c t.startResponseTimeout(c) + //t.send(r.destID, r.destAddr, r.msg, nil) case p := <-t.packetInCh: t.handlePacket(p.Data, p.Addr) diff --git a/portalnetwork/history/history_network_test.go b/portalnetwork/history/history_network_test.go index e1bba686313a..b701bcacff84 100644 --- a/portalnetwork/history/history_network_test.go +++ b/portalnetwork/history/history_network_test.go @@ -140,6 +140,11 @@ func TestGetContentByKey(t *testing.T) { contentId := historyNetwork1.portalProtocol.ToContentId(headerEntry.key) err = historyNetwork1.portalProtocol.Put(headerEntry.key, contentId, headerEntry.value) require.NoError(t, err) + + header, err = historyNetwork1.GetBlockHeader(headerEntry.key[1:]) + require.NoError(t, err) + require.NotNil(t, header) + // get content from historyNetwork1 header, err = historyNetwork2.GetBlockHeader(headerEntry.key[1:]) require.NoError(t, err) diff --git a/portalnetwork/history/storage.go b/portalnetwork/history/storage.go index 7878093f6b52..488d2ff3622c 100644 --- a/portalnetwork/history/storage.go +++ b/portalnetwork/history/storage.go @@ -452,7 +452,7 @@ func (p *ContentStorage) deleteContentOutOfRadius(radius *uint256.Int) error { return err } count, _ := res.RowsAffected() - p.log.Trace("delete %d items", count) + p.log.Trace("delete items", "count", count) return err } From 7f99d4ff11debffb009abc82ec76fed66d556fb3 Mon Sep 17 00:00:00 2001 From: thinkAfCod Date: Tue, 22 Oct 2024 16:23:19 +0800 Subject: [PATCH 6/6] fix: cache id and addr in discv5 revert --- p2p/discover/portal_protocol.go | 60 +++---------------- p2p/discover/portal_protocol_test.go | 9 --- p2p/discover/v5_udp.go | 44 ++++++++------ portalnetwork/history/history_network_test.go | 5 -- 4 files changed, 35 insertions(+), 83 deletions(-) diff --git a/p2p/discover/portal_protocol.go b/p2p/discover/portal_protocol.go index d9a39fa452a9..d2b26d09c6a9 100644 --- a/p2p/discover/portal_protocol.go +++ b/p2p/discover/portal_protocol.go @@ -167,10 +167,7 @@ func DefaultPortalProtocolConfig() *PortalProtocolConfig { } type PortalProtocol struct { - table *Table - cachedIdsLock sync.Mutex - cachedNodes map[string]*enode.Node - cachedIds map[string]enode.ID + table *Table protocolId string protocolName string @@ -214,8 +211,6 @@ func NewPortalProtocol(config *PortalProtocolConfig, protocolId portalwire.Proto closeCtx, cancelCloseCtx := context.WithCancel(context.Background()) protocol := &PortalProtocol{ - cachedNodes: make(map[string]*enode.Node), - cachedIds: make(map[string]enode.ID), protocolId: string(protocolId), protocolName: protocolId.Name(), ListenAddr: config.ListenAddr, @@ -330,19 +325,11 @@ func (p *PortalProtocol) setupUDPListening() 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)) - p.cachedIdsLock.Lock() - defer p.cachedIdsLock.Unlock() - if n, ok := p.cachedNodes[addr.String()]; ok { + if n, ok := p.DiscV5.cachedAddrNode[addr.String()]; ok { //_, err := p.DiscV5.TalkRequestToID(id, addr, string(portalwire.UTPNetwork), buf) req := &v5wire.TalkRequest{Protocol: string(portalwire.Utp), Message: buf} p.DiscV5.sendFromAnotherThreadWithNode(n, netip.AddrPortFrom(netutil.IPToAddr(addr.IP), uint16(addr.Port)), req) - return len(buf), err - } else if id, ok := p.cachedIds[addr.String()]; ok { - //_, err := p.DiscV5.TalkRequestToID(id, addr, string(portalwire.UTPNetwork), buf) - req := &v5wire.TalkRequest{Protocol: string(portalwire.Utp), Message: buf} - p.DiscV5.sendFromAnotherThread(id, netip.AddrPortFrom(netutil.IPToAddr(addr.IP), uint16(addr.Port)), req) - return len(buf), err } else { p.Log.Warn("not found target node info", "ip", addr.IP.To4().String(), "port", addr.Port, "bufLength", len(buf)) @@ -396,35 +383,6 @@ func (p *PortalProtocol) setupDiscV5AndTable() error { return nil } -func (p *PortalProtocol) cacheNode(node *enode.Node) { - p.cachedIdsLock.Lock() - defer p.cachedIdsLock.Unlock() - addr := &net.UDPAddr{IP: node.IP(), Port: node.UDP()} - if _, ok := p.cachedNodes[addr.String()]; !ok && node != nil { - p.cachedNodes[addr.String()] = node - } -} - -func (p *PortalProtocol) cacheNodeId(id enode.ID, addr *net.UDPAddr) { - p.cachedIdsLock.Lock() - defer p.cachedIdsLock.Unlock() - if (id != enode.ID{}) { - p.cachedIds[addr.String()] = id - } -} - -func (p *PortalProtocol) cacheNodeById(id enode.ID, addr *net.UDPAddr) { - go func() { - p.cacheNodeId(id, addr) - if _, ok := p.cachedNodes[addr.String()]; !ok { - n := p.ResolveNodeId(id) - if n != nil { - p.cacheNode(n) - } - } - }() -} - func (p *PortalProtocol) ping(node *enode.Node) (uint64, error) { pong, err := p.pingInner(node) if err != nil { @@ -586,7 +544,6 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request * } p.Log.Info("will process Offer", "id", target.ID(), "ip", target.IP().To4().String(), "port", target.UDP()) - p.cacheNode(target) accept := &portalwire.Accept{} err = accept.UnmarshalSSZ(resp[1:]) @@ -724,7 +681,6 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, } p.Log.Info("will process content", "id", target.ID(), "ip", target.IP().To4().String(), "port", target.UDP()) - p.cacheNode(target) switch resp[1] { case portalwire.ContentRawSelector: @@ -930,14 +886,18 @@ func (p *PortalProtocol) processPong(target *enode.Node, resp []byte) (*portalwi } func (p *PortalProtocol) handleUtpTalkRequest(id enode.ID, addr *net.UDPAddr, msg []byte) []byte { - p.cacheNodeById(id, addr) + if n := p.DiscV5.getNode(id); n != nil { + p.table.addInboundNode(n) + } p.Log.Trace("receive utp data", "addr", addr, "msg-length", len(msg)) p.packetRouter.ReceiveMessage(msg, addr) return []byte("") } func (p *PortalProtocol) handleTalkRequest(id enode.ID, addr *net.UDPAddr, msg []byte) []byte { - p.cacheNodeById(id, addr) + if n := p.DiscV5.getNode(id); n != nil { + p.table.addInboundNode(n) + } msgCode := msg[0] @@ -1120,8 +1080,6 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque return nil, err } - p.cacheNodeById(id, addr) - if errors.Is(err, ContentNotFound) { closestNodes := p.findNodesCloseToContent(contentId, portalFindnodesResultLimit) for i, n := range closestNodes { @@ -1316,8 +1274,6 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po } } - p.cacheNodeById(id, addr) - idBuffer := make([]byte, 2) if contentKeyBitlist.Count() != 0 { connId := p.connIdGen.GenCid(id, false) diff --git a/p2p/discover/portal_protocol_test.go b/p2p/discover/portal_protocol_test.go index efc49ae7c314..7f42ee71f4c6 100644 --- a/p2p/discover/portal_protocol_test.go +++ b/p2p/discover/portal_protocol_test.go @@ -114,15 +114,6 @@ func TestPortalWireProtocolUdp(t *testing.T) { assert.NoError(t, err) time.Sleep(12 * time.Second) - node1.cacheNode(node2.localNode.Node()) - node1.cacheNode(node3.localNode.Node()) - - node2.cacheNode(node1.localNode.Node()) - node2.cacheNode(node3.localNode.Node()) - - node3.cacheNode(node1.localNode.Node()) - node3.cacheNode(node2.localNode.Node()) - udpAddrStr1 := fmt.Sprintf("%s:%d", node1.localNode.Node().IP(), node1.localNode.Node().UDP()) udpAddrStr2 := fmt.Sprintf("%s:%d", node2.localNode.Node().IP(), node2.localNode.Node().UDP()) diff --git a/p2p/discover/v5_udp.go b/p2p/discover/v5_udp.go index 0422e9009dac..92a97929ea4a 100644 --- a/p2p/discover/v5_udp.go +++ b/p2p/discover/v5_udp.go @@ -62,15 +62,17 @@ type codecV5 interface { // UDPv5 is the implementation of protocol version 5. type UDPv5 struct { // static fields - conn UDPConn - tab *Table - netrestrict *netutil.Netlist - priv *ecdsa.PrivateKey - localNode *enode.LocalNode - db *enode.DB - log log.Logger - clock mclock.Clock - validSchemes enr.IdentityScheme + conn UDPConn + tab *Table + cachedIds map[enode.ID]*enode.Node + cachedAddrNode map[string]*enode.Node + netrestrict *netutil.Netlist + priv *ecdsa.PrivateKey + localNode *enode.LocalNode + db *enode.DB + log log.Logger + clock mclock.Clock + validSchemes enr.IdentityScheme // misc buffers used during message handling logcontext []interface{} @@ -151,14 +153,16 @@ func newUDPv5(conn UDPConn, ln *enode.LocalNode, cfg Config) (*UDPv5, error) { cfg = cfg.withDefaults() t := &UDPv5{ // static fields - conn: newMeteredConn(conn), - localNode: ln, - db: ln.Database(), - netrestrict: cfg.NetRestrict, - priv: cfg.PrivateKey, - log: cfg.Log, - validSchemes: cfg.ValidSchemes, - clock: cfg.Clock, + conn: newMeteredConn(conn), + cachedAddrNode: make(map[string]*enode.Node), + cachedIds: make(map[enode.ID]*enode.Node), + localNode: ln, + db: ln.Database(), + netrestrict: cfg.NetRestrict, + priv: cfg.PrivateKey, + log: cfg.Log, + validSchemes: cfg.ValidSchemes, + clock: cfg.Clock, // channels into dispatch packetInCh: make(chan ReadPacket, 1), readNextCh: make(chan struct{}, 1), @@ -724,6 +728,10 @@ func (t *UDPv5) send(toID enode.ID, toAddr netip.AddrPort, packet v5wire.Packet, t.log.Warn(">> "+packet.Name(), t.logcontext...) return nonce, err } + if c != nil && c.Node != nil { + t.cachedIds[toID] = c.Node + t.cachedAddrNode[toAddr.String()] = c.Node + } _, err = t.conn.WriteToUDPAddrPort(enc, toAddr) t.log.Trace(">> "+packet.Name(), t.logcontext...) @@ -785,6 +793,8 @@ func (t *UDPv5) handlePacket(rawpacket []byte, fromAddr netip.AddrPort) error { if fromNode != nil { // Handshake succeeded, add to table. t.tab.addInboundNode(fromNode) + t.cachedIds[fromID] = fromNode + t.cachedAddrNode[fromAddr.String()] = fromNode } if packet.Kind() != v5wire.WhoareyouPacket { // WHOAREYOU logged separately to report errors. diff --git a/portalnetwork/history/history_network_test.go b/portalnetwork/history/history_network_test.go index b701bcacff84..e1bba686313a 100644 --- a/portalnetwork/history/history_network_test.go +++ b/portalnetwork/history/history_network_test.go @@ -140,11 +140,6 @@ func TestGetContentByKey(t *testing.T) { contentId := historyNetwork1.portalProtocol.ToContentId(headerEntry.key) err = historyNetwork1.portalProtocol.Put(headerEntry.key, contentId, headerEntry.value) require.NoError(t, err) - - header, err = historyNetwork1.GetBlockHeader(headerEntry.key[1:]) - require.NoError(t, err) - require.NotNil(t, header) - // get content from historyNetwork1 header, err = historyNetwork2.GetBlockHeader(headerEntry.key[1:]) require.NoError(t, err)