net interface
Juan Batiz-Benet committed
Sep 14, 2014 at 01:24 UTC
2b03664ae4cc95c81222014580fd94dce7f42b33
5 files changed
+153
-33
net/interface.go
new
+33
@@ -0,0 +1,33 @@
1
+package net
2
+
3
+import (
4
+ msg "github.com/jbenet/go-ipfs/net/message"
5
+ mux "github.com/jbenet/go-ipfs/net/mux"
6
+ peer "github.com/jbenet/go-ipfs/peer"
7
+)
8
+
9
+// Network is the interface IPFS uses for connecting to the world.
10
+type Network interface {
11
+
12
+ // Listen handles incoming connections on given Multiaddr.
13
+ // Listen(*ma.Muliaddr) error
14
+ // TODO: for now, only listen on addrs in local peer when initializing.
15
+
16
+ // DialPeer attempts to establish a connection to a given peer
17
+ DialPeer(*peer.Peer) error
18
+
19
+ // ClosePeer connection to peer
20
+ ClosePeer(*peer.Peer) error
21
+
22
+ // IsConnected returns whether a connection to given peer exists.
23
+ IsConnected(*peer.Peer) (bool, error)
24
+
25
+ // GetProtocols returns the protocols registered in the network.
26
+ GetProtocols() *mux.ProtocolMap
27
+
28
+ // SendMessage sends given Message out
29
+ SendMessage(*msg.Message) error
30
+
31
+ // Close terminates all network operation
32
+ Close() error
33
+}
net/net.go
new
+106
@@ -0,0 +1,106 @@
1
+package net
2
+
3
+import (
4
+ "errors"
5
+
6
+ msg "github.com/jbenet/go-ipfs/net/message"
7
+ mux "github.com/jbenet/go-ipfs/net/mux"
8
+ swarm "github.com/jbenet/go-ipfs/net/swarm"
9
+ peer "github.com/jbenet/go-ipfs/peer"
10
+
11
+ context "code.google.com/p/go.net/context"
12
+)
13
+
14
+// IpfsNetwork implements the Network interface,
15
+type IpfsNetwork struct {
16
+
17
+ // local peer
18
+ local *peer.Peer
19
+
20
+ // protocol multiplexing
21
+ muxer *mux.Muxer
22
+
23
+ // peer connection multiplexing
24
+ swarm *swarm.Swarm
25
+
26
+ // network context
27
+ ctx context.Context
28
+ cancel context.CancelFunc
29
+}
30
+
31
+// NewIpfsNetwork is the structure that implements the network interface
32
+func NewIpfsNetwork(ctx context.Context, local *peer.Peer,
33
+ pmap *mux.ProtocolMap) (*IpfsNetwork, error) {
34
+
35
+ ctx, cancel := context.WithCancel(ctx)
36
+
37
+ in := &IpfsNetwork{
38
+ local: local,
39
+ muxer: &mux.Muxer{Protocols: *pmap},
40
+ ctx: ctx,
41
+ cancel: cancel,
42
+ }
43
+
44
+ err := in.muxer.Start(ctx)
45
+ if err != nil {
46
+ cancel()
47
+ return nil, err
48
+ }
49
+
50
+ in.swarm, err = swarm.NewSwarm(ctx, local)
51
+ if err != nil {
52
+ cancel()
53
+ return nil, err
54
+ }
55
+
56
+ return in, nil
57
+}
58
+
59
+// Listen handles incoming connections on given Multiaddr.
60
+// func (n *IpfsNetwork) Listen(*ma.Muliaddr) error {}
61
+
62
+// DialPeer attempts to establish a connection to a given peer
63
+func (n *IpfsNetwork) DialPeer(p *peer.Peer) error {
64
+ _, err := n.swarm.Dial(p)
65
+ return err
66
+}
67
+
68
+// ClosePeer connection to peer
69
+func (n *IpfsNetwork) ClosePeer(p *peer.Peer) error {
70
+ return n.swarm.CloseConnection(p)
71
+}
72
+
73
+// IsConnected returns whether a connection to given peer exists.
74
+func (n *IpfsNetwork) IsConnected(p *peer.Peer) (bool, error) {
75
+ return n.swarm.GetConnection(p.ID) != nil, nil
76
+}
77
+
78
+// GetProtocols returns the protocols registered in the network.
79
+func (n *IpfsNetwork) GetProtocols() *mux.ProtocolMap {
80
+ // copy over because this map should be read only.
81
+ pmap := mux.ProtocolMap{}
82
+ for id, proto := range n.muxer.Protocols {
83
+ pmap[id] = proto
84
+ }
85
+ return &pmap
86
+}
87
+
88
+// SendMessage sends given Message out
89
+func (n *IpfsNetwork) SendMessage(m *msg.Message) error {
90
+ n.swarm.Outgoing <- m
91
+ return nil
92
+}
93
+
94
+// Close terminates all network operation
95
+func (n *IpfsNetwork) Close() error {
96
+ if n.cancel == nil {
97
+ return errors.New("Network already closed.")
98
+ }
99
+
100
+ n.swarm.Close()
101
+ n.muxer.Stop()
102
+
103
+ n.cancel()
104
+ n.cancel = nil
105
+ return nil
106
+}
net/net_test.go
new
+1
@@ -0,0 +1 @@
1
+package net
net/service/service.go
+5
@@ -68,6 +68,11 @@ func (s *Service) Stop() {
68
s.cancel = context.CancelFunc(nil)
69
}
70
71
+// GetPipe implements the mux.Protocol interface
72
+func (s *Service) GetPipe() *msg.Pipe {
73
+ return s.Pipe
74
+}
75
+
76
// SendMessage sends a message out
77
func (s *Service) SendMessage(ctx context.Context, m *msg.Message, rid RequestID) error {
78
net/swarm/swarm.go
+8
-33
@@ -110,15 +110,8 @@ func (s *Swarm) Dial(peer *peer.Peer) (*conn.Conn, error) {
110
return nil, errors.New("Attempted connection to self!")
111
}
112
113
- k := peer.Key()
114
-
113
// check if we already have an open connection first
116
- s.connsLock.RLock()
117
- c, found := s.conns[k]
118
- s.connsLock.RUnlock()
119
- if found {
120
- return c, nil
121
- }
114
+ c := s.GetConnection(peer.ID)
115
116
// open connection to peer
117
c, err := conn.Dial("tcp", peer)
@@ -158,40 +151,22 @@ func (s *Swarm) DialAddr(addr *ma.Multiaddr) (*conn.Conn, error) {
151
return c, err
152
}
153
161
-// GetPeer returns the peer in the swarm with given key id.
162
-func (s *Swarm) GetPeer(key u.Key) *peer.Peer {
154
+// GetConnection returns the connection in the swarm to given peer.ID
155
+func (s *Swarm) GetConnection(pid peer.ID) *conn.Conn {
156
s.connsLock.RLock()
164
- conn, found := s.conns[key]
157
+ c, found := s.conns[u.Key(pid)]
158
s.connsLock.RUnlock()
159
160
if !found {
161
return nil
162
}
170
- return conn.Peer
171
-}
172
-
173
-// GetConnection will check if we are already connected to the peer in question
174
-// and only open a new connection if we arent already
175
-func (s *Swarm) GetConnection(id peer.ID, addr *ma.Multiaddr) (*peer.Peer, error) {
176
- p := &peer.Peer{
177
- ID: id,
178
- Addresses: []*ma.Multiaddr{addr},
179
- }
180
-
181
- c, err := s.Dial(p)
182
- if err != nil {
183
- return nil, err
184
- }
185
-
186
- return c.Peer, nil
163
+ return c
164
}
165
166
// CloseConnection removes a given peer from swarm + closes the connection
167
func (s *Swarm) CloseConnection(p *peer.Peer) error {
191
- s.connsLock.RLock()
192
- conn, found := s.conns[u.Key(p.ID)]
193
- s.connsLock.RUnlock()
194
- if !found {
168
+ c := s.GetConnection(p.ID)
169
+ if c == nil {
170
return u.ErrNotFound
171
}
172
@@ -199,7 +174,7 @@ func (s *Swarm) CloseConnection(p *peer.Peer) error {
174
delete(s.conns, u.Key(p.ID))
175
s.connsLock.Unlock()
176
202
- return conn.Close()
177
+ return c.Close()
178
}
179
180
func (s *Swarm) Error(e error) {