@cryptotaxi247 / kubo / commits / 453a66709

moved stuff

Juan Batiz-Benet committed Sep 13, 2014 at 18:36 UTC 453a66709f027e0a1cccb582c6dcb71a9da6a2db
2 files changed +123 -121
net/swarm/conn.go
+123 -1
@@ -7,10 +7,72 @@ import (
7
8 ident "github.com/jbenet/go-ipfs/identify"
9 conn "github.com/jbenet/go-ipfs/net/conn"
10 -
10 + msg "github.com/jbenet/go-ipfs/net/message"
11 u "github.com/jbenet/go-ipfs/util"
12 +
13 + ma "github.com/jbenet/go-multiaddr"
14 )
15
16 +// Open listeners for each network the swarm should listen on
17 +func (s *Swarm) listen() error {
18 + hasErr := false
19 + retErr := &ListenErr{
20 + Errors: make([]error, len(s.local.Addresses)),
21 + }
22 +
23 + // listen on every address
24 + for i, addr := range s.local.Addresses {
25 + err := s.connListen(addr)
26 + if err != nil {
27 + hasErr = true
28 + retErr.Errors[i] = err
29 + u.PErr("Failed to listen on: %s [%s]", addr, err)
30 + }
31 + }
32 +
33 + if hasErr {
34 + return retErr
35 + }
36 + return nil
37 +}
38 +
39 +// Listen for new connections on the given multiaddr
40 +func (s *Swarm) connListen(maddr *ma.Multiaddr) error {
41 + netstr, addr, err := maddr.DialArgs()
42 + if err != nil {
43 + return err
44 + }
45 +
46 + list, err := net.Listen(netstr, addr)
47 + if err != nil {
48 + return err
49 + }
50 +
51 + // NOTE: this may require a lock around it later. currently, only run on setup
52 + s.listeners = append(s.listeners, list)
53 +
54 + // Accept and handle new connections on this listener until it errors
55 + go func() {
56 + for {
57 + nconn, err := list.Accept()
58 + if err != nil {
59 + e := fmt.Errorf("Failed to accept connection: %s - %s [%s]",
60 + netstr, addr, err)
61 + s.errChan <- e
62 +
63 + // if cancel is nil, we're closed.
64 + if s.cancel == nil {
65 + return
66 + }
67 + } else {
68 + go s.handleIncomingConn(nconn)
69 + }
70 + }
71 + }()
72 +
73 + return nil
74 +}
75 +
76 // Handle getting ID from this peer, handshake, and adding it into the map
77 func (s *Swarm) handleIncomingConn(nconn net.Conn) {
78
@@ -67,3 +129,63 @@ func (s *Swarm) connHandshake(c *conn.Conn) error {
129 // needs cleanup. needs context. use msg.Pipe.
130 return ident.Handshake(s.local, c.Peer, c.Incoming.MsgChan, c.Outgoing.MsgChan)
131 }
132 +
133 +// Handles the unwrapping + sending of messages to the right connection.
134 +func (s *Swarm) fanOut() {
135 + for {
136 + select {
137 + case <-s.ctx.Done():
138 + return // told to close.
139 +
140 + case msg, ok := <-s.Outgoing:
141 + if !ok {
142 + return
143 + }
144 +
145 + s.connsLock.RLock()
146 + conn, found := s.conns[msg.Peer.Key()]
147 + s.connsLock.RUnlock()
148 +
149 + if !found {
150 + e := fmt.Errorf("Sent msg to peer without open conn: %v",
151 + msg.Peer)
152 + s.errChan <- e
153 + continue
154 + }
155 +
156 + // queue it in the connection's buffer
157 + conn.Outgoing.MsgChan <- msg.Data
158 + }
159 + }
160 +}
161 +
162 +// Handles the receiving + wrapping of messages, per conn.
163 +// Consider using reflect.Select with one goroutine instead of n.
164 +func (s *Swarm) fanIn(c *conn.Conn) {
165 + for {
166 + select {
167 + case <-s.ctx.Done():
168 + // close Conn.
169 + c.Close()
170 + goto out
171 +
172 + case <-c.Closed:
173 + goto out
174 +
175 + case data, ok := <-c.Incoming.MsgChan:
176 + if !ok {
177 + e := fmt.Errorf("Error retrieving from conn: %v", c.Peer.Key().Pretty())
178 + s.errChan <- e
179 + goto out
180 + }
181 +
182 + msg := &msg.Message{Peer: c.Peer, Data: data}
183 + s.Incoming <- msg
184 + }
185 + }
186 +
187 +out:
188 + s.connsLock.Lock()
189 + delete(s.conns, c.Peer.Key())
190 + s.connsLock.Unlock()
191 +}
net/swarm/swarm.go
-120
@@ -78,66 +78,6 @@ func NewSwarm(ctx context.Context, local *peer.Peer) (*Swarm, error) {
78 return s, s.listen()
79 }
80
81 -// Open listeners for each network the swarm should listen on
82 -func (s *Swarm) listen() error {
83 - hasErr := false
84 - retErr := &ListenErr{
85 - Errors: make([]error, len(s.local.Addresses)),
86 - }
87 -
88 - // listen on every address
89 - for i, addr := range s.local.Addresses {
90 - err := s.connListen(addr)
91 - if err != nil {
92 - hasErr = true
93 - retErr.Errors[i] = err
94 - u.PErr("Failed to listen on: %s [%s]", addr, err)
95 - }
96 - }
97 -
98 - if hasErr {
99 - return retErr
100 - }
101 - return nil
102 -}
103 -
104 -// Listen for new connections on the given multiaddr
105 -func (s *Swarm) connListen(maddr *ma.Multiaddr) error {
106 - netstr, addr, err := maddr.DialArgs()
107 - if err != nil {
108 - return err
109 - }
110 -
111 - list, err := net.Listen(netstr, addr)
112 - if err != nil {
113 - return err
114 - }
115 -
116 - // NOTE: this may require a lock around it later. currently, only run on setup
117 - s.listeners = append(s.listeners, list)
118 -
119 - // Accept and handle new connections on this listener until it errors
120 - go func() {
121 - for {
122 - nconn, err := list.Accept()
123 - if err != nil {
124 - e := fmt.Errorf("Failed to accept connection: %s - %s [%s]",
125 - netstr, addr, err)
126 - s.errChan <- e
127 -
128 - // if cancel is nil, we're closed.
129 - if s.cancel == nil {
130 - return
131 - }
132 - } else {
133 - go s.handleIncomingConn(nconn)
134 - }
135 - }
136 - }()
137 -
138 - return nil
139 -}
140 -
81 // Close stops a swarm.
82 func (s *Swarm) Close() error {
83 if s.cancel == nil {
@@ -218,66 +158,6 @@ func (s *Swarm) DialAddr(addr *ma.Multiaddr) (*conn.Conn, error) {
158 return c, err
159 }
160
221 -// Handles the unwrapping + sending of messages to the right connection.
222 -func (s *Swarm) fanOut() {
223 - for {
224 - select {
225 - case <-s.ctx.Done():
226 - return // told to close.
227 -
228 - case msg, ok := <-s.Outgoing:
229 - if !ok {
230 - return
231 - }
232 -
233 - s.connsLock.RLock()
234 - conn, found := s.conns[msg.Peer.Key()]
235 - s.connsLock.RUnlock()
236 -
237 - if !found {
238 - e := fmt.Errorf("Sent msg to peer without open conn: %v",
239 - msg.Peer)
240 - s.errChan <- e
241 - continue
242 - }
243 -
244 - // queue it in the connection's buffer
245 - conn.Outgoing.MsgChan <- msg.Data
246 - }
247 - }
248 -}
249 -
250 -// Handles the receiving + wrapping of messages, per conn.
251 -// Consider using reflect.Select with one goroutine instead of n.
252 -func (s *Swarm) fanIn(c *conn.Conn) {
253 - for {
254 - select {
255 - case <-s.ctx.Done():
256 - // close Conn.
257 - c.Close()
258 - goto out
259 -
260 - case <-c.Closed:
261 - goto out
262 -
263 - case data, ok := <-c.Incoming.MsgChan:
264 - if !ok {
265 - e := fmt.Errorf("Error retrieving from conn: %v", c.Peer.Key().Pretty())
266 - s.errChan <- e
267 - goto out
268 - }
269 -
270 - msg := &msg.Message{Peer: c.Peer, Data: data}
271 - s.Incoming <- msg
272 - }
273 - }
274 -
275 -out:
276 - s.connsLock.Lock()
277 - delete(s.conns, c.Peer.Key())
278 - s.connsLock.Unlock()
279 -}
280 -
161 // GetPeer returns the peer in the swarm with given key id.
162 func (s *Swarm) GetPeer(key u.Key) *peer.Peer {
163 s.connsLock.RLock()