p2p: fix remote/local listener races
License: MIT Signed-off-by: Łukasz Magiera <magik6k@gmail.com>
Łukasz Magiera committed
Jun 25, 2018 at 17:25 UTC
4c98edaff6b392962ff863b79d92a56df57a7ec3
4 files changed
+28
-22
p2p/listener.go
+19
-8
@@ -31,26 +31,33 @@ type ListenerRegistry struct {
31
sync.Mutex
32
33
Listeners map[listenerKey]Listener
34
+ starting map[listenerKey]struct{}
35
}
36
37
// Register registers listenerInfo into this registry and starts it
38
func (r *ListenerRegistry) Register(l Listener) error {
39
r.Lock()
40
+ k := getListenerKey(l)
41
40
- if _, ok := r.Listeners[getListenerKey(l)]; ok {
42
+ if _, ok := r.Listeners[k]; ok {
43
r.Unlock()
44
return errors.New("listener already registered")
45
}
46
45
- r.Listeners[getListenerKey(l)] = l
47
+ r.Listeners[k] = l
48
+ r.starting[k] = struct{}{}
49
50
r.Unlock()
51
49
- if err := l.start(); err != nil {
50
- r.Lock()
51
- defer r.Lock()
52
+ err := l.start()
53
53
- delete(r.Listeners, getListenerKey(l))
54
+ r.Lock()
55
+ defer r.Unlock()
56
+
57
+ delete(r.starting, k)
58
+
59
+ if err != nil {
60
+ delete(r.Listeners, k)
61
return err
62
}
63
@@ -58,13 +65,17 @@ func (r *ListenerRegistry) Register(l Listener) error {
65
}
66
67
// Deregister removes p2p listener from this registry
61
-func (r *ListenerRegistry) Deregister(k listenerKey) bool {
68
+func (r *ListenerRegistry) Deregister(k listenerKey) (bool, error) {
69
r.Lock()
70
defer r.Unlock()
71
72
+ if _, ok := r.starting[k]; ok {
73
+ return false, errors.New("listener didn't start yet")
74
+ }
75
+
76
_, ok := r.Listeners[k]
77
delete(r.Listeners, k)
67
- return ok
78
+ return ok, nil
79
}
80
81
func getListenerKey(l Listener) listenerKey {
p2p/local.go
+4
-5
@@ -2,7 +2,6 @@ package p2p
2
3
import (
4
"context"
5
- "errors"
5
"time"
6
7
"gx/ipfs/QmV6FjemM1K8oXjrvuq3wuVWWoU2TLDPmNnKrxHzY3v6Ai/go-multiaddr-net"
@@ -105,11 +104,11 @@ func (l *localListener) start() error {
104
}
105
106
func (l *localListener) Close() error {
108
- if l.listener == nil {
109
- return errors.New("uninitialized")
107
+ ok, err := l.p2p.Listeners.Deregister(getListenerKey(l))
108
+ if err != nil {
109
+ return err
110
}
111
-
112
- if l.p2p.Listeners.Deregister(getListenerKey(l)) {
111
+ if ok {
112
l.listener.Close()
113
l.listener = nil
114
}
p2p/p2p.go
+1
@@ -28,6 +28,7 @@ func NewP2P(identity peer.ID, peerHost p2phost.Host, peerstore pstore.Peerstore)
28
29
Listeners: &ListenerRegistry{
30
Listeners: map[listenerKey]Listener{},
31
+ starting: map[listenerKey]struct{}{},
32
},
33
Streams: &StreamRegistry{
34
Streams: map[uint64]*Stream{},
p2p/remote.go
+4
-9
@@ -2,7 +2,6 @@ package p2p
2
3
import (
4
"context"
5
- "errors"
5
6
manet "gx/ipfs/QmV6FjemM1K8oXjrvuq3wuVWWoU2TLDPmNnKrxHzY3v6Ai/go-multiaddr-net"
7
ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
@@ -21,8 +20,6 @@ type remoteListener struct {
20
21
// Address to proxy the incoming connections to
22
addr ma.Multiaddr
24
-
25
- initialized bool
23
}
24
25
// ForwardRemote creates new p2p listener
@@ -72,7 +69,6 @@ func (l *remoteListener) start() error {
69
stream.startStreaming()
70
})
71
75
- l.initialized = true
72
return nil
73
}
74
@@ -93,13 +89,12 @@ func (l *remoteListener) TargetAddress() ma.Multiaddr {
89
}
90
91
func (l *remoteListener) Close() error {
96
- if !l.initialized {
97
- return errors.New("uninitialized")
92
+ ok, err := l.p2p.Listeners.Deregister(getListenerKey(l))
93
+ if err != nil {
94
+ return err
95
}
99
-
100
- if l.p2p.Listeners.Deregister(getListenerKey(l)) {
96
+ if ok {
97
l.p2p.peerHost.RemoveStreamHandler(l.proto)
102
- l.initialized = false
98
}
99
return nil
100
}