Peer remote read WaitGroup and add readErr chan in struct

This commit is contained in:
root 2018-10-28 12:47:15 +03:30
parent aeb733623e
commit 3a74aec06b

View file

@ -108,6 +108,7 @@ type Peer struct {
wg sync.WaitGroup wg sync.WaitGroup
protoErr chan error protoErr chan error
readErr chan error
closed chan struct{} closed chan struct{}
disc chan DiscReason disc chan DiscReason
@ -183,6 +184,7 @@ func newPeer(conn *conn, protocols []Protocol) *Peer {
running: protomap, running: protomap,
created: mclock.Now(), created: mclock.Now(),
disc: make(chan DiscReason), disc: make(chan DiscReason),
readErr: make(chan error),
protoErr: make(chan error, len(protomap)+1), // protocols + pingLoop protoErr: make(chan error, len(protomap)+1), // protocols + pingLoop
closed: make(chan struct{}), closed: make(chan struct{}),
log: log.New("id", conn.node.ID(), "conn", conn.flags), log: log.New("id", conn.node.ID(), "conn", conn.flags),
@ -190,6 +192,7 @@ func newPeer(conn *conn, protocols []Protocol) *Peer {
return p return p
} }
//Log return logger
func (p *Peer) Log() log.Logger { func (p *Peer) Log() log.Logger {
return p.log return p.log
} }
@ -198,11 +201,10 @@ func (p *Peer) run() (remoteRequested bool, err error) {
var ( var (
writeStart = make(chan struct{}, 1) writeStart = make(chan struct{}, 1)
writeErr = make(chan error, 1) writeErr = make(chan error, 1)
readErr = make(chan error, 1)
reason DiscReason // sent to the peer reason DiscReason // sent to the peer
) )
p.wg.Add(2) p.wg.Add(1)
go p.readLoop(readErr) go p.readLoop()
go p.pingLoop() go p.pingLoop()
// Start all protocol handlers. // Start all protocol handlers.
@ -221,7 +223,7 @@ loop:
break loop break loop
} }
writeStart <- struct{}{} writeStart <- struct{}{}
case err = <-readErr: case err = <-p.readErr:
if r, ok := err.(DiscReason); ok { if r, ok := err.(DiscReason); ok {
remoteRequested = true remoteRequested = true
reason = r reason = r
@ -246,8 +248,8 @@ loop:
func (p *Peer) pingLoop() { func (p *Peer) pingLoop() {
ping := time.NewTimer(pingInterval) ping := time.NewTimer(pingInterval)
defer p.wg.Done()
defer ping.Stop() defer ping.Stop()
defer p.wg.Done()
for { for {
select { select {
case <-ping.C: case <-ping.C:
@ -262,17 +264,16 @@ func (p *Peer) pingLoop() {
} }
} }
func (p *Peer) readLoop(errc chan<- error) { func (p *Peer) readLoop() {
defer p.wg.Done()
for { for {
msg, err := p.rw.ReadMsg() msg, err := p.rw.ReadMsg()
if err != nil { if err != nil {
errc <- err p.readErr <- err
return return
} }
msg.ReceivedAt = time.Now() msg.ReceivedAt = time.Now()
if err = p.handle(msg); err != nil { if err = p.handle(msg); err != nil {
errc <- err p.readErr <- err
return return
} }
} }