mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-27 07:06:42 +00:00
p2p/discover: send a packet and not wait for a response
Signed-off-by: thinkAfCod <q315xia@163.com>
This commit is contained in:
parent
d2176f463b
commit
4a5923668a
1 changed files with 72 additions and 7 deletions
|
|
@ -96,12 +96,14 @@ type UDPv5 struct {
|
||||||
callDoneCh chan *callV5
|
callDoneCh chan *callV5
|
||||||
respTimeoutCh chan *callTimeout
|
respTimeoutCh chan *callTimeout
|
||||||
sendCh chan sendRequest
|
sendCh chan sendRequest
|
||||||
|
sendNoRespCh chan *sendNoRespRequest
|
||||||
unhandled chan<- ReadPacket
|
unhandled chan<- ReadPacket
|
||||||
|
|
||||||
// state of dispatch
|
// state of dispatch
|
||||||
codec codecV5
|
codec codecV5
|
||||||
activeCallByNode map[enode.ID]*callV5
|
activeCallByNode map[enode.ID]*callV5
|
||||||
activeCallByAuth map[v5wire.Nonce]*callV5
|
activeCallByAuth map[v5wire.Nonce]*callV5
|
||||||
|
noRespCallByAuth map[v5wire.Nonce]*callV5
|
||||||
callQueue map[enode.ID][]*callV5
|
callQueue map[enode.ID][]*callV5
|
||||||
|
|
||||||
// shutdown stuff
|
// shutdown stuff
|
||||||
|
|
@ -115,6 +117,13 @@ type sendRequest struct {
|
||||||
destID enode.ID
|
destID enode.ID
|
||||||
destAddr netip.AddrPort
|
destAddr netip.AddrPort
|
||||||
msg v5wire.Packet
|
msg v5wire.Packet
|
||||||
|
destNode *enode.Node
|
||||||
|
}
|
||||||
|
|
||||||
|
type sendNoRespRequest struct {
|
||||||
|
destNode *enode.Node
|
||||||
|
destAddr netip.AddrPort
|
||||||
|
msg v5wire.Packet
|
||||||
}
|
}
|
||||||
|
|
||||||
// callV5 represents a remote procedure call against another node.
|
// callV5 represents a remote procedure call against another node.
|
||||||
|
|
@ -176,12 +185,14 @@ func newUDPv5(conn UDPConn, ln *enode.LocalNode, cfg Config) (*UDPv5, error) {
|
||||||
callCh: make(chan *callV5),
|
callCh: make(chan *callV5),
|
||||||
callDoneCh: make(chan *callV5),
|
callDoneCh: make(chan *callV5),
|
||||||
sendCh: make(chan sendRequest),
|
sendCh: make(chan sendRequest),
|
||||||
|
sendNoRespCh: make(chan *sendNoRespRequest),
|
||||||
respTimeoutCh: make(chan *callTimeout),
|
respTimeoutCh: make(chan *callTimeout),
|
||||||
unhandled: cfg.Unhandled,
|
unhandled: cfg.Unhandled,
|
||||||
// state of dispatch
|
// state of dispatch
|
||||||
codec: v5wire.NewCodec(ln, cfg.PrivateKey, cfg.Clock, cfg.V5ProtocolID),
|
codec: v5wire.NewCodec(ln, cfg.PrivateKey, cfg.Clock, cfg.V5ProtocolID),
|
||||||
activeCallByNode: make(map[enode.ID]*callV5),
|
activeCallByNode: make(map[enode.ID]*callV5),
|
||||||
activeCallByAuth: make(map[v5wire.Nonce]*callV5),
|
activeCallByAuth: make(map[v5wire.Nonce]*callV5),
|
||||||
|
noRespCallByAuth: make(map[v5wire.Nonce]*callV5),
|
||||||
callQueue: make(map[enode.ID][]*callV5),
|
callQueue: make(map[enode.ID][]*callV5),
|
||||||
// shutdown
|
// shutdown
|
||||||
closeCtx: closeCtx,
|
closeCtx: closeCtx,
|
||||||
|
|
@ -581,12 +592,18 @@ func (t *UDPv5) dispatch() {
|
||||||
case c := <-t.callCh:
|
case c := <-t.callCh:
|
||||||
t.callQueue[c.id] = append(t.callQueue[c.id], c)
|
t.callQueue[c.id] = append(t.callQueue[c.id], c)
|
||||||
t.sendNextCall(c.id)
|
t.sendNextCall(c.id)
|
||||||
|
case cnr := <-t.sendNoRespCh:
|
||||||
|
// send a TalkReq call but not waiting for TalkResp
|
||||||
|
// may more used by portal network
|
||||||
|
t.sendCallNotWaitResp(cnr)
|
||||||
|
|
||||||
case ct := <-t.respTimeoutCh:
|
case ct := <-t.respTimeoutCh:
|
||||||
active := t.activeCallByNode[ct.c.id]
|
active := t.activeCallByNode[ct.c.id]
|
||||||
if ct.c == active && ct.timer == active.timeout {
|
if ct.c == active && ct.timer == active.timeout {
|
||||||
ct.c.err <- errTimeout
|
ct.c.err <- errTimeout
|
||||||
}
|
}
|
||||||
|
delete(t.noRespCallByAuth, ct.c.nonce)
|
||||||
|
ct.c.timeout.Stop()
|
||||||
|
|
||||||
case c := <-t.callDoneCh:
|
case c := <-t.callDoneCh:
|
||||||
active := t.activeCallByNode[c.id]
|
active := t.activeCallByNode[c.id]
|
||||||
|
|
@ -600,7 +617,6 @@ func (t *UDPv5) dispatch() {
|
||||||
|
|
||||||
case r := <-t.sendCh:
|
case r := <-t.sendCh:
|
||||||
t.send(r.destID, r.destAddr, r.msg, nil)
|
t.send(r.destID, r.destAddr, r.msg, nil)
|
||||||
|
|
||||||
case p := <-t.packetInCh:
|
case p := <-t.packetInCh:
|
||||||
t.handlePacket(p.Data, p.Addr)
|
t.handlePacket(p.Data, p.Addr)
|
||||||
// Arm next read.
|
// Arm next read.
|
||||||
|
|
@ -619,6 +635,9 @@ func (t *UDPv5) dispatch() {
|
||||||
delete(t.activeCallByNode, id)
|
delete(t.activeCallByNode, id)
|
||||||
delete(t.activeCallByAuth, c.nonce)
|
delete(t.activeCallByAuth, c.nonce)
|
||||||
}
|
}
|
||||||
|
for nonce := range t.noRespCallByAuth {
|
||||||
|
delete(t.activeCallByAuth, nonce)
|
||||||
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -644,6 +663,36 @@ func (t *UDPv5) startResponseTimeout(c *callV5) {
|
||||||
close(done)
|
close(done)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// sendCallNotWaitResp send a talk request contains utp packet by call, call will not insert into queue
|
||||||
|
func (t *UDPv5) sendCallNotWaitResp(r *sendNoRespRequest) {
|
||||||
|
// send out a TalkRequest that is a UTP packet for portal network
|
||||||
|
// todo If the destination node has been handshaked, it may not be necessary to use call
|
||||||
|
// request should be cached to handle WHOAREYOU
|
||||||
|
c := &callV5{id: r.destNode.ID(), 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)
|
||||||
|
|
||||||
|
nonce, _ := t.send(c.id, c.addr, c.packet, nil)
|
||||||
|
c.nonce = nonce
|
||||||
|
t.noRespCallByAuth[nonce] = c
|
||||||
|
t.startResponseTimeout(c)
|
||||||
|
}
|
||||||
|
|
||||||
|
// sendNoRespData send a data from a call of no wait resp.
|
||||||
|
// The handshake just has been successful.
|
||||||
|
func (t *UDPv5) sendNoRespData(c *callV5) {
|
||||||
|
// Just resend the data and not use call again
|
||||||
|
delete(t.noRespCallByAuth, c.nonce)
|
||||||
|
c.timeout.Stop()
|
||||||
|
t.send(c.node.ID(), c.addr, c.packet, nil)
|
||||||
|
}
|
||||||
|
|
||||||
// sendNextCall sends the next call in the call queue if there is no active call.
|
// sendNextCall sends the next call in the call queue if there is no active call.
|
||||||
func (t *UDPv5) sendNextCall(id enode.ID) {
|
func (t *UDPv5) sendNextCall(id enode.ID) {
|
||||||
queue := t.callQueue[id]
|
queue := t.callQueue[id]
|
||||||
|
|
@ -684,7 +733,15 @@ 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) {
|
func (t *UDPv5) sendFromAnotherThread(toID enode.ID, toAddr netip.AddrPort, packet v5wire.Packet) {
|
||||||
select {
|
select {
|
||||||
case t.sendCh <- sendRequest{toID, toAddr, packet}:
|
case t.sendCh <- sendRequest{destID: toID, destAddr: toAddr, msg: packet}:
|
||||||
|
case <-t.closeCtx.Done():
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// SendNoResp Send a packet but will not wait for a response
|
||||||
|
func (t *UDPv5) SendNoResp(n *enode.Node, toAddr netip.AddrPort, packet v5wire.Packet) {
|
||||||
|
select {
|
||||||
|
case t.sendNoRespCh <- &sendNoRespRequest{n, toAddr, packet}:
|
||||||
case <-t.closeCtx.Done():
|
case <-t.closeCtx.Done():
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -877,17 +934,25 @@ func (t *UDPv5) handleWhoareyou(p *v5wire.Whoareyou, fromID enode.ID, fromAddr n
|
||||||
c.err <- errors.New("remote wants handshake, but call has no ENR")
|
c.err <- errors.New("remote wants handshake, but call has no ENR")
|
||||||
return
|
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)
|
||||||
c.handshakeCount++
|
if _, ok := t.noRespCallByAuth[p.Nonce]; !ok {
|
||||||
c.challenge = p
|
// Resend the call that was answered by WHOAREYOU.
|
||||||
p.Node = c.node
|
c.handshakeCount++
|
||||||
t.sendCall(c)
|
c.challenge = p
|
||||||
|
p.Node = c.node
|
||||||
|
t.sendCall(c)
|
||||||
|
} else {
|
||||||
|
t.sendNoRespData(c)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// matchWithCall checks whether a handshake attempt matches the active call.
|
// matchWithCall checks whether a handshake attempt matches the active call.
|
||||||
func (t *UDPv5) matchWithCall(fromID enode.ID, nonce v5wire.Nonce) (*callV5, error) {
|
func (t *UDPv5) matchWithCall(fromID enode.ID, nonce v5wire.Nonce) (*callV5, error) {
|
||||||
c := t.activeCallByAuth[nonce]
|
c := t.activeCallByAuth[nonce]
|
||||||
|
if c == nil {
|
||||||
|
c = t.noRespCallByAuth[nonce]
|
||||||
|
}
|
||||||
if c == nil {
|
if c == nil {
|
||||||
return nil, errChallengeNoCall
|
return nil, errChallengeNoCall
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue