@cryptotaxi247 / kubo / commits / b84a71de8

transport refactor update

License: MIT Signed-off-by: Steven Allen <steven@stebalien.com>

Steven Allen committed Mar 7, 2018 at 22:06 UTC b84a71de8cf58643faa0ffdaafa8ba4e0da2e2e3
11 files changed +125 -152
cmd/ipfs/daemon.go
+3 -4
@@ -25,7 +25,6 @@ import (
25 cmds "gx/ipfs/QmSKYWC84fqkKB54Te5JMcov2MBVzucXaRGxFqByzzCbHe/go-ipfs-cmds"
26 ma "gx/ipfs/QmWWQ2Txc2c6tqjsBpzg5Ar652cHPGNsQQp2SejkNmkUMb/go-multiaddr"
27 "gx/ipfs/QmX3QZ5jHEPidwUrymXV1iSCSUhdGxj15sm2gP4jKMef7B/client_golang/prometheus"
28 - iconn "gx/ipfs/QmYDNqBAMWVMHKndYR35Sd8PfEVWBiDmpHYkuRJTunJDeJ/go-libp2p-interface-conn"
28 mprome "gx/ipfs/Qma63DWYgaK1snYcNEv1dBfrZGc961V6frGQiVBGc4TU6h/go-metrics-prometheus"
29 "gx/ipfs/QmceUdzxkimdYsgtX733uNgzf1DLHyBKN6ehGSp85ayppM/go-ipfs-cmdkit"
30 )
@@ -215,7 +214,6 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
214 if unencrypted {
215 log.Warningf(`Running with --%s: All connections are UNENCRYPTED.
216 You will not be able to connect to regular encrypted networks.`, unencryptTransportKwd)
218 - iconn.EncryptConnections = false
217 }
218
219 // first, whether user has provided the initialization flag. we may be
@@ -292,6 +290,7 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
290 Repo: repo,
291 Permanent: true, // It is temporary way to signify that node is permanent
292 Online: !offline,
293 + DisableEncryptedConnections: unencrypted,
294 ExtraOpts: map[string]bool{
295 "pubsub": pubsub,
296 "ipnsps": ipnsps,
@@ -474,7 +473,7 @@ func serveHTTPApi(req *cmds.Request, cctx *oldcmds.Context) (<-chan error, error
473
474 errc := make(chan error)
475 go func() {
477 - errc <- corehttp.Serve(node, apiLis.NetListener(), opts...)
476 + errc <- corehttp.Serve(node, manet.NetListener(apiLis), opts...)
477 close(errc)
478 }()
479 return errc, nil
@@ -561,7 +560,7 @@ func serveHTTPGateway(req *cmds.Request, cctx *oldcmds.Context) (<-chan error, e
560
561 errc := make(chan error)
562 go func() {
564 - errc <- corehttp.Serve(node, gwLis.NetListener(), opts...)
563 + errc <- corehttp.Serve(node, manet.NetListener(gwLis), opts...)
564 close(errc)
565 }()
566 return errc, nil
cmd/seccat/seccat.go
+8 -5
@@ -152,17 +152,20 @@ func connect(args args) error {
152 }
153
154 // log everything that goes through conn
155 - rwc := &logRW{n: "conn", rw: conn}
155 + rwc := &logConn{n: "conn", Conn: conn}
156
157 // OK, let's setup the channel.
158 sk := ps.PrivKey(p)
159 - sg := secio.SessionGenerator{LocalID: p, PrivateKey: sk}
160 - sess, err := sg.NewSession(context.TODO(), rwc)
159 + sg, err := secio.New(sk)
160 if err != nil {
161 return err
162 }
164 - out("remote peer id: %s", sess.RemotePeer())
165 - netcat(sess.ReadWriter().(io.ReadWriteCloser))
163 + sconn, err := sg.SecureInbound(context.TODO(), rwc)
164 + if err != nil {
165 + return err
166 + }
167 + out("remote peer id: %s", sconn.RemotePeer())
168 + netcat(sconn)
169 return nil
170 }
171
cmd/seccat/util.go
+10 -14
@@ -2,7 +2,7 @@ package main
2
3 import (
4 "fmt"
5 - "io"
5 + "net"
6 "os"
7
8 logging "gx/ipfs/QmTG23dvpBCBjqQwyDxV8CQT6jmS4PSftNr1VqHhE3MLy7/go-log"
@@ -24,28 +24,24 @@ func out(format string, vals ...interface{}) {
24 }
25 }
26
27 -type logRW struct {
28 - n string
29 - rw io.ReadWriter
27 +type logConn struct {
28 + net.Conn
29 + n string
30 }
31
32 -func (r *logRW) Read(buf []byte) (int, error) {
33 - n, err := r.rw.Read(buf)
32 +func (r *logConn) Read(buf []byte) (int, error) {
33 + n, err := r.Conn.Read(buf)
34 if n > 0 {
35 log.Debugf("%s read: %v", r.n, buf)
36 }
37 return n, err
38 }
39
40 -func (r *logRW) Write(buf []byte) (int, error) {
40 +func (r *logConn) Write(buf []byte) (int, error) {
41 log.Debugf("%s write: %v", r.n, buf)
42 - return r.rw.Write(buf)
42 + return r.Conn.Write(buf)
43 }
44
45 -func (r *logRW) Close() error {
46 - c, ok := r.rw.(io.Closer)
47 - if ok {
48 - return c.Close()
49 - }
50 - return nil
45 +func (r *logConn) Close() error {
46 + return r.Conn.Close()
47 }
core/builder.go
+18 -1
@@ -21,7 +21,9 @@ import (
21
22 metrics "gx/ipfs/QmRg1gKTHzc3CZXSKzem8aR4E3TubFhbgXwfVuWnSK5CC5/go-metrics-interface"
23 goprocessctx "gx/ipfs/QmSF8fPo3jgVBAy8fpdjjYqgG87dkJgUprRBHRd2tmfgpP/goprocess/context"
24 + libp2p "gx/ipfs/QmY6iAoG9DVgZwh5ZRcQEpa2uErAe1Hbei8qXPCjpDS9Ge/go-libp2p"
25 offline "gx/ipfs/QmYk9mQ4iByLLFzZPGWMnjJof3DQ3QneFFR6ZtNAXd8UvS/go-ipfs-exchange-offline"
26 + p2phost "gx/ipfs/QmaSfSMvc1VPZ8JbMponFs4WHvF9FgEruF56opm5E1RgQA/go-libp2p-host"
27 bstore "gx/ipfs/QmayRSLCiM2gWR7Kay8vqu3Yy5mf7yPqocF9ZRgDUPYMcc/go-ipfs-blockstore"
28 peer "gx/ipfs/QmcJukH2sAFjY3HdBKq35WDzWoL3UUu2gt9wdfqZTUyM74/go-libp2p-peer"
29 pstore "gx/ipfs/QmdeiKhUy1TVGBaKxt7y1QmBDLBdisSrLJ1x58Eoj4PXUh/go-libp2p-peerstore"
@@ -42,6 +44,10 @@ type BuildCfg struct {
44 // that will improve performance in long run
45 Permanent bool
46
47 + // DisableEncryptedConnections disables connection encryption *entirely*.
48 + // DO NOT SET THIS UNLESS YOU'RE TESTING.
49 + DisableEncryptedConnections bool
50 +
51 // If NilRepo is set, a repo backed by a nil datastore will be constructed
52 NilRepo bool
53
@@ -126,6 +132,7 @@ func NewNode(ctx context.Context, cfg *BuildCfg) (*IpfsNode, error) {
132 if err != nil {
133 return nil, err
134 }
135 +
136 ctx = metrics.CtxScope(ctx, "ipfs")
137
138 n := &IpfsNode{
@@ -214,9 +221,19 @@ func setupNode(ctx context.Context, n *IpfsNode, cfg *BuildCfg) error {
221 bs.HashOnRead(true)
222 }
223
224 + hostOption := cfg.Host
225 + if cfg.DisableEncryptedConnections {
226 + innerHostOption := hostOption
227 + hostOption = func(ctx context.Context, id peer.ID, ps pstore.Peerstore, options ...libp2p.Option) (p2phost.Host, error) {
228 + return innerHostOption(ctx, id, ps, append(options, libp2p.NoSecurity)...)
229 + }
230 + log.Warningf(`Your IPFS node has been configured to run WITHOUT ENCRYPTED CONNECTIONS.
231 + You will not be able to connect to any nodes configured to use encrypted connections`)
232 + }
233 +
234 if cfg.Online {
235 do := setupDiscoveryOption(rcfg.Discovery)
219 - if err := n.startOnlineServices(ctx, cfg.Routing, cfg.Host, do, cfg.getOpt("pubsub"), cfg.getOpt("ipnsps"), cfg.getOpt("mplex")); err != nil {
236 + if err := n.startOnlineServices(ctx, cfg.Routing, hostOption, do, cfg.getOpt("pubsub"), cfg.getOpt("ipnsps"), cfg.getOpt("mplex")); err != nil {
237 return err
238 }
239 } else {
core/commands/swarm.go
+20 -20
@@ -90,10 +90,13 @@ var swarmPeersCmd = &cmds.Command{
90 Peer: pid.Pretty(),
91 }
92
93 - swcon, ok := c.(*swarm.Conn)
94 - if ok {
95 - ci.Muxer = fmt.Sprintf("%T", swcon.StreamConn().Conn())
96 - }
93 + /*
94 + // FIXME(steb):
95 + swcon, ok := c.(*swarm.Conn)
96 + if ok {
97 + ci.Muxer = fmt.Sprintf("%T", swcon.StreamConn().Conn())
98 + }
99 + */
100
101 if verbose || latency {
102 lat := n.Peerstore.LatencyEWMA(pid)
@@ -104,11 +107,7 @@ var swarmPeersCmd = &cmds.Command{
107 }
108 }
109 if verbose || streams {
107 - strs, err := c.GetStreams()
108 - if err != nil {
109 - res.SetError(err, cmdkit.ErrNormal)
110 - return
111 - }
110 + strs := c.GetStreams()
111
112 for _, s := range strs {
113 ci.Streams = append(ci.Streams, streamInfo{Protocol: string(s.Protocol())})
@@ -384,14 +383,13 @@ ipfs swarm connect /ip4/104.131.131.82/tcp/4001/ipfs/QmaCpDMGvV2BGHeYERUEnRQAwe3
383 return
384 }
385
387 - snet, ok := n.PeerHost.Network().(*swarm.Network)
386 + // FIXME(steb): Nasty
387 + swrm, ok := n.PeerHost.Network().(*swarm.Swarm)
388 if !ok {
389 res.SetError(fmt.Errorf("peerhost network was not swarm"), cmdkit.ErrNormal)
390 return
391 }
392
393 - swrm := snet.Swarm()
394 -
393 pis, err := peersWithAddresses(addrs)
394 if err != nil {
395 res.SetError(err, cmdkit.ErrNormal)
@@ -574,14 +572,15 @@ Filters default to those specified under the "Swarm.AddrFilters" config key.
572 return
573 }
574
577 - snet, ok := n.PeerHost.Network().(*swarm.Network)
575 + // FIXME(steb)
576 + swrm, ok := n.PeerHost.Network().(*swarm.Swarm)
577 if !ok {
578 res.SetError(errors.New("failed to cast network to swarm network"), cmdkit.ErrNormal)
579 return
580 }
581
582 var output []string
584 - for _, f := range snet.Filters.Filters() {
583 + for _, f := range swrm.Filters.Filters() {
584 s, err := mafilter.ConvertIPNet(f)
585 if err != nil {
586 res.SetError(err, cmdkit.ErrNormal)
@@ -621,7 +620,8 @@ add your filters to the ipfs config file.
620 return
621 }
622
624 - snet, ok := n.PeerHost.Network().(*swarm.Network)
623 + // FIXME(steb)
624 + swrm, ok := n.PeerHost.Network().(*swarm.Swarm)
625 if !ok {
626 res.SetError(errors.New("failed to cast network to swarm network"), cmdkit.ErrNormal)
627 return
@@ -651,7 +651,7 @@ add your filters to the ipfs config file.
651 return
652 }
653
654 - snet.Filters.AddDialFilter(mask)
654 + swrm.Filters.AddDialFilter(mask)
655 }
656
657 added, err := filtersAdd(r, cfg, req.Arguments())
@@ -693,7 +693,7 @@ remove your filters from the ipfs config file.
693 return
694 }
695
696 - snet, ok := n.PeerHost.Network().(*swarm.Network)
696 + swrm, ok := n.PeerHost.Network().(*swarm.Swarm)
697 if !ok {
698 res.SetError(errors.New("failed to cast network to swarm network"), cmdkit.ErrNormal)
699 return
@@ -712,9 +712,9 @@ remove your filters from the ipfs config file.
712 }
713
714 if req.Arguments()[0] == "all" || req.Arguments()[0] == "*" {
715 - fs := snet.Filters.Filters()
715 + fs := swrm.Filters.Filters()
716 for _, f := range fs {
717 - snet.Filters.Remove(f)
717 + swrm.Filters.Remove(f)
718 }
719
720 removed, err := filtersRemoveAll(r, cfg)
@@ -735,7 +735,7 @@ remove your filters from the ipfs config file.
735 return
736 }
737
738 - snet.Filters.Remove(mask)
738 + swrm.Filters.Remove(mask)
739 }
740
741 removed, err := filtersRemove(r, cfg, req.Arguments())
core/core.go
+54 -87
@@ -16,7 +16,6 @@ import (
16 "fmt"
17 "io"
18 "io/ioutil"
19 - "net"
19 "os"
20 "strings"
21 "time"
@@ -41,18 +40,16 @@ import (
40 u "gx/ipfs/QmNiJuT8Ja3hMVpBHXv3Q6dwmperaQ6JjLtpMQgMCD7xvx/go-ipfs-util"
41 circuit "gx/ipfs/QmR5sXZi68rm9m2E3KiXj6hE5m3GeLaDjbLPUeV6W3MLR8/go-libp2p-circuit"
42 floodsub "gx/ipfs/QmRMgHdiLHJvySrXbtLBehr1W1yTQyuNmZG8HghG54ZPDz/go-libp2p-floodsub"
44 - swarm "gx/ipfs/QmRpKdg1xs4Yyrn9yrVYRBp7AQqyRxMLpD6Jgp1eZAGqEr/go-libp2p-swarm"
43 goprocess "gx/ipfs/QmSF8fPo3jgVBAy8fpdjjYqgG87dkJgUprRBHRd2tmfgpP/goprocess"
44 pnet "gx/ipfs/QmSGoP33Ufev1UDsUuHco8rfhVTzxfq6smXhwhN16c5CWd/go-libp2p-pnet"
45 mamask "gx/ipfs/QmSMZwvs3n4GBikZ7hKzT17c3bk65FmyZo2JqtJ16swqCv/multiaddr-filter"
46 logging "gx/ipfs/QmTG23dvpBCBjqQwyDxV8CQT6jmS4PSftNr1VqHhE3MLy7/go-log"
49 - addrutil "gx/ipfs/QmTGSre9j1otFgsr1opCUQDXTPSM6BTZnMWwPeA5nYJM7w/go-addr-util"
47 record "gx/ipfs/QmTUyK82BVPA6LmSzEJpfEunk9uBaQzWtMsNP917tVj4sT/go-libp2p-record"
48 routing "gx/ipfs/QmUHRKTeaoASDvDj7cTAXsmjAY7KQ13ErtzkQHZQq6uFUz/go-libp2p-routing"
49 metrics "gx/ipfs/QmVvu4bS5QLfS19ePkp5Wgzn2ZUma5oXTT9BgDFyQLxUZF/go-libp2p-metrics"
50 psrouter "gx/ipfs/QmWKLW1C2jmGAEzX8jNpCTii6n2ScGxytnoRMRdNTK5Knt/go-libp2p-pubsub-router"
51 ma "gx/ipfs/QmWWQ2Txc2c6tqjsBpzg5Ar652cHPGNsQQp2SejkNmkUMb/go-multiaddr"
55 - mssmux "gx/ipfs/QmWzjXAyBTygw6CeCTUnhJzhFucfxY5FJivSoiGuiSbPjS/go-smux-multistream"
52 + libp2p "gx/ipfs/QmY6iAoG9DVgZwh5ZRcQEpa2uErAe1Hbei8qXPCjpDS9Ge/go-libp2p"
53 discovery "gx/ipfs/QmY6iAoG9DVgZwh5ZRcQEpa2uErAe1Hbei8qXPCjpDS9Ge/go-libp2p/p2p/discovery"
54 p2pbhost "gx/ipfs/QmY6iAoG9DVgZwh5ZRcQEpa2uErAe1Hbei8qXPCjpDS9Ge/go-libp2p/p2p/host/basic"
55 rhost "gx/ipfs/QmY6iAoG9DVgZwh5ZRcQEpa2uErAe1Hbei8qXPCjpDS9Ge/go-libp2p/p2p/host/routed"
@@ -68,7 +65,6 @@ import (
65 peer "gx/ipfs/QmcJukH2sAFjY3HdBKq35WDzWoL3UUu2gt9wdfqZTUyM74/go-libp2p-peer"
66 cid "gx/ipfs/QmcZfnkapfECQGcLZaf9B79NRg7cRa9EnZh4LSbkCzwNvY/go-cid"
67 yamux "gx/ipfs/QmcsgrV3nCAKjiHKZhKVXWc4oY3WBECJCqahXEMpHeMrev/go-smux-yamux"
71 - ipnet "gx/ipfs/Qmd3oYWVLCVWryDV6Pobv6whZcvDXAHqS3chemZ658y4a8/go-libp2p-interface-pnet"
68 exchange "gx/ipfs/QmdcAXgEHUueP4A7b5hjabKn2EooeHgMreMvFC249dGCgc/go-ipfs-exchange-interface"
69 pstore "gx/ipfs/QmdeiKhUy1TVGBaKxt7y1QmBDLBdisSrLJ1x58Eoj4PXUh/go-libp2p-peerstore"
70 ic "gx/ipfs/Qme1knMqwt1hKZbc1BmQFmnm9f36nyQGwXxPGVpVJ9rMK5/go-libp2p-crypto"
@@ -173,30 +169,29 @@ func (n *IpfsNode) startOnlineServices(ctx context.Context, routingOption Routin
169 if err != nil {
170 return err
171 }
176 - var addrfilter []*net.IPNet
172 +
173 + var libp2pOpts []libp2p.Option
174 for _, s := range cfg.Swarm.AddrFilters {
175 f, err := mamask.NewMask(s)
176 if err != nil {
177 return fmt.Errorf("incorrectly formatted address filter in config: %s", s)
178 }
182 - addrfilter = append(addrfilter, f)
179 + libp2pOpts = append(libp2pOpts, libp2p.FilterAddresses(f))
180 }
181
182 if !cfg.Swarm.DisableBandwidthMetrics {
183 // Set reporter
184 n.Reporter = metrics.NewBandwidthCounter()
185 + libp2pOpts = append(libp2pOpts, libp2p.BandwidthReporter(n.Reporter))
186 }
187
190 - tpt := makeSmuxTransport(mplex)
191 -
188 swarmkey, err := n.Repo.SwarmKey()
189 if err != nil {
190 return err
191 }
192
197 - var protec ipnet.Protector
193 if swarmkey != nil {
199 - protec, err = pnet.NewProtector(bytes.NewReader(swarmkey))
194 + protec, err := pnet.NewProtector(bytes.NewReader(swarmkey))
195 if err != nil {
196 return fmt.Errorf("failed to configure private network: %s", err)
197 }
@@ -219,27 +214,39 @@ func (n *IpfsNode) startOnlineServices(ctx context.Context, routingOption Routin
214 }
215 }
216 }()
217 +
218 + libp2pOpts = append(libp2pOpts, libp2p.PrivateNetwork(protec))
219 }
220
221 addrsFactory, err := makeAddrsFactory(cfg.Addresses)
222 if err != nil {
223 return err
224 }
225 + if !cfg.Swarm.DisableRelay {
226 + addrsFactory = composeAddrsFactory(addrsFactory, filterRelayAddrs)
227 + }
228 + libp2pOpts = append(libp2pOpts, libp2p.AddrsFactory(addrsFactory))
229
229 - connmgr, err := constructConnMgr(cfg.Swarm.ConnMgr)
230 + connm, err := constructConnMgr(cfg.Swarm.ConnMgr)
231 if err != nil {
232 return err
233 }
234 + libp2pOpts = append(libp2pOpts, libp2p.ConnectionManager(connm))
235 +
236 + libp2pOpts = append(libp2pOpts, makeSmuxTransportOption(mplex))
237
234 - hostopts := &ConstructPeerHostOpts{
235 - AddrsFactory: addrsFactory,
236 - DisableNatPortMap: cfg.Swarm.DisableNatPortMap,
237 - DisableRelay: cfg.Swarm.DisableRelay,
238 - EnableRelayHop: cfg.Swarm.EnableRelayHop,
239 - ConnectionManager: connmgr,
238 + if !cfg.Swarm.DisableNatPortMap {
239 + libp2pOpts = append(libp2pOpts, libp2p.NATPortMap())
240 }
241 - peerhost, err := hostOption(ctx, n.Identity, n.Peerstore, n.Reporter,
242 - addrfilter, tpt, protec, hostopts)
241 + if !cfg.Swarm.DisableRelay {
242 + var opts []circuit.RelayOpt
243 + if cfg.Swarm.EnableRelayHop {
244 + opts = append(opts, circuit.OptHop)
245 + }
246 + libp2pOpts = append(libp2pOpts, libp2p.EnableRelay(opts...))
247 + }
248 +
249 + peerhost, err := hostOption(ctx, n.Identity, n.Peerstore, libp2pOpts...)
250
251 if err != nil {
252 return err
@@ -372,8 +379,9 @@ func makeAddrsFactory(cfg config.Addresses) (p2pbhost.AddrsFactory, error) {
379 }, nil
380 }
381
375 -func makeSmuxTransport(mplexExp bool) smux.Transport {
376 - mstpt := mssmux.NewBlankTransport()
382 +func makeSmuxTransportOption(mplexExp bool) libp2p.Option {
383 + const yamuxID = "/yamux/1.0.0"
384 + const mplexID = "/mplex/6.7.0"
385
386 ymxtpt := &yamux.Transport{
387 AcceptBacklog: 512,
@@ -388,18 +396,29 @@ func makeSmuxTransport(mplexExp bool) smux.Transport {
396 ymxtpt.LogOutput = os.Stderr
397 }
398
391 - mstpt.AddTransport("/yamux/1.0.0", ymxtpt)
392 -
399 + muxers := map[string]smux.Transport{yamuxID: ymxtpt}
400 if mplexExp {
394 - mstpt.AddTransport("/mplex/6.7.0", mplex.DefaultTransport)
401 + muxers[mplexID] = mplex.DefaultTransport
402 }
403
404 // Allow muxer preference order overriding
405 + order := []string{yamuxID, mplexID}
406 if prefs := os.Getenv("LIBP2P_MUX_PREFS"); prefs != "" {
399 - mstpt.OrderPreference = strings.Fields(prefs)
407 + order = strings.Fields(prefs)
408 + }
409 +
410 + opts := make([]libp2p.Option, 0, len(order))
411 + for _, id := range order {
412 + tpt, ok := muxers[id]
413 + if !ok {
414 + log.Warning("unknown or duplicate muxer in LIBP2P_MUX_PREFS: %s", id)
415 + continue
416 + }
417 + delete(muxers, id)
418 + opts = append(opts, libp2p.Muxer(id, tpt))
419 }
420
402 - return mstpt
421 + return libp2p.ChainOptions(opts...)
422 }
423
424 func setupDiscoveryOption(d config.Discovery) DiscoveryOption {
@@ -853,62 +872,18 @@ type ConstructPeerHostOpts struct {
872 ConnectionManager ifconnmgr.ConnManager
873 }
874
856 -type HostOption func(ctx context.Context, id peer.ID, ps pstore.Peerstore, bwr metrics.Reporter, fs []*net.IPNet, tpt smux.Transport, protc ipnet.Protector, opts *ConstructPeerHostOpts) (p2phost.Host, error)
875 +type HostOption func(ctx context.Context, id peer.ID, ps pstore.Peerstore, options ...libp2p.Option) (p2phost.Host, error)
876
877 var DefaultHostOption HostOption = constructPeerHost
878
879 // isolates the complex initialization steps
861 -func constructPeerHost(ctx context.Context, id peer.ID, ps pstore.Peerstore, bwr metrics.Reporter, fs []*net.IPNet, tpt smux.Transport, protec ipnet.Protector, opts *ConstructPeerHostOpts) (p2phost.Host, error) {
862 -
863 - // no addresses to begin with. we'll start later.
864 - swrm, err := swarm.NewSwarmWithProtector(ctx, nil, id, ps, protec, tpt, bwr)
865 - if err != nil {
866 - return nil, err
867 - }
868 -
869 - network := (*swarm.Network)(swrm)
870 -
871 - for _, f := range fs {
872 - network.Swarm().Filters.AddDialFilter(f)
880 +func constructPeerHost(ctx context.Context, id peer.ID, ps pstore.Peerstore, options ...libp2p.Option) (p2phost.Host, error) {
881 + pkey := ps.PrivKey(id)
882 + if pkey == nil {
883 + return nil, fmt.Errorf("missing private key for node ID: %s", id.Pretty())
884 }
874 -
875 - hostOpts := []interface{}{bwr}
876 - if !opts.DisableNatPortMap {
877 - hostOpts = append(hostOpts, p2pbhost.NATPortMap)
878 - }
879 - if opts.ConnectionManager != nil {
880 - hostOpts = append(hostOpts, opts.ConnectionManager)
881 - }
882 -
883 - addrsFactory := opts.AddrsFactory
884 - if !opts.DisableRelay {
885 - if addrsFactory != nil {
886 - addrsFactory = composeAddrsFactory(addrsFactory, filterRelayAddrs)
887 - } else {
888 - addrsFactory = filterRelayAddrs
889 - }
890 - }
891 -
892 - if addrsFactory != nil {
893 - hostOpts = append(hostOpts, addrsFactory)
894 - }
895 -
896 - host := p2pbhost.New(network, hostOpts...)
897 -
898 - if !opts.DisableRelay {
899 - var relayOpts []circuit.RelayOpt
900 - if opts.EnableRelayHop {
901 - relayOpts = append(relayOpts, circuit.OptHop)
902 - }
903 -
904 - err := circuit.AddRelayTransport(ctx, host, relayOpts...)
905 - if err != nil {
906 - host.Close()
907 - return nil, err
908 - }
909 - }
910 -
911 - return host, nil
885 + options = append([]libp2p.Option{libp2p.Identity(pkey), libp2p.Peerstore(ps)}, options...)
886 + return libp2p.New(ctx, options...)
887 }
888
889 func filterRelayAddrs(addrs []ma.Multiaddr) []ma.Multiaddr {
@@ -936,16 +911,8 @@ func startListening(host p2phost.Host, cfg *config.Config) error {
911 return err
912 }
913
939 - // make sure we error out if our config does not have addresses we can use
940 - log.Debugf("Config.Addresses.Swarm:%s", listenAddrs)
941 - filteredAddrs := addrutil.FilterUsableAddrs(listenAddrs)
942 - log.Debugf("Config.Addresses.Swarm:%s (filtered)", filteredAddrs)
943 - if len(filteredAddrs) < 1 {
944 - return fmt.Errorf("addresses in config not usable: %s", listenAddrs)
945 - }
946 -
914 // Actually start listening:
948 - if err := host.Network().Listen(filteredAddrs...); err != nil {
915 + if err := host.Network().Listen(listenAddrs...); err != nil {
916 return err
917 }
918
core/corehttp/corehttp.go
+1 -1
@@ -61,7 +61,7 @@ func ListenAndServe(n *core.IpfsNode, listeningMultiAddr string, options ...Serv
61 addr = list.Multiaddr()
62 fmt.Printf("API server listening on %s\n", addr)
63
64 - return Serve(n, list.NetListener(), options...)
64 + return Serve(n, manet.NetListener(list), options...)
65 }
66
67 func Serve(node *core.IpfsNode, lis net.Listener, options ...ServeOption) error {
core/corehttp/metrics_test.go
+3 -3
@@ -7,9 +7,9 @@ import (
7
8 core "github.com/ipfs/go-ipfs/core"
9
10 + swarmt "gx/ipfs/QmRpKdg1xs4Yyrn9yrVYRBp7AQqyRxMLpD6Jgp1eZAGqEr/go-libp2p-swarm/testing"
11 inet "gx/ipfs/QmXoz9o2PT3tEzf7hicegwex5UgVP54n3k82K7jrWFyN86/go-libp2p-net"
12 bhost "gx/ipfs/QmY6iAoG9DVgZwh5ZRcQEpa2uErAe1Hbei8qXPCjpDS9Ge/go-libp2p/p2p/host/basic"
12 - testutil "gx/ipfs/Qma2UuHusnaFV24DgeZ5hyrM9uc4UdyVaZbtn2FQsPRhES/go-libp2p-netutil"
13 )
14
15 // This test is based on go-libp2p/p2p/net/swarm.TestConnectednessCorrect
@@ -20,11 +20,11 @@ func TestPeersTotal(t *testing.T) {
20
21 hosts := make([]*bhost.BasicHost, 4)
22 for i := 0; i < 4; i++ {
23 - hosts[i] = bhost.New(testutil.GenSwarmNetwork(t, ctx))
23 + hosts[i] = bhost.New(swarmt.GenSwarm(t, ctx))
24 }
25
26 dial := func(a, b inet.Network) {
27 - testutil.DivulgeAddresses(b, a)
27 + swarmt.DivulgeAddresses(b, a)
28 if _, err := a.DialPeer(ctx, b.LocalPeer()); err != nil {
29 t.Fatalf("Failed to dial: %s", err)
30 }
core/mock/mock.go
+2 -5
@@ -2,7 +2,6 @@ package coremock
2
3 import (
4 "context"
5 - "net"
5
6 commands "github.com/ipfs/go-ipfs/commands"
7 core "github.com/ipfs/go-ipfs/core"
@@ -10,12 +9,10 @@ import (
9 config "github.com/ipfs/go-ipfs/repo/config"
10
11 testutil "gx/ipfs/QmUJzxQQ2kzwQubsMqBTr1NGDpLfh7pGA2E1oaJULcKDPq/go-testutil"
13 - metrics "gx/ipfs/QmVvu4bS5QLfS19ePkp5Wgzn2ZUma5oXTT9BgDFyQLxUZF/go-libp2p-metrics"
12 + libp2p "gx/ipfs/QmY6iAoG9DVgZwh5ZRcQEpa2uErAe1Hbei8qXPCjpDS9Ge/go-libp2p"
13 mocknet "gx/ipfs/QmY6iAoG9DVgZwh5ZRcQEpa2uErAe1Hbei8qXPCjpDS9Ge/go-libp2p/p2p/net/mock"
15 - smux "gx/ipfs/QmY9JXR3FupnYAYJWK9aMr9bCpqWKcToQ1tz8DVGTrHpHw/go-stream-muxer"
14 host "gx/ipfs/QmaSfSMvc1VPZ8JbMponFs4WHvF9FgEruF56opm5E1RgQA/go-libp2p-host"
15 peer "gx/ipfs/QmcJukH2sAFjY3HdBKq35WDzWoL3UUu2gt9wdfqZTUyM74/go-libp2p-peer"
18 - ipnet "gx/ipfs/Qmd3oYWVLCVWryDV6Pobv6whZcvDXAHqS3chemZ658y4a8/go-libp2p-interface-pnet"
16 pstore "gx/ipfs/QmdeiKhUy1TVGBaKxt7y1QmBDLBdisSrLJ1x58Eoj4PXUh/go-libp2p-peerstore"
17 datastore "gx/ipfs/QmeiCcJfDW1GJnWUArudsv5rQsihpi4oyddPhdqo3CfX6i/go-datastore"
18 syncds "gx/ipfs/QmeiCcJfDW1GJnWUArudsv5rQsihpi4oyddPhdqo3CfX6i/go-datastore/sync"
@@ -33,7 +30,7 @@ func NewMockNode() (*core.IpfsNode, error) {
30 }
31
32 func MockHostOption(mn mocknet.Mocknet) core.HostOption {
36 - return func(ctx context.Context, id peer.ID, ps pstore.Peerstore, bwr metrics.Reporter, fs []*net.IPNet, _ smux.Transport, _ ipnet.Protector, _ *core.ConstructPeerHostOpts) (host.Host, error) {
33 + return func(ctx context.Context, id peer.ID, ps pstore.Peerstore, _ ...libp2p.Option) (host.Host, error) {
34 return mn.AddPeerWithPeerstore(id, ps)
35 }
36 }
exchange/bitswap/network/ipfs_impl.go
+6 -6
@@ -54,7 +54,7 @@ type streamMessageSender struct {
54 }
55
56 func (s *streamMessageSender) Close() error {
57 - return s.s.Close()
57 + return inet.FullClose(s.s)
58 }
59
60 func (s *streamMessageSender) Reset() error {
@@ -119,13 +119,13 @@ func (bsnet *impl) SendMessage(
119 return err
120 }
121
122 - err = msgToStream(ctx, s, outgoing)
123 - if err != nil {
122 + if err = msgToStream(ctx, s, outgoing); err != nil {
123 s.Reset()
125 - } else {
126 - s.Close()
124 + return err
125 }
128 - return err
126 + // Yes, return this error. We have no reason to believe that the block
127 + // was actually *sent* unless we see the EOF.
128 + return inet.FullClose(s)
129 }
130
131 func (bsnet *impl) SetDelegate(r Receiver) {
package.json
-6
@@ -311,12 +311,6 @@
311 "name": "go-multibase",
312 "version": "0.2.6"
313 },
314 - {
315 - "author": "whyrusleeping",
316 - "hash": "QmYDNqBAMWVMHKndYR35Sd8PfEVWBiDmpHYkuRJTunJDeJ",
317 - "name": "go-libp2p-interface-conn",
318 - "version": "0.4.13"
319 - },
314 {
315 "author": "multiformats",
316 "hash": "QmWWQ2Txc2c6tqjsBpzg5Ar652cHPGNsQQp2SejkNmkUMb",