@cryptotaxi247 / kubo / commits / b53936a70

p2p: Close on Listeners

License: MIT Signed-off-by: Łukasz Magiera <magik6k@gmail.com>

Łukasz Magiera committed Sep 12, 2018 at 23:27 UTC b53936a706905d72711d935180df00df4fd5531c
5 files changed +30 -55
core/commands/p2p.go
+3 -29
@@ -340,36 +340,10 @@ var p2pCloseCmd = &cmds.Command{
340 return true
341 }
342
343 - todo := make([]p2p.Listener, 0)
344 - n.P2P.ListenersLocal.Lock()
345 - for _, l := range n.P2P.ListenersLocal.Listeners {
346 - if !match(l) {
347 - continue
348 - }
349 - todo = append(todo, l)
350 - }
351 - n.P2P.ListenersLocal.Unlock()
352 - n.P2P.ListenersP2P.Lock()
353 - for _, l := range n.P2P.ListenersP2P.Listeners {
354 - if !match(l) {
355 - continue
356 - }
357 - todo = append(todo, l)
358 - }
359 - n.P2P.ListenersP2P.Unlock()
360 -
361 - var errs []string
362 - for _, l := range todo {
363 - if err := l.Close(); err != nil {
364 - errs = append(errs, err.Error())
365 - }
366 - }
367 - if len(errs) != 0 {
368 - res.SetError(fmt.Errorf("errors when closing streams: %s", strings.Join(errs, "; ")), cmdkit.ErrNormal)
369 - return
370 - }
343 + done := n.P2P.ListenersLocal.Close(match)
344 + done += n.P2P.ListenersP2P.Close(match)
345
372 - res.SetOutput(len(todo))
346 + res.SetOutput(done)
347 },
348 Type: int(0),
349 Marshalers: cmds.MarshalerMap{
p2p/listener.go
+20 -8
@@ -18,8 +18,8 @@ type Listener interface {
18
19 key() string
20
21 - // Close closes the listener. Does not affect child streams
22 - Close() error
21 + // close closes the listener. Does not affect child streams
22 + close()
23 }
24
25 // Listeners manages a group of Listener implementations,
@@ -73,12 +73,24 @@ func (r *Listeners) Register(l Listener) error {
73 return nil
74 }
75
76 -// Deregister removes p2p listener from this registry
77 -func (r *Listeners) Deregister(k string) (bool, error) {
76 +func (r *Listeners) Close(matchFunc func(listener Listener) bool) int {
77 + todo := make([]Listener, 0)
78 r.Lock()
79 - defer r.Unlock()
79 + for _, l := range r.Listeners {
80 + if !matchFunc(l) {
81 + continue
82 + }
83 +
84 + if _, ok := r.Listeners[l.key()]; ok {
85 + delete(r.Listeners, l.key())
86 + todo = append(todo, l)
87 + }
88 + }
89 + r.Unlock()
90 +
91 + for _, l := range todo {
92 + l.close()
93 + }
94
81 - _, ok := r.Listeners[k]
82 - delete(r.Listeners, k)
83 - return ok, nil
95 + return len(todo)
96 }
p2p/local.go
+2 -9
@@ -98,15 +98,8 @@ func (l *localListener) setupStream(local manet.Conn) {
98 l.p2p.Streams.Register(stream)
99 }
100
101 -func (l *localListener) Close() error {
102 - ok, err := l.p2p.ListenersLocal.Deregister(l.laddr.String())
103 - if err != nil {
104 - return err
105 - }
106 - if ok {
107 - return l.listener.Close()
108 - }
109 - return nil
101 +func (l *localListener) close() {
102 + l.listener.Close()
103 }
104
105 func (l *localListener) Protocol() protocol.ID {
p2p/remote.go
+1 -7
@@ -85,13 +85,7 @@ func (l *remoteListener) TargetAddress() ma.Multiaddr {
85 return l.addr
86 }
87
88 -func (l *remoteListener) Close() error {
89 - _, err := l.p2p.ListenersP2P.Deregister(string(l.proto))
90 - if err != nil {
91 - return err
92 - }
93 - return nil
94 -}
88 +func (l *remoteListener) close() {}
89
90 func (l *remoteListener) key() string {
91 return string(l.proto)
test/sharness/t0180-p2p.sh
+4 -2
@@ -285,11 +285,13 @@ test_expect_success 'S->C Setup client side (custom proto)' '
285 test_server_to_client
286
287 test_expect_success 'C->S Close local listener' '
288 - ipfsi 0 p2p close -p /p2p-test
289 - ipfsi 0 p2p ls > actual &&
288 + ipfsi 1 p2p close -p /p2p-test
289 + ipfsi 1 p2p ls > actual &&
290 test_must_be_empty actual
291 '
292
293 +check_test_ports
294 +
295 test_expect_success 'stop iptb' '
296 iptb stop
297 '