From 33a000cb3eedc5ede05a0ecb822724f20bc1ce0b Mon Sep 17 00:00:00 2001 From: thinkAfCod Date: Mon, 19 May 2025 14:44:55 +0800 Subject: [PATCH] remove write channel --- p2p/discover/v5_udp.go | 58 +++++++++++++----------------------------- 1 file changed, 17 insertions(+), 41 deletions(-) diff --git a/p2p/discover/v5_udp.go b/p2p/discover/v5_udp.go index 08903a0d97..06e0eefe7f 100644 --- a/p2p/discover/v5_udp.go +++ b/p2p/discover/v5_udp.go @@ -97,7 +97,6 @@ type UDPv5 struct { sendCh chan sendRequest sendNoRespCh chan *sendNoRespRequest unhandled chan<- ReadPacket - writeCh chan pendingWrite // New channel for outgoing packets // state of dispatch codec codecV5 @@ -113,12 +112,6 @@ type UDPv5 struct { wg sync.WaitGroup } -// pendingWrite holds data for a packet to be sent by the writeLoop. -type pendingWrite struct { - toAddr netip.AddrPort - data []byte -} - type sendRequest struct { destID enode.ID destAddr netip.AddrPort @@ -166,7 +159,6 @@ func ListenV5(conn UDPConn, ln *enode.LocalNode, cfg Config) (*UDPv5, error) { t.wg.Add(3) go t.readLoop() go t.dispatch() - go t.writeLoop() return t, nil } @@ -186,14 +178,13 @@ func newUDPv5(conn UDPConn, ln *enode.LocalNode, cfg Config) (*UDPv5, error) { clock: cfg.Clock, respTimeout: cfg.V5RespTimeout, // channels into dispatch - packetInCh: make(chan ReadPacket, 64), + packetInCh: make(chan ReadPacket, 1), callCh: make(chan *callV5), callDoneCh: make(chan *callV5), sendCh: make(chan sendRequest), sendNoRespCh: make(chan *sendNoRespRequest), respTimeoutCh: make(chan *callTimeout), unhandled: cfg.Unhandled, - writeCh: make(chan pendingWrite, 32), // Buffered channel for outgoing packets // state of dispatch codec: v5wire.NewCodec(ln, cfg.PrivateKey, cfg.Clock, cfg.V5ProtocolID), activeCallByNode: make(map[enode.ID]*callV5), @@ -625,7 +616,6 @@ func (t *UDPv5) dispatch() { t.handlePacket(p.Data, p.Addr) case <-t.closeCtx.Done(): - close(t.writeCh) for id, queue := range t.callQueue { for _, c := range queue { c.err <- errClosed @@ -766,40 +756,26 @@ func (t *UDPv5) send(toID enode.ID, toAddr netip.AddrPort, packet v5wire.Packet, return nonce, err } + _, err = t.conn.WriteToUDPAddrPort(enc, toAddr) t.log.Trace(">> "+packet.Name(), t.logcontext...) - - dataForSend := make([]byte, len(enc)) - copy(dataForSend, enc) - - pw := pendingWrite{ - toAddr: toAddr, - data: dataForSend, // codec.Encode should return a new slice, safe to pass directly. - } - - select { - case t.writeCh <- pw: - // Packet successfully queued. - return nonce, nil - case <-t.closeCtx.Done(): - return nonce, errClosed - } + return nonce, err } // writeLoop runs in its own goroutine and sends packets from the writeCh to the network. -func (t *UDPv5) writeLoop() { - defer t.wg.Done() - for pw := range t.writeCh { // Loop continues until writeCh is closed and empty. - _, err := t.conn.WriteToUDPAddrPort(pw.data, pw.toAddr) - if netutil.IsTemporaryError(err) { - t.log.Debug("Temporary UDP write error", "addr", pw.toAddr, "err", err) - } else if err != nil { - if !errors.Is(err, net.ErrClosed) || !errors.Is(err, io.EOF) { - t.log.Warn("UDP write error", "addr", pw.toAddr, "err", err) - } - return - } - } -} +//func (t *UDPv5) writeLoop() { +// defer t.wg.Done() +// for pw := range t.writeCh { // Loop continues until writeCh is closed and empty. +// _, err := t.conn.WriteToUDPAddrPort(pw.data, pw.toAddr) +// if netutil.IsTemporaryError(err) { +// t.log.Debug("Temporary UDP write error", "addr", pw.toAddr, "err", err) +// } else if err != nil { +// if !errors.Is(err, net.ErrClosed) || !errors.Is(err, io.EOF) { +// t.log.Warn("UDP write error", "addr", pw.toAddr, "err", err) +// } +// return +// } +// } +//} // readLoop runs in its own goroutine and reads packets from the network. func (t *UDPv5) readLoop() {