p2p: tag connections in connection manager
License: MIT Signed-off-by: Łukasz Magiera <magik6k@gmail.com>
Łukasz Magiera committed
Jul 30, 2018 at 17:15 UTC
4badcdc340f809046966fad75707bbb2182151c1
5 files changed
+28
-8
core/commands/p2p.go
-1
@@ -14,7 +14,6 @@ import (
14
core "github.com/ipfs/go-ipfs/core"
15
p2p "github.com/ipfs/go-ipfs/p2p"
16
17
- "gx/ipfs/Qme4QgoVPyQqxVc4G1c2L2wc9TDa6o294rtspGMnBNRujm/go-ipfs-addr"
17
ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
18
"gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
19
pstore "gx/ipfs/QmZR2XWVVBCtbgBWnQhWk2xcQfaR3W8faQPriAiaaj7rsr/go-libp2p-peerstore"
p2p/local.go
+10
-3
@@ -4,12 +4,12 @@ import (
4
"context"
5
"time"
6
7
+ "gx/ipfs/QmPjvxTpVH8qJyQDnxnsxF9kv9jezKD1kozz1hs3fCGsNh/go-libp2p-net"
8
"gx/ipfs/QmV6FjemM1K8oXjrvuq3wuVWWoU2TLDPmNnKrxHzY3v6Ai/go-multiaddr-net"
8
- ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
9
- "gx/ipfs/QmdVrMn1LhB4ybb8hMVaMLXnA8XRSewMnK6YqXKXoTcRvN/go-libp2p-peer"
9
tec "gx/ipfs/QmWHgLqrghM9zw77nF6gdvT9ExQ2RB9pLxkd8sDHZf1rWb/go-temp-err-catcher"
11
- "gx/ipfs/QmPjvxTpVH8qJyQDnxnsxF9kv9jezKD1kozz1hs3fCGsNh/go-libp2p-net"
10
+ ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
11
"gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
12
+ "gx/ipfs/QmdVrMn1LhB4ybb8hMVaMLXnA8XRSewMnK6YqXKXoTcRvN/go-libp2p-peer"
13
)
14
15
// localListener manet streams and proxies them to libp2p services
@@ -77,6 +77,9 @@ func (l *localListener) setupStream(local manet.Conn) {
77
return
78
}
79
80
+ cmgr := l.p2p.peerHost.ConnManager()
81
+ cmgr.TagPeer(l.peer, CMGR_TAG, 20)
82
+
83
stream := &Stream{
84
Protocol: l.proto,
85
@@ -87,6 +90,10 @@ func (l *localListener) setupStream(local manet.Conn) {
90
Remote: remote,
91
92
Registry: l.p2p.Streams,
93
+
94
+ cleanup: func() {
95
+ cmgr.UntagPeer(l.peer, CMGR_TAG)
96
+ },
97
}
98
99
l.p2p.Streams.Register(stream)
p2p/p2p.go
+1
-1
@@ -1,9 +1,9 @@
1
package p2p
2
3
import (
4
- logging "gx/ipfs/QmcVVHfdyv15GVPk7NrxdWjh2hLVccXnoD8j2tyQShiXJb/go-log"
4
pstore "gx/ipfs/QmZR2XWVVBCtbgBWnQhWk2xcQfaR3W8faQPriAiaaj7rsr/go-libp2p-peerstore"
5
p2phost "gx/ipfs/Qmb8T6YBBsjYsVGfrihQLfCJveczZnneSBqBKkYEBWDjge/go-libp2p-host"
6
+ logging "gx/ipfs/QmcVVHfdyv15GVPk7NrxdWjh2hLVccXnoD8j2tyQShiXJb/go-log"
7
peer "gx/ipfs/QmdVrMn1LhB4ybb8hMVaMLXnA8XRSewMnK6YqXKXoTcRvN/go-libp2p-peer"
8
)
9
p2p/remote.go
+11
-2
@@ -3,9 +3,9 @@ package p2p
3
import (
4
"context"
5
6
+ net "gx/ipfs/QmPjvxTpVH8qJyQDnxnsxF9kv9jezKD1kozz1hs3fCGsNh/go-libp2p-net"
7
manet "gx/ipfs/QmV6FjemM1K8oXjrvuq3wuVWWoU2TLDPmNnKrxHzY3v6Ai/go-multiaddr-net"
8
ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
8
- net "gx/ipfs/QmPjvxTpVH8qJyQDnxnsxF9kv9jezKD1kozz1hs3fCGsNh/go-libp2p-net"
9
protocol "gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
10
)
11
@@ -47,12 +47,17 @@ func (l *remoteListener) start() error {
47
return
48
}
49
50
- peerMa, err := ma.NewMultiaddr(maPrefix + remote.Conn().RemotePeer().Pretty())
50
+ peer := remote.Conn().RemotePeer()
51
+
52
+ peerMa, err := ma.NewMultiaddr(maPrefix + peer.Pretty())
53
if err != nil {
54
remote.Reset()
55
return
56
}
57
58
+ cmgr := l.p2p.peerHost.ConnManager()
59
+ cmgr.TagPeer(peer, CMGR_TAG, 20)
60
+
61
stream := &Stream{
62
Protocol: l.proto,
63
@@ -63,6 +68,10 @@ func (l *remoteListener) start() error {
68
Remote: remote,
69
70
Registry: l.p2p.Streams,
71
+
72
+ cleanup: func() {
73
+ cmgr.UntagPeer(peer, CMGR_TAG)
74
+ },
75
}
76
77
l.p2p.Streams.Register(stream)
p2p/stream.go
+6
-1
@@ -4,12 +4,14 @@ import (
4
"io"
5
"sync"
6
7
+ net "gx/ipfs/QmPjvxTpVH8qJyQDnxnsxF9kv9jezKD1kozz1hs3fCGsNh/go-libp2p-net"
8
manet "gx/ipfs/QmV6FjemM1K8oXjrvuq3wuVWWoU2TLDPmNnKrxHzY3v6Ai/go-multiaddr-net"
9
ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
9
- net "gx/ipfs/QmPjvxTpVH8qJyQDnxnsxF9kv9jezKD1kozz1hs3fCGsNh/go-libp2p-net"
10
"gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
11
)
12
13
+const CMGR_TAG = "stream-fwd"
14
+
15
// Stream holds information on active incoming and outgoing p2p streams.
16
type Stream struct {
17
id uint64
@@ -23,12 +25,15 @@ type Stream struct {
25
Remote net.Stream
26
27
Registry *StreamRegistry
28
+
29
+ cleanup func()
30
}
31
32
// Close closes stream endpoints and deregisters it
33
func (s *Stream) Close() error {
34
s.Local.Close()
35
s.Remote.Close()
36
+ s.cleanup()
37
s.Registry.Deregister(s.id)
38
return nil
39
}