From 6ee96d46326acf579fb15bd65600c7c53a5226b5 Mon Sep 17 00:00:00 2001 From: Atsushi Watanabe Date: Sat, 4 Apr 2020 11:49:56 +0900 Subject: [PATCH] Fix candidate/state callback order Call callbacks serially instead of running asynchronously. --- agent.go | 24 +++++++++++++++++++++--- gather.go | 2 +- 2 files changed, 22 insertions(+), 4 deletions(-) diff --git a/agent.go b/agent.go index 360af5c..3dc1e4e 100644 --- a/agent.go +++ b/agent.go @@ -151,6 +151,9 @@ type Agent struct { done chan struct{} err atomicError + chanCandidateCallback chan func() + chanStateCallback chan func() + loggerFactory logging.LoggerFactory log logging.LeveledLogger @@ -372,6 +375,8 @@ func NewAgent(config *AgentConfig) (*Agent, error) { onConnected: make(chan struct{}), buffer: packetio.NewBuffer(), done: make(chan struct{}), + chanCandidateCallback: make(chan func(), 1), + chanStateCallback: make(chan func(), 1), portmin: config.PortMin, portmax: config.PortMax, trickle: config.Trickle, @@ -423,6 +428,17 @@ func NewAgent(config *AgentConfig) (*Agent, error) { return nil, err } + go func() { + for f := range a.chanCandidateCallback { + f() + } + }() + go func() { + for f := range a.chanStateCallback { + f() + } + }() + // Initialize local candidates if !a.trickle { a.gatherCandidates() @@ -628,9 +644,9 @@ func (a *Agent) updateConnectionState(newState ConnectionState) { a.connectionState = newState hdlr := a.onConnectionStateChangeHdlr if hdlr != nil { - // Call handler async since we may be holding the agent lock + // Call handler in different routine since we may be holding the agent lock // and the handler may also require it - go hdlr(newState) + a.chanStateCallback <- func() { hdlr(newState) } } } } @@ -868,7 +884,7 @@ func (a *Agent) addCandidate(c Candidate, candidateConn net.PacketConn) error { a.requestConnectivityCheck() if a.onCandidateHdlr != nil { - go a.onCandidateHdlr(c) + a.chanCandidateCallback <- func() { a.onCandidateHdlr(c) } } }) } @@ -901,6 +917,8 @@ func (a *Agent) Close() error { done := make(chan struct{}) err := a.run(func(agent *Agent) { defer func() { + close(agent.chanCandidateCallback) + close(agent.chanStateCallback) close(done) }() agent.err.Store(ErrClosed) diff --git a/gather.go b/gather.go index 3b2638d..2dff4ca 100644 --- a/gather.go +++ b/gather.go @@ -97,7 +97,7 @@ func (a *Agent) gatherCandidates() { } if err := a.run(func(agent *Agent) { if a.onCandidateHdlr != nil { - go a.onCandidateHdlr(nil) + a.chanCandidateCallback <- func() { a.onCandidateHdlr(nil) } } }); err != nil { a.log.Warnf("Failed to run onCandidateHdlr task: %v\n", err)