@cryptotaxi247 / kubo / commits / 0061f0c15

new swarm -- it's so simple

Juan Batiz-Benet committed Dec 14, 2014 at 15:08 UTC 0061f0c151dcd0b195787f55aba37da0b14a1f09
3 files changed +337
net/swarm2/addr.go new
+124
@@ -0,0 +1,124 @@
1 +package swarm
2 +
3 +import (
4 + conn "github.com/jbenet/go-ipfs/net/conn"
5 + eventlog "github.com/jbenet/go-ipfs/util/eventlog"
6 +
7 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8 + ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
9 + manet "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr-net"
10 +)
11 +
12 +// ListenAddresses returns a list of addresses at which this swarm listens.
13 +func (s *Swarm) ListenAddresses() []ma.Multiaddr {
14 + listeners := s.swarm.Listeners()
15 + addrs := make([]ma.Multiaddr, 0, len(listeners))
16 + for _, l := range listeners {
17 + if l2, ok := l.NetListener().(conn.Listener); ok {
18 + addrs = append(addrs, l2.Multiaddr())
19 + }
20 + }
21 + return addrs
22 +}
23 +
24 +// InterfaceListenAddresses returns a list of addresses at which this swarm
25 +// listens. It expands "any interface" addresses (/ip4/0.0.0.0, /ip6/::) to
26 +// use the known local interfaces.
27 +func InterfaceListenAddresses(s *Swarm) ([]ma.Multiaddr, error) {
28 + return resolveUnspecifiedAddresses(s.ListenAddresses())
29 +}
30 +
31 +// resolveUnspecifiedAddresses expands unspecified ip addresses (/ip4/0.0.0.0, /ip6/::) to
32 +// use the known local interfaces.
33 +func resolveUnspecifiedAddresses(unspecifiedAddrs []ma.Multiaddr) ([]ma.Multiaddr, error) {
34 + var outputAddrs []ma.Multiaddr
35 +
36 + // todo optimize: only fetch these if we have a "any" addr.
37 + ifaceAddrs, err := interfaceAddresses()
38 + if err != nil {
39 + return nil, err
40 + }
41 +
42 + for _, a := range unspecifiedAddrs {
43 +
44 + // split address into its components
45 + split := ma.Split(a)
46 +
47 + // if first component (ip) is not unspecified, use it as is.
48 + if !manet.IsIPUnspecified(split[0]) {
49 + outputAddrs = append(outputAddrs, a)
50 + continue
51 + }
52 +
53 + // unspecified? add one address per interface.
54 + for _, ia := range ifaceAddrs {
55 + split[0] = ia
56 + joined := ma.Join(split...)
57 + outputAddrs = append(outputAddrs, joined)
58 + }
59 + }
60 +
61 + log.Event(context.TODO(), "interfaceListenAddresses", func() eventlog.Loggable {
62 + var addrs []string
63 + for _, addr := range outputAddrs {
64 + addrs = append(addrs, addr.String())
65 + }
66 + return eventlog.Metadata{"addresses": addrs}
67 + }())
68 + log.Debug("InterfaceListenAddresses:", outputAddrs)
69 + return outputAddrs, nil
70 +}
71 +
72 +// interfaceAddresses returns a list of addresses associated with local machine
73 +func interfaceAddresses() ([]ma.Multiaddr, error) {
74 + maddrs, err := manet.InterfaceMultiaddrs()
75 + if err != nil {
76 + return nil, err
77 + }
78 +
79 + var nonLoopback []ma.Multiaddr
80 + for _, a := range maddrs {
81 + if !manet.IsIPLoopback(a) {
82 + nonLoopback = append(nonLoopback, a)
83 + }
84 + }
85 +
86 + return nonLoopback, nil
87 +}
88 +
89 +// addrInList returns whether or not an address is part of a list.
90 +// this is useful to check if NAT is happening (or other bugs?)
91 +func addrInList(addr ma.Multiaddr, list []ma.Multiaddr) bool {
92 + for _, addr2 := range list {
93 + if addr.Equal(addr2) {
94 + return true
95 + }
96 + }
97 + return false
98 +}
99 +
100 +// checkNATWarning checks if our observed addresses differ. if so,
101 +// informs the user that certain things might not work yet
102 +func checkNATWarning(s *Swarm, observed ma.Multiaddr, expected ma.Multiaddr) {
103 + if observed.Equal(expected) {
104 + return
105 + }
106 +
107 + listen, err := InterfaceListenAddresses(s)
108 + if err != nil {
109 + log.Errorf("Error retrieving swarm.InterfaceListenAddresses: %s", err)
110 + return
111 + }
112 +
113 + if !addrInList(observed, listen) { // probably a nat
114 + log.Warningf(natWarning, observed, listen)
115 + }
116 +}
117 +
118 +const natWarning = `Remote peer observed our address to be: %s
119 +The local addresses are: %s
120 +Thus, connection is going through NAT, and other connections may fail.
121 +
122 +IPFS NAT traversal is still under development. Please bug us on github or irc to fix this.
123 +Baby steps: http://jbenet.static.s3.amazonaws.com/271dfcf/baby-steps.gif
124 +`
net/swarm2/swarm.go new
+105
@@ -0,0 +1,105 @@
1 +// package swarm implements a connection muxer with a pair of channels
2 +// to synchronize all network communication.
3 +package swarm
4 +
5 +import (
6 + conn "github.com/jbenet/go-ipfs/net/conn"
7 + peer "github.com/jbenet/go-ipfs/peer"
8 + eventlog "github.com/jbenet/go-ipfs/util/eventlog"
9 +
10 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
11 + ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
12 + ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
13 + ps "github.com/jbenet/go-peerstream"
14 +)
15 +
16 +var log = eventlog.Logger("swarm2")
17 +
18 +// Swarm is a connection muxer, allowing connections to other peers to
19 +// be opened and closed, while still using the same Chan for all
20 +// communication. The Chan sends/receives Messages, which note the
21 +// destination or source Peer.
22 +//
23 +// Uses peerstream.Swarm
24 +type Swarm struct {
25 + swarm *ps.Swarm
26 + local peer.Peer
27 + peers peer.Peerstore
28 +
29 + cg ctxgroup.ContextGroup
30 +}
31 +
32 +// NewSwarm constructs a Swarm, with a Chan.
33 +func NewSwarm(ctx context.Context, listenAddrs []ma.Multiaddr,
34 + local peer.Peer, peers peer.Peerstore) (*Swarm, error) {
35 +
36 + s := &Swarm{
37 + swarm: ps.NewSwarm(),
38 + local: local,
39 + peers: peers,
40 + cg: ctxgroup.WithContext(ctx),
41 + }
42 +
43 + // configure Swarm
44 + s.cg.SetTeardown(s.teardown)
45 + s.swarm.SetConnHandler(s.connHandler)
46 +
47 + return s, s.listen(listenAddrs)
48 +}
49 +
50 +func (s *Swarm) teardown() error {
51 + return s.swarm.Close()
52 +}
53 +
54 +func (s *Swarm) Close() error {
55 + return s.cg.Close()
56 +}
57 +
58 +func (s *Swarm) StreamSwarm() *ps.Swarm {
59 + return s.swarm
60 +}
61 +
62 +// Connections returns a slice of all connections.
63 +func (s *Swarm) Connections() []conn.Conn {
64 + conns1 := s.swarm.Conns()
65 + conns2 := make([]conn.Conn, len(conns1))
66 + for i, c1 := range conns1 {
67 + conns2[i] = UnwrapConn(c1)
68 + }
69 + return conns2
70 +}
71 +
72 +// CloseConnection removes a given peer from swarm + closes the connection
73 +func (s *Swarm) CloseConnection(p peer.Peer) error {
74 + conns := s.swarm.ConnsWithGroup(p) // boom.
75 + for _, c := range conns {
76 + c.Close()
77 + }
78 + return nil
79 +}
80 +
81 +// GetPeerList returns a copy of the set of peers swarm is connected to.
82 +func (s *Swarm) GetPeerList() []peer.Peer {
83 + conns := s.swarm.Conns()
84 +
85 + seen := make(map[peer.Peer]struct{})
86 + peers := make([]peer.Peer, 0, len(conns))
87 + for _, c := range conns {
88 + c2 := UnwrapConn(c)
89 + p := c2.RemotePeer()
90 + if _, found := seen[p]; found {
91 + continue
92 + }
93 + peers = append(peers, p)
94 + }
95 + return peers
96 +}
97 +
98 +// LocalPeer returns the local peer swarm is associated to.
99 +func (s *Swarm) LocalPeer() peer.Peer {
100 + return s.local
101 +}
102 +
103 +func UnwrapConn(c *ps.Conn) conn.Conn {
104 + return c.NetConn().(conn.Conn)
105 +}
net/swarm2/swarm_listen.go new
+108
@@ -0,0 +1,108 @@
1 +package swarm
2 +
3 +import (
4 + "fmt"
5 +
6 + conn "github.com/jbenet/go-ipfs/net/conn"
7 +
8 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
9 + ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
10 + multierr "github.com/jbenet/go-ipfs/util/multierr"
11 + ps "github.com/jbenet/go-peerstream"
12 +)
13 +
14 +// Open listeners for each network the swarm should listen on
15 +func (s *Swarm) listen(addrs []ma.Multiaddr) error {
16 + retErr := multierr.New()
17 +
18 + // listen on every address
19 + for i, addr := range addrs {
20 + err := s.setupListener(addr)
21 + if err != nil {
22 + retErr.Errors[i] = err
23 + log.Errorf("Failed to listen on: %s - %s", addr, err)
24 + }
25 + }
26 +
27 + if len(retErr.Errors) > 0 {
28 + return retErr
29 + }
30 + return nil
31 +}
32 +
33 +// Listen for new connections on the given multiaddr
34 +func (s *Swarm) setupListener(maddr ma.Multiaddr) error {
35 +
36 + resolved, err := resolveUnspecifiedAddresses([]ma.Multiaddr{maddr})
37 + if err != nil {
38 + return err
39 + }
40 +
41 + list, err := conn.Listen(s.cg.Context(), maddr, s.local, s.peers)
42 + if err != nil {
43 + return err
44 + }
45 +
46 + // add resolved local addresses to peer
47 + for _, addr := range resolved {
48 + s.local.AddAddress(addr)
49 + }
50 +
51 + // AddListener to the peerstream Listener. this will begin accepting connections
52 + // and streams!
53 + _, err = s.swarm.AddListener(list)
54 + return err
55 +}
56 +
57 +// connHandler is called by the StreamSwarm whenever a new connection is added
58 +// here we configure it slightly. Note that this is sequential, so if anything
59 +// will take a while do it in a goroutine.
60 +// See https://godoc.org/github.com/jbenet/go-peerstream for more information
61 +func (s *Swarm) connHandler(c1 *ps.Conn) {
62 +
63 + // grab the underlying connection.
64 + if c2, ok := c1.NetConn().(conn.Conn); ok {
65 +
66 + // set the RemotePeer as a group on the conn. this lets us group
67 + // connections in the StreamSwarm by peer, and get a streams from
68 + // any available connection in the group (better multiconn):
69 + // swarm.StreamSwarm().NewStreamWithGroup(remotePeer)
70 + c1.AddGroup(c2.RemotePeer())
71 +
72 + go func() {
73 + ctx := context.Background()
74 + err := runHandshake3(ctx, s, c1, c2)
75 + if err != nil {
76 + log.Error("Handshake3 failed. disconnecting", err)
77 + log.Event(ctx, "Handshake3FailureDisconnect", c2.LocalPeer(), c2.RemotePeer())
78 + c1.Close() // boom.
79 + }
80 + }()
81 + }
82 +}
83 +
84 +func runHandshake3(ctx context.Context, s *Swarm, sc *ps.Conn, c conn.Conn) error {
85 + log.Event(ctx, "newConnection", c.LocalPeer(), c.RemotePeer())
86 +
87 + stream, err := sc.NewStream()
88 + if err != nil {
89 + return err
90 + }
91 +
92 + // handshake3
93 + h3result, err := conn.Handshake3(ctx, stream, c)
94 + if err != nil {
95 + return fmt.Errorf("Handshake3 failed: %s", err)
96 + }
97 +
98 + // check for nats. you know, just in case.
99 + if h3result.LocalObservedAddress != nil {
100 + checkNATWarning(s, h3result.LocalObservedAddress, c.LocalMultiaddr())
101 + } else {
102 + log.Warningf("Received nil observed address from %s", c.RemotePeer())
103 + }
104 +
105 + stream.Close()
106 + log.Event(ctx, "handshake3Succeeded", c.LocalPeer(), c.RemotePeer())
107 + return nil
108 +}